fix(eval): forked subagent kernel resets away from shared sessions
- Subagents inherit the parent's eval session id, so a child's reset: true destroyed the co-owned kernel and every sibling's interpreter state mid-session. - resolveOwnerScopedSessionKey now routes a reset from a non-exclusive owner onto a deterministic per-owner fork key: the requester gets a fresh private kernel, co-owners keep the shared one, and the fork stays sticky for that owner until its teardown reaps it. - Applied across Python, JavaScript, Ruby, and Julia executors; JS contexts gained an owner registry plus disposeVmContextsByOwner, wired into EvalRunner.disposeKernels and SDK session teardown. - Covered by pure key-resolution contracts and an end-to-end JS test: co-owner reset forks, shared state survives, fork is sticky, and per-owner dispose reaps only the fork.
This commit is contained in:
@@ -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)).
|
||||
|
||||
@@ -313,6 +313,45 @@ export function attachSessionOwner(
|
||||
}
|
||||
}
|
||||
|
||||
/** Owner registry shared by every language's live/starting session records. */
|
||||
export interface SessionOwners {
|
||||
ownerIds: Set<string>;
|
||||
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
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
@@ -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<void> {
|
||||
|
||||
async function executeOnSession(code: string, cwd: string, options: JuliaExecutorOptions): Promise<JuliaResult> {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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<string, PendingRun>;
|
||||
ownerIds: Set<string>;
|
||||
hasFallbackOwner: boolean;
|
||||
}
|
||||
|
||||
interface StartingJsSession extends SessionOwners {
|
||||
promise: Promise<JsSession>;
|
||||
}
|
||||
|
||||
const sessions = new Map<string, JsSession>();
|
||||
const startingSessions = new Map<string, Promise<JsSession>>();
|
||||
const startingSessions = new Map<string, StartingJsSession>();
|
||||
const resettingSessions = new Map<string, Promise<void>>();
|
||||
// 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<string, string>;
|
||||
@@ -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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<JsSession> {
|
||||
async function acquireSession(
|
||||
sessionKey: string,
|
||||
snapshot: SessionSnapshot,
|
||||
timeoutMs?: number,
|
||||
ownerId?: string,
|
||||
): Promise<JsSession> {
|
||||
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<JsSession> => {
|
||||
// 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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<PythonSession>;
|
||||
}
|
||||
|
||||
const sessions = new Map<string, PythonSession>();
|
||||
const startingSessions = new Map<string, Promise<PythonSession>>();
|
||||
const startingSessions = new Map<string, StartingPythonSession>();
|
||||
const resettingSessions = new Map<string, Promise<void>>();
|
||||
|
||||
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<void> {
|
||||
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<void> {
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
export async function disposeAllKernelSessions(): Promise<void> {
|
||||
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<voi
|
||||
}
|
||||
session.ownerIds.delete(ownerId);
|
||||
}
|
||||
const startingToShutdown: StartingPythonSession[] = [];
|
||||
for (const [sessionKey, starting] of [...startingSessions.entries()]) {
|
||||
if (sessions.has(sessionKey) || !starting.ownerIds.has(ownerId)) continue;
|
||||
if (starting.ownerIds.size === 1) {
|
||||
startingSessions.delete(sessionKey);
|
||||
startingToShutdown.push(starting);
|
||||
continue;
|
||||
}
|
||||
starting.ownerIds.delete(ownerId);
|
||||
}
|
||||
for (const session of toShutdown) {
|
||||
session.generation += 1;
|
||||
if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey);
|
||||
}
|
||||
const started = await Promise.allSettled(startingToShutdown.map(starting => 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<PythonResult> {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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<void> {
|
||||
|
||||
async function executeOnSession(code: string, cwd: string, options: RubyExecutorOptions): Promise<RubyResult> {
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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, SessionOwners> }): 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<void>();
|
||||
const releaseFirstStart = Promise.withResolvers<void>();
|
||||
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);
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user