Skip to content

Commit 3349fe2

Browse files
authored
Fix embedded session file ownership race (#87159)
* fix: serialize embedded session file attempts * test: update reply runtime mock for session file lookup * fix: thread session files into diagnostic recovery * fix: attach causes to session owner abort errors
1 parent ebe09be commit 3349fe2

19 files changed

Lines changed: 560 additions & 19 deletions

src/agents/pi-embedded-runner.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ export {
2020
queueEmbeddedPiMessage,
2121
queueEmbeddedPiMessage as queueEmbeddedAgentMessage,
2222
queueEmbeddedPiMessageWithOutcome,
23+
resolveActiveEmbeddedRunSessionIdBySessionFile,
2324
resolveActiveEmbeddedRunSessionId,
2425
resolveActiveEmbeddedRunSessionId as resolveActiveEmbeddedAgentRunSessionId,
2526
waitForEmbeddedPiRunEnd,

src/agents/pi-embedded-runner/run-state.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ const embeddedRunState = resolveGlobalSingleton(EMBEDDED_RUN_STATE_KEY, () => ({
4949
activeRuns: new Map<string, EmbeddedPiQueueHandle>(),
5050
snapshots: new Map<string, ActiveEmbeddedRunSnapshot>(),
5151
sessionIdsByKey: new Map<string, string>(),
52+
sessionIdsByFile: new Map<string, string>(),
5253
waiters: new Map<string, Set<EmbeddedRunWaiter>>(),
5354
modelSwitchRequests: new Map<string, EmbeddedRunModelSwitchRequest>(),
5455
}));
@@ -62,6 +63,9 @@ export const ACTIVE_EMBEDDED_RUN_SNAPSHOTS =
6263
export const ACTIVE_EMBEDDED_RUN_SESSION_IDS_BY_KEY =
6364
embeddedRunState.sessionIdsByKey ??
6465
(embeddedRunState.sessionIdsByKey = new Map<string, string>());
66+
export const ACTIVE_EMBEDDED_RUN_SESSION_IDS_BY_FILE =
67+
embeddedRunState.sessionIdsByFile ??
68+
(embeddedRunState.sessionIdsByFile = new Map<string, string>());
6569
export const EMBEDDED_RUN_WAITERS =
6670
embeddedRunState.waiters ??
6771
(embeddedRunState.waiters = new Map<string, Set<EmbeddedRunWaiter>>());
@@ -93,6 +97,7 @@ export function listActiveEmbeddedRunSessionIds(): string[] {
9397
...new Set([
9498
...ACTIVE_EMBEDDED_RUNS.keys(),
9599
...ACTIVE_EMBEDDED_RUN_SESSION_IDS_BY_KEY.values(),
100+
...ACTIVE_EMBEDDED_RUN_SESSION_IDS_BY_FILE.values(),
96101
...listActiveReplyRunSessionIds(),
97102
]),
98103
].toSorted((a, b) => a.localeCompare(b));

src/agents/pi-embedded-runner/run/attempt.session-lock.test.ts

Lines changed: 51 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,11 +14,13 @@ import {
1414
resetSessionWriteLockStateForTest,
1515
} from "../../session-write-lock.js";
1616
import {
17+
acquireEmbeddedAttemptSessionFileOwner,
1718
createEmbeddedAttemptSessionLockController,
1819
EmbeddedAttemptSessionTakeoverError,
1920
installPromptSubmissionLockRelease,
2021
installSessionEventWriteLock,
2122
installSessionExternalHookWriteLock,
23+
resetEmbeddedAttemptSessionFileOwnersForTest,
2224
} from "./attempt.session-lock.js";
2325

2426
const lockOptions = {
@@ -31,6 +33,7 @@ const lockOptions = {
3133
const tempDirs: string[] = [];
3234

3335
afterEach(async () => {
36+
resetEmbeddedAttemptSessionFileOwnersForTest();
3437
resetSessionWriteLockStateForTest();
3538
for (const dir of tempDirs.splice(0)) {
3639
await fs.rm(dir, { recursive: true, force: true });
@@ -46,6 +49,50 @@ async function createTempSessionFile(): Promise<string> {
4649
}
4750

4851
describe("embedded attempt session lock lifecycle", () => {
52+
it("serializes embedded attempts that share a session file owner", async () => {
53+
const sessionFile = await createTempSessionFile();
54+
const firstOwner = await acquireEmbeddedAttemptSessionFileOwner({ sessionFile });
55+
56+
let secondOwnerAcquired = false;
57+
const secondOwnerPromise = acquireEmbeddedAttemptSessionFileOwner({ sessionFile }).then(
58+
(owner) => {
59+
secondOwnerAcquired = true;
60+
return owner;
61+
},
62+
);
63+
64+
await Promise.resolve();
65+
expect(secondOwnerAcquired).toBe(false);
66+
67+
firstOwner.release();
68+
const secondOwner = await secondOwnerPromise;
69+
expect(secondOwnerAcquired).toBe(true);
70+
secondOwner.release();
71+
});
72+
73+
it("uses the same embedded attempt owner for symlinked session file paths", async () => {
74+
const sessionFile = await createTempSessionFile();
75+
const symlinkFile = path.join(path.dirname(sessionFile), "session-link.jsonl");
76+
await fs.symlink(sessionFile, symlinkFile);
77+
const firstOwner = await acquireEmbeddedAttemptSessionFileOwner({ sessionFile });
78+
79+
let symlinkOwnerAcquired = false;
80+
const symlinkOwnerPromise = acquireEmbeddedAttemptSessionFileOwner({
81+
sessionFile: symlinkFile,
82+
}).then((owner) => {
83+
symlinkOwnerAcquired = true;
84+
return owner;
85+
});
86+
87+
await Promise.resolve();
88+
expect(symlinkOwnerAcquired).toBe(false);
89+
90+
firstOwner.release();
91+
const symlinkOwner = await symlinkOwnerPromise;
92+
expect(symlinkOwnerAcquired).toBe(true);
93+
symlinkOwner.release();
94+
});
95+
4996
it("releases the coarse attempt lock before prompt submission and reacquires for cleanup", async () => {
5097
const releases: string[] = [];
5198
const acquireSessionWriteLock = vi
@@ -874,7 +921,7 @@ describe("embedded attempt session lock lifecycle", () => {
874921
expect(releaseForPrompt).toHaveBeenCalledTimes(1);
875922
expect(reacquireAfterPrompt).toHaveBeenCalledTimes(1);
876923
expect(streamFn).toHaveBeenCalledWith("model", "context");
877-
expect(events).toEqual(["drain", "release", "stream", "reacquire"]);
924+
expect(events).toEqual(["drain", "release", "stream", "drain", "reacquire"]);
878925
});
879926

880927
it("rewraps provider stream submission after the stream function is rebuilt", async () => {
@@ -921,17 +968,19 @@ describe("embedded attempt session lock lifecycle", () => {
921968

922969
expect(firstStreamFn).toHaveBeenCalledTimes(1);
923970
expect(secondStreamFn).toHaveBeenCalledTimes(1);
924-
expect(waitForSessionEvents).toHaveBeenCalledTimes(2);
971+
expect(waitForSessionEvents).toHaveBeenCalledTimes(4);
925972
expect(releaseForPrompt).toHaveBeenCalledTimes(2);
926973
expect(reacquireAfterPrompt).toHaveBeenCalledTimes(2);
927974
expect(events).toEqual([
928975
"drain",
929976
"release",
930977
"first-stream",
978+
"drain",
931979
"reacquire",
932980
"drain",
933981
"release",
934982
"second-stream",
983+
"drain",
935984
"reacquire",
936985
]);
937986
});

src/agents/pi-embedded-runner/run/attempt.session-lock.ts

Lines changed: 161 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,13 @@
11
import { AsyncLocalStorage } from "node:async_hooks";
22
import { statSync } from "node:fs";
33
import fs from "node:fs/promises";
4-
import path from "node:path";
54
import { isDeepStrictEqual } from "node:util";
65
import { withOwnedSessionTranscriptWrites } from "../../../config/sessions/transcript-write-context.js";
6+
import { resolveGlobalSingleton } from "../../../shared/global-singleton.js";
77
import { normalizeStringEntries } from "../../../shared/string-normalization.js";
88
import { isSessionWriteLockTimeoutError } from "../../session-write-lock-error.js";
99
import type { acquireSessionWriteLock } from "../../session-write-lock.js";
10+
import { resolveEmbeddedSessionFileKey } from "../session-file-key.js";
1011

1112
type SessionLock = Awaited<ReturnType<typeof acquireSessionWriteLock>>;
1213
type AcquireSessionWriteLock = typeof acquireSessionWriteLock;
@@ -403,7 +404,164 @@ const trustedSessionFileStates = new Map<string, TrustedSessionFileState>();
403404
let ownedSessionFileWriteGeneration = 0;
404405

405406
function resolveSessionFileFenceKey(sessionFile: string): string {
406-
return path.resolve(sessionFile);
407+
return resolveEmbeddedSessionFileKey(sessionFile);
408+
}
409+
410+
type SessionFileOwnerWaiter = {
411+
resolve: () => void;
412+
reject: (error: unknown) => void;
413+
timer?: NodeJS.Timeout;
414+
abortListener?: () => void;
415+
signal?: AbortSignal;
416+
};
417+
418+
type SessionFileOwnerEntry = {
419+
ownerId: symbol;
420+
waiters: Set<SessionFileOwnerWaiter>;
421+
};
422+
423+
type SessionFileOwnerState = {
424+
owners: Map<string, SessionFileOwnerEntry>;
425+
};
426+
427+
const EMBEDDED_ATTEMPT_SESSION_FILE_OWNER_STATE_KEY = Symbol.for(
428+
"openclaw.embeddedAttemptSessionFileOwnerState",
429+
);
430+
431+
const sessionFileOwnerState = resolveGlobalSingleton(
432+
EMBEDDED_ATTEMPT_SESSION_FILE_OWNER_STATE_KEY,
433+
(): SessionFileOwnerState => ({
434+
owners: new Map<string, SessionFileOwnerEntry>(),
435+
}),
436+
);
437+
438+
export type EmbeddedAttemptSessionFileOwner = {
439+
sessionFileKey: string;
440+
release(): void;
441+
};
442+
443+
export class EmbeddedAttemptSessionFileOwnerTimeoutError extends Error {
444+
constructor(sessionFile: string, timeoutMs: number) {
445+
super(`timed out waiting for embedded session file owner after ${timeoutMs}ms: ${sessionFile}`);
446+
this.name = "EmbeddedAttemptSessionFileOwnerTimeoutError";
447+
}
448+
}
449+
450+
function abortReason(signal: AbortSignal): unknown {
451+
return "reason" in signal ? (signal as { reason?: unknown }).reason : undefined;
452+
}
453+
454+
function abortOwnerWaitReason(signal: AbortSignal): unknown {
455+
return abortReason(signal) ?? new Error("operation aborted", { cause: signal });
456+
}
457+
458+
function waitForSessionFileOwnerRelease(params: {
459+
sessionFile: string;
460+
entry: SessionFileOwnerEntry;
461+
timeoutMs?: number;
462+
signal?: AbortSignal;
463+
}): Promise<void> {
464+
if (params.signal?.aborted) {
465+
return Promise.reject(abortOwnerWaitReason(params.signal));
466+
}
467+
return new Promise<void>((resolve, reject) => {
468+
const waiter: SessionFileOwnerWaiter = {
469+
resolve,
470+
reject,
471+
signal: params.signal,
472+
};
473+
const cleanup = () => {
474+
params.entry.waiters.delete(waiter);
475+
if (waiter.timer) {
476+
clearTimeout(waiter.timer);
477+
}
478+
if (waiter.signal && waiter.abortListener) {
479+
waiter.signal.removeEventListener("abort", waiter.abortListener);
480+
}
481+
};
482+
waiter.resolve = () => {
483+
cleanup();
484+
resolve();
485+
};
486+
waiter.reject = (error) => {
487+
cleanup();
488+
reject(error);
489+
};
490+
if (params.timeoutMs !== undefined && Number.isFinite(params.timeoutMs)) {
491+
waiter.timer = setTimeout(
492+
() => {
493+
waiter.reject(
494+
new EmbeddedAttemptSessionFileOwnerTimeoutError(
495+
params.sessionFile,
496+
params.timeoutMs ?? 0,
497+
),
498+
);
499+
},
500+
Math.max(1, Math.floor(params.timeoutMs)),
501+
);
502+
waiter.timer.unref?.();
503+
}
504+
if (params.signal) {
505+
waiter.abortListener = () => {
506+
waiter.reject(abortOwnerWaitReason(params.signal!));
507+
};
508+
params.signal.addEventListener("abort", waiter.abortListener, { once: true });
509+
}
510+
params.entry.waiters.add(waiter);
511+
});
512+
}
513+
514+
export async function acquireEmbeddedAttemptSessionFileOwner(params: {
515+
sessionFile: string;
516+
timeoutMs?: number;
517+
signal?: AbortSignal;
518+
}): Promise<EmbeddedAttemptSessionFileOwner> {
519+
const sessionFileKey = resolveEmbeddedSessionFileKey(params.sessionFile);
520+
const ownerId = Symbol(sessionFileKey);
521+
while (true) {
522+
if (params.signal?.aborted) {
523+
throw abortOwnerWaitReason(params.signal);
524+
}
525+
const entry = sessionFileOwnerState.owners.get(sessionFileKey);
526+
if (!entry) {
527+
sessionFileOwnerState.owners.set(sessionFileKey, {
528+
ownerId,
529+
waiters: new Set(),
530+
});
531+
return {
532+
sessionFileKey,
533+
release() {
534+
const current = sessionFileOwnerState.owners.get(sessionFileKey);
535+
if (!current || current.ownerId !== ownerId) {
536+
return;
537+
}
538+
sessionFileOwnerState.owners.delete(sessionFileKey);
539+
for (const waiter of current.waiters) {
540+
waiter.resolve();
541+
}
542+
},
543+
};
544+
}
545+
await waitForSessionFileOwnerRelease({
546+
sessionFile: params.sessionFile,
547+
entry,
548+
timeoutMs: params.timeoutMs,
549+
signal: params.signal,
550+
});
551+
}
552+
}
553+
554+
export function resetEmbeddedAttemptSessionFileOwnersForTest(): void {
555+
for (const entry of sessionFileOwnerState.owners.values()) {
556+
for (const waiter of entry.waiters) {
557+
waiter.reject(
558+
new Error("embedded attempt session file owners reset", {
559+
cause: "resetEmbeddedAttemptSessionFileOwnersForTest",
560+
}),
561+
);
562+
}
563+
}
564+
sessionFileOwnerState.owners.clear();
407565
}
408566

409567
function recordOwnedSessionFileWrite(
@@ -924,6 +1082,7 @@ export function installPromptSubmissionLockRelease(params: {
9241082
}
9251083
return await originalStreamFn(...args);
9261084
} finally {
1085+
await params.waitForSessionEvents(params.session);
9271086
await params.reacquireAfterPrompt();
9281087
}
9291088
};

0 commit comments

Comments
 (0)