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>
370 lines
11 KiB
Python
370 lines
11 KiB
Python
"""Tests for CLI serve/channel glue behavior."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import os
|
|
import signal
|
|
from types import SimpleNamespace
|
|
|
|
from EvoScientist.cli import commands
|
|
from EvoScientist.config import MemoryObservationWriter
|
|
from EvoScientist.runtime import AsyncRuntime
|
|
|
|
|
|
def _make_config(
|
|
*,
|
|
default_workdir: str = "",
|
|
channel_send_thinking: bool = True,
|
|
log_level: str = "warning",
|
|
channel_debug_tracing: bool = False,
|
|
auto_approve: bool = False,
|
|
auto_mode: bool = False,
|
|
enable_ask_user: bool = True,
|
|
dangerous_mode: bool = False,
|
|
):
|
|
return SimpleNamespace(
|
|
channel_enabled="telegram",
|
|
default_workdir=default_workdir,
|
|
channel_send_thinking=channel_send_thinking,
|
|
log_level=log_level,
|
|
channel_debug_tracing=channel_debug_tracing,
|
|
auto_approve=auto_approve,
|
|
auto_mode=auto_mode,
|
|
enable_ask_user=enable_ask_user,
|
|
dangerous_mode=dangerous_mode,
|
|
enable_async_subagents=False,
|
|
enable_scheduler=False,
|
|
memory_profile_enabled=True,
|
|
memory_observations_enabled=True,
|
|
memory_observation_writer=MemoryObservationWriter.ALL,
|
|
memory_workers_enabled=False,
|
|
memory_skill_synthesis_enabled=False,
|
|
model="test-model",
|
|
provider="anthropic",
|
|
anthropic_auth_mode="api_key",
|
|
openai_auth_mode="api_key",
|
|
)
|
|
|
|
|
|
def _run_serve_once(
|
|
monkeypatch,
|
|
config,
|
|
*,
|
|
workdir: str | None = None,
|
|
no_thinking: bool = False,
|
|
debug: bool = False,
|
|
cwd: str | None = None,
|
|
auto_approve: bool = False,
|
|
auto_mode: bool = False,
|
|
ask_user: bool = False,
|
|
dangerous: bool = False,
|
|
message_queue=None,
|
|
process_message=None,
|
|
):
|
|
import EvoScientist.config as config_mod
|
|
|
|
order: list[tuple[str, str | None]] = []
|
|
captured: dict[str, object] = {}
|
|
|
|
def _fake_set_workspace_root(path):
|
|
order.append(("set_workspace_root", str(path)))
|
|
|
|
def _fake_ensure_dirs():
|
|
order.append(("ensure_dirs", None))
|
|
|
|
def _fake_load_agent(
|
|
workspace_dir=None, checkpointer=None, config=None, *, runtime=None
|
|
):
|
|
captured["workspace_dir"] = workspace_dir
|
|
captured["async_runtime"] = runtime
|
|
return object()
|
|
|
|
def _fake_start_channels_bus_mode(cfg, agent, thread_id, *, send_thinking=None):
|
|
captured["started"] = True
|
|
captured["send_thinking"] = send_thinking
|
|
captured["thread_id"] = thread_id
|
|
|
|
def _fake_channels_stop(channel_type=None, *, runtime=None):
|
|
captured["stopped"] = True
|
|
|
|
class _InterruptQueue:
|
|
"""A fake queue whose get() immediately raises KeyboardInterrupt."""
|
|
|
|
def get(self, timeout=None):
|
|
raise KeyboardInterrupt()
|
|
|
|
monkeypatch.setattr(commands, "set_workspace_root", _fake_set_workspace_root)
|
|
monkeypatch.setattr(commands, "ensure_dirs", _fake_ensure_dirs)
|
|
monkeypatch.setattr(commands, "_load_agent", _fake_load_agent)
|
|
monkeypatch.setattr(
|
|
commands, "_start_channels_bus_mode", _fake_start_channels_bus_mode
|
|
)
|
|
monkeypatch.setattr(commands, "_channels_stop", _fake_channels_stop)
|
|
monkeypatch.setattr(
|
|
commands,
|
|
"_message_queue",
|
|
message_queue if message_queue is not None else _InterruptQueue(),
|
|
)
|
|
if process_message is not None:
|
|
monkeypatch.setattr(commands, "_serve_process_message", process_message)
|
|
|
|
def _fake_get_effective_config(cli_overrides=None):
|
|
captured["cli_overrides"] = dict(cli_overrides or {})
|
|
merged = vars(config).copy()
|
|
merged.update(cli_overrides or {})
|
|
return SimpleNamespace(**merged)
|
|
|
|
monkeypatch.setattr(config_mod, "get_effective_config", _fake_get_effective_config)
|
|
monkeypatch.setattr(config_mod, "apply_config_to_env", lambda _cfg: None)
|
|
|
|
if cwd is not None:
|
|
monkeypatch.setattr(commands.os, "getcwd", lambda: cwd)
|
|
|
|
with AsyncRuntime(thread_name="test-serve-runtime") as runtime:
|
|
monkeypatch.setattr(commands, "_get_cli_async_runtime", lambda _ctx: runtime)
|
|
commands.serve(
|
|
object(),
|
|
no_thinking=no_thinking,
|
|
workdir=workdir,
|
|
debug=debug,
|
|
auto_approve=auto_approve,
|
|
auto_mode=auto_mode,
|
|
ask_user=ask_user,
|
|
dangerous=dangerous,
|
|
)
|
|
return order, captured
|
|
|
|
|
|
def test_serve_workdir_has_highest_priority_and_sets_root_before_ensure(
|
|
monkeypatch, tmp_path
|
|
):
|
|
cfg_ws = tmp_path / "cfg_ws"
|
|
cli_ws = tmp_path / "cli_ws"
|
|
config = _make_config(default_workdir=str(cfg_ws), channel_send_thinking=True)
|
|
|
|
order, captured = _run_serve_once(
|
|
monkeypatch,
|
|
config,
|
|
workdir=str(cli_ws),
|
|
)
|
|
|
|
expected = str(cli_ws.resolve())
|
|
assert captured["workspace_dir"] == expected
|
|
assert any(step == ("set_workspace_root", expected) for step in order)
|
|
set_idx = next(i for i, step in enumerate(order) if step[0] == "set_workspace_root")
|
|
ensure_idx = next(i for i, step in enumerate(order) if step[0] == "ensure_dirs")
|
|
assert set_idx < ensure_idx
|
|
|
|
|
|
def test_serve_uses_config_default_workdir_when_no_cli_workdir(monkeypatch, tmp_path):
|
|
cfg_ws = tmp_path / "cfg_ws"
|
|
config = _make_config(default_workdir=str(cfg_ws), channel_send_thinking=True)
|
|
|
|
order, captured = _run_serve_once(monkeypatch, config)
|
|
|
|
expected = str(cfg_ws.resolve())
|
|
assert captured["workspace_dir"] == expected
|
|
assert ("set_workspace_root", expected) in order
|
|
|
|
|
|
def test_serve_uses_cwd_when_no_workdir_config(monkeypatch, tmp_path):
|
|
cwd = str(tmp_path.resolve())
|
|
config = _make_config(default_workdir="", channel_send_thinking=True)
|
|
|
|
order, captured = _run_serve_once(
|
|
monkeypatch,
|
|
config,
|
|
cwd=cwd,
|
|
)
|
|
|
|
assert captured["workspace_dir"] == cwd
|
|
assert ("set_workspace_root", cwd) in order
|
|
|
|
|
|
def test_serve_channel_thinking_respects_config_and_no_thinking(monkeypatch, tmp_path):
|
|
ws = str((tmp_path / "ws").resolve())
|
|
|
|
_, captured_default = _run_serve_once(
|
|
monkeypatch,
|
|
_make_config(default_workdir=ws, channel_send_thinking=True),
|
|
)
|
|
assert captured_default["send_thinking"] is True
|
|
|
|
_, captured_cfg_off = _run_serve_once(
|
|
monkeypatch,
|
|
_make_config(default_workdir=ws, channel_send_thinking=False),
|
|
)
|
|
assert captured_cfg_off["send_thinking"] is False
|
|
|
|
_, captured_cli_off = _run_serve_once(
|
|
monkeypatch,
|
|
_make_config(default_workdir=ws, channel_send_thinking=True),
|
|
no_thinking=True,
|
|
)
|
|
assert captured_cli_off["send_thinking"] is False
|
|
|
|
|
|
def test_serve_debug_sets_log_level_and_channel_trace(monkeypatch, tmp_path):
|
|
ws = str((tmp_path / "ws").resolve())
|
|
config = _make_config(default_workdir=ws, channel_send_thinking=True)
|
|
configure_calls: list[tuple[str | None, str | None]] = []
|
|
monkeypatch.delenv("EVOSCIENTIST_LOG_LEVEL", raising=False)
|
|
monkeypatch.delenv("EVOSCIENTIST_CHANNEL_DEBUG_TRACING", raising=False)
|
|
|
|
monkeypatch.setattr(
|
|
commands,
|
|
"_configure_logging",
|
|
lambda: configure_calls.append(
|
|
(
|
|
os.environ.get("EVOSCIENTIST_LOG_LEVEL"),
|
|
os.environ.get("EVOSCIENTIST_CHANNEL_DEBUG_TRACING"),
|
|
)
|
|
),
|
|
)
|
|
|
|
_, captured = _run_serve_once(
|
|
monkeypatch,
|
|
config,
|
|
workdir=ws,
|
|
debug=True,
|
|
)
|
|
|
|
assert captured["cli_overrides"] == {
|
|
"log_level": "DEBUG",
|
|
"channel_debug_tracing": True,
|
|
}
|
|
assert configure_calls == [("DEBUG", "true")]
|
|
|
|
|
|
def test_serve_auto_approve_only_sets_auto_approve(monkeypatch, tmp_path):
|
|
ws = str((tmp_path / "ws").resolve())
|
|
config = _make_config(default_workdir=ws, enable_ask_user=True)
|
|
|
|
_, captured = _run_serve_once(
|
|
monkeypatch,
|
|
config,
|
|
workdir=ws,
|
|
auto_approve=True,
|
|
)
|
|
|
|
assert captured["cli_overrides"] == {"auto_approve": True}
|
|
|
|
|
|
def test_serve_auto_mode_implies_auto_approve_and_disables_ask_user(
|
|
monkeypatch, tmp_path
|
|
):
|
|
ws = str((tmp_path / "ws").resolve())
|
|
config = _make_config(default_workdir=ws, enable_ask_user=True)
|
|
|
|
_, captured = _run_serve_once(
|
|
monkeypatch,
|
|
config,
|
|
workdir=ws,
|
|
auto_mode=True,
|
|
ask_user=True,
|
|
)
|
|
|
|
assert captured["cli_overrides"] == {
|
|
"auto_mode": True,
|
|
"auto_approve": True,
|
|
"enable_ask_user": False,
|
|
}
|
|
|
|
|
|
def test_serve_dangerous_sets_dangerous_mode(monkeypatch, tmp_path):
|
|
ws = str((tmp_path / "ws").resolve())
|
|
config = _make_config(default_workdir=ws)
|
|
|
|
_, captured = _run_serve_once(
|
|
monkeypatch,
|
|
config,
|
|
workdir=ws,
|
|
dangerous=True,
|
|
)
|
|
|
|
assert captured["cli_overrides"] == {"dangerous_mode": True}
|
|
|
|
|
|
def test_serve_sigterm_cancels_active_message_scope_before_shutdown(
|
|
monkeypatch, tmp_path
|
|
):
|
|
handlers = {}
|
|
message = commands.ChannelMessage(
|
|
msg_id="message-1",
|
|
content="run a long command",
|
|
sender="user",
|
|
channel_type="telegram",
|
|
chat_id="chat-1",
|
|
)
|
|
|
|
class _OneMessageQueue:
|
|
def get(self, timeout=None):
|
|
return message
|
|
|
|
def _fake_signal(signum, handler):
|
|
previous = handlers.get(signum, signal.SIG_DFL)
|
|
handlers[signum] = handler
|
|
return previous
|
|
|
|
cancelled_scopes = []
|
|
|
|
def _fake_process_message(*args, **kwargs):
|
|
handlers[signal.SIGTERM](signal.SIGTERM, None)
|
|
|
|
monkeypatch.setattr(signal, "signal", _fake_signal)
|
|
monkeypatch.setattr(
|
|
"EvoScientist.stream.display.request_stream_cancel",
|
|
cancelled_scopes.append,
|
|
)
|
|
|
|
_run_serve_once(
|
|
monkeypatch,
|
|
_make_config(default_workdir=str(tmp_path)),
|
|
message_queue=_OneMessageQueue(),
|
|
process_message=_fake_process_message,
|
|
)
|
|
|
|
assert cancelled_scopes == ["channel:telegram:chat-1:message-1"]
|
|
|
|
|
|
def test_serve_sigint_cancels_active_message_scope_before_interrupt(
|
|
monkeypatch, tmp_path
|
|
):
|
|
handlers = {}
|
|
message = commands.ChannelMessage(
|
|
msg_id="message-2",
|
|
content="run a long command",
|
|
sender="user",
|
|
channel_type="telegram",
|
|
chat_id="chat-2",
|
|
)
|
|
|
|
class _OneMessageQueue:
|
|
def get(self, timeout=None):
|
|
return message
|
|
|
|
def _fake_signal(signum, handler):
|
|
previous = handlers.get(signum, signal.SIG_DFL)
|
|
handlers[signum] = handler
|
|
return previous
|
|
|
|
cancelled_scopes = []
|
|
|
|
def _fake_process_message(*args, **kwargs):
|
|
handlers[signal.SIGINT](signal.SIGINT, None)
|
|
|
|
monkeypatch.setattr(signal, "signal", _fake_signal)
|
|
monkeypatch.setattr(
|
|
"EvoScientist.stream.display.request_stream_cancel",
|
|
cancelled_scopes.append,
|
|
)
|
|
|
|
_run_serve_once(
|
|
monkeypatch,
|
|
_make_config(default_workdir=str(tmp_path)),
|
|
message_queue=_OneMessageQueue(),
|
|
process_message=_fake_process_message,
|
|
)
|
|
|
|
assert cancelled_scopes == ["channel:telegram:chat-2:message-2"]
|