refactor(runtime): centralize async bridges under an owned runtime (#376)
* feat(runtime): add application-scoped async runtime * refactor(cli): use owned runtime for session stats * refactor(onboard): use the owned async runtime * docs(runtime): record async bridge ownership * refactor(middleware): keep sync fallback synchronous * refactor(mcp): load tools on an owned runtime * refactor(cli): share owned runtime across entry points * refactor(channels): make inbound sync bridge explicit * refactor(stream): run Rich streaming on owned runtime * chore(runtime): remove nest-asyncio dependency * refactor(asyncio): require active loops in async code * docs(runtime): document final event loop ownership * fix(stream): cancel stalled owned streams * fix(cli): recover cleanly from stream cancellation * fix(runtime): drain executor work before shutdown * fix(runtime): terminate cancelled shell process trees * fix(models): let fallback bypass selector failures * fix(cli): reset interrupt handling between turns * docs: rm implementation spec * fix(serve): cancel active turns during shutdown * fix(runtime): protect settlement from waiter cancellation * fix(backends): reject empty shell commands * fix(runtime): terminate descendants after shell exit * fix(mcp): keep standalone discovery off channel loop * fix(cli): own and settle interactive prompt cancellation * fix(serve): keep channel sends off runtime loop * fix(stream): scope cancel context to iterator steps * refactor(serve): require the owned async runtime * fix(channels): keep interactive sends off runtime loop * fix(selector): surface fallback without log spam * test(runtime): normalize Windows shell marker * fix(cli): serialize interactive session turns * fix(shell): bound output drain after termination * fix(ui): do not retry owned runtime failures * fix(shell): allow signal-safe registry reentry * fix(shell): avoid terminating reused process ids * fix(channels): preserve streaming send order * fix(cli): report runtime shutdown timeouts cleanly * fix(mcp): guide async callers to async loader * docs(runtime): clarify reserved async bridge APIs * fix(runtime): bound code interpreter cleanup * test(shell): use active Python for drain regression --------- Co-authored-by: Xi Zhang <106144707+X-iZhang@users.noreply.github.com>
This commit is contained in:
+126
-302
@@ -1,325 +1,149 @@
|
||||
"""Tests for event loop management in streaming display."""
|
||||
"""Owned-runtime tests for the synchronous Rich streaming adapter."""
|
||||
|
||||
import asyncio
|
||||
import threading
|
||||
from unittest.mock import Mock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
from EvoScientist.stream.display import _create_event_loop, _get_event_loop
|
||||
from EvoScientist.runtime import AsyncRuntime, AsyncRuntimeError
|
||||
from EvoScientist.stream.display import _run_streaming
|
||||
from tests.fakes import FakeGraphGateway
|
||||
|
||||
|
||||
class _TrackingEventLoopPolicy(asyncio.DefaultEventLoopPolicy):
|
||||
"""Event loop policy that records loops created by one test."""
|
||||
def _text_stream(loops, response="test response"):
|
||||
async def _stream(_request):
|
||||
loops.append(asyncio.get_running_loop())
|
||||
yield {"type": "text", "content": response}
|
||||
yield {"type": "done", "response": response}
|
||||
|
||||
def __init__(self):
|
||||
super().__init__()
|
||||
self.created_loops: list[asyncio.AbstractEventLoop] = []
|
||||
|
||||
def new_event_loop(self) -> asyncio.AbstractEventLoop:
|
||||
loop = super().new_event_loop()
|
||||
self.created_loops.append(loop)
|
||||
return loop
|
||||
return _stream
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def isolated_event_loop_policy():
|
||||
previous_policy = asyncio.get_event_loop_policy()
|
||||
test_policy = _TrackingEventLoopPolicy()
|
||||
asyncio.set_event_loop_policy(test_policy)
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
try:
|
||||
for loop in test_policy.created_loops:
|
||||
if not loop.is_closed():
|
||||
loop.close()
|
||||
finally:
|
||||
asyncio.set_event_loop_policy(previous_policy)
|
||||
def test_sequential_streams_reuse_application_runtime_loop():
|
||||
loops: list[asyncio.AbstractEventLoop] = []
|
||||
gateway = FakeGraphGateway(stream=_text_stream(loops))
|
||||
|
||||
|
||||
class TestCreateEventLoop:
|
||||
"""Tests for _create_event_loop helper."""
|
||||
|
||||
def test_creates_new_loop(self):
|
||||
"""Should create a new event loop and set it as current."""
|
||||
# Get initial loop (if any)
|
||||
try:
|
||||
initial_loop = asyncio.get_event_loop()
|
||||
initial_loop.close()
|
||||
except RuntimeError:
|
||||
pass
|
||||
|
||||
# Create new loop
|
||||
loop = _create_event_loop()
|
||||
|
||||
assert loop is not None
|
||||
assert not loop.is_closed()
|
||||
assert asyncio.get_event_loop() is loop
|
||||
|
||||
# Cleanup
|
||||
loop.close()
|
||||
|
||||
def test_replaces_closed_loop(self):
|
||||
"""Should replace a closed loop."""
|
||||
old_loop = asyncio.new_event_loop()
|
||||
asyncio.set_event_loop(old_loop)
|
||||
old_loop.close()
|
||||
|
||||
new_loop = _create_event_loop()
|
||||
|
||||
assert new_loop is not old_loop
|
||||
assert not new_loop.is_closed()
|
||||
assert asyncio.get_event_loop() is new_loop
|
||||
|
||||
# Cleanup
|
||||
new_loop.close()
|
||||
|
||||
|
||||
class TestGetEventLoop:
|
||||
"""Tests for _get_event_loop helper."""
|
||||
|
||||
def test_returns_existing_open_loop(self):
|
||||
"""Should return existing event loop if it's open."""
|
||||
loop = asyncio.new_event_loop()
|
||||
asyncio.set_event_loop(loop)
|
||||
|
||||
result = _get_event_loop()
|
||||
|
||||
assert result is loop
|
||||
assert not result.is_closed()
|
||||
|
||||
# Cleanup
|
||||
loop.close()
|
||||
|
||||
def test_creates_new_loop_when_closed(self):
|
||||
"""Should create new event loop if current one is closed."""
|
||||
old_loop = asyncio.new_event_loop()
|
||||
asyncio.set_event_loop(old_loop)
|
||||
old_loop.close()
|
||||
|
||||
result = _get_event_loop()
|
||||
|
||||
assert result is not old_loop
|
||||
assert not result.is_closed()
|
||||
|
||||
# Cleanup
|
||||
result.close()
|
||||
|
||||
def test_handles_no_event_loop(self):
|
||||
"""Should handle RuntimeError when no event loop exists (edge case)."""
|
||||
# This test simulates what happens in a worker thread
|
||||
# In practice, get_event_loop() returns a closed loop, not RuntimeError
|
||||
# But we handle the RuntimeError case defensively
|
||||
loop = _get_event_loop()
|
||||
assert loop is not None
|
||||
assert not loop.is_closed()
|
||||
|
||||
# Cleanup
|
||||
loop.close()
|
||||
|
||||
|
||||
class TestMultipleStreamingCalls:
|
||||
"""Tests for the main bug fix: multiple _run_streaming calls."""
|
||||
|
||||
def test_sequential_streaming_calls(self):
|
||||
"""Multiple sequential calls should work without 'Event loop is closed' error."""
|
||||
from EvoScientist.stream.display import _run_streaming
|
||||
|
||||
# Mock agent that returns simple events
|
||||
mock_agent = Mock()
|
||||
|
||||
async def mock_stream(_request):
|
||||
"""Mock event stream."""
|
||||
yield {"type": "text", "content": "test response"}
|
||||
yield {"type": "done", "response": "test response"}
|
||||
|
||||
# Clean up any existing event loop to start fresh
|
||||
try:
|
||||
existing_loop = asyncio.get_event_loop()
|
||||
if not existing_loop.is_closed():
|
||||
existing_loop.close()
|
||||
except RuntimeError:
|
||||
pass
|
||||
|
||||
gateway = FakeGraphGateway(stream=mock_stream)
|
||||
|
||||
# Patch Live to avoid terminal output during tests
|
||||
with patch("EvoScientist.stream.display.Live"):
|
||||
# First call
|
||||
_run_streaming(
|
||||
agent=mock_agent,
|
||||
message="test message 1",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
gateway=gateway,
|
||||
)
|
||||
|
||||
# Second call - this would fail with "Event loop is closed" before the fix
|
||||
_run_streaming(
|
||||
agent=mock_agent,
|
||||
message="test message 2",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
gateway=gateway,
|
||||
)
|
||||
|
||||
# Third call for good measure
|
||||
_run_streaming(
|
||||
agent=mock_agent,
|
||||
message="test message 3",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
gateway=gateway,
|
||||
)
|
||||
|
||||
def test_loop_reused_across_calls(self):
|
||||
"""Event loop should be reused across multiple calls."""
|
||||
# Create a fresh loop
|
||||
loop = _create_event_loop()
|
||||
|
||||
# Simulate multiple calls
|
||||
for _ in range(3):
|
||||
current_loop = _get_event_loop()
|
||||
assert not current_loop.is_closed()
|
||||
|
||||
# Run a simple coroutine
|
||||
async def dummy():
|
||||
return "ok"
|
||||
|
||||
result = current_loop.run_until_complete(dummy())
|
||||
assert result == "ok"
|
||||
|
||||
# Loop should still be open
|
||||
assert not loop.is_closed()
|
||||
|
||||
# Cleanup
|
||||
loop.close()
|
||||
|
||||
def test_closed_loop_recovery(self):
|
||||
"""If loop gets closed, next call should create a new one."""
|
||||
# Create and close a loop
|
||||
loop1 = _create_event_loop()
|
||||
loop1.close()
|
||||
|
||||
# Next call should detect closed loop and create new one
|
||||
loop2 = _get_event_loop()
|
||||
|
||||
assert loop2 is not loop1
|
||||
assert not loop2.is_closed()
|
||||
|
||||
# Should be able to use the new loop
|
||||
async def dummy():
|
||||
return "success"
|
||||
|
||||
result = loop2.run_until_complete(dummy())
|
||||
assert result == "success"
|
||||
|
||||
# Cleanup
|
||||
loop2.close()
|
||||
|
||||
def test_recursive_streaming_does_not_resend_same_thinking(self):
|
||||
"""Resumed runs should not replay the original thinking to channels."""
|
||||
from EvoScientist.stream.display import _run_streaming
|
||||
|
||||
mock_agent = Mock()
|
||||
thinking = "Initial plan. " * 20
|
||||
stream_calls = 0
|
||||
|
||||
async def mock_stream(_request):
|
||||
nonlocal stream_calls
|
||||
stream_calls += 1
|
||||
if stream_calls == 1:
|
||||
yield {"type": "thinking", "content": thinking}
|
||||
yield {
|
||||
"type": "ask_user",
|
||||
"interrupt_id": "ask-1",
|
||||
"tool_call_id": "tc-1",
|
||||
"questions": [{"question": "Continue?"}],
|
||||
}
|
||||
return
|
||||
|
||||
yield {"type": "text", "content": "final answer"}
|
||||
yield {"type": "done", "response": "final answer"}
|
||||
|
||||
sent_thinking: list[str] = []
|
||||
|
||||
with patch("EvoScientist.stream.display.Live"):
|
||||
with (
|
||||
AsyncRuntime(thread_name="test-stream-runtime") as runtime,
|
||||
patch("EvoScientist.stream.display.Live"),
|
||||
):
|
||||
for index in range(3):
|
||||
result = _run_streaming(
|
||||
agent=mock_agent,
|
||||
message="test message",
|
||||
agent=Mock(),
|
||||
message=f"message {index}",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
on_thinking=sent_thinking.append,
|
||||
ask_user_prompt_fn=lambda _data: {
|
||||
"answers": ["yes"],
|
||||
"status": "answered",
|
||||
},
|
||||
gateway=FakeGraphGateway(stream=mock_stream),
|
||||
gateway=gateway,
|
||||
runtime=runtime,
|
||||
)
|
||||
assert result == "test response"
|
||||
|
||||
assert result == "final answer"
|
||||
assert sent_thinking == [thinking.rstrip()]
|
||||
|
||||
def test_recursive_streaming_sends_new_thinking_after_resume(self):
|
||||
"""Genuinely new thinking in resumed rounds should be relayed."""
|
||||
from EvoScientist.stream.display import _run_streaming
|
||||
|
||||
mock_agent = Mock()
|
||||
thinking_r1 = "Initial plan. " * 20
|
||||
thinking_r2 = "Revised plan. " * 20
|
||||
stream_calls = 0
|
||||
|
||||
async def mock_stream(_request):
|
||||
nonlocal stream_calls
|
||||
stream_calls += 1
|
||||
if stream_calls == 1:
|
||||
yield {"type": "thinking", "content": thinking_r1}
|
||||
yield {
|
||||
"type": "ask_user",
|
||||
"interrupt_id": "ask-1",
|
||||
"tool_call_id": "tc-1",
|
||||
"questions": [{"question": "Continue?"}],
|
||||
}
|
||||
return
|
||||
|
||||
yield {"type": "thinking", "content": thinking_r2}
|
||||
yield {"type": "text", "content": "final answer"}
|
||||
yield {"type": "done", "response": "final answer"}
|
||||
|
||||
sent_thinking: list[str] = []
|
||||
|
||||
with patch("EvoScientist.stream.display.Live"):
|
||||
result = _run_streaming(
|
||||
agent=mock_agent,
|
||||
message="test message",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
on_thinking=sent_thinking.append,
|
||||
ask_user_prompt_fn=lambda _data: {
|
||||
"answers": ["yes"],
|
||||
"status": "answered",
|
||||
},
|
||||
gateway=FakeGraphGateway(stream=mock_stream),
|
||||
)
|
||||
|
||||
assert result == "final answer"
|
||||
assert sent_thinking == [thinking_r1.rstrip(), thinking_r2.rstrip()]
|
||||
assert len(loops) == 3
|
||||
assert loops[0] is loops[1] is loops[2]
|
||||
|
||||
|
||||
class TestEventLoopThreadSafety:
|
||||
"""Tests for thread safety edge cases."""
|
||||
def test_direct_streaming_call_scopes_and_closes_runtime():
|
||||
execution: dict[str, object] = {}
|
||||
|
||||
def test_main_thread_normal_case(self):
|
||||
"""Normal case in main thread should work."""
|
||||
loop = _get_event_loop()
|
||||
assert loop is not None
|
||||
assert not loop.is_closed()
|
||||
async def stream(_request):
|
||||
execution["thread"] = threading.current_thread().name
|
||||
execution["loop"] = asyncio.get_running_loop()
|
||||
yield {"type": "done", "response": "ok"}
|
||||
|
||||
# Cleanup
|
||||
loop.close()
|
||||
with patch("EvoScientist.stream.display.Live"):
|
||||
result = _run_streaming(
|
||||
agent=Mock(),
|
||||
message="message",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
gateway=FakeGraphGateway(stream=stream),
|
||||
)
|
||||
|
||||
assert result == "ok"
|
||||
assert execution["thread"] == "evosci-stream-runtime"
|
||||
assert isinstance(execution["loop"], asyncio.AbstractEventLoop)
|
||||
assert not any(
|
||||
thread.name == "evosci-stream-runtime" and thread.is_alive()
|
||||
for thread in threading.enumerate()
|
||||
)
|
||||
|
||||
|
||||
async def test_async_caller_must_offload_synchronous_renderer():
|
||||
gateway = FakeGraphGateway(stream=_text_stream([]))
|
||||
|
||||
with (
|
||||
AsyncRuntime(thread_name="test-stream-runtime") as runtime,
|
||||
patch("EvoScientist.stream.display.Live"),
|
||||
pytest.raises(AsyncRuntimeError, match="running event loop"),
|
||||
):
|
||||
_run_streaming(
|
||||
agent=Mock(),
|
||||
message="message",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
gateway=gateway,
|
||||
runtime=runtime,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("second_thinking", "expected_count"),
|
||||
[(None, 1), ("Revised plan. " * 20, 2)],
|
||||
)
|
||||
def test_recursive_streaming_reuses_runtime_and_deduplicates_thinking(
|
||||
second_thinking, expected_count
|
||||
):
|
||||
initial_thinking = "Initial plan. " * 20
|
||||
stream_calls = 0
|
||||
loops: list[asyncio.AbstractEventLoop] = []
|
||||
|
||||
async def stream(_request):
|
||||
nonlocal stream_calls
|
||||
loops.append(asyncio.get_running_loop())
|
||||
stream_calls += 1
|
||||
if stream_calls == 1:
|
||||
yield {"type": "thinking", "content": initial_thinking}
|
||||
yield {
|
||||
"type": "ask_user",
|
||||
"interrupt_id": "ask-1",
|
||||
"tool_call_id": "tc-1",
|
||||
"questions": [{"question": "Continue?"}],
|
||||
}
|
||||
return
|
||||
if second_thinking is not None:
|
||||
yield {"type": "thinking", "content": second_thinking}
|
||||
else:
|
||||
yield {"type": "thinking", "content": initial_thinking}
|
||||
yield {"type": "text", "content": "final answer"}
|
||||
yield {"type": "done", "response": "final answer"}
|
||||
|
||||
sent_thinking: list[str] = []
|
||||
with (
|
||||
AsyncRuntime(thread_name="test-stream-runtime") as runtime,
|
||||
patch("EvoScientist.stream.display.Live"),
|
||||
):
|
||||
result = _run_streaming(
|
||||
agent=Mock(),
|
||||
message="test message",
|
||||
thread_id="thread1",
|
||||
show_thinking=False,
|
||||
interactive=True,
|
||||
on_thinking=sent_thinking.append,
|
||||
ask_user_prompt_fn=lambda _data: {
|
||||
"answers": ["yes"],
|
||||
"status": "answered",
|
||||
},
|
||||
gateway=FakeGraphGateway(stream=stream),
|
||||
runtime=runtime,
|
||||
)
|
||||
|
||||
assert result == "final answer"
|
||||
assert len(sent_thinking) == expected_count
|
||||
assert sent_thinking[0] == initial_thinking.rstrip()
|
||||
if second_thinking is not None:
|
||||
assert sent_thinking[1] == second_thinking.rstrip()
|
||||
assert loops[0] is loops[1]
|
||||
|
||||
Reference in New Issue
Block a user