diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-forward.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-forward.ts index a73e8b2111..b9559868f4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-forward.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-forward.ts @@ -12,7 +12,10 @@ import type { StructuredAgentSessionHandoffFlowContext, StructuredTuiOwner } from './structured-agent-session-handoff-types' -import { StructuredTuiLaunchCleanupError } from './structured-agent-session-handoff-types' +import { + StructuredTuiCatchupStoppedError, + StructuredTuiLaunchCleanupError +} from './structured-agent-session-handoff-types' export async function handoffStructuredSessionToTui( context: StructuredAgentSessionHandoffFlowContext, @@ -75,7 +78,8 @@ export async function handoffStructuredSessionToTui( let owner: StructuredTuiOwner | null = null let processIdentityCommitted = false try { - await deps.prepareTuiHistoryCatchup?.(sessionId, record.lease.runtimeFence) + const prepared = await deps.prepareTuiHistoryCatchup?.(sessionId, record.lease.runtimeFence) + prepared?.throwIfAborted() owner = await deps.transport!.launchTui({ record, fence: record.lease.runtimeFence, @@ -91,6 +95,7 @@ export async function handoffStructuredSessionToTui( processIdentityCommitted = true } }) + prepared?.throwIfAborted() if (!processIdentityCommitted) { await deps.store.commitProcessIdentity({ sessionId, @@ -128,6 +133,16 @@ export async function handoffStructuredSessionToTui( ) } } + if (error instanceof StructuredTuiCatchupStoppedError && (owner || !processIdentityCommitted)) { + await abandonStoredAgentSessionHandoffAttempt(deps.store, { + sessionId, + expectedFence: record.lease.runtimeFence, + operationId, + recoverableRuntimeKind: 'native', + now: deps.now() + }) + throw error + } await recoverNativeAfterTuiFailure(context, sessionId, operationId) throw error } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart-tui.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart-tui.ts index c16ac68122..3bf0bc08e8 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart-tui.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart-tui.ts @@ -96,3 +96,16 @@ export async function persistReprovedTuiOwner( }) } } + +export async function startRecoveredTuiCatchup( + input: StructuredAgentSessionRestartAccess, + record: AgentSessionRecord +): Promise { + const prepared = await input.deps.recoverTuiHistoryCatchup?.( + record.sessionId, + record.lease.runtimeFence + ) + prepared?.throwIfAborted() + await input.deps.activateTuiHistoryCatchup?.(record.sessionId) + prepared?.throwIfAborted() +} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart.ts index 13a26f2a7a..a6ad92d91e 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-handoff-restart.ts @@ -12,10 +12,12 @@ import { structuredTuiRecoveryProofIsAdmissible } from './structured-agent-session-handoff-status' import type { StructuredTuiOwner } from './structured-agent-session-handoff-types' +import { StructuredTuiCatchupStoppedError } from './structured-agent-session-handoff-types' import { persistReprovedTuiOwner, recoverTuiOwnerOrContinue, recoverUnavailableTuiAsNative, + startRecoveredTuiCatchup, type StructuredAgentSessionRestartAccess } from './structured-agent-session-handoff-restart-tui' @@ -57,6 +59,15 @@ export async function restoreStructuredAgentSessionHandoff( } return } catch (error) { + if (error instanceof StructuredTuiCatchupStoppedError) { + if (operationId) { + await input.deps.store.recordOperationOutcome({ + operationId, + outcome: { status: 'failed', code: 'agent_session_handoff_failed' } + }) + } + throw error + } lastError = error if (attempt < 2) { await new Promise((resolve) => setTimeout(resolve, 100 * 2 ** attempt)) @@ -278,14 +289,6 @@ async function restoreProving(input: RestartAccess, record: AgentSessionRecord): await continueHandoff(input, stopped) } -async function startRecoveredTuiCatchup( - input: RestartAccess, - record: AgentSessionRecord -): Promise { - await input.deps.recoverTuiHistoryCatchup?.(record.sessionId, record.lease.runtimeFence) - await input.deps.activateTuiHistoryCatchup?.(record.sessionId) -} - async function continueHandoff(input: RestartAccess, record: AgentSessionRecord): Promise { const direction = record.lease.runtimeKind === 'native' ? 'to-tui' : 'to-native' const operationId = record.lease.handoffOperationId! diff --git a/src/main/native-chat/agent-session-wire/structured-tui-transcript-catchup.ts b/src/main/native-chat/agent-session-wire/structured-tui-transcript-catchup.ts index cc343c9231..10ce2416a3 100644 --- a/src/main/native-chat/agent-session-wire/structured-tui-transcript-catchup.ts +++ b/src/main/native-chat/agent-session-wire/structured-tui-transcript-catchup.ts @@ -18,12 +18,15 @@ import { type NativeChatTranscriptSubscription } from '../transcript-watch' import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types' +import { StructuredTuiCatchupStoppedError } from './structured-agent-session-handoff-types' import { readStructuredTuiTranscriptBoundary, writeStructuredTuiTranscriptBoundary } from './structured-tui-transcript-boundary' type CatchupState = { + controller: AbortController + initialReady: (() => void) | null active: boolean fence: number agent: AgentSessionHandleProvider @@ -35,6 +38,7 @@ type CatchupState = { export class StructuredTuiTranscriptCatchup { private readonly states = new Map() + private readonly teardown = new AbortController() constructor( private readonly input: { @@ -47,15 +51,16 @@ export class StructuredTuiTranscriptCatchup { } ) {} - async prepare(sessionId: string, fence: number): Promise { - await this.start(sessionId, fence, false) + async prepare(sessionId: string, fence: number): Promise { + return this.start(sessionId, fence, false) } - async recover(sessionId: string, fence: number): Promise { - await this.start(sessionId, fence, true) + async recover(sessionId: string, fence: number): Promise { + return this.start(sessionId, fence, true) } - private async start(sessionId: string, fence: number, recovering: boolean): Promise { + private async start(sessionId: string, fence: number, recovering: boolean): Promise { + this.teardown.signal.throwIfAborted() this.stop(sessionId) const record = this.input.store.getRecord(sessionId) const head = record?.providerHandleChain.at(-1) @@ -64,7 +69,7 @@ export class StructuredTuiTranscriptCatchup { !head || (head.handle.provider !== 'codex' && head.handle.provider !== 'claude') ) { - return + return this.teardown.signal } const agent = head.handle.provider const providerSessionId = agent === 'claude' ? head.handle.sessionId : head.handle.threadId @@ -73,17 +78,9 @@ export class StructuredTuiTranscriptCatchup { agent === 'claude' ? { claudeProjectsDir: join(record.accountHome.path, 'projects') } : { codexSessionsDirs: [join(record.accountHome.path, 'sessions')] } - const boundary = recovering - ? await readStructuredTuiTranscriptBoundary(journal.directory) - : null - const filePath = await resolveSessionFilePath(agent, providerSessionId, { - ...transcriptOptions, - ...(boundary?.filePath ? { transcriptPath: boundary.filePath } : {}) - }) - let initialReady: (() => void) | null = null - let baselineOffset = 0 - const ready = filePath ? new Promise((resolve) => (initialReady = resolve)) : null const state: CatchupState = { + controller: new AbortController(), + initialReady: null, active: false, fence, agent, @@ -95,20 +92,42 @@ export class StructuredTuiTranscriptCatchup { const receive = (messages: NativeChatMessage[]) => this.receive(sessionId, state, messages) this.states.set(sessionId, state) try { - state.subscription = await subscribeNativeChatTranscript({ + const signal = state.controller.signal + const boundary = recovering + ? await readStructuredTuiTranscriptBoundary(journal.directory) + : null + signal.throwIfAborted() + const filePath = await resolveSessionFilePath( agent, - sessionId: providerSessionId, - ...transcriptOptions, - ...(filePath ? { filePath, initialLimit: 0 } : {}), - onInitialSnapshot: (messages, _hasMore, beforeOffset) => { - baselineOffset = beforeOffset - receive(messages) - initialReady?.() - initialReady = null + providerSessionId, + { + ...transcriptOptions, + ...(boundary?.filePath ? { transcriptPath: boundary.filePath } : {}) }, - onAppend: receive - }) + signal + ) + signal.throwIfAborted() + let baselineOffset = 0 + const ready = filePath ? new Promise((resolve) => (state.initialReady = resolve)) : null + state.subscription = await subscribeNativeChatTranscript( + { + agent, + sessionId: providerSessionId, + ...transcriptOptions, + ...(filePath ? { filePath, initialLimit: 0 } : {}), + onInitialSnapshot: (messages, _hasMore, beforeOffset) => { + baselineOffset = beforeOffset + receive(messages) + state.initialReady?.() + state.initialReady = null + }, + onAppend: receive + }, + signal + ) + signal.throwIfAborted() await ready + signal.throwIfAborted() if (!recovering) { await writeStructuredTuiTranscriptBoundary(journal.directory, { providerSessionId, @@ -134,13 +153,21 @@ export class StructuredTuiTranscriptCatchup { if (!imported.ok) { throw new Error(imported.error) } + signal.throwIfAborted() this.input.reset(sessionId, fence) } + signal.throwIfAborted() + return signal } catch (error) { + const stopped = state.controller.signal.aborted if (this.states.get(sessionId) === state) { - this.states.delete(sessionId) + this.stop(sessionId) + } else { + state.subscription?.unsubscribe() + } + if (stopped) { + state.controller.signal.throwIfAborted() } - state.subscription?.unsubscribe() throw error } } @@ -197,10 +224,16 @@ export class StructuredTuiTranscriptCatchup { stop(sessionId: string): void { const state = this.states.get(sessionId) this.states.delete(sessionId) + state?.controller.abort(new StructuredTuiCatchupStoppedError()) + state?.initialReady?.() + if (state) { + state.initialReady = null + } state?.subscription?.unsubscribe() } stopAll(): void { + this.teardown.abort(new StructuredTuiCatchupStoppedError()) for (const sessionId of this.states.keys()) { this.stop(sessionId) }