From 85c3ce3775bb17c95798796c1be4a9615ffc6d5e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Igor=20=C5=A0=C4=87eki=C4=87?= Date: Wed, 23 Sep 2026 00:12:29 +0200 Subject: [PATCH] fix(cloud-agent-next): recover an accepted no-output turn once before failing it https://github.com/Kilo-Org/cloud/pull/6582 --- services/cloud-agent-next/AGENTS.md | 2 +- .../src/sandbox-session/SandboxSession.ts | 118 ++++- .../accepted-stored-settlement.test.ts | 101 +++++ .../accepted-stored-settlement.ts | 89 ++++ .../session-message-queue.test.ts | 425 +++++++++++++++--- .../sandbox-session/session-message-queue.ts | 71 +++ .../src/session/preparation-test-helpers.ts | 56 +++ .../src/session/session-message-queue.test.ts | 87 ++++ .../src/session/session-message-queue.ts | 81 ++++ .../src/session/session-message-state.test.ts | 102 +++++ .../src/session/session-message-state.ts | 65 +++ .../src/session/wrapper-supervisor.test.ts | 155 ++++++- .../src/session/wrapper-supervisor.ts | 150 +++++-- ...sandbox-session-no-output-recovery.test.ts | 149 ++++++ .../session/execute-directly-failure.test.ts | 74 ++- 15 files changed, 1614 insertions(+), 111 deletions(-) create mode 100644 services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.test.ts create mode 100644 services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.ts create mode 100644 services/cloud-agent-next/test/integration/sandbox-session-no-output-recovery.test.ts diff --git a/services/cloud-agent-next/AGENTS.md b/services/cloud-agent-next/AGENTS.md index 1a36855900..ed001fd9cb 100644 --- a/services/cloud-agent-next/AGENTS.md +++ b/services/cloud-agent-next/AGENTS.md @@ -127,7 +127,7 @@ This pattern blocks API endpoints from running for external contributors who don - New pending-message writes use a versioned record containing one nested immutable `SessionMessageIntent`, delivery retry state, and callback snapshot; flat pending rows are decoded only inside `src/session/pending-messages.ts`. - `SessionMessageState` owns lifecycle/outbox status, terminal effect accounting, and a named immutable `admissionSnapshot` only for post-pending replay validation and recovery; predecessor records normalize into partial `legacyAdmissionConstraints` and never fabricate missing immutable input. Terminal and accepted/sent effects are repairable from pending/alarm replay and events use deterministic uniqueness. - Wrapper handoff is currently at-least-once under ambiguous delivery failures: the wrapper forwards prompt/command submissions directly to Kilo and does not query Kilo to suppress or recover duplicate `messageId` submissions. Duplicate prompt/command processing is an accepted edge-case trade-off until Kilo provides an atomic submit-or-return-existing contract. -- When accepted work has no pending residue and its fenced wrapper runtime/socket is gone, disconnect or liveness expiry first reconciles each accepted message against the DO's stored kilocode events (`getAssistantMessageForUserMessage`): positive terminal evidence (assistant `time.completed` or a terminal assistant error) settles the message as `idle_reconciliation`; anything else terminalizes as wrapper failure without redispatch. There is still no live authoritative Kilo terminal query for redispatch; adding one remains separate lifecycle capability work. +- When accepted work has no pending residue and its fenced wrapper runtime/socket is gone, disconnect or liveness expiry first reconciles each accepted message against the DO's stored kilocode events (`getAssistantMessageForUserMessage`): positive terminal evidence (assistant `time.completed` or a terminal assistant error) settles the message as `idle_reconciliation`; anything else terminalizes as wrapper failure. A detection of no output on an accepted turn is recoverable once: the accepted-message inactivity path here and the wrapper no-output watchdog in the legacy plane each re-dispatch the same durable intent on a fresh runtime, and only a second identical detection terminalizes as wrapper failure with the attempt count recorded in the failure payload (`attempts`). There is still no live authoritative Kilo terminal query; adding one remains separate lifecycle capability work. - Legacy wrapper-supervisor physical cleanup exhaustion (`WRAPPER_STOP_MAX_ATTEMPTS` reached) is fenced but recoverable: the lease re-observes the sandbox on a slow cadence (`WRAPPER_CLEANUP_EXHAUSTED_RECHECK_MS`) and releases to `none` only after a confirmed `absent` observation, so a wedged sandbox reaped later by the container runtime does not brick the session. Recovery is observation-only (`observeWrappersWithoutWaking`) and never issues another stop: the attempt budget and its rollback fence still hold, and the probe must not wake a stopped container to ask about a process that cannot outlive it. Background rechecks stop after `WRAPPER_CLEANUP_EXHAUSTED_RECHECK_WINDOW_MS` so an unrecoverable exhaustion stops re-arming the DO alarm; explicit sends still force one probe afterwards. - A legacy pending-message flush blocked on exhaustion forces one out-of-cadence recheck (`recoverExhaustedDeliveryBlock` → `recheckExhaustedCleanup`) because a user is actively waiting, then retries on the `WRAPPER_CLEANUP_EXHAUSTED` budget before failing closed. The retry budget is what keeps the two halves consistent — recovery takes minutes, so terminalizing on the first blocked attempt would discard messages a later probe would have delivered — and failing closed at the end of it is what keeps a message from sitting `queued` with no terminal signal. The flush failure code must stay authoritative: `INTERNAL` is treated as non-authoritative by `recordPendingFlushFailure` and would terminalize the message under whatever earlier cause it carried. - Callback delivery retry policy is paired with `wrangler.jsonc`: `CALLBACK_DELIVERY_MAX_ATTEMPTS` includes the initial attempt, and each Cloud Agent Next callback queue consumer must configure `max_retries` for the remaining redeliveries. diff --git a/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts b/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts index 3d3d78b0a5..d6fd16af95 100644 --- a/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts +++ b/services/cloud-agent-next/src/sandbox-session/SandboxSession.ts @@ -226,13 +226,16 @@ import { failWaitingMessages as applyFailWaitingMessages, failedMessageSnapshot, freezeLegacyQueuedMessages, + getSessionMessageTurn, hasAcceptedMessage, hasUnreleasedOperationProof, incrementDeliveryFailure, markSessionOperationRejection, matchesSessionMessageReplay, nextQueuedMessageId, + noOutputRecoveryAllowed, recordAcceptedMessageActivity, + redispatchAcceptedMessage, releaseCompletedRetryableAttach, releaseUnadmittedWaitingMessages, replacePreparationAttemptId, @@ -244,6 +247,11 @@ import { type SessionMessageRecord, } from './session-message-queue.js'; import { createMessageCallbacks, type MessageCallbacks } from './message-callbacks.js'; +import { + applyStoredAssistantSettlement, + projectStoredAssistantSettlement, + type AcceptedStoredSettlement, +} from './accepted-stored-settlement.js'; import { createReportOutbox, readReportAnchor, @@ -3027,6 +3035,11 @@ export class SandboxSession extends DurableObject { report('inactivity'); return; } + if (!this.acceptedTurnIsDrivable(accepted.messageId)) { + diagnostic.reason = 'accepted_message_changed'; + report('superseded'); + return; + } diagnostic.healthy = true; report('healthy'); await this.scheduleAcceptedRecheck(epoch, accepted.messageId); @@ -3059,6 +3072,11 @@ export class SandboxSession extends DurableObject { report('inactivity'); return; } + if (!this.acceptedTurnIsDrivable(accepted.messageId)) { + diagnostic.reason = 'accepted_message_changed'; + report('superseded'); + return; + } if (!waiting) { diagnostic.reason = 'inactive_snapshot'; throw new Error('Accepted execution is no longer active'); @@ -3066,7 +3084,10 @@ export class SandboxSession extends DurableObject { report('healthy'); await this.scheduleAcceptedRecheck(epoch, accepted.messageId); } catch { - if (this.isCurrentAcceptedMessage(accepted, epoch)) { + if ( + this.isCurrentAcceptedMessage(accepted, epoch) && + this.acceptedTurnIsDrivable(accepted.messageId) + ) { diagnostic.reason ??= 'sync_failed'; report('runtime_unhealthy'); await this.failDelivery( @@ -3090,11 +3111,37 @@ export class SandboxSession extends DurableObject { } /** - * The inactivity fail path. Re-reads the accepted row and re-evaluates the - * bound synchronously before persisting, so real progress that landed during - * the preceding observe/sync keeps the turn alive. Message-only: no - * quarantine and no native-runtime retirement. The follow-up abort is - * best-effort and fenced with the captured wrapper instance. + * Positive terminal evidence for an accepted turn that already answered over + * the ingest channel, read from the DO's stored kilocode events. The record + * stays `accepted` because settlement here comes from the prompt operation + * result, so recovery must settle (not re-run) a turn whose answer already + * arrived. Mirrors the legacy wrapper-death reconciliation. + */ + private storedAssistantSettlement(current: MessageRecord): AcceptedStoredSettlement | undefined { + const metadata = this.terminalLifecycle.getStoredMetadata(); + const kiloSessionId = metadata?.auth.kiloSessionId; + if (!metadata || !kiloSessionId) return undefined; + return projectStoredAssistantSettlement( + this.eventQueries.getAssistantMessageForUserMessage( + metadata.identity.sessionId, + kiloSessionId, + current.messageId + ) + ); + } + + /** + * The inactivity path. Re-reads the accepted row and re-evaluates the bound + * synchronously before persisting, so real progress that landed during the + * preceding observe/sync keeps the turn alive. + * + * The first no-output detection is recoverable: the accepted turn is + * re-queued under its durable intent and the wedged runtime is retired, so + * the alarm re-dispatches it once on a replacement runtime. A second + * identical detection keeps the terminal path (message-only: no quarantine + * and no native-runtime retirement) and records the attempt count. The + * follow-up abort is best-effort and fenced with the captured wrapper + * instance. */ private async failOverdueAcceptedMessage( accepted: MessageRecord, @@ -3106,6 +3153,10 @@ export class SandboxSession extends DurableObject { if ( !current || current.state !== 'accepted' || + // A Stop admitted while this alarm awaited its observation keeps the + // record `accepted` and only adds `cancellation`. Recovering it would + // re-queue the turn the user just stopped. + current.cancellation !== undefined || activityAt === undefined || !acceptedInactivityDue(activityAt, Date.now()) ) @@ -3113,6 +3164,46 @@ export class SandboxSession extends DurableObject { diagnostic.stage = 'inactivity'; diagnostic.reason = 'inactivity_due'; diagnostic.lastActivityAt = activityAt; + // Reconcile before recovering: settlement here normally comes from the + // prompt operation result, so a runtime that went silent after streaming + // its answer still leaves the record `accepted`. Re-dispatching it would run + // the turn and its side effects a second time. + const settlement = this.storedAssistantSettlement(current); + if (settlement) { + const settled = applyStoredAssistantSettlement( + this.loadMessages(), + current.messageId, + settlement, + Date.now() + ); + if (!settled || !this.saveMessages(settled, epoch)) return false; + diagnostic.recovery = 'reconciled'; + logControlDiagnostic('accepted_no_output_reconciled', { + ...diagnostic, + producer: 'accepted_inactivity', + result: settlement.state, + }); + return true; + } + const recoveryAttempts = current.recoveryAttempts ?? 0; + if (noOutputRecoveryAllowed(recoveryAttempts)) { + const next = redispatchAcceptedMessage(this.loadMessages(), current.messageId); + if (!next || !this.saveMessages(next, epoch)) return false; + diagnostic.recovery = 'redispatched'; + logControlDiagnostic('accepted_no_output_recovery', { + ...diagnostic, + producer: 'accepted_inactivity', + recoveryAttempts: recoveryAttempts + 1, + }); + const metadata = this.terminalLifecycle.getStoredMetadata(); + if (metadata && current.wrapperInstanceId) + this.retainRuntimeCleanup(metadata, current.wrapperInstanceId, 'no_output_recovery'); + const turn = getSessionMessageTurn(current); + this.broadcastQueuedMessage(current.messageId, turn ? renderExecutionTurnContent(turn) : ''); + await this.armQueueRetry(); + return true; + } + diagnostic.attempts = recoveryAttempts + 1; const wrapperInstanceId = current.wrapperInstanceId; const messages = failAcceptedMessage( this.loadMessages(), @@ -3126,6 +3217,20 @@ export class SandboxSession extends DurableObject { return true; } + /** + * Whether the accepted turn is still the watchdog's to drive. A Stop admitted + * during an awaited observation keeps the record `accepted` and only adds + * `cancellation`; the stop lifecycle settles the turn from there, so the + * watchdog must release it instead of re-arming it or failing it as a runtime + * failure. `failOverdueAcceptedMessage` already refuses to recover such a + * record, which is why the callers re-check here before treating "not overdue" + * as healthy. + */ + private acceptedTurnIsDrivable(messageId: string): boolean { + const current = this.loadMessages().find(item => item.messageId === messageId); + return current?.state === 'accepted' && current.cancellation === undefined; + } + /** * Schedule the next non-failing check from the freshly read clock: * `min(now + acceptedAlarmCap, activityAt + kiloInactivity)`. Never rearm at @@ -3158,6 +3263,7 @@ export class SandboxSession extends DurableObject { if (!this.terminalLifecycle.isCurrent(epoch)) return; if (this.isCurrentAcceptedMessage(accepted, epoch)) { if (await this.failOverdueAcceptedMessage(accepted, epoch, diagnostic)) return; + if (!this.acceptedTurnIsDrivable(accepted.messageId)) return; await this.scheduleAcceptedRecheck(epoch, accepted.messageId); return; } diff --git a/services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.test.ts b/services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.test.ts new file mode 100644 index 0000000000..3b77c9b595 --- /dev/null +++ b/services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.test.ts @@ -0,0 +1,101 @@ +import { describe, expect, it } from 'vitest'; +import type { LatestAssistantMessage } from '../session/types.js'; +import type { SessionMessageRecord } from './session-message-queue.js'; +import { + applyStoredAssistantSettlement, + projectStoredAssistantSettlement, +} from './accepted-stored-settlement.js'; + +function assistant( + info: Record, + parts: LatestAssistantMessage['parts'] = [] +): LatestAssistantMessage { + return { + eventId: 1, + timestamp: 2, + info: { id: 'ase_1', role: 'assistant', ...info } as LatestAssistantMessage['info'], + parts, + }; +} + +function accepted(messageId = 'msg_1'): SessionMessageRecord { + return { messageId, state: 'accepted', acceptedAt: 5 }; +} + +describe('projectStoredAssistantSettlement', () => { + it('settles a stored assistant answer with a completion marker', () => { + expect( + projectStoredAssistantSettlement(assistant({ time: { created: 1, completed: 2 } })) + ).toEqual({ state: 'completed' }); + }); + + it('settles a terminal assistant error with its bounded classification', () => { + expect( + projectStoredAssistantSettlement(assistant({ error: 'Provider request timed out' })) + ).toMatchObject({ + state: 'failed', + failedReason: 'assistant_error', + assistantReason: 'timeout', + }); + }); + + it('settles an interrupted assistant error as a cancellation', () => { + expect( + projectStoredAssistantSettlement( + assistant({ error: { name: 'MessageAbortedError', message: 'interrupted' } }) + ) + ).toEqual({ state: 'cancelled', failedReason: 'interrupted' }); + }); + + it('does not settle a partial answer with no terminal evidence', () => { + expect(projectStoredAssistantSettlement(assistant({ time: { created: 1 } }))).toBeUndefined(); + expect(projectStoredAssistantSettlement(null)).toBeUndefined(); + }); +}); + +describe('applyStoredAssistantSettlement', () => { + it('only settles a record that is still accepted', () => { + const messages = [{ messageId: 'msg_1', state: 'accepted' } as SessionMessageRecord]; + + expect( + applyStoredAssistantSettlement(messages, 'msg_1', { state: 'completed' }, 9)?.[0] + ).toMatchObject({ state: 'completed', terminalAt: 9, terminalSource: 'coordinator' }); + expect( + applyStoredAssistantSettlement(messages, 'msg_1', { state: 'completed' }, 9) + ).toHaveLength(1); + expect( + applyStoredAssistantSettlement( + [{ messageId: 'msg_1', state: 'queued' }], + 'msg_1', + { state: 'completed' }, + 9 + ) + ).toBeUndefined(); + expect(applyStoredAssistantSettlement([], 'msg_1', { state: 'completed' }, 9)).toBeUndefined(); + }); + + it('carries the assistant failure facts onto the terminal record', () => { + const [settled] = + applyStoredAssistantSettlement( + [accepted()], + 'msg_1', + { + state: 'failed', + failedReason: 'assistant_error', + failedDetail: 'Assistant request timed out', + assistantReason: 'timeout', + providerOwnership: 'unknown', + }, + 9 + ) ?? []; + + expect(settled).toMatchObject({ + state: 'failed', + failedReason: 'assistant_error', + failedDetail: 'Assistant request timed out', + assistantReason: 'timeout', + providerOwnership: 'unknown', + terminalSource: 'coordinator', + }); + }); +}); diff --git a/services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.ts b/services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.ts new file mode 100644 index 0000000000..83f08196fa --- /dev/null +++ b/services/cloud-agent-next/src/sandbox-session/accepted-stored-settlement.ts @@ -0,0 +1,89 @@ +import type { + CloudAgentAssistantFailureReason, + CloudAgentProviderOwnership, +} from '@kilocode/worker-utils/cloud-agent-failure'; +import { + classifyAssistantFailure, + isAssistantInterrupt, + projectSafeAssistantError, +} from '../shared/assistant-failure.js'; +import type { AssistantMessageInfo, LatestAssistantMessage } from '../session/types.js'; +import type { SessionMessageRecord } from './session-message-queue.js'; + +/** + * A settlement for an accepted turn read from the DO's stored kilocode events. + * Positive terminal evidence only: a silent runtime can have streamed a partial + * answer, and that is not terminal. + */ +export type AcceptedStoredSettlement = { + state: 'completed' | 'failed' | 'cancelled'; + failedReason?: string; + failedDetail?: string; + assistantReason?: CloudAgentAssistantFailureReason; + providerOwnership?: CloudAgentProviderOwnership; +}; + +function hasAssistantCompletionMarker(info: AssistantMessageInfo): boolean { + const time = info.time; + if (typeof time !== 'object' || time === null || !('completed' in time)) return false; + return typeof time.completed === 'number'; +} + +/** + * Project the stored assistant answer for a user message into a terminal + * settlement, mirroring the legacy wrapper-death reconciliation + * (`projectWrapperDeathReconciliation`). `undefined` means no positive terminal + * evidence, so the caller may still recover the turn. + */ +export function projectStoredAssistantSettlement( + assistant: LatestAssistantMessage | null +): AcceptedStoredSettlement | undefined { + if (!assistant) return undefined; + const error = assistant.info.error; + if (error !== undefined && error !== null) { + if (isAssistantInterrupt(error)) return { state: 'cancelled', failedReason: 'interrupted' }; + const failure = classifyAssistantFailure(error); + return { + state: 'failed', + failedReason: 'assistant_error', + failedDetail: projectSafeAssistantError(error) ?? 'Assistant request failed', + assistantReason: failure.reason, + providerOwnership: failure.providerOwnership, + }; + } + if (!hasAssistantCompletionMarker(assistant.info)) return undefined; + return { state: 'completed' }; +} + +/** + * Apply a stored-assistant settlement to the matching accepted record. Returns + * `undefined` when the record is gone or no longer accepted, so a stale read + * cannot override a newer transition. The terminal source stays `coordinator`: + * the wrapper never reported this settlement. + */ +export function applyStoredAssistantSettlement( + messages: readonly SessionMessageRecord[], + messageId: string, + settlement: AcceptedStoredSettlement, + now: number +): SessionMessageRecord[] | undefined { + const message = messages.find(item => item.messageId === messageId); + if (!message || message.state !== 'accepted') return undefined; + return messages.map(item => + item.messageId !== messageId + ? item + : { + ...item, + state: settlement.state, + unresolvedDispatch: undefined, + terminalAt: now, + terminalSource: 'coordinator', + ...(settlement.failedReason ? { failedReason: settlement.failedReason } : {}), + ...(settlement.failedDetail ? { failedDetail: settlement.failedDetail } : {}), + ...(settlement.assistantReason ? { assistantReason: settlement.assistantReason } : {}), + ...(settlement.providerOwnership + ? { providerOwnership: settlement.providerOwnership } + : {}), + } + ); +} diff --git a/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts b/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts index c9311e823e..9e071880df 100644 --- a/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts +++ b/services/cloud-agent-next/src/sandbox-session/session-message-queue.test.ts @@ -54,14 +54,17 @@ import { completeSessionOperationAttachment, createSessionMessageRecord, failAcceptedMessage, + failedMessageSnapshot, failQueuedMessage, freezeLegacyQueuedMessages, getSessionMessageTurn, hasAcceptedMessage, matchesSessionMessageReplay, nextQueuedMessageId, + noOutputRecoveryAllowed, recordAcceptedMessageActivity, recordSessionOperationDispatch, + redispatchAcceptedMessage, releaseCompletedRetryableAttach, releaseUnadmittedWaitingMessages, releaseUnconfirmedAttach, @@ -846,6 +849,140 @@ describe('rotateLostPreparationAttempt', () => { }); }); +describe('noOutputRecoveryAllowed', () => { + it('allows exactly one no-output recovery per accepted turn', () => { + expect(noOutputRecoveryAllowed(0)).toBe(true); + expect(noOutputRecoveryAllowed(1)).toBe(false); + expect(noOutputRecoveryAllowed(2)).toBe(false); + }); +}); + +describe('redispatchAcceptedMessage', () => { + const promptAuthorization: SessionOperationAuthorization = { + operation: 'session.prompt', + operationId: 'prompt-a', + messageId: promptTurn.messageId, + session: { sessionId: SESSION_ID, kiloSessionId: 'kilo_root', directory: DIRECTORY }, + wrapperInstanceId: RUNTIME_ID, + dispatchDeadlineAt: 100, + }; + + function acceptedMessage(): SessionMessageRecord { + return { + ...createSessionMessageRecord({ turn: promptTurn, agent: defaultAgent }), + state: 'accepted', + acceptedAt: 10, + lastActivityAt: 20, + wrapperInstanceId: RUNTIME_ID, + preparationAttemptId: 'attempt-1', + preparationWait: { step: 'sync', message: 'Waiting for the session' }, + deliveryDeadlineAt: 500, + executionDeadlineAt: 900, + unresolvedDispatch: true, + retryNotBefore: 40, + terminalAt: 30, + terminalSource: 'coordinator', + failedReason: 'accepted_overdue', + failedDetail: 'Turn did not complete', + operations: { + attach: { + authorization: { ...promptAuthorization, operation: 'session.attach' }, + dispatched: true, + }, + prompt: { authorization: promptAuthorization, dispatched: true }, + }, + }; + } + + it('re-queues only the accepted turn, spends one recovery, and clears its dispatch state', () => { + const accepted = acceptedMessage(); + const queued: SessionMessageRecord = { messageId: 'b', state: 'queued' }; + const messages = [accepted, queued]; + + const next = redispatchAcceptedMessage(messages, promptTurn.messageId); + + expect(next?.[0]).toStrictEqual({ + version: 2, + messageId: promptTurn.messageId, + state: 'queued', + intent: { turn: promptTurn, agent: { ...defaultAgent } }, + acceptedAt: undefined, + lastActivityAt: undefined, + deliveryDeadlineAt: undefined, + executionDeadlineAt: undefined, + terminalAt: undefined, + terminalSource: undefined, + failedReason: undefined, + failedDetail: undefined, + wrapperInstanceId: undefined, + preparationAttemptId: undefined, + preparationWait: undefined, + unresolvedDispatch: undefined, + retryNotBefore: undefined, + operations: undefined, + recoveryAttempts: 1, + }); + // The typed turn survives the recovery under its durable identity, and the + // array is rebuilt instead of mutating the caller's records. + expect(next?.[0]?.intent).toStrictEqual(accepted.intent); + expect(next?.[1]).toBe(queued); + expect(messages[0]).toBe(accepted); + }); + + it('accumulates the recovery count for the next accepted state', () => { + const first = redispatchAcceptedMessage([acceptedMessage()], promptTurn.messageId) ?? []; + expect(first[0]?.recoveryAttempts).toBe(1); + // A queued record is not accepted any more, so a replayed detection cannot + // spend a second recovery on it. + expect(redispatchAcceptedMessage(first, promptTurn.messageId)).toBeUndefined(); + const acceptedAgain: SessionMessageRecord = { + ...first[0], + state: 'accepted', + acceptedAt: 50, + lastActivityAt: 60, + wrapperInstanceId: RUNTIME_ID, + }; + expect( + redispatchAcceptedMessage([acceptedAgain], promptTurn.messageId)?.[0]?.recoveryAttempts + ).toBe(2); + }); + + it('returns undefined for a missing or non-accepted record', () => { + expect(redispatchAcceptedMessage([acceptedMessage()], 'missing')).toBeUndefined(); + expect(redispatchAcceptedMessage([msg('a', 'queued')], 'a')).toBeUndefined(); + expect(redispatchAcceptedMessage([msg('a', 'failed')], 'a')).toBeUndefined(); + }); +}); + +describe('failedMessageSnapshot', () => { + it('records the attempt count once a no-output recovery was spent', () => { + const record: SessionMessageRecord = { + ...createSessionMessageRecord({ turn: promptTurn, agent: defaultAgent }), + state: 'failed', + acceptedAt: 10, + recoveryAttempts: 1, + failedReason: 'accepted_overdue', + failedDetail: 'Turn did not complete', + }; + + expect(failedMessageSnapshot(record, 99)).toMatchObject({ + attempts: 2, + reason: 'accepted_overdue', + status: 'failed', + }); + }); + + it('omits the attempt count when the turn never spent a recovery', () => { + const record: SessionMessageRecord = { + ...createSessionMessageRecord({ turn: promptTurn, agent: defaultAgent }), + state: 'failed', + failedReason: 'preparation_failed', + }; + + expect(failedMessageSnapshot(record, 99)).not.toHaveProperty('attempts'); + }); +}); + describe('applyMessageOutcome', () => { it('settles only the message identified by the matching runtime', () => { const before = [{ ...msg('a', 'accepted'), wrapperInstanceId: 'runtime' }, msg('b', 'queued')]; @@ -1267,6 +1404,28 @@ function sessionFixture( return createSessionFixture(fixtureDeps, overrides, sharedControl, callbackQueue); } +/** + * Mark the accepted turn as having already spent its one automatic no-output + * recovery, so the next inactivity detection is the second identical detection + * and reaches the existing terminal path. The recovery itself is driven by the + * accepted-inactivity tests plus `test/integration/sandbox-session-no-output-recovery.test.ts` + * and the `redispatchAcceptedMessage` unit tests. + */ +function markNoOutputRecoverySpent( + fixture: ReturnType, + messageId: string +): void { + const messages = fixture.storage.kv.get('session_messages') ?? []; + if (!messages.some(message => message.messageId === messageId && message.state === 'accepted')) + throw new Error(`No accepted message ${messageId} to recover`); + fixture.storage.kv.put( + 'session_messages', + messages.map(message => + message.messageId === messageId ? { ...message, recoveryAttempts: 1 } : message + ) + ); +} + function controlDiagnostics( fields: { mock: { calls: Parameters[] } }, diagnosticEvent: string @@ -3051,7 +3210,7 @@ describe('SandboxSession orchestration', () => { ); it.each(['running', 'completed'] as const)( - 'settles a late accepted prompt from its original %s operation result without redispatch', + 'reconciles a late accepted prompt against its original %s operation result', async state => { const fixture = sessionFixture(); fixture.setStatus({ @@ -3119,15 +3278,164 @@ describe('SandboxSession orchestration', () => { ([input]) => input.operation === 'session.operation.ack' ) ).toHaveLength(state === 'completed' ? 1 : 0); - expect(fixture.record('a')?.state).toBe(state === 'completed' ? 'completed' : 'failed'); + expect(fixture.record('a')?.state).toBe(state === 'completed' ? 'completed' : 'queued'); if (state === 'running') { - // An operation receipt is liveness, not progress: the 5-minute - // inactivity bound settles the turn instead of refreshing its clock. - expect(fixture.record('a')?.failedReason).toBe('accepted_overdue'); + // An operation receipt is liveness, not progress: the inactivity bound + // spends the turn's one no-output recovery instead of refreshing its + // clock. A second identical detection would terminalize. + expect(fixture.record('a')?.recoveryAttempts).toBe(1); + expect(fixture.record('a')?.failedReason).toBeUndefined(); } } ); + it('settles an accepted no-output turn from a stored completed reply instead of recovering it', async () => { + const fields = vi.spyOn(logger, 'withFields').mockReturnValue(logger); + const fixture = sessionFixture(); + fixture.setStatus({ + physical: 'running', + connection: 'ready', + wrapperInstanceId: RUNTIME_ID, + operationResults: true, + }); + await fixture.admit('a'); + await fixture.flush(); + const authorization = fixture.record('a')?.operations?.prompt?.authorization; + if (!authorization) throw new Error('Missing prompt operation authorization'); + // The runtime went silent, but the answer already reached the DO over the + // ingest channel. Re-dispatching would run the turn and its side effects a + // second time. + fixture.eventQueries.upsert({ + executionId: '', + sessionId: SESSION_ID, + streamEventType: 'kilocode', + entityId: 'message/ase_stored_reply', + payload: JSON.stringify({ + event: 'message.updated', + properties: { + info: { + id: 'ase_stored_reply', + role: 'assistant', + sessionID: 'kilo_root', + parentID: 'a', + time: { created: 1, completed: 2 }, + }, + }, + }), + timestamp: 2, + }); + delegateRequest(fixture, 'session.operation.get', async () => + controlResponse({ state: 'running', authorization }) + ); + + vi.setSystemTime(authorization.dispatchDeadlineAt + 1); + fixture.reload(); + await fixture.fireAlarm(); + await fixture.flush(); + + expect(fixture.record('a')).toMatchObject({ + state: 'completed', + terminalSource: 'coordinator', + }); + expect(fixture.record('a')?.recoveryAttempts).toBeUndefined(); + expect(fixture.record('a')?.failedReason).toBeUndefined(); + expect( + fixture.control.request.mock.calls.filter(([input]) => input.operation === 'session.prompt') + ).toHaveLength(1); + expect(controlDiagnostics(fields, 'accepted_no_output_reconciled')).toHaveLength(1); + }); + + it('does not recover an accepted turn whose stop was admitted while the alarm observed it', async () => { + const fixture = sessionFixture(); + fixture.setStatus({ + physical: 'running', + connection: 'ready', + wrapperInstanceId: RUNTIME_ID, + operationResults: true, + }); + await fixture.admit('a'); + await fixture.flush(); + const authorization = fixture.record('a')?.operations?.prompt?.authorization; + if (!authorization) throw new Error('Missing prompt operation authorization'); + // The user's Stop lands while the alarm awaits the observation: the record + // stays `accepted` and only gains `cancellation`. + delegateRequest(fixture, 'session.operation.get', async () => { + const messages = fixture.storage.kv.get('session_messages') ?? []; + fixture.storage.kv.put( + 'session_messages', + messages.map(message => + message.messageId === 'a' + ? { + ...message, + cancellation: { operationId: 'stop_1', deadlineAt: Date.now() + 60_000 }, + } + : message + ) + ); + return controlResponse({ state: 'running', authorization }); + }); + + vi.setSystemTime(authorization.dispatchDeadlineAt + 1); + fixture.reload(); + await fixture.fireAlarm(); + await fixture.flush(); + + const record = fixture.record('a'); + expect(record?.state).toBe('accepted'); + expect(record?.cancellation).toEqual({ + operationId: 'stop_1', + deadlineAt: expect.any(Number), + }); + expect(record?.recoveryAttempts).toBeUndefined(); + expect(record?.failedReason).toBeUndefined(); + expect( + fixture.control.request.mock.calls.filter(([input]) => input.operation === 'session.prompt') + ).toHaveLength(1); + expect(fixture.control.quarantineRuntime).not.toHaveBeenCalled(); + }); + + it('releases an accepted turn whose stop was admitted while the alarm awaited its health sync', async () => { + const fixture = sessionFixture(); + await fixture.admit('a'); + await fixture.flush(); + const acceptedAt = fixture.record('a')?.acceptedAt; + if (acceptedAt === undefined) throw new Error('Missing accepted timestamp'); + // The liveness snapshot cannot keep the turn alive, so the inactivity bound + // is due. The user's Stop lands while the alarm awaits that snapshot: the + // record stays `accepted` and only gains `cancellation`. The stop lifecycle + // owns the turn now, so the watchdog must release it instead of treating it + // as lost execution and failing it, which would quarantine the wrapper. + delegateRequest(fixture, 'session.sync', async () => { + const messages = fixture.storage.kv.get('session_messages') ?? []; + fixture.storage.kv.put( + 'session_messages', + messages.map(message => + message.messageId === 'a' + ? { + ...message, + cancellation: { operationId: 'stop_1', deadlineAt: Date.now() + 60_000 }, + } + : message + ) + ); + return controlResponse({ status: { type: 'idle' }, questions: [], permissions: [] }); + }); + + vi.setSystemTime(acceptedAt + DEADLINE_MS.kiloInactivity); + await fixture.fireAlarm(); + await fixture.flush(); + + const record = fixture.record('a'); + expect(record?.state).toBe('accepted'); + expect(record?.cancellation).toEqual({ + operationId: 'stop_1', + deadlineAt: expect.any(Number), + }); + expect(record?.failedReason).toBeUndefined(); + expect(fixture.terminalEvents()).toHaveLength(0); + expect(fixture.control.quarantineRuntime).not.toHaveBeenCalled(); + }); + describe('native startup attach authority', () => { it.each([false, true])( 'applies preparing from its pending attach proof without changing prior fence=%s', @@ -6968,7 +7276,7 @@ describe('SandboxSession orchestration', () => { ).toBe(true); }); - it('fails a retrying prompt at seven minutes without real events and fences the abort', async () => { + it('fails a retrying prompt at seven minutes on the second no-output detection and fences the abort', async () => { const fixture = sessionFixture(); const reports: CloudAgentQueueReport[] = []; ( @@ -7007,6 +7315,9 @@ describe('SandboxSession orchestration', () => { // threshold that already elapsed. expect(fixture.alarmAt()).toBe(acceptedAt + DEADLINE_MS.kiloInactivity); + // The turn already spent its one automatic no-output recovery, so this + // second identical detection reaches the existing terminal path. + markNoOutputRecoverySpent(fixture, 'a'); vi.setSystemTime(acceptedAt + DEADLINE_MS.kiloInactivity); await fixture.fireAlarm(); await fixture.flush(); @@ -7245,7 +7556,7 @@ describe('SandboxSession orchestration', () => { properties: { id: 'permission_1', sessionID: 'kilo_root' }, }, ])( - 'keeps a watchdog when a $name event supersedes the health-check sync', + 'keeps a watchdog and recovers once when a $name event supersedes the health-check sync', async ({ type, properties }) => { const fixture = sessionFixture(); const sync = deferred(); @@ -7271,13 +7582,16 @@ describe('SandboxSession orchestration', () => { expect(fixture.alarmAt()).not.toBeNull(); expect(fixture.alarmAt()!).toBeLessThanOrEqual(acceptedAt + DEADLINE_MS.kiloInactivity); + // At the bound the surviving watchdog spends the turn's one no-output + // recovery instead of terminalizing it. vi.setSystemTime(acceptedAt + DEADLINE_MS.kiloInactivity); await fixture.fireAlarm(); await fixture.flush(); expect(fixture.record('a')).toMatchObject({ - state: 'failed', - failedReason: 'accepted_overdue', + state: 'queued', + recoveryAttempts: 1, }); + expect(fixture.record('a')?.failedReason).toBeUndefined(); } ); @@ -7299,6 +7613,9 @@ describe('SandboxSession orchestration', () => { return null; }); + // The turn already spent its one no-output recovery, so this second + // identical detection persists the failure and dispatches the abort. + markNoOutputRecoverySpent(fixture, 'a'); vi.setSystemTime(acceptedAt + DEADLINE_MS.kiloInactivity); const alarm = fixture.fireAlarm(); await entered.promise; @@ -7337,6 +7654,9 @@ describe('SandboxSession orchestration', () => { const acceptedAt = fixture.record('a')?.acceptedAt; if (acceptedAt === undefined) throw new Error('Missing accepted timestamp'); + // Second identical detection: the turn already spent its one no-output + // recovery, so the existing terminal path applies. + markNoOutputRecoverySpent(fixture, 'a'); vi.setSystemTime(acceptedAt + DEADLINE_MS.kiloInactivity); await fixture.fireAlarm(); expect(fixture.record('a')).toMatchObject({ @@ -7367,53 +7687,52 @@ describe('SandboxSession orchestration', () => { questions: [], permissions: [{ id: 'permission_1', sessionID: 'kilo_root' }], }, - ])( - 'treats $name as waiting at 90s but fails at the five-minute inactivity bound', - async snapshot => { - const fixture = sessionFixture(); - const result = { - status: snapshot.status, - questions: snapshot.questions, - permissions: snapshot.permissions, - }; - delegateRequest(fixture, 'session.sync', async input => { - expect(input.session).toEqual({ - sessionId: SESSION_ID, - kiloSessionId: 'kilo_root', - directory: DIRECTORY, - }); - return controlResponse(result); + ])('treats $name as waiting at 90s and recovers once at the inactivity bound', async snapshot => { + const fixture = sessionFixture(); + const result = { + status: snapshot.status, + questions: snapshot.questions, + permissions: snapshot.permissions, + }; + delegateRequest(fixture, 'session.sync', async input => { + expect(input.session).toEqual({ + sessionId: SESSION_ID, + kiloSessionId: 'kilo_root', + directory: DIRECTORY, }); - await fixture.admit('a'); - await fixture.admit('b'); - await fixture.flush(); - const acceptedAt = fixture.record('a')?.acceptedAt; - if (acceptedAt === undefined) throw new Error('Missing accepted timestamp'); + return controlResponse(result); + }); + await fixture.admit('a'); + await fixture.admit('b'); + await fixture.flush(); + const acceptedAt = fixture.record('a')?.acceptedAt; + if (acceptedAt === undefined) throw new Error('Missing accepted timestamp'); - // 90s: a snapshot is liveness, not progress, so the turn is waiting. The - // next wake is capped toward the five-minute inactivity bound. - vi.setSystemTime(acceptedAt + DEADLINE_MS.acceptedOverdue); - await fixture.fireAlarm(); - expect(fixture.record('a')?.state).toBe('accepted'); - expect(fixture.record('b')?.state).toBe('queued'); - expect(fixture.alarmAt()).toBe( - Math.min(Date.now() + DEADLINE_MS.acceptedAlarmCap, acceptedAt + DEADLINE_MS.kiloInactivity) - ); - expect(fixture.control.quarantineRuntime).not.toHaveBeenCalled(); + // 90s: a snapshot is liveness, not progress, so the turn is waiting. The + // next wake is capped toward the five-minute inactivity bound. + vi.setSystemTime(acceptedAt + DEADLINE_MS.acceptedOverdue); + await fixture.fireAlarm(); + expect(fixture.record('a')?.state).toBe('accepted'); + expect(fixture.record('b')?.state).toBe('queued'); + expect(fixture.alarmAt()).toBe( + Math.min(Date.now() + DEADLINE_MS.acceptedAlarmCap, acceptedAt + DEADLINE_MS.kiloInactivity) + ); + expect(fixture.control.quarantineRuntime).not.toHaveBeenCalled(); - // Seven minutes with no real event: the snapshot cannot keep it alive. - vi.setSystemTime(acceptedAt + DEADLINE_MS.kiloInactivity); - await fixture.fireAlarm(); - expect(fixture.record('a')).toMatchObject({ - state: 'failed', - failedReason: 'accepted_overdue', - failedDetail: 'Turn did not complete', - }); - expect(fixture.record('b')?.state).toBe('queued'); - expect(fixture.control.quarantineRuntime).not.toHaveBeenCalled(); - expect(fixture.terminalEvents()).toHaveLength(1); - } - ); + // At the inactivity bound the snapshot cannot keep it alive, so the turn + // spends its one no-output recovery and is re-queued for a fresh runtime + // instead of terminalizing. + vi.setSystemTime(acceptedAt + DEADLINE_MS.kiloInactivity); + await fixture.fireAlarm(); + expect(fixture.record('a')).toMatchObject({ + state: 'queued', + recoveryAttempts: 1, + }); + expect(fixture.record('a')?.failedReason).toBeUndefined(); + expect(fixture.record('b')?.state).toBe('queued'); + expect(fixture.control.quarantineRuntime).not.toHaveBeenCalled(); + expect(fixture.terminalEvents()).toHaveLength(0); + }); it.each(['idle', 'error', 'hang'] as const)( 'fails lost accepted execution on %s health and fences cleanup across reset', diff --git a/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts b/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts index 82a8b0a693..7246892d53 100644 --- a/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts +++ b/services/cloud-agent-next/src/sandbox-session/session-message-queue.ts @@ -71,6 +71,12 @@ type SessionMessageLifecycle = { providerOwnership?: CloudAgentProviderOwnership; attachFailures?: number; promptFailures?: number; + /** + * No-output recoveries already spent on this turn. Bumped by + * `redispatchAcceptedMessage`, so a second identical no-output detection + * terminalizes instead of recovering again. + */ + recoveryAttempts?: number; preparationAttemptId?: string; /** * Durable wait reason for a head whose preparation attempt is finalized but @@ -116,6 +122,21 @@ export type SessionMessageRecord = SessionMessageRecordV2 | LegacySessionMessage export const ATTACH_FAILURE_LIMIT = 2; export const PROMPT_FAILURE_LIMIT = 5; +const NO_OUTPUT_RECOVERY_LIMIT = 1; + +/** + * Whether one more no-output recovery may be spent on an accepted turn. + * + * Both producers of `wrapper_no_output` share this bound: the accepted-message + * inactivity timeout in `SandboxSession.failOverdueAcceptedMessage` here, and + * the wrapper no-output watchdog in `src/session/wrapper-supervisor.ts`. Each + * plane owns its recovery lifecycle, so the same bound is duplicated on purpose: + * a turn gets at most one automatic re-dispatch before a second identical + * detection terminalizes with the attempt count recorded. + */ +export function noOutputRecoveryAllowed(recoveryAttempts: number): boolean { + return recoveryAttempts < NO_OUTPUT_RECOVERY_LIMIT; +} export function resolveSessionMessageIntent( input: ControlSessionMessageInput, @@ -383,6 +404,52 @@ export function releaseUnadmittedWaitingMessages( }; } +/** + * Re-queue an accepted turn whose runtime produced no output, so the alarm can + * re-dispatch the same durable turn on a fresh runtime. The immutable `intent` + * (the user's typed message) is preserved, which is what makes the recovery + * lossless. Every ambiguous dispatch proof is dropped: a late operation result + * for the retired authorization must not revive it, and the cleared prompt + * proof is what lets the delivery start a new acquisition instead of + * reconciling the old one. Only the matching record changes; every other record + * is returned untouched. Returns `undefined` when no accepted record matches. + * + * `deliveryDeadlineAt` and `executionDeadlineAt` are cleared with the other + * dispatch state: both were measured against the runtime that just went silent, + * so keeping either would fail the replacement delivery immediately or hand the + * fresh prompt a bound that belongs to the retired authorization. + */ +export function redispatchAcceptedMessage( + messages: readonly SessionMessageRecord[], + messageId: string +): SessionMessageRecord[] | undefined { + const message = messages.find(item => item.messageId === messageId); + if (!message || message.state !== 'accepted') return undefined; + return messages.map(item => + item.messageId !== messageId + ? item + : { + ...item, + state: 'queued', + acceptedAt: undefined, + lastActivityAt: undefined, + deliveryDeadlineAt: undefined, + executionDeadlineAt: undefined, + terminalAt: undefined, + terminalSource: undefined, + failedReason: undefined, + failedDetail: undefined, + wrapperInstanceId: undefined, + preparationAttemptId: undefined, + preparationWait: undefined, + unresolvedDispatch: undefined, + retryNotBefore: undefined, + operations: undefined, + recoveryAttempts: (item.recoveryAttempts ?? 0) + 1, + } + ); +} + /** * Replace a finalized preparation attempt with a fresh one so later wait * progress is visible. Preparation may resume on an environment rebuild, so a @@ -667,6 +734,10 @@ export function failedMessageSnapshot( delivery: accepted ? 'sent' : 'queued', accepted, reason: cancelled ? 'interrupted' : message.failedReason, + // Record the attempt count only once a no-output recovery was spent on this + // turn: a first-attempt failure has no recovery to account for, and emitting + // `attempts: 1` there would change every other terminal failure payload. + ...(message.recoveryAttempts !== undefined ? { attempts: message.recoveryAttempts + 1 } : {}), ...(cancelled ? { error: 'The message was interrupted' } : message.failedDetail || message.failedReason diff --git a/services/cloud-agent-next/src/session/preparation-test-helpers.ts b/services/cloud-agent-next/src/session/preparation-test-helpers.ts index 21536a2b00..be5eb32094 100644 --- a/services/cloud-agent-next/src/session/preparation-test-helpers.ts +++ b/services/cloud-agent-next/src/session/preparation-test-helpers.ts @@ -1,8 +1,27 @@ import { materializePreparationEvent } from './preparation-history.js'; import type { PreparationAttempt, PreparationStepSnapshot } from '../shared/protocol.js'; import type { EventQueries } from './queries/index.js'; +import type { + AssistantMessageInfo, + AssistantMessagePart, + LatestAssistantMessage, +} from './types.js'; import type { StoredEvent } from '../websocket/types.js'; +type KilocodePayload = { + event?: string; + properties?: { info?: AssistantMessageInfo; part?: AssistantMessagePart }; +}; + +function parseKilocodePayload(payload: string): KilocodePayload | null { + try { + const parsed: unknown = JSON.parse(payload); + return typeof parsed === 'object' && parsed !== null ? (parsed as KilocodePayload) : null; + } catch { + return null; + } +} + /** Entity-keyed in-memory stand-in for the preparation slice of EventQueries. */ export function createMemoryEventQueries(): EventQueries { const rows = new Map(); @@ -34,6 +53,43 @@ export function createMemoryEventQueries(): EventQueries { .filter(([entityId]) => entityId.startsWith(prefix)) .map(([, row]) => row) .sort((a, b) => a.timestamp - b.timestamp || a.id - b.id), + getAssistantMessageForUserMessage: ( + sessionId: string, + kiloSessionId: string, + parentMessageId: string + ): LatestAssistantMessage | null => { + const updated = [...rows.values()] + .filter(row => row.stream_event_type === 'kilocode' && row.session_id === sessionId) + .flatMap(row => { + const payload = parseKilocodePayload(row.payload); + const info = payload?.event === 'message.updated' ? payload.properties?.info : undefined; + if ( + info?.role !== 'assistant' || + info.sessionID !== kiloSessionId || + info.parentID !== parentMessageId + ) + return []; + return [{ row, info }]; + }) + .sort( + (left, right) => left.row.timestamp - right.row.timestamp || left.row.id - right.row.id + ); + const latest = updated[updated.length - 1]; + if (!latest) return null; + const partsById = new Map(); + for (const row of rows.values()) { + const payload = parseKilocodePayload(row.payload); + const part = + payload?.event === 'message.part.updated' ? payload.properties?.part : undefined; + if (part?.messageID === latest.info.id) partsById.set(part.id, part); + } + return { + eventId: latest.row.id, + timestamp: latest.row.timestamp, + info: latest.info, + parts: [...partsById.values()].sort((left, right) => left.id.localeCompare(right.id)), + } satisfies LatestAssistantMessage; + }, } as unknown as EventQueries; } diff --git a/services/cloud-agent-next/src/session/session-message-queue.test.ts b/services/cloud-agent-next/src/session/session-message-queue.test.ts index d29f9314f5..50b60bad82 100644 --- a/services/cloud-agent-next/src/session/session-message-queue.test.ts +++ b/services/cloud-agent-next/src/session/session-message-queue.test.ts @@ -12,12 +12,14 @@ import { buildCloudMessageFailedPayload } from './message-settlement-outbox.js'; import { createSessionMessageQueue, flushNextPendingSessionMessage, + requeueAcceptedMessageForRecovery, type SessionMessageQueueDependencies, type SessionMessageQueueStorage, } from './session-message-queue.js'; import { createPendingSessionMessage, createPendingSessionMessageFromIntent, + findPendingSessionMessageByMessageId, listPendingSessionMessages, PENDING_FLUSH_RETRY_BASE_DELAY_MS, PENDING_SESSION_MESSAGE_LIMIT, @@ -29,6 +31,7 @@ import { createQueuedSessionMessageState, getSessionMessageState, putSessionMessageState, + type SessionMessageState, type TerminalizeParams, } from './session-message-state.js'; import { @@ -2559,3 +2562,87 @@ describe('SessionMessageQueue', () => { } ); }); + +describe('requeueAcceptedMessageForRecovery', () => { + const attachments = { + path: '11111111-1111-4111-8111-111111111111', + files: ['22222222-2222-4222-8222-222222222222.txt'], + }; + + function legacyAcceptedState(): SessionMessageState { + return { + messageId: FIRST_MESSAGE_ID, + status: 'accepted', + prompt: 'review these files', + createdAt: 1, + queuedAt: 2, + acceptedAt: 5, + dispatchAcceptanceKind: 'observed', + wrapperRunId: 'wrapper_run_1', + callbackRequired: false, + legacyAdmissionConstraints: { + turn: { + type: 'prompt', + messageId: FIRST_MESSAGE_ID, + prompt: 'review these files', + attachments, + }, + agent: { mode: 'code', model: 'default-model' }, + }, + }; + } + + it('re-queues a predecessor prompt turn with its attachments intact', async () => { + const storage = createMemoryStorage(); + const state = legacyAcceptedState(); + await putSessionMessageState(storage, state); + + await expect(requeueAcceptedMessageForRecovery(storage, state)).resolves.toBe(true); + + await expect(getSessionMessageState(storage, FIRST_MESSAGE_ID)).resolves.toMatchObject({ + status: 'queued', + recoveryAttempts: 1, + }); + const pending = await listPendingSessionMessages(storage); + expect(pending).toHaveLength(1); + // The legacy pending shape cannot encode prompt attachments, so the + // recovered turn must be rebuilt from the immutable constraints. + expect(pending[0]?.intent?.turn).toEqual({ + type: 'prompt', + messageId: FIRST_MESSAGE_ID, + prompt: 'review these files', + attachments, + }); + expect(pending[0]?.intent?.agent).toEqual({ mode: 'code', model: 'default-model' }); + }); + + it('writes the pending row before committing the queued state', async () => { + const storage = createMemoryStorage([], { failPutPrefix: 'pending_message:' }); + const state = legacyAcceptedState(); + await putSessionMessageState(storage, state); + + await expect(requeueAcceptedMessageForRecovery(storage, state)).rejects.toThrow( + 'failed to put pending_message:' + ); + + // A failed pending write must not leave a `queued` record that neither the + // pending drain nor the accepted-only repair can find: the message stays + // accepted, so the next detection retries the recovery. + await expect(getSessionMessageState(storage, FIRST_MESSAGE_ID)).resolves.toMatchObject({ + status: 'accepted', + }); + }); + + it('drops the pending row when the accepted state no longer applies', async () => { + const storage = createMemoryStorage(); + const state = legacyAcceptedState(); + await putSessionMessageState(storage, state); + await putSessionMessageState(storage, { ...state, status: 'completed', terminalAt: 9 }); + + await expect(requeueAcceptedMessageForRecovery(storage, state)).resolves.toBe(false); + + await expect( + findPendingSessionMessageByMessageId(storage, FIRST_MESSAGE_ID) + ).resolves.toBeUndefined(); + }); +}); diff --git a/services/cloud-agent-next/src/session/session-message-queue.ts b/services/cloud-agent-next/src/session/session-message-queue.ts index b142a5d6d8..9bd91f04c8 100644 --- a/services/cloud-agent-next/src/session/session-message-queue.ts +++ b/services/cloud-agent-next/src/session/session-message-queue.ts @@ -34,6 +34,7 @@ import { resolvePendingSessionMessageIntent, shouldSkipPendingFlush, storePendingSessionMessage, + type LegacyPendingSessionMessage, type PendingFlushFailureResult, type PendingFlushPolicy, type PendingSessionExecutionDefaults, @@ -46,6 +47,7 @@ import { getSessionMessageState, listReconnectVisibleTerminalQueuedMessages, markMessageInterrupted, + markMessageQueuedForRecovery, putSessionMessageState, type SessionMessageFailureCode, type SessionMessageStorage, @@ -1305,3 +1307,82 @@ export async function getQueuedMessageByMessageId( ): Promise { return findPendingSessionMessageByMessageId(storage, messageId); } + +/** + * Re-queue an accepted turn that an unhealthy wrapper left without output so the + * next pending drain re-dispatches it on a fresh runtime. The immutable + * `admissionSnapshot` (or the normalized `legacyAdmissionConstraints` predecessor + * shape) holds the user's typed turn across the recovery. Returns false when no + * intent can be resolved, leaving the caller to terminalize as today. + */ +export async function requeueAcceptedMessageForRecovery( + storage: SessionMessageQueueStorage, + state: SessionMessageState +): Promise { + const intent = resolveRecoveryIntent(state); + if (!intent) return false; + + const callbackSnapshot = state.callbackRequired + ? { required: true, target: state.callbackTarget } + : undefined; + // Write the durable pending row before the queued state. The reverse order + // could strand a `queued` record with no pending row — invisible to both the + // pending drain and the accepted-only repair — if the write throws. With this + // order a failed pending write leaves the record accepted, so the next + // detection retries the recovery, and a state transition that no longer + // applies (the message already terminalized) compensates by dropping the row + // just written. + await enqueuePendingSessionMessageIntent(storage, intent, Date.now(), callbackSnapshot); + const queued = await markMessageQueuedForRecovery(storage, state.messageId); + if (!queued) { + await deletePendingSessionMessageByMessageId(storage, state.messageId); + return false; + } + return true; +} + +function resolveRecoveryIntent(state: SessionMessageState): SessionMessageIntent | undefined { + if (state.admissionSnapshot) return state.admissionSnapshot; + const legacy = state.legacyAdmissionConstraints; + if (!legacy?.turn || !legacy.agent?.model) return undefined; + if (legacy.turn.type === 'prompt') { + // The legacy pending shape cannot encode a prompt turn's attachments, so a + // predecessor prompt is rebuilt from the immutable constraints directly. + // Going through `decodeLegacyPendingMessage` would drop the user's files + // and re-run the recovered turn without them. + return { + turn: { + type: 'prompt', + messageId: legacy.turn.messageId, + prompt: legacy.turn.prompt, + ...(legacy.turn.attachments ? { attachments: legacy.turn.attachments } : {}), + }, + agent: { + mode: legacy.agent.mode ?? 'code', + model: legacy.agent.model, + ...(legacy.agent.variant !== undefined ? { variant: legacy.agent.variant } : {}), + }, + ...(legacy.finalization ? { finalization: legacy.finalization } : {}), + }; + } + const legacyMessage: PendingSessionMessage & { legacy: LegacyPendingSessionMessage } = { + messageId: state.messageId, + content: state.prompt, + createdAt: state.createdAt, + legacy: { + messageId: state.messageId, + role: 'user', + content: state.prompt, + createdAt: state.createdAt, + turn: { type: 'command', command: legacy.turn.command, arguments: legacy.turn.arguments }, + executionOptions: { + mode: legacy.agent.mode, + model: legacy.agent.model, + variant: legacy.agent.variant, + autoCommit: legacy.finalization?.autoCommit, + condenseOnComplete: legacy.finalization?.condenseOnComplete, + }, + }, + }; + return resolvePendingSessionMessageIntent(legacyMessage, {}); +} diff --git a/services/cloud-agent-next/src/session/session-message-state.test.ts b/services/cloud-agent-next/src/session/session-message-state.test.ts index faf59d9e4e..482f6fe087 100644 --- a/services/cloud-agent-next/src/session/session-message-state.test.ts +++ b/services/cloud-agent-next/src/session/session-message-state.test.ts @@ -8,6 +8,8 @@ import { markMessageCompleted, markMessageFailed, markMessageInterrupted, + markMessageQueuedForRecovery, + noOutputRecoveryAllowed, terminalizeMessageOnce, listNonTerminalAcceptedMessages, listMessagesForWrapperRun, @@ -428,6 +430,106 @@ describe('markAgentActivityObserved', () => { }); }); +describe('noOutputRecoveryAllowed', () => { + it('allows exactly one recovery per accepted turn', () => { + expect(noOutputRecoveryAllowed(0)).toBe(true); + expect(noOutputRecoveryAllowed(1)).toBe(false); + expect(noOutputRecoveryAllowed(2)).toBe(false); + }); +}); + +describe('markMessageQueuedForRecovery', () => { + const target = { url: 'https://example.com/callback' }; + + async function putAcceptedMessageForRecovery( + storage: SessionMessageStorage, + overrides?: Partial + ): Promise { + const intent = createIntent(VALID_MESSAGE_ID, 'recover me'); + await putSessionMessageState(storage, { + ...createQueuedSessionMessageState(intent, { required: true, target }, 1000), + status: 'accepted', + acceptedAt: 2000, + dispatchAcceptanceKind: 'observed', + agentActivityObservedAt: 2500, + wrapperRunId: 'wr_recovery', + terminalAt: 3000, + failureReason: 'wrapper_failure', + error: 'Wrapper made no execution progress during the watchdog window', + failureStage: 'post_dispatch_no_activity', + failureCode: 'workspace_setup_failed', + failureSubtype: 'sandbox_storage_full', + assistantFailureReason: 'provider_unavailable', + providerOwnership: 'managed', + safeFailureMessage: 'Assistant service is unavailable', + ...overrides, + }); + return intent; + } + + it('returns the accepted turn to queued and clears terminal fields', async () => { + const storage = createFakeStorage(); + const intent = await putAcceptedMessageForRecovery(storage); + + const updated = await markMessageQueuedForRecovery(storage, VALID_MESSAGE_ID, 4000); + const loaded = await getSessionMessageState(storage, VALID_MESSAGE_ID); + + expect(updated?.status).toBe('queued'); + expect(loaded).toMatchObject({ + status: 'queued', + recoveryAttempts: 1, + prompt: 'recover me', + createdAt: 1000, + callbackRequired: true, + callbackTarget: target, + }); + expect(loaded?.admissionSnapshot).toEqual(intent); + expect(loaded?.acceptedAt).toBeUndefined(); + expect(loaded?.wrapperRunId).toBeUndefined(); + expect(loaded?.dispatchAcceptanceKind).toBeUndefined(); + expect(loaded?.agentActivityObservedAt).toBeUndefined(); + for (const field of [ + 'terminalAt', + 'failureReason', + 'error', + 'failureStage', + 'failureCode', + 'failureSubtype', + 'assistantFailureReason', + 'providerOwnership', + 'safeFailureMessage', + ] as const) { + expect(loaded).not.toHaveProperty(field); + } + }); + + it('increments an existing recovery attempt count', async () => { + const storage = createFakeStorage(); + await putAcceptedMessageForRecovery(storage, { recoveryAttempts: 1 }); + + const updated = await markMessageQueuedForRecovery(storage, VALID_MESSAGE_ID, 4000); + + expect(updated?.recoveryAttempts).toBe(2); + }); + + it('returns null unless the message is accepted', async () => { + const storage = createFakeStorage(); + await putSessionMessageState( + storage, + createQueuedSessionMessageState( + createIntent(VALID_MESSAGE_ID, 'still queued'), + undefined, + 1000 + ) + ); + + expect(await markMessageQueuedForRecovery(storage, VALID_MESSAGE_ID)).toBeNull(); + expect( + await markMessageQueuedForRecovery(storage, 'msg_unknown00000000ABCDEFGHIJKLMN') + ).toBeNull(); + }); +}); + describe('markMessageCompleted', () => { it('transitions accepted to completed', async () => { const storage = createFakeStorage(); diff --git a/services/cloud-agent-next/src/session/session-message-state.ts b/services/cloud-agent-next/src/session/session-message-state.ts index c4883aaf58..659ba76612 100644 --- a/services/cloud-agent-next/src/session/session-message-state.ts +++ b/services/cloud-agent-next/src/session/session-message-state.ts @@ -122,6 +122,12 @@ export type SessionMessageState = { error?: string; failureReason?: string; attempts?: number; + /** + * Number of automatic recoveries already spent on this accepted turn. It is + * the shared one-recovery budget for both no-output producers; a message whose + * `recoveryAttempts` reaches `NO_OUTPUT_RECOVERY_LIMIT` is terminalized. + */ + recoveryAttempts?: number; gateResult?: 'pass' | 'fail'; callbackRequired?: boolean; callbackTarget?: CallbackTarget; @@ -230,6 +236,7 @@ export const SessionMessageStateSchema = z error: z.string().optional(), failureReason: z.string().optional(), attempts: z.number().int().nonnegative().optional(), + recoveryAttempts: z.number().int().nonnegative().optional(), gateResult: z.enum(['pass', 'fail']).optional(), callbackRequired: z.boolean().optional(), callbackTarget: z @@ -498,6 +505,64 @@ export async function markAgentActivityObserved( return updated; } +/** One automatic re-dispatch per accepted turn, shared by both no-output producers. */ +const NO_OUTPUT_RECOVERY_LIMIT = 1; + +/** + * The one-recovery bound for the two producers that raise `wrapper_no_output` + * for an accepted turn that produced no execution progress: + * + * - the control plane's accepted-message inactivity deadline + * (`SandboxSession.failOverdueAcceptedMessage`, per s1), and + * - the legacy wrapper-supervisor no-output watchdog + * (`handleUnhealthyWrapper`, this producer). + * + * Both share the same `recoveryAttempts` counter on `SessionMessageState`, so a + * turn recovered by one producer is not recovered again by the other. The + * harness ends a silenced accepted turn as "Response failed" for the user; the + * first detection re-queues the turn onto a fresh runtime instead. + */ +export function noOutputRecoveryAllowed(recoveryAttempts: number): boolean { + return recoveryAttempts < NO_OUTPUT_RECOVERY_LIMIT; +} + +/** + * Return an accepted message to `queued` so the pending drain re-dispatches it + * on a fresh runtime. The immutable `admissionSnapshot` (or the normalized + * `legacyAdmissionConstraints` predecessor shape) keeps the user's typed turn + * across the recovery; every terminal and dispatch field is cleared and the + * recovery attempt is recorded. + */ +export async function markMessageQueuedForRecovery( + storage: SessionMessageStorage, + messageId: string, + now = Date.now() +): Promise { + const state = await getSessionMessageState(storage, messageId); + if (!state || state.status !== 'accepted') return null; + const updated: SessionMessageState = { + ...state, + status: 'queued', + queuedAt: now, + recoveryAttempts: (state.recoveryAttempts ?? 0) + 1, + }; + delete updated.acceptedAt; + delete updated.wrapperRunId; + delete updated.dispatchAcceptanceKind; + delete updated.agentActivityObservedAt; + delete updated.terminalAt; + delete updated.failureStage; + delete updated.failureCode; + delete updated.failureSubtype; + delete updated.assistantFailureReason; + delete updated.providerOwnership; + delete updated.safeFailureMessage; + delete updated.error; + delete updated.failureReason; + await putSessionMessageState(storage, updated); + return updated; +} + export type MarkMessageCompletedParams = { assistantMessageId?: string; completionSource: SessionMessageCompletionSource; diff --git a/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts b/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts index de4a6c610a..065af54db4 100644 --- a/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts +++ b/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts @@ -9,7 +9,10 @@ import { createMessageSettlementOutbox, type MessageSettlementOutboxStorage, } from './message-settlement-outbox.js'; -import { storePendingSessionMessage } from './pending-messages.js'; +import { + findPendingSessionMessageByMessageId, + storePendingSessionMessage, +} from './pending-messages.js'; import { getSessionMessageState, putSessionMessageState, @@ -114,6 +117,16 @@ function acceptedMessage(messageId = MESSAGE_ID): SessionMessageState { }; } +function acceptedMessageWithIntent(messageId = MESSAGE_ID): SessionMessageState { + return { + ...acceptedMessage(messageId), + admissionSnapshot: { + turn: { type: 'prompt', messageId, prompt: 'supervise this wrapper' }, + agent: { mode: 'code', model: 'test-model' }, + }, + }; +} + function createHarness( initialEntries?: Array<[string, unknown]>, options?: { @@ -1535,6 +1548,146 @@ describe('WrapperSupervisor', () => { expect(harness.events.map(event => event.streamEventType)).toEqual(['cloud.message.failed']); }); + it('re-queues an accepted turn once at the no-output deadline instead of failing it', async () => { + const acceptedAt = 2_000; + const noOutputDeadlineAt = acceptedAt + WRAPPER_NO_OUTPUT_TIMEOUT_MS; + const harness = createHarness([ + liveRuntimeState({ noOutputDeadlineAt, nextPingAt: noOutputDeadlineAt + 1 }), + OWNED_WRAPPER_LEASE, + ]); + await putSessionMessageState(harness.storage, acceptedMessageWithIntent()); + + await harness.supervisor.runMaintenance(noOutputDeadlineAt); + + const state = await getSessionMessageState(harness.storage, MESSAGE_ID); + expect(state).toMatchObject({ status: 'queued', recoveryAttempts: 1 }); + expect(state?.acceptedAt).toBeUndefined(); + expect(state?.wrapperRunId).toBeUndefined(); + await expect( + findPendingSessionMessageByMessageId(harness.storage, MESSAGE_ID) + ).resolves.toMatchObject({ messageId: MESSAGE_ID }); + expect(harness.events.map(event => event.streamEventType)).not.toContain( + 'cloud.message.failed' + ); + expect(harness.requestPendingDrainIfNeeded).toHaveBeenCalledTimes(1); + }); + + it('reconciles a completed reply instead of re-dispatching a no-output turn', async () => { + const acceptedAt = 2_000; + const noOutputDeadlineAt = acceptedAt + WRAPPER_NO_OUTPUT_TIMEOUT_MS; + const assistantMessageId = 'ase_no_output_reconcile'; + const harness = createHarness( + [ + liveRuntimeState({ noOutputDeadlineAt, nextPingAt: noOutputDeadlineAt + 1 }), + OWNED_WRAPPER_LEASE, + ], + { + getAssistantMessageForUserMessage: () => + ({ + info: { + id: assistantMessageId, + role: 'assistant', + time: { created: 2_500, completed: 2_600 }, + }, + parts: [], + }) as unknown as LatestAssistantMessage, + } + ); + await putSessionMessageState(harness.storage, acceptedMessageWithIntent()); + + await harness.supervisor.runMaintenance(noOutputDeadlineAt); + + // The answer already reached the DO over the ingest channel while the + // wrapper went silent, so the turn is settled from that positive evidence. + // Re-dispatching it would run the turn twice and duplicate the reply and its + // side effects. + await expect(getSessionMessageState(harness.storage, MESSAGE_ID)).resolves.toMatchObject({ + status: 'completed', + completionSource: 'idle_reconciliation', + assistantMessageId, + }); + await expect( + findPendingSessionMessageByMessageId(harness.storage, MESSAGE_ID) + ).resolves.toBeUndefined(); + expect(harness.events.map(event => event.streamEventType)).toEqual(['cloud.message.completed']); + expect(harness.requestPendingDrainIfNeeded).not.toHaveBeenCalled(); + }); + + it('fails the second identical no-output detection with the attempt count', async () => { + const acceptedAt = 2_000; + const noOutputDeadlineAt = acceptedAt + WRAPPER_NO_OUTPUT_TIMEOUT_MS; + const harness = createHarness([ + liveRuntimeState({ noOutputDeadlineAt, nextPingAt: noOutputDeadlineAt + 1 }), + OWNED_WRAPPER_LEASE, + ]); + await putSessionMessageState(harness.storage, { + ...acceptedMessageWithIntent(), + recoveryAttempts: 1, + }); + + await harness.supervisor.runMaintenance(noOutputDeadlineAt); + + const state = await getSessionMessageState(harness.storage, MESSAGE_ID); + expect(state).toMatchObject({ + status: 'failed', + failureCode: 'wrapper_no_output', + attempts: 2, + recoveryAttempts: 1, + }); + expect(harness.events.map(event => event.streamEventType)).toEqual(['cloud.message.failed']); + expect(harness.events.map(event => JSON.parse(event.payload))).toContainEqual( + expect.objectContaining({ messageId: MESSAGE_ID, attempts: 2 }) + ); + expect(harness.requestPendingDrainIfNeeded).not.toHaveBeenCalled(); + }); + + it('omits the attempt count when the first no-output detection spent no recovery', async () => { + const acceptedAt = 2_000; + const noOutputDeadlineAt = acceptedAt + WRAPPER_NO_OUTPUT_TIMEOUT_MS; + const harness = createHarness([ + liveRuntimeState({ noOutputDeadlineAt, nextPingAt: noOutputDeadlineAt + 1 }), + OWNED_WRAPPER_LEASE, + ]); + // A predecessor record with no resolvable intent: the recovery cannot be + // spent, so this first detection is terminal and no retry ran to count. + await putSessionMessageState(harness.storage, acceptedMessage()); + + await harness.supervisor.runMaintenance(noOutputDeadlineAt); + + const state = await getSessionMessageState(harness.storage, MESSAGE_ID); + expect(state).toMatchObject({ status: 'failed', failureCode: 'wrapper_no_output' }); + expect(state?.attempts).toBeUndefined(); + const payloads = harness.events.map( + event => JSON.parse(event.payload) as Record + ); + expect(payloads).not.toHaveLength(0); + expect(payloads.every(payload => !('attempts' in payload))).toBe(true); + }); + + it('keeps terminalizing a ping timeout on the first detection', async () => { + const pingDeadlineAt = 92_000; + const noOutputDeadlineAt = 332_000; + const harness = createHarness([ + liveRuntimeState({ pingDeadlineAt, noOutputDeadlineAt }), + OWNED_WRAPPER_LEASE, + ]); + await putSessionMessageState(harness.storage, acceptedMessageWithIntent()); + + await harness.supervisor.runMaintenance(pingDeadlineAt); + + const state = await getSessionMessageState(harness.storage, MESSAGE_ID); + expect(state).toMatchObject({ + status: 'failed', + failureCode: 'wrapper_ping_timeout', + }); + expect(state?.attempts).toBeUndefined(); + expect(state?.recoveryAttempts).toBeUndefined(); + await expect( + findPendingSessionMessageByMessageId(harness.storage, MESSAGE_ID) + ).resolves.toBeUndefined(); + expect(harness.requestPendingDrainIfNeeded).not.toHaveBeenCalled(); + }); + it('terminates an unresponsive wrapper on ping timeout before no-output expires', async () => { const pingDeadlineAt = 92_000; const noOutputDeadlineAt = 332_000; diff --git a/services/cloud-agent-next/src/session/wrapper-supervisor.ts b/services/cloud-agent-next/src/session/wrapper-supervisor.ts index a1cfeb8045..7e95278939 100644 --- a/services/cloud-agent-next/src/session/wrapper-supervisor.ts +++ b/services/cloud-agent-next/src/session/wrapper-supervisor.ts @@ -17,10 +17,14 @@ import { isAssistantInterrupt, } from './safe-failure-projection.js'; import { countPendingSessionMessages, type SessionQueueStorage } from './pending-messages.js'; -import type { SessionMessageQueue } from './session-message-queue.js'; +import { + requeueAcceptedMessageForRecovery, + type SessionMessageQueue, +} from './session-message-queue.js'; import { listMessagesForWrapperRun, listNonTerminalAcceptedMessages, + noOutputRecoveryAllowed, type SessionMessageState, type SessionMessageStorage, type TerminalizeParams, @@ -941,45 +945,64 @@ export function createWrapperSupervisor( } /** - * When the wrapper dies before its terminal report, the DO's event store - * still holds every kilocode event that arrived over the FIFO ingest - * channel — including, by ingest ordering, the completed assistant state - * whenever the wrapper had already started finalizing. Settle from that - * positive terminal evidence when present; otherwise fall back to - * wrapper-failure handling. + * Positive terminal evidence for accepted messages whose wrapper never + * reported a terminal state. The DO's event store still holds every kilocode + * event that arrived over the FIFO ingest channel — including, by ingest + * ordering, the completed assistant state whenever the wrapper had already + * started finalizing. + * + * Reconcile before any no-output recovery: a turn whose answer already + * arrived must be settled from that evidence, never re-dispatched, or it runs + * twice and duplicates the reply and its side effects. Messages with no + * terminal evidence are absent from the returned map. + */ + async function reconcileAcceptedMessages( + acceptedMessages: SessionMessageState[] + ): Promise> { + const metadata = await getMetadata(); + const kiloSessionId = metadata?.auth.kiloSessionId; + const reconciledParams = new Map(); + if (!metadata || !kiloSessionId) return reconciledParams; + const codeReviewSession = metadata.identity.createdOnPlatform === 'code-review'; + for (const message of acceptedMessages) { + const reconciled = projectWrapperDeathReconciliation( + getAssistantMessageForUserMessage( + metadata.identity.sessionId, + kiloSessionId, + message.messageId + ), + codeReviewSession + ); + if (reconciled) reconciledParams.set(message.messageId, reconciled); + } + return reconciledParams; + } + + /** + * When the wrapper dies before its terminal report, settle accepted work from + * the DO's stored positive terminal evidence when present; otherwise fall + * back to wrapper-failure handling. A caller that already reconciled the same + * messages passes its map so the evidence is read once. */ async function terminalizeAcceptedMessagesForDeadWrapper( acceptedMessages: SessionMessageState[], - fallbackParams: (message: SessionMessageState) => TerminalizeParams + fallbackParams: (message: SessionMessageState) => TerminalizeParams, + reconciledParams?: Map ): Promise { - const metadata = await getMetadata(); - const kiloSessionId = metadata?.auth.kiloSessionId; - const codeReviewSession = metadata?.identity.createdOnPlatform === 'code-review'; - let reconciledCount = 0; + const reconciled = reconciledParams ?? (await reconcileAcceptedMessages(acceptedMessages)); for (const message of acceptedMessages) { - const reconciled = - metadata && kiloSessionId - ? projectWrapperDeathReconciliation( - getAssistantMessageForUserMessage( - metadata.identity.sessionId, - kiloSessionId, - message.messageId - ), - codeReviewSession - ) - : null; - if (reconciled) reconciledCount += 1; await messageSettlementOutbox.terminalizeSessionMessageOnce( message.messageId, - reconciled ?? fallbackParams(message) + reconciled.get(message.messageId) ?? fallbackParams(message) ); } - if (reconciledCount > 0) { + if (reconciled.size > 0) { logger .withFields({ sessionId: getSessionIdForLogs(), - reconciledCount, - fallbackCount: acceptedMessages.length - reconciledCount, + reconciledCount: reconciled.size, + fallbackCount: acceptedMessages.filter(message => !reconciled.has(message.messageId)) + .length, }) .warn('Settled accepted wrapper work from stored assistant events after wrapper death'); } @@ -1002,17 +1025,66 @@ export function createWrapperSupervisor( await requestPhysicalWrapperStop('unhealthy-wrapper'); const acceptedMessages = await listNonTerminalAcceptedMessages(storage, state.wrapperRunId); - await terminalizeAcceptedMessagesForDeadWrapper(acceptedMessages, message => { - const activityObserved = message.agentActivityObservedAt !== undefined; - return { - kind: 'failed', - reason: 'wrapper_failure', - error, - completionSource: 'wrapper_failure', - failureStage: activityObserved ? 'agent_activity' : 'post_dispatch_no_activity', - failureCode, - }; - }); + // Recovery is gated strictly on the no-output code: a wrapper that stopped + // answering pings is gone, and re-dispatching its in-flight turn would only + // duplicate work. A no-output silence is the harness stalling on an accepted + // turn, so re-dispatch it once on a fresh runtime before failing the user's + // message. See s1 for the matching control-plane producer. + const recoveredMessageIds = new Set(); + let reconciledParams: Map | undefined; + if (failureCode === 'wrapper_no_output') { + // A silent wrapper can still have delivered the answer over the ingest + // channel before it stopped reporting. Reconcile those turns against the + // DO's stored assistant events first and settle them from that evidence: + // recovering one would run the turn a second time. + reconciledParams = await reconcileAcceptedMessages(acceptedMessages); + for (const message of acceptedMessages) { + if (reconciledParams.has(message.messageId)) continue; + if (!noOutputRecoveryAllowed(message.recoveryAttempts ?? 0)) continue; + const recovered = await requeueAcceptedMessageForRecovery(storage, message); + if (!recovered) continue; + recoveredMessageIds.add(message.messageId); + logger + .withFields({ + sessionId: getSessionIdForLogs(), + wrapperRunId: state.wrapperRunId, + messageId: message.messageId, + producer: 'wrapper_no_output_watchdog', + recoveryAttempts: (message.recoveryAttempts ?? 0) + 1, + logTag: 'wrapper_no_output_recovery', + }) + .info('Wrapper liveness no-output recovery re-queued the accepted turn'); + } + } + const terminalMessages = + recoveredMessageIds.size === 0 + ? acceptedMessages + : acceptedMessages.filter(message => !recoveredMessageIds.has(message.messageId)); + await terminalizeAcceptedMessagesForDeadWrapper( + terminalMessages, + message => { + const activityObserved = message.agentActivityObservedAt !== undefined; + return { + kind: 'failed', + reason: 'wrapper_failure', + error, + completionSource: 'wrapper_failure', + failureStage: activityObserved ? 'agent_activity' : 'post_dispatch_no_activity', + failureCode, + // Record the attempt count only once a recovery was actually spent on + // this turn. A first detection that could not be re-queued (no + // resolvable intent) spent nothing, and any non-null `attempts` maps + // to retry exhaustion on the client even though no retry ran. + ...(failureCode === 'wrapper_no_output' && message.recoveryAttempts !== undefined + ? { attempts: message.recoveryAttempts + 1 } + : {}), + }; + }, + reconciledParams + ); + if (recoveredMessageIds.size > 0) { + await sessionMessageQueue.requestPendingDrainIfNeeded(); + } await messageSettlementOutbox.releaseWrapperTerminalWaitForIdleBatch(); if (isWrapperRunFinalizing(state) && state.wrapperRunId) { await messageSettlementOutbox.finalizeTerminalWrapperRunCallbackIfReady(state.wrapperRunId); diff --git a/services/cloud-agent-next/test/integration/sandbox-session-no-output-recovery.test.ts b/services/cloud-agent-next/test/integration/sandbox-session-no-output-recovery.test.ts new file mode 100644 index 0000000000..b08c00dc3e --- /dev/null +++ b/services/cloud-agent-next/test/integration/sandbox-session-no-output-recovery.test.ts @@ -0,0 +1,149 @@ +import { env, runInDurableObject } from 'cloudflare:test'; +import { drizzle } from 'drizzle-orm/durable-sqlite'; +import { describe, expect, it, vi } from 'vitest'; +import { DEADLINE_MS } from '../../src/sandbox-control/deadlines'; +import { events } from '../../src/db/sqlite-schema'; +import type { SandboxSession } from '../../src/sandbox-session/SandboxSession'; +import type { + SessionMessageRecord, + SessionOperationProof, +} from '../../src/sandbox-session/session-message-queue'; +import type { SessionOperationAuthorization } from '../../src/shared/sandbox-control-protocol'; + +const rootKiloSessionId = 'ses_00000000000000000000000009'; +const directory = '/workspace/shared'; +const messageId = 'msg_no_output'; + +/** + * Producer 1 of 2: the accepted-message inactivity timeout + * (`SandboxSession.failOverdueAcceptedMessage`). This starves an accepted turn + * of activity and drives the real DO alarm twice through Miniflare: the first + * detection re-dispatches the same durable turn, the second terminalizes with + * the attempt count recorded. Producer 2 (the wrapper no-output watchdog) is + * owned by the sibling slice. + */ +describe('sandbox session no-output recovery', () => { + it('re-dispatches an accepted no-output turn once, then terminalizes with the attempt count', async () => { + const stub = env.SANDBOX_SESSION.getByName(`user_no_output:workspace_${crypto.randomUUID()}`); + await runInDurableObject(stub, async (instance: SandboxSession, state) => { + const sessionId = instance['requireSessionId'](); + await instance.registerSession({ + identity: { sessionId, userId: 'user_no_output' }, + auth: { kiloSessionId: rootKiloSessionId, kilocodeToken: 'fixture-token' }, + agent: { mode: 'code', model: 'test-model' }, + repository: { type: 'github', repo: 'Kilo-Org/cloud' }, + workspace: { sandboxId: 'usr-abcdef123419', workspacePath: directory }, + }); + const wrapperInstanceId = crypto.randomUUID(); + const authorization: SessionOperationAuthorization = { + operation: 'session.prompt', + operationId: messageId, + messageId, + session: { sessionId, kiloSessionId: rootKiloSessionId, directory }, + wrapperInstanceId, + dispatchDeadlineAt: Date.now() + 60_000, + }; + const promptProof = (): SessionOperationProof => ({ authorization, dispatched: true }); + const staleActivityAt = () => Date.now() - DEADLINE_MS.kiloInactivity - 1; + const record = (recoveryAttempts?: number): SessionMessageRecord => ({ + version: 2, + messageId, + state: 'accepted', + acceptedAt: staleActivityAt(), + lastActivityAt: staleActivityAt(), + wrapperInstanceId, + operations: { prompt: promptProof() }, + ...(recoveryAttempts !== undefined ? { recoveryAttempts } : {}), + intent: { + turn: { type: 'prompt', messageId, prompt: 'keep the typed message' }, + agent: { mode: 'code', model: 'test-model' }, + }, + }); + + const control = { + getStatus: vi.fn(async () => ({ + connection: 'ready', + physical: 'running', + wrapperInstanceId, + })), + request: vi.fn(async () => ({ + type: 'response', + requestId: 'no_output_recovery', + ok: true, + result: { status: 'aborted' }, + })), + quarantineRuntime: vi.fn(async () => ({ + quarantined: true, + disposition: 'native_retired', + })), + }; + const originalEnv = instance['env']; + Object.assign(instance, { + env: { ...originalEnv, SANDBOX_CONTROL: { getByName: () => control } }, + // The operation receipt is liveness, not progress: hold it at `running` + // so the alarm takes the accepted-operation branch and reaches the + // inactivity bound. + observeAcceptedOperation: vi.fn(async () => 'running' as const), + }); + const messages = () => state.storage.kv.get('session_messages') ?? []; + const failedEvents = () => + drizzle(state.storage) + .select() + .from(events) + .all() + .filter(event => event.stream_event_type === 'cloud.message.failed'); + try { + // First detection: the accepted turn has produced no activity past the bound. + state.storage.kv.put('session_messages', [record()]); + await instance.alarm(); + await state.storage.deleteAlarm(); + + const recovered = messages()[0]; + if (!recovered) throw new Error('Missing recovered record'); + expect(recovered).toMatchObject({ + messageId, + state: 'queued', + recoveryAttempts: 1, + failedReason: undefined, + }); + // The typed turn survives under its durable identity, while the + // ambiguous dispatch state is dropped so a fresh runtime can take it. + expect(recovered.intent).toEqual({ + turn: { type: 'prompt', messageId, prompt: 'keep the typed message' }, + agent: { mode: 'code', model: 'test-model' }, + }); + expect(recovered.wrapperInstanceId).toBeUndefined(); + expect(recovered.operations).toBeUndefined(); + expect(recovered.acceptedAt).toBeUndefined(); + expect(recovered.lastActivityAt).toBeUndefined(); + expect(failedEvents()).toHaveLength(0); + + // Second identical detection on the replacement runtime: terminal. + state.storage.kv.put('session_messages', [record(1)]); + await instance.alarm(); + await state.storage.deleteAlarm(); + + const terminal = messages()[0]; + if (!terminal) throw new Error('Missing terminal record'); + expect(terminal).toMatchObject({ + messageId, + state: 'failed', + failedReason: 'accepted_overdue', + failedDetail: 'Turn did not complete', + recoveryAttempts: 1, + }); + const persistedFailed = failedEvents(); + expect(persistedFailed).toHaveLength(1); + expect(JSON.parse(persistedFailed[0]?.payload ?? '{}')).toMatchObject({ + messageId, + status: 'failed', + reason: 'accepted_overdue', + attempts: 2, + }); + } finally { + Object.assign(instance, { env: originalEnv }); + await state.storage.deleteAlarm(); + } + }); + }); +}); diff --git a/services/cloud-agent-next/test/integration/session/execute-directly-failure.test.ts b/services/cloud-agent-next/test/integration/session/execute-directly-failure.test.ts index b1f9322ddc..8100aba226 100644 --- a/services/cloud-agent-next/test/integration/session/execute-directly-failure.test.ts +++ b/services/cloud-agent-next/test/integration/session/execute-directly-failure.test.ts @@ -10,8 +10,14 @@ import { env, runInDurableObject, listDurableObjectIds } from 'cloudflare:test'; import { afterEach, beforeEach, describe, it, expect } from 'vitest'; import { drizzle } from 'drizzle-orm/durable-sqlite'; import { createEventQueries } from '../../../src/session/queries/events.js'; -import type { FencedWrapperDispatchRequest } from '../../../src/execution/types.js'; -import { listPendingSessionMessages } from '../../../src/session/pending-messages.js'; +import type { + FencedWrapperDispatchRequest, + SessionMessageIntent, +} from '../../../src/execution/types.js'; +import { + deletePendingSessionMessageByMessageId, + listPendingSessionMessages, +} from '../../../src/session/pending-messages.js'; import { getWrapperLease, getWrapperRuntimeState, @@ -528,10 +534,11 @@ describe('new-path liveness without executionId', () => { ); }); - it('schedules liveness deadlines for accepted messages and fails them on no-output timeout', async () => { + it('recovers an accepted turn once on no-output timeout, then fails the second identical detection', async () => { const userId = 'user_newpath_liveness'; const sessionId = 'agent_newpath_liveness'; const stub = sessionStub(userId, sessionId); + const messageId = 'msg_018f1e2d3c4bnewlivabcdefgh'; const result = await runInDurableObject(stub, async (instance, state) => { await registerReadySession(instance, { @@ -549,26 +556,51 @@ describe('new-path liveness without executionId', () => { const { state: wrapperState } = await allocateWrapperRuntimeState(instance.ctx.storage); const { wrapperRunId, wrapperConnectionId } = wrapperState; - // Store an accepted (non-terminal) session message state + // Store an accepted (non-terminal) session message state carrying the + // immutable admission snapshot the recovery re-dispatches. + const admissionSnapshot: SessionMessageIntent = { + turn: { type: 'prompt', messageId, prompt: 'hello' }, + agent: { mode: 'code', model: 'test-model' }, + }; const acceptedMessage: SessionMessageState = { - messageId: 'msg_018f1e2d3c4bnewlivabcdefgh', + messageId, status: 'accepted', prompt: 'hello', createdAt: Date.now(), acceptedAt: Date.now(), wrapperRunId: wrapperRunId!, + admissionSnapshot, }; await putSessionMessageState(instance.ctx.storage, acceptedMessage); // Set expired liveness deadlines — new path has no executionId const expiredAt = Date.now() - 1; - await instance.ctx.storage.put('wrapper_runtime_state', { + const expiredRuntimeState = { wrapperGeneration: wrapperState.wrapperGeneration, wrapperConnectionId, wrapperRunId, noOutputDeadlineAt: expiredAt, lastHeartbeatUpdate: expiredAt - 10 * 60_000, + }; + await instance.ctx.storage.put('wrapper_runtime_state', expiredRuntimeState); + + await instance.alarm(); + + const recoveredMessage = await getSessionMessageState(instance.ctx.storage, messageId); + const pendingAfterRecovery = await listPendingSessionMessages(instance.ctx.storage); + const db = drizzle(state.storage, { logger: false }); + const eventQueries = createEventQueries(db, state.storage.sql); + const eventsAfterRecovery = eventQueries.findByFilters({}); + + // Second identical detection: the recovery budget is spent, so the turn + // must take the existing terminal path and record the attempt count. + await instance.ctx.storage.delete('wrapper_lease'); + await deletePendingSessionMessageByMessageId(instance.ctx.storage, messageId); + await putSessionMessageState(instance.ctx.storage, { + ...acceptedMessage, + recoveryAttempts: 1, }); + await instance.ctx.storage.put('wrapper_runtime_state', expiredRuntimeState); await instance.alarm(); @@ -576,28 +608,48 @@ describe('new-path liveness without executionId', () => { instance.ctx.storage, wrapperRunId! ); - const db = drizzle(state.storage, { logger: false }); - const eventQueries = createEventQueries(db, state.storage.sql); + const terminalMessage = await getSessionMessageState(instance.ctx.storage, messageId); const allEvents = eventQueries.findByFilters({}); return { + recoveredMessage, + pendingAfterRecovery, + eventsAfterRecovery, nonTerminalMessages, + terminalMessage, allEvents, wrapperRuntimeState: await getWrapperRuntimeState(instance.ctx.storage), }; }); - // Message must be terminalized as failed + // First detection: the accepted turn is recovered, not failed. + expect(result.recoveredMessage).toMatchObject({ + status: 'queued', + recoveryAttempts: 1, + }); + expect(result.recoveredMessage?.acceptedAt).toBeUndefined(); + expect(result.recoveredMessage?.wrapperRunId).toBeUndefined(); + expect(result.pendingAfterRecovery.map(message => message.messageId)).toEqual([messageId]); + expect( + result.eventsAfterRecovery.filter(event => event.stream_event_type === 'cloud.message.failed') + ).toHaveLength(0); + + // Second detection: terminal, with the attempt count recorded. expect(result.nonTerminalMessages).toHaveLength(0); + expect(result.terminalMessage).toMatchObject({ + status: 'failed', + attempts: 2, + failureCode: 'wrapper_no_output', + }); - // A cloud.message.failed event must be persisted const failedEvents = result.allEvents.filter( event => event.stream_event_type === 'cloud.message.failed' ); expect(failedEvents).toHaveLength(1); const failedPayload = JSON.parse(failedEvents[0].payload); expect(failedPayload).toMatchObject({ - messageId: 'msg_018f1e2d3c4bnewlivabcdefgh', + messageId, status: 'failed', + attempts: 2, error: 'Agent wrapper made no execution progress during the watchdog window', delivery: 'sent', accepted: true,