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
This commit is contained in:
@@ -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)
|
||||
|
||||
+45
-33
@@ -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 <PID> (or: ps aux | grep {grep_hint}) · "
|
||||
f"Read: cat /output.log · Stop: kill <PID>"
|
||||
),
|
||||
exit_code=response.exit_code,
|
||||
truncated=response.truncated,
|
||||
|
||||
@@ -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 ``<cwd>/.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)
|
||||
@@ -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
|
||||
|
||||
@@ -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: <prompt preview> 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)
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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]
|
||||
+11
-13
@@ -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(<id>)` for status + recent output, `stop_process(<id>)` 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 <PID> · stop: kill <PID> · 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.
|
||||
"""
|
||||
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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"))
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user