Skip to content

Commit da655ab

Browse files
committed
fix(core,sdk): stop chat.agent losing messages during recovery
A continuation boot replays unacknowledged user messages off session.in and dispatches them itself. Suppressing the tail's re-answer by folding every recovered message into the resume cursor let the cursor advance past a message the run had not answered yet, so a crash mid-recovery dropped the rest. Recovered messages are now claimed on the session-stream router: a claim drops a late tail re-delivery (no double answer) and holds the resume cursor behind the message until the boot settles it (no drop). The boot settles each as it dispatches it, or immediately for ones folded into the seed chain or skipped.
1 parent 0350e00 commit da655ab

4 files changed

Lines changed: 304 additions & 33 deletions

File tree

packages/core/src/v3/sessionStreams/router.test.ts

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -573,3 +573,67 @@ describe("SessionChannelRouter: untake", () => {
573573
expect(r.pendingCount("messages")).toBe(1);
574574
});
575575
});
576+
577+
describe("SessionChannelRouter: recovered claim/settle floor", () => {
578+
it("drops a claimed record on ingest instead of queueing it", () => {
579+
const r = router();
580+
r.restore({ resumeFrom: 0 });
581+
r.markRecovered([1, 2]);
582+
expect(r.ingest(rec(1, "message"))).toEqual({ action: "drop", reason: "recovered" });
583+
expect(r.hasPending("messages")).toBe(false);
584+
});
585+
586+
it("holds the resume floor below the earliest owed record until it is settled", () => {
587+
const r = router();
588+
r.restore({ resumeFrom: 0 });
589+
r.markRecovered([1, 2]);
590+
expect(r.resumeFloor()).toBe(0);
591+
r.settleRecovered(1);
592+
expect(r.resumeFloor()).toBe(1);
593+
r.settleRecovered(2);
594+
expect(r.resumeFloor()).toBe(2);
595+
});
596+
597+
it("advances the floor after settling even when the tail never re-delivers", () => {
598+
const r = router();
599+
r.restore({ resumeFrom: 0 });
600+
r.markRecovered([1, 2]);
601+
r.settleRecovered(1);
602+
r.settleRecovered(2);
603+
expect(r.resumeFloor()).toBe(2);
604+
expect(r.appliedThrough()).toBe(2);
605+
});
606+
607+
it("keeps dropping a claimed record after it is settled, so a late tail re-read is never answered", () => {
608+
const r = router();
609+
r.restore({ resumeFrom: 0 });
610+
r.markRecovered([1]);
611+
r.settleRecovered(1);
612+
expect(r.ingest(rec(1, "message"))).toEqual({ action: "drop", reason: "recovered" });
613+
expect(r.hasPending("messages")).toBe(false);
614+
});
615+
616+
it("queues a live record whose sequence was never claimed", () => {
617+
const r = router();
618+
r.restore({ resumeFrom: 0 });
619+
r.markRecovered([1, 2]);
620+
expect(r.ingest(rec(3, "message"))).toEqual({ action: "queue", route: "messages" });
621+
});
622+
623+
it("advances only over the contiguous claimed run, holding the floor below a gap", () => {
624+
const r = router();
625+
r.restore({ resumeFrom: 0 });
626+
r.markRecovered([1, 3]);
627+
r.settleRecovered(1);
628+
r.settleRecovered(3);
629+
expect(r.resumeFloor()).toBe(1);
630+
});
631+
632+
it("clears claims and owed records on reset", () => {
633+
const r = router();
634+
r.restore({ resumeFrom: 0 });
635+
r.markRecovered([1, 2]);
636+
r.reset();
637+
expect(r.ingest(rec(1, "message"))).toEqual({ action: "queue", route: "messages" });
638+
});
639+
});

packages/core/src/v3/sessionStreams/router.ts

Lines changed: 70 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,14 @@ export type RouterDropReason =
6161
*/
6262
| "replayed"
6363
/** An `at-arrival` record with no handler attached right now. */
64-
| "no-handler";
64+
| "no-handler"
65+
/**
66+
* A record the boot took responsibility for over HTTP (see
67+
* {@link SessionChannelRouter.markRecovered}). It is dropped however late the
68+
* live tail re-delivers it, so a recovered message is never also answered as
69+
* a router turn.
70+
*/
71+
| "recovered";
6572

6673
export type RouterDecision =
6774
/** Handed to a consumer that was already waiting, or to a live handler. */
@@ -143,6 +150,8 @@ export class SessionChannelRouter {
143150
#highestSeq: number | undefined;
144151
#resumeFrom: number | undefined;
145152
#appliedThrough: number | undefined;
153+
#claimed = new Set<number>();
154+
#owed = new Set<number>();
146155
#onDrop?: (record: SessionStreamRecord, reason: RouterDropReason, route?: string) => void;
147156

148157
constructor(
@@ -197,6 +206,57 @@ export class SessionChannelRouter {
197206
return this.#resumeFrom;
198207
}
199208

209+
/**
210+
* Declare the sequences a continuation boot already read over HTTP and took
211+
* responsibility for dispatching itself, so the live tail's re-read of the
212+
* same records does not answer them a second time.
213+
*
214+
* Two effects, deliberately separate:
215+
*
216+
* - Every claimed sequence is dropped in {@link ingest} however late it
217+
* arrives (`"recovered"`), so a record the boot owns never also becomes a
218+
* router turn. This survives the boot settling it — the boot owns its
219+
* disposition for the whole run, and the tail may re-deliver at any time.
220+
* - Each claimed sequence is also *owed*: it holds the resume floor exactly
221+
* like a queued record would, until the boot {@link settleRecovered}s it.
222+
* Suppressing the second delivery without this would let a turn boundary
223+
* publish a floor past a message the boot has not answered yet, turning a
224+
* duplicate into a dropped message.
225+
*
226+
* `#highestSeq` is advanced over the contiguous run of claimed sequences from
227+
* the current high water, so once every claim is settled the floor can move
228+
* past them even if the tail never re-delivers (a silent tail then degrades
229+
* to a duplicate, never a loss). The advance stops at the first gap so a
230+
* record the boot left for the router is never skipped.
231+
*/
232+
markRecovered(seqNums: Iterable<number>): void {
233+
const nums = [...seqNums].filter((n) => Number.isFinite(n)).sort((a, b) => a - b);
234+
if (nums.length === 0) return;
235+
for (const n of nums) {
236+
this.#claimed.add(n);
237+
this.#owed.add(n);
238+
}
239+
let base = this.#highestSeq ?? nums[0]! - 1;
240+
while (this.#claimed.has(base + 1)) base++;
241+
if (this.#highestSeq === undefined || base > this.#highestSeq) {
242+
this.#highestSeq = base;
243+
}
244+
const maxClaimed = nums[nums.length - 1]!;
245+
if (this.#appliedThrough === undefined || maxClaimed > this.#appliedThrough) {
246+
this.#appliedThrough = maxClaimed;
247+
}
248+
}
249+
250+
/**
251+
* Release a claimed sequence's hold on the resume floor once the boot has
252+
* decided its disposition (dispatched it as a turn, folded it into the seed
253+
* chain, or deliberately dropped it). It stays claimed, so a late tail
254+
* re-delivery is still dropped rather than answered again.
255+
*/
256+
settleRecovered(seqNum: number): void {
257+
this.#owed.delete(seqNum);
258+
}
259+
200260
/**
201261
* Classify one record and act on it. The record's destination is decided
202262
* here, once, and never by whichever consumer happens to be waiting.
@@ -213,6 +273,10 @@ export class SessionChannelRouter {
213273
}
214274
}
215275

276+
if (this.#claimed.has(record.seqNum)) {
277+
return this.#drop(record, "recovered");
278+
}
279+
216280
const kind = this.#kindOf(record.data);
217281
if (kind === undefined) {
218282
return this.#drop(record, "malformed");
@@ -308,6 +372,9 @@ export class SessionChannelRouter {
308372
const pending = state.earliestUnrecovered();
309373
if (pending !== undefined) earliestPending = Math.min(earliestPending, pending);
310374
}
375+
for (const owed of this.#owed) {
376+
if (owed < earliestPending) earliestPending = owed;
377+
}
311378

312379
if (earliestPending === Infinity) return this.#highestSeq;
313380

@@ -518,5 +585,7 @@ export class SessionChannelRouter {
518585
this.#highestSeq = undefined;
519586
this.#resumeFrom = undefined;
520587
this.#appliedThrough = undefined;
588+
this.#claimed.clear();
589+
this.#owed.clear();
521590
}
522591
}

packages/trigger-sdk/src/v3/ai.ts

Lines changed: 55 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -2472,7 +2472,11 @@ async function findSessionInReplayWindowEnd(
24722472
*/
24732473
async function installChatInputRouter(
24742474
chatId: string,
2475-
options?: { fallbackResumeFrom?: number; recoveredThrough?: number; resuming?: boolean }
2475+
options?: {
2476+
fallbackResumeFrom?: number;
2477+
recoveredSeqNums?: readonly number[];
2478+
resuming?: boolean;
2479+
}
24762480
): Promise<SessionChannelRouter> {
24772481
const entry = chatInputRouterEntry(chatId);
24782482
if (entry.attached) return entry.router;
@@ -2503,20 +2507,13 @@ async function installChatInputRouter(
25032507
}
25042508
}
25052509

2506-
// A boot that replayed `.in` itself has already answered everything up to
2507-
// `recoveredThrough`, so the floor has to cover it before the tail opens.
2508-
if (options?.recoveredThrough !== undefined) {
2509-
const recovered = options.recoveredThrough;
2510-
checkpoint.resumeFrom = Math.max(checkpoint.resumeFrom ?? recovered, recovered);
2511-
checkpoint.appliedThrough = Math.max(
2512-
checkpoint.appliedThrough ?? checkpoint.resumeFrom,
2513-
checkpoint.resumeFrom
2514-
);
2515-
}
2516-
25172510
const router = entry.router;
25182511
router.restore(checkpoint);
25192512

2513+
if (options?.recoveredSeqNums && options.recoveredSeqNums.length > 0) {
2514+
router.markRecovered(options.recoveredSeqNums);
2515+
}
2516+
25202517
const floor = router.resumeFrom();
25212518
if (floor !== undefined) {
25222519
sessionStreams.setLastSeqNum(chatId, "in", floor);
@@ -7211,6 +7208,19 @@ function chatAgent<
72117208
// `messagesInput.waitWithIdleTimeout` so recovered turns fire first.
72127209
const bootInjectedQueue: ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>[] =
72137210
[];
7211+
const recoveredSeqByPayload = new WeakMap<
7212+
ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>,
7213+
number
7214+
>();
7215+
const dispatchBootInjected = (): ChatTaskWirePayload<
7216+
TUIMessage,
7217+
inferSchemaIn<TClientDataSchema>
7218+
> => {
7219+
const injected = bootInjectedQueue.shift()!;
7220+
const settledSeq = recoveredSeqByPayload.get(injected);
7221+
if (settledSeq !== undefined) chatInputRouter().settleRecovered(settledSeq);
7222+
return injected;
7223+
};
72147224
const couldHavePriorState = payload.continuation === true || ctx.attempt.number > 1;
72157225

72167226
// `.in` resume cursor, computed at most once per boot. The boot
@@ -7355,18 +7365,11 @@ function chatAgent<
73557365

73567366
// ── session.in router ──────────────────────────────────────────
73577367
//
7358-
// Reads the turn boundary and subscribes in one call. `bootInCursor` is
7359-
// only a fallback: the boot block above may already have resolved a
7360-
// cursor from the snapshot, which is used when the boundary itself
7361-
// carries none. Everything the boot replayed off `.in` is dispatched from
7362-
// `bootInjectedQueue` below, so it goes into the floor here — folded in
7363-
// after the subscription opens, the live tail re-delivers it as a turn.
7364-
const lastRecoveredInSeq =
7365-
replayedInTail.length > 0 ? replayedInTail[replayedInTail.length - 1]!.seqNum : undefined;
7368+
const recoveredSeqNums = replayedInTail.map((r) => r.seqNum);
73667369

73677370
await installChatInputRouter(payload.chatId, {
73687371
fallbackResumeFrom: bootInCursorResolved ? bootInCursor : undefined,
7369-
recoveredThrough: lastRecoveredInSeq,
7372+
recoveredSeqNums,
73707373
resuming: Boolean(payload.continuation) || ctx.attempt.number > 1,
73717374
});
73727375

@@ -7456,7 +7459,7 @@ function chatAgent<
74567459
// branches: at n=1 the orphan partial is dropped and the interrupted
74577460
// user is re-dispatched as a fresh turn instead.
74587461
let seedChain: TUIMessage[];
7459-
let recoveredTurns: TUIMessage[];
7462+
let recoveredEntries: { message: TUIMessage; seqNum: number | undefined }[];
74607463
if (hookChain !== undefined) {
74617464
seedChain = hookChain;
74627465
} else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
@@ -7465,11 +7468,20 @@ function chatAgent<
74657468
seedChain = settledMessages;
74667469
}
74677470
if (hookRecoveredTurns !== undefined) {
7468-
recoveredTurns = hookRecoveredTurns;
7471+
const seqNumByRecoveredId = new Map<string, number>();
7472+
for (const entry of replayedInTail) {
7473+
seqNumByRecoveredId.set(entry.message.id, entry.seqNum);
7474+
}
7475+
recoveredEntries = hookRecoveredTurns.map((message) => ({
7476+
message,
7477+
seqNum: seqNumByRecoveredId.get(message.id),
7478+
}));
74697479
} else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
7470-
recoveredTurns = inFlightUsers.slice(1);
7480+
recoveredEntries = replayedInTail
7481+
.slice(1)
7482+
.map((r) => ({ message: r.message, seqNum: r.seqNum }));
74717483
} else {
7472-
recoveredTurns = inFlightUsers;
7484+
recoveredEntries = replayedInTail.map((r) => ({ message: r.message, seqNum: r.seqNum }));
74737485
}
74747486
// `beforeBoot` errors bubble — the customer opted into blocking
74757487
// persistence and a failure there should fail the run rather than
@@ -7500,12 +7512,13 @@ function chatAgent<
75007512
for (const entry of replayedInTail) {
75017513
metadataById.set(entry.message.id, entry.metadata);
75027514
}
7503-
for (const msg of recoveredTurns) {
7515+
const dispatchedRecoveredSeqs = new Set<number>();
7516+
for (const { message: msg, seqNum } of recoveredEntries) {
75047517
if (wireMessageId && msg.id === wireMessageId) continue;
75057518
const recoveredMetadata = metadataById.has(msg.id)
75067519
? metadataById.get(msg.id)
75077520
: payload.metadata;
7508-
bootInjectedQueue.push({
7521+
const injectedPayload = {
75097522
chatId: payload.chatId,
75107523
sessionId: payload.sessionId,
75117524
metadata: recoveredMetadata,
@@ -7514,7 +7527,17 @@ function chatAgent<
75147527
messageId: msg.id,
75157528
continuation: payload.continuation,
75167529
previousRunId: payload.previousRunId,
7517-
} as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>);
7530+
} as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>;
7531+
bootInjectedQueue.push(injectedPayload);
7532+
if (seqNum !== undefined) {
7533+
recoveredSeqByPayload.set(injectedPayload, seqNum);
7534+
dispatchedRecoveredSeqs.add(seqNum);
7535+
}
7536+
}
7537+
for (const entry of replayedInTail) {
7538+
if (!dispatchedRecoveredSeqs.has(entry.seqNum)) {
7539+
chatInputRouter().settleRecovered(entry.seqNum);
7540+
}
75187541
}
75197542

75207543
accumulatedUIMessages = seedChain;
@@ -7683,7 +7706,7 @@ function chatAgent<
76837706
*/
76847707
let dispatchedRecoveredFirstTurn = false;
76857708
if (preloaded && bootInjectedQueue.length > 0) {
7686-
currentWirePayload = bootInjectedQueue.shift()!;
7709+
currentWirePayload = dispatchBootInjected();
76877710
dispatchedRecoveredFirstTurn = true;
76887711
}
76897712

@@ -7934,7 +7957,7 @@ function chatAgent<
79347957
// waiting on the live session.in. Subsequent recovered turns
79357958
// get drained by the end-of-turn picker below.
79367959
if (bootInjectedQueue.length > 0) {
7937-
currentWirePayload = bootInjectedQueue.shift()!;
7960+
currentWirePayload = dispatchBootInjected();
79387961
} else {
79397962
const effectiveIdleTimeout = idleTimeoutInSeconds ?? payload.idleTimeoutInSeconds;
79407963
const effectiveTurnTimeout =
@@ -9478,7 +9501,7 @@ function chatAgent<
94789501
// produced these from in-flight user messages on session.in
94799502
// that the dead predecessor never acknowledged.
94809503
if (bootInjectedQueue.length > 0) {
9481-
currentWirePayload = bootInjectedQueue.shift()!;
9504+
currentWirePayload = dispatchBootInjected();
94829505
return "continue";
94839506
}
94849507

@@ -9849,7 +9872,7 @@ function chatAgent<
98499872
// recovered turn shouldn't strand the rest of the boot queue
98509873
// until an unrelated live message arrives.
98519874
if (bootInjectedQueue.length > 0) {
9852-
currentWirePayload = bootInjectedQueue.shift()!;
9875+
currentWirePayload = dispatchBootInjected();
98539876
continue;
98549877
}
98559878

0 commit comments

Comments
 (0)