Skip to content

Commit 52addef

Browse files
committed
fix(core,sdk): stop chat.agent losing messages during recovery
Rebased onto the transcript-storage stack. A continuation boot claims recovered session.in seqNums on the router and holds the resume cursor behind each until the boot settles it, so suppressing the tail's re-answer no longer advances the cursor past an un-answered message. Adds a changeset and router/boot tests.
1 parent 1c08e6c commit 52addef

5 files changed

Lines changed: 319 additions & 33 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
"@trigger.dev/sdk": patch
3+
"@trigger.dev/core": patch
4+
---
5+
6+
`chat.agent`: a run that recovers a session with more than one in-flight user message no longer drops the unanswered ones if it restarts mid-recovery. Recovered messages now hold the resume cursor until each has been answered, so a restart re-answers the rest instead of resuming past them. Previously the cursor could advance past messages that were only held in memory, so a crash before they were dispatched lost them.

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

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -573,3 +573,76 @@ 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("holds the floor below an unclaimed gap while later claims are owed and the tail is silent", () => {
633+
const r = router();
634+
r.restore({ resumeFrom: 0 });
635+
r.markRecovered([1, 2, 5, 6]);
636+
r.settleRecovered(1);
637+
r.settleRecovered(2);
638+
expect(r.resumeFloor()).toBe(2);
639+
});
640+
641+
it("clears claims and owed records on reset", () => {
642+
const r = router();
643+
r.restore({ resumeFrom: 0 });
644+
r.markRecovered([1, 2]);
645+
r.reset();
646+
expect(r.ingest(rec(1, "message"))).toEqual({ action: "queue", route: "messages" });
647+
});
648+
});

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
@@ -2319,7 +2319,11 @@ async function findSessionInReplayWindowEnd(
23192319
*/
23202320
async function installChatInputRouter(
23212321
chatId: string,
2322-
options?: { fallbackResumeFrom?: number; recoveredThrough?: number; resuming?: boolean }
2322+
options?: {
2323+
fallbackResumeFrom?: number;
2324+
recoveredSeqNums?: readonly number[];
2325+
resuming?: boolean;
2326+
}
23232327
): Promise<SessionChannelRouter> {
23242328
const entry = chatInputRouterEntry(chatId);
23252329
if (entry.attached) return entry.router;
@@ -2350,20 +2354,13 @@ async function installChatInputRouter(
23502354
}
23512355
}
23522356

2353-
// A boot that replayed `.in` itself has already answered everything up to
2354-
// `recoveredThrough`, so the floor has to cover it before the tail opens.
2355-
if (options?.recoveredThrough !== undefined) {
2356-
const recovered = options.recoveredThrough;
2357-
checkpoint.resumeFrom = Math.max(checkpoint.resumeFrom ?? recovered, recovered);
2358-
checkpoint.appliedThrough = Math.max(
2359-
checkpoint.appliedThrough ?? checkpoint.resumeFrom,
2360-
checkpoint.resumeFrom
2361-
);
2362-
}
2363-
23642357
const router = entry.router;
23652358
router.restore(checkpoint);
23662359

2360+
if (options?.recoveredSeqNums && options.recoveredSeqNums.length > 0) {
2361+
router.markRecovered(options.recoveredSeqNums);
2362+
}
2363+
23672364
const floor = router.resumeFrom();
23682365
if (floor !== undefined) {
23692366
sessionStreams.setLastSeqNum(chatId, "in", floor);
@@ -7271,6 +7268,19 @@ function chatAgent<
72717268
// `messagesInput.waitWithIdleTimeout` so recovered turns fire first.
72727269
const bootInjectedQueue: ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>[] =
72737270
[];
7271+
const recoveredSeqByPayload = new WeakMap<
7272+
ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>,
7273+
number
7274+
>();
7275+
const dispatchBootInjected = (): ChatTaskWirePayload<
7276+
TUIMessage,
7277+
inferSchemaIn<TClientDataSchema>
7278+
> => {
7279+
const injected = bootInjectedQueue.shift()!;
7280+
const settledSeq = recoveredSeqByPayload.get(injected);
7281+
if (settledSeq !== undefined) chatInputRouter().settleRecovered(settledSeq);
7282+
return injected;
7283+
};
72747284
const couldHavePriorState = payload.continuation === true || ctx.attempt.number > 1;
72757285

72767286
// `.in` resume cursor, computed at most once per boot. The boot
@@ -7436,18 +7446,11 @@ function chatAgent<
74367446

74377447
// ── session.in router ──────────────────────────────────────────
74387448
//
7439-
// Reads the turn boundary and subscribes in one call. `bootInCursor` is
7440-
// only a fallback: the boot block above may already have resolved a
7441-
// cursor from the snapshot, which is used when the boundary itself
7442-
// carries none. Everything the boot replayed off `.in` is dispatched from
7443-
// `bootInjectedQueue` below, so it goes into the floor here — folded in
7444-
// after the subscription opens, the live tail re-delivers it as a turn.
7445-
const lastRecoveredInSeq =
7446-
replayedInTail.length > 0 ? replayedInTail[replayedInTail.length - 1]!.seqNum : undefined;
7449+
const recoveredSeqNums = replayedInTail.map((r) => r.seqNum);
74477450

74487451
await installChatInputRouter(payload.chatId, {
74497452
fallbackResumeFrom: bootInCursorResolved ? bootInCursor : undefined,
7450-
recoveredThrough: lastRecoveredInSeq,
7453+
recoveredSeqNums,
74517454
resuming: Boolean(payload.continuation) || ctx.attempt.number > 1,
74527455
});
74537456

@@ -7537,7 +7540,7 @@ function chatAgent<
75377540
// branches: at n=1 the orphan partial is dropped and the interrupted
75387541
// user is re-dispatched as a fresh turn instead.
75397542
let seedChain: TUIMessage[];
7540-
let recoveredTurns: TUIMessage[];
7543+
let recoveredEntries: { message: TUIMessage; seqNum: number | undefined }[];
75417544
if (hookChain !== undefined) {
75427545
seedChain = hookChain;
75437546
} else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
@@ -7546,11 +7549,20 @@ function chatAgent<
75467549
seedChain = settledMessages;
75477550
}
75487551
if (hookRecoveredTurns !== undefined) {
7549-
recoveredTurns = hookRecoveredTurns;
7552+
const seqNumByRecoveredId = new Map<string, number>();
7553+
for (const entry of replayedInTail) {
7554+
seqNumByRecoveredId.set(entry.message.id, entry.seqNum);
7555+
}
7556+
recoveredEntries = hookRecoveredTurns.map((message) => ({
7557+
message,
7558+
seqNum: seqNumByRecoveredId.get(message.id),
7559+
}));
75507560
} else if (partialAssistant !== undefined && inFlightUsers.length > 1) {
7551-
recoveredTurns = inFlightUsers.slice(1);
7561+
recoveredEntries = replayedInTail
7562+
.slice(1)
7563+
.map((r) => ({ message: r.message, seqNum: r.seqNum }));
75527564
} else {
7553-
recoveredTurns = inFlightUsers;
7565+
recoveredEntries = replayedInTail.map((r) => ({ message: r.message, seqNum: r.seqNum }));
75547566
}
75557567
// `beforeBoot` errors bubble — the customer opted into blocking
75567568
// persistence and a failure there should fail the run rather than
@@ -7581,12 +7593,13 @@ function chatAgent<
75817593
for (const entry of replayedInTail) {
75827594
metadataById.set(entry.message.id, entry.metadata);
75837595
}
7584-
for (const msg of recoveredTurns) {
7596+
const dispatchedRecoveredSeqs = new Set<number>();
7597+
for (const { message: msg, seqNum } of recoveredEntries) {
75857598
if (wireMessageId && msg.id === wireMessageId) continue;
75867599
const recoveredMetadata = metadataById.has(msg.id)
75877600
? metadataById.get(msg.id)
75887601
: payload.metadata;
7589-
bootInjectedQueue.push({
7602+
const injectedPayload = {
75907603
chatId: payload.chatId,
75917604
sessionId: payload.sessionId,
75927605
metadata: recoveredMetadata,
@@ -7595,7 +7608,17 @@ function chatAgent<
75957608
messageId: msg.id,
75967609
continuation: payload.continuation,
75977610
previousRunId: payload.previousRunId,
7598-
} as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>);
7611+
} as ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>;
7612+
bootInjectedQueue.push(injectedPayload);
7613+
if (seqNum !== undefined) {
7614+
recoveredSeqByPayload.set(injectedPayload, seqNum);
7615+
dispatchedRecoveredSeqs.add(seqNum);
7616+
}
7617+
}
7618+
for (const entry of replayedInTail) {
7619+
if (!dispatchedRecoveredSeqs.has(entry.seqNum)) {
7620+
chatInputRouter().settleRecovered(entry.seqNum);
7621+
}
75997622
}
76007623

76017624
accumulatedUIMessages = seedChain;
@@ -7779,7 +7802,7 @@ function chatAgent<
77797802
*/
77807803
let dispatchedRecoveredFirstTurn = false;
77817804
if (preloaded && bootInjectedQueue.length > 0) {
7782-
currentWirePayload = bootInjectedQueue.shift()!;
7805+
currentWirePayload = dispatchBootInjected();
77837806
dispatchedRecoveredFirstTurn = true;
77847807
}
77857808

@@ -8030,7 +8053,7 @@ function chatAgent<
80308053
// waiting on the live session.in. Subsequent recovered turns
80318054
// get drained by the end-of-turn picker below.
80328055
if (bootInjectedQueue.length > 0) {
8033-
currentWirePayload = bootInjectedQueue.shift()!;
8056+
currentWirePayload = dispatchBootInjected();
80348057
} else {
80358058
const effectiveIdleTimeout = idleTimeoutInSeconds ?? payload.idleTimeoutInSeconds;
80368059
const effectiveTurnTimeout =
@@ -9611,7 +9634,7 @@ function chatAgent<
96119634
// produced these from in-flight user messages on session.in
96129635
// that the dead predecessor never acknowledged.
96139636
if (bootInjectedQueue.length > 0) {
9614-
currentWirePayload = bootInjectedQueue.shift()!;
9637+
currentWirePayload = dispatchBootInjected();
96159638
return "continue";
96169639
}
96179640

@@ -9987,7 +10010,7 @@ function chatAgent<
998710010
// recovered turn shouldn't strand the rest of the boot queue
998810011
// until an unrelated live message arrives.
998910012
if (bootInjectedQueue.length > 0) {
9990-
currentWirePayload = bootInjectedQueue.shift()!;
10013+
currentWirePayload = dispatchBootInjected();
999110014
continue;
999210015
}
999310016

0 commit comments

Comments
 (0)