From 7130ac83a10f74183a36be152cda40197199fcf4 Mon Sep 17 00:00:00 2001 From: Aster Date: Thu, 17 Sep 2026 01:30:20 -0700 Subject: [PATCH] fix(code-execution): separate script lifetime from bounded observation Reserve interpreter startup, isolate late replies, serialize wake delivery, retain execution results, and support ending an observation with completion wake armed. Co-Authored-By: Codex (GPT-6) --- changelog.d/code-execution-lifecycle.added.md | 3 + .../code-execution-lifecycle.breaking.md | 5 + changelog.d/code-execution-lifecycle.fixed.md | 5 + docs/code-execution-lifecycle.md | 121 +++++++ src/code-execution/py-runner.ts | 89 ++--- src/code-execution/runtime-py.ts | 2 +- src/code-execution/script-run.ts | 54 ++++ src/code-execution/tool-definition.ts | 87 ++--- src/framework.ts | 305 ++++++++++++------ src/types/framework.ts | 2 + test/code-execution.test.ts | 205 ++++++++++++ test/script-run.test.ts | 64 ++++ 12 files changed, 739 insertions(+), 203 deletions(-) create mode 100644 changelog.d/code-execution-lifecycle.added.md create mode 100644 changelog.d/code-execution-lifecycle.breaking.md create mode 100644 changelog.d/code-execution-lifecycle.fixed.md create mode 100644 docs/code-execution-lifecycle.md create mode 100644 src/code-execution/script-run.ts create mode 100644 test/script-run.test.ts diff --git a/changelog.d/code-execution-lifecycle.added.md b/changelog.d/code-execution-lifecycle.added.md new file mode 100644 index 0000000..664e2a5 --- /dev/null +++ b/changelog.d/code-execution-lifecycle.added.md @@ -0,0 +1,3 @@ +- Retain execution IDs/results for bounded waits, listing and cancellation. Expose + a framework method for the host's admin-only release command to unblock an + observation without killing its script, then wake on completion. diff --git a/changelog.d/code-execution-lifecycle.breaking.md b/changelog.d/code-execution-lifecycle.breaking.md new file mode 100644 index 0000000..0cd4013 --- /dev/null +++ b/changelog.d/code-execution-lifecycle.breaking.md @@ -0,0 +1,5 @@ +- **Callers of `code_execution`:** foreground calls now return a running `script_id` + after a bounded observation budget (10 seconds by default). Check `status` and + use `action=wait` for the result; timeout no longer monopolizes inference until + script termination. `on_timeout=end_turn` arms completion notification and + releases the turn. Explicit cancellation remains separate. diff --git a/changelog.d/code-execution-lifecycle.fixed.md b/changelog.d/code-execution-lifecycle.fixed.md new file mode 100644 index 0000000..29de738 --- /dev/null +++ b/changelog.d/code-execution-lifecycle.fixed.md @@ -0,0 +1,5 @@ +- Prevent concurrent interpreter startup from losing an execution; settle startup + cancellation and reject stale tool/wake replies after interpreter replacement. + Preserve sub-second Python tool timeouts. +- Serialize background wake delivery across rate limits and caps; report delivery + failure to Python. Keep deferred end-turn effects scoped to their execution. diff --git a/docs/code-execution-lifecycle.md b/docs/code-execution-lifecycle.md new file mode 100644 index 0000000..02af216 --- /dev/null +++ b/docs/code-execution-lifecycle.md @@ -0,0 +1,121 @@ +# Code execution: execution, observation, and attention + +A running operation and a waiting agent have different lifetimes. Python's +`await` suspends the coroutine; it should not oblige the agent to spend its +whole turn waiting. Keep Python semantics intact and put a bounded observation +around execution in the harness. + +This branch implements that separation while preserving existing Python +contexts and background watchers. It is an initial lifecycle change, not a +durable job service or a replacement for the terminal execution protocol. + +## Agent interface implemented here + +One tool, `code_execution`, retains `code`, `background`, `action`, and +`script_id`, and adds `wait_ms` and `on_timeout`: + +```json +{"code":"import asyncio\nawait asyncio.sleep(30)\nprint('done')","wait_ms":1000} +``` + +Quick results return stdout, stderr, and return_code as before, with an +execution ID and status. If still running after the observation budget, the +tool returns that ID with status `running`. The script continues, and its +completion notifies the owner. `action=wait` retrieves the retained result or +observes the same operation for another bounded interval. `wait_ms=0` returns +immediately; the maximum is 60 seconds. The default is 10 seconds, configurable +through agent-framework's `codeExecution.foregroundWaitMs`. + +```json +{"action":"wait","script_id":"py-1","wait_ms":1000,"on_timeout":"end_turn"} +``` + +If the result arrives within the budget, the agent receives it and continues. +Otherwise, completion notification is armed **before** the tool result requests +`endTurn`. This uses the scheduler's existing result boundary: the stream ends +without cancelling the script. Completion is queued even if it arrives during +turn teardown. It bypasses ambient event gating because the agent armed it. +This ends a turn, not a timed gate sleep: other authorized events can still wake +the agent before the script finishes. + +Repeated waits refer to one execution and never replay its side effects. A +wait already present at completion receives the result directly, suppressing +the additional completion wake. Multiple expired waits arm one completion +notice. The notice contains a bounded tail; the full captured result can be +retrieved with `wait` and uses the normal spill policy. Ordinary foreground +stdout is currently available at completion, not streamed during execution. + +The existing foreground interpreter remains persistent and serial. A second +run while it is busy fails with the running ID and recovery instructions. +`background=true` uses an independent interpreter, retains `wake_agent`, and +journals output to the workspace. Its default remains immediate return and +silent clean exit; specifying a wait budget uses the same observation policy. +Background crashes still notify. `list` reports both modes; the legacy +`background_scripts` field is preserved. `cancel` stops Python explicitly. + +Ephemeral agents cannot request `on_timeout=end_turn`: their owner is destroyed +when the turn ends and cannot receive the promised wake. Owner disposal cancels +its executions and releases their interpreters. Background watchers retain the +existing primary-agent restriction. + +## Operator interface implemented here + +The connectome-host companion adds the admin command `/release-wait [script_id]` +alongside `/undo`. With no ID it releases all current code-execution observations +for the selected agent. With an ID it releases only that execution's observers. +It does not abort Python or reset the agent. The underlying framework method is +`releaseCodeExecutionWait(agentName, scriptId?)`. + +This is an admin operation, not an agent tool. The host requires trusted command +provenance: the local operator CLI/TUI or a full-authority web client. Generic +headless/fleet IPC and read-only web observers cannot invoke it. A caller cannot +grant itself admin status with an argument in the slash command. The framework's +generic socket API does not expose a release command. + +A release after completion or when no observation is active returns `released: 0`; +it does not end an unrelated turn. Releasing a code observation does not settle +unrelated tools in the same batch. This is not a general-purpose inference reset. + +## Correctness changes + +- Reserve the Python runner before interpreter startup, so simultaneous cold + calls cannot overwrite one another and startup can be cancelled. +- Bind asynchronous tool replies and wake acknowledgements to the originating + interpreter and execution. A late reply must not satisfy a replacement + interpreter's reused `t1` or `w1` ID. +- Preserve fractional seconds when sending inner-tool timeouts to Python. +- Serialize each script's wake requests, enforcing the interval and cap under + `asyncio.gather`; stop rate-limit waits when execution ends. +- Acknowledge `wake_agent` only after context delivery and inference enqueue + succeed. Delivery failures are errors visible to Python. +- Scope inner-tool end-turn requests to their execution, preventing background + or late results from ending an unrelated foreground call. + +## Limits and the next design steps + +Execution IDs and the five most recently settled results per owner are held in +memory. Host restart loses them and stops Python. Completion delivery is not a +durable receipt/outbox: storage failure is logged, and an operator can retrieve +the result while the host remains alive. Explicit script wakes report delivery +failure to Python. Cancellation is not rollback, and does not cancel an already +dispatched inner tool. There is no claim of exactly-once external effects. + +The next increment should give operations a bounded, cursor-addressed journal +and persistent terminal outcomes. Associate explicit notifications with delivery +receipts: queued, delivered at a tool boundary, or used to start a fresh turn. +That would avoid an unnecessary subsequent turn when an active inference has +already consumed the completion event. The existing wake scheduler can currently +queue such a follow-up. Crash recovery should report an interrupted operation, +not automatically rerun side-effectful code. + +Keep routine output in the journal; only explicit signals, failures, and +requested completions should request attention. Waiting on several operations +should eventually offer “any” and “all” without polling code. Add named Python +contexts only when workflows need several persistent contexts; arbitrary Python +globals and stdout cannot safely be shared by concurrent top-level scripts. + +For MCPL, this should be an operation lifecycle with ownership and cancellation +capabilities, separate from inference lifecycle metadata. A terminal operation +should retain its own run ID and completion even if Python or the model turn +ends. The harness observes that operation; a server should not need to guess +whether a timeout meant “stop waiting,” “stop executing,” or “go idle.” diff --git a/src/code-execution/py-runner.ts b/src/code-execution/py-runner.ts index b9f2cf8..1e210a5 100644 --- a/src/code-execution/py-runner.ts +++ b/src/code-execution/py-runner.ts @@ -139,36 +139,21 @@ export class PyRunner { }; } this.clearIdleTimer(); - - try { - await this.ensureChild(); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - this.reclaim('spawn-failed'); - return { - stdout: '', - stderr: `Failed to start python runtime (${this.pythonPath}): ${message}`, - returnCode: 1, - aborted: true, - }; - } - const execId = `e${++this.execCounter}`; const deadlineMs = background?.lifetimeMs ?? this.scriptTimeoutMs; this.onWake = background?.onWake ?? null; - const result = await new Promise((resolve) => { - const pending: PendingExec = { - id: execId, - resolve, - deadlineTimer: null, - killTimer: null, - settled: false, - }; - this.pending = pending; - + let resolve!: (result: ExecResult) => void; + const completion = new Promise((r) => { resolve = r; }); + // Reserve before the first await. Startup is part of the execution: + // concurrent calls must not overwrite its result resolver, and abort / + // dispose must be able to settle it even before Python says ready. + const pending: PendingExec = { + id: execId, resolve, deadlineTimer: null, killTimer: null, settled: false, + }; + this.pending = pending; + void this.ensureChild().then(() => { + if (this.pending !== pending || pending.settled) return; pending.deadlineTimer = setTimeout(() => { - // Deadline: ask politely first (script sees CancelledError and its - // exec_result still flows back), then kill on unresponsiveness. this.send({ op: 'cancel', id: execId, reason: 'deadline' }); pending.killTimer = setTimeout(() => { this.settlePending({ @@ -180,22 +165,31 @@ export class PyRunner { this.reclaim('deadline-kill'); }, CANCEL_GRACE_MS); }, deadlineMs); - // A day-scale background deadline must not hold the process open. if (background) pending.deadlineTimer.unref?.(); - this.send({ op: 'init', tools: tools.map((t) => ({ py_name: t.pyName, tool_name: t.toolName })), - call_timeout_s: Math.round(this.toolCallTimeoutMs / 1000), - ...(background - ? { background: true, log_path: background.logPath ?? null } - : {}), + call_timeout_s: this.toolCallTimeoutMs / 1000, + ...(background ? { background: true, log_path: background.logPath ?? null } : {}), }); this.send({ op: 'exec', id: execId, code }); + }).catch((err) => { + if (this.pending !== pending || pending.settled) return; + const message = err instanceof Error ? err.message : String(err); + this.settlePending({ + stdout: '', + stderr: `Failed to start python runtime (${this.pythonPath}): ${message}`, + returnCode: 1, + aborted: true, + }); + this.reclaim('spawn-failed'); }); + const result = await completion; - this.onWake = null; - if (!background) this.armIdleTimer(); + if (!this.pending) { + this.onWake = null; + if (!background) this.armIdleTimer(); + } return result; } @@ -261,6 +255,7 @@ export class PyRunner { }); child.on('exit', (exitCode, signal) => { + if (this.child !== child) return; if (this.pending) { this.settlePending({ stdout: '', @@ -275,7 +270,9 @@ export class PyRunner { }); this.reader = createInterface({ input: child.stdout }); - this.reader.on('line', (line) => this.handleLine(line)); + this.reader.on('line', (line) => { + if (this.child === child) this.handleLine(line); + }); this.childReady = new Promise((resolve, reject) => { const onReady = () => { @@ -297,10 +294,12 @@ export class PyRunner { const cleanup = () => { clearTimeout(timeout); this.readyResolver = null; + this.readyRejecter = null; child.off('exit', onExit); child.off('error', onError); }; this.readyResolver = onReady; + this.readyRejecter = onError; child.on('exit', onExit); child.on('error', onError); }); @@ -308,9 +307,10 @@ export class PyRunner { } private readyResolver: (() => void) | null = null; + private readyRejecter: ((error: Error) => void) | null = null; private handleLine(line: string): void { - let msg: { op?: string; id?: string; name?: string; args?: unknown; stdout?: string; stderr?: string; return_code?: number }; + let msg: { op?: string; id?: string; exec_id?: string; name?: string; args?: unknown; stdout?: string; stderr?: string; return_code?: number }; try { msg = JSON.parse(line); } catch { @@ -324,6 +324,9 @@ export class PyRunner { return; case 'tool_call': { + const child = this.child; + const pending = this.pending; + if (!pending || msg.exec_id !== pending.id) return; const callId = msg.id; const toolName = msg.name; if (!callId || !toolName) return; @@ -334,12 +337,19 @@ export class PyRunner { this.onToolCall(toolName, args) .catch((err) => `Error: ${err instanceof Error ? err.message : String(err)}`) .then((result) => { - this.send({ op: 'tool_result', id: callId, result }); + // A reclaimed interpreter starts call ids again at t1. Late + // results from its predecessor must never satisfy the new call. + if (this.child === child && this.pending === pending) { + this.send({ op: 'tool_result', id: callId, result }); + } }); return; } case 'wake': { + const child = this.child; + const pending = this.pending; + if (!pending || msg.exec_id !== pending.id) return; const wakeId = msg.id; if (!wakeId) return; const line = typeof (msg as { line?: unknown }).line === 'number' @@ -352,7 +362,9 @@ export class PyRunner { `wake handler failed: ${err instanceof Error ? err.message : String(err)}`) : Promise.resolve('this script is not allowed to wake the agent'); void refuse.then((error) => { - this.send({ op: 'wake_ack', id: wakeId, ...(error ? { error } : {}) }); + if (this.child === child && this.pending === pending) { + this.send({ op: 'wake_ack', id: wakeId, ...(error ? { error } : {}) }); + } }); return; } @@ -408,6 +420,7 @@ export class PyRunner { } private teardownChild(): void { + this.readyRejecter?.(new Error('python runtime reclaimed during startup')); if (this.reader) { this.reader.close(); this.reader = null; diff --git a/src/code-execution/runtime-py.ts b/src/code-execution/runtime-py.ts index f29c15f..79d10b0 100644 --- a/src/code-execution/runtime-py.ts +++ b/src/code-execution/runtime-py.ts @@ -196,7 +196,7 @@ def _make_tool_fn(tool_name, py_name): except asyncio.TimeoutError: raise TimeoutError( "Calling tool ['" + tool_name + "'] timed out (no response after " - + str(int(CALL_TIMEOUT_S)) + "s)." + + str(CALL_TIMEOUT_S) + "s)." ) finally: _pending_tool_futures.pop(call_id, None) diff --git a/src/code-execution/script-run.ts b/src/code-execution/script-run.ts new file mode 100644 index 0000000..3a3cf9c --- /dev/null +++ b/src/code-execution/script-run.ts @@ -0,0 +1,54 @@ +import type { ExecResult } from './py-runner.js'; + +export type TimeoutPolicy = 'continue' | 'end_turn'; +export interface ScriptObservation { + result?: ExecResult; + endTurn: boolean; +} + +/** A script's lifetime is independent of any one caller's observation budget. + * One completion notice is armed when an observation times out or is released. + * An observer present at completion receives the result directly instead. + */ +export class ScriptRun { + result: ExecResult | undefined; + private waiters = new Set<(endTurn?: boolean) => void>(); + private notifyOnCompletion = false; + + constructor( + completion: Promise, + onComplete: (result: ExecResult, notify: boolean, observed: boolean) => void, + ) { + void completion.then((result) => { + this.result = result; + const observed = this.waiters.size > 0; + for (const finish of [...this.waiters]) finish(); + onComplete(result, this.notifyOnCompletion && !observed, observed); + }); + } + + get observing(): boolean { return this.waiters.size > 0; } + + observe(waitMs: number, onTimeout: TimeoutPolicy): Promise { + if (this.result) return Promise.resolve({ result: this.result, endTurn: false }); + return new Promise((resolve) => { + let timer: ReturnType | undefined; + const finish = (endTurn = false) => { + if (!this.waiters.delete(finish)) return; + clearTimeout(timer); + if (!this.result) this.notifyOnCompletion = true; + resolve({ result: this.result, endTurn: !this.result && endTurn }); + }; + this.waiters.add(finish); + if (waitMs === 0) finish(onTimeout === 'end_turn'); + else timer = setTimeout(() => finish(onTimeout === 'end_turn'), waitMs); + }); + } + + /** Operator rescue: end only the observation/turn, never the script. */ + release(): number { + const count = this.waiters.size; + for (const finish of [...this.waiters]) finish(true); + return count; + } +} diff --git a/src/code-execution/tool-definition.ts b/src/code-execution/tool-definition.ts index 451733a..d4047bb 100644 --- a/src/code-execution/tool-definition.ts +++ b/src/code-execution/tool-definition.ts @@ -1,13 +1,4 @@ -/** - * Synthesized `code_execution` tool definition. - * - * The name, input shape ({code}), and description framing deliberately track - * Anthropic's server-side programmatic tool calling so models' trained - * priors transfer: python scripts, tools as async functions taking a single - * dict and returning a string, stdout as the return channel, persistent - * interpreter state, asyncio.gather for fan-out. - */ - +/** Compact agent-facing surface for Python orchestration and its execution lifecycle. */ import type { ToolDefinition } from '../types/index.js'; export const CODE_EXECUTION_TOOL_NAME = 'code_execution'; @@ -15,63 +6,39 @@ export const CODE_EXECUTION_TOOL_NAME = 'code_execution'; export function buildCodeExecutionToolDefinition(opts?: { idleReclaimMs?: number; toolCallTimeoutMs?: number; + foregroundWaitMs?: number; }): ToolDefinition { - const idleMinutes = Math.max(1, Math.round((opts?.idleReclaimMs ?? 300_000) / 60_000)); - const callTimeoutSeconds = Math.round((opts?.toolCallTimeoutMs ?? 270_000) / 1000); + const idleMs = opts?.idleReclaimMs ?? 300_000; + const reuse = idleMs === 0 ? 'until cancelled or the host stops' : `until ${idleMs}ms idle`; return { name: CODE_EXECUTION_TOOL_NAME, description: - 'Run Python code that can call your other tools programmatically. ' + - 'Every tool you have is available inside the script as an async Python function: ' + - "the function name is the tool name with '--' replaced by '__' and any other " + - "non-identifier character replaced by '_' " + - "(e.g. tool 'mcpl--discord--fetch_history' is the function mcpl__discord__fetch_history, " + - "and tool 'mcpl--dog-events--status' is mcpl__dog_events__status). " + - 'Each function takes a single dict of arguments and returns a string — the same text ' + - 'the tool would have returned to you directly; parse structured results with json.loads. ' + - "An exact-name lookup dict is also available: tools['mcpl--discord--fetch_history']({...}). " + - 'Use top-level await; run independent calls in parallel with asyncio.gather. ' + - 'Only what you print() (plus stderr and the exit code) comes back to you — intermediate ' + - 'tool results stay in the script, so filter or aggregate large data there and print only ' + - 'what you need. Interpreter state (variables, imports) persists across code_execution ' + - `calls but is reclaimed after ~${idleMinutes} minutes idle. A tool call that receives no ` + - `response within ~${callTimeoutSeconds}s raises TimeoutError inside the script. ` + - 'Use this when fanning out across many items, looping over tool calls, or when tool ' + - 'results are large and you only need a slice or summary. Call tools directly (not via ' + - 'code) when a single call answers the question or when you need to reason about each ' + - 'result before deciding the next step. ' + - 'BACKGROUND MODE: pass background=true to run the script as a detached watcher that ' + - 'outlives this turn — the tool returns immediately with a script_id and you can end ' + - 'your turn (e.g. sleep). Inside a background script, await wake_agent(payload) wakes ' + - 'you: the payload plus provenance (script id, the line number in your script, elapsed ' + - 'time) is delivered into your context and starts a turn for you. A script that ends ' + - 'without calling wake_agent wakes nobody — that silence is the point (poll cheaply, ' + - 'wake only on signal). If your background script CRASHES you are woken with the error. ' + - 'Its print() output streams to a workspace log file you can read any time. Wakes are ' + - 'rate-limited (early wakes are delayed, not dropped) and capped per script. ' + - 'CAUTION: background scripts die silently if the host process restarts — for a wake ' + - 'you absolutely must not miss, also arm a wake rule as backup. ' + - 'Manage your scripts with {"action": "list"} and {"action": "cancel", "script_id": "..."}.', + 'Run Python with top-level await. Tools are async Python functions taking one dict and returning text; ' + + "use tools['exact-tool-name']({...}) or the sanitized name ('--' becomes '__', other punctuation becomes '_'). " + + 'Use asyncio.gather for parallel calls. Tool errors return Error: text; image results become placeholders. ' + + 'Only printed stdout, stderr and return_code reach you; intermediate results stay in Python. ' + + 'Choose direct calls or code as useful, including for a single call. ' + + `The tool waits at most ${opts?.foregroundWaitMs ?? 10_000}ms by default (override with wait_ms), then returns a running script_id. ` + + 'The script continues and completion notifies you. Use action=wait with script_id to inspect or wait again; ' + + 'wait_ms=0 inspects immediately. on_timeout=end_turn ends your turn if still running, with completion waking you. ' + + 'A wait timeout does not cancel execution. ' + + `Ordinary calls share one Python interpreter; variables/imports persist ${reuse}. ` + + 'While a script is running that context is busy. background=true starts an independent interpreter, returns immediately, ' + + 'and makes await wake_agent(payload) available for explicit notifications. Clean background exits are silent unless ' + + 'you wait and yield; crashes notify you. Output is journaled to the reported workspace log when available. ' + + 'Background scripts are primary-agent-only; wakes are rate limited. action=list lists your scripts; action=cancel stops one. ' + + `Inner tool calls time out after ${opts?.toolCallTimeoutMs ?? 270_000}ms. ` + + 'Cancellation stops Python, but already dispatched tools may still finish. ' + + 'Scripts and retained results do not survive host restart; the five most recently settled results are retained per agent.', inputSchema: { type: 'object' as const, properties: { - code: { - type: 'string', - description: 'Python code to execute. Top-level await is allowed.', - }, - background: { - type: 'boolean', - description: 'Run detached as a background watcher with wake_agent() available (default false).', - }, - action: { - type: 'string', - enum: ['run', 'list', 'cancel'], - description: 'run (default) executes `code`; list shows your background scripts; cancel stops one.', - }, - script_id: { - type: 'string', - description: 'Background script id (for action: cancel).', - }, + code: { type: 'string', description: 'Python code; top-level await is supported.' }, + action: { type: 'string', enum: ['run', 'wait', 'list', 'cancel'], description: 'run (default), wait/inspect a script, list scripts, or cancel one.' }, + script_id: { type: 'string', description: 'Execution id for wait or cancel.' }, + wait_ms: { type: 'integer', description: 'Observation budget, 0–60000 milliseconds; 0 returns immediately. Does not limit script lifetime.' }, + on_timeout: { type: 'string', enum: ['continue', 'end_turn'], description: 'After the observation budget expires, continue thinking (default) or end this turn. Completion will notify/wake you.' }, + background: { type: 'boolean', description: 'Independent interpreter with wake_agent and output journal. Returns immediately unless wait_ms or on_timeout is supplied.' }, }, required: [], }, diff --git a/src/framework.ts b/src/framework.ts index db55fc9..b4c4f93 100644 --- a/src/framework.ts +++ b/src/framework.ts @@ -1,3 +1,4 @@ +import { ScriptRun, type TimeoutPolicy } from './code-execution/script-run.js'; import { join } from 'node:path'; import { INLINE_WITHHELD_TEXT, classifyBlock, isInlineContradiction, referenceRegistry, referenceStubOrNull } from './mcpl/references.js'; import { ReferenceFetcher, DEFAULT_FETCH_MAX_BYTES, EAGER_FETCH_TIMEOUT_MS } from './mcpl/reference-fetcher.js'; @@ -640,9 +641,14 @@ interface Deferred { * the caller awaits plus the liveness/progress bookkeeping the watchdogs and * the stream driver share. One map entry per run — created and torn down in * runEphemeralToCompletion — so the pieces cannot desync. */ -/** One model-authored background (daemon) script — see runCodeExecution. */ -interface BackgroundScriptRecord { +/** A Python execution, independent of the tool call currently observing it. */ +interface CodeExecutionRecord { id: string; + mode: 'foreground' | 'background'; + run: ScriptRun; + endTurn: boolean; + wakeQueue: Promise; + wakeAbort: AbortController; agentName: string; runner: PyRunner; /** The agent-authored source (line numbers in wake envelopes index into it). */ @@ -955,15 +961,11 @@ export class AgentFramework { private codeExecutionConfig: import('./types/index.js').CodeExecutionConfig | null = null; private codeExecutionRunners: Map = new Map(); private scriptToolWaiters: Map void> = new Map(); - /** Agents whose running script hit an endTurn-carrying inner result — - * deferred and applied to the final code_execution result instead of - * cancelling the stream mid-script (which would wedge the turn). */ - private scriptDeferredEndTurn: Set = new Set(); - /** Background (daemon) scripts: model-authored watchers that outlive their - * spawning turn. Each gets a DEDICATED PyRunner; wake_agent() injects a - * provenance envelope + payload and requests inference. Keyed by script id. */ - private backgroundScripts: Map = new Map(); - private backgroundScriptCounter = 0; + /** Running executions and the five most recently settled results per owner. + * Foreground runs share their owner's interpreter; background runs have a + * dedicated one. Observation and wake state belong to the run, not a turn. */ + private codeExecutionScripts: Map = new Map(); + private codeExecutionScriptCounter = 0; /** Resident-set per-agent tool-result inline cap (chars), via * agent_settings `tool_result_inline_max_chars`. DURABLE: persisted in * framework state like the core runtime settings and restored at create @@ -1419,21 +1421,15 @@ export class AgentFramework { // Kill running code_execution scripts before cancelling streams: a // zombie script must not keep firing side-effectful tool calls into a // framework that is shutting down. - for (const runner of this.codeExecutionRunners.values()) { - runner.dispose(); - } - this.codeExecutionRunners.clear(); - // Background daemons die with the process (documented limitation) — - // mark cancelled FIRST so their settle path stays silent (no crash-wake - // into a framework that is shutting down). - for (const record of this.backgroundScripts.values()) { - if (record.status === 'running') { - record.cancelled = true; - record.status = 'cancelled'; - } + for (const record of this.codeExecutionScripts.values()) { + record.cancelled = true; + if (record.status === 'running') record.status = 'cancelled'; + record.wakeAbort.abort(); record.runner.dispose(); } - this.backgroundScripts.clear(); + this.codeExecutionScripts.clear(); + for (const runner of this.codeExecutionRunners.values()) runner.dispose(); + this.codeExecutionRunners.clear(); // Retained tool images are per-process by design; a retained stopped // framework must not keep up to the whole ledger budget referenced. this.toolImageLedgers.clear(); @@ -2743,6 +2739,7 @@ export class AgentFramework { // The late stream finalizer intentionally refuses name-keyed cleanup once // deregistered, so disposal must release EventGate liveness itself. this.eventGate?.onInferenceEnded(agent.name); + this.disposeAgentScripts(agent.name); this.agents.delete(agent.name); } this.toolImageLedgers.delete(agent.name); @@ -5156,6 +5153,7 @@ export class AgentFramework { this.closingConversationAgents.delete(agentName); const channelId = this.conversationAgentHomes.get(agentName); const agent = this.agents.get(agentName); + this.disposeAgentScripts(agentName); this.agents.delete(agentName); this.toolImageLedgers.delete(agentName); this.agentConfigs.delete(agentName); @@ -8090,36 +8088,54 @@ export class AgentFramework { background?: unknown; action?: unknown; script_id?: unknown; + wait_ms?: unknown; + on_timeout?: unknown; }; - // Management surface: the agent's own daemon fleet is inspectable and + const waitMs = input.wait_ms ?? this.codeExecutionConfig?.foregroundWaitMs ?? 10_000; + if (typeof waitMs !== 'number' || !Number.isInteger(waitMs) || waitMs < 0 || waitMs > 60_000) { + return { success: false, error: 'wait_ms must be an integer from 0 to 60000', isError: true }; + } + const onTimeout = input.on_timeout ?? 'continue'; + if (onTimeout !== 'continue' && onTimeout !== 'end_turn') { + return { success: false, error: 'on_timeout must be continue or end_turn', isError: true }; + } + + if (onTimeout === 'end_turn' && (!this.agents.has(agentName) || this.ephemeralRuns.has(agentName))) { + return { success: false, isError: true, error: 'end_turn requires a persistent registered agent that can receive the completion wake' }; + } + + // Management surface: the agent's own execution fleet is inspectable and // killable through the same tool that spawns it. if (input.action === 'list') { - const scripts = [...this.backgroundScripts.values()] + const scripts = [...this.codeExecutionScripts.values()] .filter((s) => s.agentName === agentName) .map((s) => ({ script_id: s.id, status: s.status, + mode: s.mode, started_at: new Date(s.startedAt).toISOString(), wakes: s.wakes, last_wake_at: s.lastWakeAt ? new Date(s.lastWakeAt).toISOString() : null, log: s.logPath, })); - return { success: true, data: { background_scripts: scripts } }; + return { success: true, data: { scripts, background_scripts: scripts.filter(s => s.mode === 'background') } }; } - if (input.action === 'cancel') { + if (input.action === 'cancel' || input.action === 'wait') { if (typeof input.script_id !== 'string') { - return { success: false, error: 'cancel requires `script_id`', isError: true }; + return { success: false, error: `${input.action} requires \`script_id\``, isError: true }; } - const record = this.backgroundScripts.get(input.script_id); + const record = this.codeExecutionScripts.get(input.script_id); if (!record || record.agentName !== agentName) { - return { success: false, error: `no background script '${input.script_id}'`, isError: true }; + return { success: false, error: `no script '${input.script_id}'`, isError: true }; } + if (input.action === 'wait') return this.observeCodeExecution(record, waitMs, onTimeout); if (record.status === 'running') { + record.wakeAbort.abort(); record.cancelled = true; record.status = 'cancelled'; record.runner.abort('cancelled by agent'); - record.runner.dispose(); + if (record.mode === 'background') record.runner.dispose(); } return { success: true, @@ -8147,22 +8163,73 @@ export class AgentFramework { ); if (input.background === true) { - return this.startBackgroundScript(agentName, input.code, injected); + const started = this.startBackgroundScript(agentName, input.code, injected); + if (!started.success || (input.wait_ms === undefined && input.on_timeout === undefined)) return started; + const id = (started.data as { script_id: string }).script_id; + return this.observeCodeExecution(this.codeExecutionScripts.get(id)!, waitMs, onTimeout); } + const active = [...this.codeExecutionScripts.values()].find(s => + s.agentName === agentName && s.mode === 'foreground' && s.status === 'running'); + if (active) { + return { success: false, isError: true, + error: `Python context is busy with ${active.id}; wait/cancel it, or use background=true for an independent interpreter`, + data: { script_id: active.id, status: active.status } }; + } const runner = this.getOrCreateScriptRunner(agentName); - this.scriptDeferredEndTurn.delete(agentName); - const exec = await runner.exec(input.code, injected); - const endTurn = this.scriptDeferredEndTurn.delete(agentName); + const record: CodeExecutionRecord = { + id: `py-${++this.codeExecutionScriptCounter}`, agentName, runner, + mode: 'foreground', code: input.code, startedAt: Date.now(), + wakes: 0, lastWakeAt: null, logPath: null, status: 'running', cancelled: false, + endTurn: false, wakeQueue: Promise.resolve(), wakeAbort: new AbortController(), + run: undefined!, + }; + this.codeExecutionScripts.set(record.id, record); + this.emitTrace({ type: 'tool:started', module: 'code_execution', tool: `foreground:${record.id}`, + callId: record.id, input: { lines: record.code.split('\n').length } }); + record.run = new ScriptRun(runner.exec(input.code, injected), (exec, notify, observed) => + this.settleCodeExecution(record, exec, notify, observed)); + return this.observeCodeExecution(record, waitMs, onTimeout); + } + private async observeCodeExecution( + record: CodeExecutionRecord, waitMs: number, onTimeout: TimeoutPolicy, + ): Promise { + const observed = await record.run.observe(waitMs, onTimeout); + const exec = observed.result; + const endTurn = observed.endTurn || record.endTurn; + record.endTurn = false; return { - success: true, - data: { stdout: exec.stdout, stderr: exec.stderr, return_code: exec.returnCode }, - isError: false, + success: true, isError: false, + data: { + script_id: record.id, status: exec ? record.status : 'running', mode: record.mode, + ...(exec ? { stdout: exec.stdout, stderr: exec.stderr, return_code: exec.returnCode, + ...(exec.aborted ? { aborted: true } : {}), ...(exec.tail ? { tail: exec.tail } : {}) } + : { note: 'Still running. Completion will notify you; use action=wait to retrieve the result. Waiting does not cancel execution.' }), + ...(record.logPath ? { log: record.logPath } : {}), + }, ...(endTurn ? { endTurn: true } : {}), }; } + /** Release a stuck observation through its normal tool-result boundary. + * The script survives; its completion wakes its owner. No stream surgery. + */ + releaseCodeExecutionWait(agentName: string, scriptId?: string): { released: number } { + if (!this.agents.has(agentName) || this.ephemeralRuns.has(agentName)) { + throw new Error('Releasing a wait requires a persistent registered agent'); + } + const records = [...this.codeExecutionScripts.values()].filter(s => s.agentName === agentName); + if (scriptId && !records.some(s => s.id === scriptId)) { + throw new Error(`No script '${scriptId}' for '${agentName}'`); + } + let released = 0; + for (const record of records) { + if (!scriptId || record.id === scriptId) released += record.run.release(); + } + return { released }; + } + // --------------------------------------------------------------------------- // Background (daemon) scripts — model-authored watchers that outlive the turn // --------------------------------------------------------------------------- @@ -8181,8 +8248,8 @@ export class AgentFramework { injected: import('./code-execution/py-runner.js').InjectedTool[], ): ToolResult { // v1: primary-agent-only. A conversation fork's or ephemeral's daemon - // would outlive its owner and its wake would land in the primary - // conversation anyway — refuse loudly instead of surprising anyone. + // would outlive its owner. Keep the existing restriction rather than + // promising a watcher whose lifetime depends on a short-lived fork. if (this.primaryAgentName && agentName !== this.primaryAgentName) { return { success: false, @@ -8192,8 +8259,8 @@ export class AgentFramework { } const cfg = this.codeExecutionConfig; const maxScripts = cfg?.maxBackgroundScripts ?? 3; - const runningCount = [...this.backgroundScripts.values()] - .filter((s) => s.agentName === agentName && s.status === 'running').length; + const runningCount = [...this.codeExecutionScripts.values()] + .filter((s) => s.agentName === agentName && s.mode === 'background' && s.status === 'running').length; if (runningCount >= maxScripts) { return { success: false, @@ -8203,7 +8270,7 @@ export class AgentFramework { }; } - const scriptId = `bg-${++this.backgroundScriptCounter}`; + const scriptId = `bg-${++this.codeExecutionScriptCounter}`; const lifetimeMs = cfg?.backgroundMaxLifetimeMs ?? 86_400_000; // Journal: a file under the agent's first read-write workspace mount so @@ -8230,8 +8297,13 @@ export class AgentFramework { onToolCall: (toolName, args) => this.handleScriptToolCall(agentName, toolName, args), }); - const record: BackgroundScriptRecord = { + const record: CodeExecutionRecord = { id: scriptId, + mode: 'background', + run: undefined!, + endTurn: false, + wakeQueue: Promise.resolve(), + wakeAbort: new AbortController(), agentName, runner, code, @@ -8242,26 +8314,17 @@ export class AgentFramework { status: 'running', cancelled: false, }; - this.backgroundScripts.set(scriptId, record); + this.codeExecutionScripts.set(scriptId, record); this.emitTrace({ type: 'tool:started', module: 'code_execution', tool: `background:${scriptId}`, callId: scriptId, input: { lines: code.split('\n').length }, }); - void runner - .exec(code, injected, { - logPath: logAbsPath, - lifetimeMs, - onWake: (line, payload) => this.handleScriptWake(record, line, payload), - }) - .then((exec) => this.settleBackgroundScript(record, exec)) - .catch((err) => { - // exec never rejects by contract; this is the belt-and-suspenders. - console.error(`[pytc:${agentName}:${scriptId}] background exec rejected: ${String(err)}`); - this.settleBackgroundScript(record, { - stdout: '', stderr: String(err), returnCode: 1, aborted: true, - }); - }); + record.run = new ScriptRun(runner.exec(code, injected, { + logPath: logAbsPath, + lifetimeMs, + onWake: (line, payload) => this.handleScriptWake(record, line, payload), + }), (exec, notify, observed) => this.settleCodeExecution(record, exec, notify, observed)); return { success: true, @@ -8295,11 +8358,17 @@ export class AgentFramework { * RuntimeError inside the script). */ private async handleScriptWake( - record: BackgroundScriptRecord, + record: CodeExecutionRecord, line: number, payload: unknown, ): Promise { - if (record.cancelled) return 'script was cancelled'; + const delivery = record.wakeQueue.then(() => this.deliverScriptWake(record, line, payload)); + record.wakeQueue = delivery.catch(() => {}); + return delivery; + } + + private async deliverScriptWake(record: CodeExecutionRecord, line: number, payload: unknown): Promise { + if (record.cancelled || record.status !== 'running') return 'script is no longer running'; const cfg = this.codeExecutionConfig; const maxWakes = cfg?.maxWakesPerScript ?? 100; if (record.wakes >= maxWakes) { @@ -8309,15 +8378,15 @@ export class AgentFramework { if (record.lastWakeAt !== null) { const wait = record.lastWakeAt + floorMs - Date.now(); if (wait > 0) { - await new Promise((r) => { - const t = setTimeout(r, wait); - (t as { unref?: () => void }).unref?.(); + await new Promise((resolve) => { + const finish = () => { clearTimeout(timer); record.wakeAbort.signal.removeEventListener('abort', finish); resolve(); }; + const timer = setTimeout(finish, wait); + timer.unref?.(); + record.wakeAbort.signal.addEventListener('abort', finish, { once: true }); }); - if (record.cancelled) return 'script was cancelled'; + if (record.cancelled || record.status !== 'running') return 'script is no longer running'; } } - record.wakes += 1; - record.lastWakeAt = Date.now(); const elapsedMin = Math.round((Date.now() - record.startedAt) / 60_000); const payloadJson = payload === undefined || payload === null @@ -8330,51 +8399,62 @@ export class AgentFramework { ? this.resolveToolResultInlineCap(agent).cap : DEFAULT_TOOL_RESULT_INLINE_MAX_CHARS; const { text: payloadText } = await this.spillOrTruncate( - payloadJson, maxChars, `${record.id}-wake${record.wakes}`, + payloadJson, maxChars, `${record.id}-wake${record.wakes + 1}`, ); const envelope = `[background script ${record.id}] Woke you: wake_agent() called at line ${line} of your script, ` + - `${elapsedMin}m after you started it (wake ${record.wakes} of ${maxWakes}).\n` + + `${elapsedMin}m after you started it (wake ${record.wakes + 1} of ${maxWakes}).\n` + `Arguments:\n${payloadText}\n` + (record.logPath ? `Script output so far: workspace file ${record.logPath}\n` : '') + `Script status: still running`; - this.injectScriptWake(record, envelope); - return null; + if (record.cancelled || record.status !== 'running') return 'script is no longer running'; + const error = this.injectScriptWake(record, envelope); + if (!error) { record.wakes += 1; record.lastWakeAt = Date.now(); } + return error; } - /** Background script ended (clean, crash, deadline, cancel/dispose). */ - private settleBackgroundScript( - record: BackgroundScriptRecord, + /** Execution ended (clean, crash, deadline, cancel/dispose). */ + private settleCodeExecution( + record: CodeExecutionRecord, exec: import('./code-execution/py-runner.js').ExecResult, + notify = false, + observed = false, ): void { - record.runner.dispose(); + record.wakeAbort.abort(); + if (record.mode === 'background') record.runner.dispose(); if (record.status === 'running') { record.status = exec.returnCode === 0 ? 'finished' : 'died'; } this.emitTrace({ - type: 'tool:completed', module: 'code_execution', tool: `background:${record.id}`, + type: 'tool:completed', module: 'code_execution', tool: `${record.mode}:${record.id}`, callId: record.id, durationMs: Date.now() - record.startedAt, }); - // Crash honesty: a resident sleeping on a promise must never have that - // promise die silently. Deliberate cancels/disposes stay silent (the - // agent or host chose them); clean exits stay silent (the null path). - if (record.status === 'died' && !record.cancelled) { - const tail = (exec.tail ?? exec.stderr ?? '').slice(-2000); + // An unobserved completion after yielding wakes once. Legacy watchers + // remain silent on clean exit unless observed/yielded; crashes wake. + // Explicit cancellation and owner disposal suppress notifications. + if ((notify || (record.mode === 'background' && record.status === 'died' && !observed)) && !record.cancelled) { + const tail = (exec.tail || exec.stderr || exec.stdout).slice(-2000); const elapsedMin = Math.round((Date.now() - record.startedAt) / 60_000); const envelope = - `[background script ${record.id}] Your background script DIED ${elapsedMin}m after start ` + - `(it did not call wake_agent and is no longer watching).\n` + + `[${record.mode} script ${record.id}] Script ${record.status === 'died' ? 'DIED' : 'finished'} ${elapsedMin}m after start ` + + `(return_code=${exec.returnCode}). Retrieve the retained result with code_execution action=wait, script_id=${record.id}.\n` + `Last output:\n${tail || '(none)'}\n` + (record.logPath ? `Full journal: workspace file ${record.logPath}` : ''); this.injectScriptWake(record, envelope); } - // Keep a short memory of settled scripts for {"action":"list"}, then drop. - const settled = [...this.backgroundScripts.values()] + // Refresh insertion order on settlement so a long-running older job + // isn't immediately evicted by five newer jobs that finished before it. + if (this.codeExecutionScripts.has(record.id)) { + this.codeExecutionScripts.delete(record.id); + this.codeExecutionScripts.set(record.id, record); + } + // Keep a short memory of settled scripts for list/wait, then drop. + const settled = [...this.codeExecutionScripts.values()] .filter((s) => s.agentName === record.agentName && s.status !== 'running'); for (const old of settled.slice(0, Math.max(0, settled.length - 5))) { - this.backgroundScripts.delete(old.id); + this.codeExecutionScripts.delete(old.id); } } @@ -8385,13 +8465,15 @@ export class AgentFramework { * authority); it enters tagged so gate policies could be taught about it * later if that ever needs revisiting. */ - private injectScriptWake(record: BackgroundScriptRecord, envelope: string): void { + private injectScriptWake(record: CodeExecutionRecord, envelope: string): string | null { + if (!this.agents.has(record.agentName)) return 'script owner is no longer registered'; try { const messageId = this.addMessage('user', [{ type: 'text', text: envelope }], { source: 'background-script', scriptId: record.id, tags: ['script:wake'], - }); + system: true, + }, { forAgent: record.agentName }); this.pendingRequests.push({ agentName: record.agentName, reason: 'script:wake', @@ -8403,10 +8485,11 @@ export class AgentFramework { messageId, source: `background-script:${record.id}`, }); + return null; } catch (err) { - console.error( - `[pytc:${record.agentName}:${record.id}] wake injection failed: ${err instanceof Error ? err.message : String(err)}`, - ); + const error = `wake injection failed: ${err instanceof Error ? err.message : String(err)}`; + console.error(`[pytc:${record.agentName}:${record.id}] ${error}`); + return error; } } @@ -8550,7 +8633,11 @@ export class AgentFramework { scriptTimeoutMs: cfg?.scriptTimeoutMs, idleReclaimMs: cfg?.idleReclaimMs, label: agentName, - onToolCall: (toolName, args) => this.handleScriptToolCall(agentName, toolName, args), + onToolCall: (toolName, args) => { + const record = [...this.codeExecutionScripts.values()].find(s => + s.agentName === agentName && s.mode === 'foreground' && s.status === 'running'); + return this.handleScriptToolCall(agentName, toolName, args, () => { if (record) record.endTurn = true; }); + }, }); this.codeExecutionRunners.set(agentName, runner); } @@ -8570,17 +8657,15 @@ export class AgentFramework { agentName: string, toolName: string, args: Record, + onEndTurn?: () => void, ): Promise { if (toolName === CODE_EXECUTION_TOOL_NAME) { return 'Error: code_execution cannot be called from within a script'; } const result = await this.dispatchScriptToolCall(agentName, toolName, args); - if (result.endTurn) { - // Deferred: applied to the final code_execution result (see - // scriptDeferredEndTurn) — ending the turn mid-script would cancel the - // stream while the script still runs and wedge the tool round. - this.scriptDeferredEndTurn.add(agentName); - } + // Bound to this execution; background or late inner results must never + // request endTurn on an unrelated later foreground script. + if (result.endTurn) onEndTurn?.(); if (result.isError) { return `Error: ${result.error ?? 'tool call failed'}`; } @@ -8651,12 +8736,24 @@ export class AgentFramework { }); } - /** Abort a running script when the agent's turn dies underneath it. */ - private abortAgentScript(agentName: string, reason: string): void { - const runner = this.codeExecutionRunners.get(agentName); - if (runner?.busy) { - runner.abort(reason); + private disposeAgentScripts(agentName: string): void { + for (const record of this.codeExecutionScripts.values()) { + if (record.agentName !== agentName) continue; + record.cancelled = true; + record.status = 'cancelled'; + record.wakeAbort.abort(); + record.runner.dispose(); + this.codeExecutionScripts.delete(record.id); } + this.codeExecutionRunners.get(agentName)?.dispose(); + this.codeExecutionRunners.delete(agentName); + } + + /** Abort a failed observation; unrelated later stream failures leave yielded work alone. */ + private abortAgentScript(agentName: string, reason: string): void { + const record = [...this.codeExecutionScripts.values()].find(s => + s.agentName === agentName && s.mode === 'foreground' && s.status === 'running'); + if (record?.run.observing) record.runner.abort(reason); } private toMembraneToolResult( diff --git a/src/types/framework.ts b/src/types/framework.ts index e5c73ae..8938ac7 100644 --- a/src/types/framework.ts +++ b/src/types/framework.ts @@ -50,6 +50,8 @@ export interface CodeExecutionConfig { * script (default 270_000 ms, mirroring the managed runtime's message). */ toolCallTimeoutMs?: number; + /** Default observation budget; then return a running script_id (default 10_000, max 60_000 ms). */ + foregroundWaitMs?: number; /** Whole-script deadline: cancel → grace → SIGKILL (default 600_000 ms). */ scriptTimeoutMs?: number; /** diff --git a/test/code-execution.test.ts b/test/code-execution.test.ts index c9bbd29..e8b9f0d 100644 --- a/test/code-execution.test.ts +++ b/test/code-execution.test.ts @@ -261,6 +261,73 @@ describe('PyRunner (real python3)', () => { } }); + it('reserves the runner during cold startup and allows immediate cancellation', async () => { + const runner = new PyRunner({ onToolCall: async () => '' }); + try { + const first = runner.exec('import asyncio\nawait asyncio.sleep(60)', []); + assert.strictEqual(runner.busy, true, 'startup must count as busy'); + const second = await runner.exec('print("must not run")', []); + assert.match(second.stderr, /already running/); + runner.abort('startup cancelled'); + assert.strictEqual((await first).aborted, true); + const next = await runner.exec('print("fresh")', []); + assert.strictEqual(next.returnCode, 0, next.stderr); + assert.match(next.stdout, /fresh/); + } finally { runner.dispose(); } + }); + + it('disposal during startup settles promptly', async () => { + const runner = new PyRunner({ onToolCall: async () => '' }); + const result = runner.exec('print("must not run")', []); + runner.dispose(); + assert.strictEqual((await result).aborted, true); + assert.strictEqual(runner.busy, false); + }); + + it('keeps sub-second inner-tool timeouts instead of rounding them to zero', async () => { + const runner = new PyRunner({ + toolCallTimeoutMs: 50, + onToolCall: async () => { await new Promise(r => setTimeout(r, 300)); return 'too late'; }, + }); + try { + const result = await runner.exec('print(await test__echo({}))', ECHO_TOOLS); + assert.strictEqual(result.returnCode, 1); + assert.match(result.stderr, /TimeoutError/); + } finally { runner.dispose(); } + }); + + it('does not deliver a late tool result to a replacement interpreter', async () => { + let resolveOld!: (value: string) => void; + let resolveNew!: (value: string) => void; + let startedOld!: () => void; + let startedNew!: () => void; + const oldStarted = new Promise(r => { startedOld = r; }); + const newStarted = new Promise(r => { startedNew = r; }); + let calls = 0; + const runner = new PyRunner({ onToolCall: async () => { + if (++calls === 1) { + startedOld(); + return new Promise(r => { resolveOld = r; }); + } + startedNew(); + return new Promise(r => { resolveNew = r; }); + } }); + try { + const old = runner.exec('print(await test__echo({}))', ECHO_TOOLS); + await oldStarted; + runner.abort('replace interpreter'); + await old; + const fresh = runner.exec('print(await test__echo({}))', ECHO_TOOLS); + await newStarted; + resolveOld('STALE RESULT'); + // Give the old reply a chance to reach the new interpreter's t1 call. + await new Promise(r => setTimeout(r, 100)); + resolveNew('CURRENT RESULT'); + const result = await fresh; + assert.strictEqual(result.stdout.trim(), 'CURRENT RESULT'); + } finally { runner.dispose(); } + }); + it('fails gracefully when the python binary is missing', async () => { const runner = new PyRunner({ pythonPath: '/definitely/not/a/python', @@ -582,6 +649,144 @@ describe('framework background scripts + spill (real python3)', () => { return framework; } + it('yields foreground await, retains its result, and preserves interpreter state', async () => { + const { tempDir, storePath } = tempStorePath('pytc-yield-'); + const framework = await createBgFramework(storePath, new MockMembrane(), { + codeExecution: { foregroundWaitMs: 5 }, + }); + const call = (input: Record) => framework.executeToolCall({ + id: `test-${Math.random()}`, name: 'code_execution', callerAgentName: 'prime', input, + }); + try { + const started = await call({ code: 'import asyncio\nx = 41\nawait asyncio.sleep(0.25)\nprint("finished")' }); + const id = (started.data as { script_id: string }).script_id; + assert.equal((started.data as { status: string }).status, 'running'); + assert.equal(started.endTurn, undefined); + const busy = await call({ code: 'print("must not run")' }); + assert.equal(busy.isError, true); + assert.match(busy.error ?? '', new RegExp(id)); + const done = await call({ action: 'wait', script_id: id, wait_ms: 3000 }); + assert.equal((done.data as { stdout: string }).stdout.trim(), 'finished'); + assert.equal((done.data as { return_code: number }).return_code, 0); + assert.equal((done.data as { status: string }).status, 'finished'); + const reused = await call({ code: 'print(x + 1)', wait_ms: 3000 }); + assert.equal((reused.data as { stdout: string }).stdout.trim(), '42'); + const reread = await call({ action: 'wait', script_id: id, wait_ms: 0 }); + assert.equal((reread.data as { stdout: string }).stdout.trim(), 'finished'); + const wrongOwner = await framework.executeToolCall({ + id: 'wrong-owner', name: 'code_execution', callerAgentName: 'someone-else', + input: { action: 'wait', script_id: id }, + }); + assert.equal(wrongOwner.isError, true); + } finally { await framework.stop(); rmSync(tempDir, { recursive: true, force: true }); } + }); + + for (const releaseBy of ['agent', 'operator'] as const) { + it(`${releaseBy} can end a waiting turn while the script survives to wake it`, async () => { + const { tempDir, storePath } = tempStorePath('pytc-wait-wake-'); + const membrane = new MockMembrane(); + membrane.pushResponse(createMockResponse([{ + type: 'tool_use', id: 'wait-call', name: 'code_execution', input: { + code: 'import asyncio\nawait asyncio.sleep(0.4)\nprint("completion payload")', + wait_ms: releaseBy === 'agent' ? 0 : 60_000, + on_timeout: releaseBy === 'agent' ? 'end_turn' : 'continue', + }, + }], 'tool_use')); + membrane.pushResponse(createMockResponse([{ type: 'text', text: 'I saw completion.' }])); + const framework = await createBgFramework(storePath, membrane); + try { + framework.start(); + framework.pushEvent({ type: 'external-message', source: 'test', + content: [{ type: 'text', text: 'run and rest' }], metadata: {}, triggerInference: true } as ProcessEvent); + const deadline = Date.now() + 5000; + if (releaseBy === 'operator') { + let released = false; + while (!released && Date.now() < deadline) { + const list = await framework.executeToolCall({ id: 'list', name: 'code_execution', + callerAgentName: 'prime', input: { action: 'list' } }); + const scripts = (list.data as { scripts: { script_id: string; status: string }[] }).scripts; + if (scripts[0]?.status === 'running') { + assert.deepEqual(framework.releaseCodeExecutionWait('prime'), { released: 1 }); + released = true; + } else await new Promise(r => setTimeout(r, 5)); + } + assert.ok(released, 'operator should find an active wait'); + } + while (membrane.calls.length < 2 && Date.now() < deadline) await new Promise(r => setTimeout(r, 10)); + assert.equal(membrane.calls.length, 2, 'completion must start a second inference'); + assert.match(JSON.stringify(membrane.calls[1]), /completion payload/); + const messages = framework.getAgent('prime')!.getContextManager().queryMessages({}).messages; + const results = messages.flatMap(m => m.content).filter(b => b.type === 'tool_result'); + assert.ok(results.some(b => JSON.parse((b as { content: string }).content).status === 'running')); + // The second response must not be consumed by a continuation BEFORE completion. + assert.match(JSON.stringify(messages), /I saw completion|completion payload/); + } finally { await framework.stop(); rmSync(tempDir, { recursive: true, force: true }); } + }); + } + + it('cancels a yielded foreground run and reports its retained cancellation', async () => { + const { tempDir, storePath } = tempStorePath('pytc-yield-cancel-'); + const framework = await createBgFramework(storePath, new MockMembrane()); + const call = (input: Record) => framework.executeToolCall({ + id: 'test', name: 'code_execution', callerAgentName: 'prime', input, + }); + try { + const started = await call({ code: 'import asyncio\nawait asyncio.sleep(60)', wait_ms: 0 }); + const id = (started.data as { script_id: string }).script_id; + await call({ action: 'cancel', script_id: id }); + const result = await call({ action: 'wait', script_id: id }); + assert.equal((result.data as { status: string }).status, 'cancelled'); + assert.equal((result.data as { aborted: boolean }).aborted, true); + const next = await call({ code: 'print("fresh context")' }); + assert.match((next.data as { stdout: string }).stdout, /fresh context/); + } finally { await framework.stop(); rmSync(tempDir, { recursive: true, force: true }); } + }); + + it('serializes concurrent wake requests across the rate floor and cap', async () => { + const { tempDir, storePath } = tempStorePath('pytc-wake-limit-'); + const framework = await createBgFramework(storePath, new MockMembrane(), { + codeExecution: { wakeMinIntervalMs: 30, maxWakesPerScript: 2 }, + }); + const times: number[] = []; + const delivery = framework as unknown as { addMessage: (...args: unknown[]) => string }; + const add = delivery.addMessage.bind(framework); + delivery.addMessage = (...args) => { times.push(Date.now()); return add(...args); }; + try { + const result = await framework.executeToolCall({ id: 'wake-limit', name: 'code_execution', callerAgentName: 'prime', + input: { background: true, wait_ms: 3000, code: + 'import asyncio\nawait wake_agent("first")\nr = await asyncio.gather(wake_agent("second"), wake_agent("third"), return_exceptions=True)\nprint(r)' } }); + assert.match((result.data as { tail: string }).tail, /wake limit reached/); + assert.equal(times.length, 2); + assert.ok(times[1] - times[0] >= 25, `wakes were only ${times[1] - times[0]}ms apart`); + } finally { await framework.stop(); rmSync(tempDir, { recursive: true, force: true }); } + }); + + it('refuses a wake if context delivery failed instead of falsely acknowledging it', async () => { + const { tempDir, storePath } = tempStorePath('pytc-wake-fail-'); + const framework = await createBgFramework(storePath, new MockMembrane()); + (framework as unknown as { addMessage: () => string }).addMessage = () => { throw new Error('storage unavailable'); }; + try { + const result = await framework.executeToolCall({ id: 'wake-fail', name: 'code_execution', callerAgentName: 'prime', + input: { background: true, wait_ms: 3000, code: + 'try:\n await wake_agent("signal")\nexcept RuntimeError as e:\n print(str(e))' } }); + assert.match((result.data as { tail: string }).tail, /wake injection failed: storage unavailable/); + } finally { await framework.stop(); rmSync(tempDir, { recursive: true, force: true }); } + }); + + it('returns an observed background failure without an extra crash wake', async () => { + const { tempDir, storePath } = tempStorePath('pytc-observed-crash-'); + const framework = await createBgFramework(storePath, new MockMembrane()); + let notifications = 0; + (framework as unknown as { addMessage: () => string }).addMessage = () => { notifications++; return ''; }; + try { + const result = await framework.executeToolCall({ id: 'observed-crash', name: 'code_execution', callerAgentName: 'prime', + input: { background: true, wait_ms: 3000, code: 'raise ValueError("observed failure")' } }); + assert.equal((result.data as { return_code: number }).return_code, 1); + assert.match((result.data as { tail: string }).tail, /observed failure/); + assert.equal(notifications, 0); + } finally { await framework.stop(); rmSync(tempDir, { recursive: true, force: true }); } + }); + it('background script detaches, wakes the agent with provenance, and triggers inference', async () => { const { tempDir, storePath } = tempStorePath('pytc-bgfw-'); const mountDir = join(tempDir, 'mount'); diff --git a/test/script-run.test.ts b/test/script-run.test.ts new file mode 100644 index 0000000..e37fc9a --- /dev/null +++ b/test/script-run.test.ts @@ -0,0 +1,64 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { ScriptRun } from '../src/code-execution/script-run.js'; +import type { ExecResult } from '../src/code-execution/py-runner.js'; + +function harness() { + let complete!: (result: ExecResult) => void; + const notifications: boolean[] = []; + const result = { stdout: 'done', stderr: '', returnCode: 0 }; + const run = new ScriptRun(new Promise(r => { complete = r; }), + (_result, notify) => { notifications.push(notify); }); + return { run, complete: () => complete(result), notifications, result }; +} + +describe('ScriptRun observation lifecycle', () => { + it('completion within budget returns directly and does not request another turn', async () => { + const h = harness(); + const wait = h.run.observe(10_000, 'end_turn'); + h.complete(); + assert.deepEqual(await wait, { result: h.result, endTurn: false }); + assert.deepEqual(h.notifications, [false]); + }); + + it('zero-budget yield arms completion before returning endTurn', async () => { + const h = harness(); + assert.deepEqual(await h.run.observe(0, 'end_turn'), { result: undefined, endTurn: true }); + h.complete(); + await Promise.resolve(); + assert.deepEqual(h.notifications, [true]); + assert.deepEqual(await h.run.observe(0, 'continue'), { result: h.result, endTurn: false }); + assert.deepEqual(h.notifications, [true], 'retrieval must not wake again'); + }); + + it('a timed wait releases inference without settling execution', async () => { + const h = harness(); + assert.deepEqual(await h.run.observe(10, 'continue'), { result: undefined, endTurn: false }); + assert.equal(h.run.result, undefined); + h.complete(); + await Promise.resolve(); + assert.deepEqual(h.notifications, [true]); + }); + + it('rejoining before completion consumes the result without an extra wake', async () => { + const h = harness(); + await h.run.observe(0, 'continue'); + const wait = h.run.observe(10_000, 'end_turn'); + h.complete(); + assert.equal((await wait).result, h.result); + assert.deepEqual(h.notifications, [false]); + }); + + it('operator release ends all active observations and leaves execution alive', async () => { + const h = harness(); + const a = h.run.observe(60_000, 'continue'); + const b = h.run.observe(60_000, 'continue'); + assert.equal(h.run.release(), 2); + assert.equal(h.run.release(), 0); + assert.equal((await a).endTurn, true); + assert.equal((await b).endTurn, true); + h.complete(); + await Promise.resolve(); + assert.deepEqual(h.notifications, [true]); + }); +});