diff --git a/libs/sdk-js/src/client.mts b/libs/sdk-js/src/client.mts index 459e9e6b1..9081749af 100644 --- a/libs/sdk-js/src/client.mts +++ b/libs/sdk-js/src/client.mts @@ -654,6 +654,29 @@ export class RunsClient extends BaseClient { }); } + /** + * Create a batch of stateless background runs. + * + * @param payloads An array of payloads for creating runs. + * @returns An array of created runs. + */ + async createBatch( + payloads: (RunsCreatePayload & { assistantId: string })[], + ): Promise { + const filteredPayloads = payloads + .map((payload) => ({ ...payload, assistant_id: payload.assistantId })) + .map((payload) => { + return Object.fromEntries( + Object.entries(payload).filter(([_, v]) => v !== undefined), + ); + }); + + return this.fetch("/runs/batch", { + method: "POST", + json: filteredPayloads, + }); + } + async wait( threadId: null, assistantId: string, @@ -775,6 +798,71 @@ export class RunsClient extends BaseClient { return this.fetch(`/threads/${threadId}/runs/${runId}/join`); } + /** + * Stream output from a run in real-time, until the run is done. + * Output is not buffered, so any output produced before this call will + * not be received here. + * + * @param threadId The ID of the thread. + * @param runId The ID of the run. + * @param signal An optional abort signal. + * @returns An async generator yielding stream parts. + */ + async *joinStream( + threadId: string, + runId: string, + signal?: AbortSignal, + ): AsyncGenerator<{ event: StreamEvent; data: any }> { + const response = await this.asyncCaller.fetch( + ...this.prepareFetchOptions(`/threads/${threadId}/runs/${runId}/stream`, { + method: "GET", + signal, + }), + ); + + let parser: EventSourceParser; + let onEndEvent: () => void; + const textDecoder = new TextDecoder(); + + const stream: ReadableStream<{ event: string; data: any }> = ( + response.body || new ReadableStream({ start: (ctrl) => ctrl.close() }) + ).pipeThrough( + new TransformStream({ + async start(ctrl) { + parser = createParser((event) => { + if ( + (signal && signal.aborted) || + (event.type === "event" && event.data === "[DONE]") + ) { + ctrl.terminate(); + return; + } + + if ("data" in event) { + ctrl.enqueue({ + event: event.event ?? "message", + data: JSON.parse(event.data), + }); + } + }); + onEndEvent = () => { + ctrl.enqueue({ event: "end", data: undefined }); + }; + }, + async transform(chunk) { + const payload = textDecoder.decode(chunk); + parser.feed(payload); + + // eventsource-parser will ignore events + // that are not terminated by a newline + if (payload.trim() === "event: end") onEndEvent(); + }, + }), + ); + + yield* IterableReadableStream.fromReadableStream(stream); + } + /** * Delete a run. * diff --git a/libs/sdk-py/langgraph_sdk/client.py b/libs/sdk-py/langgraph_sdk/client.py index 5d6779c5d..1a7766a92 100644 --- a/libs/sdk-py/langgraph_sdk/client.py +++ b/libs/sdk-py/langgraph_sdk/client.py @@ -1234,7 +1234,7 @@ class RunsClient: return await self.http.post("/runs", json=payload) async def create_batch(self, payloads: list[RunCreate]) -> list[Run]: - """Create a batch of background runs.""" + """Create a batch of stateless background runs.""" def filter_payload(payload: RunCreate): return {k: v for k, v in payload.items() if v is not None} @@ -1484,7 +1484,7 @@ class RunsClient: Example Usage: - await client.runs.join( + await client.runs.join_stream( thread_id="thread_id_to_join", run_id="run_id_to_join" )