8b1451cdda
* 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>
76 lines
2.4 KiB
Python
76 lines
2.4 KiB
Python
"""Behavioral tests for channel sends crossing frontend event loops."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
import threading
|
|
|
|
import pytest
|
|
|
|
from EvoScientist.cli.channel_sends import PendingChannelSends
|
|
from EvoScientist.runtime import AsyncRuntime
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_pending_send_does_not_stall_owned_runtime() -> None:
|
|
"""A blocked channel transport must not block unrelated runtime work."""
|
|
bus_loop = asyncio.get_running_loop()
|
|
send_started = asyncio.Event()
|
|
release_send = asyncio.Event()
|
|
runtime_progressed = threading.Event()
|
|
sends = PendingChannelSends(bus_loop, logging.getLogger(__name__))
|
|
|
|
async def _blocked_send() -> None:
|
|
send_started.set()
|
|
await release_send.wait()
|
|
|
|
async def _stream_callback_and_probe() -> None:
|
|
sends.submit(_blocked_send(), "Thinking")
|
|
await asyncio.sleep(0)
|
|
runtime_progressed.set()
|
|
|
|
with AsyncRuntime(thread_name="test-channel-send-runtime") as runtime:
|
|
callback = runtime.submit(_stream_callback_and_probe)
|
|
await asyncio.wait_for(send_started.wait(), timeout=1)
|
|
assert runtime_progressed.wait(timeout=1)
|
|
callback.result(timeout=1)
|
|
|
|
settle = asyncio.create_task(sends.settle_async())
|
|
await asyncio.sleep(0)
|
|
assert not settle.done()
|
|
|
|
release_send.set()
|
|
await asyncio.wait_for(settle, timeout=1)
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_async_settlement_waits_for_every_scheduled_send() -> None:
|
|
"""The channel response can wait for all callback delivery off-loop."""
|
|
first_started = asyncio.Event()
|
|
first_release = asyncio.Event()
|
|
events: list[str] = []
|
|
sends = PendingChannelSends(asyncio.get_running_loop(), logging.getLogger(__name__))
|
|
|
|
async def _first() -> None:
|
|
events.append("first-started")
|
|
first_started.set()
|
|
await first_release.wait()
|
|
events.append("first-finished")
|
|
|
|
async def _second() -> None:
|
|
events.append("second-finished")
|
|
|
|
sends.submit(_first(), "First")
|
|
sends.submit(_second(), "Second")
|
|
|
|
settle = asyncio.create_task(sends.settle_async())
|
|
await asyncio.wait_for(first_started.wait(), timeout=1)
|
|
assert not settle.done()
|
|
assert events == ["first-started"]
|
|
|
|
first_release.set()
|
|
await asyncio.wait_for(settle, timeout=1)
|
|
|
|
assert events == ["first-started", "first-finished", "second-finished"]
|