Skip to content

Commit 794a336

Browse files
psh4607codex
andcommitted
fix(server-utils): Preserve LangGraph pipe lifecycle
Record source chunks and complete spans across pipeTo and pipeThrough consumption. Co-Authored-By: OpenAI Codex <codex@openai.com>
1 parent bbd119d commit 794a336

2 files changed

Lines changed: 129 additions & 18 deletions

File tree

‎packages/server-utils/src/ai/langgraph/streaming.ts‎

Lines changed: 41 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -172,31 +172,54 @@ function instrumentReadableStream(stream: InstrumentableReadableStream, span: Sp
172172

173173
if (stream.pipeTo) {
174174
const originalPipeTo = stream.pipeTo.bind(stream);
175-
const originalPipeThrough = stream.pipeThrough?.bind(stream);
176-
stream.pipeTo = (destination: WritableStream<unknown>, options?: StreamPipeOptions): Promise<void> => {
175+
const instrumentedPipeTo = (destination: WritableStream<unknown>, options?: StreamPipeOptions): Promise<void> => {
176+
let destinationWriter: WritableStreamDefaultWriter<unknown>;
177+
try {
178+
destinationWriter = destination.getWriter();
179+
} catch (error) {
180+
lifecycle.fail();
181+
return Promise.reject(error);
182+
}
183+
184+
const recordingDestination = new WritableStream<unknown>({
185+
write(chunk) {
186+
lifecycle.recordChunk(chunk);
187+
return destinationWriter.write(chunk);
188+
},
189+
close() {
190+
return destinationWriter.close();
191+
},
192+
abort(reason) {
193+
return destinationWriter.abort(reason);
194+
},
195+
});
196+
177197
let pipePromise: Promise<void>;
178198
try {
179-
pipePromise = withActiveSpan(span, () => {
180-
if (!originalPipeThrough) {
181-
return originalPipeTo(destination, options);
182-
}
183-
184-
const passthrough = new TransformStream<unknown, unknown>({
185-
transform(chunk, controller) {
186-
lifecycle.recordChunk(chunk);
187-
controller.enqueue(chunk);
188-
},
189-
});
190-
const outputStream: ReadableStream<unknown> = originalPipeThrough(passthrough);
191-
return outputStream.pipeTo(destination, options);
192-
});
199+
pipePromise = withActiveSpan(span, () => originalPipeTo(recordingDestination, options));
193200
} catch (error) {
201+
destinationWriter.releaseLock();
194202
lifecycle.fail();
195-
throw error;
203+
return Promise.reject(error);
196204
}
197205

198-
return completeWithLifecycle(pipePromise, lifecycle);
206+
return completeWithLifecycle(pipePromise, lifecycle).finally(() => {
207+
destinationWriter.releaseLock();
208+
});
199209
};
210+
211+
stream.pipeTo = instrumentedPipeTo;
212+
213+
if (stream.pipeThrough) {
214+
stream.pipeThrough = (
215+
transform: ReadableWritablePair<unknown, unknown>,
216+
options?: StreamPipeOptions,
217+
): ReadableStream<unknown> => {
218+
// pipeThrough exposes pipeline failures through the returned readable instead of its internal promise.
219+
void instrumentedPipeTo(transform.writable, options).catch(() => {});
220+
return transform.readable;
221+
};
222+
}
200223
}
201224
}
202225

‎packages/server-utils/test/ai/lib/tracing/langgraph-stream.test.ts‎

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -197,6 +197,94 @@ describe('instrumentStateGraphCompile stream instrumentation', () => {
197197
expect(spanEnd).toHaveBeenCalledTimes(1);
198198
});
199199

200+
it('records response attributes when pipeThrough is unavailable', async () => {
201+
const responseMessage = { role: 'assistant', content: 'Clear skies' };
202+
const stream = new TestLangGraphStream([{ agent: { messages: [responseMessage] } }]);
203+
Object.defineProperty(stream, 'pipeThrough', { value: undefined });
204+
const compiledGraph = { stream: vi.fn().mockResolvedValue(stream) };
205+
const compile = instrumentStateGraphCompile(() => compiledGraph, { recordOutputs: true });
206+
const graph = compile() as { stream: () => Promise<TestLangGraphStream<unknown>> };
207+
208+
await (await graph.stream()).pipeTo(new WritableStream());
209+
210+
expect(spanSetAttribute).toHaveBeenCalledWith(
211+
GEN_AI_RESPONSE_TEXT,
212+
'[{"role":"assistant","content":"Clear skies"}]',
213+
);
214+
expect(spanEnd).toHaveBeenCalledTimes(1);
215+
});
216+
217+
it('does not re-enter custom pipeThrough implementations from pipeTo', async () => {
218+
const chunks: string[] = [];
219+
const stream = new TestLangGraphStream(['first update', 'final update']);
220+
Object.defineProperty(stream, 'pipeThrough', {
221+
configurable: true,
222+
value(
223+
this: TestLangGraphStream<string>,
224+
transform: ReadableWritablePair<unknown, string>,
225+
options?: StreamPipeOptions,
226+
): ReadableStream<unknown> {
227+
void this.pipeTo(transform.writable, options).catch(() => {});
228+
return transform.readable;
229+
},
230+
writable: true,
231+
});
232+
const compiledGraph = { stream: vi.fn().mockResolvedValue(stream) };
233+
const compile = instrumentStateGraphCompile(() => compiledGraph, {});
234+
const graph = compile() as { stream: () => Promise<TestLangGraphStream<string>> };
235+
236+
await (
237+
await graph.stream()
238+
).pipeTo(
239+
new WritableStream({
240+
write(chunk) {
241+
chunks.push(chunk);
242+
},
243+
}),
244+
);
245+
246+
expect(chunks).toEqual(['first update', 'final update']);
247+
expect(spanEnd).toHaveBeenCalledTimes(1);
248+
});
249+
250+
it('forwards pipe options to the source pipeTo operation', async () => {
251+
const stream = new TestLangGraphStream(['first update']);
252+
const originalPipeTo = stream.pipeTo.bind(stream);
253+
const sourcePipeTo = vi.fn((destination: WritableStream<string>, options?: StreamPipeOptions) =>
254+
originalPipeTo(destination, options),
255+
);
256+
Object.defineProperty(stream, 'pipeTo', { configurable: true, value: sourcePipeTo, writable: true });
257+
const compiledGraph = { stream: vi.fn().mockResolvedValue(stream) };
258+
const compile = instrumentStateGraphCompile(() => compiledGraph, {});
259+
const graph = compile() as { stream: () => Promise<TestLangGraphStream<string>> };
260+
const destination = new WritableStream<string>();
261+
const options = { preventClose: true, signal: new AbortController().signal };
262+
263+
await (await graph.stream()).pipeTo(destination, options);
264+
265+
expect(sourcePipeTo).toHaveBeenCalledWith(expect.any(WritableStream), options);
266+
expect(destination.locked).toBe(false);
267+
expect(spanEnd).toHaveBeenCalledTimes(1);
268+
});
269+
270+
it('ends the span when pipeThrough output is consumed', async () => {
271+
const responseMessage = { role: 'assistant', content: 'Clear skies' };
272+
const stream = new TestLangGraphStream([{ agent: { messages: [responseMessage] } }]);
273+
const compiledGraph = { stream: vi.fn().mockResolvedValue(stream) };
274+
const compile = instrumentStateGraphCompile(() => compiledGraph, { recordOutputs: true });
275+
const graph = compile() as { stream: () => Promise<TestLangGraphStream<unknown>> };
276+
277+
const transformed = (await graph.stream()).pipeThrough(new TransformStream());
278+
await transformed.pipeTo(new WritableStream());
279+
await Promise.resolve();
280+
281+
expect(spanSetAttribute).toHaveBeenCalledWith(
282+
GEN_AI_RESPONSE_TEXT,
283+
'[{"role":"assistant","content":"Clear skies"}]',
284+
);
285+
expect(spanEnd).toHaveBeenCalledTimes(1);
286+
});
287+
200288
it('ends the span when direct next calls consume the stream', async () => {
201289
const stream = new TestLangGraphStream(['first update']);
202290
const compiledGraph = { stream: vi.fn().mockResolvedValue(stream) };

0 commit comments

Comments
 (0)