Skip to content

Commit 18c2b11

Browse files
committed
fix(core): prevent streamText abort unhandled rejections
1 parent 41cad53 commit 18c2b11

3 files changed

Lines changed: 143 additions & 26 deletions

File tree

.changeset/quiet-stream-abort.md

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
---
2+
"@voltagent/core": patch
3+
---
4+
5+
fix: prevent unhandled rejections when aborting `Agent.streamText()` streams
6+
7+
## The Problem
8+
9+
`Agent.streamText()` eagerly read the AI SDK result getters for `text`, `usage`, and `finishReason` while constructing VoltAgent's wrapped result. In AI SDK v6 these fields are lazy promises, so reading them early could materialize promises that the caller never consumes.
10+
11+
When a caller only consumed the UI/full stream and aborted the run, those unconsumed promises could reject globally as `unhandledRejection` events.
12+
13+
## The Solution
14+
15+
VoltAgent now preserves the lazy getter behavior for `text`, `usage`, and `finishReason`. The sanitized text promise is also created only when `result.text` is accessed.
16+
17+
## Impact
18+
19+
- Aborting a consumed `streamText()` stream no longer emits unhandled rejections for unconsumed result fields
20+
- Callers using only `toUIMessageStream()`, `toUIMessageStreamResponse()`, `fullStream`, or `textStream` do not need to attach defensive `.catch()` handlers to `text`, `usage`, or `finishReason`
21+
- Matches AI SDK v6's lazy stream result contract more closely

packages/core/src/agent/agent.spec.ts

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1008,6 +1008,82 @@ Use pandas and summarize findings.`.split("\n"),
10081008
expect(text).toBe("Streamed response");
10091009
});
10101010

1011+
it("does not eagerly materialize lazy stream result promises", async () => {
1012+
const agent = new Agent({
1013+
name: "TestAgent",
1014+
instructions: "You are a helpful assistant",
1015+
model: mockModel as any,
1016+
});
1017+
1018+
const getterAccesses = {
1019+
text: 0,
1020+
usage: 0,
1021+
finishReason: 0,
1022+
};
1023+
1024+
const mockStream = {
1025+
get text() {
1026+
getterAccesses.text += 1;
1027+
return Promise.resolve("Streamed response");
1028+
},
1029+
textStream: (async function* () {
1030+
yield "Streamed response";
1031+
})(),
1032+
fullStream: toAsyncIterableStream(
1033+
convertArrayToReadableStream([
1034+
{
1035+
type: "text-delta" as const,
1036+
id: "text-1",
1037+
delta: "Streamed response",
1038+
text: "Streamed response",
1039+
},
1040+
]),
1041+
),
1042+
get usage() {
1043+
getterAccesses.usage += 1;
1044+
return Promise.resolve({
1045+
inputTokens: 10,
1046+
outputTokens: 5,
1047+
totalTokens: 15,
1048+
});
1049+
},
1050+
get finishReason() {
1051+
getterAccesses.finishReason += 1;
1052+
return Promise.resolve("stop");
1053+
},
1054+
warnings: [],
1055+
toUIMessageStream: vi.fn(),
1056+
toUIMessageStreamResponse: vi.fn(),
1057+
pipeUIMessageStreamToResponse: vi.fn(),
1058+
pipeTextStreamToResponse: vi.fn(),
1059+
toTextStreamResponse: vi.fn(),
1060+
partialOutputStream: undefined,
1061+
};
1062+
1063+
vi.mocked(ai.streamText).mockReturnValue(mockStream as any);
1064+
1065+
const result = await agent.streamText("Stream this");
1066+
1067+
expect(getterAccesses).toEqual({
1068+
text: 0,
1069+
usage: 0,
1070+
finishReason: 0,
1071+
});
1072+
1073+
await expect(result.text).resolves.toBe("Streamed response");
1074+
await expect(result.usage).resolves.toEqual({
1075+
inputTokens: 10,
1076+
outputTokens: 5,
1077+
totalTokens: 15,
1078+
});
1079+
await expect(result.finishReason).resolves.toBe("stop");
1080+
expect(getterAccesses).toEqual({
1081+
text: 1,
1082+
usage: 1,
1083+
finishReason: 1,
1084+
});
1085+
});
1086+
10111087
it("pre-creates streaming message ids and forwards them to UI streams", async () => {
10121088
const agent = new Agent({
10131089
name: "TestAgent",

packages/core/src/agent/agent.ts

Lines changed: 46 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -1961,7 +1961,7 @@ export class Agent {
19611961
const guardrailStreamingEnabled = guardrailSet.output.length > 0;
19621962

19631963
let guardrailPipeline: GuardrailPipeline | null = null;
1964-
let sanitizedTextPromise!: PromiseLike<string>;
1964+
let sanitizedTextPromise: Promise<string> | undefined;
19651965
const { result, modelName: effectiveModelName } = await this.executeWithModelFallback({
19661966
oc,
19671967
operation: "streamText",
@@ -2190,7 +2190,7 @@ export class Agent {
21902190
finalText = bailedResult.response;
21912191
}
21922192
} else if (guardrailPipeline) {
2193-
finalText = await sanitizedTextPromise;
2193+
finalText = await getSanitizedTextPromise();
21942194
} else if (guardrailSet.output.length > 0) {
21952195
finalText = await executeOutputGuardrails({
21962196
output: finalResult.text,
@@ -2512,31 +2512,28 @@ export class Agent {
25122512
? createBaseFullStream()
25132513
: undefined;
25142514

2515-
if (guardrailStreamingEnabled) {
2516-
guardrailPipeline = createGuardrailPipeline(
2517-
baseFullStreamForPipeline as AsyncIterable<VoltAgentTextStreamPart>,
2518-
result.textStream,
2519-
guardrailContext,
2520-
);
2521-
sanitizedTextPromise = guardrailPipeline.finalizePromise.then(async () => {
2522-
const sanitized = guardrailPipeline?.runner?.getSanitizedText();
2523-
if (typeof sanitized === "string" && sanitized.length > 0) {
2524-
return sanitized;
2525-
}
2526-
// Wait for AI SDK text first (stream must complete)
2527-
const aiSdkText = await result.text;
2515+
const createSanitizedTextPromise = (): Promise<string> => {
2516+
if (guardrailPipeline) {
2517+
return guardrailPipeline.finalizePromise.then(async () => {
2518+
const sanitized = guardrailPipeline?.runner?.getSanitizedText();
2519+
if (typeof sanitized === "string" && sanitized.length > 0) {
2520+
return sanitized;
2521+
}
2522+
// Wait for AI SDK text first (stream must complete)
2523+
const aiSdkText = await result.text;
2524+
2525+
// NOW check for bailed result (set during stream processing)
2526+
const bailedResult = oc.systemContext.get("bailedResult") as
2527+
| { agentName: string; response: string }
2528+
| undefined;
2529+
return bailedResult?.response || aiSdkText;
2530+
});
2531+
}
25282532

2529-
// NOW check for bailed result (set during stream processing)
2530-
const bailedResult = oc.systemContext.get("bailedResult") as
2531-
| { agentName: string; response: string }
2532-
| undefined;
2533-
return bailedResult?.response || aiSdkText;
2534-
});
2535-
} else {
25362533
// Wrap result.text with a bail check
25372534
// IMPORTANT: Wait for AI SDK text first (stream must complete/abort)
25382535
// This ensures createStepHandler has processed tool results and set bailedResult
2539-
sanitizedTextPromise = result.text.then((aiSdkText) => {
2536+
return Promise.resolve(result.text).then((aiSdkText) => {
25402537
// NOW check if bailed (set by createStepHandler during stream processing)
25412538
const bailedResult = oc.systemContext.get("bailedResult") as
25422539
| { agentName: string; response: string }
@@ -2545,6 +2542,23 @@ export class Agent {
25452542
// Return bailed subagent's result instead of supervisor's (if bailed)
25462543
return bailedResult?.response || aiSdkText;
25472544
});
2545+
};
2546+
2547+
const getSanitizedTextPromise = (): Promise<string> => {
2548+
sanitizedTextPromise ??= createSanitizedTextPromise();
2549+
return sanitizedTextPromise;
2550+
};
2551+
2552+
if (guardrailStreamingEnabled) {
2553+
guardrailPipeline = createGuardrailPipeline(
2554+
baseFullStreamForPipeline as AsyncIterable<VoltAgentTextStreamPart>,
2555+
result.textStream,
2556+
guardrailContext,
2557+
);
2558+
void guardrailPipeline.finalizePromise.catch(() => {
2559+
// The guarded streams surface this error to their consumers. Keep the
2560+
// internal finalizer promise from leaking when text is never requested.
2561+
});
25482562
}
25492563

25502564
const getGuardrailAwareFullStream = (): AsyncIterable<VoltAgentTextStreamPart> => {
@@ -2676,15 +2690,21 @@ export class Agent {
26762690

26772691
// Create a wrapper that includes context and delegates to the original result
26782692
const resultWithContext: StreamTextResultWithContext = {
2679-
text: sanitizedTextPromise,
2693+
get text() {
2694+
return getSanitizedTextPromise();
2695+
},
26802696
get textStream() {
26812697
return getGuardrailAwareTextStream();
26822698
},
26832699
get fullStream() {
26842700
return getGuardrailAwareFullStream();
26852701
},
2686-
usage: result.usage,
2687-
finishReason: result.finishReason,
2702+
get usage() {
2703+
return result.usage;
2704+
},
2705+
get finishReason() {
2706+
return result.finishReason;
2707+
},
26882708
get partialOutputStream() {
26892709
return result.partialOutputStream;
26902710
},

0 commit comments

Comments
 (0)