Conversation
…ecting 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.
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
| const jobQueue = db.getJobQueueRepo?.(); | ||
| if (jobQueue) { | ||
| jobQueue.cancelForSession(session.id); | ||
| } |
There was a problem hiding this comment.
🔴 Reset drops batched user messages
cancelForSession deletes batch metadata before buildRequeueContentFor reads it. A transcript-dropping reset then omits accepted batch members from the replacement conversation.
Prompt for agents
In packages/daemon/src/lib/agent/query-lifecycle-manager.ts, QueryLifecycleManager.reset cancels message_delivery jobs before resolving the last consumed UUID and rebuilding its batch content. cancelForSession deletes the only active job payload containing batchUuids, so buildRequeueContentFor can no longer reconstruct consumed batch members when restartAfter is true and the provider transcript was dropped. Capture the consumed UUID and complete reconstructed content before canceling jobs, as restart already does, then conditionally enqueue that snapshot after cancellation. Add a reset test with a multi-member active batch and no preserved transcript that verifies every admitted member appears in the replay content.
Was this helpful? React with 👍 or 👎 to provide feedback.
| 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' }; |
There was a problem hiding this comment.
🟡 Completed reclaims retain stale observers
When alreadyConsumed is true, the return bypasses cleanup for the admission's idle waiter and response observer. Later activity can trigger this stale delivery state.
| 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' }; | |
| 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; | |
| } | |
| return { outcome: 'completed' }; | |
| } |
Was this helpful? React with 👍 or 👎 to provide feedback.
|
|
||
| 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); |
There was a problem hiding this comment.
🔍 Lifecycle cleanup bypasses pipeline convention
restart and reset add imperative transcript-replay flows. Repository guidance requires new daemon business paths to use one direct superpipe pipeline where applicable.
(Refers to this code)
Was this helpful? React with 👍 or 👎 to provide feedback.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 2d0adb4ded
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| const { messages } = repo.getUserMessagesByStatus(session.id, 'consumed', 1, 'desc'); | ||
| const last = messages[0]; | ||
| return last?.uuid ?? null; |
There was a problem hiding this comment.
Resolve the active batch from any consumed member
When a transcript-dropping restart follows a batched delivery, resolveLastConsumedUuid() can select the newest batch member rather than the kickoff because markDeliveryConsumedByUuids() gives the members a shared timestamp and the query breaks ties by descending rowid. However, getActiveDeliveryBatchUuids() only matches the job's kickoff messageUuid, so this lookup returns no batch for a selected member and the replacement query receives only that member's content, silently omitting the rest of the accepted batch. The same loss occurs if the now-consumption-complete job is removed before this lookup; preserve the batch before completion/cancellation or resolve it by membership.
Useful? React with 👍 / 👎.
| for (const message of Array.from(this.yielded)) { | ||
| if (this.requeueYielded(message.id, options)) { | ||
| requeued.push(message.id); |
There was a problem hiding this comment.
Preserve FIFO order when requeueing yielded messages
When more than one message has been yielded but not acknowledged during a restart, iterating the Set in yield order while requeueYielded() prepends each entry reverses them (A, B becomes B, A). The replacement provider therefore receives pending user messages in the opposite order, which can change their meaning; iterate the snapshot in reverse or otherwise prepend the group while retaining FIFO order.
Useful? React with 👍 / 👎.
| const jobQueue = db.getJobQueueRepo?.(); | ||
| if (jobQueue) { | ||
| jobQueue.cancelForSession(session.id); | ||
| } |
There was a problem hiding this comment.
Requeue unaccepted deliveries before canceling reset jobs
When a config-triggered reset({ restartAfter: true }) races a prompt that has been yielded or submitted but not accepted by the provider, messageQueue.clear() has already removed its in-memory admission and this new cancellation deletes its durable delivery job. Unlike restart(), this path never calls requeueAllYielded; a pre-yield callback may also have marked the row consumed even though the aborted SDK never committed it, causing transcript-preservation checks to skip re-enqueue and leaving no reconciler-visible enqueued row. The replacement query can therefore silently omit the user's current prompt.
Useful? React with 👍 / 👎.
| 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' }; |
There was a problem hiding this comment.
Release admission resources before completing consumed reclaims
For an already-consumed reclaim, runDeliveryTurnAdmission() has installed a turn-end waiter and may have armed deliveryResponseObserver, but this return occurs before the common finally that cancels and disarms them. Each reclaimed job can consequently leave an idle-transition waiter behind and attribute a later SDK response to an already-completed delivery observer; explicitly release started.turnEnd and started.responseObserver before returning, as the removed reclaim branch previously did.
Useful? React with 👍 / 👎.
| this.ctx.messageQueue.requeueAllYielded({ durable: true }); | ||
| this.reenqueueLastConsumedIfTranscriptDropped(lastConsumedUuid, lastConsumedContent); |
There was a problem hiding this comment.
Restrict transcript-drop replay to the interrupted delivery
When a persistent query is idle after a completed turn and the user switches across the ACP/non-ACP boundary, the provider session ID is cleared before restart(), so this unconditional global-last-consumed lookup re-enqueues the previous prompt even though no delivery was interrupted. The new provider then answers an already-completed user message again and may repeat its tool side effects. Use the canceled active-job UUIDs or the yielded/pending queue snapshot to decide whether a consumed row actually needs recovery instead of replaying the latest historical row whenever the transcript is absent.
Useful? React with 👍 / 👎.
|
Closing as superseded: the same fix landed re-sliced as MDR 3/N (#3436). Branch retained. |
Summary
Fixes the model-switch re-injection loop (#1686): switching models on a running session re-injected the last already-delivered message and retried forever.
Root cause: the
message_deliveryjob treated turn end as its completion boundary. A torn-down turn (model switch, reset) hitclassifyTurnCompletionwithreopenForRetry: true,reopenDeliveryForRetryflipped the consumed row back toenqueued, and the retried job re-fed it.Changes
driveDeliveryTurnis delivery-only now (agent-session.ts): non-ACP turns complete at the SDK admission ack (markDeliveryBatchConsumed+signalDeliveryConsumed, batch members included); ACP turns complete atmarkMessageAccepted→onDeliveryTurnAccepted. The turn-end wait loop, spurious-rearm grace, and the recovery-intercepted reclaim that reopened consumed rows are gone. The failure tail (query ended before consumption, invalidated acknowledgment) still classifies and throws recoverable, but never reopens a row that was already acknowledged.QueryLifecycleManager.restart(): cancels in-flightmessage_deliveryjobs, durably requeues unconsumed yielded prompts (MessageQueue.requeueAllYielded), and re-feeds the last consumed message only when the replacement query does not preserve the provider transcript (nosdkSessionId/acpSessionId) — so cross-provider/ACP state drops still re-feed, same-generation restarts do not.QueryLifecycleManager.reset(): same cancel + conditional re-feed forrestartAfter: true(ACP reset clearsacpSessionId, so the transcript really is dropped there).Tests
query-lifecycle-manager.test.ts: restart/reset re-injection cleanup matrix (cancel called; re-feed only when transcript dropped; ACP/SDK preservation).message-queue.test.ts:requeueAllYieldeddurable requeue, empty-set, and re-yield behavior.agent-session(consumption completes delivery; already-consumed reclaims complete without reopening), livelock convergence, transcript harness (nodb:recordTurnEndon the delivery path), admission pipeline.bun run checkgreen; all 6 daemon shards green (1486 tests, 0 failures).Note: the delivery-turn stall watchdog now only covers the ACP pre-acceptance window; turn-level wedged-query self-healing for non-ACP sessions was previously a side effect of the delivery job waiting for turn end and is intentionally decoupled along with it.