From 59aaeae2bca61b2808aaaac56790827ec164f99c Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Wed, 16 Sep 2026 14:42:00 -0700 Subject: [PATCH] fix(native-chat): release session lane during commands --- ...uctured-agent-session-control-lane.test.ts | 65 --- .../structured-agent-session-control-lane.ts | 37 -- ...structured-agent-session-host-mutations.ts | 31 +- .../structured-agent-session-host.ts | 17 +- ...ctured-agent-session-mutation-admission.ts | 84 +++- .../structured-compaction-recovery.ts | 17 +- ...uctured-conversation-command-controller.ts | 193 +++++--- ...ructured-conversation-command-execution.ts | 300 ++++++++++++ ...ructured-conversation-command-lifecycle.ts | 53 ++ .../structured-conversation-command.test.ts | 231 ++++++--- .../structured-conversation-command.ts | 454 ++++++++---------- ...uctured-conversation-command-claim.test.ts | 234 +++------ .../structured-conversation-command-claim.ts | 201 +++----- ...e-structured-conversation-command.test.tsx | 48 +- .../use-structured-conversation-command.ts | 10 +- 15 files changed, 1057 insertions(+), 918 deletions(-) delete mode 100644 src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.test.ts delete mode 100644 src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.ts create mode 100644 src/main/native-chat/agent-session-wire/structured-conversation-command-execution.ts create mode 100644 src/main/native-chat/agent-session-wire/structured-conversation-command-lifecycle.ts diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.test.ts deleted file mode 100644 index 991134926f3..00000000000 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.test.ts +++ /dev/null @@ -1,65 +0,0 @@ -import { describe, expect, it } from 'vitest' -import type { AgentSessionRecord } from '../../../shared/agent-session-record' -import { structuredAgentSessionControlLaneFor } from './structured-agent-session-control-lane' - -const SESSION = 'session-1' - -function record(conversationCommand?: { - command?: 'clear' | 'compact' - phase: 'prepared' | 'committed' - state: 'unknown' | 'completed' -}): Pick { - return conversationCommand - ? { - conversationCommand: { - ...conversationCommand, - command: conversationCommand.command ?? 'compact', - operationId: 'op-1', - callerKey: 'desktop' - } - } - : { conversationCommand: undefined } -} - -describe('structuredAgentSessionControlLaneFor', () => { - it('keeps user controls in order on the main lane when nothing parks it', () => { - expect(structuredAgentSessionControlLaneFor(SESSION, record())).toBe(SESSION) - expect( - structuredAgentSessionControlLaneFor( - SESSION, - record({ phase: 'committed', state: 'completed' }) - ) - ).toBe(SESSION) - expect(structuredAgentSessionControlLaneFor(SESSION, null)).toBe(SESSION) - }) - - it('moves user controls off the main lane while a command awaits its terminal frame', () => { - const lane = structuredAgentSessionControlLaneFor( - SESSION, - record({ phase: 'prepared', state: 'unknown' }) - ) - - expect(lane).not.toBe(SESSION) - expect(lane).toContain(SESSION) - }) - - it('keeps clear serialized with close while the replacement is prepared', () => { - expect( - structuredAgentSessionControlLaneFor( - SESSION, - record({ command: 'clear', phase: 'prepared', state: 'unknown' }) - ) - ).toBe(SESSION) - }) - - it('uses the live owner decision before the durable record catches up', () => { - expect(structuredAgentSessionControlLaneFor(SESSION, record(), true)).not.toBe(SESSION) - expect( - structuredAgentSessionControlLaneFor( - SESSION, - record({ phase: 'prepared', state: 'unknown' }), - false - ) - ).toBe(SESSION) - }) -}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.ts deleted file mode 100644 index 7c1e3c531a8..00000000000 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-control-lane.ts +++ /dev/null @@ -1,37 +0,0 @@ -import type { AgentSessionRecord } from '../../../shared/agent-session-record' - -/** - * A session's mutations serialize on one lane keyed by session id. A conversation command holds that - * lane for as long as it waits on the provider's terminal frame, so the controls a user can always - * reach -- interrupt, cancel, close -- take a lane of their own while that wait is outstanding. - * Closing is what stops the provider child, which is what ends the wait; parking it behind the wait - * makes the escape depend on the thing it escapes. - */ -export function structuredAgentSessionControlLane(sessionId: string): string { - return `session-control:${sessionId}` -} - -/** - * Re-derived from the durable record, never timed: only compaction holds the main lane across a - * provider round trip. Clear stays serialized with close because it creates and commits a - * replacement before the source may be retired. - */ -export function structuredAgentSessionMainLaneParked( - record: Pick | null | undefined -): boolean { - const command = record?.conversationCommand - return ( - command?.command === 'compact' && command.phase === 'prepared' && command.state === 'unknown' - ) -} - -/** The lane a user control should run on: its own while a command parks the main one, else the main one. */ -export function structuredAgentSessionControlLaneFor( - sessionId: string, - record: Pick | null | undefined, - liveMainLaneParked?: boolean -): string { - return (liveMainLaneParked ?? structuredAgentSessionMainLaneParked(record)) - ? structuredAgentSessionControlLane(sessionId) - : sessionId -} diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts index 853eb395965..d9d43fd600e 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-mutations.ts @@ -19,7 +19,6 @@ import type { AgentSessionPromptResult, AgentSessionSendResult } from '../../../shared/agent-session-wire' -import { structuredAgentSessionControlLaneFor } from './structured-agent-session-control-lane' import { admitAndRunAgentSessionMutation } from './structured-agent-session-mutation-admission' import { cancelPlan, @@ -41,8 +40,7 @@ export type StructuredAgentSessionMutationContext = { flushStreamedEvents: (sessionId: string) => Promise requireSession: (sessionId: string) => StructuredAgentSessionHostSession serialize: (sessionId: string, task: () => Promise) => Promise - conversationCommandMainLaneParked: (sessionId: string) => boolean | undefined - requestConversationCommandControl: (sessionId: string, turnId?: string) => boolean | undefined + abandonConversationCommand: (sessionId: string, turnId?: string) => Promise now: () => number } @@ -117,28 +115,15 @@ export function cancelStructuredAgentSessionTurn( prompt?: { itemId: string; expectedRevision: number } } ): Promise> { - const liveMainLaneParked = context.conversationCommandMainLaneParked(params.envelope.sessionId) - // Interrupts must reach a provider while a command awaits its terminal frame. - const cancellationContext = { - ...context, - serialize: (sessionId: string, task: () => Promise) => - context.serialize( - structuredAgentSessionControlLaneFor( - sessionId, - context.deps.store.getRecord(sessionId), - liveMainLaneParked - ), - task - ) - } const plan = cancelPlan(params) - return mutate(cancellationContext, caller, params.envelope, { + return mutate(context, caller, params.envelope, { ...plan, - run: (ctx) => { - // Mutating the pending command before admission lets a stale or conflicting request cancel - // work even though the request itself is refused. - context.requestConversationCommandControl(params.envelope.sessionId, params.turnId) - return plan.run(ctx) + run: async (ctx) => { + const outcome = await plan.run(ctx) + if (outcome.ok) { + await context.abandonConversationCommand(params.envelope.sessionId, params.turnId) + } + return outcome } }) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts index 20d43fe003b..401bbeec9ae 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts @@ -9,7 +9,6 @@ import type { AgentSessionExecutionLocation } from '../../../shared/agent-sessio import type * as SessionWire from '../../../shared/agent-session-wire' import type { AgentSessionAttachParams } from './structured-agent-session-attach' import { AGENT_SESSION_NOT_ATTACHED } from './structured-agent-session-mutation-admission' -import { structuredAgentSessionControlLaneFor } from './structured-agent-session-control-lane' import { createRestartReconciler } from './structured-agent-session-restart-reconcile' import type { AgentSessionSubscribeInput } from './structured-agent-session-subscribers' import { StructuredAgentSessionTaskQueue } from './structured-agent-session-task-queue' @@ -192,14 +191,8 @@ export class StructuredAgentSessionHost { /** Releases a session's resources without ending the conversation: the record and journal stay * on disk, so the same session can be attached again. */ close(sessionId: string): Promise { - // Closing stops the provider child, which is what ends a command's wait for its terminal frame; - // it must not queue behind that wait. - const lane = structuredAgentSessionControlLaneFor( - sessionId, - this.deps.store.getRecord(sessionId), - this.conversationCommands.requestControl(sessionId) - ) - return this.serialize(lane, async () => { + return this.serialize(sessionId, async () => { + await this.conversationCommands.abandon(sessionId) await this.handoffs.closeRetainedTuiOwner(sessionId) await evictHeldStructuredAgentSession(this.lifetimeContext(), sessionId) this.clientDelivery.closeSession(sessionId) @@ -278,10 +271,8 @@ export class StructuredAgentSessionHost { flushStreamedEvents: this.flushStreamedEvents, requireSession: (sessionId) => this.requireSession(sessionId), serialize: (sessionId, task) => this.serialize(sessionId, task), - conversationCommandMainLaneParked: (sessionId) => - this.conversationCommands.mainLaneParked(sessionId), - requestConversationCommandControl: (sessionId, turnId) => - this.conversationCommands.requestControl(sessionId, turnId), + abandonConversationCommand: (sessionId, turnId) => + this.conversationCommands.abandon(sessionId, turnId), now: () => this.now() } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts index 34e8004cc81..6df4fbb6b31 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts @@ -36,12 +36,22 @@ export function refuseAgentSessionMutation(refusal: AgentSessionWireRefusal): { return { ok: false, refusal } } -export type AgentSessionMutationRequest = { +type AgentSessionMutationAdmissionPlan = Pick< + MutationPlan, + | 'method' + | 'fields' + | 'operationIdScope' + | 'replay' + | 'rerunWhenReplayMissing' + | 'recoverUnknownFromDurableState' +> + +export type AgentSessionMutationAdmissionRequest = { store: AgentSessionRecordStore adapter: StructuredAgentSessionAdapter callerKey: string envelope: AgentSessionMutationEnvelope - plan: MutationPlan + plan: AgentSessionMutationAdmissionPlan /** Journal of the attached session; absent when this host holds none. */ journal: AgentSessionJournal | undefined publish: (journal: AgentSessionJournal) => void @@ -49,12 +59,27 @@ export type AgentSessionMutationRequest = { now: () => number } -export async function admitAndRunAgentSessionMutation( - request: AgentSessionMutationRequest -): Promise> { +export type AgentSessionMutationRequest = Omit< + AgentSessionMutationAdmissionRequest, + 'plan' +> & { plan: MutationPlan } + +export type AgentSessionMutationAdmission = + | { decision: 'return'; result: AgentSessionMutationResult } + | { + decision: 'run' + context: AgentSessionTurnContext + operationCallerKey: string + } + +/** Admit a mutation without running it. Long-running host work uses this split so the session queue + * protects durable admission without remaining held across a provider round trip. */ +export async function admitAgentSessionMutationRequest( + request: AgentSessionMutationAdmissionRequest +): Promise> { const { envelope, plan, journal } = request if (!journal) { - return refuseAgentSessionMutation(AGENT_SESSION_NOT_ATTACHED) + return { decision: 'return', result: refuseAgentSessionMutation(AGENT_SESSION_NOT_ATTACHED) } } const hostFingerprint = computeAgentSessionPayloadFingerprint({ method: plan.method, @@ -63,7 +88,7 @@ export async function admitAndRunAgentSessionMutation( }) const conflict = agentSessionFingerprintConflict(envelope, hostFingerprint) if (conflict) { - return refuseAgentSessionMutation(conflict) + return { decision: 'return', result: refuseAgentSessionMutation(conflict) } } const admitted = await request.store.admitMutationOperation({ callerKey: request.callerKey, @@ -73,11 +98,11 @@ export async function admitAndRunAgentSessionMutation( ...(plan.operationIdScope ? { operationIdScope: plan.operationIdScope } : {}) }) if (!admitted) { - return refuseAgentSessionMutation(AGENT_SESSION_NOT_ATTACHED) + return { decision: 'return', result: refuseAgentSessionMutation(AGENT_SESSION_NOT_ATTACHED) } } const { admission, record } = admitted if (admission.decision === 'refused') { - return refuseAgentSessionMutation(admission.refusal) + return { decision: 'return', result: refuseAgentSessionMutation(admission.refusal) } } const fence = record.lease.runtimeFence @@ -91,15 +116,14 @@ export async function admitAndRunAgentSessionMutation( recoverUnknownFromDurableState: plan.recoverUnknownFromDurableState }) if (replay.decision === 'refuse') { - return refuseAgentSessionMutation(replay.refusal) + return { decision: 'return', result: refuseAgentSessionMutation(replay.refusal) } } if (replay.decision === 'replay') { - return { ok: true, replayed: true, fence, cursor: journal.cursor(), value: replay.value } + return { + decision: 'return', + result: { ok: true, replayed: true, fence, cursor: journal.cursor(), value: replay.value } + } } - // Nothing durable landed, so this id is about to run for the first time. A - // refused call leaves its ledger row behind, and replaying past the lease and - // the fence would let a resend act under an owner that has since changed — so - // a first run pays the full admission price either way. const rerun = admitAgentSessionMutation({ envelope, hostFingerprint, @@ -107,26 +131,42 @@ export async function admitAndRunAgentSessionMutation( lease: record.lease }) if (rerun.decision === 'refused') { - return refuseAgentSessionMutation(rerun.refusal) + return { decision: 'return', result: refuseAgentSessionMutation(rerun.refusal) } } } + return { decision: 'run', context, operationCallerKey: admission.row.callerKey } +} + +export async function admitAndRunAgentSessionMutation( + request: AgentSessionMutationRequest +): Promise> { + const admitted = await admitAgentSessionMutationRequest(request) + if (admitted.decision === 'return') { + return admitted.result + } const outcome = await runSettledAgentSessionMutation({ store: request.store, // A global send replay can cross caller identities. Settlement still owns // the durable row admitted by the original caller. - operationCallerKey: admission.row.callerKey, - envelope, - plan, - context + operationCallerKey: admitted.operationCallerKey, + envelope: request.envelope, + plan: request.plan, + context: admitted.context }) return outcome.ok - ? { ok: true, replayed: false, fence, cursor: journal.cursor(), value: outcome.value } + ? { + ok: true, + replayed: false, + fence: admitted.context.fence, + cursor: admitted.context.journal.cursor(), + value: outcome.value + } : refuseAgentSessionMutation(outcome.refusal) } function turnContext( - request: AgentSessionMutationRequest, + request: AgentSessionMutationAdmissionRequest, journal: AgentSessionJournal, fence: number ): AgentSessionTurnContext { diff --git a/src/main/native-chat/agent-session-wire/structured-compaction-recovery.ts b/src/main/native-chat/agent-session-wire/structured-compaction-recovery.ts index 2199da05392..73bd559f184 100644 --- a/src/main/native-chat/agent-session-wire/structured-compaction-recovery.ts +++ b/src/main/native-chat/agent-session-wire/structured-compaction-recovery.ts @@ -18,11 +18,6 @@ export async function recoverInterruptedCompaction( return } const error = 'Previous compaction completion could not be confirmed after session recovery.' - await journal.appendItem( - { provider: 'orca', clientMessageId: `compact:${command.operationId}` }, - { kind: 'status', text: error }, - { fence } - ) const recovered = { ...command, phase: 'committed' as const, state: 'unknown' as const, error } await store.setConversationCommand(sessionId, fence, recovered) await store.recordOperationOutcome({ @@ -30,4 +25,16 @@ export async function recoverInterruptedCompaction( operationId: command.operationId, outcome: { status: 'succeeded', sessionId, conversationCommand: recovered } }) + await journal.appendItem( + { provider: 'orca', clientMessageId: `compact:${command.operationId}` }, + { + kind: 'status', + text: error, + turnLifecycle: { + turnId: `compact:${command.operationId}`, + state: 'unverifiable' + } + }, + { fence } + ) } diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts b/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts index 25586016974..a8f082876a2 100644 --- a/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts +++ b/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts @@ -1,52 +1,31 @@ import { sendStructuredAgentSessionTurn } from './structured-agent-session-host-mutations' import { - runStructuredConversationCommand, + prepareStructuredConversationCommand, type ConversationCommandParams } from './structured-conversation-command' +import { + StructuredConversationCommandExecution, + type ConversationCommandResult, + type PendingConversationCommand +} from './structured-conversation-command-execution' import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' import type { StructuredAgentSessionCaller } from './structured-agent-session-host-types' import type { StructuredAgentSessionHost } from './structured-agent-session-host' -type PendingConversationCommand = { - key: string - command: ConversationCommandParams['command'] - operationId: string - count: number - providerCallStarted: boolean - cancelBeforeProvider: boolean -} - export class StructuredConversationCommandController { private readonly pending = new Map() + private readonly execution: StructuredConversationCommandExecution + constructor( private readonly context: () => StructuredAgentSessionMutationContext, - private readonly host: Pick - ) {} - - /** Read-only lane ownership for a control request that has not passed mutation admission yet. */ - mainLaneParked(sessionId: string): boolean | undefined { - const entry = this.pending.get(sessionId) - return entry ? entry.command === 'compact' : undefined - } - - /** Resolve the lane synchronously with admission. False means the queued command will yield - * before provider execution; true means the provider call must be stopped from the control lane. */ - requestControl(sessionId: string, turnId?: string): boolean | undefined { - const entry = this.pending.get(sessionId) - if (!entry) { - return undefined - } - if (entry.command === 'clear') { - return false - } - if (entry.providerCallStarted) { - return true - } - if (turnId === undefined || turnId === `compact:${entry.operationId}`) { - entry.cancelBeforeProvider = true - return false - } - return true + host: Pick + ) { + this.execution = new StructuredConversationCommandExecution(context, host, { + isCurrent: (entry) => this.isCurrent(entry), + finish: (entry, result) => this.finish(entry, result), + settleWaiter: (entry, result) => this.settleWaiter(entry, result), + report: (entry, error) => this.report(entry, error) + }) } send = ( @@ -63,47 +42,68 @@ export class StructuredConversationCommandController { }) : sendStructuredAgentSessionTurn(this.context(), caller, params) - run = (caller: StructuredAgentSessionCaller, params: ConversationCommandParams) => { - const key = JSON.stringify([caller.callerKey, params.envelope.clientOperationId]) - const pending = this.pending.get(params.envelope.sessionId) - if (pending && pending.key !== key) { - return Promise.resolve({ - ok: false as const, - refusal: { - code: 'agent_session_operation_invalid' as const, - message: 'Wait for the conversation operation to finish.' - } - }) + run = ( + caller: StructuredAgentSessionCaller, + params: ConversationCommandParams + ): Promise => { + const sessionId = params.envelope.sessionId + const key = JSON.stringify([ + caller.callerKey, + params.envelope.clientOperationId, + params.command + ]) + const pending = this.pending.get(sessionId) + if (pending) { + return pending.key === key + ? pending.promise + : Promise.resolve({ + ok: false, + refusal: { + code: 'agent_session_operation_invalid', + message: 'Wait for the conversation operation to finish.' + } + }) } - const entry = - pending ?? - ({ - key, - command: params.command, - operationId: params.envelope.clientOperationId, - count: 0, - providerCallStarted: false, - cancelBeforeProvider: false - } satisfies PendingConversationCommand) - entry.count++ - this.pending.set(params.envelope.sessionId, entry) - return runStructuredConversationCommand(this.context(), this.host, caller, params, { - isCancelled: () => entry.cancelBeforeProvider, - beginProviderCall: () => { - if (entry.cancelBeforeProvider) { + const waiter = Promise.withResolvers() + const entry: PendingConversationCommand = { + key, + command: params.command, + operationId: params.envelope.clientOperationId, + execution: null, + promise: waiter.promise, + resolve: waiter.resolve, + waiterSettled: false + } + this.pending.set(sessionId, entry) + void this.context() + .serialize(sessionId, async () => { + const prepared = await prepareStructuredConversationCommand(this.context(), caller, params) + if (prepared.decision === 'return') { + this.finish(entry, prepared.result) return false } - entry.providerCallStarted = true + entry.execution = prepared.execution return true - }, - endProviderCall: () => { - entry.providerCallStarted = false - } - }).finally(() => { - if (--entry.count === 0 && this.pending.get(params.envelope.sessionId) === entry) { - this.pending.delete(params.envelope.sessionId) - } - }) + }) + .then((execute) => { + if (execute) { + void this.execution.run(entry).catch((error) => this.report(entry, error)) + } + }) + .catch((error) => this.failAdmission(entry, error)) + return entry.promise + } + + /** Called while the session lane is held. It terminalizes the host token before teardown. */ + abandon = async (sessionId: string, turnId?: string): Promise => { + const entry = this.pending.get(sessionId) + if ( + !entry?.execution || + (turnId !== undefined && turnId !== `${entry.command}:${entry.operationId}`) + ) { + return + } + this.finish(entry, await this.execution.abandon(entry)) } replacements = () => { @@ -137,7 +137,6 @@ export class StructuredConversationCommandController { return records.flatMap((record) => { const target = destination(record.sessionId) const sessionId = target !== record.sessionId ? target : null - // Explicit history reveals remain readable; closed replacements stay closed. return sessionId && visible.has(sessionId) && !visible.has(record.sessionId) ? [ { @@ -150,4 +149,48 @@ export class StructuredConversationCommandController { : [] }) } + + private isCurrent(entry: PendingConversationCommand): boolean { + const sessionId = entry.execution?.turn.sessionId + return sessionId !== undefined && this.pending.get(sessionId) === entry + } + + private settleWaiter(entry: PendingConversationCommand, result: ConversationCommandResult): void { + if (!entry.waiterSettled) { + entry.waiterSettled = true + entry.resolve(result) + } + } + + private finish(entry: PendingConversationCommand, result: ConversationCommandResult): void { + const sessionId = entry.execution?.turn.sessionId + if (sessionId && this.pending.get(sessionId) === entry) { + this.pending.delete(sessionId) + } else { + for (const [candidate, pending] of this.pending) { + if (pending === entry) { + this.pending.delete(candidate) + break + } + } + } + this.settleWaiter(entry, result) + } + + private failAdmission(entry: PendingConversationCommand, error: unknown): void { + this.finish(entry, { + ok: false, + refusal: { + code: 'agent_session_operation_invalid', + message: error instanceof Error ? error.message : 'Conversation operation failed.' + } + }) + } + + private report(entry: PendingConversationCommand, error: unknown): void { + this.context().deps.onEventSinkError?.({ + sessionId: entry.execution?.turn.sessionId ?? 'unknown', + error + }) + } } diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-command-execution.ts b/src/main/native-chat/agent-session-wire/structured-conversation-command-execution.ts new file mode 100644 index 00000000000..057847cb7cf --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-conversation-command-execution.ts @@ -0,0 +1,300 @@ +import type { + AgentSessionConversationCommandRecord, + AgentSessionConversationCommandResult +} from '../../../shared/agent-session-conversation-command' +import type { AgentSessionMutationResult } from '../../../shared/agent-session-wire' +import { attachConversationClearReplacement } from './structured-conversation-clear-replacement' +import { + conversationCommandResult, + persistConversationCommandResult, + type ConversationCommandParams, + type PreparedConversationCommand +} from './structured-conversation-command' +import { + COMPACTION_UNCONFIRMED, + CONVERSATION_COMMAND_ABANDONED, + publishConversationCommandLifecycle +} from './structured-conversation-command-lifecycle' +import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' +import type { StructuredAgentSessionHost } from './structured-agent-session-host' + +export type ConversationCommandResult = + AgentSessionMutationResult + +export type PendingConversationCommand = { + key: string + command: ConversationCommandParams['command'] + operationId: string + execution: PreparedConversationCommand | null + promise: Promise + resolve: (result: ConversationCommandResult) => void + waiterSettled: boolean +} + +type ExecutionOwner = { + isCurrent: (entry: PendingConversationCommand) => boolean + finish: (entry: PendingConversationCommand, result: ConversationCommandResult) => void + settleWaiter: (entry: PendingConversationCommand, result: ConversationCommandResult) => void + report: (entry: PendingConversationCommand, error: unknown) => void +} + +export class StructuredConversationCommandExecution { + constructor( + private readonly context: () => StructuredAgentSessionMutationContext, + private readonly host: Pick< + StructuredAgentSessionHost, + 'attach' | 'close' | 'flushStreamedEvents' + >, + private readonly owner: ExecutionOwner + ) {} + + async run(entry: PendingConversationCommand): Promise { + const execution = entry.execution + if (!execution || !this.owner.isCurrent(entry)) { + return + } + try { + await this.publishLifecycle(entry, execution.prepared, 'running') + await this.host.flushStreamedEvents(execution.turn.sessionId) + if (!this.owner.isCurrent(entry)) { + return + } + await (entry.command === 'clear' + ? this.executeClear(entry, execution) + : this.executeCompact(entry, execution)) + } catch (error) { + await this.markUnknown(entry, error) + } + } + + async abandon(entry: PendingConversationCommand): Promise { + const execution = entry.execution + if (!execution) { + throw new Error('Conversation command was not prepared.') + } + const value: AgentSessionConversationCommandRecord = { + ...execution.prepared, + phase: 'committed', + state: 'unknown', + error: CONVERSATION_COMMAND_ABANDONED, + ...(entry.command === 'clear' ? { replacementSessionId: undefined } : {}) + } + try { + await persistConversationCommandResult(this.context(), execution, value) + await this.publishLifecycle(entry, value, 'interrupted') + } catch (error) { + this.owner.report(entry, error) + } + return conversationCommandResult(execution, value) + } + + private async executeCompact( + entry: PendingConversationCommand, + execution: PreparedConversationCommand + ): Promise { + const compact = execution.turn.adapter.compact + if (!compact) { + await this.complete(entry, 'Compaction is unavailable for this provider.') + return + } + const result = await compact({ + turnId: `compact:${entry.operationId}`, + sessionId: execution.turn.sessionId, + fence: execution.turn.fence, + onLateResult: (late) => this.complete(entry, late.error) + }) + await this.host.flushStreamedEvents(execution.turn.sessionId) + await this.complete(entry, result.error) + } + + private async executeClear( + entry: PendingConversationCommand, + execution: PreparedConversationCommand + ): Promise { + let effectiveOptions = execution.source.options + if (!execution.supersededOperation) { + try { + const options = await execution.turn.adapter.readOptions?.({ + sessionId: execution.turn.sessionId, + fence: execution.turn.fence + }) + effectiveOptions = { + ...effectiveOptions, + ...(options + ? { + model: options.current.model, + ...(options.current.effort ? { effort: options.current.effort } : {}) + } + : {}) + } + } catch { + await this.complete( + entry, + 'Could not read the current session configuration. Try again when the provider is connected.', + true + ) + return + } + } + if (!this.owner.isCurrent(entry)) { + return + } + if (effectiveOptions) { + await this.context().serialize(execution.turn.sessionId, async () => { + if (this.canSettle(entry, execution)) { + await execution.turn.persistOptions(effectiveOptions) + } + }) + } + if (!this.owner.isCurrent(entry)) { + return + } + const replacementSessionId = execution.prepared.replacementSessionId! + let attachError: string | null + try { + attachError = await attachConversationClearReplacement({ + host: this.host, + store: this.context().deps.store, + sourceSessionId: execution.turn.sessionId, + replacementSessionId, + callerKey: execution.supersededOperation?.callerKey ?? execution.prepared.callerKey, + operationId: execution.supersededOperation?.operationId ?? execution.prepared.operationId, + source: { ...execution.source, options: effectiveOptions } + }) + } catch (error) { + const replacement = this.context().deps.store.getRecord(replacementSessionId) + if (!replacement || replacement.lease.claimStatus === 'released') { + throw error + } + attachError = null + } + if (!this.owner.isCurrent(entry)) { + if (!attachError) { + await this.host + .close(replacementSessionId) + .catch((error) => this.owner.report(entry, error)) + } + return + } + await this.complete(entry, attachError ?? undefined, Boolean(attachError)) + } + + private async complete( + entry: PendingConversationCommand, + error?: string, + discardReplacement = false + ): Promise { + const execution = entry.execution + if (!execution) { + return + } + await this.context().serialize(execution.turn.sessionId, async () => { + if (!this.canSettle(entry, execution)) { + return + } + const value: AgentSessionConversationCommandRecord = { + ...execution.prepared, + phase: 'committed', + state: 'completed', + ...(error ? { error: error.slice(0, 4096) } : {}), + ...(discardReplacement ? { replacementSessionId: undefined } : {}) + } + try { + await persistConversationCommandResult(this.context(), execution, value) + } catch (cause) { + await this.markUnknownInLane(entry, cause, true) + return + } + try { + await this.publishLifecycle(entry, value, 'completed') + } catch (cause) { + this.owner.report(entry, cause) + } + this.owner.finish(entry, conversationCommandResult(execution, value)) + }) + } + + private async markUnknown(entry: PendingConversationCommand, cause: unknown): Promise { + const execution = entry.execution + if (!execution) { + return + } + await this.context().serialize(execution.turn.sessionId, () => + this.markUnknownInLane(entry, cause) + ) + } + + private async markUnknownInLane( + entry: PendingConversationCommand, + cause: unknown, + retire = false + ): Promise { + const execution = entry.execution + if (!execution || !this.ownsExecution(entry, execution)) { + return + } + const error = cause instanceof Error ? cause.message : COMPACTION_UNCONFIRMED + try { + await this.context().deps.store.recordOperationOutcome({ + callerKey: execution.operationCallerKey, + operationId: execution.prepared.operationId, + outcome: { status: 'unknown' } + }) + await this.publishLifecycle(entry, { ...execution.prepared, error }, 'unverifiable') + } catch (persistError) { + this.owner.report(entry, persistError) + } + const result = conversationCommandResult(execution, { + command: entry.command, + state: 'unknown', + error: entry.command === 'compact' ? COMPACTION_UNCONFIRMED : error + }) + if (retire) { + this.owner.finish(entry, result) + } else { + this.owner.settleWaiter(entry, result) + } + } + + private publishLifecycle( + entry: PendingConversationCommand, + value: AgentSessionConversationCommandResult, + state: Parameters[0]['state'] + ): Promise { + return publishConversationCommandLifecycle({ + context: this.context, + command: entry.command, + operationId: entry.operationId, + execution: entry.execution!, + value, + state + }) + } + + private canSettle( + entry: PendingConversationCommand, + execution: PreparedConversationCommand + ): boolean { + const command = this.context().deps.store.getRecord( + execution.turn.sessionId + )?.conversationCommand + return this.ownsExecution(entry, execution) && command?.phase === 'prepared' + } + + private ownsExecution( + entry: PendingConversationCommand, + execution: PreparedConversationCommand + ): boolean { + const sessionId = execution.turn.sessionId + const session = this.context().sessions.get(sessionId) + const command = this.context().deps.store.getRecord(sessionId)?.conversationCommand + return ( + this.owner.isCurrent(entry) && + session?.journal === execution.turn.journal && + session.fence === execution.turn.fence && + command?.runtimeFence === execution.turn.fence && + command.operationId === execution.prepared.operationId && + command.callerKey === execution.prepared.callerKey + ) + } +} diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-command-lifecycle.ts b/src/main/native-chat/agent-session-wire/structured-conversation-command-lifecycle.ts new file mode 100644 index 00000000000..b26007394ef --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-conversation-command-lifecycle.ts @@ -0,0 +1,53 @@ +import type { AgentSessionConversationCommandResult } from '../../../shared/agent-session-conversation-command' +import type { + ConversationCommandParams, + PreparedConversationCommand +} from './structured-conversation-command' +import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' + +export const COMPACTION_UNCONFIRMED = 'Compaction completion is unconfirmed.' +export const CONVERSATION_COMMAND_ABANDONED = + 'Conversation operation was interrupted before completion.' + +export type ConversationCommandLifecycleState = + | 'running' + | 'completed' + | 'interrupted' + | 'unverifiable' + +export async function publishConversationCommandLifecycle(input: { + context: () => StructuredAgentSessionMutationContext + command: ConversationCommandParams['command'] + operationId: string + execution: PreparedConversationCommand + value: AgentSessionConversationCommandResult + state: ConversationCommandLifecycleState +}): Promise { + const { command, context, execution, operationId, state, value } = input + const turnId = `${command}:${operationId}` + const text = + state === 'running' + ? command === 'compact' + ? 'Compacting conversation…' + : 'Clearing conversation…' + : state === 'unverifiable' + ? command === 'compact' + ? COMPACTION_UNCONFIRMED + : 'Conversation clear completion is unconfirmed.' + : (value.error ?? + (command === 'compact' ? 'Conversation compacted.' : 'Conversation cleared.')) + await execution.turn.journal.appendItem( + { provider: 'orca', clientMessageId: turnId }, + { + kind: 'status', + text, + turnLifecycle: { + turnId, + state, + ...(state === 'running' ? { startedAt: context().now() } : { completedAt: context().now() }) + } + }, + { fence: execution.turn.fence } + ) + execution.turn.publish() +} diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts b/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts index edfe93f93c5..2a1a11d5e47 100644 --- a/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts @@ -138,19 +138,22 @@ describe('host conversation commands', () => { it('keeps an unknown compaction from being executed again', async () => { compact.mockRejectedValue(new Error('connection lost')) const params = commandParams('compact') - await expect(host.conversationCommand(caller, params)).rejects.toThrow('connection lost') expect(await host.conversationCommand(caller, params)).toMatchObject({ - ok: false, - refusal: { code: 'agent_session_operation_unknown' } + ok: true, + value: { state: 'unknown' } + }) + expect(await host.conversationCommand(caller, params)).toMatchObject({ + ok: true, + value: { state: 'unknown' } }) expect(compact).toHaveBeenCalledTimes(1) const status = host .history({ sessionId: HOST_TEST_SESSION, direction: 'tail' }) .page.items.find((item) => item.body.kind === 'status') expect(status?.body).toMatchObject({ - text: 'Compaction completion is unconfirmed.' + text: 'Compaction completion is unconfirmed.', + turnLifecycle: { state: 'unverifiable' } }) - expect(status?.body).not.toHaveProperty('turnLifecycle') }) /** The replacement seeds from what the provider reports now, not from what the @@ -266,8 +269,8 @@ describe('host conversation commands', () => { }) expect(cancel).toMatchObject({ ok: true, value: { cancelled: true } }) expect(adapter.cancelTurn).toHaveBeenCalled() + await expect(running).resolves.toMatchObject({ ok: true, value: { state: 'unknown' } }) finish({}) - await running }) it('does not let a stale cancellation stop an admitted compaction', async () => { @@ -297,14 +300,8 @@ describe('host conversation commands', () => { expect(compact).toHaveBeenCalledTimes(1) }) - it('closes a session while a compaction still awaits its terminal frame', async () => { - let finish!: (value: {}) => void - compact.mockImplementation( - () => - new Promise((resolve) => { - finish = resolve - }) - ) + it('closes and reattaches while a provider compaction never settles', async () => { + compact.mockImplementation(() => new Promise(() => {})) const running = host.conversationCommand(caller, commandParams('compact')) await vi.waitFor(() => expect(compact).toHaveBeenCalled()) @@ -319,25 +316,81 @@ describe('host conversation commands', () => { expect(outcome).toBe('closed') expect(host.hasSession(HOST_TEST_SESSION)).toBe(false) + await expect(running).resolves.toMatchObject({ ok: true, value: { state: 'unknown' } }) - finish({}) - await running.catch(() => undefined) + const fence = store.getRecord(HOST_TEST_SESSION)!.lease.runtimeFence + expect(await host.attach(caller, hostTestAttachParams(fence))).toMatchObject({ ok: true }) + compact.mockResolvedValue({}) + await expect(host.conversationCommand(caller, commandParams('compact'))).resolves.toMatchObject( + { + ok: true, + value: { state: 'completed' } + } + ) }) - it('cancels a compaction admitted immediately before close instead of parking close', async () => { + it('settles a compaction admitted immediately before close', async () => { const running = host.conversationCommand(caller, commandParams('compact')) const closed = host.close(HOST_TEST_SESSION) await expect(closed).resolves.toBeUndefined() await expect(running).resolves.toMatchObject({ - ok: false, - refusal: { message: 'Conversation operation was cancelled before provider execution.' } + ok: true, + value: { state: 'unknown' } }) - expect(compact).not.toHaveBeenCalled() expect(host.hasSession(HOST_TEST_SESSION)).toBe(false) }) - it('finishes clear before a concurrent close retires the source session', async () => { + it('closes and reattaches when the pre-provider event drain never settles', async () => { + const flush = vi + .spyOn(host, 'flushStreamedEvents') + .mockImplementationOnce(() => new Promise(() => {})) + const running = host.conversationCommand(caller, commandParams('compact')) + await vi.waitFor(() => expect(flush).toHaveBeenCalled()) + + await expect(host.close(HOST_TEST_SESSION)).resolves.toBeUndefined() + await expect(running).resolves.toMatchObject({ ok: true, value: { state: 'unknown' } }) + expect(compact).not.toHaveBeenCalled() + + flush.mockRestore() + const fence = store.getRecord(HOST_TEST_SESSION)!.lease.runtimeFence + expect(await host.attach(caller, hostTestAttachParams(fence))).toMatchObject({ ok: true }) + await expect(host.conversationCommand(caller, commandParams('compact'))).resolves.toMatchObject( + { + ok: true, + value: { state: 'completed' } + } + ) + }) + + it('interrupts a command whose pre-provider event drain never settles', async () => { + const flush = vi + .spyOn(host, 'flushStreamedEvents') + .mockImplementationOnce(() => new Promise(() => {})) + const params = commandParams('compact') + const running = host.conversationCommand(caller, params) + await vi.waitFor(() => expect(flush).toHaveBeenCalled()) + const turnId = `compact:${params.envelope.clientOperationId}` + + await expect( + host.cancel(caller, { + turnId, + envelope: { + ...params.envelope, + clientOperationId: hostTestOperationId(), + payloadFingerprint: computeAgentSessionPayloadFingerprint({ + method: 'agentSession.cancel', + sessionId: HOST_TEST_SESSION, + fields: { turnId } + }) + } + }) + ).resolves.toMatchObject({ ok: true }) + await expect(running).resolves.toMatchObject({ ok: true, value: { state: 'unknown' } }) + expect(compact).not.toHaveBeenCalled() + }) + + it('abandons a blocked replacement attach and cleans up its stale late completion', async () => { let finishReplacement!: () => void const originalAcquire = vi.mocked(adapter.acquire).getMockImplementation()! vi.mocked(adapter.acquire).mockImplementation(async (input) => { @@ -350,43 +403,85 @@ describe('host conversation commands', () => { }) const running = host.conversationCommand(caller, commandParams('clear')) await vi.waitFor(() => expect(finishReplacement).toBeTypeOf('function')) - let closed = false - const close = host.close(HOST_TEST_SESSION).then(() => { - closed = true - }) - await Promise.resolve() - expect(closed).toBe(false) + const replacementSessionId = + store.getRecord(HOST_TEST_SESSION)?.conversationCommand?.replacementSessionId + const close = host.close(HOST_TEST_SESSION) + await expect(close).resolves.toBeUndefined() + await expect(running).resolves.toMatchObject({ ok: true, value: { state: 'unknown' } }) finishReplacement() - const result = await running - expect(result).toMatchObject({ ok: true, value: { state: 'completed' } }) - await close + await vi.waitFor(() => expect(host.hasSession(replacementSessionId!)).toBe(false)) expect(store.getRecord(HOST_TEST_SESSION)?.conversationCommand).toMatchObject({ command: 'clear', phase: 'committed', - state: 'completed' + state: 'unknown', + replacementSessionId: undefined }) }) - it('reconstructs a committed replacement after the ledger settlement is lost', async () => { + it('does not publish terminal success before its durable command commit', async () => { + const persist = store.setConversationCommand.bind(store) + let failed = false + vi.spyOn(store, 'setConversationCommand').mockImplementation(async (...input) => { + if (!failed && input[2].phase === 'committed') { + failed = true + throw new Error('disk full') + } + return persist(...input) + }) + const params = commandParams('compact') + expect(await host.conversationCommand(caller, params)).toMatchObject({ + ok: true, + value: { state: 'unknown' } + }) + const status = host + .history({ sessionId: HOST_TEST_SESSION, direction: 'tail' }) + .page.items.find((item) => item.body.kind === 'status') + expect(status?.body).toMatchObject({ turnLifecycle: { state: 'unverifiable' } }) + expect(status?.body).not.toMatchObject({ turnLifecycle: { state: 'completed' } }) + expect(store.getRecord(HOST_TEST_SESSION)?.conversationCommand).toMatchObject({ + phase: 'prepared', + state: 'unknown' + }) + }) + + it('does not publish terminal success before its durable operation outcome', async () => { + const params = commandParams('compact') const persist = store.recordOperationOutcome.bind(store) + let failed = false vi.spyOn(store, 'recordOperationOutcome').mockImplementation(async (input) => { - if (input.outcome.status === 'succeeded' && input.outcome.conversationCommand) { - throw new Error('crash') + if ( + !failed && + input.operationId === params.envelope.clientOperationId && + input.outcome.status === 'succeeded' + ) { + failed = true + throw new Error('ledger write failed') } return persist(input) }) - const params = commandParams('clear') - await expect(host.conversationCommand(caller, params)).rejects.toThrow('crash') + + expect(await host.conversationCommand(caller, params)).toMatchObject({ + ok: true, + value: { state: 'unknown' } + }) + const status = host + .history({ sessionId: HOST_TEST_SESSION, direction: 'tail' }) + .page.items.find((item) => item.body.kind === 'status') + expect(status?.body).toMatchObject({ turnLifecycle: { state: 'unverifiable' } }) + expect(status?.body).not.toMatchObject({ turnLifecycle: { state: 'completed' } }) + expect(store.getRecord(HOST_TEST_SESSION)?.conversationCommand).toMatchObject({ + phase: 'committed', + state: 'completed' + }) expect(await host.conversationCommand(caller, params)).toMatchObject({ ok: true, replayed: true, value: { state: 'completed' } }) - expect(acquisitions).toBe(2) }) - it('adopts an interrupted clear under a fresh operation after verified reacquisition', async () => { + it('settles clear from the durable replacement when its attach reply is lost', async () => { const originalAttach = host.attach.bind(host) vi.spyOn(host, 'attach').mockImplementationOnce(async (...args) => { const attached = await originalAttach(...args) @@ -396,29 +491,18 @@ describe('host conversation commands', () => { throw new Error('response lost after replacement attach') }) const interrupted = commandParams('clear') - await expect(host.conversationCommand(caller, interrupted)).rejects.toThrow('response lost') + const result = await host.conversationCommand(caller, interrupted) + expect(result).toMatchObject({ + ok: true, + value: { state: 'completed' } + }) const replacementSessionId = store.getRecord(HOST_TEST_SESSION)?.conversationCommand?.replacementSessionId expect(store.getRecord(HOST_TEST_SESSION)?.conversationCommand).toMatchObject({ - phase: 'prepared', + phase: 'committed', operationId: interrupted.envelope.clientOperationId, replacementSessionId }) - - await host.close(HOST_TEST_SESSION) - const fence = store.getRecord(HOST_TEST_SESSION)!.lease.runtimeFence - expect(await host.attach(caller, hostTestAttachParams(fence))).toMatchObject({ ok: true }) - const recovery = commandParams('clear') - const recovered = await host.conversationCommand({ callerKey: 'mobile' }, recovery) - - expect(recovered).toMatchObject({ - ok: true, - value: { - phase: 'committed', - operationId: recovery.envelope.clientOperationId, - replacementSessionId - } - }) expect(await host.conversationCommand(caller, interrupted)).toMatchObject({ ok: true, replayed: true, @@ -429,7 +513,10 @@ describe('host conversation commands', () => { it('repairs an unknown receipt when the provider completes late', async () => { compact.mockRejectedValue(new Error('connection lost')) const params = commandParams('compact') - await expect(host.conversationCommand(caller, params)).rejects.toThrow() + await expect(host.conversationCommand(caller, params)).resolves.toMatchObject({ + ok: true, + value: { state: 'unknown' } + }) await compact.mock.calls[0]![0].onLateResult?.({}) expect(await host.conversationCommand(caller, params)).toMatchObject({ ok: true, @@ -439,6 +526,33 @@ describe('host conversation commands', () => { expect(compact).toHaveBeenCalledTimes(1) }) + it('ignores a stale late completion after close and a new generation', async () => { + compact.mockRejectedValueOnce(new Error('connection lost')) + const oldParams = commandParams('compact') + await expect(host.conversationCommand(caller, oldParams)).resolves.toMatchObject({ + ok: true, + value: { state: 'unknown' } + }) + const late = compact.mock.calls[0]![0].onLateResult! + + await host.close(HOST_TEST_SESSION) + const fence = store.getRecord(HOST_TEST_SESSION)!.lease.runtimeFence + expect(await host.attach(caller, hostTestAttachParams(fence))).toMatchObject({ ok: true }) + compact.mockResolvedValue({}) + const current = commandParams('compact') + await expect(host.conversationCommand(caller, current)).resolves.toMatchObject({ + ok: true, + value: { state: 'completed' } + }) + + await late({ error: 'stale failure' }) + expect(store.getRecord(HOST_TEST_SESSION)?.conversationCommand).toMatchObject({ + operationId: current.envelope.clientOperationId, + state: 'completed' + }) + expect(store.getRecord(HOST_TEST_SESSION)?.conversationCommand?.error).toBeUndefined() + }) + it('keeps explicitly revealed history and closed replacement tabs out of automatic restoration', async () => { const result = await host.conversationCommand(caller, commandParams('clear')) if (!result.ok) { @@ -454,7 +568,10 @@ describe('host conversation commands', () => { it('keeps the old compact outcome unknown but restores usability after verified reacquisition', async () => { compact.mockRejectedValue(new Error('lost response')) const params = commandParams('compact') - await expect(host.conversationCommand(caller, params)).rejects.toThrow() + await expect(host.conversationCommand(caller, params)).resolves.toMatchObject({ + ok: true, + value: { state: 'unknown' } + }) await host.close(HOST_TEST_SESSION) const fence = store.getRecord(HOST_TEST_SESSION)!.lease.runtimeFence expect(await host.attach(caller, hostTestAttachParams(fence))).toMatchObject({ ok: true }) diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-command.ts b/src/main/native-chat/agent-session-wire/structured-conversation-command.ts index f071b0c7cf9..f96c7829613 100644 --- a/src/main/native-chat/agent-session-wire/structured-conversation-command.ts +++ b/src/main/native-chat/agent-session-wire/structured-conversation-command.ts @@ -1,23 +1,30 @@ import { createHash } from 'node:crypto' import type { AgentSessionConversationCommand, + AgentSessionConversationCommandRecord, AgentSessionConversationCommandResult } from '../../../shared/agent-session-conversation-command' +import type { AgentSessionRecord } from '../../../shared/agent-session-record' import type { AgentSessionMutationEnvelope, - AgentSessionMutationResult + AgentSessionMutationResult, + AgentSessionWireRefusal } from '../../../shared/agent-session-wire' -import { admitAndRunAgentSessionMutation } from './structured-agent-session-mutation-admission' +import { + AGENT_SESSION_NOT_ATTACHED, + admitAgentSessionMutationRequest, + type AgentSessionMutationAdmissionRequest +} from './structured-agent-session-mutation-admission' import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations' import type { StructuredAgentSessionCaller } from './structured-agent-session-host-types' -import type { StructuredAgentSessionHost } from './structured-agent-session-host' +import type { AgentSessionTurnContext } from './structured-agent-session-turns' import { conversationCommandBlocked } from './structured-conversation-command-admission' -import { attachConversationClearReplacement } from './structured-conversation-clear-replacement' export type ConversationCommandParams = { envelope: AgentSessionMutationEnvelope command: AgentSessionConversationCommand } + export type ConversationReplacement = { sourceSessionId: string sessionId: string @@ -25,276 +32,199 @@ export type ConversationReplacement = { agent: 'claude' | 'codex' } -type ConversationCommandControl = { - isCancelled: () => boolean - beginProviderCall: () => boolean - endProviderCall: () => void +export type PreparedConversationCommand = { + turn: AgentSessionTurnContext + prepared: AgentSessionConversationCommandRecord + source: AgentSessionRecord + operationCallerKey: string + supersededOperation: AgentSessionConversationCommandRecord | null } -const CANCELLED_BEFORE_PROVIDER = 'Conversation operation was cancelled before provider execution.' +export type ConversationCommandPreparation = + | { + decision: 'return' + result: AgentSessionMutationResult + } + | { decision: 'execute'; execution: PreparedConversationCommand } -function cancelledBeforeProvider() { +function matchingCommand( + context: StructuredAgentSessionMutationContext, + callerKey: string, + params: ConversationCommandParams +): AgentSessionConversationCommandRecord | null { + const command = context.deps.store.getRecord(params.envelope.sessionId)?.conversationCommand + return command?.operationId === params.envelope.clientOperationId && + command.callerKey === callerKey + ? command + : null +} + +function mutationRequest( + context: StructuredAgentSessionMutationContext, + caller: StructuredAgentSessionCaller, + params: ConversationCommandParams +): AgentSessionMutationAdmissionRequest { + const { command, envelope } = params + const { sessionId } = envelope return { - ok: false as const, - refusal: { - code: 'agent_session_operation_invalid' as const, - message: CANCELLED_BEFORE_PROVIDER + store: context.deps.store, + adapter: context.deps.adapter, + callerKey: caller.callerKey, + envelope, + journal: context.sessions.get(sessionId)?.journal, + publish: (journal) => context.publish(sessionId, journal), + flushStreamedEvents: context.flushStreamedEvents, + now: context.now, + plan: { + method: 'agentSession.conversationCommand', + fields: { command }, + recoverUnknownFromDurableState: true, + replay: (_ctx, outcome) => { + if (outcome.status === 'succeeded' && outcome.conversationCommand) { + return outcome.conversationCommand + } + const prior = matchingCommand(context, caller.callerKey, params) + if (prior?.phase === 'committed') { + return prior + } + if (command === 'compact' && prior && outcome.status !== 'unknown') { + return { + command, + state: 'unknown', + error: 'Compaction completion is unconfirmed; it was not run again.' + } + } + return outcome.status === 'succeeded' && command === 'compact' + ? { command, state: 'completed' } + : null + }, + rerunWhenReplayMissing: () => + command === 'clear' && + matchingCommand(context, caller.callerKey, params)?.phase === 'prepared' } } } -export function runStructuredConversationCommand( +async function recordRefusal( + context: StructuredAgentSessionMutationContext, + operationCallerKey: string, + operationId: string, + refusal: AgentSessionWireRefusal +): Promise { + await context.deps.store.recordOperationOutcome({ + callerKey: operationCallerKey, + operationId, + outcome: { status: 'failed', code: refusal.code, message: refusal.message } + }) +} + +/** Admit and durably prepare on the session lane. Provider work starts only after this returns. */ +export async function prepareStructuredConversationCommand( context: StructuredAgentSessionMutationContext, - host: Pick, caller: StructuredAgentSessionCaller, - params: ConversationCommandParams, - control?: ConversationCommandControl -): Promise> { - const { envelope, command } = params + params: ConversationCommandParams +): Promise { + const admitted = await admitAgentSessionMutationRequest(mutationRequest(context, caller, params)) + if (admitted.decision === 'return') { + return { decision: 'return', result: admitted.result } + } + const { command, envelope } = params const { sessionId, clientOperationId } = envelope const store = context.deps.store - const matching = () => { - const record = store.getRecord(sessionId)?.conversationCommand - return record?.operationId === clientOperationId && record.callerKey === caller.callerKey - ? record - : null + const source = store.getRecord(sessionId) + if (!source) { + const refusal = AGENT_SESSION_NOT_ATTACHED + await recordRefusal(context, admitted.operationCallerKey, clientOperationId, refusal) + return { decision: 'return', result: { ok: false, refusal } } } - return context.serialize(sessionId, () => - admitAndRunAgentSessionMutation({ - store, - adapter: context.deps.adapter, - callerKey: caller.callerKey, - envelope, - journal: context.sessions.get(sessionId)?.journal, - publish: (journal) => context.publish(sessionId, journal), - flushStreamedEvents: context.flushStreamedEvents, - now: context.now, - plan: { - method: 'agentSession.conversationCommand', - fields: { command }, - recoverUnknownFromDurableState: true, - settledOutcome: (value) => ({ status: 'succeeded', sessionId, conversationCommand: value }), - replay: (_ctx, outcome) => { - if (outcome.status === 'succeeded' && outcome.conversationCommand) { - return outcome.conversationCommand - } - const prior = matching() - if (prior?.phase === 'committed') { - return prior - } - if (command === 'compact' && prior && outcome.status !== 'unknown') { - return { - command, - state: 'unknown', - error: 'Compaction completion is unconfirmed; it was not run again.' - } - } - return outcome.status === 'succeeded' && command === 'compact' - ? { command, state: 'completed' } - : null - }, - rerunWhenReplayMissing: () => command === 'clear' && matching()?.phase === 'prepared', - run: async (ctx) => { - await host.flushStreamedEvents(sessionId) - if (control?.isCancelled()) { - return cancelledBeforeProvider() - } - const record = store.getRecord(sessionId)! - const interruptedClear = record.conversationCommand - const prior = - matching() ?? - (command === 'clear' && - interruptedClear?.command === 'clear' && - interruptedClear.phase === 'prepared' && - interruptedClear.runtimeFence !== ctx.fence - ? interruptedClear - : null) - const supersededOperation = - prior && prior.operationId !== clientOperationId ? prior : null - const settleSupersededOperation = async ( - value: AgentSessionConversationCommandResult - ): Promise => { - if (!supersededOperation) { - return - } - await store.recordOperationOutcome({ - callerKey: supersededOperation.callerKey, - operationId: supersededOperation.operationId, - outcome: { status: 'succeeded', sessionId, conversationCommand: value } - }) - } - const blocked = - prior?.phase === 'prepared' && command === 'clear' - ? null - : conversationCommandBlocked(ctx, record) - if (blocked) { - return { - ok: false, - refusal: { code: 'agent_session_operation_invalid', message: blocked } - } - } - const replacementSessionId = - command === 'clear' - ? (prior?.replacementSessionId ?? - `clear-${createHash('sha256') - .update(JSON.stringify([sessionId, caller.callerKey, clientOperationId])) - .digest('hex') - .slice(0, 40)}`) - : undefined - const prepared = { - command, - runtimeFence: ctx.fence, - operationId: clientOperationId, - callerKey: caller.callerKey, - phase: 'prepared' as const, - state: 'unknown' as const, - ...(replacementSessionId ? { replacementSessionId } : {}) - } - let effectiveOptions = record.options - if (command === 'clear' && !prior) { - try { - const options = await ctx.adapter.readOptions?.({ sessionId, fence: ctx.fence }) - effectiveOptions = { - ...record.options, - ...(options - ? { - model: options.current.model, - ...(options.current.effort ? { effort: options.current.effort } : {}) - } - : {}) - } - } catch { - return { - ok: false, - refusal: { - code: 'agent_session_operation_invalid', - message: - 'Could not read the current session configuration. Try again when the provider is connected.' - } - } - } - } - if (effectiveOptions && command === 'clear') { - await ctx.persistOptions(effectiveOptions) - } - if (control?.isCancelled()) { - return cancelledBeforeProvider() - } - await store.setConversationCommand(sessionId, ctx.fence, prepared) - let error: string | undefined - if (command === 'clear' && replacementSessionId) { - const attachError = await attachConversationClearReplacement({ - host, - store, - sourceSessionId: sessionId, - replacementSessionId, - callerKey: prior?.callerKey ?? caller.callerKey, - operationId: prior?.operationId ?? clientOperationId, - source: { ...record, options: effectiveOptions } - }) - if (attachError) { - const failed = { - ...prepared, - replacementSessionId: undefined, - phase: 'committed' as const, - state: 'completed' as const, - error: attachError.slice(0, 4096) - } - await store.setConversationCommand(sessionId, ctx.fence, failed) - await settleSupersededOperation(failed) - return { ok: true, value: failed } - } - } else { - if (!ctx.adapter.compact) { - throw new Error('Compaction is unavailable for this provider.') - } - const identity = { - provider: 'orca' as const, - clientMessageId: `compact:${clientOperationId}` - } - await ctx.journal.appendItem( - identity, - { - kind: 'status', - text: 'Compacting conversation…', - turnLifecycle: { turnId: `compact:${clientOperationId}`, state: 'running' } - }, - { fence: ctx.fence } - ) - ctx.publish() - if (control && !control.beginProviderCall()) { - error = CANCELLED_BEFORE_PROVIDER - } else { - try { - error = ( - await ctx.adapter.compact({ - turnId: `compact:${clientOperationId}`, - sessionId, - fence: ctx.fence, - onLateResult: (result) => - context.serialize(sessionId, async () => { - if ( - matching()?.phase !== 'prepared' || - context.sessions.get(sessionId)?.journal !== ctx.journal - ) { - return - } - await host.flushStreamedEvents(sessionId) - await ctx.journal.appendItem( - identity, - { kind: 'status', text: result.error ?? 'Conversation compacted.' }, - { fence: ctx.fence } - ) - await store.setConversationCommand(sessionId, ctx.fence, { - ...prepared, - phase: 'committed', - state: 'completed', - ...(result.error ? { error: result.error.slice(0, 4096) } : {}) - }) - await store.recordOperationOutcome({ - callerKey: caller.callerKey, - operationId: clientOperationId, - outcome: { - status: 'succeeded', - sessionId, - conversationCommand: matching()! - } - }) - ctx.publish() - }) - }) - ).error - control?.endProviderCall() - await host.flushStreamedEvents(sessionId) - } catch (cause) { - control?.endProviderCall() - await ctx.journal.appendItem( - identity, - { kind: 'status', text: 'Compaction completion is unconfirmed.' }, - { fence: ctx.fence } - ) - ctx.publish() - throw cause - } - } - await ctx.journal.appendItem( - identity, - { kind: 'status', text: error ?? 'Conversation compacted.' }, - { fence: ctx.fence } - ) - ctx.publish() - } - const completed = { - ...prepared, - phase: 'committed' as const, - state: 'completed' as const, - ...(error ? { error: error.slice(0, 4096) } : {}) - } - await store.setConversationCommand(sessionId, ctx.fence, completed) - await settleSupersededOperation(completed) - return { ok: true, value: completed } - } - } + const interruptedClear = source.conversationCommand + const matching = matchingCommand(context, caller.callerKey, params) + const prior = + matching ?? + (command === 'clear' && + interruptedClear?.command === 'clear' && + interruptedClear.phase === 'prepared' && + interruptedClear.runtimeFence !== admitted.context.fence + ? interruptedClear + : null) + const blocked = + prior?.phase === 'prepared' && command === 'clear' + ? null + : conversationCommandBlocked(admitted.context, source) + if (blocked) { + const refusal = { code: 'agent_session_operation_invalid' as const, message: blocked } + await recordRefusal(context, admitted.operationCallerKey, clientOperationId, refusal) + return { decision: 'return', result: { ok: false, refusal } } + } + const replacementSessionId = + command === 'clear' + ? (prior?.replacementSessionId ?? + `clear-${createHash('sha256') + .update(JSON.stringify([sessionId, caller.callerKey, clientOperationId])) + .digest('hex') + .slice(0, 40)}`) + : undefined + const prepared: AgentSessionConversationCommandRecord = { + command, + runtimeFence: admitted.context.fence, + operationId: clientOperationId, + callerKey: caller.callerKey, + phase: 'prepared', + state: 'unknown', + ...(replacementSessionId ? { replacementSessionId } : {}) + } + try { + await store.setConversationCommand(sessionId, admitted.context.fence, prepared) + } catch (error) { + await store.recordOperationOutcome({ + callerKey: admitted.operationCallerKey, + operationId: clientOperationId, + outcome: { status: 'unknown' } }) - ) + throw error + } + return { + decision: 'execute', + execution: { + turn: admitted.context, + prepared, + source, + operationCallerKey: admitted.operationCallerKey, + supersededOperation: prior && prior.operationId !== clientOperationId ? prior : null + } + } +} + +export async function persistConversationCommandResult( + context: StructuredAgentSessionMutationContext, + execution: PreparedConversationCommand, + value: AgentSessionConversationCommandRecord +): Promise { + const { prepared, operationCallerKey, supersededOperation, turn } = execution + await context.deps.store.setConversationCommand(turn.sessionId, turn.fence, value) + await context.deps.store.recordOperationOutcome({ + callerKey: operationCallerKey, + operationId: prepared.operationId, + outcome: { status: 'succeeded', sessionId: turn.sessionId, conversationCommand: value } + }) + if (supersededOperation) { + await context.deps.store.recordOperationOutcome({ + callerKey: supersededOperation.callerKey, + operationId: supersededOperation.operationId, + outcome: { status: 'succeeded', sessionId: turn.sessionId, conversationCommand: value } + }) + } +} + +export function conversationCommandResult( + execution: PreparedConversationCommand, + value: AgentSessionConversationCommandResult +): AgentSessionMutationResult { + return { + ok: true, + replayed: false, + fence: execution.turn.fence, + cursor: execution.turn.journal.cursor(), + value + } } diff --git a/src/renderer/src/components/native-chat/structured-conversation-command-claim.test.ts b/src/renderer/src/components/native-chat/structured-conversation-command-claim.test.ts index be85211d891..287ef80fb59 100644 --- a/src/renderer/src/components/native-chat/structured-conversation-command-claim.test.ts +++ b/src/renderer/src/components/native-chat/structured-conversation-command-claim.test.ts @@ -1,36 +1,33 @@ -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { describe, expect, it, vi } from 'vitest' import { agentJournalSubmissionKey } from '../../../../shared/agent-session-journal-item-key' import type { AgentJournalRenderItem } from '../../../../shared/agent-session-journal-types' import { - CONVERSATION_COMMAND_DEADLINE_MS, StructuredConversationCommandClaim, type ConversationCommandReply } from './structured-conversation-command-claim' const OPERATION_ID = 'op-1' -function compactItem( - operationId: string, - state: 'running' | 'completed' | 'unconfirmed' | 'failed' +function lifecycleItem( + command: 'clear' | 'compact', + state: 'running' | 'completed' | 'unverifiable', + text?: string ): AgentJournalRenderItem { return { - itemId: agentJournalSubmissionKey(`compact:${operationId}`), + itemId: agentJournalSubmissionKey(`${command}:${OPERATION_ID}`), revision: state === 'running' ? 1 : 2, sequence: 1, observedAt: 1, body: { kind: 'status', text: - state === 'running' - ? 'Compacting conversation…' - : state === 'completed' + text ?? + (state === 'running' + ? `${command === 'compact' ? 'Compacting' : 'Clearing'} conversation…` + : command === 'compact' ? 'Conversation compacted.' - : state === 'unconfirmed' - ? 'Compaction completion is unconfirmed.' - : 'Provider refused compaction.', - ...(state === 'running' - ? { turnLifecycle: { turnId: `compact:${operationId}`, state: 'running' as const } } - : {}) + : 'Conversation cleared.'), + turnLifecycle: { turnId: `${command}:${OPERATION_ID}`, state } } } } @@ -39,191 +36,88 @@ function neverReplies(): Promise { return new Promise(() => {}) } -type TrackedOutcome = { settled: boolean; accepted: boolean; error: string | null } - -function track(promise: Promise<{ accepted: boolean; error: string | null }>): TrackedOutcome { - const outcome: TrackedOutcome = { settled: false, accepted: false, error: null } - void promise.then((value) => { - outcome.settled = true - outcome.accepted = value.accepted - outcome.error = value.error - }) - return outcome -} - -beforeEach(() => { - vi.useFakeTimers() -}) - -afterEach(() => { - vi.useRealTimers() -}) - describe('StructuredConversationCommandClaim', () => { - it('settles a compaction on its terminal frame even when no reply ever arrives', async () => { - const claim = new StructuredConversationCommandClaim() - - const outcome = track( - claim.run({ - command: 'compact', + it.each(['compact', 'clear'] as const)( + 'settles %s from typed host lifecycle when the reply is lost', + async (command) => { + const claim = new StructuredConversationCommandClaim() + const outcome = claim.run({ + command, operationId: OPERATION_ID, blocked: false, send: neverReplies }) - ) - await vi.advanceTimersByTimeAsync(1) - claim.applyStreamSnapshot([compactItem(OPERATION_ID, 'running')]) - expect(outcome.settled).toBe(false) + expect(claim.applyStreamSnapshot([lifecycleItem(command, 'running')])).toBe(false) + expect(claim.applyStreamSnapshot([lifecycleItem(command, 'completed')])).toBe(true) + await expect(outcome).resolves.toEqual({ accepted: true, error: null }) + } + ) + + it('keeps an unverifiable host lifecycle pending until interrupt or restart', async () => { + const claim = new StructuredConversationCommandClaim() + const outcome = claim.run({ + command: 'compact', + operationId: OPERATION_ID, + blocked: false, + send: async () => ({ result: { command: 'compact', state: 'unknown' }, unresolved: false }) + }) + await Promise.resolve() + expect(claim.applyStreamSnapshot([lifecycleItem('compact', 'unverifiable')])).toBe(false) expect(claim.isRunning).toBe(true) - claim.applyStreamSnapshot([compactItem(OPERATION_ID, 'completed')]) - await vi.advanceTimersByTimeAsync(0) - - expect(outcome).toMatchObject({ settled: true, accepted: true }) - expect(claim.isRunning).toBe(false) - }) - - it('refuses a second attempt after the deadline instead of racing the first', async () => { - const claim = new StructuredConversationCommandClaim() - const send = vi.fn(neverReplies) - - const first = track( - claim.run({ command: 'compact', operationId: OPERATION_ID, blocked: false, send }) - ) - await vi.advanceTimersByTimeAsync(CONVERSATION_COMMAND_DEADLINE_MS + 1) - - expect(first).toMatchObject({ settled: true, accepted: false }) - expect(first.error).toContain('may still be running') - expect(claim.isRunning).toBe(false) - - const second = track( - claim.run({ command: 'compact', operationId: 'op-2', blocked: false, send }) - ) - await vi.advanceTimersByTimeAsync(0) - - expect(second).toMatchObject({ settled: true, accepted: false }) - expect(second.error).toContain('may still be running') - // The refusal is the point: only one attempt ever reached the host. - expect(send).toHaveBeenCalledTimes(1) - }) - - it('retires the marker when the terminal frame lands late', async () => { - const claim = new StructuredConversationCommandClaim() - const send = vi.fn(neverReplies) - - track(claim.run({ command: 'compact', operationId: OPERATION_ID, blocked: false, send })) - await vi.advanceTimersByTimeAsync(CONVERSATION_COMMAND_DEADLINE_MS + 1) - claim.applyStreamSnapshot([compactItem(OPERATION_ID, 'completed')]) - - track(claim.run({ command: 'compact', operationId: 'op-2', blocked: false, send })) - await vi.advanceTimersByTimeAsync(0) - - expect(send).toHaveBeenCalledTimes(2) - }) - - it('retires the marker on a restart or an interrupt', async () => { - const claim = new StructuredConversationCommandClaim() - const send = vi.fn(neverReplies) - - track(claim.run({ command: 'compact', operationId: OPERATION_ID, blocked: false, send })) - await vi.advanceTimersByTimeAsync(CONVERSATION_COMMAND_DEADLINE_MS + 1) claim.reset() - - track(claim.run({ command: 'compact', operationId: 'op-2', blocked: false, send })) - await vi.advanceTimersByTimeAsync(0) - - expect(send).toHaveBeenCalledTimes(2) + await expect(outcome).resolves.toMatchObject({ accepted: false }) }) - it('keeps waiting on a reply the host could not confirm', async () => { + it('preserves a provider failure carried by terminal lifecycle', async () => { const claim = new StructuredConversationCommandClaim() + const outcome = claim.run({ + command: 'compact', + operationId: OPERATION_ID, + blocked: false, + send: neverReplies + }) + claim.applyStreamSnapshot([ + lifecycleItem('compact', 'completed', 'Provider refused compaction.') + ]) + await expect(outcome).resolves.toEqual({ + accepted: false, + error: 'Provider refused compaction.' + }) + }) - const outcome = track( + it('retains the operation id for retry when transport returns no reply', async () => { + const claim = new StructuredConversationCommandClaim() + await expect( claim.run({ command: 'compact', operationId: OPERATION_ID, blocked: false, - send: async () => ({ result: { command: 'compact', state: 'unknown' }, unresolved: false }) + send: async () => ({ result: null, unresolved: true }) }) - ) - await vi.advanceTimersByTimeAsync(1) - - expect(outcome.settled).toBe(false) - - claim.applyStreamSnapshot([compactItem(OPERATION_ID, 'completed')]) - await vi.advanceTimersByTimeAsync(0) - - expect(outcome).toMatchObject({ settled: true, accepted: true }) + ).resolves.toMatchObject({ accepted: false, retrySameOperation: true }) + expect(claim.isRunning).toBe(false) }) - it('settles a clear on its committed reply', async () => { + it('uses replies from older hosts that publish no typed lifecycle', async () => { const claim = new StructuredConversationCommandClaim() - - const outcome = track( + await expect( claim.run({ command: 'clear', operationId: OPERATION_ID, blocked: false, send: async () => ({ result: { command: 'clear', state: 'completed' }, unresolved: false }) }) - ) - await vi.advanceTimersByTimeAsync(0) - - expect(outcome).toMatchObject({ settled: true, accepted: true, error: null }) - expect(claim.isRunning).toBe(false) + ).resolves.toEqual({ accepted: true, error: null }) }) - it('keeps an unconfirmed host frame pending until the client deadline', async () => { - const claim = new StructuredConversationCommandClaim() - const outcome = track( - claim.run({ - command: 'compact', - operationId: OPERATION_ID, - blocked: false, - send: neverReplies - }) - ) - - claim.applyStreamSnapshot([compactItem(OPERATION_ID, 'unconfirmed')]) - await vi.advanceTimersByTimeAsync(CONVERSATION_COMMAND_DEADLINE_MS - 1) - expect(outcome.settled).toBe(false) - - await vi.advanceTimersByTimeAsync(2) - expect(outcome).toMatchObject({ settled: true, accepted: false }) - expect(outcome.error).toContain('may still be running') - }) - - it('preserves a provider failure carried by the terminal frame', async () => { - const claim = new StructuredConversationCommandClaim() - const outcome = track( - claim.run({ - command: 'compact', - operationId: OPERATION_ID, - blocked: false, - send: neverReplies - }) - ) - - claim.applyStreamSnapshot([compactItem(OPERATION_ID, 'failed')]) - await vi.advanceTimersByTimeAsync(0) - - expect(outcome).toMatchObject({ - settled: true, - accepted: false, - error: 'Provider refused compaction.' - }) - }) - - it('refuses a concurrent command while one is still outstanding', async () => { + it('refuses a concurrent command without sending it', async () => { const claim = new StructuredConversationCommandClaim() const send = vi.fn(neverReplies) - - track(claim.run({ command: 'compact', operationId: OPERATION_ID, blocked: false, send })) - const second = track(claim.run({ command: 'clear', operationId: 'op-2', blocked: false, send })) - await vi.advanceTimersByTimeAsync(0) - - expect(second).toMatchObject({ settled: true, accepted: false }) - expect(second.error).toContain('Wait for the conversation operation to finish') + void claim.run({ command: 'compact', operationId: OPERATION_ID, blocked: false, send }) + await expect( + claim.run({ command: 'clear', operationId: 'op-2', blocked: false, send }) + ).resolves.toMatchObject({ accepted: false }) expect(send).toHaveBeenCalledTimes(1) }) }) diff --git a/src/renderer/src/components/native-chat/structured-conversation-command-claim.ts b/src/renderer/src/components/native-chat/structured-conversation-command-claim.ts index c22cce27c82..661ce24b0c7 100644 --- a/src/renderer/src/components/native-chat/structured-conversation-command-claim.ts +++ b/src/renderer/src/components/native-chat/structured-conversation-command-claim.ts @@ -1,13 +1,3 @@ -// One conversation command per session at a time, resolved against the command's own terminal frame. -// -// A conversation command runs on the host long after the request that started it; the reply is a -// receipt, not a completion. So the claim below waits for the frame the host revises when the -// command finishes -- carried on the session's own journal stream -- and treats the reply as one way -// that frame can be learned rather than as the answer itself. When neither arrives inside the -// deadline the claim does not disappear: it becomes a marker that refuses the next attempt, because -// elapsed time is not evidence the first attempt stopped. The marker retires when the frame lands -// late, when the session restarts, or when the user interrupts. - import { isAgentSessionConversationCommandResult, type AgentSessionConversationCommand, @@ -18,90 +8,47 @@ import type { AgentJournalRenderItem } from '../../../../shared/agent-session-jo import { readAgentJournalTurn } from '../../../../shared/agent-session-turn-record' import { translate } from '@/i18n/i18n' -/** The host settles a compaction on the provider's terminal frame inside its own 180s window - * (`structured-session-compaction.ts`). A client that gave up sooner would call a command - * unresolved while the host still knew the answer. */ -export const CONVERSATION_COMMAND_DEADLINE_MS = 195_000 - export type ConversationCommandOutcome = { accepted: boolean; error: string | null } -type ConversationCommandClaimOutcome = ConversationCommandOutcome & { retrySameOperation?: true } - -/** What the host said, as far as the client can tell. `unresolved` means the reply never arrived, - * which leaves the command possibly still running. */ +export type ConversationCommandClaimOutcome = ConversationCommandOutcome & { + retrySameOperation?: true +} export type ConversationCommandReply = { result: AgentSessionConversationCommandResult | null unresolved: boolean } -type Obligation = { +type LiveClaim = { command: AgentSessionConversationCommand operationId: string - /** The journal item the host revises when the command reaches its terminal frame. Compaction - * only: a clear writes no journal item, so it has nothing on the stream to wait for. */ - terminalItemId: string | null -} - -type LiveClaim = Obligation & { - deadline: ReturnType settle: (outcome: ConversationCommandClaimOutcome) => void - onLateReply: () => void } -const COMPACTION_COMPLETED = 'Conversation compacted.' -const COMPACTION_UNCONFIRMED = 'Compaction completion is unconfirmed.' - -function runningMessage(): string { - return translate( - 'components.native-chat.conversationCommand.running', - 'Wait for the conversation operation to finish.' - ) +function message(key: string, fallback: string): string { + return translate(`components.native-chat.conversationCommand.${key}`, fallback) } -function unresolvedMessage(command: AgentSessionConversationCommand): string { - return translate( - 'components.native-chat.conversationCommand.mayStillBeRunning', - 'The previous /{{value0}} may still be running. Restart the session before running it again.', - { value0: command } - ) -} - -function pendingWorkMessage(): string { - return translate( - 'components.native-chat.conversationCommand.pendingWork', - 'Wait for pending work and messages to finish before using this command.' - ) -} - -function unconfirmedMessage(): string { - return translate( - 'components.native-chat.conversationCommand.unconfirmed', - 'Conversation operation was not confirmed.' - ) -} - -/** The command row is a keyed host projection. Keep old-host unconfirmed rows pending, and preserve - * provider failures instead of turning every non-running revision into success. */ function terminalFrameOutcome( items: readonly AgentJournalRenderItem[], - itemId: string | null + claim: LiveClaim ): ConversationCommandOutcome | null { - if (itemId === null) { + const item = items.find( + (entry) => entry.itemId === agentJournalSubmissionKey(`${claim.command}:${claim.operationId}`) + ) + if (item?.body.kind !== 'status') { return null } - const item = items.find((entry) => entry.itemId === itemId) - if ( - item?.body.kind !== 'status' || - readAgentJournalTurn(item.body)?.state === 'running' || - item.body.text === COMPACTION_UNCONFIRMED - ) { + const lifecycle = readAgentJournalTurn(item.body) + if (lifecycle?.state === 'running' || lifecycle?.state === 'unverifiable') { return null } - return item.body.text === COMPACTION_COMPLETED + const completedText = + claim.command === 'compact' ? 'Conversation compacted.' : 'Conversation cleared.' + return (lifecycle === null || lifecycle.state === 'completed') && item.body.text === completedText ? { accepted: true, error: null } : { accepted: false, error: item.body.text } } -/** A reply the host could not confirm; the operation id must stay reusable for the same attempt. */ +/** A reply the host could not confirm; retries must keep the same durable operation id. */ export function isUnconfirmedConversationCommand(method: string, value: unknown): boolean { return ( method === 'agentSession.conversationCommand' && @@ -110,23 +57,20 @@ export function isUnconfirmedConversationCommand(method: string, value: unknown) ) } +/** Correlates one in-flight request with host lifecycle. Durable ownership remains on the host. */ export class StructuredConversationCommandClaim { private live: LiveClaim | null = null - private unresolved: Obligation | null = null - constructor(private readonly deadlineMs: number = CONVERSATION_COMMAND_DEADLINE_MS) {} - - /** A command is outstanding; sends stay blocked until it settles. */ get isRunning(): boolean { return this.live !== null } get hasObligation(): boolean { - return this.live !== null || this.unresolved !== null + return this.live !== null } isOperationOutstanding(operationId: string): boolean { - return this.live?.operationId === operationId || this.unresolved?.operationId === operationId + return this.live?.operationId === operationId } run(input: { @@ -134,109 +78,80 @@ export class StructuredConversationCommandClaim { operationId: string blocked: boolean send: () => Promise - onLateReply?: () => void }): Promise { - if (this.live) { - return Promise.resolve({ accepted: false, error: runningMessage() }) + if (this.live || input.blocked) { + return Promise.resolve({ + accepted: false, + error: message( + this.live ? 'running' : 'pendingWork', + this.live + ? 'Wait for the conversation operation to finish.' + : 'Wait for pending work and messages to finish before using this command.' + ) + }) } - if (this.unresolved) { - return Promise.resolve({ accepted: false, error: unresolvedMessage(this.unresolved.command) }) - } - if (input.blocked) { - return Promise.resolve({ accepted: false, error: pendingWorkMessage() }) - } - const { promise, resolve } = Promise.withResolvers() + const waiter = Promise.withResolvers() const claim: LiveClaim = { command: input.command, operationId: input.operationId, - terminalItemId: - input.command === 'compact' - ? agentJournalSubmissionKey(`compact:${input.operationId}`) - : null, - deadline: setTimeout(() => this.expire(claim), this.deadlineMs), - settle: resolve, - onLateReply: input.onLateReply ?? (() => {}) + settle: waiter.resolve } this.live = claim void input.send().then( (reply) => this.applyReply(claim, reply), - // A thrown send is the same as no reply: the request may still be running. - () => {} + () => this.finishUnconfirmed(claim) ) - return promise + return waiter.promise } - /** Fold in one snapshot of the session's own stream. Returns true when host truth retired work. */ applyStreamSnapshot(items: readonly AgentJournalRenderItem[]): boolean { - const liveOutcome = this.live ? terminalFrameOutcome(items, this.live.terminalItemId) : null - if (this.live && liveOutcome) { - this.finish(this.live, liveOutcome) - return true + if (!this.live) { + return false } - if (this.unresolved && terminalFrameOutcome(items, this.unresolved.terminalItemId) !== null) { - this.unresolved = null - return true + const outcome = terminalFrameOutcome(items, this.live) + if (!outcome) { + return false } - return false + this.finish(this.live, outcome) + return true } - /** A restart or an interrupt supersedes the obligation: nothing is owed any more. */ - reset(retryPreparedClear = false): void { + reset(): void { if (this.live) { this.finish(this.live, { accepted: false, - error: unconfirmedMessage(), - ...(retryPreparedClear && this.live.command === 'clear' - ? { retrySameOperation: true as const } - : {}) + error: message('unconfirmed', 'Conversation operation was not confirmed.') }) } - this.unresolved = null } private applyReply(claim: LiveClaim, reply: ConversationCommandReply): void { - if (reply.unresolved || reply.result?.state === 'unknown') { - // The host either never answered or answered that it cannot confirm. Either way the frame, - // not the reply, decides. + if (reply.result?.state === 'unknown') { return } - if (this.unresolved?.operationId === claim.operationId) { - this.unresolved = null - claim.onLateReply() + if (reply.unresolved || !reply.result) { + this.finishUnconfirmed(claim) return } - this.finish( - claim, - reply.result - ? { accepted: !reply.result.error, error: reply.result.error ?? null } - : { accepted: false, error: unconfirmedMessage() } - ) + this.finish(claim, { + accepted: !reply.result.error, + error: reply.result.error ?? null + }) + } + + private finishUnconfirmed(claim: LiveClaim): void { + this.finish(claim, { + accepted: false, + error: message('unconfirmed', 'Conversation operation was not confirmed.'), + retrySameOperation: true + }) } private finish(claim: LiveClaim, outcome: ConversationCommandClaimOutcome): void { if (this.live !== claim) { return } - clearTimeout(claim.deadline) this.live = null - this.unresolved = null claim.settle(outcome) } - - private expire(claim: LiveClaim): void { - if (this.live !== claim) { - return - } - this.live = null - this.unresolved = { - command: claim.command, - operationId: claim.operationId, - terminalItemId: claim.terminalItemId - } - claim.settle({ - accepted: false, - error: unresolvedMessage(claim.command), - retrySameOperation: true - }) - } } diff --git a/src/renderer/src/components/native-chat/use-structured-conversation-command.test.tsx b/src/renderer/src/components/native-chat/use-structured-conversation-command.test.tsx index b7aad73d48f..a5b97d6a7b5 100644 --- a/src/renderer/src/components/native-chat/use-structured-conversation-command.test.tsx +++ b/src/renderer/src/components/native-chat/use-structured-conversation-command.test.tsx @@ -2,7 +2,7 @@ import { act, renderHook, waitFor } from '@testing-library/react' import { useEffect, useRef } from 'react' -import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { beforeEach, describe, expect, it, vi } from 'vitest' import { agentJournalSubmissionKey } from '../../../../shared/agent-session-journal-item-key' import type { AgentJournalRenderItem } from '../../../../shared/agent-session-journal-types' @@ -12,7 +12,6 @@ vi.mock('@/runtime/structured-agent-session-client', () => ({ callStructuredAgentSession: mocks.call })) -import { CONVERSATION_COMMAND_DEADLINE_MS } from './structured-conversation-command-claim' import { useStructuredAgentSessionMutate } from './use-structured-agent-session-mutate' import { useStructuredConversationCommand } from './use-structured-conversation-command' @@ -54,12 +53,7 @@ describe('useStructuredConversationCommand', () => { vi.clearAllMocks() }) - afterEach(() => { - vi.useRealTimers() - }) - - it('reuses an unresolved clear identity after a host restart', async () => { - vi.useFakeTimers() + it('retires an unresolved clear identity after a host restart', async () => { mocks.call.mockRejectedValue(new Error('connection lost')) const initialProps: { fence: number; items: AgentJournalRenderItem[] } = { fence: 1, @@ -69,50 +63,28 @@ describe('useStructuredConversationCommand', () => { let first!: Awaited> await act(async () => { - const pending = view.result.current.run('clear') - await vi.advanceTimersByTimeAsync(CONVERSATION_COMMAND_DEADLINE_MS + 1) - first = await pending + first = await view.result.current.run('clear') }) expect(first.accepted).toBe(false) const firstOperationId = mocks.call.mock.calls[0]![2].envelope.clientOperationId view.rerender({ fence: 2, items: [] }) - let second!: Promise - act(() => { - second = view.result.current.run('clear') + await act(async () => { + await view.result.current.run('clear') }) - expect(mocks.call.mock.calls[1]![2].envelope.clientOperationId).toBe(firstOperationId) - - act(() => view.result.current.retire()) - await second + expect(mocks.call.mock.calls[1]![2].envelope.clientOperationId).not.toBe(firstOperationId) }) - it('retires an expired clear on its late reply and gives the next command a fresh identity', async () => { - vi.useFakeTimers() - const firstReply = Promise.withResolvers<{ - ok: true - value: { command: 'clear'; state: 'completed' } - }>() + it('retires a completed clear and gives the next command a fresh identity', async () => { mocks.call - .mockImplementationOnce(() => firstReply.promise) + .mockResolvedValueOnce({ ok: true, value: { command: 'clear', state: 'completed' } }) .mockImplementation(() => new Promise(() => {})) const view = renderHook(() => useCommandHarness({ fence: 1, items: [] })) - let first!: ReturnType - act(() => { - first = view.result.current.run('clear') + await act(async () => { + await view.result.current.run('clear') }) const firstOperationId = mocks.call.mock.calls[0]![2].envelope.clientOperationId - await act(async () => { - await vi.advanceTimersByTimeAsync(CONVERSATION_COMMAND_DEADLINE_MS + 1) - }) - await expect(first).resolves.toMatchObject({ accepted: false }) - - await act(async () => { - firstReply.resolve({ ok: true, value: { command: 'clear', state: 'completed' } }) - await firstReply.promise - await Promise.resolve() - }) let second!: ReturnType act(() => { diff --git a/src/renderer/src/components/native-chat/use-structured-conversation-command.ts b/src/renderer/src/components/native-chat/use-structured-conversation-command.ts index a744e913763..58891cfa85e 100644 --- a/src/renderer/src/components/native-chat/use-structured-conversation-command.ts +++ b/src/renderer/src/components/native-chat/use-structured-conversation-command.ts @@ -48,8 +48,8 @@ export function useStructuredConversationCommand(args: { const current = claim.current const ids = operationIds.current return () => { - ids.delete('compact') - current.reset(true) + ids.clear() + current.reset() } }, [fence]) @@ -77,12 +77,6 @@ export function useStructuredConversationCommand(args: { command, operationId, blocked, - onLateReply: () => { - if (operationIds.current.get(command) === operationId) { - operationIds.current.delete(command) - } - onReconciled(operationId) - }, send: async (): Promise => { let unresolved = false const result = await mutate(