Skip to content

Commit 10a1d68

Browse files
committed
feat(sdk): run tail recovery for every chat.agent and let a transcript storage own the model's context
One condition used to decide three things at boot: whether to read the persisted transcript, whether to replay the session's output tail, and whether to replay unacknowledged input. Registering hydrateMessages switched all three off, so an app that owned its own context also lost crash recovery, and no application can rebuild the tail its dead run had already emitted. The replays and onRecoveryBoot now run for every agent; only the transcript read is skipped for hydrateMessages. The storage can now declare loadContext, which the runtime calls on every turn and action in place of the accumulated transcript, the role hydrateMessages played, while save keeps receiving every change. hydrateMessages is deprecated with a one-time warning, and configuring it together with a storage that has loadContext is an error.
1 parent 7a020c4 commit 10a1d68

3 files changed

Lines changed: 280 additions & 26 deletions

File tree

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

Lines changed: 73 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -5128,6 +5128,18 @@ function isUIMessageStreamable(value: unknown): value is UIMessageStreamable {
51285128
);
51295129
}
51305130

5131+
const warnedHydrateMessagesDeprecated = new Set<string>();
5132+
function warnHydrateMessagesDeprecatedOnce(agentId: string) {
5133+
if (warnedHydrateMessagesDeprecated.has(agentId)) return;
5134+
warnedHydrateMessagesDeprecated.add(agentId);
5135+
console.warn(
5136+
`[chat.agent] \`hydrateMessages\` on "${agentId}" is deprecated. Give the agent a transcript ` +
5137+
"storage instead: `save` receives every change to the conversation and `loadContext` " +
5138+
"lets the application own the model's context, with crash recovery and durable " +
5139+
"compaction that `hydrateMessages` never had."
5140+
);
5141+
}
5142+
51315143
let warnedMissingOnAction = false;
51325144
function warnMissingOnActionOnce() {
51335145
if (warnedMissingOnAction) return;
@@ -5338,8 +5350,9 @@ export type RecoveryPendingToolCall = {
53385350
* `chat.endRun()` with no buffered user messages, fresh chat, OOM retry
53395351
* after a successful turn-complete with no in-flight tail).
53405352
*
5341-
* Does NOT fire when `hydrateMessages` is registered (the customer owns
5342-
* persistence; recovery decisions live in their own DB query).
5353+
* Fires regardless of who owns the model's context. With `hydrateMessages`
5354+
* or a storage `loadContext`, the recovered tail reaches that hook in
5355+
* `previousMessages` on the next turn.
53435356
*/
53445357
export type RecoveryBootEvent<TUIM extends UIMessage = UIMessage> = {
53455358
/** Task run context — same as `task({ run })` second-argument `ctx`. */
@@ -5409,8 +5422,9 @@ export type RecoveryBootResult<TUIM extends UIMessage = UIMessage> = {
54095422
* context, mutate its tool parts to inject synthesized results,
54105423
* collapse history, etc.
54115424
*
5412-
* Ignored when `hydrateMessages` is registered (the hydrate hook
5413-
* runs per-turn and overwrites the chain).
5425+
* With `hydrateMessages` or a storage `loadContext`, this chain is what
5426+
* the hook receives as `previousMessages` on the next turn; the hook's
5427+
* return value is the chain the model sees.
54145428
*/
54155429
chain?: TUIM[];
54165430
/**
@@ -5983,9 +5997,9 @@ export type ChatAgentOptions<
59835997
* continuation after `chat.endRun()` with no buffered user, a fresh
59845998
* chat, or an OOM retry on top of a complete snapshot.
59855999
*
5986-
* Does NOT fire when `hydrateMessages` is registered — that hook owns
5987-
* the per-turn chain and overlapping recovery decisions belong in the
5988-
* customer's DB.
6000+
* Fires regardless of who owns the model's context; a `hydrateMessages`
6001+
* hook or a storage `loadContext` receives the recovered tail in
6002+
* `previousMessages` on the next turn.
59896003
*
59906004
* Defaults (returned when the hook is omitted or returns no field):
59916005
* - With two or more in-flight users, the partial and the user it
@@ -6740,6 +6754,18 @@ function chatAgent<
67406754
...restOptions
67416755
} = options;
67426756

6757+
if (hydrateMessages) {
6758+
const storageAtDefinition = transcriptStorageOverride ?? defaultStorage;
6759+
if (typeof storageAtDefinition.loadContext === "function") {
6760+
throw new Error(
6761+
`chat.agent: "${options.id}" sets \`hydrateMessages\` and uses a transcript storage with ` +
6762+
"`loadContext`. Both would own the model's context; keep one. `hydrateMessages` is " +
6763+
"deprecated, so prefer `loadContext` on the storage."
6764+
);
6765+
}
6766+
warnHydrateMessagesDeprecatedOnce(options.id);
6767+
}
6768+
67436769
const parseClientData = clientDataSchema ? getSchemaParseFn(clientDataSchema) : undefined;
67446770
const parseAction = actionSchema ? getSchemaParseFn(actionSchema) : undefined;
67456771

@@ -6890,6 +6916,22 @@ function chatAgent<
68906916
// swallow errors internally; the agent stays available either way.
68916917
const sessionIdForSnapshot = payload.sessionId ?? payload.chatId;
68926918
const transcriptStorage = transcriptStorageOverride ?? defaultStorage;
6919+
const storageLoadContext = transcriptStorage.loadContext?.bind(transcriptStorage);
6920+
/**
6921+
* Who supplies the model's context each turn: the deprecated
6922+
* `hydrateMessages` hook, the storage's `loadContext`, or (undefined)
6923+
* the runtime's own transcript.
6924+
*/
6925+
const loadContextHook = hydrateMessages
6926+
? (event: HydrateMessagesEvent<inferSchemaOut<TClientDataSchema>, TUIMessage>) =>
6927+
hydrateMessages(event)
6928+
: storageLoadContext
6929+
? (event: HydrateMessagesEvent<inferSchemaOut<TClientDataSchema>, TUIMessage>) =>
6930+
storageLoadContext<TUIMessage>(
6931+
{ chatId: event.chatId, clientData: event.clientData },
6932+
event
6933+
)
6934+
: undefined;
68936935
let transcriptShadow: TranscriptShadow = createTranscriptShadow([]);
68946936
let bootTranscriptState: unknown = null;
68956937
/**
@@ -7063,7 +7105,7 @@ function chatAgent<
70637105
let bootInCursor: number | undefined;
70647106
let bootInCursorResolved = false;
70657107

7066-
if (!hydrateMessages && couldHavePriorState) {
7108+
if (couldHavePriorState) {
70677109
// Single parent span for the whole boot read phase — snapshot
70687110
// read, session.out replay, session.in replay. Per-phase timing
70697111
// + result counts are attributes on the span.
@@ -7073,18 +7115,22 @@ function chatAgent<
70737115
// snapshot read
70747116
const snapStart = Date.now();
70757117
try {
7076-
const loaded = await transcriptStorage.load<TUIMessage>({
7077-
chatId: payload.chatId,
7078-
clientData: bootClientData,
7079-
});
7080-
transcriptShadow = createTranscriptShadow(loaded.messages);
7081-
bootTranscriptState = loaded.state;
7082-
persistedStateSet = loaded.state !== null && loaded.state !== undefined;
7083-
bootSnapshot = {
7084-
messages: loaded.messages,
7085-
lastOutEventId: loaded.cursors?.lastOutEventId,
7086-
lastInEventId: loaded.cursors?.lastInEventId,
7087-
};
7118+
const loaded = hydrateMessages
7119+
? undefined
7120+
: await transcriptStorage.load<TUIMessage>({
7121+
chatId: payload.chatId,
7122+
clientData: bootClientData,
7123+
});
7124+
if (loaded) {
7125+
transcriptShadow = createTranscriptShadow(loaded.messages);
7126+
bootTranscriptState = loaded.state;
7127+
persistedStateSet = loaded.state !== null && loaded.state !== undefined;
7128+
bootSnapshot = {
7129+
messages: loaded.messages,
7130+
lastOutEventId: loaded.cursors?.lastOutEventId,
7131+
lastInEventId: loaded.cursors?.lastInEventId,
7132+
};
7133+
}
70887134
} catch (error) {
70897135
logger.warn("chat.agent: transcript load failed; continuing from the stream tail", {
70907136
error: error instanceof Error ? error.message : String(error),
@@ -7222,7 +7268,7 @@ function chatAgent<
72227268
});
72237269

72247270
// ── Recovery boot + chain reconstruction ────────────────────────
7225-
if (!hydrateMessages) {
7271+
{
72267272
const settledMessages = mergeByIdReplaceWins<TUIMessage>(
72277273
(bootSnapshot?.messages as TUIMessage[]) ?? [],
72287274
replayedSettled
@@ -7391,6 +7437,7 @@ function chatAgent<
73917437
// and it's safe because the route handler isn't subject to the
73927438
// `/in/append` 512 KiB cap.
73937439
if (
7440+
!hydrateMessages &&
73947441
accumulatedUIMessages.length === 0 &&
73957442
payload.trigger === "handover-prepare" &&
73967443
Array.isArray(payload.headStartMessages) &&
@@ -8077,11 +8124,11 @@ function chatAgent<
80778124
: currentWirePayload.action;
80788125

80798126
// Hydrate messages from backend if configured
8080-
if (hydrateMessages) {
8127+
if (loadContextHook) {
80818128
const hydrated = await tracer.startActiveSpan(
80828129
"hydrateMessages()",
80838130
async () => {
8084-
return hydrateMessages({
8131+
return loadContextHook({
80858132
chatId: currentWirePayload.chatId,
80868133
turn,
80878134
trigger: "action",
@@ -8169,7 +8216,7 @@ function chatAgent<
81698216
// incoming messages instead (gated on the pending handover).
81708217
if (
81718218
turn === 0 &&
8172-
hydrateMessages &&
8219+
loadContextHook &&
81738220
cleanedUIMessages.length === 0 &&
81748221
(locals.get(chatHandoverPartialKey)?.length ?? 0) > 0 &&
81758222
Array.isArray(payload.headStartMessages) &&
@@ -8206,7 +8253,7 @@ function chatAgent<
82068253
)) as TUIMessage[];
82078254
}
82088255

8209-
if (hydrateMessages) {
8256+
if (loadContextHook) {
82108257
// Snapshot the ids the accumulator knew BEFORE this
82118258
// turn ran — used below to decide whether an
82128259
// incoming wire message is genuinely new or just a
@@ -8229,7 +8276,7 @@ function chatAgent<
82298276
const hydrated = await tracer.startActiveSpan(
82308277
"hydrateMessages()",
82318278
async () => {
8232-
return hydrateMessages({
8279+
return loadContextHook({
82338280
chatId: currentWirePayload.chatId,
82348281
turn,
82358282
trigger: currentWirePayload.trigger as

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

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,18 +77,42 @@ type TranscriptLoadResult<TUIMessage extends UIMessage = UIMessage> = {
7777
nextCursor?: string;
7878
};
7979

80+
/** What `loadContext` receives on every turn and action. */
81+
type LoadContextEvent<TClientData = unknown, TUIMessage extends UIMessage = UIMessage> = {
82+
chatId: string;
83+
/** The turn number (0-indexed). */
84+
turn: number;
85+
trigger: "submit-message" | "regenerate-message" | "action";
86+
/** The messages the frontend sent for this turn. Empty for actions. */
87+
incomingMessages: TUIMessage[];
88+
/** The runtime's transcript before this turn, including any tail it recovered. */
89+
previousMessages: TUIMessage[];
90+
clientData?: TClientData;
91+
continuation: boolean;
92+
previousRunId?: string;
93+
};
94+
8095
/**
8196
* A persistence adapter for a `chat.agent` transcript. The runtime calls
8297
* `load` once at a continuation boot and `save` after every change to the
8398
* conversation. Both are best-effort from the runtime's point of view: an
8499
* error is logged and the turn continues.
100+
*
101+
* `loadContext` is optional. Its presence declares that the application
102+
* owns the model's context: the runtime calls it on every turn and action
103+
* and uses what it returns as the conversation, instead of the transcript
104+
* it accumulated. Tail recovery still runs and `save` is still called.
85105
*/
86106
export type TranscriptStorage<TClientData = unknown> = {
87107
load<TUIMessage extends UIMessage = UIMessage>(
88108
scope: TranscriptScope<TClientData>,
89109
opts?: TranscriptLoadOptions
90110
): Promise<TranscriptLoadResult<TUIMessage>>;
91111
save(ctx: TranscriptStorageContext<TClientData>, changeset: TranscriptChangeset): Promise<void>;
112+
loadContext?<TUIMessage extends UIMessage = UIMessage>(
113+
scope: TranscriptScope<TClientData>,
114+
event: LoadContextEvent<TClientData, TUIMessage>
115+
): Promise<TUIMessage[]> | TUIMessage[];
92116
};
93117

94118
/** An in-memory transcript: ordered entries plus the opaque state record. */

0 commit comments

Comments
 (0)