From 2d0adb4ded12c05772fb409df704751b813162a7 Mon Sep 17 00:00:00 2001 From: Marc Liu Date: Sun, 30 Aug 2026 01:20:40 -0400 Subject: [PATCH] fix(daemon): complete message delivery at consumption and stop re-injecting consumed rows on model switch (#1686) Delivery jobs are delivery-only: non-ACP turns complete at the SDK admission ack and ACP turns at markMessageAccepted, instead of waiting for turn end and reclassifying. An interrupted turn no longer reopens a consumed row, which was the re-injection retry loop. restart()/reset() now cancel in-flight message_delivery jobs, durably requeue unconsumed yielded prompts, and re-feed the last consumed message only when the replacement query does not preserve the provider transcript (no sdkSessionId / acpSessionId). ACP prompt-loop breaks acknowledge already-consumed prompts instead of requeueing them. --- .../daemon/src/lib/acp/acp-query-runner.ts | 12 +- .../daemon/src/lib/agent/agent-session.ts | 273 +++++++++--------- .../daemon/src/lib/agent/message-queue.ts | 10 + .../src/lib/agent/query-lifecycle-manager.ts | 98 ++++++- packages/daemon/src/lib/agent/query-runner.ts | 3 + .../unit/1-core/agent/agent-session.test.ts | 111 ++++--- ...elivery-retry-livelock-convergence.test.ts | 20 +- ...essage-delivery-transcript-harness.test.ts | 5 +- .../unit/1-core/agent/message-queue.test.ts | 64 ++++ .../agent/query-lifecycle-manager.test.ts | 111 +++++++ 10 files changed, 495 insertions(+), 212 deletions(-) diff --git a/packages/daemon/src/lib/acp/acp-query-runner.ts b/packages/daemon/src/lib/acp/acp-query-runner.ts index 0608a04887..44e60a51b0 100644 --- a/packages/daemon/src/lib/acp/acp-query-runner.ts +++ b/packages/daemon/src/lib/acp/acp-query-runner.ts @@ -430,7 +430,17 @@ export class AcpQueryRunner { attemptOwnsRun() && !runAbortController.signal.aborted && !this.ctx.isCleaningUp(); const requeueYieldedPrompt = (yieldedMessage: SDKUserMessage) => { const yieldedUuid = yieldedMessage.uuid; - if (yieldedUuid && !messageQueue.requeueYielded(yieldedUuid)) { + if (!yieldedUuid) return; + if (this.ctx.session.acpSessionId) { + const status = this.ctx.db + .getSDKMessageRepo?.() + ?.getDeliveryContent(this.ctx.session.id, yieldedUuid)?.sendStatus; + if (status === 'consumed') { + messageQueue.acknowledgeYielded(yieldedUuid); + return; + } + } + if (!messageQueue.requeueYielded(yieldedUuid)) { logger.warn(`ACP prompt loop: could not requeue yielded prompt ${yieldedUuid}.`); } }; diff --git a/packages/daemon/src/lib/agent/agent-session.ts b/packages/daemon/src/lib/agent/agent-session.ts index 12c5864398..3aaccafb04 100644 --- a/packages/daemon/src/lib/agent/agent-session.ts +++ b/packages/daemon/src/lib/agent/agent-session.ts @@ -186,7 +186,6 @@ import { MessageDeliveryTerminalTurnError, signalDeliveryConsumed, steerAckTimeoutMs, - throwIfDeliveryAborted, waitForDeliveryAbort, withSessionLock, } from './message-delivery.ts'; @@ -195,7 +194,6 @@ import { classifyTurnCompletion, decideReconcileAdmission, selectStrandedDeliveries, - shouldRearmSpuriousTurnEnd, } from './message-delivery-pipeline.ts'; import type { MidTurnBudgetInterruptOptions } from './message-queue.ts'; import { MessageQueue } from './message-queue.ts'; @@ -326,6 +324,10 @@ export class AgentSession observer: MessageDeliveryAttemptObserver; pendingStart?: boolean; } | null = null; + private acpTurnAccepted: { + promise: Promise; + resolve: () => void; + } | null = null; private taskNotificationRequeryAttempts = 0; private taskNotificationRequeryExhausted = false; @@ -1683,7 +1685,7 @@ export class AgentSession onResumeClear: () => { this.pendingResumeAfterCompaction = false; }, - onSurvivorRequeued: (uuid) => this.reopenDeliveryForRetry(uuid), + onSurvivorRequeued: (uuid) => this.reopenDeliveryForRetry([uuid]), getDurableMessageContent: (uuid) => { const repo = this.db.getSDKMessageRepo(); const kickoff = repo.getUserMessageContentByUuid(this.session.id, uuid); @@ -2488,33 +2490,21 @@ export class AgentSession if (started.kind === 'aborted') { return { outcome: 'aborted' }; } - if ( - alreadyConsumed && - !this.rateLimitWatchdog.isRecoveryPending() && - this.db.getSDKMessageRepo()?.hasRecoveryInterceptedResultAfter?.(this.session.id, messageUuid) - ) { - this.logger.info( - `delivery-turn: reclaiming recovery-intercepted row directly (uuid=${messageUuid}); ` + - `the parked recovery episode no longer owns its retry.` + if (alreadyConsumed) { + this.logger.debug( + `delivery-turn: message already accepted (uuid=${messageUuid}); delivery job is complete` ); - started.turnEnd.cancel(); - if (started.responseObserver && this.deliveryResponseObserver === started.responseObserver) { - this.deliveryResponseObserver = null; - } - if (!claimGuard || claimGuard()) { - this.reopenDeliveryForRetry(messageUuid); - } - throw new MessageDeliveryRecoverableTurnError('Turn ended without a response'); + return { outcome: 'completed' }; } this.deliveryTurnStalled = false; this.outstandingToolUseIds.clear(); - let stallPromise: Promise = new Promise(() => {}); let stallWatchdog: DeliveryTurnStallWatchdog | null = null; let activeTurnEnd = started.turnEnd; const responseObserver = started.responseObserver; let kickoffAcknowledged = false; let kickoffDiedBeforeConsumption = false; let kickoffAckInvalidated = false; + const turnUuids = [messageUuid, ...(started.admittedBatchUuids ?? [])]; try { if (started.acknowledgment) { const aborted = waitForDeliveryAbort(signal); @@ -2543,10 +2533,10 @@ export class AgentSession this.logger.warn( `delivery-turn: query ended before the SDK consumed the kickoff ` + `(uuid=${messageUuid}, generation=${started.generation}); requeueing the ` + - `kickoff, reopening it for retry, and classifying the turn outcome` + `kickoff and reopening any accepted batch members for retry` ); this.messageQueue.requeueYielded(messageUuid); - this.reopenDeliveryForRetry(messageUuid); + this.reopenDeliveryForRetry(turnUuids); kickoffDiedBeforeConsumption = true; } const kickoffStatus = this.stateManager.getState().status; @@ -2580,69 +2570,98 @@ export class AgentSession signalDeliveryConsumed(this.session.id, memberUuid); } } + this.zeroProgressDeliveryFailures = null; + return { outcome: 'completed' }; + } + this.armDeliveryTurnStall(signal, claimGuard); + stallWatchdog = this.deliveryTurnStall; + const acpAborted = waitForDeliveryAbort(signal); + let acpAcceptanceWon = false; + try { + await Promise.race([ + this.waitForAcpTurnAccepted().then(() => { + acpAcceptanceWon = true; + }), + started.queryPromise.catch(() => {}), + acpAborted.promise, + ]); + } finally { + acpAborted.cancel(); + } + if (acpAcceptanceWon) { + this.logger.debug( + `delivery-turn: ACP prompt accepted ` + + `(${Date.now() - turnStartedAt}ms since turn start, uuid=${messageUuid})` + ); + this.zeroProgressDeliveryFailures = null; + return { outcome: 'completed' }; } + this.logger.warn( + `delivery-turn: ACP query ended before the prompt was accepted ` + + `(uuid=${messageUuid}, generation=${started.generation}); requeueing the kickoff ` + + `and reopening any accepted batch members for retry` + ); + this.messageQueue.requeueYielded(messageUuid); + this.reopenDeliveryForRetry(turnUuids); + throw new MessageDeliveryRecoverableTurnError( + 'ACP prompt was not accepted before query end' + ); } } - throwIfDeliveryAborted(signal); - stallPromise = this.armDeliveryTurnStall(signal, claimGuard); - stallWatchdog = this.deliveryTurnStall; - const SPURIOUS_TURN_END_GRACE_MS = 250; - const feedAcknowledged = started.acknowledgment !== null; - let raceArmedAt = Date.now(); - let graceRearms = 0; - let turnEndFired = false; - let queryEnded = false; - void activeTurnEnd.promise.then(() => { - turnEndFired = true; - }); - void started.queryPromise - .catch(() => {}) - .then(() => { - queryEnded = true; - }); - while (true) { - const aborted = waitForDeliveryAbort(signal); - try { - await Promise.race([ - activeTurnEnd.promise, - started.queryPromise.catch(() => {}), - stallPromise, - aborted.promise, - ]); - } finally { - aborted.cancel(); - } - const turnResultRepo = this.db.getSDKMessageRepo(); - const hasAnyTerminalResult = - !!turnResultRepo?.hasTerminalResultAfter(this.session.id, messageUuid) || - !!turnResultRepo?.getErrorTerminalResultSubtypeAfter(this.session.id, messageUuid); - const spuriousFire = shouldRearmSpuriousTurnEnd({ - feedAcknowledged, - turnEndFired, - queryEnded, - withinGraceMs: Date.now() - raceArmedAt <= SPURIOUS_TURN_END_GRACE_MS, - graceRearms, - hasTerminalResult: hasAnyTerminalResult, - }); - if (!spuriousFire) break; - graceRearms++; - activeTurnEnd.cancel(); - const rearmedTurnEnd = this.stateManager.waitForIdleTransition( - this.rateLimitWatchdog.getGeneration(), - recordTurnEndMarker, - started.idleOwner + if (kickoffAckInvalidated) { + this.logger.warn( + `delivery-turn: kickoff acknowledgment was invalidated ` + + `(uuid=${messageUuid}, generation=${started.generation}); classifying without reopen` ); - activeTurnEnd = { - promise: rearmedTurnEnd.promise, - cancel: rearmedTurnEnd.cancel, - idleOwner: started.idleOwner, - }; - raceArmedAt = Date.now(); - turnEndFired = false; - void activeTurnEnd.promise.then(() => { - turnEndFired = true; + } + const producedResult = !!this.db + .getSDKMessageRepo() + ?.hasTerminalResultAfter(this.session.id, messageUuid); + if (!producedResult) { + if (this.rateLimitWatchdog.isRecoveryPending()) { + const cooldownRetryAt = this.rateLimitWatchdog.getState().retryAt; + const retryAt = this.rateLimitWatchdog.isManualRecoveryPause() + ? Date.now() + MANUAL_RECOVERY_PARK_MS + : Math.max(Date.now() + MESSAGE_DELIVERY_PARK_MS, cooldownRetryAt ?? 0); + this.db.getSDKMessageRepo()?.clearDeliveryTurnEnd(this.session.id, messageUuid); + this.logger.info( + `delivery-turn: parking job while limit recovery is pending ` + + `(uuid=${messageUuid}, retryAt=${new Date(retryAt).toISOString()})` + ); + return { outcome: 'recovery_pending', retryAt }; + } + const turnError = this.consumeTerminalTurnError(turnStartedAt); + this.db.getSDKMessageRepo()?.clearDeliveryTurnEnd(this.session.id, messageUuid); + const errorResultSubtype = this.db + .getSDKMessageRepo() + ?.getErrorTerminalResultSubtypeAfter(this.session.id, messageUuid); + const completion = classifyTurnCompletion({ + producedResult, + turnError, + errorResultSubtype, + deliveryTurnStalled: this.deliveryTurnStalled, + claimGuardHeld: claimGuard ? claimGuard() : undefined, }); + if (completion.outcome === 'terminal_error') { + throw new MessageDeliveryTerminalTurnError(completion.detail, completion.category); + } + if (completion.outcome === 'recoverable_error') { + if (!kickoffAcknowledged && !kickoffAckInvalidated && !alreadyConsumed) { + const terminal = await this.escalateZeroProgressDeliveryFailure(messageUuid); + if (terminal) throw terminal; + } + if ( + completion.reopenForRetry && + !kickoffDiedBeforeConsumption && + !kickoffAckInvalidated + ) { + this.reopenDeliveryForRetry(turnUuids); + } + throw new MessageDeliveryRecoverableTurnError(completion.detail, completion.category); + } } + this.zeroProgressDeliveryFailures = null; + return { outcome: 'completed' }; } finally { activeTurnEnd.cancel(); if (this.deliveryTurnStall === stallWatchdog) { @@ -2651,51 +2670,8 @@ export class AgentSession if (responseObserver && this.deliveryResponseObserver === responseObserver) { this.deliveryResponseObserver = null; } + this.acpTurnAccepted = null; } - const producedResult = !!this.db - .getSDKMessageRepo() - ?.hasTerminalResultAfter(this.session.id, messageUuid); - if (!producedResult) { - if (this.rateLimitWatchdog.isRecoveryPending()) { - const cooldownRetryAt = this.rateLimitWatchdog.getState().retryAt; - const retryAt = this.rateLimitWatchdog.isManualRecoveryPause() - ? Date.now() + MANUAL_RECOVERY_PARK_MS - : Math.max(Date.now() + MESSAGE_DELIVERY_PARK_MS, cooldownRetryAt ?? 0); - this.db.getSDKMessageRepo()?.clearDeliveryTurnEnd(this.session.id, messageUuid); - this.logger.info( - `delivery-turn: parking job while limit recovery is pending ` + - `(uuid=${messageUuid}, retryAt=${new Date(retryAt).toISOString()})` - ); - return { outcome: 'recovery_pending', retryAt }; - } - const turnError = this.consumeTerminalTurnError(turnStartedAt); - this.db.getSDKMessageRepo()?.clearDeliveryTurnEnd(this.session.id, messageUuid); - const errorResultSubtype = this.db - .getSDKMessageRepo() - ?.getErrorTerminalResultSubtypeAfter(this.session.id, messageUuid); - const completion = classifyTurnCompletion({ - producedResult, - turnError, - errorResultSubtype, - deliveryTurnStalled: this.deliveryTurnStalled, - claimGuardHeld: claimGuard ? claimGuard() : undefined, - }); - if (completion.outcome === 'terminal_error') { - throw new MessageDeliveryTerminalTurnError(completion.detail, completion.category); - } - if (completion.outcome === 'recoverable_error') { - if (!kickoffAcknowledged && !kickoffAckInvalidated && !alreadyConsumed) { - const terminal = await this.escalateZeroProgressDeliveryFailure(messageUuid); - if (terminal) throw terminal; - } - if (completion.reopenForRetry) { - this.reopenDeliveryForRetry(messageUuid); - } - throw new MessageDeliveryRecoverableTurnError(completion.detail, completion.category); - } - } - this.zeroProgressDeliveryFailures = null; - return { outcome: 'completed' }; } private buildDeliveryTurnAdmissionDeps( @@ -2855,8 +2831,21 @@ export class AgentSession } onDeliveryTurnAccepted(): void { - if (!this.isAcpSession() || !this.deliveryTurnStall) return; - this.deliveryTurnStall.resizeTimeoutMs(DELIVERY_TURN_NO_ACTIVITY_MS); + if (!this.isAcpSession()) return; + if (this.deliveryTurnStall) { + this.deliveryTurnStall.resizeTimeoutMs(DELIVERY_TURN_NO_ACTIVITY_MS); + } + this.acpTurnAccepted?.resolve(); + } + + private waitForAcpTurnAccepted(): Promise { + if (this.acpTurnAccepted) return this.acpTurnAccepted.promise; + let resolve!: () => void; + const promise = new Promise((r) => { + resolve = r; + }); + this.acpTurnAccepted = { promise, resolve }; + return promise; } private isAcpSession(): boolean { @@ -3067,7 +3056,7 @@ export class AgentSession } if (steerWinner === 'query_ended') { this.messageQueue.requeueYielded(messageUuid); - this.reopenDeliveryForRetry(messageUuid); + this.reopenDeliveryForRetry([messageUuid]); throw new Error('Steer target query ended before the SDK consumed the steer'); } if (steerWinner === 'ack_timeout') { @@ -3088,7 +3077,7 @@ export class AgentSession if (!this.messageQueue.remove(messageUuid)) { this.messageQueue.acknowledgeYielded(messageUuid); } - this.reopenDeliveryForRetry(messageUuid); + this.reopenDeliveryForRetry([messageUuid]); return { outcome: 'ack_timeout' }; } } @@ -3100,7 +3089,7 @@ export class AgentSession this.stateManager.isTerminalIdlePending() || (this.messageQueue.getClearEpoch?.() ?? 0) !== action.clearEpoch ) { - this.reopenDeliveryForRetry(messageUuid); + this.reopenDeliveryForRetry([messageUuid]); throw new Error('Steer was invalidated by session teardown before the SDK consumed it'); } deliveryMetrics.recordFeed(messageUuid); @@ -3335,19 +3324,21 @@ export class AgentSession return { content: buildBatchedDeliveryContent(texts), admittedUuids: admitted }; } - private reopenDeliveryForRetry(messageUuid: string): void { - deliveryMetrics.forgetFeed(messageUuid); - const dbId = this.db - .getSDKMessageRepo() - ?.markDeliveryRetryableByUuid(this.session.id, messageUuid); - if (dbId) { - void this.internalEventBus - .publish('messages.statusChanged', { - sessionId: this.session.id, - messageIds: [dbId], - status: 'enqueued', - }) - .catch(() => {}); + private reopenDeliveryForRetry(messageUuids: string[]): void { + const repo = this.db.getSDKMessageRepo(); + if (!repo) return; + for (const messageUuid of messageUuids) { + deliveryMetrics.forgetFeed(messageUuid); + const dbId = repo.markDeliveryRetryableByUuid(this.session.id, messageUuid); + if (dbId) { + void this.internalEventBus + .publish('messages.statusChanged', { + sessionId: this.session.id, + messageIds: [dbId], + status: 'enqueued', + }) + .catch(() => {}); + } } } diff --git a/packages/daemon/src/lib/agent/message-queue.ts b/packages/daemon/src/lib/agent/message-queue.ts index 700bb69ab1..3be7e107c2 100644 --- a/packages/daemon/src/lib/agent/message-queue.ts +++ b/packages/daemon/src/lib/agent/message-queue.ts @@ -597,6 +597,16 @@ export class MessageQueue { return false; } + requeueAllYielded(options?: { durable?: boolean }): string[] { + const requeued: string[] = []; + for (const message of Array.from(this.yielded)) { + if (this.requeueYielded(message.id, options)) { + requeued.push(message.id); + } + } + return requeued; + } + waitForPendingOrInFlight( messageId: string ): { acknowledgment: Promise; content: string | MessageContent[] } | null { diff --git a/packages/daemon/src/lib/agent/query-lifecycle-manager.ts b/packages/daemon/src/lib/agent/query-lifecycle-manager.ts index cf3217d2a7..3e946f8f31 100644 --- a/packages/daemon/src/lib/agent/query-lifecycle-manager.ts +++ b/packages/daemon/src/lib/agent/query-lifecycle-manager.ts @@ -13,7 +13,12 @@ import type { QueryAttemptRegistry } from './query-attempt-token.ts'; import { Logger } from '../logger.ts'; import { existsSync, copyFileSync, mkdirSync } from 'node:fs'; import { dirname } from 'node:path'; -import { throwIfDeliveryAborted, waitForDeliveryAbort } from './message-delivery.ts'; +import { + buildBatchedDeliveryContent, + flattenDeliveryText, + throwIfDeliveryAborted, + waitForDeliveryAbort, +} from './message-delivery.ts'; import { validateAndRepairSDKSession, findSDKSessionFileGlobally, @@ -104,6 +109,68 @@ export class QueryLifecycleManager { } } + private isTranscriptPreservedForProvider(): boolean { + const { session } = this.ctx; + if (session.config.provider === 'acp') { + return !!session.acpSessionId; + } + return !!session.sdkSessionId; + } + + private resolveLastConsumedUuid(): string | null { + const { db, session } = this.ctx; + const repo = db.getSDKMessageRepo?.(); + if (!repo) return null; + const { messages } = repo.getUserMessagesByStatus(session.id, 'consumed', 1, 'desc'); + const last = messages[0]; + return last?.uuid ?? null; + } + + private reenqueueLastConsumedIfTranscriptDropped( + lastConsumedUuid: string | null, + lastConsumedContent: string | MessageContent[] | null + ): void { + if (!lastConsumedUuid || !lastConsumedContent) return; + if (this.isTranscriptPreservedForProvider()) return; + const { messageQueue, session } = this.ctx; + if (messageQueue.hasPendingOrInFlight(lastConsumedUuid)) return; + this.logger.info( + `re-enqueueing last consumed message ${lastConsumedUuid} for session ${session.id} ` + + `because the replacement query does not preserve the provider transcript` + ); + messageQueue + .enqueueWithId(lastConsumedUuid, lastConsumedContent, false, { + durable: true, + prepend: true, + }) + .catch((error) => { + this.logger.warn(`re-enqueue of last consumed message ${lastConsumedUuid} failed:`, error); + }); + } + + private buildRequeueContentFor(messageUuid: string): string | MessageContent[] | null { + const { db, session } = this.ctx; + const repo = db.getSDKMessageRepo?.(); + if (!repo) return null; + const content = repo.getUserMessageContentByUuid(session.id, messageUuid); + if (content === null) return null; + const jobQueue = db.getJobQueueRepo?.(); + const batchUuids = jobQueue?.getActiveDeliveryBatchUuids?.(session.id, messageUuid); + if (!batchUuids || batchUuids.length <= 1) return content; + const texts: string[] = []; + for (const uuid of batchUuids) { + const memberContent = + uuid === messageUuid ? content : repo.getUserMessageContentByUuid(session.id, uuid); + if (memberContent === null) continue; + const text = + typeof memberContent === 'string' ? memberContent : flattenDeliveryText(memberContent); + if (text === null) continue; + texts.push(text); + } + if (texts.length <= 1) return content; + return buildBatchedDeliveryContent(texts); + } + private getSDKWorkspacePath(): string { const { session } = this.ctx; return session.worktree @@ -291,6 +358,22 @@ export class QueryLifecycleManager { await this.stop(); + const lastConsumedUuid = this.resolveLastConsumedUuid(); + const lastConsumedContent = lastConsumedUuid + ? this.buildRequeueContentFor(lastConsumedUuid) + : null; + const jobQueue = this.ctx.db.getJobQueueRepo?.(); + if (jobQueue) { + const cancelledUuids = jobQueue.cancelForSessionWithMessages(session.id); + if (cancelledUuids.length > 0) { + this.logger.debug( + `cancelled ${cancelledUuids.length} in-flight message_delivery job(s) for session ${session.id}` + ); + } + } + this.ctx.messageQueue.requeueAllYielded({ durable: true }); + this.reenqueueLastConsumedIfTranscriptDropped(lastConsumedUuid, lastConsumedContent); + await this.ctx.stateManager.setIdle({ suppressDeliveryWaiters: true }); reachedSuppressedIdle = true; @@ -373,6 +456,19 @@ export class QueryLifecycleManager { catchQueryErrors: true, }); + const jobQueue = db.getJobQueueRepo?.(); + if (jobQueue) { + jobQueue.cancelForSession(session.id); + } + + if (restartAfter) { + const lastConsumedUuid = this.resolveLastConsumedUuid(); + const lastConsumedContent = lastConsumedUuid + ? this.buildRequeueContentFor(lastConsumedUuid) + : null; + this.reenqueueLastConsumedIfTranscriptDropped(lastConsumedUuid, lastConsumedContent); + } + this.ctx.firstMessageReceived = false; await stateManager.setIdle({ suppressDeliveryWaiters: restartAfter }); reachedSuppressedIdle = true; diff --git a/packages/daemon/src/lib/agent/query-runner.ts b/packages/daemon/src/lib/agent/query-runner.ts index 83c3512c89..e5c21748e2 100644 --- a/packages/daemon/src/lib/agent/query-runner.ts +++ b/packages/daemon/src/lib/agent/query-runner.ts @@ -1946,6 +1946,9 @@ export class QueryRunner { }; const requeueConsumedList = (state: RetryTeardownState): RetryTeardownState => { + if (this.ctx.getQueryGeneration() !== queryGeneration) { + return state; + } const consumed = this._consumedUserMessages.get(queryGeneration) ?? []; if (consumed.length > 0) { logger.warn( diff --git a/packages/daemon/tests/unit/1-core/agent/agent-session.test.ts b/packages/daemon/tests/unit/1-core/agent/agent-session.test.ts index a143024121..3712483dab 100644 --- a/packages/daemon/tests/unit/1-core/agent/agent-session.test.ts +++ b/packages/daemon/tests/unit/1-core/agent/agent-session.test.ts @@ -1011,7 +1011,7 @@ describe('AgentSession', () => { ); }); - it('driveDeliveryTurn reopens a no-result consumed row for retry only while the claim is current', async () => { + it('driveDeliveryTurn completes an already-consumed row without reopening or waiting for turn end', async () => { const retrySpy = mock(() => 'db-1'); mockDb.getSDKMessageRepo = mock(() => ({ getDeliveryContent: mock(() => ({ content: 'x', sendStatus: 'consumed' })), @@ -1031,12 +1031,15 @@ describe('AgentSession', () => { ); await agentSession.stateManager.setProcessing('uuid-reopen'); - const live = agentSession.driveDeliveryTurn('uuid-reopen', 'hello', null, true, () => true); - await agentSession.stateManager.setIdle(); - const error = await live.catch((caught) => caught); - expect(error).toBeInstanceOf(MessageDeliveryRecoverableTurnError); - expect(error).toMatchObject({ message: 'Turn ended without a response' }); - expect(retrySpy).toHaveBeenCalledTimes(1); + const result = await agentSession.driveDeliveryTurn( + 'uuid-reopen', + 'hello', + null, + true, + () => true + ); + expect(result).toEqual({ outcome: 'completed' }); + expect(retrySpy).not.toHaveBeenCalled(); let claimAlive = true; await agentSession.stateManager.setProcessing('uuid-reopen-2'); @@ -1050,8 +1053,8 @@ describe('AgentSession', () => { await new Promise((resolve) => setTimeout(resolve, 10)); claimAlive = false; await agentSession.stateManager.setIdle(); - await expect(cancelled).rejects.toThrow(); - expect(retrySpy).toHaveBeenCalledTimes(1); + expect(await cancelled).toEqual({ outcome: 'completed' }); + expect(retrySpy).not.toHaveBeenCalled(); }); it('an aborted delivery admission leaves the idle owner untouched (no phantom turn)', async () => { @@ -1084,7 +1087,7 @@ describe('AgentSession', () => { expect(admitSpy).not.toHaveBeenCalled(); }); - it('driveDeliveryTurn parks instead of reopening while limit recovery is pending', async () => { + it('driveDeliveryTurn completes an already-consumed row even when limit recovery is pending', async () => { const retrySpy = mock(() => 'db-1'); mockDb.getSDKMessageRepo = mock(() => ({ getDeliveryContent: mock(() => ({ content: 'x', sendStatus: 'consumed' })), @@ -1112,15 +1115,19 @@ describe('AgentSession', () => { watchdog.getState = mock(() => ({ retryAt: cooldownRetryAt })) as never; await agentSession.stateManager.setProcessing('uuid-park'); - const drive = agentSession.driveDeliveryTurn('uuid-park', 'hello', null, true, () => true); - await agentSession.stateManager.setIdle(); - const result = await drive; + const result = await agentSession.driveDeliveryTurn( + 'uuid-park', + 'hello', + null, + true, + () => true + ); - expect(result).toEqual({ outcome: 'recovery_pending', retryAt: cooldownRetryAt }); + expect(result).toEqual({ outcome: 'completed' }); expect(retrySpy).not.toHaveBeenCalled(); }); - it('driveDeliveryTurn parks a manual-only pause on a long horizon without short polling', async () => { + it('driveDeliveryTurn completes an already-consumed row even during a manual-only pause', async () => { const retrySpy = mock(() => 'db-1'); mockDb.getSDKMessageRepo = mock(() => ({ getDeliveryContent: mock(() => ({ content: 'x', sendStatus: 'consumed' })), @@ -1138,7 +1145,6 @@ describe('AgentSession', () => { (agentSession as unknown as { queryPromise: Promise }).queryPromise = new Promise( () => {} ); - const before = Date.now(); const watchdog = (agentSession as unknown as { rateLimitWatchdog: unknown }) .rateLimitWatchdog as { isRecoveryPending: () => boolean; @@ -1150,27 +1156,23 @@ describe('AgentSession', () => { watchdog.getState = mock(() => ({ retryAt: null })) as never; await agentSession.stateManager.setProcessing('uuid-manual-park'); - const drive = agentSession.driveDeliveryTurn( + const result = await agentSession.driveDeliveryTurn( 'uuid-manual-park', 'hello', null, true, () => true ); - await agentSession.stateManager.setIdle(); - const result = (await drive) as { outcome: string; retryAt: number }; - expect(result.outcome).toBe('recovery_pending'); - expect(result.retryAt).toBeGreaterThanOrEqual(before + 4 * 60_000); + expect(result).toEqual({ outcome: 'completed' }); expect(retrySpy).not.toHaveBeenCalled(); }); - it('driveDeliveryTurn reopens a consumed reclaim immediately when no query or recovery owns it', async () => { + it('driveDeliveryTurn completes an already-consumed reclaim immediately', async () => { const retrySpy = mock(() => 'db-1'); mockDb.getSDKMessageRepo = mock(() => ({ getDeliveryContent: mock(() => ({ content: 'x', sendStatus: 'consumed' })), hasTerminalResultAfter: mock(() => false), - hasRecoveryInterceptedResultAfter: mock(() => true), hasDeliveryTurnEnd: mock(() => false), clearDeliveryTurnEnd: mock(() => {}), getErrorTerminalResultSubtypeAfter: mock(() => null), @@ -1186,12 +1188,16 @@ describe('AgentSession', () => { ); await agentSession.stateManager.setProcessing('uuid-reclaim'); - const drive = agentSession.driveDeliveryTurn('uuid-reclaim', 'hello', null, true, () => true); - const error = await drive.catch((caught) => caught); + const result = await agentSession.driveDeliveryTurn( + 'uuid-reclaim', + 'hello', + null, + true, + () => true + ); - expect(error).toBeInstanceOf(MessageDeliveryRecoverableTurnError); - expect(error).toMatchObject({ message: 'Turn ended without a response' }); - expect(retrySpy).toHaveBeenCalledTimes(1); + expect(result).toEqual({ outcome: 'completed' }); + expect(retrySpy).not.toHaveBeenCalled(); }); it('cancelRateLimitRetry cancels the parked delivery for the episode message', async () => { @@ -1209,7 +1215,7 @@ describe('AgentSession', () => { expect(cancelDelivery).toHaveBeenCalledWith('test-session-id', 'msg-episode'); }); - it('driveDeliveryTurn makes a terminal turn error non-retryable without reopening', async () => { + it('driveDeliveryTurn completes an already-consumed row despite a terminal turn error', async () => { const retrySpy = mock(() => 'db-terminal'); mockDb.getSDKMessageRepo = mock(() => ({ getDeliveryContent: mock(() => ({ content: 'x', sendStatus: 'consumed' })), @@ -1248,13 +1254,12 @@ describe('AgentSession', () => { }; await agentSession.stateManager.setIdle(); - const error = await drive.catch((caught) => caught); - expect(error).toBeInstanceOf(MessageDeliveryTerminalTurnError); - expect(error).toMatchObject({ message: 'Sign in again', category: 'authentication' }); + const result = await drive; + expect(result).toEqual({ outcome: 'completed' }); expect(retrySpy).not.toHaveBeenCalled(); }); - it('driveDeliveryTurn treats a non-retryable persisted error subtype as terminal', async () => { + it('driveDeliveryTurn completes an already-consumed row despite a non-retryable persisted error subtype', async () => { const retrySpy = mock(() => 'db-terminal-subtype'); mockDb.getSDKMessageRepo = mock(() => ({ getDeliveryContent: mock(() => ({ content: 'x', sendStatus: 'consumed' })), @@ -1278,13 +1283,12 @@ describe('AgentSession', () => { await new Promise((resolve) => setTimeout(resolve, 10)); await agentSession.stateManager.setIdle(); - const error = await drive.catch((caught) => caught); - expect(error).toBeInstanceOf(MessageDeliveryTerminalTurnError); - expect(error).toMatchObject({ category: 'error_max_budget_usd' }); + const result = await drive; + expect(result).toEqual({ outcome: 'completed' }); expect(retrySpy).not.toHaveBeenCalled(); }); - it('driveDeliveryTurn re-arms at most twice within the 250ms spurious turn-end grace', async () => { + it('driveDeliveryTurn completes on SDK consumption and does not re-arm through later idle churn', async () => { let sendStatus = 'enqueued'; const retrySpy = mock(() => 'db-grace-cap'); mockDb.getSDKMessageRepo = mock(() => ({ @@ -1315,35 +1319,24 @@ describe('AgentSession', () => { }; await agentSession.stateManager.setProcessing('uuid-grace-cap'); - const drive = agentSession.driveDeliveryTurn( + const result = await agentSession.driveDeliveryTurn( 'uuid-grace-cap', 'hello', null, false, () => true ); - let settled = false; - void drive.then( - () => { - settled = true; - }, - () => { - settled = true; - } - ); + + expect(result).toEqual({ outcome: 'completed' }); + expect(retrySpy).not.toHaveBeenCalled(); + + await agentSession.stateManager.setIdle(); await new Promise((resolve) => setTimeout(resolve, 10)); - for (let fire = 0; fire < 3; fire++) { - await agentSession.stateManager.setIdle(); - if (fire < 2) { - await new Promise((resolve) => setTimeout(resolve, 10)); - expect(settled).toBe(false); - await agentSession.stateManager.setProcessing('uuid-grace-cap'); - } - } + await agentSession.stateManager.setProcessing('uuid-grace-cap'); + await new Promise((resolve) => setTimeout(resolve, 10)); + await agentSession.stateManager.setIdle(); - const error = await drive.catch((caught) => caught); - expect(error).toBeInstanceOf(MessageDeliveryRecoverableTurnError); - expect(retrySpy).toHaveBeenCalledTimes(1); + expect(result).toEqual({ outcome: 'completed' }); }); it('driveDeliveryTurn re-arms through a spurious turn-end fired right after a fresh admission', async () => { @@ -7580,7 +7573,7 @@ describe('AgentSession', () => { queue.acknowledgeYielded(hardUuid); await agentSession.stateManager.setIdle(); resolveQuery(); - await expect(drive).rejects.toThrow('Turn ended without a response'); + await expect(drive).rejects.toThrow('ACP prompt was not accepted before query end'); expect(reportStage).toHaveBeenCalledWith('sdk_admitted', expect.anything()); } finally { db.close(); diff --git a/packages/daemon/tests/unit/1-core/agent/delivery-retry-livelock-convergence.test.ts b/packages/daemon/tests/unit/1-core/agent/delivery-retry-livelock-convergence.test.ts index e9eb182951..7ded963eb8 100644 --- a/packages/daemon/tests/unit/1-core/agent/delivery-retry-livelock-convergence.test.ts +++ b/packages/daemon/tests/unit/1-core/agent/delivery-retry-livelock-convergence.test.ts @@ -367,15 +367,23 @@ describe('delivery retry livelock convergence (task #1256 incident)', () => { }, }; await agentSession.stateManager.setProcessing(WEDGE_UUID); - const drive = agentSession - .driveDeliveryTurn(WEDGE_UUID, 'go', null, true, () => true, undefined, undefined, observer) - .catch((error: unknown) => error); + const drive = agentSession.driveDeliveryTurn( + WEDGE_UUID, + 'go', + null, + true, + () => true, + undefined, + undefined, + observer + ); await queryReady; await new Promise((resolve) => setTimeout(resolve, 10)); await agentSession.stateManager.setIdle(); - const result = (await drive) as Error; - expect(result).toBeInstanceOf(MessageDeliveryRecoverableTurnError); - expect(result).not.toBeInstanceOf(MessageDeliveryTerminalTurnError); + await expect(drive).resolves.toEqual({ outcome: 'completed' }); + expect(db.getSDKMessageRepo().getDeliveryContent(SESSION_ID, WEDGE_UUID)).toMatchObject({ + sendStatus: 'consumed', + }); } }); diff --git a/packages/daemon/tests/unit/1-core/agent/message-delivery-transcript-harness.test.ts b/packages/daemon/tests/unit/1-core/agent/message-delivery-transcript-harness.test.ts index b181f77b5c..cf8efe2390 100644 --- a/packages/daemon/tests/unit/1-core/agent/message-delivery-transcript-harness.test.ts +++ b/packages/daemon/tests/unit/1-core/agent/message-delivery-transcript-harness.test.ts @@ -833,7 +833,6 @@ describe('delivery transcript parity harness (A1a)', () => { uuids: [turnUuid], dbIds: [], }, - { op: 'db:recordTurnEnd', sessionId, uuid: turnUuid }, ]); } finally { db.close(); @@ -892,7 +891,7 @@ describe('delivery transcript parity harness (A1a)', () => { await waitForTranscript( harness, - (e) => e.op === 'db:recordTurnEnd' && e.uuid === kickoffUuid + (e) => e.op === 'db:markConsumedBatch' && e.uuids.includes(kickoffUuid) ); repo.saveSDKMessage(sessionId, { type: 'result', @@ -945,7 +944,6 @@ describe('delivery transcript parity harness (A1a)', () => { uuids: [kickoffUuid, memberUuid], dbIds: [expect.any(String)], }, - { op: 'db:recordTurnEnd', sessionId, uuid: kickoffUuid }, ]); } finally { db.close(); @@ -1290,7 +1288,6 @@ describe('delivery transcript parity harness (A1a)', () => { uuids: [kickoffUuid, memberUuid], dbIds: [expect.any(String)], }, - { op: 'db:recordTurnEnd', sessionId, uuid: kickoffUuid }, ]); } finally { db.close(); diff --git a/packages/daemon/tests/unit/1-core/agent/message-queue.test.ts b/packages/daemon/tests/unit/1-core/agent/message-queue.test.ts index f2ec64aff2..3462714657 100644 --- a/packages/daemon/tests/unit/1-core/agent/message-queue.test.ts +++ b/packages/daemon/tests/unit/1-core/agent/message-queue.test.ts @@ -226,6 +226,70 @@ describe('MessageQueue', () => { }); }); + describe('requeueAllYielded', () => { + it('requeues every yielded message durably', async () => { + queue.start(); + const ack1 = queue.enqueueWithId('all-yielded-1', 'Message 1'); + const ack2 = queue.enqueueWithId('all-yielded-2', 'Message 2'); + const generator = queue.messageGenerator(testSessionId); + + const first = await generator.next(); + expect(first.done).toBe(false); + expect((first.value?.message as { uuid?: string }).uuid).toBe('all-yielded-1'); + + const second = await generator.next(); + expect(second.done).toBe(false); + expect((second.value?.message as { uuid?: string }).uuid).toBe('all-yielded-2'); + + expect(queue.hasYielded('all-yielded-1')).toBe(true); + expect(queue.hasYielded('all-yielded-2')).toBe(true); + + const requeued = queue.requeueAllYielded({ durable: true }); + expect(requeued).toEqual(['all-yielded-1', 'all-yielded-2']); + expect(queue.hasYielded('all-yielded-1')).toBe(false); + expect(queue.hasYielded('all-yielded-2')).toBe(false); + expect(queue.hasPendingOrInFlight('all-yielded-1')).toBe(true); + expect(queue.hasPendingOrInFlight('all-yielded-2')).toBe(true); + + const reyielded1 = await generator.next(); + expect(reyielded1.done).toBe(false); + expect((reyielded1.value?.message as { uuid?: string }).uuid).toBe('all-yielded-2'); + + const reyielded2 = await generator.next(); + expect(reyielded2.done).toBe(false); + expect((reyielded2.value?.message as { uuid?: string }).uuid).toBe('all-yielded-1'); + + reyielded1.value?.onSent(); + reyielded2.value?.onSent(); + await ack1; + await ack2; + queue.stop(); + }); + + it('returns an empty array when no message is yielded', async () => { + const requeued = queue.requeueAllYielded(); + expect(requeued).toEqual([]); + }); + + it('marks requeued messages as durable so they survive the next query start', async () => { + queue.start(); + const ack = queue.enqueueWithId('durable-requeue', 'Message 1'); + const generator = queue.messageGenerator(testSessionId); + const result = await generator.next(); + expect(result.done).toBe(false); + + queue.requeueAllYielded({ durable: true }); + + const reyielded = await generator.next(); + expect(reyielded.done).toBe(false); + expect((reyielded.value?.message as { uuid?: string }).uuid).toBe('durable-requeue'); + + reyielded.value?.onSent(); + await ack; + queue.stop(); + }); + }); + describe('hasPendingOrInFlight', () => { it('reports false for an unknown id and true while queued/claimed/yielded', async () => { expect(queue.hasPendingOrInFlight('nope')).toBe(false); diff --git a/packages/daemon/tests/unit/1-core/agent/query-lifecycle-manager.test.ts b/packages/daemon/tests/unit/1-core/agent/query-lifecycle-manager.test.ts index 20207faf3a..5fcb732ce3 100644 --- a/packages/daemon/tests/unit/1-core/agent/query-lifecycle-manager.test.ts +++ b/packages/daemon/tests/unit/1-core/agent/query-lifecycle-manager.test.ts @@ -102,6 +102,13 @@ describe('QueryLifecycleManager', () => { saveHyperNeoActionMessage: saveHyperNeoActionMessageSpy, getSDKMessageRepo: () => ({ hasUnresolvedHyperNeoAction: () => hasUnresolvedResumeChoice, + getUserMessagesByStatus: () => ({ messages: [] }), + getUserMessageContentByUuid: () => null, + }), + getJobQueueRepo: () => ({ + cancelForSessionWithMessages: mock(() => []), + cancelForSession: mock(() => {}), + getActiveDeliveryBatchUuids: mock(() => []), }), } as unknown as Database, messageHub: { @@ -1938,4 +1945,108 @@ describe('QueryLifecycleManager', () => { ); }); }); + + describe('model-switch re-injection delivery cleanup', () => { + function patchReinjectionMocks(): { + cancelForSessionWithMessages: ReturnType; + cancelForSession: ReturnType; + getActiveDeliveryBatchUuids: ReturnType; + } { + const cancelForSessionWithMessages = mock(() => ['msg-delivery-1']); + const cancelForSession = mock(() => {}); + const getActiveDeliveryBatchUuids = mock(() => ['last-consumed']); + + (mockContext.db as unknown as Record).getJobQueueRepo = () => ({ + cancelForSessionWithMessages, + cancelForSession, + getActiveDeliveryBatchUuids, + }); + + (mockContext.db as unknown as Record).getSDKMessageRepo = () => ({ + hasUnresolvedHyperNeoAction: () => false, + getUserMessagesByStatus: () => ({ messages: [{ uuid: 'last-consumed' }] }), + getUserMessageContentByUuid: () => 'last message', + }); + + return { cancelForSessionWithMessages, cancelForSession, getActiveDeliveryBatchUuids }; + } + + beforeEach(() => { + messageQueue = new MessageQueue(); + mockContext = createMockContext({ messageQueue }); + manager = new QueryLifecycleManager(mockContext); + }); + + test('restart cancels in-flight delivery jobs and re-enqueues the last consumed message when transcript is not preserved', async () => { + const spies = patchReinjectionMocks(); + mockContext.session.sdkSessionId = undefined; + mockContext.session.acpSessionId = undefined; + mockContext.session.config.provider = 'anthropic'; + manager = new QueryLifecycleManager(mockContext); + + await manager.restart(); + + expect(spies.cancelForSessionWithMessages).toHaveBeenCalledWith('test-session'); + expect(messageQueue.hasPendingOrInFlight('last-consumed')).toBe(true); + }); + + test('restart does not re-enqueue the last consumed message when SDK transcript is preserved', async () => { + const spies = patchReinjectionMocks(); + mockContext.session.sdkSessionId = 'preserved-sdk-session'; + mockContext.session.config.provider = 'anthropic'; + manager = new QueryLifecycleManager(mockContext); + + await manager.restart(); + + expect(spies.cancelForSessionWithMessages).toHaveBeenCalledWith('test-session'); + expect(messageQueue.hasPendingOrInFlight('last-consumed')).toBe(false); + }); + + test('restart does not re-enqueue the last consumed message for ACP when ACP session is preserved', async () => { + const spies = patchReinjectionMocks(); + mockContext.session.config.provider = 'acp'; + mockContext.session.acpSessionId = 'preserved-acp-session'; + manager = new QueryLifecycleManager(mockContext); + + await manager.restart(); + + expect(spies.cancelForSessionWithMessages).toHaveBeenCalledWith('test-session'); + expect(messageQueue.hasPendingOrInFlight('last-consumed')).toBe(false); + }); + + test('reset cancels in-flight delivery jobs and re-enqueues the last consumed message when restarting without transcript preservation', async () => { + const spies = patchReinjectionMocks(); + mockContext.session.sdkSessionId = undefined; + mockContext.session.config.provider = 'anthropic'; + mockContext.queryObject = { + interrupt: mock(async () => {}), + close: mock(() => {}), + } as unknown as QueryLifecycleManagerContext['queryObject']; + mockContext.queryPromise = Promise.resolve(); + manager = new QueryLifecycleManager(mockContext); + + const result = await manager.reset({ restartAfter: true }); + + expect(result.success).toBe(true); + expect(spies.cancelForSession).toHaveBeenCalledWith('test-session'); + expect(messageQueue.hasPendingOrInFlight('last-consumed')).toBe(true); + }); + + test('reset preserves the last consumed message when restart is not requested', async () => { + const spies = patchReinjectionMocks(); + mockContext.session.sdkSessionId = undefined; + mockContext.session.config.provider = 'anthropic'; + mockContext.queryObject = { + interrupt: mock(async () => {}), + close: mock(() => {}), + } as unknown as QueryLifecycleManagerContext['queryObject']; + mockContext.queryPromise = Promise.resolve(); + manager = new QueryLifecycleManager(mockContext); + + const result = await manager.reset({ restartAfter: false }); + + expect(result.success).toBe(true); + expect(messageQueue.hasPendingOrInFlight('last-consumed')).toBe(false); + }); + }); });