@@ -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);
@@ -7273,6 +7270,19 @@ function chatAgent<
72737270 // `messagesInput.waitWithIdleTimeout` so recovered turns fire first.
72747271 const bootInjectedQueue: ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>[] =
72757272 [];
7273+ const recoveredSeqByPayload = new WeakMap<
7274+ ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>,
7275+ number
7276+ >();
7277+ const dispatchBootInjected = (): ChatTaskWirePayload<
7278+ TUIMessage,
7279+ inferSchemaIn<TClientDataSchema>
7280+ > => {
7281+ const injected = bootInjectedQueue.shift()!;
7282+ const settledSeq = recoveredSeqByPayload.get(injected);
7283+ if (settledSeq !== undefined) chatInputRouter().settleRecovered(settledSeq);
7284+ return injected;
7285+ };
72767286 const couldHavePriorState = payload.continuation === true || ctx.attempt.number > 1;
72777287
72787288 // `.in` resume cursor, computed at most once per boot. The boot
@@ -7438,18 +7448,11 @@ function chatAgent<
74387448
74397449 // ── session.in router ──────────────────────────────────────────
74407450 //
7441- // Reads the turn boundary and subscribes in one call. `bootInCursor` is
7442- // only a fallback: the boot block above may already have resolved a
7443- // cursor from the snapshot, which is used when the boundary itself
7444- // carries none. Everything the boot replayed off `.in` is dispatched from
7445- // `bootInjectedQueue` below, so it goes into the floor here — folded in
7446- // after the subscription opens, the live tail re-delivers it as a turn.
7447- const lastRecoveredInSeq =
7448- replayedInTail.length > 0 ? replayedInTail[replayedInTail.length - 1]!.seqNum : undefined;
7451+ const recoveredSeqNums = replayedInTail.map((r) => r.seqNum);
74497452
74507453 await installChatInputRouter(payload.chatId, {
74517454 fallbackResumeFrom: bootInCursorResolved ? bootInCursor : undefined,
7452- recoveredThrough: lastRecoveredInSeq ,
7455+ recoveredSeqNums ,
74537456 resuming: Boolean(payload.continuation) || ctx.attempt.number > 1,
74547457 });
74557458
@@ -7539,7 +7542,7 @@ function chatAgent<
75397542 // branches: at n=1 the orphan partial is dropped and the interrupted
75407543 // user is re-dispatched as a fresh turn instead.
75417544 let seedChain: TUIMessage[];
7542- let recoveredTurns: TUIMessage[];
7545+ let recoveredEntries: { message: TUIMessage; seqNum: number | undefined } [];
75437546 if (hookChain !== undefined) {
75447547 seedChain = hookChain;
75457548 } else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
@@ -7548,11 +7551,20 @@ function chatAgent<
75487551 seedChain = settledMessages;
75497552 }
75507553 if (hookRecoveredTurns !== undefined) {
7551- recoveredTurns = hookRecoveredTurns;
7554+ const seqNumByRecoveredId = new Map<string, number>();
7555+ for (const entry of replayedInTail) {
7556+ seqNumByRecoveredId.set(entry.message.id, entry.seqNum);
7557+ }
7558+ recoveredEntries = hookRecoveredTurns.map((message) => ({
7559+ message,
7560+ seqNum: seqNumByRecoveredId.get(message.id),
7561+ }));
75527562 } else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
7553- recoveredTurns = inFlightUsers.slice(1);
7563+ recoveredEntries = replayedInTail
7564+ .slice(1)
7565+ .map((r) => ({ message: r.message, seqNum: r.seqNum }));
75547566 } else {
7555- recoveredTurns = inFlightUsers ;
7567+ recoveredEntries = replayedInTail.map((r) => ({ message: r.message, seqNum: r.seqNum })) ;
75567568 }
75577569 // `beforeBoot` errors bubble — the customer opted into blocking
75587570 // persistence and a failure there should fail the run rather than
@@ -7583,12 +7595,13 @@ function chatAgent<
75837595 for (const entry of replayedInTail) {
75847596 metadataById.set(entry.message.id, entry.metadata);
75857597 }
7586- for (const msg of recoveredTurns) {
7598+ const dispatchedRecoveredSeqs = new Set<number>();
7599+ for (const { message: msg, seqNum } of recoveredEntries) {
75877600 if (wireMessageId && msg.id === wireMessageId) continue;
75887601 const recoveredMetadata = metadataById.has(msg.id)
75897602 ? metadataById.get(msg.id)
75907603 : payload.metadata;
7591- bootInjectedQueue.push( {
7604+ const injectedPayload = {
75927605 chatId: payload.chatId,
75937606 sessionId: payload.sessionId,
75947607 metadata: recoveredMetadata,
@@ -7597,7 +7610,17 @@ function chatAgent<
75977610 messageId: msg.id,
75987611 continuation: payload.continuation,
75997612 previousRunId: payload.previousRunId,
7600- } as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>);
7613+ } as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>;
7614+ bootInjectedQueue.push(injectedPayload);
7615+ if (seqNum !== undefined) {
7616+ recoveredSeqByPayload.set(injectedPayload, seqNum);
7617+ dispatchedRecoveredSeqs.add(seqNum);
7618+ }
7619+ }
7620+ for (const entry of replayedInTail) {
7621+ if (!dispatchedRecoveredSeqs.has(entry.seqNum)) {
7622+ chatInputRouter().settleRecovered(entry.seqNum);
7623+ }
76017624 }
76027625
76037626 accumulatedUIMessages = seedChain;
@@ -7781,7 +7804,7 @@ function chatAgent<
77817804 */
77827805 let dispatchedRecoveredFirstTurn = false;
77837806 if (preloaded && bootInjectedQueue.length > 0) {
7784- currentWirePayload = bootInjectedQueue.shift()! ;
7807+ currentWirePayload = dispatchBootInjected() ;
77857808 dispatchedRecoveredFirstTurn = true;
77867809 }
77877810
@@ -8032,7 +8055,7 @@ function chatAgent<
80328055 // waiting on the live session.in. Subsequent recovered turns
80338056 // get drained by the end-of-turn picker below.
80348057 if (bootInjectedQueue.length > 0) {
8035- currentWirePayload = bootInjectedQueue.shift()! ;
8058+ currentWirePayload = dispatchBootInjected() ;
80368059 } else {
80378060 const effectiveIdleTimeout = idleTimeoutInSeconds ?? payload.idleTimeoutInSeconds;
80388061 const effectiveTurnTimeout =
@@ -9613,7 +9636,7 @@ function chatAgent<
96139636 // produced these from in-flight user messages on session.in
96149637 // that the dead predecessor never acknowledged.
96159638 if (bootInjectedQueue.length > 0) {
9616- currentWirePayload = bootInjectedQueue.shift()! ;
9639+ currentWirePayload = dispatchBootInjected() ;
96179640 return "continue";
96189641 }
96199642
@@ -9989,7 +10012,7 @@ function chatAgent<
998910012 // recovered turn shouldn't strand the rest of the boot queue
999010013 // until an unrelated live message arrives.
999110014 if (bootInjectedQueue.length > 0) {
9992- currentWirePayload = bootInjectedQueue.shift()! ;
10015+ currentWirePayload = dispatchBootInjected() ;
999310016 continue;
999410017 }
999510018
0 commit comments