Skip to content

Commit f2a1179

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 516c8de commit f2a1179

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
/**
@@ -7064,7 +7106,7 @@ function chatAgent<
70647106
let bootInCursor: number | undefined;
70657107
let bootInCursorResolved = false;
70667108

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

72257271
// ── Recovery boot + chain reconstruction ────────────────────────
7226-
if (!hydrateMessages) {
7272+
{
72277273
const settledMessages = mergeByIdReplaceWins<TUIMessage>(
72287274
(bootSnapshot?.messages as TUIMessage[]) ?? [],
72297275
replayedSettled
@@ -7392,6 +7438,7 @@ function chatAgent<
73927438
// and it's safe because the route handler isn't subject to the
73937439
// `/in/append` 512 KiB cap.
73947440
if (
7441+
!hydrateMessages &&
73957442
accumulatedUIMessages.length === 0 &&
73967443
payload.trigger === "handover-prepare" &&
73977444
Array.isArray(payload.headStartMessages) &&
@@ -8078,11 +8125,11 @@ function chatAgent<
80788125
: currentWirePayload.action;
80798126

80808127
// Hydrate messages from backend if configured
8081-
if (hydrateMessages) {
8128+
if (loadContextHook) {
80828129
const hydrated = await tracer.startActiveSpan(
80838130
"hydrateMessages()",
80848131
async () => {
8085-
return hydrateMessages({
8132+
return loadContextHook({
80868133
chatId: currentWirePayload.chatId,
80878134
turn,
80888135
trigger: "action",
@@ -8170,7 +8217,7 @@ function chatAgent<
81708217
// incoming messages instead (gated on the pending handover).
81718218
if (
81728219
turn === 0 &&
8173-
hydrateMessages &&
8220+
loadContextHook &&
81748221
cleanedUIMessages.length === 0 &&
81758222
(locals.get(chatHandoverPartialKey)?.length ?? 0) > 0 &&
81768223
Array.isArray(payload.headStartMessages) &&
@@ -8207,7 +8254,7 @@ function chatAgent<
82078254
)) as TUIMessage[];
82088255
}
82098256

8210-
if (hydrateMessages) {
8257+
if (loadContextHook) {
82118258
// Snapshot the ids the accumulator knew BEFORE this
82128259
// turn ran — used below to decide whether an
82138260
// incoming wire message is genuinely new or just a
@@ -8230,7 +8277,7 @@ function chatAgent<
82308277
const hydrated = await tracer.startActiveSpan(
82318278
"hydrateMessages()",
82328279
async () => {
8233-
return hydrateMessages({
8280+
return loadContextHook({
82348281
chatId: currentWirePayload.chatId,
82358282
turn,
82368283
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)