Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 58 additions & 1 deletion apps/server/src/provider/Layers/ClaudeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4247,6 +4247,52 @@ describe("ClaudeAdapterLive", () => {
);
});

it.effect("rejects an echoed Claude resume cursor when the SDK starts a fresh session", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const adapter = yield* ClaudeAdapter;
const persistedResumeCursor = {
threadId: RESUME_THREAD_ID,
resume: "550e8400-e29b-41d4-a716-446655440000",
turnCount: 3,
};
const session = yield* adapter.startSession({
threadId: RESUME_THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
resumeCursor: persistedResumeCursor,
runtimeMode: "full-access",
});
assert.notEqual(adapter.isSameResumeCursor, undefined);
if (adapter.isSameResumeCursor === undefined) return;

const verification = yield* adapter
.isSameResumeCursor(RESUME_THREAD_ID, persistedResumeCursor, session.resumeCursor)
.pipe(Effect.forkChild);
harness.query.emit({
type: "system",
subtype: "init",
apiKeySource: "none",
claude_code_version: "test",
cwd: "/tmp/claude-adapter-test",
tools: [],
mcp_servers: [],
model: SYNTHETIC_CLAUDE_STANDARD_MODEL,
permissionMode: "bypassPermissions",
slash_commands: [],
output_style: "default",
skills: [],
plugins: [],
session_id: "7368d0c7-40a3-4d8a-bcc1-ac80c49f2719",
uuid: "fresh-init",
} as unknown as SDKMessage);

assert.equal(yield* Fiber.join(verification), false);
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
);
});

it.effect("preserves durable resume ids across Claude resume hooks", () => {
const harness = makeHarness();
return Effect.gen(function* () {
Expand All @@ -4259,7 +4305,7 @@ describe("ClaudeAdapterLive", () => {
Effect.forkChild,
);

yield* adapter.startSession({
const session = yield* adapter.startSession({
threadId: RESUME_THREAD_ID,
provider: ProviderDriverKind.make("claudeAgent"),
resumeCursor: {
Expand Down Expand Up @@ -4331,6 +4377,17 @@ describe("ClaudeAdapterLive", () => {
}
| undefined;
assert.equal(resumeCursor?.resume, durableSessionId);
assert.notEqual(adapter.isSameResumeCursor, undefined);
if (adapter.isSameResumeCursor !== undefined) {
assert.equal(
yield* adapter.isSameResumeCursor(
RESUME_THREAD_ID,
session.resumeCursor,
activeSessions[0]?.resumeCursor,
),
true,
);
}
}).pipe(
Effect.provideService(Random.Random, makeDeterministicRandomService()),
Effect.provide(harness.layer),
Expand Down
23 changes: 23 additions & 0 deletions apps/server/src/provider/Layers/ClaudeAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as FileSystem from "effect/FileSystem";
import * as Fiber from "effect/Fiber";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
Expand Down Expand Up @@ -291,6 +292,7 @@ interface ClaudeSessionContext {
* effort override inherit this. */
currentEffort: string | undefined;
resumeSessionId: string | undefined;
readonly confirmedResumeSessionId: Deferred.Deferred<string>;
readonly pendingApprovals: Map<ApprovalRequestId, PendingApproval>;
readonly pendingUserInputs: Map<ApprovalRequestId, PendingUserInput>;
readonly turns: Array<{
Expand Down Expand Up @@ -2047,6 +2049,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
const nextThreadId = message.session_id;
context.resumeSessionId = message.session_id;
yield* updateResumeCursor(context);
yield* Deferred.succeed(context.confirmedResumeSessionId, message.session_id);

if (context.lastThreadStartedId !== nextThreadId) {
context.lastThreadStartedId = nextThreadId;
Expand Down Expand Up @@ -3886,6 +3889,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
const runPromise = Effect.runPromiseWith(runtimeContext);

const promptQueue = yield* Queue.unbounded<PromptQueueItem>();
const confirmedResumeSessionId = yield* Deferred.make<string>();
const prompt = Stream.fromQueue(promptQueue).pipe(
Stream.filter((item) => item.type === "message"),
Stream.map((item) => item.message),
Expand Down Expand Up @@ -4461,6 +4465,7 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
currentApiModelId: apiModelId,
currentEffort: effectiveEffort ?? undefined,
resumeSessionId: sessionId,
confirmedResumeSessionId,
pendingApprovals,
pendingUserInputs,
turns: [],
Expand Down Expand Up @@ -4804,6 +4809,24 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* (
sessionModelSwitch: "in-session",
},
startSession,
isSameResumeCursor: (threadId, persisted, recovered) =>
Effect.gen(function* () {
const persistedResume = readClaudeResumeState(persisted)?.resume;
const recoveredResume = readClaudeResumeState(recovered)?.resume;
const context = sessions.get(threadId);
if (
persistedResume === undefined ||
recoveredResume === undefined ||
persistedResume !== recoveredResume ||
context === undefined
) {
return false;
}
const confirmed = yield* Deferred.await(context.confirmedResumeSessionId).pipe(
Effect.timeoutOption("10 seconds"),
);
return Option.isSome(confirmed) && confirmed.value === persistedResume;
}),
sendTurn,
interruptTurn,
readThread,
Expand Down
6 changes: 6 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2324,6 +2324,12 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
promptlessTurnContinuation: true,
},
startSession,
isSameResumeCursor: (_threadId, persisted, recovered) =>
Effect.succeed(
isCodexResumeCursorSchema(persisted) &&
isCodexResumeCursorSchema(recovered) &&
persisted.threadId === recovered.threadId,
),
sendTurn,
interruptTurn,
readThread,
Expand Down
9 changes: 9 additions & 0 deletions apps/server/src/provider/Layers/CursorAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1205,6 +1205,15 @@ export function makeCursorAdapter(
provider: PROVIDER,
capabilities: { sessionModelSwitch: "in-session" },
startSession,
isSameResumeCursor: (_threadId, persisted, recovered) => {
const persistedSession = parseCursorResume(persisted);
const recoveredSession = parseCursorResume(recovered);
return Effect.succeed(
persistedSession !== undefined &&
recoveredSession !== undefined &&
persistedSession.sessionId === recoveredSession.sessionId,
);
},
sendTurn,
interruptTurn,
readThread,
Expand Down
9 changes: 9 additions & 0 deletions apps/server/src/provider/Layers/GrokAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2114,6 +2114,15 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte
provider: PROVIDER,
capabilities: { sessionModelSwitch: "in-session" },
startSession,
isSameResumeCursor: (_threadId, persisted, recovered) => {
const persistedSession = parseGrokResume(persisted);
const recoveredSession = parseGrokResume(recovered);
return Effect.succeed(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 High Layers/GrokAdapter.ts:2120

isSameResumeCursor returns true for any syntactically valid persisted cursor, even when session/load has fallen back to a fresh conversation. resumed.resumeCursor is reconstructed from options.resumeSessionId, so the comparison cannot validate the loaded session and restart recovery proceeds in the wrong conversation instead of failing closed. Compare against the session identity returned by the load operation, or otherwise detect fallback before accepting recovery.

🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/provider/Layers/GrokAdapter.ts around line 2120:

`isSameResumeCursor` returns `true` for any syntactically valid persisted cursor, even when `session/load` has fallen back to a fresh conversation. `resumed.resumeCursor` is reconstructed from `options.resumeSessionId`, so the comparison cannot validate the loaded session and restart recovery proceeds in the wrong conversation instead of failing closed. Compare against the session identity returned by the load operation, or otherwise detect fallback before accepting recovery.

persistedSession !== undefined &&
recoveredSession !== undefined &&
persistedSession.sessionId === recoveredSession.sessionId,
);
},
sendTurn,
interruptTurn,
readThread,
Expand Down
9 changes: 9 additions & 0 deletions apps/server/src/provider/Layers/OpenCodeAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3267,6 +3267,15 @@ export function makeOpenCodeAdapter(
sessionModelSwitch: "in-session",
},
startSession,
isSameResumeCursor: (_threadId, persisted, recovered) => {
const persistedSession = parseOpenCodeResume(persisted);
const recoveredSession = parseOpenCodeResume(recovered);
return Effect.succeed(
persistedSession !== undefined &&
recoveredSession !== undefined &&
persistedSession.sessionId === recoveredSession.sessionId,
);
},
sendTurn,
interruptTurn,
respondToRequest,
Expand Down
67 changes: 62 additions & 5 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,19 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) {
...(provider === CODEX_DRIVER ? { promptlessTurnContinuation: true } : {}),
},
startSession,
isSameResumeCursor: (_threadId, persisted, recovered) => {
const readOpaque = (value: unknown) =>
typeof value === "object" &&
value !== null &&
!Array.isArray(value) &&
typeof (value as { opaque?: unknown }).opaque === "string"
? (value as { opaque: string }).opaque
: undefined;
const persistedOpaque = readOpaque(persisted);
return Effect.succeed(
persistedOpaque !== undefined && persistedOpaque === readOpaque(recovered),
);
},
sendTurn,
interruptTurn,
respondToRequest,
Expand Down Expand Up @@ -1464,11 +1477,14 @@ routing.layer("ProviderServiceLive routing", (it) => {
routing.codex.startSession.mockClear();
routing.codex.sendTurn.mockClear();

yield* provider.sendTurn({
threadId: initial.threadId,
input: "resume",
attachments: [],
});
yield* provider.sendTurn(
{
threadId: initial.threadId,
input: "resume",
attachments: [],
},
{ requireResumeCursor: initial.resumeCursor },
);

assert.equal(routing.codex.startSession.mock.calls.length, 1);
const resumedStartInput = routing.codex.startSession.mock.calls[0]?.[0];
Expand All @@ -1489,6 +1505,47 @@ routing.layer("ProviderServiceLive routing", (it) => {
}),
);

it.effect("rejects automatic continuation when recovery starts a fresh conversation", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
const threadId = asThreadId("thread-strict-resume-mismatch");
const initial = yield* provider.startSession(threadId, {
provider: CODEX_DRIVER,
providerInstanceId: codexInstanceId,
threadId,
cwd: "/tmp/project-strict-resume",
runtimeMode: "full-access",
});

yield* routing.codex.stopAll();
routing.codex.startSession.mockClear();
routing.codex.sendTurn.mockClear();
routing.codex.stopSession.mockClear();
routing.codex.startSession.mockImplementationOnce(() =>
Effect.succeed({
...initial,
resumeCursor: { opaque: "fresh-provider-conversation" },
}),
);

const failure = yield* provider
.sendTurn(
{
threadId,
input: "Continue where you left off.",
attachments: [],
},
{ requireResumeCursor: initial.resumeCursor },
)
.pipe(Effect.flip);

assert.instanceOf(failure, ProviderValidationError);
assert.match(failure.message, /did not resume the persisted conversation/u);
assert.equal(routing.codex.sendTurn.mock.calls.length, 0);
assert.deepStrictEqual(routing.codex.stopSession.mock.calls, [[threadId]]);
}),
);

it.effect("recovers stale claudeAgent sessions for sendTurn using persisted cwd", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
Expand Down
Loading
Loading