diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts index 38bdb7f1f2e7..990afedb074c 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.test.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.test.ts @@ -315,6 +315,57 @@ async function readPromptMessages( const THREAD_ID = ThreadId.make("thread-claude-1"); const RESUME_THREAD_ID = ThreadId.make("thread-claude-resume"); const SYNTHETIC_SUBAGENT_MODEL = "claude-synthetic-subagent[expanded]"; +const CLAUDE_ORIGINAL_SESSION_ID = "550e8400-e29b-41d4-a716-446655440010"; +const CLAUDE_FORK_SESSION_ID = "550e8400-e29b-41d4-a716-446655440020"; + +function claudeHistoryMessage(input: { + readonly type: "user" | "assistant" | "system"; + readonly uuid: string; + readonly sessionId?: string; + readonly content?: unknown; + readonly parentToolUseId?: string | null; +}) { + const content = + input.content ?? + (input.type === "user" ? "prompt" : input.type === "assistant" ? [] : { subtype: "init" }); + return { + type: input.type, + uuid: input.uuid, + session_id: input.sessionId ?? CLAUDE_ORIGINAL_SESSION_ID, + parent_tool_use_id: input.parentToolUseId ?? null, + parent_agent_id: null, + message: input.type === "system" ? content : { content }, + }; +} + +const sendCompletedClaudeTurn = ( + adapter: ClaudeAdapterShape, + harness: ReturnType, + threadId: ThreadId, + input: string, +) => + Effect.gen(function* () { + const turn = yield* adapter.sendTurn({ + threadId, + input, + attachments: [], + }); + const completedFiber = yield* Stream.filter( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runHead, Effect.forkChild); + harness.queries.at(-1)!.emit({ + type: "result", + subtype: "success", + is_error: false, + errors: [], + session_id: CLAUDE_ORIGINAL_SESSION_ID, + uuid: `result-${turn.turnId}`, + } as unknown as SDKMessage); + const completed = yield* Fiber.join(completedFiber); + assert.equal(completed._tag, "Some"); + return turn; + }); describe("ClaudeAdapterLive", () => { it.effect("returns validation error for non-claude provider on startSession", () => { @@ -6484,6 +6535,589 @@ describe("ClaudeAdapterLive", () => { ); }); + it.effect("rewinds Claude history when the fork omits retained system messages", () => { + const forkCalls: Array>> = []; + let firstTurnId = ""; + let secondTurnId = ""; + const harness = makeHarness({ + forkSession: async (...args) => { + forkCalls.push(args); + return { sessionId: CLAUDE_FORK_SESSION_ID }; + }, + getSessionMessages: async (sessionId) => { + if (sessionId === CLAUDE_FORK_SESSION_ID) { + return [ + claudeHistoryMessage({ + type: "user", + uuid: `fork-${firstTurnId}`, + sessionId, + content: "first", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-1", + sessionId, + }), + ]; + } + return [ + claudeHistoryMessage({ type: "system", uuid: "system-init" }), + claudeHistoryMessage({ type: "user", uuid: firstTurnId, content: "first" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1" }), + claudeHistoryMessage({ + type: "system", + uuid: "compact-boundary", + content: { subtype: "compact_boundary" }, + }), + claudeHistoryMessage({ type: "user", uuid: secondTurnId, content: "second" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-2" }), + ]; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + firstTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first")) + .turnId; + secondTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "second")) + .turnId; + + const snapshot = yield* adapter.rollbackThread(session.threadId, 1); + assert.equal(snapshot.turns.length, 1); + assert.deepEqual(forkCalls, [ + [CLAUDE_ORIGINAL_SESSION_ID, { upToMessageId: "compact-boundary" }], + ]); + assert.deepEqual((yield* adapter.listSessions())[0]?.resumeCursor, { + threadId: session.threadId, + resume: CLAUDE_FORK_SESSION_ID, + turnCount: 1, + turnStartMessageIds: [`fork-${firstTurnId}`], + }); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("rewinds Claude history when the fork inserts extra system messages", () => { + let firstTurnId = ""; + let secondTurnId = ""; + const harness = makeHarness({ + forkSession: async () => ({ sessionId: CLAUDE_FORK_SESSION_ID }), + getSessionMessages: async (sessionId) => { + if (sessionId === CLAUDE_FORK_SESSION_ID) { + return [ + claudeHistoryMessage({ + type: "system", + uuid: "fork-system-init", + sessionId, + }), + claudeHistoryMessage({ + type: "user", + uuid: `fork-${firstTurnId}`, + sessionId, + content: "first", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-1", + sessionId, + }), + ]; + } + return [ + claudeHistoryMessage({ type: "user", uuid: firstTurnId, content: "first" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1" }), + claudeHistoryMessage({ type: "user", uuid: secondTurnId, content: "second" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-2" }), + ]; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + firstTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first")) + .turnId; + secondTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "second")) + .turnId; + + yield* adapter.rollbackThread(session.threadId, 1); + assert.deepEqual((yield* adapter.listSessions())[0]?.resumeCursor, { + threadId: session.threadId, + resume: CLAUDE_FORK_SESSION_ID, + turnCount: 1, + turnStartMessageIds: [`fork-${firstTurnId}`], + }); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("rewinds Claude history when the fork keeps a compact conversation prefix", () => { + let firstTurnId = ""; + let secondTurnId = ""; + const harness = makeHarness({ + forkSession: async () => ({ sessionId: CLAUDE_FORK_SESSION_ID }), + getSessionMessages: async (sessionId) => { + if (sessionId === CLAUDE_FORK_SESSION_ID) { + return [ + claudeHistoryMessage({ + type: "user", + uuid: "fork-compacted-user", + sessionId, + content: "earlier compacted turn", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-compacted-assistant", + sessionId, + }), + claudeHistoryMessage({ + type: "user", + uuid: `fork-${firstTurnId}`, + sessionId, + content: "first", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-1", + sessionId, + }), + ]; + } + return [ + claudeHistoryMessage({ type: "user", uuid: firstTurnId, content: "first" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1" }), + claudeHistoryMessage({ type: "user", uuid: secondTurnId, content: "second" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-2" }), + ]; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + firstTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first")) + .turnId; + secondTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "second")) + .turnId; + + yield* adapter.rollbackThread(session.threadId, 1); + assert.deepEqual((yield* adapter.listSessions())[0]?.resumeCursor, { + threadId: session.threadId, + resume: CLAUDE_FORK_SESSION_ID, + turnCount: 1, + turnStartMessageIds: [`fork-${firstTurnId}`], + }); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("rewinds a Claude turn that includes tool results and a later steer", () => { + const forkCalls: Array>> = []; + let firstTurnId = ""; + let secondTurnId = ""; + const harness = makeHarness({ + forkSession: async (...args) => { + forkCalls.push(args); + return { sessionId: CLAUDE_FORK_SESSION_ID }; + }, + getSessionMessages: async (sessionId) => { + const history = [ + claudeHistoryMessage({ type: "user", uuid: firstTurnId, content: "first" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1-tool" }), + claudeHistoryMessage({ + type: "user", + uuid: "tool-result-1", + content: [{ type: "tool_result" }], + }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1-final" }), + claudeHistoryMessage({ type: "user", uuid: secondTurnId, content: "second" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-2" }), + claudeHistoryMessage({ + type: "user", + uuid: "steer", + sessionId, + content: "steer the second turn", + }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-steer", sessionId }), + ]; + return sessionId === CLAUDE_FORK_SESSION_ID + ? history.slice(0, 4).map((message) => ({ + ...message, + uuid: `fork-${message.uuid}`, + session_id: sessionId, + })) + : history; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + firstTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first")) + .turnId; + const secondTurn = yield* adapter.sendTurn({ + threadId: session.threadId, + input: "second", + attachments: [], + }); + secondTurnId = secondTurn.turnId; + const steer = yield* adapter.sendTurn({ + threadId: session.threadId, + input: "steer the second turn", + attachments: [], + }); + assert.equal(steer.turnId, secondTurn.turnId); + const secondCompletedFiber = yield* Stream.filter( + adapter.streamEvents, + (event) => event.type === "turn.completed", + ).pipe(Stream.runHead, Effect.forkChild); + harness.queries.at(-1)!.emit({ + type: "result", + subtype: "success", + is_error: false, + errors: [], + session_id: CLAUDE_ORIGINAL_SESSION_ID, + uuid: "result-second", + } as unknown as SDKMessage); + yield* Fiber.join(secondCompletedFiber); + + const snapshot = yield* adapter.rollbackThread(session.threadId, 1); + assert.equal(snapshot.turns.length, 1); + assert.deepEqual(forkCalls, [ + [CLAUDE_ORIGINAL_SESSION_ID, { upToMessageId: "assistant-1-final" }], + ]); + assert.deepEqual((yield* adapter.listSessions())[0]?.resumeCursor, { + threadId: session.threadId, + resume: CLAUDE_FORK_SESSION_ID, + turnCount: 1, + turnStartMessageIds: [`fork-${firstTurnId}`], + }); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("rewinds the last two turns of a three-turn Claude thread", () => { + let firstTurnId = ""; + let secondTurnId = ""; + let thirdTurnId = ""; + const harness = makeHarness({ + forkSession: async () => ({ sessionId: CLAUDE_FORK_SESSION_ID }), + getSessionMessages: async (sessionId) => { + const history = [ + claudeHistoryMessage({ type: "system", uuid: "system-init" }), + claudeHistoryMessage({ type: "user", uuid: firstTurnId, content: "first" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1" }), + claudeHistoryMessage({ type: "user", uuid: secondTurnId, content: "second" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-2" }), + claudeHistoryMessage({ type: "user", uuid: thirdTurnId, content: "third" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-3" }), + ]; + return sessionId === CLAUDE_FORK_SESSION_ID + ? [ + claudeHistoryMessage({ + type: "user", + uuid: `fork-${firstTurnId}`, + sessionId, + content: "first", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-1", + sessionId, + }), + ] + : history; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + firstTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first")) + .turnId; + secondTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "second")) + .turnId; + thirdTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "third")) + .turnId; + + const snapshot = yield* adapter.rollbackThread(session.threadId, 2); + assert.equal(snapshot.turns.length, 1); + assert.deepEqual((yield* adapter.listSessions())[0]?.resumeCursor, { + threadId: session.threadId, + resume: CLAUDE_FORK_SESSION_ID, + turnCount: 1, + turnStartMessageIds: [`fork-${firstTurnId}`], + }); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("rewinds only the latest turn of a three-turn Claude thread", () => { + let firstTurnId = ""; + let secondTurnId = ""; + let thirdTurnId = ""; + const harness = makeHarness({ + forkSession: async () => ({ sessionId: CLAUDE_FORK_SESSION_ID }), + getSessionMessages: async (sessionId) => { + const history = [ + claudeHistoryMessage({ type: "system", uuid: "system-init" }), + claudeHistoryMessage({ type: "user", uuid: firstTurnId, content: "first" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1" }), + claudeHistoryMessage({ type: "user", uuid: secondTurnId, content: "second" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-2" }), + claudeHistoryMessage({ type: "user", uuid: thirdTurnId, content: "third" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-3" }), + ]; + return sessionId === CLAUDE_FORK_SESSION_ID + ? [ + claudeHistoryMessage({ + type: "user", + uuid: `fork-${firstTurnId}`, + sessionId, + content: "first", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-1", + sessionId, + }), + claudeHistoryMessage({ + type: "system", + uuid: "fork-notice", + sessionId, + }), + claudeHistoryMessage({ + type: "user", + uuid: `fork-${secondTurnId}`, + sessionId, + content: "second", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-2", + sessionId, + }), + ] + : history; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + firstTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first")) + .turnId; + secondTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "second")) + .turnId; + thirdTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "third")) + .turnId; + + const snapshot = yield* adapter.rollbackThread(session.threadId, 1); + assert.equal(snapshot.turns.length, 2); + assert.deepEqual((yield* adapter.listSessions())[0]?.resumeCursor, { + threadId: session.threadId, + resume: CLAUDE_FORK_SESSION_ID, + turnCount: 2, + turnStartMessageIds: [`fork-${firstTurnId}`, `fork-${secondTurnId}`], + }); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("resets a Claude thread when rewind removes every recorded turn", () => { + const forkCalls: Array>> = []; + const harness = makeHarness({ + forkSession: async (...args) => { + forkCalls.push(args); + return { sessionId: CLAUDE_FORK_SESSION_ID }; + }, + getSessionMessages: async () => [ + claudeHistoryMessage({ type: "user", uuid: "unused", content: "first" }), + ], + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first"); + yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "second"); + + const snapshot = yield* adapter.rollbackThread(session.threadId, 2); + assert.equal(snapshot.turns.length, 0); + assert.equal(forkCalls.length, 0); + const resetOptions = harness.getLastCreateQueryInput()?.options; + assert.equal(resetOptions?.resume, undefined); + assert.ok(resetOptions?.sessionId); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + it.effect("rejects a Claude fork that drops a retained user turn", () => { + let firstTurnId = ""; + let secondTurnId = ""; + const harness = makeHarness({ + forkSession: async () => ({ sessionId: CLAUDE_FORK_SESSION_ID }), + getSessionMessages: async (sessionId) => { + if (sessionId === CLAUDE_FORK_SESSION_ID) { + return [ + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-1", + sessionId, + }), + ]; + } + return [ + claudeHistoryMessage({ type: "user", uuid: firstTurnId, content: "first" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-1" }), + claudeHistoryMessage({ type: "user", uuid: secondTurnId, content: "second" }), + claudeHistoryMessage({ type: "assistant", uuid: "assistant-2" }), + ]; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + firstTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "first")) + .turnId; + secondTurnId = (yield* sendCompletedClaudeTurn(adapter, harness, session.threadId, "second")) + .turnId; + + const error = yield* adapter.rollbackThread(session.threadId, 1).pipe(Effect.flip); + assert.match(error.message, /did not preserve the retained turn boundaries/); + assert.equal(harness.queries.at(-1)?.closeCalls, 0); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + + for (const scenario of [ + "restores conversation messages between retained turns", + "changes retained assistant content without changing its role", + ] as const) { + it.effect(`rejects a Claude fork that ${scenario}`, () => { + const turnIds: Array = []; + const harness = makeHarness({ + forkSession: async () => ({ sessionId: CLAUDE_FORK_SESSION_ID }), + getSessionMessages: async (sessionId) => { + const history = turnIds.flatMap((turnId, index) => [ + claudeHistoryMessage({ + type: "user", + uuid: turnId, + content: `prompt ${index + 1}`, + }), + claudeHistoryMessage({ + type: "assistant", + uuid: `assistant-${index + 1}`, + content: [{ type: "text", text: `reply ${index + 1}` }], + }), + ]); + if (sessionId !== CLAUDE_FORK_SESSION_ID) return history; + const forkHistory = history.slice(0, 4).map((message) => ({ + ...message, + uuid: `fork-${message.uuid}`, + session_id: sessionId, + })); + if (scenario === "restores conversation messages between retained turns") { + forkHistory.splice( + 2, + 0, + claudeHistoryMessage({ + type: "user", + uuid: "fork-restored-steer", + sessionId, + content: "earlier steer omitted by compaction", + }), + claudeHistoryMessage({ + type: "assistant", + uuid: "fork-restored-reply", + sessionId, + content: [{ type: "text", text: "earlier steering reply" }], + }), + ); + } else { + forkHistory[1] = claudeHistoryMessage({ + type: "assistant", + uuid: "fork-assistant-1", + sessionId, + content: [{ type: "text", text: "different retained reply" }], + }); + } + return forkHistory; + }, + }); + return Effect.gen(function* () { + const adapter = yield* ClaudeAdapter; + const session = yield* adapter.startSession({ + threadId: THREAD_ID, + provider: ProviderDriverKind.make("claudeAgent"), + runtimeMode: "full-access", + }); + for (let index = 0; index < 3; index++) { + const turn = yield* sendCompletedClaudeTurn( + adapter, + harness, + session.threadId, + `prompt ${index + 1}`, + ); + turnIds.push(turn.turnId); + } + const cursorBeforeRollback = (yield* adapter.listSessions())[0]?.resumeCursor; + const queryCountBeforeRollback = harness.queries.length; + + const error = yield* adapter.rollbackThread(session.threadId, 1).pipe(Effect.flip); + assert.match(error.message, /did not preserve the retained turn boundaries/); + assert.equal(harness.queries.length, queryCountBeforeRollback); + assert.equal(harness.queries.at(-1)?.closeCalls, 0); + assert.deepEqual((yield* adapter.listSessions())[0]?.resumeCursor, cursorBeforeRollback); + }).pipe( + Effect.provideService(Random.Random, makeDeterministicRandomService()), + Effect.provide(harness.layer), + ); + }); + } + it.effect("updates model on sendTurn when model override is provided", () => { const harness = makeHarness(); return Effect.gen(function* () { diff --git a/apps/server/src/provider/Layers/ClaudeAdapter.ts b/apps/server/src/provider/Layers/ClaudeAdapter.ts index 3afde2b39a7a..57f7abd8864d 100644 --- a/apps/server/src/provider/Layers/ClaudeAdapter.ts +++ b/apps/server/src/provider/Layers/ClaudeAdapter.ts @@ -8,6 +8,7 @@ * @module ClaudeAdapterLive */ +import * as NodeUtil from "node:util"; import { type CanUseTool, query, @@ -135,6 +136,94 @@ const decodeSessionMessages = Schema.decodeSync( ), ); +type ClaudeHistoryMessage = { + readonly type: string; + readonly uuid: string; + readonly parent_tool_use_id: string | null; + readonly message: unknown; +}; + +const isClaudeConversationMessage = (message: ClaudeHistoryMessage): boolean => + message.type === "user" || message.type === "assistant"; + +const isClaudeHumanTurnStart = (message: ClaudeHistoryMessage): boolean => { + if (message.type !== "user" || message.parent_tool_use_id !== null) return false; + const body = message.message; + if (typeof body !== "object" || body === null || !("content" in body)) return false; + const content = body.content; + return ( + typeof content === "string" || + (Array.isArray(content) && + content.some( + (part: unknown) => + typeof part === "object" && + part !== null && + "type" in part && + part.type !== "tool_result", + )) + ); +}; + +const conversationIndexForUuid = ( + messages: ReadonlyArray, + uuid: string, +): number => { + let index = -1; + for (const message of messages) { + if (!isClaudeConversationMessage(message)) continue; + index += 1; + if (message.uuid === uuid) return index; + } + return -1; +}; + +// Native forks rewrite every UUID. getSessionMessages then rebuilds the +// parentUuid chain, so system notices and compact metadata can change the +// raw length without dropping retained user/assistant turns. Align those +// conversation messages from the truncated end, then remap T3 turn starts. +const remapClaudeForkTurnBoundaries = ( + messages: ReadonlyArray, + forkMessages: ReadonlyArray, + firstRemoved: number, + retainedBoundaries: ReadonlyArray, +): Array | undefined => { + const retainedConversation = messages.slice(0, firstRemoved).filter(isClaudeConversationMessage); + const forkConversation = forkMessages.filter(isClaudeConversationMessage); + if (retainedConversation.length === 0) { + return retainedBoundaries.every((id) => id === null) ? [...retainedBoundaries] : undefined; + } + const offset = forkConversation.length - retainedConversation.length; + // Forks preserve message bodies. Matching roles alone can mistake a restored + // steering message for a retained turn when compaction changes the chain. + if ( + offset < 0 || + retainedConversation.some((message, index) => { + const forkMessage = forkConversation[index + offset]; + return ( + forkMessage === undefined || + forkMessage.type !== message.type || + !NodeUtil.isDeepStrictEqual(forkMessage.message, message.message) + ); + }) + ) { + return undefined; + } + const remapped = retainedBoundaries.map((originalId) => { + if (originalId === null) return null; + const originalIndex = conversationIndexForUuid(messages, originalId); + const forkIndex = originalIndex + offset; + const forkMessage = + originalIndex >= 0 && forkIndex >= 0 ? forkConversation[forkIndex] : undefined; + const originalMessage = messages.find((message) => message.uuid === originalId); + return forkMessage !== undefined && + originalMessage !== undefined && + forkMessage.type === originalMessage.type + ? forkMessage.uuid + : null; + }); + return remapped.some((id) => id === null) ? undefined : remapped; +}; + const PROVIDER = ProviderDriverKind.make("claudeAgent"); type ClaudeTextStreamKind = Extract; type ClaudeToolResultStreamKind = Extract< @@ -5190,23 +5279,9 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( }); const messages = yield* readHistory(sessionId); // Tool results are user-role messages too. Only human prompts begin a turn. - const turnStarts = messages.flatMap((message, index) => { - if (message.type !== "user" || message.parent_tool_use_id !== null) return []; - const body = message.message; - if (typeof body !== "object" || body === null || !("content" in body)) return []; - const content = body.content; - return typeof content === "string" || - (Array.isArray(content) && - content.some( - (part: unknown) => - typeof part === "object" && - part !== null && - "type" in part && - part.type !== "tool_result", - )) - ? [index] - : []; - }); + const turnStarts = messages.flatMap((message, index) => + isClaudeHumanTurnStart(message) ? [index] : [], + ); if (messages.length === 0) { return yield* new ProviderAdapterRequestError({ provider: PROVIDER, @@ -5263,20 +5338,20 @@ export const makeClaudeAdapter = Effect.fn("makeClaudeAdapter")(function* ( const retainedBoundaries = boundaries.slice(0, retainedCount); if (fork) { const forkMessages = yield* readHistory(fork.sessionId); - if (forkMessages.length !== firstRemoved) { + const remappedBoundaries = remapClaudeForkTurnBoundaries( + messages, + forkMessages, + firstRemoved, + retainedBoundaries, + ); + if (!remappedBoundaries) { return yield* new ProviderAdapterRequestError({ provider: PROVIDER, method: "thread/rollback", detail: "Claude fork history did not preserve the retained turn boundaries.", }); } - // Native forks replace every UUID while preserving transcript order. - for (let index = 0; index < retainedBoundaries.length; index++) { - const messageIndex = messages.findIndex( - (message) => message.uuid === retainedBoundaries[index], - ); - retainedBoundaries[index] = forkMessages[messageIndex]?.uuid ?? null; - } + retainedBoundaries.splice(0, retainedBoundaries.length, ...remappedBoundaries); } yield* stopSessionInternal(context, { emitExitEvent: false }); yield* startSession({