Files
dinos 8b1451cdda 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>
2026-07-27 14:17:57 +01:00

103 lines
3.1 KiB
Python

"""Agent loading and workspace helpers."""
import os
from datetime import datetime
from pathlib import Path
from typing import TYPE_CHECKING
from ..paths import new_run_dir
if TYPE_CHECKING:
from langgraph.graph.state import CompiledStateGraph
from ..runtime import AsyncRuntime
def _shorten_path(path: str) -> str:
"""Shorten absolute path to relative path from current directory."""
if not path:
return path
try:
cwd = os.getcwd()
if path.startswith(cwd):
rel = path[len(cwd) :].lstrip(os.sep)
return (
os.path.join(os.path.basename(cwd), rel)
if rel
else os.path.basename(cwd)
)
return path
except Exception:
return path
def _deduplicate_run_name(name: str, runs_dir: Path | None = None) -> str:
"""Return *name* if available, otherwise *name_1*, *name_2*, etc."""
if runs_dir is None:
from ..paths import RUNS_DIR
runs_dir = RUNS_DIR
if not (runs_dir / name).exists():
return name
i = 1
while (runs_dir / f"{name}_{i}").exists():
i += 1
return f"{name}_{i}"
def _create_session_workspace(name: str | None = None) -> str:
"""Create a per-session workspace directory and return its path.
Args:
name: Optional human-friendly run name. Duplicates are resolved
by appending ``_1``, ``_2``, etc. Falls back to a timestamp
if *name* is None.
"""
if name:
from ..paths import RUNS_DIR
session_id = _deduplicate_run_name(name, RUNS_DIR)
else:
session_id = datetime.now().strftime("%Y%m%d_%H%M%S")
workspace_dir = str(new_run_dir(session_id))
os.makedirs(workspace_dir, exist_ok=True)
return workspace_dir
def _load_agent(
workspace_dir: str | None = None,
checkpointer=None,
config=None,
chat_model=None,
*,
on_mcp_progress=None,
events=None,
runtime: "AsyncRuntime | None" = None,
) -> "CompiledStateGraph":
"""Load the CLI agent with optional persistent checkpointer.
Args:
workspace_dir: Optional per-session workspace directory.
checkpointer: Optional LangGraph checkpointer (e.g. ``AsyncSqliteSaver``).
Falls back to ``InMemorySaver`` when ``None``.
config: Optional pre-loaded ``EvoScientistConfig``. Forwarded to
``create_cli_agent`` to avoid double config loading.
chat_model: Optional pre-built chat model. Forwarded to
``create_cli_agent``; combined with an explicit ``config`` it
selects the pure (no module-global write) build path.
on_mcp_progress: Optional per-server MCP progress callback.
Signature ``(event, server_name, detail) -> None``.
runtime: Optional application-scoped runtime used for MCP discovery.
"""
from ..EvoScientist import create_cli_agent
return create_cli_agent(
workspace_dir=workspace_dir,
checkpointer=checkpointer,
config=config,
chat_model=chat_model,
on_mcp_progress=on_mcp_progress,
events=events,
runtime=runtime,
)