From a13904185d16128722b1bee3a68b76b2f9f08df9 Mon Sep 17 00:00:00 2001 From: Xi Zhang <106144707+X-iZhang@users.noreply.github.com> Date: Sun, 31 May 2026 15:11:25 +0100 Subject: [PATCH] Feat/sandbox execute timeout (#243) * feat: implement configurable sandbox execute timeout and enhance recovery instructions * feat: add background process management tools and middleware for sandbox execution * feat: enhance background process management with completion notifications and deduplication * feat: enhance sandbox execution timeout validation and update related messages * feat: enhance background process management with thread-specific completion notifications and HITL approval handling * test: assert completion notification waits for process finish timestamp --- EvoScientist/EvoScientist.py | 25 ++- EvoScientist/backends.py | 78 ++++--- EvoScientist/background.py | 303 ++++++++++++++++++++++++++ EvoScientist/channels/consumer.py | 26 ++- EvoScientist/cli/async_notifier.py | 85 +++++--- EvoScientist/cli/commands.py | 2 +- EvoScientist/config/settings.py | 34 +++ EvoScientist/middleware/background.py | 134 ++++++++++++ EvoScientist/prompts.py | 24 +- EvoScientist/stream/display.py | 6 +- tests/test_async_notifier.py | 14 +- tests/test_backends.py | 8 + tests/test_background.py | 152 +++++++++++++ tests/test_background_middleware.py | 274 +++++++++++++++++++++++ tests/test_config.py | 46 ++++ tests/test_hitl.py | 76 +++++++ tests/test_prompts.py | 4 +- 17 files changed, 1202 insertions(+), 89 deletions(-) create mode 100644 EvoScientist/background.py create mode 100644 EvoScientist/middleware/background.py create mode 100644 tests/test_background.py create mode 100644 tests/test_background_middleware.py diff --git a/EvoScientist/EvoScientist.py b/EvoScientist/EvoScientist.py index 954f8d1..9741259 100644 --- a/EvoScientist/EvoScientist.py +++ b/EvoScientist/EvoScientist.py @@ -419,6 +419,7 @@ def _get_default_backend(): from .backends import CustomSandboxBackend, MergedSkillsBackend + cfg = _ensure_config() workspace_dir = str(_paths_mod.WORKSPACE_ROOT) set_active_workspace(workspace_dir) memory_dir = str(_paths_mod.MEMORIES_DIR) @@ -428,7 +429,7 @@ def _get_default_backend(): ws_backend = CustomSandboxBackend( root_dir=workspace_dir, virtual_mode=True, - timeout=300, + timeout=cfg.sandbox_execute_timeout, ) sk_backend = MergedSkillsBackend( primary_dir=user_skills_dir, @@ -499,6 +500,14 @@ def _get_default_middleware(*, for_async_subagent: bool = False): mw.insert(0, AskUserMiddleware()) + # Background-process tools (run_in_background / check_process / stop_process / + # list_processes) — main agent only. Async sub-agents run on langgraph-dev and + # must not spawn local OS processes. + if not for_async_subagent: + from .middleware.background import BackgroundExecutionMiddleware + + mw.append(BackgroundExecutionMiddleware()) + mw.append( create_code_interpreter_middleware( timeout=cfg.code_interpreter_timeout, @@ -544,7 +553,11 @@ def _get_default_agent(): # breaks parallel execute calls (multi-pending-interrupt LangGraph # error). See PR #202. if not cfg.auto_approve: - mw.append(HumanInTheLoopMiddleware(interrupt_on={"execute": True})) + mw.append( + HumanInTheLoopMiddleware( + interrupt_on={"execute": True, "run_in_background": True} + ) + ) if os.environ.get("EVOSCIENTIST_DEPLOY_MODE", "").lower() == "stripped": kwargs = _build_base_kwargs(be, mw) @@ -634,7 +647,7 @@ def create_cli_agent( ws_backend = CustomSandboxBackend( root_dir=workspace_dir, virtual_mode=True, - timeout=300, + timeout=cfg.sandbox_execute_timeout, ) sk_backend = MergedSkillsBackend( primary_dir=_usr_skills_dir, @@ -662,7 +675,11 @@ def create_cli_agent( # would propagate it to every subagent, breaking parallel execute calls # (multi-pending-interrupt LangGraph error). if not cfg.auto_approve: - mw.append(HumanInTheLoopMiddleware(interrupt_on={"execute": True})) + mw.append( + HumanInTheLoopMiddleware( + interrupt_on={"execute": True, "run_in_background": True} + ) + ) # Re-load MCP tools from current config (picks up /mcp add changes) kwargs = load_mcp_and_build_kwargs(be, mw, on_mcp_progress=on_mcp_progress) diff --git a/EvoScientist/backends.py b/EvoScientist/backends.py index 444d242..d1e1471 100644 --- a/EvoScientist/backends.py +++ b/EvoScientist/backends.py @@ -517,6 +517,40 @@ class MergedSkillsBackend(BackendProtocol): return self._primary.upload_files(files) +def prepare_sandbox_command( + command: str, cwd: str | Path, *, virtual_mode: bool = True +) -> tuple[str, str | None]: + """Normalize workspace paths in ``command`` and validate it for the sandbox. + + Shared by :meth:`CustomSandboxBackend.execute` and the background-process tools so + both enforce *identical* workspace-path rewriting (so virtual ``/`` paths resolve to + the workspace, not the host root) and the same command validation. + + Returns ``(prepared_command, error)``: ``error`` is a message string when the command + is rejected (the caller must NOT run it), otherwise ``None``. + """ + cwd_str = str(cwd).rstrip("/") + # Replace literal workspace-root absolute paths with ./ BEFORE validation, so + # workspace paths are sanitized before the system-path check fires. + ws = cwd_str + "/" + if ws in command: + command = command.replace(ws, "./") + if virtual_mode: + command = convert_virtual_paths_in_command( + command=command, + workspace_name=Path(cwd_str).name, + ) + # Skills/memory dirs must be allowlisted: the workspace-literal replace above runs + # before the resolver, so any absolute path it later injects reaches validate unstripped. + allow_prefixes = ( + str(paths.USER_SKILLS_DIR), + str(paths.GLOBAL_SKILLS_DIR), + str(paths.MEMORIES_DIR), + str(_BUILTIN_SKILLS_DIR), + ) + return command, validate_command(command, allow_prefixes=allow_prefixes) + + class CustomSandboxBackend(LocalShellBackend): """ Custom sandbox backend - inherits LocalShellBackend with added safety. @@ -620,36 +654,11 @@ class CustomSandboxBackend(LocalShellBackend): Then delegates to LocalShellBackend.execute() for actual execution. """ - # Replace literal workspace-root absolute paths with ./ - # Must happen BEFORE validation so workspace paths (e.g. /tmp/...) - # are sanitized before the system-path check fires. - ws = str(self.cwd).rstrip("/") + "/" - if ws in command: - command = command.replace(ws, "./") - - # Convert virtual paths to relative paths - if self.virtual_mode: - command = convert_virtual_paths_in_command( - command=command, - workspace_name=Path(str(self.cwd)).name, - ) - - # USER_SKILLS_DIR must be in the allowlist: the workspace-literal - # replace above runs BEFORE the resolver, so any absolute path the - # resolver later injects reaches validate_command unstripped. - allow_prefixes = ( - str(paths.USER_SKILLS_DIR), - str(paths.GLOBAL_SKILLS_DIR), - str(paths.MEMORIES_DIR), - str(_BUILTIN_SKILLS_DIR), + command, error = prepare_sandbox_command( + command, self.cwd, virtual_mode=self.virtual_mode ) - error = validate_command(command, allow_prefixes=allow_prefixes) if error: - return ExecuteResponse( - output=error, - exit_code=1, - truncated=False, - ) + return ExecuteResponse(output=error, exit_code=1, truncated=False) # Delegate to parent for subprocess execution response = super().execute(command, timeout=timeout) @@ -658,14 +667,17 @@ class CustomSandboxBackend(LocalShellBackend): if response.exit_code == 124: cmd_words = command.split() grep_hint = cmd_words[0] if cmd_words else "process" - bg_cmd = f"{command} > /output.log 2>&1 &" + bg_cmd = f'{command} > /output.log 2>&1 & echo "PID: $!"' response = ExecuteResponse( output=( f"{response.output}\n\n" - f"Recovery: re-run in background to avoid the sandbox timeout:\n" - f" {bg_cmd}\n" - f"Then check progress: ps aux | grep {grep_hint}\n" - f"Read results: cat /output.log" + f"Recovery — pick one:\n" + f" 1. Needs more time? Re-run with a larger timeout (up to 3600s): " + f"execute(command=..., timeout=600)\n" + f" 2. Runs indefinitely? Run it in the background and keep the PID:\n" + f" {bg_cmd}\n" + f" Check: ps -p (or: ps aux | grep {grep_hint}) · " + f"Read: cat /output.log · Stop: kill " ), exit_code=response.exit_code, truncated=response.truncated, diff --git a/EvoScientist/background.py b/EvoScientist/background.py new file mode 100644 index 0000000..9ead364 --- /dev/null +++ b/EvoScientist/background.py @@ -0,0 +1,303 @@ +"""Background OS-process execution for the sandbox. + +A *process* here is a single detached OS process launched via ``run_in_background`` +(distinct from an async sub-agent *task* and a future cron *schedule* — the word +"job" is intentionally never used). + +The registry is **module-global (process-level)**: processes survive ``/new`` and +``/resume`` within the same CLI process, but are not persisted across a CLI restart. +The live ``Popen`` handle is held so ``poll()`` / ``returncode`` stay authoritative +(no PID-reuse risk). + +Command validation and cwd resolution happen at the tool layer +(``middleware/background.py``); this module is the pure execution + tracking mechanism +and is safe to unit-test on its own. A future scheduler (cron) would reuse ``launch``. +""" + +from __future__ import annotations + +import logging +import os +import signal +import subprocess +import threading +import time +import uuid +from collections.abc import Callable +from dataclasses import dataclass, field +from datetime import UTC, datetime +from pathlib import Path + +logger = logging.getLogger(__name__) + +_BG_DIRNAME = ".bg_processes" +_KILL_GRACE_SECONDS = 2.0 + + +@dataclass +class BgProcess: + """A tracked background OS process.""" + + process_id: str + name: str + command: str + popen: subprocess.Popen + pid: int + log_path: Path + started_at: str # ISO-8601 UTC (record/display) + started_ts: float # epoch seconds (elapsed computation) + origin_thread_id: str | None = None # CLI thread/session that launched it + returncode: int | None = None + finished_at: str | None = None + finished_ts: float | None = None # epoch at exit; freezes elapsed once done + stopped: bool = False # set by stop(); suppresses the completion notification + # epoch each thread last checked this process (status/list); keyed by thread_id + # so a check from one session can't dedup another session's completion ping. + last_checked_by_thread: dict[str | None, float] = field(default_factory=dict) + + +_PROCESSES: dict[str, BgProcess] = {} +_LOCK = threading.Lock() + + +def _now_iso() -> str: + return datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%SZ") + + +def _record_exit(proc: BgProcess) -> None: + """Record terminal state on first observed exit. Caller MUST hold ``_LOCK``. + + ``finished_ts`` is set when the exit is first observed. The per-process daemon + watcher (:func:`_watch`) calls this right after ``popen.wait()`` returns, so in + practice ``finished_ts`` ≈ the real exit time. Calls from ``status`` / ``list_all`` / + ``stop`` are a fallback for the brief window before the watcher runs. + """ + rc = proc.popen.poll() + if rc is not None and proc.returncode is None: + proc.returncode = rc + proc.finished_at = _now_iso() + proc.finished_ts = time.time() + + +def _elapsed(proc: BgProcess) -> int: + """Seconds the process has run — frozen at first-observed exit once it has exited.""" + end = proc.finished_ts if proc.finished_ts is not None else time.time() + return int(end - proc.started_ts) + + +def was_observed_done(process_id: str, origin_thread_id: str | None = None) -> bool: + """True if ``origin_thread_id`` already saw this process's completion itself. + + i.e. the process has exited AND was checked (``status``/``list_all``) from that thread + at or after it finished. Used to dedup the completion notification (routed to the + launching thread), so a check from a *different* session can't suppress it. + """ + with _LOCK: + proc = _PROCESSES.get(process_id) + if proc is None or proc.finished_ts is None: + return False + seen_ts = proc.last_checked_by_thread.get(origin_thread_id) + return seen_ts is not None and seen_ts >= proc.finished_ts + + +def _read_tail(log_path: Path, tail_bytes: int) -> str: + # Seek from the end so a huge log isn't fully read into memory on each status check. + try: + with log_path.open("rb") as f: + f.seek(0, os.SEEK_END) + size = f.tell() + if size == 0: + return "(no output yet)" + if size > tail_bytes: + f.seek(-tail_bytes, os.SEEK_END) + return "...(truncated)...\n" + f.read().decode("utf-8", "replace") + f.seek(0) + data = f.read() + except OSError: + return "(no output captured yet)" + return data.decode("utf-8", "replace") + + +def _watch(proc: BgProcess, on_exit: Callable[[BgProcess], None] | None) -> None: + """Block until ``proc`` exits, record the exit promptly, then fire ``on_exit``. + + Running in a daemon thread, ``popen.wait()`` lets us record ``finished_ts`` at (very + close to) the real exit time — fixing the observation-time inflation — and gives a + hook the CLI layer wires to a completion notification, without ``background.py`` + importing the notifier (kept decoupled via the callback). + """ + try: + proc.popen.wait() + except Exception: + pass + with _LOCK: + _record_exit(proc) + if on_exit is not None: + try: + on_exit(proc) + except Exception: + logger.warning("background on_exit callback failed", exc_info=True) + + +def launch( + command: str, + cwd: str, + name: str | None = None, + *, + origin_thread_id: str | None = None, + on_exit: Callable[[BgProcess], None] | None = None, +) -> str: + """Launch ``command`` detached in ``cwd``; return a short ``process_id``. + + The command is run via ``shell=True`` with output redirected to a per-process log + file under ``/.bg_processes/`` and ``start_new_session=True`` so the child is a + process-group leader (survives this call's return and can be killed as a group). + The caller is responsible for validating ``command`` first. + + ``origin_thread_id`` records the launching CLI session so ``list_all`` can scope to it. + ``on_exit`` (optional) is called with the ``BgProcess`` from a daemon watcher thread + once the process exits — used by the CLI layer to emit a completion notification. + """ + process_id = uuid.uuid4().hex[:8] + log_dir = Path(cwd) / _BG_DIRNAME + log_dir.mkdir(parents=True, exist_ok=True) + log_path = log_dir / f"{process_id}.log" + + log_file = open(log_path, "w") + try: + popen = subprocess.Popen( + command, + shell=True, + cwd=cwd, + stdout=log_file, + stderr=subprocess.STDOUT, + stdin=subprocess.DEVNULL, + start_new_session=True, + ) + finally: + # The child inherited its own dup of the fd during spawn; the parent's copy + # is no longer needed (and must be closed so the pipe/file isn't held open). + log_file.close() + + proc = BgProcess( + process_id=process_id, + name=name or command[:40], + command=command, + popen=popen, + pid=popen.pid, + log_path=log_path, + started_at=_now_iso(), + started_ts=time.time(), + origin_thread_id=origin_thread_id, + ) + with _LOCK: + _PROCESSES[process_id] = proc + # Daemon watcher: records the precise exit time and fires on_exit when done. + threading.Thread(target=_watch, args=(proc, on_exit), daemon=True).start() + return process_id + + +def status( + process_id: str, *, thread_id: str | None = None, tail_bytes: int = 16_000 +) -> str: + """Return a human-readable status + recent output tail for ``process_id``.""" + with _LOCK: + proc = _PROCESSES.get(process_id) + if proc is None: + return ( + f"No such background process: {process_id!r}. " + "Use list_processes to see tracked processes." + ) + _record_exit(proc) + proc.last_checked_by_thread[thread_id] = time.time() # this thread observed it + running = proc.returncode is None + elapsed = _elapsed(proc) + name, pid, command, returncode, log_path = ( + proc.name, + proc.pid, + proc.command, + proc.returncode, + proc.log_path, + ) + if running: + head = f"Process {process_id} (name={name!r}) RUNNING — {elapsed}s elapsed, pid {pid}." + else: + head = f"Process {process_id} (name={name!r}) EXITED code {returncode} after ~{elapsed}s." + tail = _read_tail(log_path, tail_bytes) # file IO outside the lock + return ( + f"{head}\nCommand: {command}\n--- output (last {tail_bytes} bytes) ---\n{tail}" + ) + + +def stop(process_id: str) -> str: + """Terminate ``process_id`` and its process group (SIGTERM, then SIGKILL).""" + with _LOCK: + proc = _PROCESSES.get(process_id) + if proc is None: + return f"No such background process: {process_id!r}." + if proc.popen.poll() is not None: + _record_exit(proc) + return f"Process {process_id} already finished (code {proc.returncode})." + # Mark as user-stopped so the watcher's on_exit suppresses the completion + # notification (the user already knows — no need to ping them). + proc.stopped = True + # The watcher's popen.wait() reaps without the lock, so a tiny PID-reuse race + # remains (getpgid on a recycled pid). ProcessLookupError handles the common case; + # the window is too narrow to be worth coordinating the watcher. + try: + os.killpg(os.getpgid(proc.pid), signal.SIGTERM) + except ProcessLookupError: + _record_exit(proc) + return f"Process {process_id} is no longer running." + + deadline = time.time() + _KILL_GRACE_SECONDS + while time.time() < deadline: + with _LOCK: + if proc.popen.poll() is not None: + _record_exit(proc) + break + time.sleep(0.1) + else: + with _LOCK: + if proc.popen.poll() is None: + try: + os.killpg(os.getpgid(proc.pid), signal.SIGKILL) + except ProcessLookupError: + pass + _record_exit(proc) + + with _LOCK: + _record_exit(proc) + name = proc.name + return f"Stopped background process {process_id} (name={name!r})." + + +def list_all(thread_id: str | None = None, *, include_all: bool = False) -> str: + """List tracked background processes with live statuses. + + Scoped to the launching session (``thread_id``) unless ``include_all`` is set. + """ + with _LOCK: + all_procs = list(_PROCESSES.values()) + procs = ( + all_procs + if include_all + else [p for p in all_procs if p.origin_thread_id == thread_id] + ) + if not procs: + if all_procs and not include_all: + return ( + "No background processes in this session " + f"({len(all_procs)} in other sessions — pass all_threads=True to see them)." + ) + return "No background processes tracked." + lines = [] + now = time.time() + for p in procs: + _record_exit(p) + p.last_checked_by_thread[thread_id] = now # this thread observed it + state = "RUNNING" if p.returncode is None else f"exited({p.returncode})" + lines.append( + f" {p.process_id} {state:12} {_elapsed(p)}s name={p.name!r}" + ) + return f"{len(procs)} background process(es):\n" + "\n".join(lines) diff --git a/EvoScientist/channels/consumer.py b/EvoScientist/channels/consumer.py index f593f0c..cd9438d 100644 --- a/EvoScientist/channels/consumer.py +++ b/EvoScientist/channels/consumer.py @@ -118,7 +118,7 @@ def _should_auto_approve(action_requests: list[dict]) -> bool: return True try: - from ..config.settings import load_config + from ..config.settings import HITL_SHELL_TOOLS, load_config cfg = load_config() except Exception: @@ -137,7 +137,7 @@ def _should_auto_approve(action_requests: list[dict]) -> bool: name = ( req.get("name", "") if isinstance(req, dict) else getattr(req, "name", "") ) - if name != "execute": + if name not in HITL_SHELL_TOOLS: continue args = ( req.get("args", {}) if isinstance(req, dict) else getattr(req, "args", {}) @@ -657,18 +657,34 @@ class InboundConsumer: ) self._pending_interrupts[session_key] = pending + timed_out = False try: await asyncio.wait_for( pending.event.wait(), timeout=_HITL_APPROVAL_TIMEOUT, ) except TimeoutError: - # Auto-approve on timeout - pending.decision = "approve" + timed_out = True finally: + # Unregister BEFORE any further await so a late reply can't flip + # the decision back to approve during the notification round-trip. self._pending_interrupts.pop(session_key, None) - decision = pending.decision or "approve" + if timed_out: + # Reject on timeout (fail-closed; matches cli/channel.py). Decision + # is a local constant, not pending.decision, so it can't be + # overwritten by a late reply after we unregistered above. + decision = "reject" + await self.bus.publish_outbound( + OutboundMessage( + channel=msg.channel, + chat_id=msg.chat_id, + content="⏰ Approval timed out. Action rejected.", + metadata=msg.metadata, + ) + ) + else: + decision = pending.decision or "reject" # Visible confirmation so the click/reply registers (QQ has no # message recall API for C2C). Only fires when the user diff --git a/EvoScientist/cli/async_notifier.py b/EvoScientist/cli/async_notifier.py index 19d705f..d6cebff 100644 --- a/EvoScientist/cli/async_notifier.py +++ b/EvoScientist/cli/async_notifier.py @@ -40,6 +40,7 @@ class AsyncTaskNotification: status: str # one of TERMINAL_STATUSES received_at: str # ISO-8601 UTC timestamp prompt: str = "" # original task description sent to the sub-agent + kind: str = "agent" # "agent" (sub-agent) | "bg-process" (background shell) # The CLI/main-agent thread_id under which the watcher was spawned. Used # to route the notification back to the originating CLI session so a # /new between launch and completion does not inject the synthetic @@ -374,10 +375,20 @@ def dedup_notifications( so lexicographic comparison is correct). Also skip if `last_checked_at` is empty (brand-new task where agent hasn't checked yet). """ - if not async_tasks: - return notifs + from .. import background # cli -> core import; lazy to avoid import-order issues + + async_tasks = async_tasks or {} survivors: list[AsyncTaskNotification] = [] for n in notifs: + if n.kind == "bg-process": + # Background process: skip if the launching session already inspected it + # after it finished (check_process / list_processes) — mirrors the task + # dedup below. Per-thread: another session's check doesn't suppress this. + if background.was_observed_done(n.task_id, n.origin_cli_thread_id): + logger.debug("Dedup: skipping shell notification for %s", n.task_id) + continue + survivors.append(n) + continue task = async_tasks.get(n.task_id) if ( task @@ -393,30 +404,21 @@ def dedup_notifications( return survivors -def format_notification_lines( - notifs: list[AsyncTaskNotification], +def _render_notification_group( + notifs: list[AsyncTaskNotification], title: str, label: str ) -> list[tuple[str, str]]: - """Render notifications as compact tool-result-style lines for screen display. + """Render one group of notifications inside a titled open-right frame. - Returns a list of (text, rich_style) tuples — one per notification. - Used by both Rich CLI (console.print) and TUI (_append_system). - The LLM still receives the full format_batch_message text; this is - purely a visual representation for the human operator. + Open-right compact frame; bottom matches the top's width: + ╭── ✦ Agent Teams ✦ ──── + ✔ writing Task: ... success + ╰───────────────────────── """ - if not notifs: - return [] - # Open-right compact frame: short symmetric dashes around the title. - # Bottom matches the top's width so the visual is balanced. - # ╭── ✦ Agent Teams ✦ ── - # ✔ writing Task: ... success - # ╰───────────────────── - title = " ✦ Agent Teams ✦ " top_divider = "╭──" + title + "────" # 4 dashes on the right (2x of left) bottom_divider = "╰" + "─" * (len(top_divider) - 1) lines: list[tuple[str, str]] = [(top_divider, "dim")] for n in notifs: - # Strip the "-agent" suffix so it doesn't redundantly echo the header. - # `writing-agent` → `writing`, `data-analysis-agent` → `data-analysis`. + # `writing-agent` → `writing`. name = n.agent_name.removesuffix("-agent") if n.status == "success": icon, color = "✔", "#e67e22" # carrot orange (CSS hex; Rich+Textual) @@ -424,14 +426,12 @@ def format_notification_lines( icon, color = "✗", "red" else: # cancelled, timeout, interrupted icon, color = "⚠", "yellow" - # Body format (5-space indent under "Agent:" header): - # ✔ writing Task: success - # Collapse newlines, truncate prompt to 60 chars. + # Collapse newlines, truncate prompt/command preview to 60 chars. prompt_preview = (n.prompt or "").replace("\n", " ").strip() if len(prompt_preview) > 60: prompt_preview = prompt_preview[:60] + "…" if prompt_preview: - text = f" {icon} {name:18s} Task: {prompt_preview} {n.status}" + text = f" {icon} {name:18s} {label}: {prompt_preview} {n.status}" else: # Fallback: short task_id when no prompt is available short_tid = ( @@ -445,6 +445,28 @@ def format_notification_lines( return lines +def format_notification_lines( + notifs: list[AsyncTaskNotification], +) -> list[tuple[str, str]]: + """Render notifications as compact tool-result-style lines for screen display. + + Async sub-agents and background processes get SEPARATE titled frames so a shell + background process is never mislabeled as an "Agent Team". Returns (text, rich_style) + tuples. The LLM still receives the full ``format_batch_message`` text; this is purely + the visual representation for the human operator. + """ + if not notifs: + return [] + tasks = [n for n in notifs if n.kind != "bg-process"] + shell = [n for n in notifs if n.kind == "bg-process"] + lines: list[tuple[str, str]] = [] + if tasks: + lines += _render_notification_group(tasks, " ✦ Agent Teams ✦ ", "Task") + if shell: + lines += _render_notification_group(shell, " ✦ Background ✦ ", "Cmd") + return lines + + def format_batch_message(notifs: list[AsyncTaskNotification]) -> str: """Compose the synthetic user message that wakes the supervisor. @@ -459,13 +481,24 @@ def format_batch_message(notifs: list[AsyncTaskNotification]) -> str: for n in notifs: lines.append( json.dumps( - {"agent": n.agent_name, "status": n.status, "task_id": n.task_id}, + { + "agent": n.agent_name, + "kind": n.kind, + "status": n.status, + "task_id": n.task_id, + }, ensure_ascii=False, ) ) + # bg-process is inspected with check_process; sub-agents with check_async_task. + hints: list[str] = [] + if any(n.kind != "bg-process" for n in notifs): + hints.append("check_async_task (sub-agents)") + if any(n.kind == "bg-process" for n in notifs): + hints.append("check_process (background processes)") lines.append( - "(Signal only — fetch via check_async_task if relevant to current step, " - "else acknowledge & continue.)" + f"(Signal only — fetch full result via {' or '.join(hints)} if relevant to " + "the current step, else acknowledge & continue.)" ) return "\n".join(lines) diff --git a/EvoScientist/cli/commands.py b/EvoScientist/cli/commands.py index cb5477c..172eff1 100644 --- a/EvoScientist/cli/commands.py +++ b/EvoScientist/cli/commands.py @@ -1418,7 +1418,7 @@ def config_set( if set_config_value(key, value): console.print(f"[green]Set {escape(key)}[/green]") else: - console.print(f"[red]Invalid key: {escape(key)}[/red]") + console.print(f"[red]Could not set {escape(key)}: invalid key or value[/red]") raise typer.Exit(1) diff --git a/EvoScientist/config/settings.py b/EvoScientist/config/settings.py index 9b4066c..b421c6c 100644 --- a/EvoScientist/config/settings.py +++ b/EvoScientist/config/settings.py @@ -7,6 +7,7 @@ with the following priority (highest to lowest): from __future__ import annotations +import logging import os from dataclasses import asdict, dataclass, fields from pathlib import Path @@ -15,6 +16,12 @@ from typing import Any, Literal import yaml from dotenv import find_dotenv, load_dotenv +# Tools that run shell commands and need manual HITL approval (subject to +# shell_allow_list). Single source of truth for every interrupt consumer +# (stream/display.py, channels/consumer.py) — keep aligned with the agent's +# `interrupt_on` set in EvoScientist.py. +HITL_SHELL_TOOLS = ("execute", "run_in_background") + # ============================================================================= # Configuration paths # ============================================================================= @@ -268,6 +275,11 @@ class EvoScientistConfig: code_interpreter_timeout: float = 60.0 # seconds per JS eval code_interpreter_max_result_chars: int = 10000 # truncate large JSON results + # Default per-command timeout (seconds) for the sandbox `execute` tool. + # Only the default — the agent can still override per command up to the + # deepagents max_execute_timeout cap (3600s). + sandbox_execute_timeout: int = 300 + # Checkpoint pruning (sessions.db retention per (thread_id, checkpoint_ns)) # Safety net for runaway conversations. Under DeltaChannel (deepagents 0.6+) # normal usage produces linear growth, so this default is set well above @@ -292,6 +304,19 @@ class EvoScientistConfig: stt_device: str = "cpu" # "cpu" | "cuda" stt_compute_type: str = "int8" # "int8" | "float16" | "float32" + def __post_init__(self) -> None: + # A non-positive or non-int sandbox_execute_timeout (e.g. a hand-edited + # config file value — load_config does not coerce file values — or a + # 0/negative env value) would raise inside CustomSandboxBackend.__init__ + # and crash agent/CLI startup. Fall back to the default instead, matching + # how malformed env values already degrade to defaults. + t = self.sandbox_execute_timeout + if not isinstance(t, int) or isinstance(t, bool) or t <= 0: + logging.getLogger(__name__).warning( + "Invalid sandbox_execute_timeout %r; falling back to 300.", t + ) + self.sandbox_execute_timeout = 300 + # ============================================================================= # Config file operations @@ -410,11 +435,19 @@ def set_config_value(key: str, value: Any) -> bool: field_info = next(f for f in fields(EvoScientistConfig) if f.name == key) field_type = field_info.type + # __post_init__ only clamps on load, so validate here too. Reject bool before coercion + # (_coerce_value(True, int) would turn it into 1 and slip past). + if key == "sandbox_execute_timeout" and isinstance(value, bool): + return False + try: value = _coerce_value(value, field_type) except (ValueError, TypeError): return False + if key == "sandbox_execute_timeout" and value <= 0: + return False + setattr(config, key, value) save_config(config) return True @@ -472,6 +505,7 @@ _ENV_MAPPINGS = { "langgraph_dev_port": "EVOSCIENTIST_LANGGRAPH_DEV_PORT", "code_interpreter_timeout": "EVOSCIENTIST_CODE_INTERPRETER_TIMEOUT", "code_interpreter_max_result_chars": "EVOSCIENTIST_CODE_INTERPRETER_MAX_RESULT_CHARS", + "sandbox_execute_timeout": "EVOSCIENTIST_SANDBOX_EXECUTE_TIMEOUT", "langgraph_dev_file_persistence": "EVOSCIENTIST_LANGGRAPH_DEV_FILE_PERSISTENCE", "langgraph_dev_jobs_per_worker": "EVOSCIENTIST_LANGGRAPH_DEV_JOBS_PER_WORKER", "recursion_limit": "EVOSCIENTIST_RECURSION_LIMIT", diff --git a/EvoScientist/middleware/background.py b/EvoScientist/middleware/background.py new file mode 100644 index 0000000..9fda1ab --- /dev/null +++ b/EvoScientist/middleware/background.py @@ -0,0 +1,134 @@ +"""``BackgroundExecutionMiddleware`` — background-process tools for the main agent. + +Mirrors deepagents' ``AsyncSubAgentMiddleware`` shape (a middleware that owns a set of +tools). The tools are stateless wrappers over :mod:`EvoScientist.background`, which holds +the live, process-level registry. They reuse the sandbox's ``validate_command`` so a +background launch cannot bypass the same safety checks as ``execute``. + +Naming: these manage OS *processes* (never "job" — that word is reserved-free; async +sub-agents are *tasks*, future cron is *schedules*). +""" + +from __future__ import annotations + +from datetime import UTC, datetime + +from langchain.agents.middleware import AgentMiddleware +from langchain.tools import ToolRuntime +from langchain_core.tools import tool + +from .. import background, paths +from ..backends import prepare_sandbox_command + + +def _origin_thread_id(runtime: ToolRuntime | None) -> str | None: + """Best-effort current CLI thread_id, used to route the completion notification.""" + try: + return (runtime.config or {}).get("configurable", {}).get("thread_id") + except Exception: + return None + + +def _notify_done(proc: background.BgProcess, origin_thread_id: str | None) -> None: + """Watcher ``on_exit`` hook: enqueue a completion notification (reuses async_notifier). + + Skipped for user-stopped processes (the user already knows). The notifier is imported + lazily to keep this module free of a load-time dependency on the CLI layer. + """ + if proc.stopped: + return + rc = proc.returncode + if rc == 0: + status = "success" + elif rc is not None and rc < 0: + status = "interrupted" # terminated by a signal + else: + status = "error" + from ..cli import async_notifier + + async_notifier._enqueue( + async_notifier.AsyncTaskNotification( + task_id=proc.process_id, + agent_name=proc.name, + status=status, + received_at=datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%SZ"), + prompt=proc.command, + kind="bg-process", + origin_cli_thread_id=origin_thread_id, + ) + ) + + +@tool(parse_docstring=True) +def run_in_background( + command: str, name: str | None = None, runtime: ToolRuntime = None +) -> str: + """Launch a long-running shell command in the background and return immediately. + + Use for unbounded or very long tasks (model training, large downloads, servers) + that should not block the conversation. Output streams to a log file; poll it with + check_process and stop it with stop_process. For a bounded command that just needs + more time, prefer execute(..., timeout=N) instead of backgrounding. + + Args: + command: The shell command to run in the background. + name: Optional short label to recognize the process later. + """ + cwd = str(paths.resolve_virtual_path("/")) + # Same path-rewriting + validation as execute (shared helper) so virtual paths + # resolve to the workspace and the command can't bypass the sandbox checks. + command, error = prepare_sandbox_command(command, cwd) + if error: + return error + tid = _origin_thread_id(runtime) + process_id = background.launch( + command, cwd, name, origin_thread_id=tid, on_exit=lambda p: _notify_done(p, tid) + ) + label = f" (name={name!r})" if name else "" + return ( + f"Started background process {process_id}{label}. " + f"Output -> /.bg_processes/{process_id}.log. " + f"Poll with check_process('{process_id}'), stop with stop_process('{process_id}')." + ) + + +@tool(parse_docstring=True) +def check_process(process_id: str, runtime: ToolRuntime = None) -> str: + """Check a background process's status and recent output. + + Args: + process_id: The id returned by run_in_background. + """ + return background.status(process_id, thread_id=_origin_thread_id(runtime)) + + +@tool(parse_docstring=True) +def stop_process(process_id: str) -> str: + """Stop (kill) a running background process and its child process group. + + Args: + process_id: The id returned by run_in_background. + """ + return background.stop(process_id) + + +@tool(parse_docstring=True) +def list_processes(all_threads: bool = False, runtime: ToolRuntime = None) -> str: + """List background processes launched this session with their live statuses. + + Args: + all_threads: List processes from every session, not just the current one. + """ + return background.list_all(_origin_thread_id(runtime), include_all=all_threads) + + +class BackgroundExecutionMiddleware(AgentMiddleware): + """Adds run_in_background / check_process / stop_process / list_processes. + + Modelled on ``AsyncSubAgentMiddleware``: the middleware simply exposes the tool set. + Attached to the main agent only (async sub-agents must not spawn local processes). + """ + + def __init__(self) -> None: + super().__init__() + self.tools = [run_in_background, check_process, stop_process, list_processes] diff --git a/EvoScientist/prompts.py b/EvoScientist/prompts.py index 52c2257..6942d78 100644 --- a/EvoScientist/prompts.py +++ b/EvoScientist/prompts.py @@ -225,11 +225,17 @@ WRITING_GUIDELINES = """# Writing Guidelines # Shell execution guidelines (rules for the `execute` tool) # ============================================================================= +# NOTE: the "300s" default below is intentionally hardcoded static text, not +# templated from config. get_system_prompt() must stay byte-stable for prompt +# caching, so the configured value is NOT injected here. The actually-enforced +# timeout is cfg.sandbox_execute_timeout (CustomSandboxBackend); this number is +# just the documented default, and the per-command `timeout` override is the +# mechanism that matters to the agent. SHELL_GUIDELINES = """# Shell Execution Guidelines When using the `execute` tool for shell commands: -**Sandbox limits**: Commands time out after 300 seconds (exit code 124) and output is truncated at 100 KB. Plan accordingly. +**Sandbox limits**: Commands default to a 300s timeout (a deployment may override this default) and 100 KB output. For a known long command (e.g. a download), pass `timeout` (up to 3600s): `execute(command="wget ...", timeout=600)`. For unbounded tasks, use background execution (below). **Short commands** (< 30 seconds): Run directly ```bash @@ -237,24 +243,16 @@ python script.py pip install pandas ``` -**Long-running commands** (> 30 seconds): Run in background, then check results +**Long-running commands** (> 30 seconds): prefer the `run_in_background` tool — it launches the command detached, streams output to a log, and returns a process id immediately. Then use `check_process()` for status + recent output, `stop_process()` to kill it, and `list_processes()` to see all background processes. + +If you must background manually instead, you MUST redirect output to a file (otherwise the call blocks) and capture the PID: ```bash -# Step 1: Start in background, redirect output to log python long_task.py > /output.log 2>&1 & - -# Step 2: Check if still running -ps aux | grep long_task - -# Step 3: Read results when done -cat /output.log +echo "PID: $!" # check: ps -p · stop: kill · read: cat /output.log ``` **Before heavy compute**: Estimate runtime. If likely > 5 minutes, use background execution from the start. If GPU memory is uncertain, start with a small test run (1 epoch, small batch) before the full run. -**After a timeout (exit code 124)**: Do NOT re-run the same command. Instead: -1. Re-launch in background with output logging -2. Or reduce the workload (fewer epochs, smaller model, subset of data) - This prevents blocking the conversation during long operations. """ diff --git a/EvoScientist/stream/display.py b/EvoScientist/stream/display.py index 97c9f56..66d20b2 100644 --- a/EvoScientist/stream/display.py +++ b/EvoScientist/stream/display.py @@ -962,7 +962,7 @@ def _resolve_hitl_approval( return [{"type": "approve"} for _ in action_requests] # Config-level auto-approve - from ..config.settings import load_config + from ..config.settings import HITL_SHELL_TOOLS, load_config cfg = load_config() if cfg.auto_approve: @@ -984,8 +984,8 @@ def _resolve_hitl_approval( req.get("args", {}) if isinstance(req, dict) else getattr(req, "args", {}) ) - if name != "execute": - continue # Non-execute tools auto-approve + if name not in HITL_SHELL_TOOLS: + continue # Only shell-running tools need manual approval command = args.get("command", "") if isinstance(args, dict) else "" if not _matches_shell_allow_list(command, shell_allow_list): diff --git a/tests/test_async_notifier.py b/tests/test_async_notifier.py index d1f8e54..b470d21 100644 --- a/tests/test_async_notifier.py +++ b/tests/test_async_notifier.py @@ -394,8 +394,18 @@ def test_format_batch_message_multiple(): assert lines[0] == "[Async tasks update]" obj1 = __import__("json").loads(lines[1]) obj2 = __import__("json").loads(lines[2]) - assert obj1 == {"agent": "writing-agent", "status": "success", "task_id": "t1"} - assert obj2 == {"agent": "data-analysis-agent", "status": "error", "task_id": "t2"} + assert obj1 == { + "agent": "writing-agent", + "kind": "agent", + "status": "success", + "task_id": "t1", + } + assert obj2 == { + "agent": "data-analysis-agent", + "kind": "agent", + "status": "error", + "task_id": "t2", + } assert "check_async_task" in msg.lower() # hint to LLM diff --git a/tests/test_backends.py b/tests/test_backends.py index 332c8f0..4203633 100644 --- a/tests/test_backends.py +++ b/tests/test_backends.py @@ -1033,6 +1033,14 @@ class TestExecuteTimeoutRecovery: assert "sleep 10" in resp.output assert "> /output.log 2>&1 &" in resp.output + def test_timeout_recovery_captures_pid_and_offers_timeout(self, tmp_workspace): + backend = CustomSandboxBackend(root_dir=tmp_workspace, timeout=1) + resp = backend.execute("sleep 10") + # Background recovery captures the PID so the job can be managed later. + assert "PID: $!" in resp.output + # Recovery also offers re-running with a larger per-command timeout. + assert "timeout=600" in resp.output + def test_timeout_preserves_original_error(self, tmp_workspace): backend = CustomSandboxBackend(root_dir=tmp_workspace, timeout=1) resp = backend.execute("sleep 10") diff --git a/tests/test_background.py b/tests/test_background.py new file mode 100644 index 0000000..b883ed1 --- /dev/null +++ b/tests/test_background.py @@ -0,0 +1,152 @@ +"""Tests for EvoScientist.background — the background-process manager.""" + +import time + +import pytest + +from EvoScientist import background as bg + + +def _wait_until(predicate, timeout=4.0, interval=0.05): + """Poll ``predicate`` until true or ``timeout`` — avoids flaky fixed sleeps on slow CI.""" + deadline = time.time() + timeout + while time.time() < deadline: + if predicate(): + return True + time.sleep(interval) + return False + + +@pytest.fixture(autouse=True) +def _clean_registry(): + """Isolate each test: clear the module-global registry and reap leftovers.""" + bg._PROCESSES.clear() + yield + for proc in list(bg._PROCESSES.values()): + try: + proc.popen.kill() + except Exception: + pass + bg._PROCESSES.clear() + + +def test_launch_returns_id_and_creates_log(tmp_path): + pid = bg.launch("echo hi", str(tmp_path)) + assert pid in bg._PROCESSES + assert (tmp_path / ".bg_processes" / f"{pid}.log").exists() + + +def test_status_running_then_exited(tmp_path): + pid = bg.launch("sleep 1", str(tmp_path)) + assert "RUNNING" in bg.status(pid) + assert _wait_until(lambda: "EXITED" in bg.status(pid)) + out = bg.status(pid) + assert "EXITED" in out + assert "code 0" in out + + +def test_output_captured_in_status(tmp_path): + pid = bg.launch("echo hello-from-bg", str(tmp_path)) + assert _wait_until(lambda: "hello-from-bg" in bg.status(pid)) + + +def test_large_log_returns_truncated_tail(tmp_path): + """status() preserves the truncation contract for a large log (output shape, not I/O).""" + pid = bg.launch("true", str(tmp_path)) + log_path = tmp_path / ".bg_processes" / f"{pid}.log" + log_path.write_bytes(b"A" * 5000 + b"TAIL_MARKER") + out = bg.status(pid, tail_bytes=64) + assert "...(truncated)..." in out + assert "TAIL_MARKER" in out + assert "A" * 5000 not in out # the head was not loaded + + +def test_stop_kills_running_process(tmp_path): + pid = bg.launch("sleep 600", str(tmp_path)) + assert "RUNNING" in bg.status(pid) + out = bg.stop(pid) + assert "Stopped" in out + assert bg._PROCESSES[pid].popen.poll() is not None # actually terminated + + +def test_stop_already_finished_is_graceful(tmp_path): + pid = bg.launch("true", str(tmp_path)) + assert _wait_until(lambda: bg._PROCESSES[pid].popen.poll() is not None) + assert "already finished" in bg.stop(pid) + + +def test_exited_elapsed_is_frozen(tmp_path): + """Elapsed for an exited process freezes at its runtime, it must not keep growing.""" + pid = bg.launch("true", str(tmp_path)) + assert _wait_until(lambda: bg._PROCESSES[pid].finished_ts is not None) + bg.status(pid) # observe exit -> records finished_ts + proc = bg._PROCESSES[pid] + assert proc.finished_ts is not None + first = bg._elapsed(proc) + time.sleep(1.1) # intentional: prove elapsed stays frozen, not ticking up + assert bg._elapsed(proc) == first + + +def test_watcher_records_exit_without_polling(tmp_path): + """The daemon watcher records exit on its own (no status() call needed).""" + pid = bg.launch("true", str(tmp_path)) + assert _wait_until(lambda: bg._PROCESSES[pid].finished_ts is not None) + proc = bg._PROCESSES[pid] + assert proc.finished_ts is not None + assert proc.returncode == 0 + + +def test_on_exit_callback_fires(tmp_path): + """on_exit is invoked with the BgProcess once the process exits.""" + fired = {} + + def cb(proc): + fired["pid"] = proc.process_id + fired["rc"] = proc.returncode + + pid = bg.launch("true", str(tmp_path), on_exit=cb) + assert _wait_until(lambda: fired.get("pid") == pid and fired.get("rc") == 0) + assert fired.get("pid") == pid + assert fired.get("rc") == 0 + + +def test_unknown_id_errors_gracefully(): + assert "No such background process" in bg.status("deadbeef") + assert "No such background process" in bg.stop("deadbeef") + + +def test_list_all(tmp_path): + assert "No background processes" in bg.list_all() + pid = bg.launch("sleep 1", str(tmp_path)) + listing = bg.list_all() + assert pid in listing + assert "RUNNING" in listing + + +def test_list_all_scopes_to_origin_thread(tmp_path): + """list_all defaults to the launching session; include_all sees every session.""" + pid_a = bg.launch("sleep 1", str(tmp_path), origin_thread_id="A") + pid_b = bg.launch("sleep 1", str(tmp_path), origin_thread_id="B") + listing_a = bg.list_all("A") + assert pid_a in listing_a + assert pid_b not in listing_a # B's process is hidden from session A + everything = bg.list_all("A", include_all=True) + assert pid_a in everything + assert pid_b in everything + + +def test_list_all_hints_at_other_sessions(tmp_path): + """A session with no processes of its own is told others exist.""" + bg.launch("sleep 1", str(tmp_path), origin_thread_id="A") + out = bg.list_all("B") # a different session + assert "other sessions" in out + assert "all_threads=True" in out + + +def test_dedup_is_per_thread(tmp_path): + """A check from one session must not suppress another session's completion ping.""" + pid = bg.launch("true", str(tmp_path), origin_thread_id="A") + assert _wait_until(lambda: bg._PROCESSES[pid].finished_ts is not None) + bg.status(pid, thread_id="B") # a DIFFERENT session inspects it + assert bg.was_observed_done(pid, "B") is True # B saw it + assert bg.was_observed_done(pid, "A") is False # launcher A did not -> still notify diff --git a/tests/test_background_middleware.py b/tests/test_background_middleware.py new file mode 100644 index 0000000..d2a8602 --- /dev/null +++ b/tests/test_background_middleware.py @@ -0,0 +1,274 @@ +"""Tests for BackgroundExecutionMiddleware and its tools.""" + +import time + +import pytest + +from EvoScientist import background as bg +from EvoScientist.middleware.background import ( + BackgroundExecutionMiddleware, + check_process, + list_processes, + run_in_background, + stop_process, +) + + +def _wait_until(predicate, timeout=4.0, interval=0.05): + """Poll ``predicate`` until true or ``timeout`` — avoids flaky fixed sleeps on slow CI.""" + deadline = time.time() + timeout + while time.time() < deadline: + if predicate(): + return True + time.sleep(interval) + return False + + +@pytest.fixture(autouse=True) +def _clean_registry(): + from EvoScientist.cli import async_notifier + + bg._PROCESSES.clear() + async_notifier.drain_notifications(None) + yield + for proc in list(bg._PROCESSES.values()): + try: + proc.popen.kill() + except Exception: + pass + bg._PROCESSES.clear() + async_notifier.drain_notifications(None) + + +def test_middleware_registers_four_tools(): + mw = BackgroundExecutionMiddleware() + names = {t.name for t in mw.tools} + assert names == { + "run_in_background", + "check_process", + "stop_process", + "list_processes", + } + + +def test_no_job_in_tool_names(): + """Naming ADR: the word 'job' must not appear in the tool surface.""" + mw = BackgroundExecutionMiddleware() + assert not any("job" in t.name.lower() for t in mw.tools) + + +def test_run_rejects_dangerous_command_without_launching(monkeypatch): + launched = {"called": False} + + def _spy(*args, **kwargs): + launched["called"] = True + return "should-not-happen" + + monkeypatch.setattr(bg, "launch", _spy) + out = run_in_background.invoke({"command": "sudo rm -rf /"}) + assert launched["called"] is False + assert "blocked" in out.lower() + + +def test_run_launches_valid_command(tmp_path, monkeypatch): + # Pin the workspace cwd to a temp dir so the launch is isolated. + monkeypatch.setattr("EvoScientist.paths.resolve_virtual_path", lambda _vp: tmp_path) + out = run_in_background.invoke({"command": "echo ok", "name": "demo"}) + assert "Started background process" in out + assert "check_process" in out + assert len(bg._PROCESSES) == 1 + + +def test_run_applies_virtual_path_rewriting(tmp_path, monkeypatch): + """run_in_background must rewrite virtual paths like execute (shared preprocessing).""" + monkeypatch.setattr("EvoScientist.paths.resolve_virtual_path", lambda _vp: tmp_path) + captured = {} + + def _spy(command, cwd, name=None, *, origin_thread_id=None, on_exit=None): + captured["command"] = command + return "pidX" + + monkeypatch.setattr(bg, "launch", _spy) + run_in_background.invoke({"command": "python /train.py"}) + # virtual absolute path -> workspace-relative, same as execute would produce + assert captured["command"] == "python ./train.py" + + +def test_run_enqueues_completion_notification(tmp_path, monkeypatch): + """A finished background process enqueues a shell completion notification.""" + from EvoScientist.cli import async_notifier + + monkeypatch.setattr("EvoScientist.paths.resolve_virtual_path", lambda _vp: tmp_path) + run_in_background.invoke({"command": "true", "name": "quick"}) + # drain consumes, so accumulate across polls until the watcher's on_exit enqueues. + notifs = [] + deadline = time.time() + 4.0 + while time.time() < deadline: + notifs.extend(async_notifier.drain_notifications(None)) + if any(n.kind == "bg-process" for n in notifs): + break + time.sleep(0.05) + assert any(n.kind == "bg-process" and n.status == "success" for n in notifs) + + +def test_origin_thread_id_reads_runtime_config(): + """thread_id is read from runtime.config['configurable'] (graph-injected).""" + from types import SimpleNamespace + + from EvoScientist.middleware.background import _origin_thread_id + + runtime = SimpleNamespace(config={"configurable": {"thread_id": "T-7"}}) + assert _origin_thread_id(runtime) == "T-7" + assert _origin_thread_id(None) is None # direct .invoke() / no runtime + + +def test_notify_done_routes_to_origin_thread(tmp_path): + """_notify_done enqueues the completion notification to the launching thread.""" + from EvoScientist.cli import async_notifier + from EvoScientist.middleware.background import _notify_done + + pid = bg.launch("true", str(tmp_path)) # no on_exit -> no auto-notify here + assert _wait_until(lambda: bg._PROCESSES[pid].finished_ts is not None) + _notify_done(bg._PROCESSES[pid], "T-123") + routed = async_notifier.drain_notifications("T-123") + assert any(n.task_id == pid and n.origin_cli_thread_id == "T-123" for n in routed) + + +def test_stopped_process_suppresses_notification(tmp_path, monkeypatch): + """A user-stopped process must NOT emit a completion notification.""" + from EvoScientist.cli import async_notifier + + monkeypatch.setattr("EvoScientist.paths.resolve_virtual_path", lambda _vp: tmp_path) + run_in_background.invoke({"command": "sleep 600"}) + (pid,) = list(bg._PROCESSES.keys()) + stop_process.invoke({"process_id": pid}) + # Wait until the watcher observed the exit — it would have enqueued here if the + # process weren't user-stopped. _notify_done is a no-op for stopped processes. + assert _wait_until(lambda: bg._PROCESSES[pid].finished_ts is not None) + notifs = async_notifier.drain_notifications(None) + assert not any(n.task_id == pid for n in notifs) + + +def test_checked_after_exit_dedups_notification(tmp_path): + """Agent checking a finished process suppresses its completion notification.""" + from EvoScientist.cli.async_notifier import ( + AsyncTaskNotification, + dedup_notifications, + ) + + pid = bg.launch("true", str(tmp_path)) + assert _wait_until(lambda: bg._PROCESSES[pid].finished_ts is not None) + bg.status(pid) # agent checks AFTER exit + assert bg.was_observed_done(pid) is True + n = AsyncTaskNotification( + task_id=pid, + agent_name="x", + status="success", + received_at="t", + kind="bg-process", + ) + assert dedup_notifications([n], {}) == [] # deduped + + +def test_not_checked_after_exit_keeps_notification(tmp_path): + """A finished process the agent never checked still notifies.""" + from EvoScientist.cli.async_notifier import ( + AsyncTaskNotification, + dedup_notifications, + ) + + pid = bg.launch("true", str(tmp_path)) + assert _wait_until( + lambda: bg._PROCESSES[pid].finished_ts is not None + ) # exit, but do NOT check + assert bg.was_observed_done(pid) is False + n = AsyncTaskNotification( + task_id=pid, + agent_name="x", + status="success", + received_at="t", + kind="bg-process", + ) + assert dedup_notifications([n], {}) == [n] # survives + + +def test_shell_notification_renders_own_background_frame(): + """Shell notifications render under '✦ Background ✦', not 'Agent Teams'.""" + from EvoScientist.cli.async_notifier import ( + AsyncTaskNotification, + format_notification_lines, + ) + + n = AsyncTaskNotification( + task_id="fe60ce9c", + agent_name="test-20s", + status="success", + received_at="", + prompt="python train.py", + kind="bg-process", + ) + lines = format_notification_lines([n]) + top, body = lines[0][0], lines[1][0] + assert "Background" in top + assert "Agent Teams" not in top + assert "test-20s" in body + assert "Cmd:" in body + + +def test_mixed_notifications_render_two_frames(): + """A mixed batch shows both an Agent Teams frame and a Background frame.""" + from EvoScientist.cli.async_notifier import ( + AsyncTaskNotification, + format_notification_lines, + ) + + task = AsyncTaskNotification("t1", "writing-agent", "success", "", "") + shell = AsyncTaskNotification("p1", "demo", "success", "", "", kind="bg-process") + blob = "\n".join(t for t, _ in format_notification_lines([task, shell])) + assert "Agent Teams" in blob + assert "Background" in blob + + +def test_shell_notification_hints_check_process(): + """format_batch_message points shell processes to check_process, not check_async_task.""" + from EvoScientist.cli.async_notifier import ( + AsyncTaskNotification, + format_batch_message, + ) + + n = AsyncTaskNotification( + task_id="ab12", + agent_name="demo", + status="success", + received_at="x", + kind="bg-process", + ) + msg = format_batch_message([n]) + assert "check_process" in msg + assert "check_async_task" not in msg # shell-only batch -> no sub-agent hint + + +def test_check_and_list_route_to_manager(tmp_path, monkeypatch): + monkeypatch.setattr("EvoScientist.paths.resolve_virtual_path", lambda _vp: tmp_path) + run_in_background.invoke({"command": "sleep 1"}) + (pid,) = bg._PROCESSES.keys() + assert pid in check_process.invoke({"process_id": pid}) + assert pid in list_processes.invoke({}) + assert "Stopped" in stop_process.invoke( + {"process_id": pid} + ) or "finished" in stop_process.invoke({"process_id": pid}) + + +def test_list_processes_forwards_all_threads(monkeypatch): + """The all_threads tool arg is forwarded to background.list_all(include_all=...).""" + captured = {} + + def _spy(thread_id=None, *, include_all=False): + captured["include_all"] = include_all + return "ok" + + monkeypatch.setattr(bg, "list_all", _spy) + list_processes.invoke({"all_threads": True}) + assert captured["include_all"] is True + list_processes.invoke({}) + assert captured["include_all"] is False diff --git a/tests/test_config.py b/tests/test_config.py index f315136..118f070 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -381,6 +381,52 @@ class TestPriorityChain: config = get_effective_config() assert config.channel_debug_tracing is True + def test_sandbox_execute_timeout_default(self, temp_config_dir, clean_env): + """Sandbox execute timeout defaults to 300 seconds.""" + assert EvoScientistConfig().sandbox_execute_timeout == 300 + assert get_effective_config().sandbox_execute_timeout == 300 + + def test_env_sandbox_execute_timeout_override(self, temp_config_dir, monkeypatch): + """Sandbox execute timeout can be set via env var and coerces to int.""" + monkeypatch.setenv("EVOSCIENTIST_SANDBOX_EXECUTE_TIMEOUT", "600") + config = get_effective_config() + assert config.sandbox_execute_timeout == 600 + assert isinstance(config.sandbox_execute_timeout, int) + + def test_sandbox_execute_timeout_invalid_falls_back(self): + """Non-positive / non-int values fall back to the default (would + otherwise crash CustomSandboxBackend construction at startup).""" + assert ( + EvoScientistConfig(sandbox_execute_timeout=0).sandbox_execute_timeout == 300 + ) + assert ( + EvoScientistConfig(sandbox_execute_timeout=-5).sandbox_execute_timeout + == 300 + ) + assert ( + EvoScientistConfig(sandbox_execute_timeout="abc").sandbox_execute_timeout + == 300 + ) + assert ( + EvoScientistConfig(sandbox_execute_timeout=True).sandbox_execute_timeout + == 300 + ) + + def test_set_sandbox_execute_timeout_rejects_invalid( + self, temp_config_dir, clean_env + ): + """set_config_value must reject (not silently persist) a non-positive timeout.""" + save_config(EvoScientistConfig(sandbox_execute_timeout=120)) + assert set_config_value("sandbox_execute_timeout", 0) is False + assert set_config_value("sandbox_execute_timeout", -5) is False + # bool is an int subclass; reject it before coercion turns True into 1. + assert set_config_value("sandbox_execute_timeout", True) is False + # The earlier valid value is untouched on disk. + assert get_config_value("sandbox_execute_timeout") == 120 + # A valid value still goes through. + assert set_config_value("sandbox_execute_timeout", 600) is True + assert get_config_value("sandbox_execute_timeout") == 600 + def test_env_api_key_override(self, temp_config_dir, monkeypatch): """Test API keys from env override file.""" save_config(EvoScientistConfig(anthropic_api_key="file-key")) diff --git a/tests/test_hitl.py b/tests/test_hitl.py index 2f7c99e..8f80df0 100644 --- a/tests/test_hitl.py +++ b/tests/test_hitl.py @@ -260,6 +260,67 @@ class TestResolveHitlApproval: finally: disp._session_auto_approve = original + def test_run_in_background_not_in_allow_list_prompts(self): + """run_in_background must NOT auto-approve — it runs shell like execute.""" + import EvoScientist.stream.display as disp + from EvoScientist.stream.display import _resolve_hitl_approval + + original = disp._session_auto_approve + try: + disp._session_auto_approve = False + mock_cfg = MagicMock() + mock_cfg.auto_approve = False + mock_cfg.shell_allow_list = "ls,cat" + with patch( + "EvoScientist.config.settings.load_config", return_value=mock_cfg + ): + with patch( + "EvoScientist.stream.display._prompt_hitl_approval" + ) as mock_prompt: + mock_prompt.return_value = [{"type": "approve"}] + result = _resolve_hitl_approval( + { + "action_requests": [ + { + "name": "run_in_background", + "args": {"command": "rm -rf /"}, + } + ], + } + ) + assert result == [{"type": "approve"}] + mock_prompt.assert_called_once() # prompted, not silently approved + finally: + disp._session_auto_approve = original + + def test_run_in_background_in_allow_list_auto_approves(self): + """An allow-listed command still auto-approves for run_in_background.""" + import EvoScientist.stream.display as disp + from EvoScientist.stream.display import _resolve_hitl_approval + + original = disp._session_auto_approve + try: + disp._session_auto_approve = False + mock_cfg = MagicMock() + mock_cfg.auto_approve = False + mock_cfg.shell_allow_list = "python" + with patch( + "EvoScientist.config.settings.load_config", return_value=mock_cfg + ): + result = _resolve_hitl_approval( + { + "action_requests": [ + { + "name": "run_in_background", + "args": {"command": "python train.py"}, + } + ], + } + ) + assert result == [{"type": "approve"}] + finally: + disp._session_auto_approve = original + # ============================================================================= # Config fields @@ -488,6 +549,21 @@ class TestConsumerHitlHelpers: ) assert result is False + def test_should_auto_approve_run_in_background_no_allowlist(self): + """Channel path must NOT auto-approve run_in_background (same as execute).""" + from EvoScientist.channels.consumer import _should_auto_approve + + mock_cfg = MagicMock() + mock_cfg.auto_approve = False + mock_cfg.shell_allow_list = "" + with patch("EvoScientist.config.settings.load_config", return_value=mock_cfg): + result = _should_auto_approve( + [ + {"name": "run_in_background", "args": {"command": "rm -rf /"}}, + ] + ) + assert result is False + def test_should_auto_approve_config_true(self): from EvoScientist.channels.consumer import _should_auto_approve diff --git a/tests/test_prompts.py b/tests/test_prompts.py index 45b04dc..66a372c 100644 --- a/tests/test_prompts.py +++ b/tests/test_prompts.py @@ -153,8 +153,8 @@ class TestShellGuidelines: assert len(SHELL_GUIDELINES) > 0 def test_mentions_timeout_limit(self): - assert "300" in SHELL_GUIDELINES - assert "124" in SHELL_GUIDELINES + assert "300" in SHELL_GUIDELINES # default timeout + assert "3600" in SHELL_GUIDELINES # per-command override ceiling def test_mentions_background_execution(self): assert "background" in SHELL_GUIDELINES.lower()