@@ -2322,7 +2322,11 @@ async function findSessionInReplayWindowEnd(
23222322 */
23232323async function installChatInputRouter(
23242324 chatId: string,
2325- options?: { fallbackResumeFrom?: number; recoveredThrough?: number; resuming?: boolean }
2325+ options?: {
2326+ fallbackResumeFrom?: number;
2327+ recoveredSeqNums?: readonly number[];
2328+ resuming?: boolean;
2329+ }
23262330): Promise<SessionChannelRouter> {
23272331 const entry = chatInputRouterEntry(chatId);
23282332 if (entry.attached) return entry.router;
@@ -2353,20 +2357,13 @@ async function installChatInputRouter(
23532357 }
23542358 }
23552359
2356- // A boot that replayed `.in` itself has already answered everything up to
2357- // `recoveredThrough`, so the floor has to cover it before the tail opens.
2358- if (options?.recoveredThrough !== undefined) {
2359- const recovered = options.recoveredThrough;
2360- checkpoint.resumeFrom = Math.max(checkpoint.resumeFrom ?? recovered, recovered);
2361- checkpoint.appliedThrough = Math.max(
2362- checkpoint.appliedThrough ?? checkpoint.resumeFrom,
2363- checkpoint.resumeFrom
2364- );
2365- }
2366-
23672360 const router = entry.router;
23682361 router.restore(checkpoint);
23692362
2363+ if (options?.recoveredSeqNums && options.recoveredSeqNums.length > 0) {
2364+ router.markRecovered(options.recoveredSeqNums);
2365+ }
2366+
23702367 const floor = router.resumeFrom();
23712368 if (floor !== undefined) {
23722369 sessionStreams.setLastSeqNum(chatId, "in", floor);
@@ -7272,6 +7269,20 @@ function chatAgent<
72727269 // `messagesInput.waitWithIdleTimeout` so recovered turns fire first.
72737270 const bootInjectedQueue: ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>[] =
72747271 [];
7272+ const recoveredSeqByPayload = new WeakMap<
7273+ ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>,
7274+ number
7275+ >();
7276+ const dispatchBootInjected = (): ChatTaskWirePayload<
7277+ TUIMessage,
7278+ inferSchemaIn<TClientDataSchema>
7279+ > => bootInjectedQueue.shift()!;
7280+ const settleRecoveredTurn = (
7281+ wirePayload: ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>
7282+ ) => {
7283+ const settledSeq = recoveredSeqByPayload.get(wirePayload);
7284+ if (settledSeq !== undefined) chatInputRouter().settleRecovered(settledSeq);
7285+ };
72757286 const couldHavePriorState = payload.continuation === true || ctx.attempt.number > 1;
72767287
72777288 // `.in` resume cursor, computed at most once per boot. The boot
@@ -7437,18 +7448,11 @@ function chatAgent<
74377448
74387449 // ── session.in router ──────────────────────────────────────────
74397450 //
7440- // Reads the turn boundary and subscribes in one call. `bootInCursor` is
7441- // only a fallback: the boot block above may already have resolved a
7442- // cursor from the snapshot, which is used when the boundary itself
7443- // carries none. Everything the boot replayed off `.in` is dispatched from
7444- // `bootInjectedQueue` below, so it goes into the floor here — folded in
7445- // after the subscription opens, the live tail re-delivers it as a turn.
7446- const lastRecoveredInSeq =
7447- replayedInTail.length > 0 ? replayedInTail[replayedInTail.length - 1]!.seqNum : undefined;
7451+ const recoveredSeqNums = replayedInTail.map((r) => r.seqNum);
74487452
74497453 await installChatInputRouter(payload.chatId, {
74507454 fallbackResumeFrom: bootInCursorResolved ? bootInCursor : undefined,
7451- recoveredThrough: lastRecoveredInSeq ,
7455+ recoveredSeqNums ,
74527456 resuming: Boolean(payload.continuation) || ctx.attempt.number > 1,
74537457 });
74547458
@@ -7538,7 +7542,7 @@ function chatAgent<
75387542 // branches: at n=1 the orphan partial is dropped and the interrupted
75397543 // user is re-dispatched as a fresh turn instead.
75407544 let seedChain: TUIMessage[];
7541- let recoveredTurns: TUIMessage[];
7545+ let recoveredEntries: { message: TUIMessage; seqNum: number | undefined } [];
75427546 if (hookChain !== undefined) {
75437547 seedChain = hookChain;
75447548 } else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
@@ -7547,11 +7551,22 @@ function chatAgent<
75477551 seedChain = settledMessages;
75487552 }
75497553 if (hookRecoveredTurns !== undefined) {
7550- recoveredTurns = hookRecoveredTurns;
7554+ const seqNumsByRecoveredId = new Map<string, number[]>();
7555+ for (const entry of replayedInTail) {
7556+ const existing = seqNumsByRecoveredId.get(entry.message.id);
7557+ if (existing) existing.push(entry.seqNum);
7558+ else seqNumsByRecoveredId.set(entry.message.id, [entry.seqNum]);
7559+ }
7560+ recoveredEntries = hookRecoveredTurns.map((message) => ({
7561+ message,
7562+ seqNum: seqNumsByRecoveredId.get(message.id)?.shift(),
7563+ }));
75517564 } else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
7552- recoveredTurns = inFlightUsers.slice(1);
7565+ recoveredEntries = replayedInTail
7566+ .slice(1)
7567+ .map((r) => ({ message: r.message, seqNum: r.seqNum }));
75537568 } else {
7554- recoveredTurns = inFlightUsers ;
7569+ recoveredEntries = replayedInTail.map((r) => ({ message: r.message, seqNum: r.seqNum })) ;
75557570 }
75567571 // `beforeBoot` errors bubble — the customer opted into blocking
75577572 // persistence and a failure there should fail the run rather than
@@ -7582,12 +7597,13 @@ function chatAgent<
75827597 for (const entry of replayedInTail) {
75837598 metadataById.set(entry.message.id, entry.metadata);
75847599 }
7585- for (const msg of recoveredTurns) {
7600+ const dispatchedRecoveredSeqs = new Set<number>();
7601+ for (const { message: msg, seqNum } of recoveredEntries) {
75867602 if (wireMessageId && msg.id === wireMessageId) continue;
75877603 const recoveredMetadata = metadataById.has(msg.id)
75887604 ? metadataById.get(msg.id)
75897605 : payload.metadata;
7590- bootInjectedQueue.push( {
7606+ const injectedPayload = {
75917607 chatId: payload.chatId,
75927608 sessionId: payload.sessionId,
75937609 metadata: recoveredMetadata,
@@ -7596,7 +7612,17 @@ function chatAgent<
75967612 messageId: msg.id,
75977613 continuation: payload.continuation,
75987614 previousRunId: payload.previousRunId,
7599- } as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>);
7615+ } as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>;
7616+ bootInjectedQueue.push(injectedPayload);
7617+ if (seqNum !== undefined) {
7618+ recoveredSeqByPayload.set(injectedPayload, seqNum);
7619+ dispatchedRecoveredSeqs.add(seqNum);
7620+ }
7621+ }
7622+ for (const entry of replayedInTail) {
7623+ if (!dispatchedRecoveredSeqs.has(entry.seqNum)) {
7624+ chatInputRouter().settleRecovered(entry.seqNum);
7625+ }
76007626 }
76017627
76027628 accumulatedUIMessages = seedChain;
@@ -7780,7 +7806,7 @@ function chatAgent<
77807806 */
77817807 let dispatchedRecoveredFirstTurn = false;
77827808 if (preloaded && bootInjectedQueue.length > 0) {
7783- currentWirePayload = bootInjectedQueue.shift()! ;
7809+ currentWirePayload = dispatchBootInjected() ;
77847810 dispatchedRecoveredFirstTurn = true;
77857811 }
77867812
@@ -8031,7 +8057,7 @@ function chatAgent<
80318057 // waiting on the live session.in. Subsequent recovered turns
80328058 // get drained by the end-of-turn picker below.
80338059 if (bootInjectedQueue.length > 0) {
8034- currentWirePayload = bootInjectedQueue.shift()! ;
8060+ currentWirePayload = dispatchBootInjected() ;
80358061 } else {
80368062 const effectiveIdleTimeout = idleTimeoutInSeconds ?? payload.idleTimeoutInSeconds;
80378063 const effectiveTurnTimeout =
@@ -8686,6 +8712,7 @@ function chatAgent<
86868712 chatId: currentWirePayload.chatId,
86878713 messageId: currentWirePayload.messageId,
86888714 });
8715+ settleRecoveredTurn(currentWirePayload);
86898716 await writeTurnCompleteChunk(currentWirePayload.chatId);
86908717 // Not a turn — don't consume an iteration.
86918718 turn--;
@@ -9504,6 +9531,8 @@ function chatAgent<
95049531 locals.set(chatResponsePartsKey, []);
95059532 }
95069533
9534+ settleRecoveredTurn(currentWirePayload);
9535+
95079536 // Write turn-complete control chunk — closes the frontend stream.
95089537 const turnCompleteResult = await writeTurnCompleteChunk(
95099538 currentWirePayload.chatId,
@@ -9634,7 +9663,7 @@ function chatAgent<
96349663 // produced these from in-flight user messages on session.in
96359664 // that the dead predecessor never acknowledged.
96369665 if (bootInjectedQueue.length > 0) {
9637- currentWirePayload = bootInjectedQueue.shift()! ;
9666+ currentWirePayload = dispatchBootInjected() ;
96389667 return "continue";
96399668 }
96409669
@@ -10012,7 +10041,7 @@ function chatAgent<
1001210041 // recovered turn shouldn't strand the rest of the boot queue
1001310042 // until an unrelated live message arrives.
1001410043 if (bootInjectedQueue.length > 0) {
10015- currentWirePayload = bootInjectedQueue.shift()! ;
10044+ currentWirePayload = dispatchBootInjected() ;
1001610045 continue;
1001710046 }
1001810047
0 commit comments