diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index ec979b10e..f65966890 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -16,6 +16,7 @@ ### Fixed - Fixed Kitty terminals crashing while rendering live or restored non-PNG tool-result images when the runtime throws synchronously during PNG conversion ([#7160](https://github.com/can1357/oh-my-pi/issues/7160)). +- Fixed a subagent's eval `reset: true` wiping the shared kernel it inherits from its parent, destroying every co-owner's interpreter state mid-session. A reset from a non-exclusive owner now forks into a private per-owner kernel (sticky for that owner and reaped on its teardown) across the Python, JavaScript, Ruby, and Julia executors, while an exclusive owner still resets in place. - Fixed the copy selector and ask dialog rendering raw key IDs instead of human-readable keybinding labels ([#7164](https://github.com/can1357/oh-my-pi/issues/7164)). - Fixed CLI positional initial messages bypassing automatic session-title generation, which left shell-launched sessions unnamed until a later editor submission ([#7166](https://github.com/can1357/oh-my-pi/issues/7166)). - Fixed the environment-variable reference omitting the Kitty Unicode placeholder controls and tmux placement caveat ([#7172](https://github.com/can1357/oh-my-pi/issues/7172)). diff --git a/packages/coding-agent/src/eval/executor-base.ts b/packages/coding-agent/src/eval/executor-base.ts index 4531d8cb3..eb12e287d 100644 --- a/packages/coding-agent/src/eval/executor-base.ts +++ b/packages/coding-agent/src/eval/executor-base.ts @@ -313,6 +313,45 @@ export function attachSessionOwner( } } +/** Owner registry shared by every language's live/starting session records. */ +export interface SessionOwners { + ownerIds: Set; + hasFallbackOwner: boolean; +} + +/** + * Resolve the session key an owner's eval cell runs on, forking `reset` away + * from shared kernels. + * + * Eval sessions are shared across agents by design (subagents inherit the + * parent's eval session id), so honoring `reset` on a co-owned kernel would + * destroy every other agent's state — including cells executing at that + * moment. When the requester does not exclusively own the live base session, + * its reset resolves to a deterministic per-owner fork key: the requester + * starts a fresh private kernel while co-owners keep the shared one. Once + * forked, the owner keeps resolving to its fork, and per-owner dispose reaps + * the fork since the requester is its only registered owner. + */ +export function resolveOwnerScopedSessionKey(options: { + baseKey: string; + ownerId: string | undefined; + reset: boolean; + /** True when a live or starting session exists under `key`. */ + hasSession: (key: string) => boolean; + /** Owner registry for the session under `key`, when inspectable. */ + getOwners: (key: string) => SessionOwners | undefined; +}): string { + const { baseKey, ownerId } = options; + if (ownerId === undefined) return baseKey; + const forkKey = `${baseKey}\0fork\0${ownerId}`; + if (options.hasSession(forkKey)) return forkKey; + if (!options.reset) return baseKey; + const base = options.getOwners(baseKey); + if (!base) return baseKey; + const exclusive = !base.hasFallbackOwner && base.ownerIds.size === 1 && base.ownerIds.has(ownerId); + return exclusive ? baseKey : forkKey; +} + // --------------------------------------------------------------------------- // Base executor implementation // --------------------------------------------------------------------------- diff --git a/packages/coding-agent/src/eval/jl/executor.ts b/packages/coding-agent/src/eval/jl/executor.ts index 7f2fac489..174b8803c 100644 --- a/packages/coding-agent/src/eval/jl/executor.ts +++ b/packages/coding-agent/src/eval/jl/executor.ts @@ -1,7 +1,12 @@ import * as path from "node:path"; import { getProjectDir, logger } from "@oh-my-pi/pi-utils"; import type { ToolSession } from "../../tools"; -import { attachSessionOwner, createCancelledKernelResult, executeWithKernelBase } from "../executor-base"; +import { + attachSessionOwner, + createCancelledKernelResult, + executeWithKernelBase, + resolveOwnerScopedSessionKey, +} from "../executor-base"; import { ensurePyToolBridge, type PyToolBridgeInfo } from "../py/tool-bridge"; import type { EvalDisplayOutput, EvalStatusEvent } from "../types"; import { @@ -451,7 +456,13 @@ async function ensureToolBridge(options: JuliaExecutorOptions): Promise { async function executeOnSession(code: string, cwd: string, options: JuliaExecutorOptions): Promise { const sessionId = options.sessionId ?? `session:${cwd}`; - const sessionKey = buildSessionKey(sessionId, cwd, options.interpreter); + const sessionKey = resolveOwnerScopedSessionKey({ + baseKey: buildSessionKey(sessionId, cwd, options.interpreter), + ownerId: options.kernelOwnerId, + reset: options.reset === true, + hasSession: key => sessions.has(key) || startingSessions.has(key), + getOwners: key => sessions.get(key) ?? startingSessions.get(key), + }); if (options.bridge && !options.bridgeSessionId) { options.bridgeSessionId = sessionId; } diff --git a/packages/coding-agent/src/eval/js/context-manager.ts b/packages/coding-agent/src/eval/js/context-manager.ts index ce7387ce4..ad936ce2a 100644 --- a/packages/coding-agent/src/eval/js/context-manager.ts +++ b/packages/coding-agent/src/eval/js/context-manager.ts @@ -8,6 +8,7 @@ import { import type { ToolSession } from "../../tools"; import { ToolAbortError, ToolError } from "../../tools/tool-errors"; import { safeSend as safeSendIpc } from "../../utils/ipc"; +import { attachSessionOwner, resolveOwnerScopedSessionKey, type SessionOwners } from "../executor-base"; import { shouldDetachKernel } from "../py/spawn-options"; import { callSessionTool, type JsStatusEvent } from "./tool-bridge"; import { WorkerCore } from "./worker-core"; @@ -57,10 +58,16 @@ interface JsSession { worker: WorkerHandle; state: "alive" | "dead"; pending: Map; + ownerIds: Set; + hasFallbackOwner: boolean; +} + +interface StartingJsSession extends SessionOwners { + promise: Promise; } const sessions = new Map(); -const startingSessions = new Map>(); +const startingSessions = new Map(); const resettingSessions = new Map>(); // Worker startup (module-graph import + WorkerCore construction) is infrastructure // cost, not user compute. Floor it independently of Bun's 5s default per-test timeout @@ -99,6 +106,8 @@ export function setJsEvalWorkerThreadForTests(enabled: boolean): boolean { export async function executeInVmContext(options: { sessionKey: string; sessionId: string; + /** Logical owner identifier; scopes `reset` on shared contexts and retained-worker cleanup. */ + ownerId?: string; cwd: string; session: ToolSession; localRoots?: Record; @@ -108,47 +117,56 @@ export async function executeInVmContext(options: { timeoutMs?: number; runState: VmRunState; }): Promise<{ value: unknown }> { + const sessionKey = resolveOwnerScopedSessionKey({ + baseKey: options.sessionKey, + ownerId: options.ownerId, + reset: options.reset === true, + hasSession: key => sessions.has(key) || startingSessions.has(key), + getOwners: key => sessions.get(key) ?? startingSessions.get(key), + }); if (options.reset) { // Coalesce concurrent resets: an existing in-flight reset already // produces a fresh context, so a follow-up `reset: true` cell should // just wait for it rather than failing the user-visible call. - const inFlight = resettingSessions.get(options.sessionKey); + const inFlight = resettingSessions.get(sessionKey); if (inFlight) await inFlight.catch(() => undefined); else { - const resetPromise = resetVmContext(options.sessionKey); + const resetPromise = resetVmContext(sessionKey); resettingSessions.set( - options.sessionKey, + sessionKey, resetPromise.then(() => undefined), ); try { await resetPromise; } finally { - resettingSessions.delete(options.sessionKey); + resettingSessions.delete(sessionKey); } } } else { // Internal coordination: wait for any in-flight reset to settle and // then run on the freshly-rebuilt context. - const inFlight = resettingSessions.get(options.sessionKey); + const inFlight = resettingSessions.get(sessionKey); if (inFlight) await inFlight.catch(() => undefined); } const session = await acquireSession( - options.sessionKey, + sessionKey, { cwd: options.cwd, sessionId: options.sessionId, localRoots: options.localRoots }, options.timeoutMs, + options.ownerId, ); return await runOnce(session, options); } export async function resetVmContext(sessionKey: string): Promise { - const session = sessions.get(sessionKey) ?? (await startingSessions.get(sessionKey)?.catch(() => undefined)); + const session = + sessions.get(sessionKey) ?? (await startingSessions.get(sessionKey)?.promise.catch(() => undefined)); if (!session) return; sessions.delete(sessionKey); await killSession(session, new ToolError("JS context reset"), { force: false }); } export async function disposeAllVmContexts(): Promise { - const pending = [...startingSessions.values()]; + const pending = [...startingSessions.values()].map(starting => starting.promise); startingSessions.clear(); const started = await Promise.allSettled(pending); const all = [...sessions.values()]; @@ -160,6 +178,45 @@ export async function disposeAllVmContexts(): Promise { await Promise.all(all.map(session => killSession(session, new ToolError("JS context disposed"), { force: false }))); } +/** + * Shut down retained JS contexts owned solely by `ownerId` (e.g. a subagent's + * private fork); shared contexts just drop the owner registration. + */ +export async function disposeVmContextsByOwner(ownerId: string): Promise { + const toKill: JsSession[] = []; + for (const session of [...sessions.values()]) { + if (!session.ownerIds.has(ownerId)) continue; + if (session.ownerIds.size === 1) { + toKill.push(session); + continue; + } + session.ownerIds.delete(ownerId); + } + const startingToKill: StartingJsSession[] = []; + for (const [sessionKey, starting] of [...startingSessions.entries()]) { + if (sessions.has(sessionKey) || !starting.ownerIds.has(ownerId)) continue; + if (starting.ownerIds.size === 1) { + startingSessions.delete(sessionKey); + startingToKill.push(starting); + continue; + } + starting.ownerIds.delete(ownerId); + } + for (const session of toKill) { + if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey); + } + const started = await Promise.allSettled(startingToKill.map(starting => starting.promise)); + for (const result of started) { + if (result.status !== "fulfilled") continue; + const session = result.value; + if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey); + toKill.push(session); + } + await Promise.all( + toKill.map(session => killSession(session, new ToolError("JS context disposed"), { force: false })), + ); +} + /** * Smoke probe: spawn the JS evaluator through the worker-host entry and prove * it answers the `init` handshake in a real isolated subprocess (not the inline @@ -177,6 +234,8 @@ export async function smokeTestJsEvalWorker(): Promise { worker, state: "alive", pending: new Map(), + ownerIds: new Set(), + hasFallbackOwner: false, }; try { await initWorker(session, { cwd: process.cwd(), sessionId: "smoke" }, WORKER_INIT_TIMEOUT_MS); @@ -243,15 +302,25 @@ async function runOnce( } } -async function acquireSession(sessionKey: string, snapshot: SessionSnapshot, timeoutMs?: number): Promise { +async function acquireSession( + sessionKey: string, + snapshot: SessionSnapshot, + timeoutMs?: number, + ownerId?: string, +): Promise { const existing = sessions.get(sessionKey); if (existing && existing.state === "alive") { existing.sessionId = snapshot.sessionId; existing.cwd = snapshot.cwd; + attachSessionOwner(existing, snapshot.sessionId, ownerId); return existing; } const starting = startingSessions.get(sessionKey); - if (starting) return await starting; + if (starting) { + attachSessionOwner(starting, snapshot.sessionId, ownerId); + return await starting.promise; + } + let startingSession!: StartingJsSession; const startup = (async (): Promise => { // Attach the message listener before sending init. Both Bun Worker messages @@ -264,6 +333,8 @@ async function acquireSession(sessionKey: string, snapshot: SessionSnapshot, tim worker, state: "alive", pending: new Map(), + ownerIds: new Set(), + hasFallbackOwner: false, }; // Init headroom is the fixed infrastructure floor; the caller's per-cell timeout // dominates when larger so users can grant more by raising `timeout` on a cell. @@ -293,14 +364,27 @@ async function acquireSession(sessionKey: string, snapshot: SessionSnapshot, tim session.state = "alive"; } } - sessions.set(sessionKey, session); + session.ownerIds = new Set(startingSession.ownerIds); + session.hasFallbackOwner = startingSession.hasFallbackOwner; + // Publish only while this startup still owns the key: owner disposal or + // a concurrent dispose-all may have already reaped the starting record, + // and publishing here would resurrect a context that was just torn down. + if (startingSessions.get(sessionKey) === startingSession) { + sessions.set(sessionKey, session); + } return session; })(); - startingSessions.set(sessionKey, startup); + startingSession = { + ownerIds: new Set(), + hasFallbackOwner: false, + promise: startup, + }; + attachSessionOwner(startingSession, snapshot.sessionId, ownerId); + startingSessions.set(sessionKey, startingSession); try { return await startup; } finally { - if (startingSessions.get(sessionKey) === startup) startingSessions.delete(sessionKey); + if (startingSessions.get(sessionKey) === startingSession) startingSessions.delete(sessionKey); } } diff --git a/packages/coding-agent/src/eval/js/executor.ts b/packages/coding-agent/src/eval/js/executor.ts index f10579d02..d515b08cb 100644 --- a/packages/coding-agent/src/eval/js/executor.ts +++ b/packages/coding-agent/src/eval/js/executor.ts @@ -19,6 +19,8 @@ export interface JsExecutorOptions { onStatus?: (event: JsStatusEvent) => void; signal?: AbortSignal; sessionId: string; + /** Logical owner identifier; scopes `reset` on shared contexts and retained-worker cleanup. */ + kernelOwnerId?: string; reset?: boolean; sessionFile?: string; artifactPath?: string; @@ -100,6 +102,7 @@ export async function executeJs(code: string, options: JsExecutorOptions): Promi await executeInVmContext({ sessionKey: options.sessionId, sessionId: options.sessionId, + ownerId: options.kernelOwnerId, cwd: options.cwd ?? options.session.cwd, session: options.session, localRoots: options.localRoots, diff --git a/packages/coding-agent/src/eval/js/index.ts b/packages/coding-agent/src/eval/js/index.ts index d6ed1b0ea..6842a261a 100644 --- a/packages/coding-agent/src/eval/js/index.ts +++ b/packages/coding-agent/src/eval/js/index.ts @@ -28,6 +28,7 @@ export default { idleTimeoutMs: opts.idleTimeoutMs, signal: opts.signal, sessionId: namespaceSessionId(opts.sessionId), + kernelOwnerId: opts.kernelOwnerId, sessionFile: opts.sessionFile, reset: opts.reset, onChunk: opts.onChunk, diff --git a/packages/coding-agent/src/eval/py/executor.ts b/packages/coding-agent/src/eval/py/executor.ts index 806e99a98..e0d895f92 100644 --- a/packages/coding-agent/src/eval/py/executor.ts +++ b/packages/coding-agent/src/eval/py/executor.ts @@ -13,6 +13,8 @@ import { getRemainingTimeoutMs, isCancellationError, isTimedOutCancellation, + resolveOwnerScopedSessionKey, + type SessionOwners, waitForPromiseWithCancellation, } from "../executor-base"; import type { JsStatusEvent } from "../js/shared/types"; @@ -154,8 +156,12 @@ interface PythonSession { hasFallbackOwner: boolean; } +interface StartingPythonSession extends SessionOwners { + promise: Promise; +} + const sessions = new Map(); -const startingSessions = new Map>(); +const startingSessions = new Map(); const resettingSessions = new Map>(); function normalizeSessionCwd(cwd: string): string { @@ -252,10 +258,10 @@ async function acquireSession( } const starting = startingSessions.get(sessionKey); if (starting) { - const session = await starting; - attachSessionOwner(session, sessionId, options.kernelOwnerId); - return session; + attachSessionOwner(starting, sessionId, options.kernelOwnerId); + return await starting.promise; } + let startingSession!: StartingPythonSession; const startup = (async () => { const kernel = await startKernel(cwd, options); const session: PythonSession = { @@ -264,19 +270,28 @@ async function acquireSession( cwd, kernel, generation: 0, - ownerIds: new Set(), - hasFallbackOwner: false, + ownerIds: new Set(startingSession.ownerIds), + hasFallbackOwner: startingSession.hasFallbackOwner, }; - sessions.set(sessionKey, session); + // Publish only while this startup still owns the key: owner disposal or + // a concurrent dispose-all may have already reaped the starting record, + // and publishing here would resurrect a kernel that was just torn down. + if (startingSessions.get(sessionKey) === startingSession) { + sessions.set(sessionKey, session); + } return session; })(); - startingSessions.set(sessionKey, startup); + startingSession = { + ownerIds: new Set(), + hasFallbackOwner: false, + promise: startup, + }; + attachSessionOwner(startingSession, sessionId, options.kernelOwnerId); + startingSessions.set(sessionKey, startingSession); try { - const session = await startup; - attachSessionOwner(session, sessionId, options.kernelOwnerId); - return session; + return await startup; } finally { - if (startingSessions.get(sessionKey) === startup) startingSessions.delete(sessionKey); + if (startingSessions.get(sessionKey) === startingSession) startingSessions.delete(sessionKey); } } @@ -370,7 +385,8 @@ async function acquireLiveSessionKernel( } async function resetSession(sessionKey: string): Promise { - const existing = sessions.get(sessionKey) ?? (await startingSessions.get(sessionKey)?.catch(() => undefined)); + const existing = + sessions.get(sessionKey) ?? (await startingSessions.get(sessionKey)?.promise.catch(() => undefined)); if (!existing) return; existing.generation += 1; sessions.delete(sessionKey); @@ -382,7 +398,7 @@ async function resetSession(sessionKey: string): Promise { // --------------------------------------------------------------------------- export async function disposeAllKernelSessions(): Promise { - const pending = [...startingSessions.values()]; + const pending = [...startingSessions.values()].map(starting => starting.promise); startingSessions.clear(); const started = await Promise.allSettled(pending); const all = [...sessions.entries()]; @@ -422,10 +438,28 @@ export async function disposeKernelSessionsByOwner(ownerId: string): Promise starting.promise)); + for (const result of started) { + if (result.status !== "fulfilled") continue; + const session = result.value; + session.generation += 1; + if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey); + toShutdown.push(session); + } const results = await Promise.allSettled(toShutdown.map(session => shutdownInvalidatedSession(session))); for (let i = 0; i < toShutdown.length; i += 1) { const session = toShutdown[i]; @@ -503,7 +537,13 @@ async function executePerCall(code: string, cwd: string, options: PythonExecutor async function executeOnSession(code: string, cwd: string, options: PythonExecutorOptions): Promise { const sessionId = options.sessionId ?? `session:${cwd}`; - const sessionKey = buildSessionKey(sessionId, cwd, options.interpreter); + const sessionKey = resolveOwnerScopedSessionKey({ + baseKey: buildSessionKey(sessionId, cwd, options.interpreter), + ownerId: options.kernelOwnerId, + reset: options.reset === true, + hasSession: key => sessions.has(key) || startingSessions.has(key), + getOwners: key => sessions.get(key) ?? startingSessions.get(key), + }); if (options.bridge && !options.bridgeSessionId) { options.bridgeSessionId = sessionId; } diff --git a/packages/coding-agent/src/eval/rb/executor.ts b/packages/coding-agent/src/eval/rb/executor.ts index 8e77448f6..740793962 100644 --- a/packages/coding-agent/src/eval/rb/executor.ts +++ b/packages/coding-agent/src/eval/rb/executor.ts @@ -13,6 +13,7 @@ import { getRemainingTimeoutMs, isCancellationError, isTimedOutCancellation, + resolveOwnerScopedSessionKey, waitForPromiseWithCancellation, } from "../executor-base"; import type { JsStatusEvent } from "../js/shared/types"; @@ -407,7 +408,13 @@ async function ensureToolBridge(options: RubyExecutorOptions): Promise { async function executeOnSession(code: string, cwd: string, options: RubyExecutorOptions): Promise { const sessionId = options.sessionId ?? `session:${cwd}`; - const sessionKey = buildSessionKey(sessionId, cwd, options.interpreter); + const sessionKey = resolveOwnerScopedSessionKey({ + baseKey: buildSessionKey(sessionId, cwd, options.interpreter), + ownerId: options.kernelOwnerId, + reset: options.reset === true, + hasSession: key => sessions.has(key) || startingSessions.has(key), + getOwners: key => sessions.get(key) ?? startingSessions.get(key), + }); if (options.bridge && !options.bridgeSessionId) { options.bridgeSessionId = sessionId; } diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index 75c315156..4528581f6 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -67,6 +67,7 @@ import { createBridgeEditTool, createBridgeGrepFactory } from "./cursor-bridge-t import "./discovery"; import { initializeWithSettings } from "./discovery"; import { disposeAllJuliaKernelSessions, disposeJuliaKernelSessionsByOwner } from "./eval/jl/executor"; +import { disposeVmContextsByOwner } from "./eval/js/context-manager"; import { disposeAllKernelSessions, disposeKernelSessionsByOwner } from "./eval/py/executor"; import { disposeAllRubyKernelSessions, disposeRubyKernelSessionsByOwner } from "./eval/rb/executor"; import { defaultEvalSessionId } from "./eval/session-id"; @@ -3722,6 +3723,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} await disposeKernelSessionsByOwner(evalKernelOwnerId); await disposeRubyKernelSessionsByOwner(evalKernelOwnerId); await disposeJuliaKernelSessionsByOwner(evalKernelOwnerId); + await disposeVmContextsByOwner(evalKernelOwnerId); if (ownsAuthStorage) authStorage.close(); } } catch (cleanupError) { diff --git a/packages/coding-agent/src/session/eval-runner.ts b/packages/coding-agent/src/session/eval-runner.ts index e357ef9d9..2289b6d88 100644 --- a/packages/coding-agent/src/session/eval-runner.ts +++ b/packages/coding-agent/src/session/eval-runner.ts @@ -2,6 +2,7 @@ import type { Agent } from "@oh-my-pi/pi-agent-core"; import { logger } from "@oh-my-pi/pi-utils"; import type { Settings } from "../config/settings"; import { disposeJuliaKernelSessionsByOwner } from "../eval/jl/executor"; +import { disposeVmContextsByOwner } from "../eval/js/context-manager"; import { namespaceSessionId as namespacePythonSessionId } from "../eval/py"; import { disposeKernelSessionsByOwner, @@ -181,6 +182,7 @@ export class EvalRunner { disposeKernelSessionsByOwner(this.#kernelOwnerId), disposeRubyKernelSessionsByOwner(this.#kernelOwnerId), disposeJuliaKernelSessionsByOwner(this.#kernelOwnerId), + disposeVmContextsByOwner(this.#kernelOwnerId), ]); const errors: unknown[] = []; for (const result of results) if (result.status === "rejected") errors.push(result.reason); diff --git a/packages/coding-agent/test/eval/kernel-owner-scoping.test.ts b/packages/coding-agent/test/eval/kernel-owner-scoping.test.ts new file mode 100644 index 000000000..c448142aa --- /dev/null +++ b/packages/coding-agent/test/eval/kernel-owner-scoping.test.ts @@ -0,0 +1,189 @@ +import { afterEach, describe, expect, it, vi } from "bun:test"; +import { TempDir } from "@oh-my-pi/pi-utils"; +import { Settings } from "../../src/config/settings"; +import { resolveOwnerScopedSessionKey, type SessionOwners } from "../../src/eval/executor-base"; +import { disposeAllVmContexts, disposeVmContextsByOwner } from "../../src/eval/js/context-manager"; +import { executeJs } from "../../src/eval/js/executor"; +import { disposeAllKernelSessions, executePython } from "../../src/eval/py/executor"; +import { PythonKernel } from "../../src/eval/py/kernel"; +import type { ToolSession } from "../../src/tools"; + +function makeSession(cwd: string): ToolSession { + return { + cwd, + hasUI: false, + settings: Settings.isolated({ + "async.enabled": false, + "task.isolation.mode": "none", + "task.enableLsp": true, + }), + taskDepth: 0, + enableLsp: true, + getSessionFile: () => null, + getSessionSpawns: () => "*", + getActiveModelString: () => "p/active", + getModelString: () => "p/fallback", + getArtifactsDir: () => null, + getSessionId: () => "test-session", + getEvalSessionId: () => "test-eval-session", + }; +} + +describe("resolveOwnerScopedSessionKey", () => { + const BASE = "sess\0/cwd\0interp"; + const FORK = `${BASE}\0fork\0owner-b`; + + function resolve(options: { ownerId?: string; reset?: boolean; live?: Record }): string { + const live = options.live ?? {}; + return resolveOwnerScopedSessionKey({ + baseKey: BASE, + ownerId: options.ownerId, + reset: options.reset === true, + hasSession: key => key in live, + getOwners: key => live[key], + }); + } + + it("keeps the base key when the caller has no owner identity", () => { + expect(resolve({ reset: true, live: { [BASE]: { ownerIds: new Set(["x"]), hasFallbackOwner: false } } })).toBe( + BASE, + ); + }); + + it("stays on an existing fork even without reset", () => { + const live = { + [BASE]: { ownerIds: new Set(["owner-a", "owner-b"]), hasFallbackOwner: false }, + [FORK]: { ownerIds: new Set(["owner-b"]), hasFallbackOwner: false }, + }; + expect(resolve({ ownerId: "owner-b", live })).toBe(FORK); + expect(resolve({ ownerId: "owner-b", reset: true, live })).toBe(FORK); + }); + + it("forks a reset away from a co-owned base session", () => { + const live = { [BASE]: { ownerIds: new Set(["owner-a", "owner-b"]), hasFallbackOwner: false } }; + expect(resolve({ ownerId: "owner-b", reset: true, live })).toBe(FORK); + }); + + it("forks a reset away from a fallback-owned base session", () => { + // Fallback ownership means some session without an explicit owner uses + // the context; a scoped reset must not destroy it. + const live = { [BASE]: { ownerIds: new Set(["session-id"]), hasFallbackOwner: true } }; + expect(resolve({ ownerId: "owner-b", reset: true, live })).toBe(FORK); + }); + + it("resets in place when the requester exclusively owns the base session", () => { + const live = { [BASE]: { ownerIds: new Set(["owner-b"]), hasFallbackOwner: false } }; + expect(resolve({ ownerId: "owner-b", reset: true, live })).toBe(BASE); + }); + + it("uses the base key for a reset when no live session exists", () => { + expect(resolve({ ownerId: "owner-b", reset: true })).toBe(BASE); + expect(resolve({ ownerId: "owner-b" })).toBe(BASE); + }); +}); + +describe("JS eval owner-scoped reset forking", () => { + afterEach(async () => { + await disposeAllVmContexts(); + }); + + it("forks a subagent reset instead of clobbering the shared context", async () => { + using tempDir = TempDir.createSync("@omp-js-owner-fork-"); + const session = makeSession(tempDir.path()); + const evalSessionId = `js-owner-fork:${crypto.randomUUID()}`; + const run = (code: string, kernelOwnerId: string, reset?: boolean) => + executeJs(code, { cwd: tempDir.path(), sessionId: evalSessionId, session, kernelOwnerId, reset }); + + await run("var shared = 41;", "agent-a"); + const joined = await run("return shared + 1;", "agent-b"); + expect(joined.output.trim()).toBe("42"); + + // agent-b resets: it must land on a private fork with fresh state... + const forked = await run("return typeof shared;", "agent-b", true); + expect(forked.output.trim()).toBe("undefined"); + // ...while agent-a's shared context keeps its state. + const preserved = await run("return shared + 1;", "agent-a"); + expect(preserved.output.trim()).toBe("42"); + + // The fork is sticky: agent-b keeps resolving to it without reset. + await run("var forkOnly = 7;", "agent-b"); + const sticky = await run("return forkOnly;", "agent-b"); + expect(sticky.output.trim()).toBe("7"); + + // Disposing agent-b reaps only the fork; the shared context survives. + await disposeVmContextsByOwner("agent-b"); + const survivor = await run("return shared + 1;", "agent-a"); + expect(survivor.output.trim()).toBe("42"); + }); + + it("resets in place for the exclusive owner of a context", async () => { + using tempDir = TempDir.createSync("@omp-js-owner-exclusive-"); + const session = makeSession(tempDir.path()); + const evalSessionId = `js-owner-exclusive:${crypto.randomUUID()}`; + const run = (code: string, reset?: boolean) => + executeJs(code, { + cwd: tempDir.path(), + sessionId: evalSessionId, + session, + kernelOwnerId: "agent-solo", + reset, + }); + + await run("var solo = 1;"); + const reset = await run("return typeof solo;", true); + expect(reset.output.trim()).toBe("undefined"); + }); +}); + +describe("Python cold-start reset race", () => { + afterEach(async () => { + await disposeAllKernelSessions(); + vi.restoreAllMocks(); + }); + + it("forks a reset issued while the shared kernel is still starting", async () => { + using tempDir = TempDir.createSync("@omp-py-owner-race-"); + const shutdowns = [0, 0]; + const kernels: PythonKernel[] = []; + const firstStartEntered = Promise.withResolvers(); + const releaseFirstStart = Promise.withResolvers(); + vi.spyOn(PythonKernel, "start").mockImplementation(async () => { + const index = kernels.length; + const kernel = { + isAlive: () => true, + execute: async () => ({ status: "ok" as const, cancelled: false, timedOut: false }), + shutdown: async () => { + shutdowns[index] += 1; + return { confirmed: true }; + }, + } as unknown as PythonKernel; + kernels.push(kernel); + if (index === 0) { + firstStartEntered.resolve(); + await releaseFirstStart.promise; + } + return kernel; + }); + + const sessionId = `py-owner-race:${crypto.randomUUID()}`; + const common = { cwd: tempDir.path(), sessionId }; + // The parent's kernel start is deferred, so its session sits in + // startingSessions when the subagent's reset arrives. + const parentRun = executePython("x = 1", { ...common, kernelOwnerId: "agent-a" }); + await firstStartEntered.promise; + + // Pre-fix, the reset resolved to the shared base key, awaited the + // parent's gated startup inside resetSession, and then shut the + // parent's brand-new kernel down. Post-fix it forks immediately and + // completes without ever touching the gated startup. + const childResult = await executePython("y = 2", { ...common, kernelOwnerId: "agent-b", reset: true }); + expect(childResult.exitCode).toBe(0); + expect(kernels.length).toBe(2); + + releaseFirstStart.resolve(); + const parentResult = await parentRun; + expect(parentResult.exitCode).toBe(0); + // The parent's kernel must never be reaped by the subagent's reset. + expect(shutdowns[0]).toBe(0); + }); +}); diff --git a/packages/coding-agent/test/telemetry-export.test.ts b/packages/coding-agent/test/telemetry-export.test.ts index ddab68694..a30fbd030 100644 --- a/packages/coding-agent/test/telemetry-export.test.ts +++ b/packages/coding-agent/test/telemetry-export.test.ts @@ -100,10 +100,12 @@ describe("initTelemetryExport signals export path", () => { // singleton into every later test. The probe stands up its own loopback // receiver and exits 0 only when a protobuf trace export actually lands. const probe = fileURLToPath(new URL("./otel-export-probe.ts", import.meta.url)); - const proc = Bun.spawn(["bun", probe], { stdout: "pipe", stderr: "pipe" }); - const [code, stdout] = await Promise.all([proc.exited, new Response(proc.stdout).text()]); - expect(stdout).toContain("PROBE: RECEIVED"); - expect(code).toBe(0); + const proc = Bun.spawn([process.execPath, probe], { + stdin: "ignore", + stdout: "ignore", + stderr: "ignore", + }); + expect(await proc.exited).toBe(0); }, 20_000); it("exports log records and metrics to OTLP/proto receivers", async () => { @@ -111,10 +113,12 @@ describe("initTelemetryExport signals export path", () => { // drives the bridged logger and the agent telemetry metric hooks, then // asserts protobuf POSTs landed at both /v1/logs and /v1/metrics. const probe = fileURLToPath(new URL("./otel-signals-probe.ts", import.meta.url)); - const proc = Bun.spawn(["bun", probe], { stdout: "pipe", stderr: "pipe" }); - const [code, stdout] = await Promise.all([proc.exited, new Response(proc.stdout).text()]); - expect(stdout).toContain("PROBE: RECEIVED"); - expect(code).toBe(0); + const proc = Bun.spawn([process.execPath, probe], { + stdin: "ignore", + stdout: "ignore", + stderr: "ignore", + }); + expect(await proc.exited).toBe(0); }, 20_000); it("merges OTEL_RESOURCE_ATTRIBUTES into the exported resource", async () => { @@ -123,9 +127,11 @@ describe("initTelemetryExport signals export path", () => { // asserts the merged attributes land and that OTEL_SERVICE_NAME wins // service.name over an OTEL_RESOURCE_ATTRIBUTES entry. const probe = fileURLToPath(new URL("./otel-resource-probe.ts", import.meta.url)); - const proc = Bun.spawn(["bun", probe], { stdout: "pipe", stderr: "pipe" }); - const [code, stdout] = await Promise.all([proc.exited, new Response(proc.stdout).text()]); - expect(stdout).toContain("PROBE: RECEIVED"); - expect(code).toBe(0); + const proc = Bun.spawn([process.execPath, probe], { + stdin: "ignore", + stdout: "ignore", + stderr: "ignore", + }); + expect(await proc.exited).toBe(0); }, 20_000); });