From ea99ce9f7ebe869ac0c9d6f42cdd5481c6037b9e Mon Sep 17 00:00:00 2001 From: m4 Date: Sun, 13 Sep 2026 15:12:17 +0800 Subject: [PATCH] feat(runtime): host execution registry, web checkpointer, stop control and sandbox cancellation --- EvoScientist/EvoScientist.py | 77 +++++- EvoScientist/langgraph_dev/http.py | 197 +++++++++++++- EvoScientist/langgraph_dev/langgraph.web.json | 19 ++ EvoScientist/langgraph_dev/worker_exit.py | 236 ++++++++++++++++ EvoScientist/middleware/provider_context.py | 8 +- EvoScientist/middleware/recoverable_tools.py | 19 +- EvoScientist/native_sandbox.py | 23 +- EvoScientist/stream/emitter.py | 13 + EvoScientist/stream/events.py | 22 +- EvoScientist/stream/stop.py | 120 ++++++++ EvoScientist/stream/v3_payloads.py | 21 ++ EvoScientist/web_checkpointer.py | 256 ++++++++++++++++++ EvoScientist/web_runtime.py | 6 +- start-langgraph.sh | 14 +- 14 files changed, 1000 insertions(+), 31 deletions(-) create mode 100644 EvoScientist/langgraph_dev/langgraph.web.json create mode 100644 EvoScientist/langgraph_dev/worker_exit.py create mode 100644 EvoScientist/stream/stop.py create mode 100644 EvoScientist/web_checkpointer.py diff --git a/EvoScientist/EvoScientist.py b/EvoScientist/EvoScientist.py index 0a51282..a3a9ec1 100644 --- a/EvoScientist/EvoScientist.py +++ b/EvoScientist/EvoScientist.py @@ -1104,6 +1104,30 @@ def __getattr__(name: str): # ============================================================================= +def _create_run_summarization_middleware(model, backend, summarizer): + """Resolve context thresholds against the consuming model, not the summarizer.""" + from deepagents.middleware.summarization import ( + SummarizationMiddleware, compute_summarization_defaults, + ) + + defaults = compute_summarization_defaults(model) + window = (model.profile or {}).get("max_input_tokens") + + def absolute(value): + if isinstance(value, tuple) and value[0] == "fraction": + if not isinstance(window, int) or window <= 0: + raise ValueError("fraction threshold requires a main-model context window") + return ("tokens", int(window * value[1])) + if isinstance(value, dict): + return {key: absolute(item) for key, item in value.items()} + return value + + return SummarizationMiddleware( + model=summarizer, backend=backend, + trim_tokens_to_summarize=None, **absolute(defaults), + ) + + def create_cli_agent( workspace_dir: str | None = None, checkpointer=None, @@ -1180,9 +1204,15 @@ def create_cli_agent( # locals and write no module globals. Otherwise keep the legacy # global-writing behavior — callers that pass config= only (CLI startup, # langgraph dev) rely on it to seat the active config/model. + is_web = getattr(execution_profile, "name", "") in {"web_v1", "web_v3"} + if is_web and any(value is None for value in ( + config, chat_model, workspace_dir, memory_dir, workspace_backend, checkpointer + )): + raise ValueError("Web execution requires explicit config/model/workspace/memory/backend/checkpointer") if config is not None and chat_model is not None: cfg = config - _apply_env_from_config(cfg) + if not is_web: + _apply_env_from_config(cfg) else: cfg = _ensure_config(config) chat_model = None @@ -1228,7 +1258,8 @@ def create_cli_agent( # Always construct fresh backends from current paths (avoids stale # module-level backend when workspace root changed at runtime). - set_active_workspace(workspace_dir) + if not is_web: + set_active_workspace(workspace_dir) ws_backend = workspace_backend if ws_backend is None: ws_backend = CustomSandboxBackend( @@ -1248,6 +1279,7 @@ def create_cli_agent( ) be = CompositeBackend( default=ws_backend, + artifacts_root="/workspace" if is_web else "/", routes={ "/skills/": sk_backend, "/memories/": mem_backend, @@ -1303,7 +1335,9 @@ def create_cli_agent( ) mw.insert( (error_index + 1) if error_index is not None else 0, - ProviderContextMediaMiddleware(be), + ProviderContextMediaMiddleware( + be, media_prefix="/workspace/artifacts/model-output" if is_web else "/artifacts/model-output" + ), ) if main_agent_route_middleware is not None: configurable_index = next( @@ -1324,7 +1358,13 @@ def create_cli_agent( # HITL on main agent only — passing `interrupt_on=` to create_deep_agent # would propagate it to every subagent, breaking parallel execute calls # (multi-pending-interrupt LangGraph error). - if not cfg.auto_approve: + if is_web: + from .middleware.dynamic_review import DynamicReviewMiddleware + + mw.append(DynamicReviewMiddleware(interrupt_on={ + "execute": True, "run_in_background": True, "schedule_task": True, + })) + elif not cfg.auto_approve: mw.append( HumanInTheLoopMiddleware( interrupt_on={ @@ -1346,8 +1386,37 @@ def create_cli_agent( ) if not enable_subagents: kwargs = {**kwargs, "subagents": []} + if is_web: + kwargs = {**kwargs, "tools": [ + tool for tool in kwargs.get("tools", []) + if getattr(tool, "name", "") != "skill_manager" + ]} kwargs = _apply_budgeted_skill_context(kwargs, be) + if agent_model_set is not None: + from types import FunctionType + def create_run_summarizer(model, backend): + return _create_run_summarization_middleware( + model, backend, agent_model_set.deepagents_summarizer, + ) + + # deepagents currently has no per-call summarizer factory parameter. + # Scope its existing assembly function to this run, never patch module + # globals or register a process-wide profile containing tenant models. + if (not isinstance(create_deep_agent, FunctionType) + or "create_summarization_middleware" not in create_deep_agent.__code__.co_names): + raise RuntimeError("Unsupported deepagents assembly; revalidate summarizer adapter") + create_deep_agent = FunctionType( + create_deep_agent.__code__, + {**create_deep_agent.__globals__, + "create_summarization_middleware": create_run_summarizer}, + create_deep_agent.__name__, + create_deep_agent.__defaults__, + create_deep_agent.__closure__, + ) + from deepagents import create_deep_agent as original_factory + create_deep_agent.__kwdefaults__ = original_factory.__kwdefaults__ + return create_deep_agent( **kwargs, checkpointer=checkpointer, diff --git a/EvoScientist/langgraph_dev/http.py b/EvoScientist/langgraph_dev/http.py index e3906ca..dbd6af6 100644 --- a/EvoScientist/langgraph_dev/http.py +++ b/EvoScientist/langgraph_dev/http.py @@ -26,7 +26,7 @@ import asyncio import json import os import secrets -from typing import Any +from typing import Any, cast from uuid import UUID from starlette.applications import Starlette @@ -104,6 +104,11 @@ async def recoverable_run_capabilities(_request: Request) -> JSONResponse: { "version": 1, "deterministic_run_id": True, + # The dev adapter uses process-local Runs storage. Checkpoint + # durability does not make execution identity restart-safe. + "durable_run_identity": False, + "worker_exit_confirmation": True, + "run_not_found_proves_absence": False, "stream_resumable": True, "durability_sync": True, "multitask_enqueue": True, @@ -270,16 +275,99 @@ async def bind_workspace_run(request: Request) -> JSONResponse: return JSONResponse(_run_payload(run)) +async def _compatible_checkpoint(conn, thread_id: str, assistant_id: str, + config: dict) -> tuple[bool, bool]: + """Read through the API-owned saver and graph factory, never execute here.""" + from langgraph_api._checkpointer import get_checkpointer + from langgraph_api.graph import get_graph, graph_exists + from langgraph_api.store import get_store + + saver = await get_checkpointer(conn=conn) + read_config = {**config, "configurable": { + **config.get("configurable", {}), "thread_id": thread_id, + "checkpoint_ns": "", + }} + # Admission always checks the current head, never a caller-selected ancestor. + read_config["configurable"].pop("checkpoint_id", None) + checkpoint = await saver.aget_tuple(read_config) + if checkpoint is None: + return False, False + graph_id = assistant_id + if not graph_exists(graph_id): + from langgraph_runtime.ops import Assistants + from langgraph_api.utils import fetchone + assistant = await fetchone(await Assistants.get(conn, UUID(assistant_id))) + graph_id = assistant["graph_id"] + if checkpoint.metadata.get("graph_id", graph_id) != graph_id: + raise ValueError("checkpoint graph mismatch") + # get_graph enters coroutine/async-context-manager factories and binds the + # same API saver used by the worker. aget_state also validates delta seeds. + async with get_graph(graph_id, read_config, checkpointer=saver, + store=await get_store(), access_context="threads.read") as graph: + state = await graph.aget_state(read_config, subgraphs=True) + pending = bool(getattr(state, "interrupts", ())) or any( + task.interrupts for task in state.tasks + ) + return True, pending + + +async def _has_legacy_history(conn, thread_id: str) -> bool: + """Resolve Web history from PostgreSQL and the Runtime registry only.""" + dsn = os.getenv("EVOSCIENTIST_WEB_CHECKPOINT_DSN", "") + if dsn: + from psycopg import AsyncConnection + async with await AsyncConnection.connect( + dsn, autocommit=True, connect_timeout=5, + options="-c default_transaction_read_only=on -c search_path=public", + ) as pg: + async with pg.cursor() as cursor: + await cursor.execute("SELECT EXISTS(SELECT 1 FROM threads WHERE id=%s::uuid)", + (thread_id,)) + row = await cursor.fetchone() + if row and row[0]: + return True + from langgraph_api.utils import fetchone + from langgraph_runtime.ops import Threads + try: + await fetchone(await Threads.get(conn, UUID(thread_id))) + except Exception as exc: + if getattr(exc, "status_code", None) != 404: + raise + return False + return True + + +async def _history_admission(conn, thread_id, assistant_id, config, operation, history): + exists, pending = await _compatible_checkpoint(conn, thread_id, assistant_id, config) + if operation == "resume": + return "resume" if exists and pending else "CHECKPOINT_RESUME_UNAVAILABLE" + if pending: + return "THREAD_AWAITING_INPUT" + if exists: + return "append" + if history is not None: + return "initialize" + return "HISTORY_REQUIRED" if await _has_legacy_history(conn, thread_id) else "new" + + async def create_recoverable_run(request: Request) -> JSONResponse: """Create a LangGraph Run with a caller-owned deterministic UUID. LangGraph's public create endpoint always generates its own UUID. This adapter performs lookup and insertion while holding the process-wide run creation lock and passes the durable request UUID to ``create_valid_run``. - Retrying after a lost HTTP response therefore cannot create another Run. + This prevents duplicate creation only while this process retains the Run. + A lost response followed by a restart MUST NOT be retried as a fresh create: + the dev backend forgets Runs, and a 404 is not evidence of non-execution. """ - if request.headers.get("x-auth-scheme") != "langsmith": + from EvoScientist.internal_service import internal_service_token + + token = internal_service_token() + if not token: + return JSONResponse({"code": "WORKSPACE_SERVICE_UNAVAILABLE"}, status_code=503) + header = request.headers.get("authorization", "") + if not header.startswith("Bearer ") or not secrets.compare_digest(header[7:], token): return JSONResponse({"code": "UNAUTHORIZED"}, status_code=401) value = await request.json() if not isinstance(value, dict): @@ -303,17 +391,25 @@ async def create_recoverable_run(request: Request) -> JSONResponse: if operation == "resume": if ( value.get("input") is not None + or value.get("history") is not None or not isinstance(command, dict) or set(command) != {"resume"} ): return JSONResponse({"code": "INVALID_RESUME_REQUEST"}, status_code=400) elif command is not None: return JSONResponse({"code": "INVALID_START_REQUEST"}, status_code=400) + history = value.get("history") + # History validation is performed under the create lock, after idempotent + # lookup. Only that branch can attest this request did not create a Run. from langgraph_api.models.run import Runs, create_valid_run from langgraph_api.utils import fetchone from langgraph_runtime.database import connect + from EvoScientist.langgraph_dev import worker_exit + + worker_exit.install() + payload = { "assistant_id": assistant_id, "input": value.get("input"), @@ -331,6 +427,13 @@ async def create_recoverable_run(request: Request) -> JSONResponse: "run_request_id": run_request_id, "request_hash": request_hash, } + # The admission read and the worker must address the same current head. + # Forking from an ancestor is not part of the recoverable-run contract. + configurable = dict(payload["config"].get("configurable", {})) + for key in ("checkpoint_id", "checkpoint_map"): + configurable.pop(key, None) + configurable.update(thread_id=thread_id, checkpoint_ns="") + payload["config"] = {**payload["config"], "configurable": configurable} async with _recoverable_run_lock: async with connect() as conn: existing_iter = await Runs.get(conn, run_id, thread_id=UUID(thread_id)) @@ -353,6 +456,50 @@ async def create_recoverable_run(request: Request) -> JSONResponse: return JSONResponse( {"run_id": str(existing["run_id"]), "status": existing["status"], "created": False} ) + try: + admission = await _history_admission( + conn, thread_id, assistant_id, payload["config"], operation, history, + ) + except Exception: + return JSONResponse({ + "code": "CHECKPOINT_UNAVAILABLE", "create_disposition": "not_created", + "run_id": str(run_id), "request_hash": request_hash, + }, status_code=503) + if admission not in {"new", "initialize", "append", "resume"}: + return JSONResponse({ + "code": admission, "create_disposition": "not_created", + "run_id": str(run_id), "request_hash": request_hash, + }, status_code=409) + if history is not None: + from EvoScientist.llm.history_rebuild import committed_history_input + + try: + if not isinstance(history, dict): + raise ValueError("history must be an object") + import hashlib + + history_hash = hashlib.sha256(json.dumps( + history, ensure_ascii=False, sort_keys=True, separators=(",", ":"), + ).encode()).hexdigest() + if value.get("history_hash") != history_hash: + raise ValueError("history hash mismatch") + payload["input"] = committed_history_input( + history, payload["input"], thread_id=thread_id, run_id=str(run_id), + checkpoint_exists=admission == "append", + ) + if admission == "initialize": + payload["metadata"].update( + history_hash=history_hash, + history_revision=history["conversation_revision"], + history_schema=history["schema"], + ) + except (ValueError, TypeError, KeyError, AttributeError): + return JSONResponse({ + "code": "INVALID_HISTORY_REQUEST", + "create_disposition": "not_created", + "run_id": str(run_id), "request_hash": request_hash, + }, status_code=400) + worker_exit.reserve(thread_id, str(run_id), request_hash) created = await create_valid_run( conn, thread_id, @@ -366,9 +513,53 @@ async def create_recoverable_run(request: Request) -> JSONResponse: ) +async def cancel_recoverable_run(request: Request) -> JSONResponse: + from EvoScientist.internal_service import internal_service_token + + token = internal_service_token() + if not token: + return JSONResponse({"code": "WORKSPACE_SERVICE_UNAVAILABLE"}, status_code=503) + header = request.headers.get("authorization", "") + if not header.startswith("Bearer ") or not secrets.compare_digest(header[7:], token): + return JSONResponse({"code": "UNAUTHORIZED"}, status_code=401) + from EvoScientist.langgraph_dev import worker_exit + from langgraph_api.models.run import Runs + from langgraph_runtime.database import connect + + thread_id = UUID(request.path_params["thread_id"]) + run_id = UUID(request.path_params["run_id"]) + # This service principal controls only pairs admitted by authenticated create. + if not worker_exit.is_reserved(str(thread_id), str(run_id)): + return JSONResponse({"code": "RUN_NOT_AUTHORIZED"}, status_code=404) + # Close worker admission before notifying the original runtime control queue. + initial = worker_exit.cancel_and_inspect(str(thread_id), str(run_id)) + if initial.get("execution_exited") is not True: + async with connect() as conn: + try: + await Runs.cancel(conn, [run_id], thread_id=thread_id, action="interrupt") + except Exception as exc: + if getattr(exc, "status_code", None) not in {404, 409}: + raise + receipt = await worker_exit.wait_for_exit(str(thread_id), str(run_id)) + if receipt.get("execution_exited") is True: + async with connect() as conn: + try: + await Runs.delete(cast(Any, conn), run_id, thread_id=thread_id) + except Exception as exc: + if getattr(exc, "status_code", None) != 404: + raise + receipt = {**receipt, "checkpoint_cleanup": "completed"} + return JSONResponse(receipt) + + app = Starlette( routes=[ Route("/api/models", get_models, methods=["GET"]), + Route( + "/api/ai4sci/recoverable-runs/{thread_id}/{run_id}/cancel", + cancel_recoverable_run, + methods=["POST"], + ), Route( "/api/ai4sci/recoverable-runs/capabilities", recoverable_run_capabilities, diff --git a/EvoScientist/langgraph_dev/langgraph.web.json b/EvoScientist/langgraph_dev/langgraph.web.json new file mode 100644 index 0000000..77fb419 --- /dev/null +++ b/EvoScientist/langgraph_dev/langgraph.web.json @@ -0,0 +1,19 @@ +{ + "dependencies": ["."], + "graphs": { + "EvoScientist": "EvoScientist.langgraph_dev.main_graph:EvoScientist_agent", + "writing-agent": "EvoScientist.langgraph_dev.graphs:writing_agent", + "data-analysis-agent": "EvoScientist.langgraph_dev.graphs:data_analysis_agent", + "scheduler": "EvoScientist.langgraph_dev.graphs:scheduler", + "evomemory-subagent-worker": "EvoScientist.langgraph_dev.graphs:evomemory_subagent_worker", + "evomemory-turn-worker": "EvoScientist.langgraph_dev.graphs:evomemory_turn_worker", + "evomemory-observation-linker": "EvoScientist.langgraph_dev.graphs:evomemory_observation_linker", + "evomemory-autoskills": "EvoScientist.langgraph_dev.graphs:evomemory_autoskills" + }, + "checkpointer": { + "backend": "custom", + "path": "EvoScientist.web_checkpointer.create_web_checkpointer" + }, + "config": {"recursion_limit": 5000}, + "http": {"app": "EvoScientist.langgraph_dev.http:app"} +} \ No newline at end of file diff --git a/EvoScientist/langgraph_dev/worker_exit.py b/EvoScientist/langgraph_dev/worker_exit.py new file mode 100644 index 0000000..1d07840 --- /dev/null +++ b/EvoScientist/langgraph_dev/worker_exit.py @@ -0,0 +1,236 @@ +"""Process-local exit receipts for the existing LangGraph worker. + +No receipt survives a restart. Missing receipts never prove execution absent. +The admission gate also covers a queued worker that has not started yet. +""" + +from __future__ import annotations + +import asyncio +from concurrent.futures import Future +from contextvars import ContextVar +from functools import wraps +from threading import RLock +from typing import Any, cast + +_lock = RLock() +_runs: dict[tuple[str, str], dict] = {} +_installed = False +_owned: ContextVar[bool] = ContextVar("recoverable_worker_owned", default=False) +_evidence: ContextVar[dict | None] = ContextVar("worker_exit_evidence", default=None) + + +class _ExitBoundary: + def __init__(self, context, *, run=False): + self.context = context + self.run = run + + async def __aenter__(self): + return await self.context.__aenter__() + + async def __aexit__(self, typ, value, tb): + evidence = _evidence.get() + try: + result = await self.context.__aexit__(typ, value, tb) + except BaseException as exc: + # An asynccontextmanager propagating its body exception is not a + # cleanup failure. Replacement exceptions are fail-closed. + if evidence is not None and exc is not value: + evidence["cleanup_failed"] = True + if evidence is not None and self.run and exc is value: + evidence["run_exited"] = True + raise + else: + if evidence is not None and self.run: + evidence["run_exited"] = True + return result + + +async def _drain(task): + while not task.done(): + try: + await asyncio.shield(task) + except asyncio.CancelledError: + continue + except BaseException: + break + if not task.cancelled(): + task.exception() + + +async def _await_remote_future(remote: Future): + """Do not let caller cancellation detach a coroutine on a worker thread.""" + async def wait(): + return await asyncio.wrap_future(remote) + + task = asyncio.create_task(wait()) + try: + return await asyncio.shield(task) + except asyncio.CancelledError: + await _drain(task) + raise + + +async def _persistent_cancellation_listener(original, queue, run_id, thread_id, done): + """The in-memory listener's idle timeout is a poll, not end-of-listening.""" + while not done.is_set(): + await original(queue, run_id, thread_id, done) + + +def install() -> None: + global _installed + from langgraph_api import worker, stream + from langgraph_runtime_inmem import ops as inmem_ops + from langgraph.pregel._loop import AsyncPregelLoop + from quickjs_rs.threading import ThreadWorker + + with _lock: + if _installed: + return + original = worker.worker + original_to_thread = asyncio.to_thread + original_enter = worker.Runs.enter + original_closing = stream.aclosing + original_stack = stream.AsyncExitStack + original_loop_exit = AsyncPregelLoop.__aexit__ + original_cancellation_listener = inmem_ops.listen_for_cancellation + + + class ObservedStack(original_stack): + async def __aexit__(self, typ, value, tb): + return await _ExitBoundary(super()).__aexit__(typ, value, tb) + + async def loop_exit(self, exc_type, exc_value, traceback): + evidence = _evidence.get() + if evidence is None: + return await original_loop_exit(self, exc_type, exc_value, traceback) + try: + return await original_loop_exit(self, exc_type, exc_value, traceback) + except asyncio.CancelledError as exc: + # Installed Pregel exposes its outstanding exit task in args. + tasks = [arg for arg in exc.args if isinstance(arg, asyncio.Task)] + for task in tasks: + await _drain(task) + if evidence is not None and ( + not tasks or any(task.cancelled() or task.exception() is not None for task in tasks) + ): + evidence["cleanup_failed"] = True + raise + except BaseException: + if evidence is not None: + evidence["cleanup_failed"] = True + raise + + cast(Any, worker.Runs).enter = staticmethod( + lambda *args, **kw: _ExitBoundary(cast(Any, original_enter)(*args, **kw), run=True) + ) + stream.aclosing = lambda iterator: _ExitBoundary(original_closing(iterator)) + stream.AsyncExitStack = ObservedStack + cast(Any, AsyncPregelLoop).__aexit__ = loop_exit + + async def persistent_cancellation_listener(queue, run_id, thread_id, done): + return await _persistent_cancellation_listener( + original_cancellation_listener, queue, run_id, thread_id, done + ) + + def quickjs_run_async(self, coro): + self._ensure_started() + remote = asyncio.run_coroutine_threadsafe(coro, self._loop) + return asyncio.create_task(_await_remote_future(remote)) + + inmem_ops.listen_for_cancellation = persistent_cancellation_listener + cast(Any, ThreadWorker).run_async = quickjs_run_async + + @wraps(original_to_thread) + async def owned_to_thread(func, /, *args, **kwargs): + if not _owned.get(): + return await original_to_thread(func, *args, **kwargs) + task = asyncio.create_task(original_to_thread(func, *args, **kwargs)) + try: + return await asyncio.shield(task) + except asyncio.CancelledError: + await _drain(task) + raise + + @wraps(original) + async def observed(run, attempt, main_loop, **kwargs): + key = (str(run["thread_id"]), str(run["run_id"])) + identity = (attempt, object()) + with _lock: + entry = _runs.get(key) + if entry is not None: + if entry["cancel_requested"]: + # Cancel closed admission before this queued attempt ran. + return None + entry["active"] += 1 + entry["exited"] = False + entry["attempts"][identity] = False + result = None + token = _owned.set(entry is not None) + evidence = {"run_exited": False, "cleanup_failed": False} + evidence_token = _evidence.set(evidence if entry is not None else None) + try: + result = await original(run, attempt, main_loop, **kwargs) + return result + finally: + _owned.reset(token) + _evidence.reset(evidence_token) + if entry is not None: + with _lock: + entry["active"] -= 1 + # Only this invocation can discharge its own evidence. + # Business error/timeout/retry is independent of cleanup. + clean = evidence["run_exited"] and not evidence["cleanup_failed"] + entry["attempts"][identity] = clean + entry["uncertain"] = not all(entry["attempts"].values()) + if clean: + entry["exited"] = True + entry["status"] = result["status"] if result else "retry" + + worker.worker = observed + asyncio.to_thread = owned_to_thread + _installed = True + + +def reserve(thread_id: str, run_id: str, request_hash: str) -> None: + with _lock: + _runs.setdefault((thread_id, run_id), { + "request_hash": request_hash, "active": 0, + "cancel_requested": False, "exited": False, + "uncertain": False, "status": "pending", "attempts": {}, + }) + + +def is_reserved(thread_id: str, run_id: str) -> bool: + with _lock: + return (thread_id, run_id) in _runs + + +def cancel_and_inspect(thread_id: str, run_id: str) -> dict: + with _lock: + entry = _runs.get((thread_id, run_id)) + if entry is None: + return {"run_id": run_id, "thread_id": thread_id, "execution_exited": False} + entry["cancel_requested"] = True + # No active worker plus closed admission is safe even for pending work. + # A worker that escaped abnormally keeps uncertainty latched. + stopped = entry["active"] == 0 and not entry["uncertain"] + return { + "run_id": run_id, "thread_id": thread_id, + "request_hash": entry["request_hash"], + "execution_exited": stopped, + "exit_kind": "worker_returned" if entry["exited"] else "admission_closed", + "status": entry["status"] if entry["exited"] else "interrupted", + } + + +async def wait_for_exit(thread_id: str, run_id: str, timeout: float = 3.0) -> dict: + deadline = asyncio.get_running_loop().time() + timeout + while True: + receipt = cancel_and_inspect(thread_id, run_id) + if receipt.get("execution_exited") is True: + return receipt + remaining = deadline - asyncio.get_running_loop().time() + if remaining <= 0: + return receipt + await asyncio.sleep(min(0.02, remaining)) \ No newline at end of file diff --git a/EvoScientist/middleware/provider_context.py b/EvoScientist/middleware/provider_context.py index 0e39106..46f5a29 100644 --- a/EvoScientist/middleware/provider_context.py +++ b/EvoScientist/middleware/provider_context.py @@ -66,7 +66,7 @@ def _provider_reference(block: Mapping[str, Any]) -> dict[str, str] | None: if ( block_type not in {"image", "audio", "video"} or not isinstance(path, str) - or not path.startswith(f"{_MEDIA_PREFIX}/") + or not path.startswith((f"{_MEDIA_PREFIX}/", f"/workspace{_MEDIA_PREFIX}/")) ): return None mime = str(block.get("mime_type") or f"{block_type}/unknown") @@ -193,16 +193,18 @@ class ProviderContextMediaMiddleware(AgentMiddleware): backend: Any, *, max_inline_media_bytes: int = _DEFAULT_MAX_INLINE_MEDIA_BYTES, + media_prefix: str = _MEDIA_PREFIX, ) -> None: self.backend = backend + self.media_prefix = media_prefix self.max_inline_media_bytes = max(1, int(max_inline_media_bytes)) - @staticmethod def _paths_for( + self, media: Mapping[str, tuple[str, bytes]], ) -> dict[str, tuple[str, str]]: return { - digest: (f"{_MEDIA_PREFIX}/{digest[:24]}.{_extension(mime)}", mime) + digest: (f"{self.media_prefix}/{digest[:24]}.{_extension(mime)}", mime) for digest, (mime, _raw) in media.items() } diff --git a/EvoScientist/middleware/recoverable_tools.py b/EvoScientist/middleware/recoverable_tools.py index 5d9d5c9..a18db65 100644 --- a/EvoScientist/middleware/recoverable_tools.py +++ b/EvoScientist/middleware/recoverable_tools.py @@ -19,8 +19,8 @@ try: except ImportError: # pragma: no cover - compatibility with older LangGraph _GraphInterrupt = None # type: ignore[assignment,misc] -from EvoScientist.llm.contracts import EvoRuntimeError from EvoScientist.internal_service import internal_service_headers +from EvoScientist.llm.contracts import EvoRuntimeError if TYPE_CHECKING: from langchain.agents.middleware.types import ToolCallRequest @@ -44,6 +44,15 @@ _IDEMPOTENT_NAMES = { logger = logging.getLogger(__name__) +def _callback_unavailable(exc: BaseException) -> bool: + if isinstance(exc, httpx.TransportError): + return True + return ( + isinstance(exc, httpx.HTTPStatusError) + and exc.response.status_code >= 500 + ) + + def _canonical(value: Any) -> str: return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"), default=str) @@ -213,7 +222,9 @@ class RecoverableToolEffectMiddleware(AgentMiddleware): "result": {}, }, ) - except httpx.TransportError: + except BaseException as callback_exc: + if not _callback_unavailable(callback_exc): + raise logger.exception( "tool-effect failure callback transport failed run_id=%s tool=%s effect_id=%s", proxy["run_id"], @@ -238,7 +249,9 @@ class RecoverableToolEffectMiddleware(AgentMiddleware): "result": payload, }, ) - except httpx.TransportError: + except BaseException as callback_exc: + if not _callback_unavailable(callback_exc): + raise logger.exception( "tool-effect terminal callback transport failed run_id=%s tool=%s effect_id=%s", proxy["run_id"], diff --git a/EvoScientist/native_sandbox.py b/EvoScientist/native_sandbox.py index 8a19ad0..39ed991 100644 --- a/EvoScientist/native_sandbox.py +++ b/EvoScientist/native_sandbox.py @@ -520,13 +520,17 @@ def _collect_process( except subprocess.TimeoutExpired: _terminate_process_group(process) return_code = process.wait(timeout=1) + except BaseException: + # Losing the output reader must not leave the command executing. + _terminate_process_group(process) + process.wait(timeout=1) + raise finally: + _kill_remaining_process_group(process) selector.close() for stream in streams: stream.close() - _kill_remaining_process_group(process) - if timed_out: return_code = 124 elif cancelled: @@ -735,10 +739,17 @@ class NativeWorkspaceBackend(ScopedFilesystemBackend, SandboxBackendProtocol): return await asyncio.shield(task) except asyncio.CancelledError: cancel_event.set() - try: - await asyncio.wait_for(asyncio.shield(task), timeout=2) - except (TimeoutError, asyncio.CancelledError): - pass + # A cancelled coroutine still owns its thread until cleanup finishes. + # Repeated cancellation must not orphan the worker either. + while not task.done(): + try: + await asyncio.shield(task) + except asyncio.CancelledError: + continue + except Exception: + break + if not task.cancelled(): + task.exception() raise diff --git a/EvoScientist/stream/emitter.py b/EvoScientist/stream/emitter.py index d9d3201..fdebc29 100644 --- a/EvoScientist/stream/emitter.py +++ b/EvoScientist/stream/emitter.py @@ -34,6 +34,19 @@ class StreamEventEmitter: "thinking", {"type": "thinking", "content": content, "id": thinking_id} ) + @staticmethod + def reasoning_summary(summary: str) -> StreamEvent: + """Provider-declared user-visible reasoning summary.""" + return StreamEvent( + "reasoning_summary", + { + "type": "reasoning_summary", + "summary": summary, + "visibility": "user_visible_summary", + "source_kind": "provider_summary", + }, + ) + @staticmethod def text(content: str) -> StreamEvent: """Text content event.""" diff --git a/EvoScientist/stream/events.py b/EvoScientist/stream/events.py index 2efd7df..e9fe213 100644 --- a/EvoScientist/stream/events.py +++ b/EvoScientist/stream/events.py @@ -34,6 +34,7 @@ from .v3_payloads import ( _event_data, _event_namespace, _reasoning_from_content, + _reasoning_summary_from_content, _split_message_event_data, _strip_legacy_thinking_tags, _text_from_content, @@ -570,6 +571,9 @@ class _V3EventProcessor: emitted_reasoning = bool(reasoning) content = msg.content + reasoning_summary = _reasoning_summary_from_content(content) + if reasoning_summary and subagent is None: + events.append(self.emitter.reasoning_summary(reasoning_summary).data) if not emitted_reasoning: events.extend( self._emit_thinking(_reasoning_from_content(content), subagent) @@ -961,6 +965,7 @@ async def stream_agent_events( callbacks: list[Any] | None = None, configurable: dict[str, Any] | None = None, error_mode: str = "emit", + stop_owner: Any = None, ) -> AsyncGenerator[dict[str, Any], None]: """Stream events from a DeepAgents/LangGraph v3 run. @@ -1036,6 +1041,8 @@ async def stream_agent_events( stream = ( await stream_result if inspect.isawaitable(stream_result) else stream_result ) + if stop_owner is not None: + stop_owner.attach(stream) async def _put_events(events: list[dict[str, Any]]) -> None: for event in events: @@ -1091,6 +1098,10 @@ async def stream_agent_events( if completion_tasks: await asyncio.gather(*completion_tasks) finally: + for task in completion_tasks: + if not task.done(): + task.cancel() + await asyncio.gather(*completion_tasks, return_exceptions=True) subagents.close() async def _run_producer(coro: Any) -> None: @@ -1146,18 +1157,23 @@ async def stream_agent_events( ).data raise finally: - if stream is not None: + abort_error = None + if stop_owner is not None: + await asyncio.shield(stop_owner.start(producers)) + elif stream is not None: try: result = stream.abort() if inspect.isawaitable(result): await result - except Exception: - pass + except Exception as exc: + abort_error = exc for task in producers: if not task.done(): task.cancel() if producers: await asyncio.gather(*producers, return_exceptions=True) + if abort_error is not None and config_values.get("checkpoint_writer_strict_close"): + raise RuntimeError("CHECKPOINT_WRITER_STOP_UNCONFIRMED") from abort_error # When the run ended with an exception the LangGraph checkpoint may be # left interrupted (``next`` non-empty). Clear it — unless it's a real # human-in-the-loop pause — so the next user message starts a fresh turn diff --git a/EvoScientist/stream/stop.py b/EvoScientist/stream/stop.py new file mode 100644 index 0000000..7ec6032 --- /dev/null +++ b/EvoScientist/stream/stop.py @@ -0,0 +1,120 @@ +"""Instance-local stop adapter for LangGraph 1.2.6, SQLite saver 3.0.3. + +No detached writers or foreign update_state/raw SQL writers are supported. +aiosqlite 0.22.1's FIFO commit is the connection barrier, not its saver lock. +""" +import asyncio +from importlib.metadata import version +from typing import Any + + +class StreamStop: + def __init__(self): + self.graph: Any = None + self.stream = None + self.pulls = set() + self.exits = set() + self.errors = [] + self.task = None + self.unsupported = False + self.exit_tasks_observed = 0 + self.unconfirmed_exit = asyncio.Event() + + def attach(self, stream): + from langgraph.stream.run_stream import AsyncGraphRunStream + if version("langgraph") != "1.2.6" or type(stream) is not AsyncGraphRunStream: + self.unsupported = True + raise RuntimeError("CHECKPOINT_STOP_ADAPTER_UNSUPPORTED") + self.stream = stream + self.graph = stream._graph_aiter + stream._graph_aiter = self + + def __aiter__(self): + return self + + def _record(self, exc): + self.exits.update(arg for arg in exc.args if isinstance(arg, asyncio.Future)) + + async def __anext__(self): + task = asyncio.current_task() + self.pulls.add(task) + try: + return await self.graph.__anext__() + except BaseException as exc: + self._record(exc) + raise + finally: + self.pulls.discard(task) + + async def _retry_close(self, close): + while True: + try: + await close() + return + except BaseException as exc: + self._record(exc) + self.errors.append(exc) + await self._settle_exits() + await asyncio.sleep(0.05) + + async def _settle_exits(self): + while self.exits: + tasks = tuple(self.exits) + self.exit_tasks_observed += len(tasks) + results = await asyncio.gather(*tasks, return_exceptions=True) + self.exits.difference_update(tasks) + for result in results: + if isinstance(result, BaseException): + self._record(result) + self.errors.append(result) + # A failed/cancelled stack exit is not executor settlement. + # There is no supported recovery handle for that case. + await self.unconfirmed_exit.wait() + + def start(self, producers): + if self.task is None: + self.task = asyncio.create_task(self._close(producers), name="evo-v3-stop") + return self.task + + async def _close(self, producers): + if self.unsupported: + await self.unconfirmed_exit.wait() + stream = self.stream + if stream is None: + return + # Freeze pumping before cancellation can discard the dependency's pull. + async with stream._pump_cond: + stream._exhausted = True + stream._aborting = True + pulls = set(self.pulls) + if stream._anext_task is not None: + pulls.add(stream._anext_task) + stream._pump_cond.notify_all() + for task in pulls: + if not task.done(): + task.cancel() + results = await asyncio.gather(*pulls, return_exceptions=True) + for result in results: + if isinstance(result, BaseException): + self._record(result) + await self._settle_exits() + await self._retry_close(self.graph.aclose) + await self._settle_exits() + await self._retry_close(stream._mux.aclose) + for task in producers: + if not task.done(): + task.cancel() + await asyncio.gather(*producers, return_exceptions=True) + + +async def sqlite_barrier(saver): + from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver + import aiosqlite + if (type(saver) is not AsyncSqliteSaver or type(saver.conn) is not aiosqlite.Connection + or version("langgraph-checkpoint-sqlite") != "3.0.3" + or version("aiosqlite") != "0.22.1"): + raise RuntimeError("CHECKPOINT_STOP_ADAPTER_UNSUPPORTED") + async with saver.lock: + # Worker executes even cancelled queued futures. Await a later commit + # after graph/executor/repair writers settle, covering those operations. + await saver.conn.commit() \ No newline at end of file diff --git a/EvoScientist/stream/v3_payloads.py b/EvoScientist/stream/v3_payloads.py index a62a667..911ae9c 100644 --- a/EvoScientist/stream/v3_payloads.py +++ b/EvoScientist/stream/v3_payloads.py @@ -82,3 +82,24 @@ def _reasoning_from_content(content: object) -> str: if isinstance(reasoning, str): parts.append(reasoning) return "".join(parts) + + +def _reasoning_summary_from_content(content: object) -> str: + if not isinstance(content, list): + return "" + parts: list[str] = [] + for block in content: + block_map = _as_raw_map(block) + if block_map is None or block_map.get("type") != "reasoning": + continue + summary = block_map.get("summary") + if not isinstance(summary, list): + continue + for part in summary: + part_map = _as_raw_map(part) + if part_map is None or part_map.get("type") != "summary_text": + continue + text = part_map.get("text") + if isinstance(text, str): + parts.append(text) + return "".join(parts) diff --git a/EvoScientist/web_checkpointer.py b/EvoScientist/web_checkpointer.py new file mode 100644 index 0000000..b380407 --- /dev/null +++ b/EvoScientist/web_checkpointer.py @@ -0,0 +1,256 @@ +"""Web-only PostgreSQL saver. CLI sessions remain in sessions.py. + +The saver also owns PostgreSQL turn fencing for Web-hosted agents. There is no +SQLite fallback. Schema setup is performed only by the reviewed Gateway and +LangGraph migration paths. +""" +from __future__ import annotations + +import asyncio +import hashlib +import os +from contextlib import asynccontextmanager +from dataclasses import dataclass + + +def _web_saver_type(): + from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver + + class WebPostgresSaver(AsyncPostgresSaver): + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self._turn_fence_lock = asyncio.Lock() + + @dataclass(frozen=True, slots=True) + class TurnLease: + thread_id: str + owner_id: str + fencing_token: int + expires_at_ms: int + checkpoint_snapshot_id: str + + @asynccontextmanager + async def _turn_lock(self, thread_id): + lock_key = str(thread_id or "") + async with self._turn_fence_lock: + async with self._cursor() as cursor: + await cursor.execute( + "SELECT pg_advisory_lock(hashtextextended(%s, 90421066))", + (lock_key,), + ) + try: + yield + finally: + async with self._cursor() as cursor: + await cursor.execute( + "SELECT pg_advisory_unlock(hashtextextended(%s, 90421066))", + (lock_key,), + ) + + async def acquire_turn_lease(self, thread_id, owner_id, *, ttl_seconds): + clean_thread = str(thread_id or "").strip() + clean_owner = str(owner_id or "").strip() + if not clean_thread or not clean_owner or ttl_seconds < 1: + raise ValueError("thread, owner and positive TTL are required") + async with self._turn_lock(clean_thread): + async with self._cursor(pipeline=True) as cursor: + await cursor.execute( + """INSERT INTO evomemory_checkpoint_turn_fences + (thread_id,generation,owner_id,expires_at,updated_at) + VALUES (%s,1,%s,NOW()+(%s*INTERVAL '1 second'),NOW()) + ON CONFLICT(thread_id) DO UPDATE SET + generation=evomemory_checkpoint_turn_fences.generation+1, + owner_id=EXCLUDED.owner_id, + expires_at=EXCLUDED.expires_at, + updated_at=NOW() + WHERE evomemory_checkpoint_turn_fences.owner_id IS NULL + OR evomemory_checkpoint_turn_fences.owner_id=EXCLUDED.owner_id + OR evomemory_checkpoint_turn_fences.expires_at=NOW() + RETURNING EXTRACT(EPOCH FROM expires_at)*1000 AS expires_at_ms""", + ( + ttl_seconds, + lease.thread_id, + lease.fencing_token, + lease.owner_id, + ), + ) + row = await cursor.fetchone() + if row is None: + raise RuntimeError("TURN_LEASE_LOST") + return self.TurnLease( + lease.thread_id, + lease.owner_id, + lease.fencing_token, + int(row["expires_at_ms"]), + lease.checkpoint_snapshot_id, + ) + + async def release_turn_lease(self, lease): + async with self._turn_lock(lease.thread_id): + async with self._cursor(pipeline=True) as cursor: + await cursor.execute( + """UPDATE evomemory_checkpoint_turn_fences + SET owner_id=NULL,expires_at=NULL,updated_at=NOW() + WHERE thread_id=%s AND generation=%s AND owner_id=%s + RETURNING 1""", + (lease.thread_id, lease.fencing_token, lease.owner_id), + ) + return await cursor.fetchone() is not None + + async def _require_write_lease(self, config): + configurable = dict(config.get("configurable") or {}) + thread_id = str(configurable.get("thread_id") or "") + if not thread_id.startswith("web:"): + return thread_id, 0 + owner_id = str(configurable.get("turn_lease_owner") or "") + token = int(configurable.get("turn_fencing_token") or 0) + async with self._cursor() as cursor: + await cursor.execute( + """SELECT 1 FROM evomemory_checkpoint_turn_fences + WHERE thread_id=%s AND generation=%s AND owner_id=%s + AND expires_at>=NOW()""", + (thread_id, token, owner_id), + ) + if await cursor.fetchone() is None: + raise RuntimeError("TURN_FENCED") + return thread_id, token + + async def aput(self, config, checkpoint, metadata, new_versions): + thread_id = str(config.get("configurable", {}).get("thread_id") or "") + async with self._turn_lock(thread_id): + thread_id, token = await self._require_write_lease(config) + result = await super().aput(config, checkpoint, metadata, new_versions) + if token: + checkpoint_id = str( + result.get("configurable", {}).get("checkpoint_id") or "" + ) + async with self._cursor(pipeline=True) as cursor: + await cursor.execute( + """INSERT INTO evomemory_checkpoint_versions + (thread_id,sequence,checkpoint_id,updated_at) + VALUES (%s,1,%s,NOW()) + ON CONFLICT(thread_id) DO UPDATE SET + sequence=evomemory_checkpoint_versions.sequence+1, + checkpoint_id=EXCLUDED.checkpoint_id, + updated_at=NOW()""", + (thread_id, checkpoint_id), + ) + return result + + async def aput_writes(self, config, writes, task_id, task_path=""): + thread_id = str(config.get("configurable", {}).get("thread_id") or "") + async with self._turn_lock(thread_id): + await self._require_write_lease(config) + await super().aput_writes(config, writes, task_id, task_path) + + async def adelete_for_runs(self, run_ids): + values = tuple(dict.fromkeys(str(run_id) for run_id in run_ids)) + if not values: + return + async with self._cursor(pipeline=True) as cursor: + await cursor.execute( + """DELETE FROM checkpoint_writes w USING checkpoints c + WHERE w.thread_id=c.thread_id + AND w.checkpoint_ns=c.checkpoint_ns + AND w.checkpoint_id=c.checkpoint_id + AND c.metadata->>'run_id'=ANY(%s)""", + (list(values),), + ) + await cursor.execute( + "DELETE FROM checkpoints WHERE metadata->>'run_id'=ANY(%s)", + (list(values),), + ) + + return WebPostgresSaver + + +@asynccontextmanager +async def create_web_checkpointer(): + from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer + from psycopg import AsyncConnection + from psycopg.rows import dict_row + + dsn = os.environ.get("EVOSCIENTIST_WEB_CHECKPOINT_DSN", "") + if not dsn.startswith(("postgresql://", "postgres://")): + raise RuntimeError("Web checkpointer requires an explicit PostgreSQL DSN") + serde = JsonPlusSerializer( + allowed_msgpack_modules=frozenset( + { + ("EvoScientist.llm.errors", "AgentControlError"), + ("EvoScientist.llm.errors", "ModelToolProtocolError"), + ("EvoScientist.llm.errors", "ProviderStreamError"), + } + ) + ) + try: + conn = await AsyncConnection.connect( + dsn, + autocommit=True, + prepare_threshold=0, + row_factory=dict_row, + connect_timeout=5, + options="-c search_path=public", + ) + except Exception: + raise RuntimeError("Web checkpoint PostgreSQL connection failed") from None + async with conn: + saver = _web_saver_type()(conn, serde=serde) + try: + async with conn.cursor() as cursor: + await cursor.execute("SELECT v FROM checkpoint_migrations ORDER BY v") + versions = [row["v"] for row in await cursor.fetchall()] + if versions != list(range(len(saver.MIGRATIONS))): + raise RuntimeError("migration version mismatch") + for table in ( + "checkpoints", + "checkpoint_blobs", + "checkpoint_writes", + "evomemory_checkpoint_turn_fences", + "evomemory_checkpoint_versions", + ): + await cursor.execute(f"SELECT 1 FROM {table} LIMIT 0") + except Exception: + raise RuntimeError( + "Web checkpoint schema is not ready; run reviewed setup separately" + ) from None + yield saver diff --git a/EvoScientist/web_runtime.py b/EvoScientist/web_runtime.py index 09dbe38..642d8b6 100644 --- a/EvoScientist/web_runtime.py +++ b/EvoScientist/web_runtime.py @@ -142,15 +142,15 @@ class _ToolRegistryFenceMiddleware(AgentMiddleware): def create_web_agent( - *, snapshot: Any, host: WebHostContext, model_set: AgentModelSet + *, snapshot: Any, host: WebHostContext, model_set: AgentModelSet, + config: Any = None, ) -> Any: """Create a `web_v3` agent without importing Gateway types.""" from .EvoScientist import create_cli_agent from .middleware.evo_route_fallback import EvoRouteFallbackMiddleware - config = copy.copy(load_config()) - config.auto_approve = True + config = copy.copy(config if config is not None else load_config()) config.auto_mode = True config.enable_ask_user = False config.enable_async_subagents = False diff --git a/start-langgraph.sh b/start-langgraph.sh index ec48527..546127d 100755 --- a/start-langgraph.sh +++ b/start-langgraph.sh @@ -3,15 +3,17 @@ set -euo pipefail PROJECT_DIR="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd)" -LANGGRAPH_CONFIG="${PROJECT_DIR}/EvoScientist/langgraph_dev/langgraph.json" +LANGGRAPH_CONFIG="${PROJECT_DIR}/EvoScientist/langgraph_dev/langgraph.web.json" HOST="${EVOSCIENTIST_LANGGRAPH_HOST:-127.0.0.1}" PORT="${EVOSCIENTIST_LANGGRAPH_DEV_PORT:-3076}" WEB_ENV="${PROJECT_DIR}/../Ai4Sci-Web/.env" +WEB_PROJECT="${PROJECT_DIR}/../Ai4Sci-Web" N_JOBS="${EVOSCIENTIST_LANGGRAPH_JOBS_PER_WORKER:-6}" -if [[ ! -x "${PROJECT_DIR}/.venv/bin/langgraph" ]]; then - echo "LangGraph executable not found: ${PROJECT_DIR}/.venv/bin/langgraph" >&2 - echo "Run 'uv sync' in ${PROJECT_DIR} first." >&2 +WEB_VENV="${PROJECT_DIR}/../Ai4Sci-Web/.venv" +if [[ ! -x "${WEB_VENV}/bin/langgraph" ]]; then + echo "Web LangGraph executable not found: ${WEB_VENV}/bin/langgraph" >&2 + echo "Run 'uv sync --group host-runtime' in Ai4Sci-Web first." >&2 exit 1 fi @@ -40,8 +42,8 @@ fi # The Web Gateway always supplies a verified conversation workspace scope. export EVOSCIENTIST_DEPLOY_MODE="${EVOSCIENTIST_DEPLOY_MODE:-full}" export EVOSCIENTIST_WORKSPACE_DIR="${EVOSCIENTIST_WORKSPACE_DIR:-${PROJECT_DIR}/../.ai4sci/workspace}" - -exec uv run --env-file "${WEB_ENV}" langgraph dev \ +exec uv run --project "${WEB_PROJECT}" --group host-runtime \ + --env-file "${WEB_ENV}" langgraph dev \ --config "${LANGGRAPH_CONFIG}" \ --host "${HOST}" \ --port "${PORT}" \