codekingpro/portable-devtools
114k
1"""Async client for managing runs in LangGraph."""2 3from __future__ import annotations4 5import builtins6import warnings7from collections.abc import AsyncIterator, Callable, Mapping, Sequence8from typing import Any, Literal, overload9 10import httpx11 12from langgraph_sdk._async.http import HttpClient13from langgraph_sdk._shared.utilities import (14 _get_run_metadata_from_response,15 _sse_to_v2_dict,16)17from langgraph_sdk.schema import (18 All,19 BulkCancelRunsStatus,20 CancelAction,21 Checkpoint,22 Command,23 Config,24 Context,25 DisconnectMode,26 Durability,27 IfNotExists,28 Input,29 LangSmithTracing,30 MultitaskStrategy,31 OnCompletionBehavior,32 QueryParamTypes,33 Run,34 RunCreate,35 RunCreateMetadata,36 RunSelectField,37 RunStatus,38 StreamMode,39 StreamPart,40 StreamPartV2,41 StreamVersion,42)43 44 45async def _wrap_stream_v2(46 raw: AsyncIterator[StreamPart],47) -> AsyncIterator[StreamPartV2]:48 """Wrap a raw SSE stream, converting each event to a v2 dict."""49 async for part in raw:50 v2 = _sse_to_v2_dict(part.event, part.data)51 if v2 is not None:52 yield v253 54 55class RunsClient:56 """Client for managing runs in LangGraph.57 58 A run is a single assistant invocation with optional input, config, context, and metadata.59 This client manages runs, which can be stateful (on threads) or stateless.60 61 ???+ example "Example"62 63 ```python64 client = get_client(url="http://localhost:2024")65 run = await client.runs.create(assistant_id="asst_123", thread_id="thread_456", input={"query": "Hello"})66 ```67 """68 69 def __init__(self, http: HttpClient) -> None:70 self.http = http71 72 @overload73 def stream(74 self,75 thread_id: str,76 assistant_id: str,77 *,78 input: Input | None = None,79 command: Command | None = None,80 stream_mode: StreamMode | Sequence[StreamMode] = "values",81 stream_subgraphs: bool = False,82 stream_resumable: bool = False,83 metadata: Mapping[str, Any] | None = None,84 config: Config | None = None,85 context: Context | None = None,86 checkpoint: Checkpoint | None = None,87 checkpoint_id: str | None = None,88 checkpoint_during: bool | None = None,89 interrupt_before: All | Sequence[str] | None = None,90 interrupt_after: All | Sequence[str] | None = None,91 feedback_keys: Sequence[str] | None = None,92 on_disconnect: DisconnectMode | None = None,93 webhook: str | None = None,94 multitask_strategy: MultitaskStrategy | None = None,95 if_not_exists: IfNotExists | None = None,96 after_seconds: int | None = None,97 langsmith_tracing: LangSmithTracing | None = None,98 headers: Mapping[str, str] | None = None,99 params: QueryParamTypes | None = None,100 on_run_created: Callable[[RunCreateMetadata], None] | None = None,101 version: Literal["v1"] = "v1",102 ) -> AsyncIterator[StreamPart]: ...103 104 @overload105 def stream(106 self,107 thread_id: str,108 assistant_id: str,109 *,110 input: Input | None = None,111 command: Command | None = None,112 stream_mode: StreamMode | Sequence[StreamMode] = "values",113 stream_subgraphs: bool = False,114 stream_resumable: bool = False,115 metadata: Mapping[str, Any] | None = None,116 config: Config | None = None,117 context: Context | None = None,118 checkpoint: Checkpoint | None = None,119 checkpoint_id: str | None = None,120 checkpoint_during: bool | None = None,121 interrupt_before: All | Sequence[str] | None = None,122 interrupt_after: All | Sequence[str] | None = None,123 feedback_keys: Sequence[str] | None = None,124 on_disconnect: DisconnectMode | None = None,125 webhook: str | None = None,126 multitask_strategy: MultitaskStrategy | None = None,127 if_not_exists: IfNotExists | None = None,128 after_seconds: int | None = None,129 langsmith_tracing: LangSmithTracing | None = None,130 headers: Mapping[str, str] | None = None,131 params: QueryParamTypes | None = None,132 on_run_created: Callable[[RunCreateMetadata], None] | None = None,133 version: Literal["v2"],134 ) -> AsyncIterator[StreamPartV2]: ...135 136 @overload137 def stream(138 self,139 thread_id: None,140 assistant_id: str,141 *,142 input: Input | None = None,143 command: Command | None = None,144 stream_mode: StreamMode | Sequence[StreamMode] = "values",145 stream_subgraphs: bool = False,146 stream_resumable: bool = False,147 metadata: Mapping[str, Any] | None = None,148 config: Config | None = None,149 checkpoint_during: bool | None = None,150 interrupt_before: All | Sequence[str] | None = None,151 interrupt_after: All | Sequence[str] | None = None,152 feedback_keys: Sequence[str] | None = None,153 on_disconnect: DisconnectMode | None = None,154 on_completion: OnCompletionBehavior | None = None,155 if_not_exists: IfNotExists | None = None,156 webhook: str | None = None,157 after_seconds: int | None = None,158 langsmith_tracing: LangSmithTracing | None = None,159 headers: Mapping[str, str] | None = None,160 params: QueryParamTypes | None = None,161 on_run_created: Callable[[RunCreateMetadata], None] | None = None,162 version: Literal["v1"] = "v1",163 ) -> AsyncIterator[StreamPart]: ...164 165 @overload166 def stream(167 self,168 thread_id: None,169 assistant_id: str,170 *,171 input: Input | None = None,172 command: Command | None = None,173 stream_mode: StreamMode | Sequence[StreamMode] = "values",174 stream_subgraphs: bool = False,175 stream_resumable: bool = False,176 metadata: Mapping[str, Any] | None = None,177 config: Config | None = None,178 checkpoint_during: bool | None = None,179 interrupt_before: All | Sequence[str] | None = None,180 interrupt_after: All | Sequence[str] | None = None,181 feedback_keys: Sequence[str] | None = None,182 on_disconnect: DisconnectMode | None = None,183 on_completion: OnCompletionBehavior | None = None,184 if_not_exists: IfNotExists | None = None,185 webhook: str | None = None,186 after_seconds: int | None = None,187 langsmith_tracing: LangSmithTracing | None = None,188 headers: Mapping[str, str] | None = None,189 params: QueryParamTypes | None = None,190 on_run_created: Callable[[RunCreateMetadata], None] | None = None,191 version: Literal["v2"],192 ) -> AsyncIterator[StreamPartV2]: ...193 194 def stream(195 self,196 thread_id: str | None,197 assistant_id: str,198 *,199 input: Input | None = None,200 command: Command | None = None,201 stream_mode: StreamMode | Sequence[StreamMode] = "values",202 stream_subgraphs: bool = False,203 stream_resumable: bool = False,204 metadata: Mapping[str, Any] | None = None,205 config: Config | None = None,206 context: Context | None = None,207 checkpoint: Checkpoint | None = None,208 checkpoint_id: str | None = None,209 checkpoint_during: bool | None = None, # deprecated210 interrupt_before: All | Sequence[str] | None = None,211 interrupt_after: All | Sequence[str] | None = None,212 feedback_keys: Sequence[str] | None = None,213 on_disconnect: DisconnectMode | None = None,214 on_completion: OnCompletionBehavior | None = None,215 webhook: str | None = None,216 multitask_strategy: MultitaskStrategy | None = None,217 if_not_exists: IfNotExists | None = None,218 after_seconds: int | None = None,219 langsmith_tracing: LangSmithTracing | None = None,220 headers: Mapping[str, str] | None = None,221 params: QueryParamTypes | None = None,222 on_run_created: Callable[[RunCreateMetadata], None] | None = None,223 durability: Durability | None = None,224 version: StreamVersion = "v1",225 ) -> AsyncIterator[StreamPart | StreamPartV2]:226 """Create a run and stream the results.227 228 Args:229 thread_id: the thread ID to assign to the thread.230 If `None` will create a stateless run.231 assistant_id: The assistant ID or graph name to stream from.232 If using graph name, will default to first assistant created from that graph.233 input: The input to the graph.234 command: A command to execute. Cannot be combined with input.235 stream_mode: The stream mode(s) to use.236 stream_subgraphs: Whether to stream output from subgraphs.237 stream_resumable: Whether the stream is considered resumable.238 If true, the stream can be resumed and replayed in its entirety even after disconnection.239 metadata: Metadata to assign to the run.240 config: The configuration for the assistant.241 context: Static context to add to the assistant.242 !!! version-added "Added in version 0.6.0"243 checkpoint: The checkpoint to resume from.244 checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).245 interrupt_before: Nodes to interrupt immediately before they get executed.246 interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.247 feedback_keys: Feedback keys to assign to run.248 on_disconnect: The disconnect mode to use.249 Must be one of 'cancel' or 'continue'.250 on_completion: Whether to delete or keep the thread created for a stateless run.251 Must be one of 'delete' or 'keep'.252 webhook: Webhook to call after LangGraph API call is done.253 multitask_strategy: Multitask strategy to use.254 Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.255 if_not_exists: How to handle missing thread. Defaults to 'reject'.256 Must be either 'reject' (raise error if missing), or 'create' (create new thread).257 after_seconds: The number of seconds to wait before starting the run.258 Use to schedule future runs.259 langsmith_tracing: LangSmith tracing configuration. Allows routing traces260 to a specific project or associating with a dataset example.261 headers: Optional custom headers to include with the request.262 params: Optional query parameters to include with the request.263 on_run_created: Callback when a run is created.264 durability: The durability to use for the run. Values are "sync", "async", or "exit".265 "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True266 "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False267 "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps268 version: Stream format version. "v1" (default) returns raw SSE StreamPart269 NamedTuples. "v2" returns typed dicts with `type`, `ns`, and `data` keys.270 271 Returns:272 Asynchronous iterator of stream results.273 274 ???+ example "Example Usage"275 276 ```python277 client = get_client(url="http://localhost:2024)278 async for chunk in client.runs.stream(279 thread_id=None,280 assistant_id="agent",281 input={"messages": [{"role": "user", "content": "how are you?"}]},282 stream_mode=["values","debug"],283 metadata={"name":"my_run"},284 context={"model_name": "anthropic"},285 interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],286 interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],287 feedback_keys=["my_feedback_key_1","my_feedback_key_2"],288 webhook="https://my.fake.webhook.com",289 multitask_strategy="interrupt"290 ):291 print(chunk)292 ```293 294 ```shell295 296 ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------297 298 StreamPart(event='metadata', data={'run_id': '1ef4a9b8-d7da-679a-a45a-872054341df2'})299 StreamPart(event='values', data={'messages': [{'content': 'how are you?', 'additional_kwargs': {}, 'response_metadata': {}, 'type': 'human', 'name': None, 'id': 'fe0a5778-cfe9-42ee-b807-0adaa1873c10', 'example': False}]})300 StreamPart(event='values', data={'messages': [{'content': 'how are you?', 'additional_kwargs': {}, 'response_metadata': {}, 'type': 'human', 'name': None, 'id': 'fe0a5778-cfe9-42ee-b807-0adaa1873c10', 'example': False}, {'content': "I'm doing well, thanks for asking! I'm an AI assistant created by Anthropic to be helpful, honest, and harmless.", 'additional_kwargs': {}, 'response_metadata': {}, 'type': 'ai', 'name': None, 'id': 'run-159b782c-b679-4830-83c6-cef87798fe8b', 'example': False, 'tool_calls': [], 'invalid_tool_calls': [], 'usage_metadata': None}]})301 StreamPart(event='end', data=None)302 ```303 304 """305 if checkpoint_during is not None:306 warnings.warn(307 "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",308 DeprecationWarning,309 stacklevel=2,310 )311 312 payload: dict[str, Any] = {313 "input": input,314 "command": (315 {k: v for k, v in command.items() if v is not None} if command else None316 ),317 "config": config,318 "context": context,319 "metadata": metadata,320 "stream_mode": stream_mode,321 "stream_subgraphs": stream_subgraphs,322 "stream_resumable": stream_resumable,323 "assistant_id": assistant_id,324 "interrupt_before": interrupt_before,325 "interrupt_after": interrupt_after,326 "feedback_keys": feedback_keys,327 "webhook": webhook,328 "checkpoint": checkpoint,329 "checkpoint_id": checkpoint_id,330 "checkpoint_during": checkpoint_during,331 "multitask_strategy": multitask_strategy,332 "if_not_exists": if_not_exists,333 "on_disconnect": on_disconnect,334 "on_completion": on_completion,335 "after_seconds": after_seconds,336 "durability": durability,337 "langsmith_tracer": langsmith_tracing,338 }339 endpoint = (340 f"/threads/{thread_id}/runs/stream"341 if thread_id is not None342 else "/runs/stream"343 )344 345 def on_response(res: httpx.Response):346 """Callback function to handle the response."""347 if on_run_created and (metadata := _get_run_metadata_from_response(res)):348 on_run_created(metadata)349 350 raw = self.http.stream(351 endpoint,352 "POST",353 json={k: v for k, v in payload.items() if v is not None},354 params=params,355 headers=headers,356 on_response=on_response if on_run_created else None,357 )358 if version == "v2":359 return _wrap_stream_v2(raw)360 return raw361 362 @overload363 async def create(364 self,365 thread_id: None,366 assistant_id: str,367 *,368 input: Input | None = None,369 command: Command | None = None,370 stream_mode: StreamMode | Sequence[StreamMode] = "values",371 stream_subgraphs: bool = False,372 stream_resumable: bool = False,373 metadata: Mapping[str, Any] | None = None,374 checkpoint_during: bool | None = None,375 config: Config | None = None,376 context: Context | None = None,377 interrupt_before: All | Sequence[str] | None = None,378 interrupt_after: All | Sequence[str] | None = None,379 webhook: str | None = None,380 on_completion: OnCompletionBehavior | None = None,381 if_not_exists: IfNotExists | None = None,382 after_seconds: int | None = None,383 langsmith_tracing: LangSmithTracing | None = None,384 headers: Mapping[str, str] | None = None,385 params: QueryParamTypes | None = None,386 on_run_created: Callable[[RunCreateMetadata], None] | None = None,387 ) -> Run: ...388 389 @overload390 async def create(391 self,392 thread_id: str,393 assistant_id: str,394 *,395 input: Input | None = None,396 command: Command | None = None,397 stream_mode: StreamMode | Sequence[StreamMode] = "values",398 stream_subgraphs: bool = False,399 stream_resumable: bool = False,400 metadata: Mapping[str, Any] | None = None,401 config: Config | None = None,402 context: Context | None = None,403 checkpoint: Checkpoint | None = None,404 checkpoint_id: str | None = None,405 checkpoint_during: bool | None = None,406 interrupt_before: All | Sequence[str] | None = None,407 interrupt_after: All | Sequence[str] | None = None,408 webhook: str | None = None,409 multitask_strategy: MultitaskStrategy | None = None,410 if_not_exists: IfNotExists | None = None,411 after_seconds: int | None = None,412 langsmith_tracing: LangSmithTracing | None = None,413 headers: Mapping[str, str] | None = None,414 params: QueryParamTypes | None = None,415 on_run_created: Callable[[RunCreateMetadata], None] | None = None,416 ) -> Run: ...417 418 async def create(419 self,420 thread_id: str | None,421 assistant_id: str,422 *,423 input: Input | None = None,424 command: Command | None = None,425 stream_mode: StreamMode | Sequence[StreamMode] = "values",426 stream_subgraphs: bool = False,427 stream_resumable: bool = False,428 metadata: Mapping[str, Any] | None = None,429 config: Config | None = None,430 context: Context | None = None,431 checkpoint: Checkpoint | None = None,432 checkpoint_id: str | None = None,433 checkpoint_during: bool | None = None, # deprecated434 interrupt_before: All | Sequence[str] | None = None,435 interrupt_after: All | Sequence[str] | None = None,436 webhook: str | None = None,437 multitask_strategy: MultitaskStrategy | None = None,438 if_not_exists: IfNotExists | None = None,439 on_completion: OnCompletionBehavior | None = None,440 after_seconds: int | None = None,441 langsmith_tracing: LangSmithTracing | None = None,442 headers: Mapping[str, str] | None = None,443 params: QueryParamTypes | None = None,444 on_run_created: Callable[[RunCreateMetadata], None] | None = None,445 durability: Durability | None = None,446 ) -> Run:447 """Create a background run.448 449 Args:450 thread_id: the thread ID to assign to the thread.451 If `None` will create a stateless run.452 assistant_id: The assistant ID or graph name to stream from.453 If using graph name, will default to first assistant created from that graph.454 input: The input to the graph.455 command: A command to execute. Cannot be combined with input.456 stream_mode: The stream mode(s) to use.457 stream_subgraphs: Whether to stream output from subgraphs.458 stream_resumable: Whether the stream is considered resumable.459 If true, the stream can be resumed and replayed in its entirety even after disconnection.460 metadata: Metadata to assign to the run.461 config: The configuration for the assistant.462 context: Static context to add to the assistant.463 !!! version-added "Added in version 0.6.0"464 checkpoint: The checkpoint to resume from.465 checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).466 interrupt_before: Nodes to interrupt immediately before they get executed.467 interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.468 webhook: Webhook to call after LangGraph API call is done.469 multitask_strategy: Multitask strategy to use.470 Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.471 on_completion: Whether to delete or keep the thread created for a stateless run.472 Must be one of 'delete' or 'keep'.473 if_not_exists: How to handle missing thread. Defaults to 'reject'.474 Must be either 'reject' (raise error if missing), or 'create' (create new thread).475 after_seconds: The number of seconds to wait before starting the run.476 Use to schedule future runs.477 langsmith_tracing: LangSmith tracing configuration. Allows routing traces478 to a specific project or associating with a dataset example.479 headers: Optional custom headers to include with the request.480 on_run_created: Optional callback to call when a run is created.481 durability: The durability to use for the run. Values are "sync", "async", or "exit".482 "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True483 "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False484 "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps485 486 Returns:487 The created background run.488 489 ???+ example "Example Usage"490 491 ```python492 493 background_run = await client.runs.create(494 thread_id="my_thread_id",495 assistant_id="my_assistant_id",496 input={"messages": [{"role": "user", "content": "hello!"}]},497 metadata={"name":"my_run"},498 context={"model_name": "openai"},499 interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],500 interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],501 webhook="https://my.fake.webhook.com",502 multitask_strategy="interrupt"503 )504 print(background_run)505 ```506 507 ```shell508 --------------------------------------------------------------------------------509 510 {511 'run_id': 'my_run_id',512 'thread_id': 'my_thread_id',513 'assistant_id': 'my_assistant_id',514 'created_at': '2024-07-25T15:35:42.598503+00:00',515 'updated_at': '2024-07-25T15:35:42.598503+00:00',516 'metadata': {},517 'status': 'pending',518 'kwargs':519 {520 'input':521 {522 'messages': [523 {524 'role': 'user',525 'content': 'how are you?'526 }527 ]528 },529 'config':530 {531 'metadata':532 {533 'created_by': 'system'534 },535 'configurable':536 {537 'run_id': 'my_run_id',538 'user_id': None,539 'graph_id': 'agent',540 'thread_id': 'my_thread_id',541 'checkpoint_id': None,542 'assistant_id': 'my_assistant_id'543 },544 },545 'context':546 {547 'model_name': 'openai'548 }549 'webhook': "https://my.fake.webhook.com",550 'temporary': False,551 'stream_mode': ['values'],552 'feedback_keys': None,553 'interrupt_after': ["node_to_stop_after_1","node_to_stop_after_2"],554 'interrupt_before': ["node_to_stop_before_1","node_to_stop_before_2"]555 },556 'multitask_strategy': 'interrupt'557 }558 ```559 """560 if checkpoint_during is not None:561 warnings.warn(562 "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",563 DeprecationWarning,564 stacklevel=2,565 )566 payload = {567 "input": input,568 "command": (569 {k: v for k, v in command.items() if v is not None} if command else None570 ),571 "stream_mode": stream_mode,572 "stream_subgraphs": stream_subgraphs,573 "stream_resumable": stream_resumable,574 "config": config,575 "context": context,576 "metadata": metadata,577 "assistant_id": assistant_id,578 "interrupt_before": interrupt_before,579 "interrupt_after": interrupt_after,580 "webhook": webhook,581 "checkpoint": checkpoint,582 "checkpoint_id": checkpoint_id,583 "checkpoint_during": checkpoint_during,584 "multitask_strategy": multitask_strategy,585 "if_not_exists": if_not_exists,586 "on_completion": on_completion,587 "after_seconds": after_seconds,588 "durability": durability,589 "langsmith_tracer": langsmith_tracing,590 }591 payload = {k: v for k, v in payload.items() if v is not None}592 593 def on_response(res: httpx.Response):594 """Callback function to handle the response."""595 if on_run_created and (metadata := _get_run_metadata_from_response(res)):596 on_run_created(metadata)597 598 return await self.http.post(599 f"/threads/{thread_id}/runs" if thread_id else "/runs",600 json=payload,601 params=params,602 headers=headers,603 on_response=on_response if on_run_created else None,604 )605 606 async def create_batch(607 self,608 payloads: builtins.list[RunCreate],609 *,610 headers: Mapping[str, str] | None = None,611 params: QueryParamTypes | None = None,612 ) -> builtins.list[Run]:613 """Create a batch of stateless background runs."""614 615 def filter_payload(payload: RunCreate):616 return {k: v for k, v in payload.items() if v is not None}617 618 filtered = [filter_payload(payload) for payload in payloads]619 return await self.http.post(620 "/runs/batch", json=filtered, headers=headers, params=params621 )622 623 @overload624 async def wait(625 self,626 thread_id: str,627 assistant_id: str,628 *,629 input: Input | None = None,630 command: Command | None = None,631 metadata: Mapping[str, Any] | None = None,632 config: Config | None = None,633 context: Context | None = None,634 checkpoint: Checkpoint | None = None,635 checkpoint_id: str | None = None,636 checkpoint_during: bool | None = None,637 interrupt_before: All | Sequence[str] | None = None,638 interrupt_after: All | Sequence[str] | None = None,639 webhook: str | None = None,640 on_disconnect: DisconnectMode | None = None,641 multitask_strategy: MultitaskStrategy | None = None,642 if_not_exists: IfNotExists | None = None,643 after_seconds: int | None = None,644 langsmith_tracing: LangSmithTracing | None = None,645 raise_error: bool = True,646 headers: Mapping[str, str] | None = None,647 params: QueryParamTypes | None = None,648 on_run_created: Callable[[RunCreateMetadata], None] | None = None,649 ) -> builtins.list[dict] | dict[str, Any]: ...650 651 @overload652 async def wait(653 self,654 thread_id: None,655 assistant_id: str,656 *,657 input: Input | None = None,658 command: Command | None = None,659 metadata: Mapping[str, Any] | None = None,660 config: Config | None = None,661 context: Context | None = None,662 checkpoint_during: bool | None = None,663 interrupt_before: All | Sequence[str] | None = None,664 interrupt_after: All | Sequence[str] | None = None,665 webhook: str | None = None,666 on_disconnect: DisconnectMode | None = None,667 on_completion: OnCompletionBehavior | None = None,668 if_not_exists: IfNotExists | None = None,669 after_seconds: int | None = None,670 langsmith_tracing: LangSmithTracing | None = None,671 raise_error: bool = True,672 headers: Mapping[str, str] | None = None,673 params: QueryParamTypes | None = None,674 on_run_created: Callable[[RunCreateMetadata], None] | None = None,675 ) -> builtins.list[dict] | dict[str, Any]: ...676 677 async def wait(678 self,679 thread_id: str | None,680 assistant_id: str,681 *,682 input: Input | None = None,683 command: Command | None = None,684 metadata: Mapping[str, Any] | None = None,685 config: Config | None = None,686 context: Context | None = None,687 checkpoint: Checkpoint | None = None,688 checkpoint_id: str | None = None,689 checkpoint_during: bool | None = None, # deprecated690 interrupt_before: All | Sequence[str] | None = None,691 interrupt_after: All | Sequence[str] | None = None,692 webhook: str | None = None,693 on_disconnect: DisconnectMode | None = None,694 on_completion: OnCompletionBehavior | None = None,695 multitask_strategy: MultitaskStrategy | None = None,696 if_not_exists: IfNotExists | None = None,697 after_seconds: int | None = None,698 langsmith_tracing: LangSmithTracing | None = None,699 raise_error: bool = True,700 headers: Mapping[str, str] | None = None,701 params: QueryParamTypes | None = None,702 on_run_created: Callable[[RunCreateMetadata], None] | None = None,703 durability: Durability | None = None,704 ) -> builtins.list[dict] | dict[str, Any]:705 """Create a run, wait until it finishes and return the final state.706 707 Args:708 thread_id: the thread ID to create the run on.709 If `None` will create a stateless run.710 assistant_id: The assistant ID or graph name to run.711 If using graph name, will default to first assistant created from that graph.712 input: The input to the graph.713 command: A command to execute. Cannot be combined with input.714 metadata: Metadata to assign to the run.715 config: The configuration for the assistant.716 context: Static context to add to the assistant.717 !!! version-added "Added in version 0.6.0"718 checkpoint: The checkpoint to resume from.719 checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).720 interrupt_before: Nodes to interrupt immediately before they get executed.721 interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.722 webhook: Webhook to call after LangGraph API call is done.723 on_disconnect: The disconnect mode to use.724 Must be one of 'cancel' or 'continue'.725 on_completion: Whether to delete or keep the thread created for a stateless run.726 Must be one of 'delete' or 'keep'.727 multitask_strategy: Multitask strategy to use.728 Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.729 if_not_exists: How to handle missing thread. Defaults to 'reject'.730 Must be either 'reject' (raise error if missing), or 'create' (create new thread).731 after_seconds: The number of seconds to wait before starting the run.732 Use to schedule future runs.733 langsmith_tracing: LangSmith tracing configuration. Allows routing traces734 to a specific project or associating with a dataset example.735 headers: Optional custom headers to include with the request.736 on_run_created: Optional callback to call when a run is created.737 durability: The durability to use for the run. Values are "sync", "async", or "exit".738 "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True739 "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False740 "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps741 742 Returns:743 The output of the run.744 745 ???+ example "Example Usage"746 747 ```python748 client = get_client(url="http://localhost:2024")749 final_state_of_run = await client.runs.wait(750 thread_id=None,751 assistant_id="agent",752 input={"messages": [{"role": "user", "content": "how are you?"}]},753 metadata={"name":"my_run"},754 context={"model_name": "anthropic"},755 interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],756 interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],757 webhook="https://my.fake.webhook.com",758 multitask_strategy="interrupt"759 )760 print(final_state_of_run)761 ```762 763 ```shell764 -------------------------------------------------------------------------------------------------------------------------------------------765 766 {767 'messages': [768 {769 'content': 'how are you?',770 'additional_kwargs': {},771 'response_metadata': {},772 'type': 'human',773 'name': None,774 'id': 'f51a862c-62fe-4866-863b-b0863e8ad78a',775 'example': False776 },777 {778 'content': "I'm doing well, thanks for asking! I'm an AI assistant created by Anthropic to be helpful, honest, and harmless.",779 'additional_kwargs': {},780 'response_metadata': {},781 'type': 'ai',782 'name': None,783 'id': 'run-bf1cd3c6-768f-4c16-b62d-ba6f17ad8b36',784 'example': False,785 'tool_calls': [],786 'invalid_tool_calls': [],787 'usage_metadata': None788 }789 ]790 }791 ```792 793 """794 if checkpoint_during is not None:795 warnings.warn(796 "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",797 DeprecationWarning,798 stacklevel=2,799 )800 payload = {801 "input": input,802 "command": (803 {k: v for k, v in command.items() if v is not None} if command else None804 ),805 "config": config,806 "context": context,807 "metadata": metadata,808 "assistant_id": assistant_id,809 "interrupt_before": interrupt_before,810 "interrupt_after": interrupt_after,811 "webhook": webhook,812 "checkpoint": checkpoint,813 "checkpoint_id": checkpoint_id,814 "multitask_strategy": multitask_strategy,815 "checkpoint_during": checkpoint_during,816 "if_not_exists": if_not_exists,817 "on_disconnect": on_disconnect,818 "on_completion": on_completion,819 "after_seconds": after_seconds,820 "durability": durability,821 "langsmith_tracer": langsmith_tracing,822 }823 endpoint = (824 f"/threads/{thread_id}/runs/wait" if thread_id is not None else "/runs/wait"825 )826 827 def on_response(res: httpx.Response):828 """Callback function to handle the response."""829 if on_run_created and (metadata := _get_run_metadata_from_response(res)):830 on_run_created(metadata)831 832 response = await self.http.request_reconnect(833 endpoint,834 "POST",835 json={k: v for k, v in payload.items() if v is not None},836 params=params,837 headers=headers,838 on_response=on_response if on_run_created else None,839 )840 if (841 raise_error842 and isinstance(response, dict)843 and "__error__" in response844 and isinstance(response["__error__"], dict)845 ):846 raise Exception(847 f"{response['__error__'].get('error')}: {response['__error__'].get('message')}"848 )849 return response850 851 async def list(852 self,853 thread_id: str,854 *,855 limit: int = 10,856 offset: int = 0,857 status: RunStatus | None = None,858 select: builtins.list[RunSelectField] | None = None,859 headers: Mapping[str, str] | None = None,860 params: QueryParamTypes | None = None,861 ) -> builtins.list[Run]:862 """List runs.863 864 Args:865 thread_id: The thread ID to list runs for.866 limit: The maximum number of results to return.867 offset: The number of results to skip.868 status: The status of the run to filter by.869 headers: Optional custom headers to include with the request.870 params: Optional query parameters to include with the request.871 872 Returns:873 The runs for the thread.874 875 ???+ example "Example Usage"876 877 ```python878 client = get_client(url="http://localhost:2024")879 await client.runs.list(880 thread_id="thread_id",881 limit=5,882 offset=5,883 )884 ```885 886 """887 query_params: dict[str, Any] = {888 "limit": limit,889 "offset": offset,890 }891 if status is not None:892 query_params["status"] = status893 if select:894 query_params["select"] = select895 if params:896 query_params.update(params)897 return await self.http.get(898 f"/threads/{thread_id}/runs", params=query_params, headers=headers899 )900 901 async def get(902 self,903 thread_id: str,904 run_id: str,905 *,906 headers: Mapping[str, str] | None = None,907 params: QueryParamTypes | None = None,908 ) -> Run:909 """Get a run.910 911 Args:912 thread_id: The thread ID to get.913 run_id: The run ID to get.914 headers: Optional custom headers to include with the request.915 params: Optional query parameters to include with the request.916 917 Returns:918 `Run` object.919 920 ???+ example "Example Usage"921 922 ```python923 client = get_client(url="http://localhost:2024")924 run = await client.runs.get(925 thread_id="thread_id_to_delete",926 run_id="run_id_to_delete",927 )928 ```929 930 """931 932 return await self.http.get(933 f"/threads/{thread_id}/runs/{run_id}", headers=headers, params=params934 )935 936 async def cancel(937 self,938 thread_id: str,939 run_id: str,940 *,941 wait: bool = False,942 action: CancelAction = "interrupt",943 headers: Mapping[str, str] | None = None,944 params: QueryParamTypes | None = None,945 ) -> None:946 """Get a run.947 948 Args:949 thread_id: The thread ID to cancel.950 run_id: The run ID to cancel.951 wait: Whether to wait until run has completed.952 action: Action to take when cancelling the run. Possible values953 are `interrupt` or `rollback`. Default is `interrupt`.954 headers: Optional custom headers to include with the request.955 params: Optional query parameters to include with the request.956 957 Returns:958 `None`959 960 ???+ example "Example Usage"961 962 ```python963 client = get_client(url="http://localhost:2024")964 await client.runs.cancel(965 thread_id="thread_id_to_cancel",966 run_id="run_id_to_cancel",967 wait=True,968 action="interrupt"969 )970 ```971 972 """973 query_params = {974 "wait": 1 if wait else 0,975 "action": action,976 }977 if params:978 query_params.update(params)979 if wait:980 return await self.http.request_reconnect(981 f"/threads/{thread_id}/runs/{run_id}/cancel",982 "POST",983 params=query_params,984 headers=headers,985 )986 else:987 return await self.http.post(988 f"/threads/{thread_id}/runs/{run_id}/cancel",989 json=None,990 params=query_params,991 headers=headers,992 )993 994 async def cancel_many(995 self,996 *,997 thread_id: str | None = None,998 run_ids: Sequence[str] | None = None,999 status: BulkCancelRunsStatus | None = None,1000 action: CancelAction = "interrupt",1001 headers: Mapping[str, str] | None = None,1002 params: QueryParamTypes | None = None,1003 ) -> None:1004 """Cancel one or more runs.1005 1006 Can cancel runs by thread ID and run IDs, or by status filter.1007 1008 Args:1009 thread_id: The ID of the thread containing runs to cancel.1010 run_ids: List of run IDs to cancel.1011 status: Filter runs by status to cancel. Must be one of1012 `"pending"`, `"running"`, or `"all"`.1013 action: Action to take when cancelling the run. Possible values1014 are `"interrupt"` or `"rollback"`. Default is `"interrupt"`.1015 headers: Optional custom headers to include with the request.1016 params: Optional query parameters to include with the request.1017 1018 Returns:1019 `None`1020 1021 ???+ example "Example Usage"1022 1023 ```python1024 client = get_client(url="http://localhost:2024")1025 # Cancel all pending runs1026 await client.runs.cancel_many(status="pending")1027 # Cancel specific runs on a thread1028 await client.runs.cancel_many(1029 thread_id="my_thread_id",1030 run_ids=["run_1", "run_2"],1031 action="rollback",1032 )1033 ```1034 1035 """1036 payload: dict[str, Any] = {}1037 if thread_id:1038 payload["thread_id"] = thread_id1039 if run_ids:1040 payload["run_ids"] = run_ids1041 if status:1042 payload["status"] = status1043 query_params: dict[str, Any] = {"action": action}1044 if params:1045 query_params.update(params)1046 await self.http.post(1047 "/runs/cancel",1048 json=payload,1049 headers=headers,1050 params=query_params,1051 )1052 1053 async def join(1054 self,1055 thread_id: str,1056 run_id: str,1057 *,1058 headers: Mapping[str, str] | None = None,1059 params: QueryParamTypes | None = None,1060 ) -> dict:1061 """Block until a run is done. Returns the final state of the thread.1062 1063 Args:1064 thread_id: The thread ID to join.1065 run_id: The run ID to join.1066 headers: Optional custom headers to include with the request.1067 params: Optional query parameters to include with the request.1068 1069 Returns:1070 `None`1071 1072 ???+ example "Example Usage"1073 1074 ```python1075 client = get_client(url="http://localhost:2024")1076 result =await client.runs.join(1077 thread_id="thread_id_to_join",1078 run_id="run_id_to_join"1079 )1080 ```1081 1082 """1083 return await self.http.request_reconnect(1084 f"/threads/{thread_id}/runs/{run_id}/join",1085 "GET",1086 headers=headers,1087 params=params,1088 )1089 1090 def join_stream(1091 self,1092 thread_id: str,1093 run_id: str,1094 *,1095 cancel_on_disconnect: bool = False,1096 stream_mode: StreamMode | Sequence[StreamMode] | None = None,1097 headers: Mapping[str, str] | None = None,1098 params: QueryParamTypes | None = None,1099 last_event_id: str | None = None,1100 ) -> AsyncIterator[StreamPart]:1101 """Stream output from a run in real-time, until the run is done.1102 Output is not buffered, so any output produced before this call will1103 not be received here.1104 1105 Args:1106 thread_id: The thread ID to join.1107 run_id: The run ID to join.1108 cancel_on_disconnect: Whether to cancel the run when the stream is disconnected.1109 stream_mode: The stream mode(s) to use. Must be a subset of the stream modes passed1110 when creating the run. Background runs default to having the union of all1111 stream modes.1112 headers: Optional custom headers to include with the request.1113 params: Optional query parameters to include with the request.1114 last_event_id: The last event ID to use for the stream.1115 1116 Returns:1117 The stream of parts.1118 1119 ???+ example "Example Usage"1120 1121 ```python1122 client = get_client(url="http://localhost:2024")1123 async for part in client.runs.join_stream(1124 thread_id="thread_id_to_join",1125 run_id="run_id_to_join",1126 stream_mode=["values", "debug"]1127 ):1128 print(part)1129 ```1130 1131 """1132 query_params = {1133 "cancel_on_disconnect": cancel_on_disconnect,1134 "stream_mode": stream_mode,1135 }1136 if params:1137 query_params.update(params)1138 return self.http.stream(1139 f"/threads/{thread_id}/runs/{run_id}/stream",1140 "GET",1141 params=query_params,1142 headers={1143 **({"Last-Event-ID": last_event_id} if last_event_id else {}),1144 **(headers or {}),1145 }1146 or None,1147 )1148 1149 async def delete(1150 self,1151 thread_id: str,1152 run_id: str,1153 *,1154 headers: Mapping[str, str] | None = None,1155 params: QueryParamTypes | None = None,1156 ) -> None:1157 """Delete a run.1158 1159 Args:1160 thread_id: The thread ID to delete.1161 run_id: The run ID to delete.1162 headers: Optional custom headers to include with the request.1163 params: Optional query parameters to include with the request.1164 1165 Returns:1166 `None`1167 1168 ???+ example "Example Usage"1169 1170 ```python1171 client = get_client(url="http://localhost:2024")1172 await client.runs.delete(1173 thread_id="thread_id_to_delete",1174 run_id="run_id_to_delete"1175 )1176 ```1177 1178 """1179 await self.http.delete(1180 f"/threads/{thread_id}/runs/{run_id}", headers=headers, params=params1181 )1182 