From 8a5b3e9552245980b9a3a049511bfb56b3d5ddcb Mon Sep 17 00:00:00 2001 From: can1357 Date: Tue, 26 May 2026 14:37:56 +0200 Subject: [PATCH] feat(eval): added shared executor inheritance for subagents with concurrent async cells - Removed per-session run queues from JS and Python backends, allowing async cells on the same session id to interleave. - Introduced `getEvalSessionId` on ToolSession so subagents spawned via `task` inherit the parent's executor id and share JS VM and Python kernel state. - Switched JS runtime state from module-level fields to AsyncLocalStorage so concurrent runs route output and tool calls to their own context. - Changed Python runner to an asyncio event loop with per-request tasks and ContextVar-based run id tracking for concurrent execution. - Added mtime-based module cache eviction to preserve singleton state across re-imports of unchanged local files. --- docs/tools/eval.md | 44 ++- packages/coding-agent/CHANGELOG.md | 44 +-- .../eval/__tests__/shared-executors.test.ts | 313 ++++++++++++++++++ .../src/eval/js/context-manager.ts | 33 +- .../src/eval/js/shared/runtime.ts | 147 +++++--- .../coding-agent/src/eval/js/worker-core.ts | 63 ++-- packages/coding-agent/src/eval/py/executor.ts | 86 +++-- packages/coding-agent/src/eval/py/kernel.ts | 3 +- packages/coding-agent/src/eval/py/prelude.py | 4 +- packages/coding-agent/src/eval/py/runner.py | 241 +++++++++----- .../coding-agent/src/eval/py/tool-bridge.ts | 27 +- packages/coding-agent/src/eval/session-id.ts | 8 + .../coding-agent/src/prompts/tools/eval.md | 2 +- packages/coding-agent/src/sdk.ts | 6 + .../coding-agent/src/session/agent-session.ts | 37 ++- packages/coding-agent/src/task/executor.ts | 3 + packages/coding-agent/src/task/index.ts | 3 + .../src/tools/browser/tab-worker.ts | 5 +- packages/coding-agent/src/tools/eval.ts | 3 +- 19 files changed, 774 insertions(+), 298 deletions(-) create mode 100644 packages/coding-agent/src/eval/__tests__/shared-executors.test.ts create mode 100644 packages/coding-agent/src/eval/session-id.ts diff --git a/docs/tools/eval.md b/docs/tools/eval.md index 573241f3b..56753e4f5 100644 --- a/docs/tools/eval.md +++ b/docs/tools/eval.md @@ -90,21 +90,22 @@ Side-channel artifacts: - `js` is gated on `eval.js !== false`. - A disabled or unavailable requested backend throws `ToolError`; there is no auto-fallback or sniffing. 3. The tool allocates an `OutputSink`, a `TailBuffer`, per-cell result objects, and a `sessionAbortController`. `session.trackEvalExecution?.(...)` can wrap the whole run for external cancellation tracking. -4. Cells execute sequentially. For each cell, `execute()`: +4. It resolves the executor session id from `session.getEvalSessionId?.()`, falling back to `defaultEvalSessionId(session)`. Subagents inherit the parent's id so both sides share the same JS VM and Python kernel for each backend. +5. Cells execute sequentially within one eval tool call. For each cell, `execute()`: - clamps `(cell.timeout ?? 30) * 1000` ms through `clampTimeout("eval", ...)` - builds a combined abort signal from the tool signal, the timeout, and the session abort controller - marks the cell `running` and emits an update - calls the backend’s `execute()` with `cwd`, `sessionId`, `sessionFile`, `kernelOwnerId`, `deadlineMs`, `reset` (defaults to `false`), artifact info, and chunk callback -5. JS cells dispatch through `packages/coding-agent/src/eval/js/index.ts` into `executeJs()`; Python cells dispatch through `packages/coding-agent/src/eval/py/index.ts` into `executePython()`. -6. Backend text chunks stream into the shared `OutputSink`; rich outputs are accumulated separately as JSON, images, markdown markers, and status events. -7. After each cell: +6. JS cells dispatch through `packages/coding-agent/src/eval/js/index.ts` into `executeJs()`; Python cells dispatch through `packages/coding-agent/src/eval/py/index.ts` into `executePython()`. +7. Backend text chunks stream into the shared `OutputSink`; rich outputs are accumulated separately as JSON, images, markdown markers, and status events. +8. After each cell: - text output is trimmed and stored on that cell result - multi-cell runs prefix text with `[i/n]` and the optional title - cancellations return early with `isError: true` and a cell-specific abort message - non-zero exit codes return early with `isError: true` and a message naming the failed cell - later cells are skipped after the first error, but earlier cell state persists in the underlying runtime -8. On success, the tool joins all cell outputs, synthesizes `(no text output)` or `(no output)` when needed, and attaches truncation metadata from `summarizeFinal()`. -9. The renderer uses `details.cells`, `details.jsonOutputs`, and `details.statusEvents` to build notebook-style output. `mergeCallAndResult = true` and `inline = true`, so call and result render together in the transcript. +9. On success, the tool joins all cell outputs, synthesizes `(no text output)` or `(no output)` when needed, and attaches truncation metadata from `summarizeFinal()`. +10. The renderer uses `details.cells`, `details.jsonOutputs`, and `details.statusEvents` to build notebook-style output. `mergeCallAndResult = true` and `inline = true`, so call and result render together in the transcript. ## Modes / Variants @@ -121,8 +122,8 @@ If the requested backend is disabled or unavailable, the tool throws `ToolError` Implemented in `packages/coding-agent/src/eval/js/context-manager.ts` and `packages/coding-agent/src/eval/js/prelude.txt`. -- Persistent `vm.Context` instances keyed by `js:${sessionId}` in `vmContexts` -- `reset: true` calls `resetVmContext(sessionKey)` before the cell executes +- Persistent worker-backed VM sessions keyed by `js:${sessionId}` +- `reset: true` calls `resetVmContext(sessionKey)` before the cell executes; reset is destructive for all live runs on that JS session - Top-level `await` and bare `return` are supported by wrapping code in an async IIFE when `wrapCode()` sees `await` or `return` - Top-level static `import ... from ...` and dynamic `import(...)` calls are routed through `rewriteImports()`, which sends them via `__omp_import__` so the specifier resolves against the session cwd - Module cache is busted for **local** imports between cells so edits to source files are picked up without restarting the runtime. `__omp_import__` deletes `require.cache[absPath]` before re-importing whenever the original specifier is a filesystem path: relative (`./x`, `../x`, `.`, `..`), POSIX-absolute (`/...`), home-prefixed (`~/...`), or Windows drive-letter (`C:\...` / `C:/...`). Bare specifiers (`react`, `lodash/x`) and URL/scheme specifiers (`node:fs`, `file://...`, `https://...`) are left in cache so package identity stays stable across cells. The cache-bust only fires when the resolved target is an absolute path — unresolved bare-package fallbacks (`resolveImportSpecifier()` returning the original specifier) skip it. @@ -136,7 +137,7 @@ Implemented in `packages/coding-agent/src/eval/js/context-manager.ts` and `packa - `{ type: "image", data, mimeType }` becomes an image output - scalars become text - The VM exposes a restricted `process` subset plus `Buffer`, `fetch`, `Blob`, `File`, `Headers`, `Request`, `Response`, `fs`, `require`, and browser-style globals -- Per-session VM runs are serialized with `runQueued()` +- Concurrent runs on the same VM are not queued end-to-end. Synchronous JS still runs on the single event loop; awaited regions can interleave with sibling runs. ### Python runtime @@ -150,9 +151,9 @@ Implemented in `packages/coding-agent/src/eval/py/executor.ts`, `packages/coding - create/connect kernel - initialize cwd / env / `sys.path` - execute `PYTHON_PRELUDE` -- Python cells run inside IPython/Jupyter, so top-level `await` works; the prompt warns not to use `asyncio.run(...)` -- The Python prelude defines synchronous helpers with the same surface as JS (except `tool.` exists only in JS) -- `display(value)` wraps dict/list/tuple values in `IPython.display.JSON`; rich display MIME bundles are preserved +- Python cells run in the runner's persistent asyncio event loop, so top-level `await` works; the prompt warns not to use `asyncio.run(...)` +- The Python prelude defines helpers with the same surface as JS where practical, including `tool.(args)` through a per-run loopback bridge +- Synchronous statement blocks run in the default executor with ContextVar state copied in; the GIL still serializes bytecode execution, but awaited regions can interleave with sibling cells - Kernel `display_data` / `execute_result` messages map to: - `application/x-omp-status` → status event - `image/png` → image output @@ -184,14 +185,14 @@ A single tool call can mix Python and JS cells. Persistence is per language runt - Session state - `session.assertEvalExecutionAllowed?.()` can block execution. - `session.trackEvalExecution?.(...)` can register cancellable eval work. - - `session.getSessionFile?.()` and `session.getEvalKernelOwnerId?.()` influence kernel reuse and artifact lookup. - - JS VM contexts persist in `vmContexts` across eval calls until reset/disposal. - - Python retained kernels persist in `kernelSessions` until reset, eviction, idle cleanup, or owner cleanup. + - `session.getSessionFile?.()`, `session.getEvalSessionId?.()`, and `session.getEvalKernelOwnerId?.()` influence VM/kernel reuse and artifact lookup. + - JS VM contexts persist across eval calls until reset/disposal. + - Python retained kernels persist until reset, owner cleanup, or process exit. - User-visible prompts / interactive UI - none; stdin requests are rejected programmatically - Background work / cancellation - Python retained kernels have heartbeat and idle cleanup timers. - - Cancellation interrupts a running Python kernel and aborts JS promise waits. + - Cancellation hard-kills/resets the shared executor for that backend: JS terminates the worker, Python sends SIGINT and may escalate to subprocess shutdown. ## Limits & Caps @@ -224,14 +225,21 @@ A single tool call can mix Python and JS cells. Persistence is per language runt - Cancellation is returned, not thrown, once backend execution has started. The tool formats it as a cell failure and sets `details.isError = true`. - If output truncates, the tool still succeeds; truncation is surfaced through `details.meta` and artifact-backed full output when available. +## Shared executor trade-offs + +- Parent agents and subagents share eval state bidirectionally when a subagent inherits the parent's executor id. Mutations in either direction are visible to the other participant. +- Async regions of concurrent runs can interleave. Synchronous JS still blocks the VM event loop; synchronous Python still contends on the GIL. +- Cancelling one run is destructive to the shared backend executor. This is intentional: JS worker termination and Python SIGINT/subprocess shutdown are the only reliable way to interrupt arbitrary user code. +- `reset: true` is destructive for every live run on that backend session id. New starts on that backend are rejected while reset is in flight. + ## Notes - Backend selection is now strictly explicit per cell: `language` must be `"py"` or `"js"`. The previous `*** Cell` header parser, the `eval.lark` constrained grammar, and the sniffer-based fallback have all been removed. - `EvalTool.customFormat` no longer exists. Tool calls flow through the standard JSON schema; there is no Lark-constrained sampling path. -- `tool.()` exists only in JS. Python prelude helpers do not call back into the full tool registry. +- `tool.()` exists in both JS and Python. Python calls route through a per-run loopback bridge keyed by the current cell id. - JS helper paths reject protocol URIs (`://`) in `resolvePath()`; the JS prelude is filesystem-only unless the code calls `tool.read(...)` or another tool explicitly. - Python helper `output(...)` depends on `PI_SESSION_FILE`; it fails outside a session-backed run. - `display()` can produce text and structured outputs from the same value; the renderer prefers markdown over `text/plain` when both exist. - JS static imports are rewritten only at top level. Nested imports stay invalid and surface normal JS syntax/runtime errors. -- `EvalTool` is `concurrency = "exclusive"`, so eval calls do not overlap within a session. +- `EvalTool` is `concurrency = "exclusive"` within one agent session, but parent and subagent sessions can run eval concurrently when they share an inherited executor id. - The tool description shown to the model is templated by backend availability (`getEvalToolDescription()`); if Python is unavailable, the prompt omits Python-specific instructions. diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 110d8e1f1..5d384c86d 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,12 +1,19 @@ # Changelog ## [Unreleased] + ### Breaking Changes - The `vim` edit mode option is no longer available; configurations using `edit.mode: vim` will be automatically mapped to `hashline` mode ### Added +- Added evaluator state inheritance for `task`-spawned subagents so JavaScript and Python variables are visible between a parent agent and its child sessions +- Added `hashline-per` edit mode to restore the legacy per-line hashline dialect alongside the default file-hash dialect +- Added file-hash computation and validation for hashline sections to detect stale edits +- Added file-read snapshot caching with multi-snapshot ring per path for recovery from agent's own writes +- Added delete operation (`!`) support to hashline grammar for explicit line deletion +- Added structural bracket/brace balance warnings when deleting lines with unclosed constructs - Added file-hash computation and validation for hashline sections to detect stale edits - Added file-read snapshot caching with multi-snapshot ring per path for recovery from agent's own writes - Added delete operation (`!`) support to hashline grammar for explicit line deletion @@ -14,6 +21,7 @@ ### Changed +- Changed JavaScript and Python `eval` execution to allow overlapping asynchronous cells on the same session ID to run concurrently instead of being strictly queued - Updated the edit mode option set to support `replace`, `patch`, `hashline`, and `apply_patch` variants - Bare `A:` / `A-B:` (no payload, no inline body) now replaces the line/range with a single blank line, symmetric with bare `A↑` / `A↓` inserting a blank line; previously rejected as ambiguous - Simplified hashline anchor format from `LINE+HASH` to bare `LINE` numbers in edit operations @@ -24,6 +32,15 @@ - Modified hashline grammar to accept optional file hash in headers and removed hash requirements from line anchors - Changed hashline diff preview format to use `LINE:content` instead of `LINE+HASH|content` - Updated prompt documentation to reflect new `¶PATH#HASH` header and bare line-number syntax +- Bare `A:` / `A-B:` (no payload, no inline body) now replaces the line/range with a single blank line, symmetric with bare `A↑` / `A↓` inserting a blank line; previously rejected as ambiguous +- Simplified hashline anchor format from `LINE+HASH` to bare `LINE` numbers in edit operations +- Updated hashline file headers to include 4-hex file hash: `¶PATH#HASH` format for anchored edits +- Changed hashline line separator from `|` to `:` in editable output (e.g., `42:content` instead of `42ab|content`) +- Removed per-line hash validation; file-level hash now validates entire section integrity +- Updated read/search output to emit file-hash headers (`¶PATH#HASH`) followed by numbered lines for hashline mode +- Modified hashline grammar to accept optional file hash in headers and removed hash requirements from line anchors +- Changed hashline diff preview format to use `LINE:content` instead of `LINE+HASH|content` +- Updated prompt documentation to reflect new `¶PATH#HASH` header and bare line-number syntax ### Removed @@ -32,31 +49,16 @@ - Removed per-line hash anchors (2-letter bigram hashes) from hashline format - Removed `RANGE_INTERIOR_HASH` constant; multi-line ranges no longer use `**` filler - Removed `HashMismatch` type and hash mismatch error reporting; replaced with file-level validation -### Added - -- Added file-hash computation and validation for hashline sections to detect stale edits -- Added file-read snapshot caching with multi-snapshot ring per path for recovery from agent's own writes -- Added delete operation (`!`) support to hashline grammar for explicit line deletion -- Added structural bracket/brace balance warnings when deleting lines with unclosed constructs - -### Changed - -- Bare `A:` / `A-B:` (no payload, no inline body) now replaces the line/range with a single blank line, symmetric with bare `A↑` / `A↓` inserting a blank line; previously rejected as ambiguous -- Simplified hashline anchor format from `LINE+HASH` to bare `LINE` numbers in edit operations -- Updated hashline file headers to include 4-hex file hash: `¶PATH#HASH` format for anchored edits -- Changed hashline line separator from `|` to `:` in editable output (e.g., `42:content` instead of `42ab|content`) -- Removed per-line hash validation; file-level hash now validates entire section integrity -- Updated read/search output to emit file-hash headers (`¶PATH#HASH`) followed by numbered lines for hashline mode -- Modified hashline grammar to accept optional file hash in headers and removed hash requirements from line anchors -- Changed hashline diff preview format to use `LINE:content` instead of `LINE+HASH|content` -- Updated prompt documentation to reflect new `¶PATH#HASH` header and bare line-number syntax - -### Removed - - Removed per-line hash anchors (2-letter bigram hashes) from hashline format - Removed `RANGE_INTERIOR_HASH` constant; multi-line ranges no longer use `**` filler - Removed `HashMismatch` type and hash mismatch error reporting; replaced with file-level validation +### Fixed + +- Fixed JavaScript `eval` imports to preserve module-level singletons across re-imports of unchanged local files and reload them only after edits +- Fixed concurrent Python evaluator tool calls to use per-run identifiers so tool responses and output are routed to the correct execution +- Fixed the `search` tool argument validation to accept a single string `paths` value as a one-path search. + ## [15.4.0] - 2026-05-26 ### Breaking Changes diff --git a/packages/coding-agent/src/eval/__tests__/shared-executors.test.ts b/packages/coding-agent/src/eval/__tests__/shared-executors.test.ts new file mode 100644 index 000000000..7711154d4 --- /dev/null +++ b/packages/coding-agent/src/eval/__tests__/shared-executors.test.ts @@ -0,0 +1,313 @@ +import { afterAll, afterEach, describe, expect, it, vi } from "bun:test"; +import * as fs from "node:fs/promises"; +import * as path from "node:path"; +import type { AssistantMessage } from "@oh-my-pi/pi-ai"; +import { TempDir } from "@oh-my-pi/pi-utils"; +import type { ModelRegistry } from "../../config/model-registry"; +import { Settings } from "../../config/settings"; +import type { LoadExtensionsResult } from "../../extensibility/extensions/types"; +import type { CreateAgentSessionOptions, CreateAgentSessionResult } from "../../sdk"; +import * as sdkModule from "../../sdk"; +import type { AgentSession, AgentSessionEvent, PromptOptions } from "../../session/agent-session"; +import { TaskTool } from "../../task"; +import * as discoveryModule from "../../task/discovery"; +import type { AgentDefinition, TaskParams } from "../../task/types"; +import type { ToolSession } from "../../tools"; +import { EventBus } from "../../utils/event-bus"; +import { disposeAllVmContexts } from "../js/context-manager"; +import { executeJs } from "../js/executor"; +import { disposeAllKernelSessions, executePython } from "../py/executor"; + +function createToolSession(cwd: string, sessionFile: string | null, evalSessionId?: string): ToolSession { + const modelRegistry = { + authStorage: undefined, + refresh: async () => {}, + getAvailable: () => [], + getApiKey: async () => null, + } as unknown as ModelRegistry; + return { + cwd, + hasUI: false, + settings: Settings.isolated({ + "async.enabled": false, + "task.isolation.mode": "none", + }), + getSessionFile: () => sessionFile, + getSessionSpawns: () => "*", + getEvalSessionId: evalSessionId ? () => evalSessionId : undefined, + modelRegistry, + } as unknown as ToolSession; +} + +function assistantStopMessage(text: string): AssistantMessage { + return { + role: "assistant", + content: text ? [{ type: "text", text }] : [], + api: "openai-responses", + provider: "openai", + model: "mock", + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: Date.now(), + }; +} + +function createYieldingSubagentSession(onPrompt: () => Promise): AgentSession { + const listeners: Array<(event: AgentSessionEvent) => void> = []; + const state = { messages: [] as AssistantMessage[] }; + const emit = (event: AgentSessionEvent) => { + for (const listener of listeners) listener(event); + }; + return { + state, + agent: { state: { systemPrompt: ["test"] } }, + model: undefined, + extensionRunner: undefined, + sessionManager: { + appendSessionInit: () => {}, + }, + getActiveToolNames: () => ["eval", "yield"], + setActiveToolsByName: async () => {}, + subscribe: (listener: (event: AgentSessionEvent) => void) => { + listeners.push(listener); + return () => { + const index = listeners.indexOf(listener); + if (index >= 0) listeners.splice(index, 1); + }; + }, + prompt: async (_text: string, _options?: PromptOptions) => { + await onPrompt(); + state.messages.push(assistantStopMessage("done")); + emit({ + type: "tool_execution_end", + toolCallId: "yield-call", + toolName: "yield", + result: { + content: [{ type: "text", text: "Result submitted." }], + details: { status: "success", data: { ok: true } }, + }, + isError: false, + }); + }, + waitForIdle: async () => {}, + getLastAssistantMessage: () => state.messages[state.messages.length - 1], + abort: async () => {}, + dispose: async () => {}, + } as unknown as AgentSession; +} + +const taskAgent: AgentDefinition = { + name: "task", + description: "Task agent", + systemPrompt: "Read eval state and yield.", + source: "bundled", + tools: ["eval", "yield"], +}; + +const taskParams: TaskParams = { + agent: "task", + tasks: [{ id: "ReadEval", description: "Read eval state", assignment: "Read parent eval state." }], +}; + +describe("shared eval executors", () => { + afterEach(() => { + vi.restoreAllMocks(); + }); + + afterAll(async () => { + await disposeAllVmContexts(); + await disposeAllKernelSessions(); + }); + + it("shares JavaScript state across executeJs calls with one session id", async () => { + using tempDir = TempDir.createSync("@omp-eval-js-shared-"); + const sessionFile = path.join(tempDir.path(), "session.jsonl"); + const sessionId = `js-shared:${crypto.randomUUID()}`; + const session = createToolSession(tempDir.path(), sessionFile); + + await executeJs("globalThis.x = 41;", { sessionId, session, sessionFile }); + const result = await executeJs("return globalThis.x + 1;", { sessionId, session, sessionFile }); + + expect(result.exitCode).toBe(0); + expect(result.output.trim()).toBe("42"); + }); + + it("shares Python state across executePython calls with one session id", async () => { + using tempDir = TempDir.createSync("@omp-eval-py-shared-"); + const sessionFile = path.join(tempDir.path(), "session.jsonl"); + const sessionId = `py-shared:${crypto.randomUUID()}`; + + await executePython("x = 41", { cwd: tempDir.path(), sessionId, sessionFile }); + const result = await executePython("print(x + 1)", { cwd: tempDir.path(), sessionId, sessionFile }); + + expect(result.exitCode).toBe(0); + expect(result.output.trim()).toBe("42"); + }); + + it("lets a subagent inherit parent JavaScript and Python eval state", async () => { + using tempDir = TempDir.createSync("@omp-eval-subagent-"); + const sessionFile = path.join(tempDir.path(), "session.jsonl"); + const evalSessionId = `session:${sessionFile}:cwd:${tempDir.path()}`; + const parentSession = createToolSession(tempDir.path(), sessionFile, evalSessionId); + let seenJs = ""; + let seenPy = ""; + let capturedOptions: CreateAgentSessionOptions | undefined; + + await executeJs('globalThis.parentSecret = "hello-js";', { + sessionId: `js:${evalSessionId}`, + session: parentSession, + sessionFile, + }); + await executePython('parent_secret = "hello-py"', { + cwd: tempDir.path(), + sessionId: `python:${evalSessionId}`, + sessionFile, + }); + + vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({ agents: [taskAgent], projectAgentsDir: null }); + vi.spyOn(sdkModule, "createAgentSession").mockImplementation(async (options = {}) => { + capturedOptions = options; + const inherited = options.parentEvalSessionId; + if (!inherited) throw new Error("Missing parent eval session id"); + return { + session: createYieldingSubagentSession(async () => { + const jsResult = await executeJs("return globalThis.parentSecret;", { + sessionId: `js:${inherited}`, + session: parentSession, + sessionFile, + }); + const pyResult = await executePython("print(parent_secret)", { + cwd: tempDir.path(), + sessionId: `python:${inherited}`, + sessionFile, + }); + seenJs = jsResult.output.trim(); + seenPy = pyResult.output.trim(); + }), + extensionsResult: {} as unknown as LoadExtensionsResult, + setToolUIContext: () => {}, + eventBus: new EventBus(), + } satisfies CreateAgentSessionResult; + }); + + const tool = await TaskTool.create(parentSession); + await tool.execute("tool-call", taskParams); + + expect(capturedOptions?.parentEvalSessionId).toBe(evalSessionId); + expect(seenJs).toBe("hello-js"); + expect(seenPy).toBe("hello-py"); + }); + + it("interleaves async JavaScript runs on one session id", async () => { + using tempDir = TempDir.createSync("@omp-eval-js-interleave-"); + const sessionFile = path.join(tempDir.path(), "session.jsonl"); + const sessionId = `js-interleave:${crypto.randomUUID()}`; + const session = createToolSession(tempDir.path(), sessionFile); + const events: string[] = []; + + const first = executeJs('await Bun.sleep(80); display("A");', { + sessionId, + session, + sessionFile, + onChunk: chunk => { + events.push(chunk.trim()); + }, + }); + await Bun.sleep(10); + const second = executeJs('display("B");', { + sessionId, + session, + sessionFile, + onChunk: chunk => { + events.push(chunk.trim()); + }, + }); + + const [firstResult, secondResult] = await Promise.all([first, second]); + expect(firstResult.exitCode).toBe(0); + expect(secondResult.exitCode).toBe(0); + expect(events.filter(Boolean)).toEqual(["B", "A"]); + }); + + it("interleaves async Python runs on one session id", async () => { + using tempDir = TempDir.createSync("@omp-eval-py-interleave-"); + const sessionFile = path.join(tempDir.path(), "session.jsonl"); + const sessionId = `py-interleave:${crypto.randomUUID()}`; + const events: string[] = []; + + const first = executePython('import asyncio\nawait asyncio.sleep(0.08)\nprint("A")', { + cwd: tempDir.path(), + sessionId, + sessionFile, + onChunk: chunk => { + events.push(chunk.trim()); + }, + }); + await Bun.sleep(10); + const second = executePython('print("B")', { + cwd: tempDir.path(), + sessionId, + sessionFile, + onChunk: chunk => { + events.push(chunk.trim()); + }, + }); + + const [firstResult, secondResult] = await Promise.all([first, second]); + expect(firstResult.exitCode).toBe(0); + expect(secondResult.exitCode).toBe(0); + expect(events.filter(Boolean)).toEqual(["B", "A"]); + }); + + it("preserves module-level singleton state across re-imports of an unchanged file", async () => { + using tempDir = TempDir.createSync("@omp-eval-js-mtime-"); + const sessionFile = path.join(tempDir.path(), "session.jsonl"); + const sessionId = `js-mtime:${crypto.randomUUID()}`; + const session = createToolSession(tempDir.path(), sessionFile); + const modulePath = path.join(tempDir.path(), "singleton.ts"); + const moduleSpec = JSON.stringify(modulePath); + await Bun.write( + modulePath, + "let value = 0;\nexport function set(v) { value = v; }\nexport function get() { return value; }\n", + ); + + const initResult = await executeJs(`const mod = await import(${moduleSpec}); mod.set(42); return mod.get();`, { + sessionId, + session, + sessionFile, + }); + expect(initResult.exitCode).toBe(0); + expect(initResult.output.trim()).toBe("42"); + + // Unchanged file: re-import must reuse the existing module namespace so the + // counter is still 42. This is the regression — the previous unconditional + // `delete require.cache[target]` reset singletons on every dynamic import. + const reuseResult = await executeJs(`const mod = await import(${moduleSpec}); return mod.get();`, { + sessionId, + session, + sessionFile, + }); + expect(reuseResult.exitCode).toBe(0); + expect(reuseResult.output.trim()).toBe("42"); + + // Bump mtime by 5s to simulate an edit; the next import must evict the cache + // and re-evaluate the file, dropping the counter back to its initializer. + const future = new Date(Date.now() + 5_000); + await fs.utimes(modulePath, future, future); + + const reloadResult = await executeJs(`const mod = await import(${moduleSpec}); return mod.get();`, { + sessionId, + session, + sessionFile, + }); + expect(reloadResult.exitCode).toBe(0); + expect(reloadResult.output.trim()).toBe("0"); + }); +}); diff --git a/packages/coding-agent/src/eval/js/context-manager.ts b/packages/coding-agent/src/eval/js/context-manager.ts index 1b77736dd..995d2a1ff 100644 --- a/packages/coding-agent/src/eval/js/context-manager.ts +++ b/packages/coding-agent/src/eval/js/context-manager.ts @@ -48,10 +48,10 @@ interface JsSession { worker: WorkerHandle; state: "alive" | "dead"; pending: Map; - queue: Promise; } const sessions = new Map(); +const resettingSessions = new Set(); const READY_TIMEOUT_MS_DEFAULT = 5_000; export async function executeInVmContext(options: { @@ -66,14 +66,24 @@ export async function executeInVmContext(options: { runState: VmRunState; }): Promise<{ value: unknown }> { if (options.reset) { - await resetVmContext(options.sessionKey); + if (resettingSessions.has(options.sessionKey)) { + throw new ToolError("JS context reset already in progress"); + } + resettingSessions.add(options.sessionKey); + try { + await resetVmContext(options.sessionKey); + } finally { + resettingSessions.delete(options.sessionKey); + } + } else if (resettingSessions.has(options.sessionKey)) { + throw new ToolError("JS context reset in progress"); } const session = await acquireSession( options.sessionKey, { cwd: options.cwd, sessionId: options.sessionId }, options.timeoutMs, ); - return await runQueued(session, () => runOnce(session, options)); + return await runOnce(session, options); } export async function resetVmContext(sessionKey: string): Promise { @@ -89,22 +99,6 @@ export async function disposeAllVmContexts(): Promise { await Promise.all(all.map(session => killSession(session, new ToolError("JS context disposed")))); } -async function runQueued(session: JsSession, work: () => Promise): Promise { - const previous = session.queue; - const { promise, resolve } = Promise.withResolvers(); - session.queue = promise; - try { - await previous; - } catch { - // Previous run's failure must not poison this one. - } - try { - return await work(); - } finally { - resolve(); - } -} - async function runOnce( session: JsSession, options: { @@ -169,7 +163,6 @@ async function acquireSession(sessionKey: string, snapshot: SessionSnapshot, tim worker, state: "alive", pending: new Map(), - queue: Promise.resolve(), }; const { promise: readyPromise, resolve: resolveReady, reject: rejectReady } = Promise.withResolvers(); let resolved = false; diff --git a/packages/coding-agent/src/eval/js/shared/runtime.ts b/packages/coding-agent/src/eval/js/shared/runtime.ts index bf8c1747e..cac92b029 100644 --- a/packages/coding-agent/src/eval/js/shared/runtime.ts +++ b/packages/coding-agent/src/eval/js/shared/runtime.ts @@ -1,3 +1,4 @@ +import { AsyncLocalStorage } from "node:async_hooks"; import { Console } from "node:console"; import * as fs from "node:fs"; import { createRequire } from "node:module"; @@ -8,7 +9,6 @@ import * as util from "node:util"; import { logger } from "@oh-my-pi/pi-utils"; -import { ToolError } from "../../../tools/tool-errors"; import { createHelpers, type HelperBundle } from "./helpers"; import { awaitMaybePromise, indirectEval } from "./indirect-eval"; import { JAVASCRIPT_PRELUDE_SOURCE } from "./prelude"; @@ -16,10 +16,8 @@ import { wrapCode } from "./rewrite-imports"; import type { JsDisplayOutput, JsStatusEvent } from "./types"; /** - * Per-run callbacks. Returned by `getHooks()` on each helper/tool/display invocation so - * the embedding worker can route emissions to the currently active run. Returning `null` - * makes status/display/tool calls reject with an error — useful for guarding against - * helpers being invoked outside a run window. + * Per-run callbacks. Runtime globals resolve these from AsyncLocalStorage so + * overlapping async cells can route output/tool calls back to their own run. */ export interface RuntimeHooks { onText(chunk: string): void; @@ -27,11 +25,17 @@ export interface RuntimeHooks { callTool(name: string, args: unknown): Promise; } +export interface RunContext { + runId: string; + hooks: RuntimeHooks; + cwd: string; + finalExpressionSet: boolean; + finalExpressionValue: unknown; +} + export interface RuntimeOptions { initialCwd: string; sessionId: string; - /** Resolve hooks for the run currently in flight, or `null` if nothing is active. */ - getHooks(): RuntimeHooks | null; /** * Extra globals installed alongside `__omp_helpers__` / prelude. Use for stable, lifetime- * of-the-worker bindings (e.g. browser's `page`, `browser`). Per-run scope should be set @@ -120,19 +124,23 @@ export class JsRuntime { #cwd: string; readonly sessionId: string; #env: Map; - #getHooks: () => RuntimeHooks | null; - #finalExpressionSet = false; - #finalExpressionValue: unknown; + #als = new AsyncLocalStorage(); + /** + * mtime (ms) of every user-owned absolute path we've routed through `__omp_import__`. + * Powers edit-aware cache eviction: an unchanged file keeps its existing module + * instance — and therefore its module-private singletons — across cells; a bumped + * mtime triggers a one-shot `require.cache` eviction so the next `import` reloads. + */ + #moduleMtimes = new Map(); constructor(opts: RuntimeOptions) { this.#cwd = opts.initialCwd; this.sessionId = opts.sessionId; this.#env = new Map(); - this.#getHooks = opts.getHooks; this.helpers = createHelpers({ - cwd: () => this.#cwd, + cwd: () => this.#activeCwd(), env: this.#env, - emitStatus: event => this.#getHooks()?.onDisplay({ type: "status", event }), + emitStatus: event => this.#activeHooks("emitStatus")?.onDisplay({ type: "status", event }), }); this.#install(opts.extraGlobals); } @@ -156,29 +164,42 @@ export class JsRuntime { Object.assign(globalThis, scope); } - async run(code: string, filename?: string): Promise { - this.#finalExpressionSet = false; - this.#finalExpressionValue = undefined; - const wrapped = wrapCode(code); - const value = indirectEval(wrapped.source, filename); - if (wrapped.finalExpressionReturned) { - const awaited = await awaitMaybePromise(value); - if (this.#finalExpressionSet) { - const finalValue = this.#finalExpressionValue; - this.#finalExpressionSet = false; - this.#finalExpressionValue = undefined; - const resolved = await awaitMaybePromise(finalValue); - return resolved; + async run( + code: string, + filename: string | undefined, + hooks: RuntimeHooks, + options: { runId?: string; cwd?: string } = {}, + ): Promise { + const context: RunContext = { + runId: options.runId ?? crypto.randomUUID(), + hooks, + cwd: options.cwd ?? this.#cwd, + finalExpressionSet: false, + finalExpressionValue: undefined, + }; + return await this.#als.run(context, async () => { + const wrapped = wrapCode(code); + const value = indirectEval(wrapped.source, filename); + if (wrapped.finalExpressionReturned) { + const awaited = await awaitMaybePromise(value); + if (context.finalExpressionSet) { + const finalValue = context.finalExpressionValue; + context.finalExpressionSet = false; + context.finalExpressionValue = undefined; + return await awaitMaybePromise(finalValue); + } + return awaited; } - return awaited; - } - return await awaitMaybePromise(value); + return await awaitMaybePromise(value); + }); } - displayValue(value: unknown): void { + displayValue(value: unknown, hooks: RuntimeHooks | undefined = this.#als.getStore()?.hooks): void { if (value === undefined) return; - const hooks = this.#getHooks(); - if (!hooks) return; + if (!hooks) { + logger.warn("js runtime display called outside an active run"); + return; + } if (value && typeof value === "object") { const record = value as Record; if (record.type === "image" && typeof record.mimeType === "string") { @@ -209,37 +230,68 @@ export class JsRuntime { hooks.onText(`${String(value)}\n`); } + #activeCwd(): string { + return this.#als.getStore()?.cwd ?? this.#cwd; + } + + #activeHooks(action: string): RuntimeHooks | undefined { + const hooks = this.#als.getStore()?.hooks; + if (!hooks) { + logger.warn("js runtime helper called outside an active run", { action }); + } + return hooks; + } + #install(extraGlobals: Record | undefined): void { const injected: Record = { __omp_session__: { cwd: this.#cwd, sessionId: this.sessionId }, __omp_helpers__: this.helpers, __omp_call_tool__: async (name: string, args: unknown) => { - const hooks = this.#getHooks(); - if (!hooks) throw new ToolError("Tool calls are only valid inside an active run"); + const hooks = this.#activeHooks("tool"); + if (!hooks) return undefined; return await hooks.callTool(name, args); }, __omp_import__: async (source: string, options?: ImportCallOptions) => { - const target = resolveImportSpecifier(this.#cwd, source); - // Always invalidate cached module records for user-owned source files so edits - // between cells are picked up. Bun ignores query-string busting on `file:` URLs - // but honors `delete require.cache[absPath]`; bare specifiers and URL schemes are - // left alone to keep package identity stable across cells. + const target = resolveImportSpecifier(this.#activeCwd(), source); + // Edit-aware module cache eviction for user-owned source files (relative or + // absolute paths). Bun's module cache otherwise pins the first evaluation for + // the lifetime of the worker, which (a) hides edits made between cells and + // (b) would force any module-private singleton state to be re-initialized on + // every re-import — breaking patterns like `Settings.init()`'s module-scoped + // `globalInstance` when one cell inits and a later cell re-imports. + // + // Strategy: stat the resolved file, compare mtime to the last value we saw, + // and only evict when the file actually changed. First sight just records + // the mtime so the current module instance — and any singleton state it + // owns — survives subsequent imports until the user edits the file. Bare + // specifiers and URL schemes are left alone: `node:` built-ins cannot be + // reloaded and busting packages would defeat module identity across cells. if (isLocalPathSpecifier(source) && path.isAbsolute(target)) { - delete require.cache[target]; + try { + const mtime = fs.statSync(target).mtimeMs; + const prev = this.#moduleMtimes.get(target); + if (prev !== undefined && prev !== mtime) { + delete require.cache[target]; + } + this.#moduleMtimes.set(target, mtime); + } catch { + // stat failure (missing file, permission error, …) — fall through and + // let the real `import` surface the underlying error. + } } return options !== undefined ? await import(target, options) : await import(target); }, __omp_emit_status__: (op: string, data: Record = {}) => { const event: JsStatusEvent = { op, ...data }; - this.#getHooks()?.onDisplay({ type: "status", event }); + this.#activeHooks("emitStatus")?.onDisplay({ type: "status", event }); }, __omp_log__: (level: string, ...args: unknown[]) => { const prefix = level === "error" ? "[error] " : level === "warn" ? "[warn] " : ""; const text = `${prefix}${formatConsoleArgs(args)}`; - this.#getHooks()?.onText(text.endsWith("\n") ? text : `${text}\n`); + this.#activeHooks("log")?.onText(text.endsWith("\n") ? text : `${text}\n`); }, __omp_table__: (...args: unknown[]) => { - const hooks = this.#getHooks(); + const hooks = this.#activeHooks("table"); if (!hooks) return; let buffer = ""; const stream = new Writable({ @@ -254,8 +306,13 @@ export class JsRuntime { }, __omp_display__: (value: unknown) => this.displayValue(value), __omp_set_final_expr__: (value: unknown) => { - this.#finalExpressionSet = true; - this.#finalExpressionValue = value; + const context = this.#als.getStore(); + if (!context) { + logger.warn("js runtime final expression set outside an active run"); + return; + } + context.finalExpressionSet = true; + context.finalExpressionValue = value; }, webcrypto: crypto, // `process` is intentionally not overridden — user code gets the host worker's real diff --git a/packages/coding-agent/src/eval/js/worker-core.ts b/packages/coding-agent/src/eval/js/worker-core.ts index 590571213..552e9af9a 100644 --- a/packages/coding-agent/src/eval/js/worker-core.ts +++ b/packages/coding-agent/src/eval/js/worker-core.ts @@ -3,6 +3,7 @@ import { JsRuntime, type RuntimeHooks } from "./shared/runtime"; import type { RunErrorPayload, SessionSnapshot, ToolReply, Transport, WorkerInbound } from "./worker-protocol"; interface PendingTool { + runId: string; resolve(value: unknown): void; reject(error: Error): void; } @@ -36,8 +37,7 @@ function errorFromPayload(payload: RunErrorPayload): Error { export class WorkerCore { #transport: Transport; #runtime: JsRuntime | null = null; - #queue: Promise = Promise.resolve(); - #active: ActiveRun | null = null; + #runs = new Map(); #unsubscribe: () => void; constructor(transport: Transport) { @@ -52,7 +52,7 @@ export class WorkerCore { this.#ensureRuntime(msg.snapshot); return; case "run": - this.#enqueueRun(msg.runId, msg.code, msg.filename, msg.snapshot); + void this.#runOne(msg.runId, msg.code, msg.filename, msg.snapshot); return; case "tool-reply": this.#deliverToolReply(msg.id, msg.reply); @@ -71,73 +71,58 @@ export class WorkerCore { this.#runtime = new JsRuntime({ initialCwd: snapshot.cwd, sessionId: snapshot.sessionId, - getHooks: () => this.#hooksForCurrentRun(), }); return this.#runtime; } - #hooksForCurrentRun(): RuntimeHooks | null { - const active = this.#active; - if (!active) return null; - const runId = active.runId; - return { - onText: chunk => this.#transport.send({ type: "text", runId, chunk }), - onDisplay: output => this.#transport.send({ type: "display", runId, output }), - callTool: (name, args) => this.#callTool(active, name, args), - }; - } - - #enqueueRun(runId: string, code: string, filename: string, snapshot: SessionSnapshot): void { - const previous = this.#queue; - const next = (async () => { - await previous.catch(() => undefined); - await this.#runOne(runId, code, filename, snapshot); - })(); - this.#queue = next.catch(() => undefined); - } - async #runOne(runId: string, code: string, filename: string, snapshot: SessionSnapshot): Promise { const runtime = this.#ensureRuntime(snapshot); runtime.setCwd(snapshot.cwd); - this.#active = { runId, pendingTools: new Map() }; + const active: ActiveRun = { runId, pendingTools: new Map() }; + this.#runs.set(runId, active); + const hooks: RuntimeHooks = { + onText: chunk => this.#transport.send({ type: "text", runId, chunk }), + onDisplay: output => this.#transport.send({ type: "display", runId, output }), + callTool: (name, args) => this.#callTool(active, name, args), + }; try { - const value = await runtime.run(code, filename); - runtime.displayValue(value); + const value = await runtime.run(code, filename, hooks, { runId, cwd: snapshot.cwd }); + runtime.displayValue(value, hooks); this.#transport.send({ type: "result", runId, ok: true }); } catch (error) { this.#transport.send({ type: "result", runId, ok: false, error: errorPayload(error) }); } finally { - this.#active = null; + this.#runs.delete(runId); } } async #callTool(active: ActiveRun, name: string, args: unknown): Promise { const id = `tc-${active.runId}-${crypto.randomUUID()}`; const { promise, resolve, reject } = Promise.withResolvers(); - active.pendingTools.set(id, { resolve, reject }); + active.pendingTools.set(id, { runId: active.runId, resolve, reject }); this.#transport.send({ type: "tool-call", id, runId: active.runId, name, args }); return await promise; } #deliverToolReply(id: string, reply: ToolReply): void { - const active = this.#active; - if (!active) return; - const pending = active.pendingTools.get(id); - if (!pending) return; - active.pendingTools.delete(id); - if (reply.ok) pending.resolve(reply.value); - else pending.reject(errorFromPayload(reply.error)); + for (const active of this.#runs.values()) { + const pending = active.pendingTools.get(id); + if (!pending) continue; + active.pendingTools.delete(id); + if (reply.ok) pending.resolve(reply.value); + else pending.reject(errorFromPayload(reply.error)); + return; + } } #close(): void { - const active = this.#active; - if (active) { + for (const active of this.#runs.values()) { for (const pending of active.pendingTools.values()) { pending.reject(new ToolError("JS worker closed")); } active.pendingTools.clear(); } - this.#active = null; + this.#runs.clear(); this.#runtime = null; this.#transport.send({ type: "closed" }); this.#unsubscribe(); diff --git a/packages/coding-agent/src/eval/py/executor.ts b/packages/coding-agent/src/eval/py/executor.ts index 42ec13b7f..00cb2feeb 100644 --- a/packages/coding-agent/src/eval/py/executor.ts +++ b/packages/coding-agent/src/eval/py/executor.ts @@ -102,10 +102,10 @@ interface PythonSession { kernel: PythonKernel; ownerIds: Set; hasFallbackOwner: boolean; - queue: Promise; } const sessions = new Map(); +const resettingSessions = new Set(); // --------------------------------------------------------------------------- // Cancellation plumbing @@ -294,7 +294,6 @@ async function acquireSession(sessionId: string, cwd: string, options: PythonExe kernel, ownerIds: new Set(), hasFallbackOwner: false, - queue: Promise.resolve(), }; attachOwner(session, sessionId, options.kernelOwnerId); sessions.set(sessionId, session); @@ -330,27 +329,6 @@ async function resetSession(sessionId: string): Promise { await existing.kernel.shutdown().catch(() => undefined); } -async function runQueued( - session: PythonSession, - options: Pick, - work: () => Promise, -): Promise { - const previous = session.queue; - const { promise: ourSlot, resolve: releaseSlot } = Promise.withResolvers(); - // Keep the queue chained even if WE bail out: future runs must still wait - // for `previous` to finish before they touch the kernel. - session.queue = previous.catch(() => undefined).then(() => ourSlot); - try { - await waitForPromiseWithCancellation( - previous.catch(() => undefined), - options, - ); - return await work(); - } finally { - releaseSlot(); - } -} - // --------------------------------------------------------------------------- // Public dispose entry points // --------------------------------------------------------------------------- @@ -424,9 +402,10 @@ async function executeWithKernel( ((event: JsStatusEvent) => { displayOutputs.push({ type: "status", event }); }); + const runId = `py-${crypto.randomUUID()}`; const unregisterBridge = options?.toolSession && options?.bridgeSessionId - ? registerPyToolBridge(options.bridgeSessionId, { + ? registerPyToolBridge(options.bridgeSessionId, runId, { toolSession: options.toolSession, signal: options.signal, emitStatus, @@ -436,6 +415,7 @@ async function executeWithKernel( try { executionTimeoutMs = requireRemainingTimeoutMs(deadlineMs); const result = await kernel.execute(code, { + id: runId, signal: options?.signal, timeoutMs: executionTimeoutMs, onChunk: text => sink.push(text), @@ -528,38 +508,46 @@ async function executeOnSession(code: string, cwd: string, options: PythonExecut options.bridgeSessionId = sessionId; } if (options.reset) { - await resetSession(sessionId); + if (resettingSessions.has(sessionId)) { + throw new Error("Python kernel reset already in progress"); + } + resettingSessions.add(sessionId); + try { + await resetSession(sessionId); + } finally { + resettingSessions.delete(sessionId); + } + } else if (resettingSessions.has(sessionId)) { + throw new Error("Python kernel reset in progress"); } const session = await acquireSession(sessionId, cwd, options); - return await runQueued(session, options, async () => { - if (options.signal?.aborted) { - throw new PythonExecutionCancelledError(isTimedOutCancellation(options.signal.reason, options.signal)); - } + if (options.signal?.aborted) { + throw new PythonExecutionCancelledError(isTimedOutCancellation(options.signal.reason, options.signal)); + } + if (sessions.get(session.sessionId) !== session) { + throw new PythonExecutionCancelledError(false); + } + if (!session.kernel.isAlive()) { + await replaceSessionKernel(session, cwd, options); if (sessions.get(session.sessionId) !== session) { throw new PythonExecutionCancelledError(false); } - if (!session.kernel.isAlive()) { - await replaceSessionKernel(session, cwd, options); - if (sessions.get(session.sessionId) !== session) { - throw new PythonExecutionCancelledError(false); - } + } + try { + return await executeWithKernel(session.kernel, code, options); + } catch (err) { + if (isCancellationError(err) || options.signal?.aborted) throw err; + if (session.kernel.isAlive()) throw err; + if (sessions.get(session.sessionId) !== session) { + throw new PythonExecutionCancelledError(false); } - try { - return await executeWithKernel(session.kernel, code, options); - } catch (err) { - if (isCancellationError(err) || options.signal?.aborted) throw err; - if (session.kernel.isAlive()) throw err; - if (sessions.get(session.sessionId) !== session) { - throw new PythonExecutionCancelledError(false); - } - // Kernel died during execute. Replace it and retry once on a fresh one. - await replaceSessionKernel(session, cwd, options); - if (sessions.get(session.sessionId) !== session) { - throw new PythonExecutionCancelledError(false); - } - return await executeWithKernel(session.kernel, code, options); + // Kernel died during execute. Replace it and retry once on a fresh one. + await replaceSessionKernel(session, cwd, options); + if (sessions.get(session.sessionId) !== session) { + throw new PythonExecutionCancelledError(false); } - }); + return await executeWithKernel(session.kernel, code, options); + } } export async function executePythonWithKernel( diff --git a/packages/coding-agent/src/eval/py/kernel.ts b/packages/coding-agent/src/eval/py/kernel.ts index 07d2eb223..a3ed9869d 100644 --- a/packages/coding-agent/src/eval/py/kernel.ts +++ b/packages/coding-agent/src/eval/py/kernel.ts @@ -52,6 +52,7 @@ const STARTUP_TIMEOUT_MS = 10_000; const INTERRUPT_ESCALATION_MS = 5_000; export interface KernelExecuteOptions { + id?: string; signal?: AbortSignal; onChunk?: (text: string) => Promise | void; onDisplay?: (output: KernelDisplayOutput) => Promise | void; @@ -260,7 +261,7 @@ export class PythonKernel { throw new Error("Python kernel is not running"); } - const msgId = Snowflake.next(); + const msgId = options?.id ?? Snowflake.next(); const { promise, resolve } = Promise.withResolvers(); const pending: PendingExecution = { resolve, diff --git a/packages/coding-agent/src/eval/py/prelude.py b/packages/coding-agent/src/eval/py/prelude.py index abf08f9da..8c0066a64 100644 --- a/packages/coding-agent/src/eval/py/prelude.py +++ b/packages/coding-agent/src/eval/py/prelude.py @@ -402,8 +402,10 @@ if "__omp_prelude_loaded__" not in globals(): merged.update(kwargs) if "_i" not in merged: merged["_i"] = "py prelude" + _run_id_getter = globals().get("__omp_current_run_id__") + _run_id = _run_id_getter() if callable(_run_id_getter) else globals().get("__omp_run_id__") payload = json.dumps( - {"session": self._proxy._session, "name": self._name, "args": merged} + {"session": self._proxy._session, "run": _run_id, "name": self._name, "args": merged} ).encode("utf-8") req = urllib.request.Request( f"{self._proxy._base}/v1/tool", diff --git a/packages/coding-agent/src/eval/py/runner.py b/packages/coding-agent/src/eval/py/runner.py index 280590c6f..8ce83b1d2 100644 --- a/packages/coding-agent/src/eval/py/runner.py +++ b/packages/coding-agent/src/eval/py/runner.py @@ -28,6 +28,7 @@ from __future__ import annotations import asyncio import ast import base64 +import contextvars import builtins import inspect import io @@ -93,7 +94,7 @@ class _StreamProxy(io.TextIOBase): data = str(data) if not data: return 0 - rid = _STATE.current_id + rid = _CURRENT_RID.get() if rid is None: _RAW_STDERR.write(data) _RAW_STDERR.flush() @@ -112,7 +113,6 @@ class _StreamProxy(io.TextIOBase): class _RunnerState: def __init__(self) -> None: - self.current_id: str | None = None self.execution_count: int = 0 self.cancel_requested: bool = False # User globals — kept across requests when running in session mode. @@ -125,6 +125,8 @@ class _RunnerState: self.loop: asyncio.AbstractEventLoop | None = None +_CURRENT_RID: contextvars.ContextVar[str | None] = contextvars.ContextVar("omp_current_rid", default=None) + _STATE = _RunnerState() @@ -284,7 +286,7 @@ def cell_magic(name: str) -> Callable[[Callable[[str, str], Any]], Callable[[str def _emit_status(op: str, **data: Any) -> None: bundle = {"application/x-omp-status": {"op": op, **data}} - rid = _STATE.current_id + rid = _CURRENT_RID.get() if rid is None: return _emit({"type": "display", "id": rid, "bundle": bundle}) @@ -424,7 +426,7 @@ def _magic_reset(_args: str) -> None: def _magic_load(args: str) -> None: path = Path(os.path.expanduser(args.strip())) source = path.read_text(encoding="utf-8") - _emit({"type": "display", "id": _STATE.current_id, "bundle": {"text/plain": source}}) + _emit({"type": "display", "id": _CURRENT_RID.get(), "bundle": {"text/plain": source}}) _exec_source(source, _STATE.user_ns) @@ -622,7 +624,7 @@ def _mime_bundle(value: Any) -> dict: def _emit_display(bundle: dict, *, kind: str = "display") -> None: - rid = _STATE.current_id + rid = _CURRENT_RID.get() if rid is None: return _emit({"type": kind, "id": rid, "bundle": bundle}) @@ -681,6 +683,7 @@ def _install_builtins(ns: dict) -> None: ns["__omp_magic"] = __omp_magic ns["__omp_magic_cell"] = __omp_magic_cell ns["__omp_shell"] = __omp_shell + ns["__omp_current_run_id__"] = lambda: _CURRENT_RID.get() _install_builtins(_STATE.user_ns) @@ -694,25 +697,20 @@ _install_builtins(_STATE.user_ns) _TLA_FLAG = getattr(ast, "PyCF_ALLOW_TOP_LEVEL_AWAIT", 0x2000) -def _get_event_loop() -> asyncio.AbstractEventLoop: - loop = _STATE.loop - if loop is None or loop.is_closed(): - loop = asyncio.new_event_loop() - asyncio.set_event_loop(loop) - _STATE.loop = loop - return loop +def _await_sync(coro) -> Any: + try: + running_loop = asyncio.get_running_loop() + except RuntimeError: + running_loop = None + if running_loop is not None and running_loop.is_running(): + raise RuntimeError("top-level await is not supported from synchronous magic execution") + return asyncio.run(coro) -def _run_compiled(code, ns: dict, *, want_value: bool) -> Any: - """Execute a code object, awaiting it if compiled as a coroutine. - - ``want_value`` is True for the trailing expression — we return ``eval``'s - result (or the awaited coroutine's value). For statement blocks the - return is always ``None``. - """ +def _run_compiled_sync(code, ns: dict, *, want_value: bool) -> Any: + """Synchronous execution path used by nested magic helpers.""" if code.co_flags & inspect.CO_COROUTINE: - coro = eval(code, ns) - result = _get_event_loop().run_until_complete(coro) + result = _await_sync(eval(code, ns)) return result if want_value else None if want_value: return eval(code, ns) @@ -720,15 +718,32 @@ def _run_compiled(code, ns: dict, *, want_value: bool) -> Any: return None -def _exec_source(source: str, ns: dict) -> None: - """Compile + execute ``source``; if the last node is an expression, route - its value through ``__omp_display`` so dataframes/figures render rich. - Top-level ``await`` / ``async for`` / ``async with`` is permitted; the - cell is driven through the runner's persistent event loop.""" - module = ast.parse(source, mode="exec") +async def _run_in_executor_with_context(fn: Callable[[], Any]) -> Any: + loop = asyncio.get_running_loop() + ctx = contextvars.copy_context() + return await loop.run_in_executor(None, ctx.run, fn) + +async def _run_compiled_async(code, ns: dict, *, want_value: bool) -> Any: + """Execute a code object in the persistent event loop. + + Coroutine code is awaited in this task so top-level ``await`` interleaves + with sibling requests. Plain statement/expression code runs in the default + executor with the current ContextVar state copied into that thread. + """ + if code.co_flags & inspect.CO_COROUTINE: + result = await eval(code, ns) + return result if want_value else None + if want_value: + return await _run_in_executor_with_context(lambda: eval(code, ns)) + await _run_in_executor_with_context(lambda: exec(code, ns)) + return None + + +def _compile_source(source: str) -> tuple[Any, Any | None, bool]: + module = ast.parse(source, mode="exec") if not module.body: - return + return None, None, False last = module.body[-1] if isinstance(last, ast.Expr): @@ -737,14 +752,36 @@ def _exec_source(source: str, ns: dict) -> None: ast.copy_location(expr_module, last) body_code = compile(body_module, "", "exec", flags=_TLA_FLAG) expr_code = compile(expr_module, "", "eval", flags=_TLA_FLAG) - _run_compiled(body_code, ns, want_value=False) - value = _run_compiled(expr_code, ns, want_value=True) + return body_code, expr_code, True + + return compile(module, "", "exec", flags=_TLA_FLAG), None, False + + +def _exec_source(source: str, ns: dict) -> None: + """Synchronous source execution for legacy magic helpers.""" + body_code, expr_code, has_expr = _compile_source(source) + if body_code is None: + return + _run_compiled_sync(body_code, ns, want_value=False) + if has_expr and expr_code is not None: + value = _run_compiled_sync(expr_code, ns, want_value=True) if value is not None: __omp_display(value, kind="result") - return - code = compile(module, "", "exec", flags=_TLA_FLAG) - _run_compiled(code, ns, want_value=False) + +async def _exec_source_async(source: str, ns: dict) -> None: + """Compile + execute ``source``; if the last node is an expression, route + its value through ``__omp_display`` so dataframes/figures render rich. + Top-level ``await`` / ``async for`` / ``async with`` is permitted; awaited + regions yield to other requests in the runner's persistent event loop.""" + body_code, expr_code, has_expr = _compile_source(source) + if body_code is None: + return + await _run_compiled_async(body_code, ns, want_value=False) + if has_expr and expr_code is not None: + value = await _run_compiled_async(expr_code, ns, want_value=True) + if value is not None: + __omp_display(value, kind="result") # --------------------------------------------------------------------------- @@ -802,61 +839,60 @@ def _start_parent_watchdog() -> None: # --------------------------------------------------------------------------- -def _handle_request(req: dict) -> None: - if req.get("type") == "exit": - sys.exit(0) - +async def _handle_request_async(req: dict) -> None: rid = str(req.get("id")) - code = req.get("code", "") - _STATE.current_id = rid + token = _CURRENT_RID.set(rid) + _STATE.user_ns["__omp_run_id__"] = rid _STATE.cancel_requested = False _STATE.execution_count += 1 + execution_count = _STATE.execution_count _emit({"type": "started", "id": rid}) status: str = "ok" cancelled = False try: - transformed = transform_cell(code) - except SyntaxError as exc: - _emit_error(rid, exc) + try: + transformed = transform_cell(req.get("code", "")) + except SyntaxError as exc: + _emit_error(rid, exc) + _emit({ + "type": "done", + "id": rid, + "status": "error", + "executionCount": execution_count, + "cancelled": False, + }) + return + + _install_exec_sigint() + try: + await _exec_source_async(transformed, _STATE.user_ns) + except KeyboardInterrupt: + cancelled = True + status = "error" + _emit_error(rid, KeyboardInterrupt("Execution interrupted")) + except SystemExit: + raise + except BaseException as exc: # noqa: BLE001 - we want to surface every user error + status = "error" + _emit_error(rid, exc) + finally: + _install_idle_sigint() + try: + _flush_matplotlib_figures() + except Exception: + pass + _emit({ "type": "done", "id": rid, - "status": "error", - "executionCount": _STATE.execution_count, - "cancelled": False, + "status": status, + "executionCount": execution_count, + "cancelled": cancelled, }) - _STATE.current_id = None - return - - _install_exec_sigint() - try: - _exec_source(transformed, _STATE.user_ns) - except KeyboardInterrupt: - cancelled = True - status = "error" - _emit_error(rid, KeyboardInterrupt("Execution interrupted")) - except SystemExit: - raise - except BaseException as exc: # noqa: BLE001 - we want to surface every user error - status = "error" - _emit_error(rid, exc) finally: - _install_idle_sigint() - try: - _flush_matplotlib_figures() - except Exception: - pass - - _emit({ - "type": "done", - "id": rid, - "status": status, - "executionCount": _STATE.execution_count, - "cancelled": cancelled, - }) - _STATE.current_id = None + _CURRENT_RID.reset(token) def _emit_error(rid: str, exc: BaseException) -> None: @@ -875,16 +911,7 @@ def _emit_error(rid: str, exc: BaseException) -> None: # --------------------------------------------------------------------------- -def main() -> None: - sys.stdout = _StreamProxy("stdout") - sys.stderr = _StreamProxy("stderr") - _install_idle_sigint() - _start_parent_watchdog() - - stdin = sys.__stdin__ - if stdin is None: - return - +def _read_stdin(loop: asyncio.AbstractEventLoop, queue: asyncio.Queue, stdin) -> None: for raw_line in stdin: line = raw_line.strip() if not line: @@ -900,10 +927,52 @@ def main() -> None: "traceback": [], }) continue + loop.call_soon_threadsafe(queue.put_nowait, req) + loop.call_soon_threadsafe(queue.put_nowait, {"type": "exit"}) + + +async def _main_async() -> None: + sys.stdout = _StreamProxy("stdout") + sys.stderr = _StreamProxy("stderr") + _install_idle_sigint() + _start_parent_watchdog() + + stdin = sys.__stdin__ + if stdin is None: + return + + loop = asyncio.get_running_loop() + _STATE.loop = loop + queue: asyncio.Queue = asyncio.Queue() + reader = threading.Thread(target=_read_stdin, args=(loop, queue, stdin), name="omp-stdin-reader", daemon=True) + reader.start() + + tasks: set[asyncio.Task] = set() + def _task_done(task: asyncio.Task) -> None: + tasks.discard(task) try: - _handle_request(req) - except SystemExit: + exc = task.exception() + except asyncio.CancelledError: return + if exc is not None: + _emit_error("", exc) + try: + while True: + req = await queue.get() + if req.get("type") == "exit": + break + task = asyncio.create_task(_handle_request_async(req)) + tasks.add(task) + task.add_done_callback(_task_done) + finally: + for task in tasks: + task.cancel() + if tasks: + await asyncio.gather(*tasks, return_exceptions=True) + + +def main() -> None: + asyncio.run(_main_async()) if __name__ == "__main__": diff --git a/packages/coding-agent/src/eval/py/tool-bridge.ts b/packages/coding-agent/src/eval/py/tool-bridge.ts index 0cbfe8f3f..7c8f27d71 100644 --- a/packages/coding-agent/src/eval/py/tool-bridge.ts +++ b/packages/coding-agent/src/eval/py/tool-bridge.ts @@ -44,21 +44,23 @@ async function startServer(): Promise { return new Response("Forbidden", { status: 403 }); } - let body: { session?: unknown; name?: unknown; args?: unknown }; + let body: { session?: unknown; run?: unknown; name?: unknown; args?: unknown }; try { - body = (await req.json()) as { session?: unknown; name?: unknown; args?: unknown }; + body = (await req.json()) as { session?: unknown; run?: unknown; name?: unknown; args?: unknown }; } catch { return Response.json({ ok: false, error: "Invalid JSON body" }, { status: 400 }); } const sessionId = typeof body.session === "string" ? body.session : ""; + const runId = typeof body.run === "string" ? body.run : ""; const name = typeof body.name === "string" ? body.name : ""; - if (!sessionId || !name) { - return Response.json({ ok: false, error: "Missing session/name" }, { status: 400 }); + if (!sessionId || !runId || !name) { + return Response.json({ ok: false, error: "Missing session/run/name" }, { status: 400 }); } - const entry = registrations.get(sessionId); + const registrationKey = bridgeRegistrationKey(sessionId, runId); + const entry = registrations.get(registrationKey) ?? registrations.get(sessionId); if (!entry) { return Response.json( - { ok: false, error: `No active Python tool bridge session: ${sessionId}` }, + { ok: false, error: `No active Python tool bridge session: ${registrationKey}` }, { status: 200 }, ); } @@ -111,11 +113,16 @@ export async function ensurePyToolBridge(): Promise { * Register a tool session for the duration of one execution. The returned * function MUST be called to remove the entry once execution finishes. */ -export function registerPyToolBridge(sessionId: string, entry: PyToolBridgeEntry): () => void { - registrations.set(sessionId, entry); +function bridgeRegistrationKey(sessionId: string, runId: string): string { + return `${sessionId}:${runId}`; +} + +export function registerPyToolBridge(sessionId: string, runId: string, entry: PyToolBridgeEntry): () => void { + const key = bridgeRegistrationKey(sessionId, runId); + registrations.set(key, entry); return () => { - if (registrations.get(sessionId) === entry) { - registrations.delete(sessionId); + if (registrations.get(key) === entry) { + registrations.delete(key); } }; } diff --git a/packages/coding-agent/src/eval/session-id.ts b/packages/coding-agent/src/eval/session-id.ts new file mode 100644 index 000000000..d339b6a9c --- /dev/null +++ b/packages/coding-agent/src/eval/session-id.ts @@ -0,0 +1,8 @@ +import type { ToolSession } from "../tools"; + +export type EvalSessionSource = Pick; + +export function defaultEvalSessionId(session: EvalSessionSource): string { + const sessionFile = session.getSessionFile?.() ?? undefined; + return sessionFile ? `session:${sessionFile}:cwd:${session.cwd}` : `cwd:${session.cwd}`; +} diff --git a/packages/coding-agent/src/prompts/tools/eval.md b/packages/coding-agent/src/prompts/tools/eval.md index 78dec551d..b95b09677 100644 --- a/packages/coding-agent/src/prompts/tools/eval.md +++ b/packages/coding-agent/src/prompts/tools/eval.md @@ -1,7 +1,7 @@ Run code in a persistent kernel using a list of cells. -Each call submits one or more cells. Cells run in array order. State persists within each language across cells **and across tool calls**. +Each call submits one or more cells. Cells run in array order. State persists within each language across cells, tool calls, and subagents spawned with `task`; variables a parent or subagent declares are visible to the other on the same shared executor. Cell fields: diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index 3ea03426c..ea20a1b76 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -51,6 +51,7 @@ import "./discovery"; import { resolveConfigValue } from "./config/resolve-config-value"; import { initializeWithSettings } from "./discovery"; import { disposeAllKernelSessions, disposeKernelSessionsByOwner } from "./eval/py/executor"; +import { defaultEvalSessionId } from "./eval/session-id"; import { TtsrManager } from "./export/ttsr"; import { type CustomCommandsLoadResult, @@ -318,6 +319,8 @@ export interface CreateAgentSessionOptions { agentRegistry?: AgentRegistry; /** Parent task ID prefix for nested artifact naming (e.g., "6-Extensions") */ parentTaskPrefix?: string; + /** Inherited eval executor session id for subagents sharing parent eval state. */ + parentEvalSessionId?: string; /** Session manager. Default: session stored under the configured agentDir sessions root */ sessionManager?: SessionManager; @@ -1177,6 +1180,8 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} taskDepth: options.taskDepth ?? 0, getSessionFile: () => sessionManager.getSessionFile() ?? null, getEvalKernelOwnerId: () => evalKernelOwnerId, + getEvalSessionId: () => + session?.getEvalSessionId() ?? options.parentEvalSessionId ?? defaultEvalSessionId(toolSession), assertEvalExecutionAllowed: () => session?.assertEvalExecutionAllowed(), trackEvalExecution: (execution, abortController) => session ? session.trackEvalExecution(execution, abortController) : execution, @@ -1994,6 +1999,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} agentId: resolvedAgentId, agentRegistry, providerSessionId: options.providerSessionId, + parentEvalSessionId: options.parentEvalSessionId, }); hasSession = true; if (asyncJobManager) { diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 3a6e953ad..f6320f140 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -108,6 +108,7 @@ import { executePython as executePythonCommand, type PythonResult, } from "../eval/py/executor"; +import { defaultEvalSessionId } from "../eval/session-id"; import { type BashResult, executeBash as executeBashCommand } from "../exec/bash-executor"; import { exportSessionToHtml } from "../export/html"; import type { TtsrManager, TtsrMatchContext } from "../export/ttsr"; @@ -139,6 +140,7 @@ import type { Skill, SkillWarning } from "../extensibility/skills"; import { expandSlashCommand, type FileSlashCommand } from "../extensibility/slash-commands"; import { GoalRuntime } from "../goals/runtime"; import type { Goal, GoalModeState } from "../goals/state"; +import { getHashlineSyntax } from "../hashline/hash"; import type { HindsightSessionState } from "../hindsight/state"; import { type LocalProtocolOptions, resolveLocalUrlToPath } from "../internal-urls"; import { @@ -312,6 +314,8 @@ export interface AgentSessionConfig { ttsrManager?: TtsrManager; /** Secret obfuscator for deobfuscating streaming edit content */ obfuscator?: SecretObfuscator; + /** Inherited eval executor session id from a parent agent. */ + parentEvalSessionId?: string; /** Logical owner for retained Python kernels created by this session. */ evalKernelOwnerId?: string; /** @@ -808,6 +812,8 @@ export class AgentSession { // Python execution state #evalAbortControllers = new Set(); #evalKernelOwnerId: string; + #parentEvalSessionId: string | undefined; + #cachedEvalSessionId: string | null | undefined; /** * AsyncJobManager owned by this session (top-level only). Subagents leave * this undefined and **MUST NOT** dispose the global instance on teardown. @@ -997,6 +1003,7 @@ export class AgentSession { this.settings = config.settings; // Power assertions are taken per turn (see #beginInFlight); nothing acquired here. this.#evalKernelOwnerId = config.evalKernelOwnerId ?? `agent-session:${Snowflake.next()}`; + this.#parentEvalSessionId = config.parentEvalSessionId; this.#ownedAsyncJobManager = config.ownedAsyncJobManager; this.#scopedModels = config.scopedModels ?? []; this.#thinkingLevel = config.thinkingLevel; @@ -3693,6 +3700,16 @@ export class AgentSession { get sessionId(): string { return this.#providerSessionId ?? this.sessionManager.getSessionId(); } + getEvalSessionId(): string | null { + if (this.#cachedEvalSessionId !== undefined) return this.#cachedEvalSessionId; + this.#cachedEvalSessionId = + this.#parentEvalSessionId ?? + defaultEvalSessionId({ + cwd: this.sessionManager.getCwd(), + getSessionFile: () => this.sessionManager.getSessionFile() ?? null, + }); + return this.#cachedEvalSessionId; + } /** Current session display name, if set */ get sessionName(): string | undefined { @@ -4120,6 +4137,7 @@ export class AgentSession { const fileMentionMessages = await generateFileMentionMessages(fileMentions, this.sessionManager.getCwd(), { autoResizeImages: this.settings.get("images.autoResize"), useHashLines: resolveFileDisplayMode(this).hashLines, + syntax: getHashlineSyntax(this.#resolveActiveEditMode()), }); messages.push(...fileMentionMessages); } @@ -7430,9 +7448,13 @@ export class AgentSession { } } - // Use the same session ID as eval's Python backend for kernel sharing - const sessionFile = this.sessionManager.getSessionFile(); - const sessionId = sessionFile ? `session:${sessionFile}:cwd:${cwd}` : `cwd:${cwd}`; + // Use the same session ID as eval's Python backend for kernel sharing. + const sessionId = + this.getEvalSessionId() ?? + defaultEvalSessionId({ + cwd, + getSessionFile: () => this.sessionManager.getSessionFile() ?? null, + }); const result = await executePythonCommand(code, { cwd, sessionId, @@ -7717,10 +7739,17 @@ export class AgentSession { // removes the surface entirely. tools: [], }; + const cacheSessionId = this.sessionId; const options = this.prepareSimpleStreamOptions( { apiKey, - sessionId: this.sessionId, + // Side-channel turns must not share OpenAI/Codex append-only + // conversation state with the main agent turn: IRC and /btw can run + // while the main turn is mid-tool-call. Keep the prompt-cache key + // stable, but give provider routing a unique request lineage. + sessionId: `${cacheSessionId}:side:${Snowflake.next()}`, + promptCacheKey: cacheSessionId, + preferWebsockets: false, reasoning: toReasoningEffort(this.thinkingLevel), hideThinkingSummary: this.agent.hideThinkingSummary, serviceTier: this.serviceTier, diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index e8ccf5af1..bc347d1dc 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -185,6 +185,8 @@ export interface ExecutorOptions { */ parentArtifactManager?: ArtifactManager; parentHindsightSessionState?: HindsightSessionState; + /** Parent agent's eval executor session id. Subagents reuse it so eval state is shared. */ + parentEvalSessionId?: string; /** * Parent agent's OpenTelemetry configuration. When defined, the subagent's * loop is started with the same tracer/hooks but its own agent identity @@ -1236,6 +1238,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise 0 ? mcpProxyTools : undefined, localProtocolOptions: options.localProtocolOptions, telemetry: subagentTelemetry, + parentEvalSessionId: options.parentEvalSessionId, }), ); diff --git a/packages/coding-agent/src/task/index.ts b/packages/coding-agent/src/task/index.ts index e75c8aff4..5adefefde 100644 --- a/packages/coding-agent/src/task/index.ts +++ b/packages/coding-agent/src/task/index.ts @@ -844,6 +844,7 @@ export class TaskTool implements AgentTool path.basename(file.path).toLowerCase() !== "agents.md", ); const promptTemplates = this.session.promptTemplates; + const parentEvalSessionId = this.session.getEvalSessionId?.() ?? undefined; // Initialize progress for all tasks for (let i = 0; i < tasksWithUniqueIds.length; i++) { @@ -911,6 +912,7 @@ export class TaskTool implements AgentTool this.#hooksForActiveRun(), }); return this.#runtime; } diff --git a/packages/coding-agent/src/tools/eval.ts b/packages/coding-agent/src/tools/eval.ts index a931c45fb..9908711f5 100644 --- a/packages/coding-agent/src/tools/eval.ts +++ b/packages/coding-agent/src/tools/eval.ts @@ -6,6 +6,7 @@ import { prompt } from "@oh-my-pi/pi-utils"; import * as z from "zod/v4"; import { jsBackend, pythonBackend } from "../eval"; import type { ExecutorBackend } from "../eval/backend"; +import { defaultEvalSessionId } from "../eval/session-id"; import type { EvalCellResult, EvalDisplayOutput, EvalLanguage, EvalStatusEvent, EvalToolDetails } from "../eval/types"; import type { RenderResultOptions } from "../extensibility/custom-tools/types"; import { truncateToVisualLines } from "../modes/components/visual-truncate"; @@ -347,7 +348,7 @@ export class EvalTool implements AgentTool { pushUpdate(); }, }); - const sessionId = sessionFile ? `session:${sessionFile}:cwd:${session.cwd}` : `cwd:${session.cwd}`; + const sessionId = session.getEvalSessionId?.() ?? defaultEvalSessionId(session); for (let i = 0; i < cells.length; i++) { const cell = cells[i];