StreamMode in Join [sdk] (#3584)

This commit is contained in:
William FH
2025-03-13 10:05:45 -07:00
committed by GitHub
4 changed files with 54 additions and 10 deletions
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@langchain/langgraph-sdk",
"version": "0.0.54",
"version": "0.0.56",
"description": "Client library for interacting with the LangGraph API",
"type": "module",
"packageManager": "yarn@1.22.19",
+15 -2
View File
@@ -1022,13 +1022,23 @@ export class RunsClient<
*
* @param threadId The ID of the thread.
* @param runId The ID of the run.
* @param options Additional options for controlling the stream behavior:
* - signal: An AbortSignal that can be used to cancel the stream request
* - cancelOnDisconnect: When true, automatically cancels the run if the client disconnects from the stream
* - streamMode: Controls what types of events to receive from the stream (can be a single mode or array of modes)
* Must be a subset of the stream modes passed when creating the run. Background runs default to having the union of all
* stream modes enabled.
* @returns An async generator yielding stream parts.
*/
async *joinStream(
threadId: string,
runId: string,
options?:
| { signal?: AbortSignal; cancelOnDisconnect?: boolean }
| {
signal?: AbortSignal;
cancelOnDisconnect?: boolean;
streamMode?: StreamMode | StreamMode[];
}
| AbortSignal,
): AsyncGenerator<{ event: StreamEvent; data: any }> {
const opts =
@@ -1043,7 +1053,10 @@ export class RunsClient<
method: "GET",
timeoutMs: null,
signal: opts?.signal,
params: { cancel_on_disconnect: opts?.cancelOnDisconnect ? "1" : "0" },
params: {
cancel_on_disconnect: opts?.cancelOnDisconnect ? "1" : "0",
stream_mode: opts?.streamMode,
},
}),
);
+37 -6
View File
@@ -1831,7 +1831,12 @@ class RunsClient:
return await self.http.get(f"/threads/{thread_id}/runs/{run_id}/join")
def join_stream(
self, thread_id: str, run_id: str, *, cancel_on_disconnect: bool = False
self,
thread_id: str,
run_id: str,
*,
cancel_on_disconnect: bool = False,
stream_mode: Optional[Union[StreamMode, Sequence[StreamMode]]] = None,
) -> AsyncIterator[StreamPart]:
"""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
@@ -1841,6 +1846,9 @@ class RunsClient:
thread_id: The thread ID to join.
run_id: The run ID to join.
cancel_on_disconnect: Whether to cancel the run when the stream is disconnected.
stream_mode: The stream mode(s) to use. Must be a subset of the stream modes passed
when creating the run. Background runs default to having the union of all
stream modes.
Returns:
None
@@ -1849,14 +1857,18 @@ class RunsClient:
await client.runs.join_stream(
thread_id="thread_id_to_join",
run_id="run_id_to_join"
run_id="run_id_to_join",
stream_mode=["values", "debug"]
)
""" # noqa: E501
return self.http.stream(
f"/threads/{thread_id}/runs/{run_id}/stream",
"GET",
params={"cancel_on_disconnect": cancel_on_disconnect},
params={
"cancel_on_disconnect": cancel_on_disconnect,
"stream_mode": stream_mode,
},
)
async def delete(self, thread_id: str, run_id: str) -> None:
@@ -3988,7 +4000,14 @@ class SyncRunsClient:
""" # noqa: E501
return self.http.get(f"/threads/{thread_id}/runs/{run_id}/join")
def join_stream(self, thread_id: str, run_id: str) -> Iterator[StreamPart]:
def join_stream(
self,
thread_id: str,
run_id: str,
*,
stream_mode: Optional[Union[StreamMode, Sequence[StreamMode]]] = None,
cancel_on_disconnect: bool = False,
) -> Iterator[StreamPart]:
"""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.
@@ -3996,6 +4015,10 @@ class SyncRunsClient:
Args:
thread_id: The thread ID to join.
run_id: The run ID to join.
stream_mode: The stream mode(s) to use. Must be a subset of the stream modes passed
when creating the run. Background runs default to having the union of all
stream modes.
cancel_on_disconnect: Whether to cancel the run when the stream is disconnected.
Returns:
None
@@ -4004,11 +4027,19 @@ class SyncRunsClient:
client.runs.join_stream(
thread_id="thread_id_to_join",
run_id="run_id_to_join"
run_id="run_id_to_join",
stream_mode=["values", "debug"]
)
""" # noqa: E501
return self.http.stream(f"/threads/{thread_id}/runs/{run_id}/stream", "GET")
return self.http.stream(
f"/threads/{thread_id}/runs/{run_id}/stream",
"GET",
params={
"stream_mode": stream_mode,
"cancel_on_disconnect": cancel_on_disconnect,
},
)
def delete(self, thread_id: str, run_id: str) -> None:
"""Delete a run.
+1 -1
View File
@@ -1,6 +1,6 @@
[tool.poetry]
name = "langgraph-sdk"
version = "0.1.56"
version = "0.1.57"
description = "SDK for interacting with LangGraph API"
authors = []
license = "MIT"