codekingpro/portable-devtools
114k
1"""Synchronous client for managing runs in LangGraph."""2 3from __future__ import annotations4 5import builtins6import warnings7from collections.abc import Callable, Iterator, Mapping, Sequence8from typing import Any, Literal, overload9 10import httpx11 12from langgraph_sdk._shared.utilities import (13 _get_run_metadata_from_response,14 _sse_to_v2_dict,15)16from langgraph_sdk._sync.http import SyncHttpClient17from 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 45def _wrap_stream_v2_sync(46 raw: Iterator[StreamPart],47) -> Iterator[StreamPartV2]:48 """Wrap a raw SSE stream, converting each event to a v2 dict."""49 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 SyncRunsClient:56 """Synchronous client for managing runs in LangGraph.57 58 This class provides methods to create, retrieve, and manage runs, which represent59 individual executions of graphs.60 61 ???+ example "Example"62 63 ```python64 client = get_sync_client(url="http://localhost:2024")65 run = client.runs.create(thread_id="thread_123", assistant_id="asst_456")66 ```67 """68 69 def __init__(self, http: SyncHttpClient) -> 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 metadata: Mapping[str, Any] | None = None,83 config: Config | None = None,84 context: Context | None = None,85 checkpoint: Checkpoint | None = None,86 checkpoint_id: str | None = None,87 checkpoint_during: bool | None = None,88 interrupt_before: All | Sequence[str] | None = None,89 interrupt_after: All | Sequence[str] | None = None,90 feedback_keys: Sequence[str] | None = None,91 on_disconnect: DisconnectMode | None = None,92 webhook: str | None = None,93 multitask_strategy: MultitaskStrategy | None = None,94 if_not_exists: IfNotExists | None = None,95 after_seconds: int | None = None,96 langsmith_tracing: LangSmithTracing | None = None,97 headers: Mapping[str, str] | None = None,98 params: QueryParamTypes | None = None,99 on_run_created: Callable[[RunCreateMetadata], None] | None = None,100 version: Literal["v1"] = "v1",101 ) -> Iterator[StreamPart]: ...102 103 @overload104 def stream(105 self,106 thread_id: str,107 assistant_id: str,108 *,109 input: Input | None = None,110 command: Command | None = None,111 stream_mode: StreamMode | Sequence[StreamMode] = "values",112 stream_subgraphs: bool = False,113 metadata: Mapping[str, Any] | None = None,114 config: Config | None = None,115 context: Context | None = None,116 checkpoint: Checkpoint | None = None,117 checkpoint_id: str | None = None,118 checkpoint_during: bool | None = None,119 interrupt_before: All | Sequence[str] | None = None,120 interrupt_after: All | Sequence[str] | None = None,121 feedback_keys: Sequence[str] | None = None,122 on_disconnect: DisconnectMode | None = None,123 webhook: str | None = None,124 multitask_strategy: MultitaskStrategy | None = None,125 if_not_exists: IfNotExists | None = None,126 after_seconds: int | None = None,127 langsmith_tracing: LangSmithTracing | None = None,128 headers: Mapping[str, str] | None = None,129 params: QueryParamTypes | None = None,130 on_run_created: Callable[[RunCreateMetadata], None] | None = None,131 version: Literal["v2"],132 ) -> Iterator[StreamPartV2]: ...133 134 @overload135 def stream(136 self,137 thread_id: None,138 assistant_id: str,139 *,140 input: Input | None = None,141 command: Command | None = None,142 stream_mode: StreamMode | Sequence[StreamMode] = "values",143 stream_subgraphs: bool = False,144 stream_resumable: bool = False,145 metadata: Mapping[str, Any] | None = None,146 config: Config | None = None,147 context: Context | None = None,148 checkpoint_during: bool | None = None,149 interrupt_before: All | Sequence[str] | None = None,150 interrupt_after: All | Sequence[str] | None = None,151 feedback_keys: Sequence[str] | None = None,152 on_disconnect: DisconnectMode | None = None,153 on_completion: OnCompletionBehavior | None = None,154 if_not_exists: IfNotExists | None = None,155 webhook: str | None = None,156 after_seconds: int | None = None,157 langsmith_tracing: LangSmithTracing | None = None,158 headers: Mapping[str, str] | None = None,159 params: QueryParamTypes | None = None,160 on_run_created: Callable[[RunCreateMetadata], None] | None = None,161 version: Literal["v1"] = "v1",162 ) -> Iterator[StreamPart]: ...163 164 @overload165 def stream(166 self,167 thread_id: None,168 assistant_id: str,169 *,170 input: Input | None = None,171 command: Command | None = None,172 stream_mode: StreamMode | Sequence[StreamMode] = "values",173 stream_subgraphs: bool = False,174 stream_resumable: bool = False,175 metadata: Mapping[str, Any] | None = None,176 config: Config | None = None,177 context: Context | 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 ) -> Iterator[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 ) -> Iterator[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: The command to execute.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 on_run_created: Optional callback to call when a run is created.263 durability: The durability to use for the run. Values are "sync", "async", or "exit".264 "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True265 "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False266 "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps267 version: Stream format version. "v1" (default) returns raw SSE StreamPart268 NamedTuples. "v2" returns typed dicts with `type`, `ns`, and `data` keys.269 270 Returns:271 Iterator of stream results.272 273 ???+ example "Example Usage"274 275 ```python276 client = get_sync_client(url="http://localhost:2024")277 async for chunk in client.runs.stream(278 thread_id=None,279 assistant_id="agent",280 input={"messages": [{"role": "user", "content": "how are you?"}]},281 stream_mode=["values","debug"],282 metadata={"name":"my_run"},283 context={"model_name": "anthropic"},284 interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],285 interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],286 feedback_keys=["my_feedback_key_1","my_feedback_key_2"],287 webhook="https://my.fake.webhook.com",288 multitask_strategy="interrupt"289 ):290 print(chunk)291 ```292 ```shell293 ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------294 295 StreamPart(event='metadata', data={'run_id': '1ef4a9b8-d7da-679a-a45a-872054341df2'})296 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}]})297 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}]})298 StreamPart(event='end', data=None)299 ```300 """301 if checkpoint_during is not None:302 warnings.warn(303 "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",304 DeprecationWarning,305 stacklevel=2,306 )307 payload: dict[str, Any] = {308 "input": input,309 "command": (310 {k: v for k, v in command.items() if v is not None} if command else None311 ),312 "config": config,313 "context": context,314 "metadata": metadata,315 "stream_mode": stream_mode,316 "stream_subgraphs": stream_subgraphs,317 "stream_resumable": stream_resumable,318 "assistant_id": assistant_id,319 "interrupt_before": interrupt_before,320 "interrupt_after": interrupt_after,321 "feedback_keys": feedback_keys,322 "webhook": webhook,323 "checkpoint": checkpoint,324 "checkpoint_id": checkpoint_id,325 "checkpoint_during": checkpoint_during,326 "multitask_strategy": multitask_strategy,327 "if_not_exists": if_not_exists,328 "on_disconnect": on_disconnect,329 "on_completion": on_completion,330 "after_seconds": after_seconds,331 "durability": durability,332 "langsmith_tracer": langsmith_tracing,333 }334 endpoint = (335 f"/threads/{thread_id}/runs/stream"336 if thread_id is not None337 else "/runs/stream"338 )339 340 def on_response(res: httpx.Response):341 """Callback function to handle the response."""342 if on_run_created and (metadata := _get_run_metadata_from_response(res)):343 on_run_created(metadata)344 345 raw = self.http.stream(346 endpoint,347 "POST",348 json={k: v for k, v in payload.items() if v is not None},349 params=params,350 headers=headers,351 on_response=on_response if on_run_created else None,352 )353 if version == "v2":354 return _wrap_stream_v2_sync(raw)355 return raw356 357 @overload358 def create(359 self,360 thread_id: None,361 assistant_id: str,362 *,363 input: Input | None = None,364 command: Command | None = None,365 stream_mode: StreamMode | Sequence[StreamMode] = "values",366 stream_subgraphs: bool = False,367 stream_resumable: bool = False,368 metadata: Mapping[str, Any] | None = None,369 config: Config | None = None,370 context: Context | None = None,371 checkpoint_during: bool | None = None,372 interrupt_before: All | Sequence[str] | None = None,373 interrupt_after: All | Sequence[str] | None = None,374 webhook: str | None = None,375 on_completion: OnCompletionBehavior | None = None,376 if_not_exists: IfNotExists | None = None,377 after_seconds: int | None = None,378 langsmith_tracing: LangSmithTracing | None = None,379 headers: Mapping[str, str] | None = None,380 params: QueryParamTypes | None = None,381 on_run_created: Callable[[RunCreateMetadata], None] | None = None,382 ) -> Run: ...383 384 @overload385 def create(386 self,387 thread_id: str,388 assistant_id: str,389 *,390 input: Input | None = None,391 command: Command | None = None,392 stream_mode: StreamMode | Sequence[StreamMode] = "values",393 stream_subgraphs: bool = False,394 stream_resumable: bool = False,395 metadata: Mapping[str, Any] | None = None,396 config: Config | None = None,397 context: Context | None = None,398 checkpoint: Checkpoint | None = None,399 checkpoint_id: str | None = None,400 checkpoint_during: bool | None = None,401 interrupt_before: All | Sequence[str] | None = None,402 interrupt_after: All | Sequence[str] | None = None,403 webhook: str | None = None,404 multitask_strategy: MultitaskStrategy | None = None,405 if_not_exists: IfNotExists | None = None,406 after_seconds: int | None = None,407 langsmith_tracing: LangSmithTracing | None = None,408 headers: Mapping[str, str] | None = None,409 params: QueryParamTypes | None = None,410 on_run_created: Callable[[RunCreateMetadata], None] | None = None,411 ) -> Run: ...412 413 def create(414 self,415 thread_id: str | None,416 assistant_id: str,417 *,418 input: Input | None = None,419 command: Command | None = None,420 stream_mode: StreamMode | Sequence[StreamMode] = "values",421 stream_subgraphs: bool = False,422 stream_resumable: bool = False,423 metadata: Mapping[str, Any] | None = None,424 config: Config | None = None,425 context: Context | None = None,426 checkpoint: Checkpoint | None = None,427 checkpoint_id: str | None = None,428 checkpoint_during: bool | None = None, # deprecated429 interrupt_before: All | Sequence[str] | None = None,430 interrupt_after: All | Sequence[str] | None = None,431 webhook: str | None = None,432 multitask_strategy: MultitaskStrategy | None = None,433 if_not_exists: IfNotExists | None = None,434 on_completion: OnCompletionBehavior | None = None,435 after_seconds: int | None = None,436 langsmith_tracing: LangSmithTracing | None = None,437 headers: Mapping[str, str] | None = None,438 params: QueryParamTypes | None = None,439 on_run_created: Callable[[RunCreateMetadata], None] | None = None,440 durability: Durability | None = None,441 ) -> Run:442 """Create a background run.443 444 Args:445 thread_id: the thread ID to assign to the thread.446 If `None` will create a stateless run.447 assistant_id: The assistant ID or graph name to stream from.448 If using graph name, will default to first assistant created from that graph.449 input: The input to the graph.450 command: The command to execute.451 stream_mode: The stream mode(s) to use.452 stream_subgraphs: Whether to stream output from subgraphs.453 stream_resumable: Whether the stream is considered resumable.454 If true, the stream can be resumed and replayed in its entirety even after disconnection.455 metadata: Metadata to assign to the run.456 config: The configuration for the assistant.457 context: Static context to add to the assistant.458 !!! version-added "Added in version 0.6.0"459 checkpoint: The checkpoint to resume from.460 checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).461 interrupt_before: Nodes to interrupt immediately before they get executed.462 interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.463 webhook: Webhook to call after LangGraph API call is done.464 multitask_strategy: Multitask strategy to use.465 Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.466 on_completion: Whether to delete or keep the thread created for a stateless run.467 Must be one of 'delete' or 'keep'.468 if_not_exists: How to handle missing thread. Defaults to 'reject'.469 Must be either 'reject' (raise error if missing), or 'create' (create new thread).470 after_seconds: The number of seconds to wait before starting the run.471 Use to schedule future runs.472 langsmith_tracing: LangSmith tracing configuration. Allows routing traces473 to a specific project or associating with a dataset example.474 headers: Optional custom headers to include with the request.475 on_run_created: Optional callback to call when a run is created.476 durability: The durability to use for the run. Values are "sync", "async", or "exit".477 "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True478 "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False479 "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps480 481 Returns:482 The created background `Run`.483 484 ???+ example "Example Usage"485 486 ```python487 client = get_sync_client(url="http://localhost:2024")488 background_run = client.runs.create(489 thread_id="my_thread_id",490 assistant_id="my_assistant_id",491 input={"messages": [{"role": "user", "content": "hello!"}]},492 metadata={"name":"my_run"},493 context={"model_name": "openai"},494 interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],495 interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],496 webhook="https://my.fake.webhook.com",497 multitask_strategy="interrupt"498 )499 print(background_run)500 ```501 502 ```shell503 --------------------------------------------------------------------------------504 505 {506 'run_id': 'my_run_id',507 'thread_id': 'my_thread_id',508 'assistant_id': 'my_assistant_id',509 'created_at': '2024-07-25T15:35:42.598503+00:00',510 'updated_at': '2024-07-25T15:35:42.598503+00:00',511 'metadata': {},512 'status': 'pending',513 'kwargs':514 {515 'input':516 {517 'messages': [518 {519 'role': 'user',520 'content': 'how are you?'521 }522 ]523 },524 'config':525 {526 'metadata':527 {528 'created_by': 'system'529 },530 'configurable':531 {532 'run_id': 'my_run_id',533 'user_id': None,534 'graph_id': 'agent',535 'thread_id': 'my_thread_id',536 'checkpoint_id': None,537 'assistant_id': 'my_assistant_id'538 }539 },540 'context':541 {542 'model_name': 'openai'543 },544 'webhook': "https://my.fake.webhook.com",545 'temporary': False,546 'stream_mode': ['values'],547 'feedback_keys': None,548 'interrupt_after': ["node_to_stop_after_1","node_to_stop_after_2"],549 'interrupt_before': ["node_to_stop_before_1","node_to_stop_before_2"]550 },551 'multitask_strategy': 'interrupt'552 }553 ```554 """555 if checkpoint_during is not None:556 warnings.warn(557 "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",558 DeprecationWarning,559 stacklevel=2,560 )561 payload = {562 "input": input,563 "command": (564 {k: v for k, v in command.items() if v is not None} if command else None565 ),566 "stream_mode": stream_mode,567 "stream_subgraphs": stream_subgraphs,568 "stream_resumable": stream_resumable,569 "config": config,570 "context": context,571 "metadata": metadata,572 "assistant_id": assistant_id,573 "interrupt_before": interrupt_before,574 "interrupt_after": interrupt_after,575 "webhook": webhook,576 "checkpoint": checkpoint,577 "checkpoint_id": checkpoint_id,578 "checkpoint_during": checkpoint_during,579 "multitask_strategy": multitask_strategy,580 "if_not_exists": if_not_exists,581 "on_completion": on_completion,582 "after_seconds": after_seconds,583 "durability": durability,584 "langsmith_tracer": langsmith_tracing,585 }586 payload = {k: v for k, v in payload.items() if v is not None}587 588 def on_response(res: httpx.Response):589 """Callback function to handle the response."""590 if on_run_created and (metadata := _get_run_metadata_from_response(res)):591 on_run_created(metadata)592 593 return self.http.post(594 f"/threads/{thread_id}/runs" if thread_id else "/runs",595 json=payload,596 params=params,597 headers=headers,598 on_response=on_response if on_run_created else None,599 )600 601 def create_batch(602 self,603 payloads: builtins.list[RunCreate],604 *,605 headers: Mapping[str, str] | None = None,606 params: QueryParamTypes | None = None,607 ) -> builtins.list[Run]:608 """Create a batch of stateless background runs."""609 610 def filter_payload(payload: RunCreate):611 return {k: v for k, v in payload.items() if v is not None}612 613 filtered = [filter_payload(payload) for payload in payloads]614 return self.http.post(615 "/runs/batch", json=filtered, headers=headers, params=params616 )617 618 @overload619 def wait(620 self,621 thread_id: str,622 assistant_id: str,623 *,624 input: Input | None = None,625 command: Command | None = None,626 metadata: Mapping[str, Any] | None = None,627 config: Config | None = None,628 context: Context | None = None,629 checkpoint: Checkpoint | None = None,630 checkpoint_id: str | None = None,631 checkpoint_during: bool | None = None,632 interrupt_before: All | Sequence[str] | None = None,633 interrupt_after: All | Sequence[str] | None = None,634 webhook: str | None = None,635 on_disconnect: DisconnectMode | None = None,636 multitask_strategy: MultitaskStrategy | None = None,637 if_not_exists: IfNotExists | None = None,638 after_seconds: int | None = None,639 langsmith_tracing: LangSmithTracing | None = None,640 raise_error: bool = True,641 headers: Mapping[str, str] | None = None,642 params: QueryParamTypes | None = None,643 on_run_created: Callable[[RunCreateMetadata], None] | None = None,644 ) -> builtins.list[dict] | dict[str, Any]: ...645 646 @overload647 def wait(648 self,649 thread_id: None,650 assistant_id: str,651 *,652 input: Input | None = None,653 command: Command | None = None,654 metadata: Mapping[str, Any] | None = None,655 config: Config | None = None,656 context: Context | None = None,657 checkpoint_during: bool | None = None,658 interrupt_before: All | Sequence[str] | None = None,659 interrupt_after: All | Sequence[str] | None = None,660 webhook: str | None = None,661 on_disconnect: DisconnectMode | None = None,662 on_completion: OnCompletionBehavior | None = None,663 if_not_exists: IfNotExists | None = None,664 after_seconds: int | None = None,665 langsmith_tracing: LangSmithTracing | None = None,666 raise_error: bool = True,667 headers: Mapping[str, str] | None = None,668 params: QueryParamTypes | None = None,669 on_run_created: Callable[[RunCreateMetadata], None] | None = None,670 ) -> builtins.list[dict] | dict[str, Any]: ...671 672 def wait(673 self,674 thread_id: str | None,675 assistant_id: str,676 *,677 input: Input | None = None,678 command: Command | None = None,679 metadata: Mapping[str, Any] | None = None,680 config: Config | None = None,681 context: Context | None = None,682 checkpoint_during: bool | None = None, # deprecated683 checkpoint: Checkpoint | None = None,684 checkpoint_id: str | None = None,685 interrupt_before: All | Sequence[str] | None = None,686 interrupt_after: All | Sequence[str] | None = None,687 webhook: str | None = None,688 on_disconnect: DisconnectMode | None = None,689 on_completion: OnCompletionBehavior | None = None,690 multitask_strategy: MultitaskStrategy | None = None,691 if_not_exists: IfNotExists | None = None,692 after_seconds: int | None = None,693 langsmith_tracing: LangSmithTracing | None = None,694 raise_error: bool = True,695 headers: Mapping[str, str] | None = None,696 params: QueryParamTypes | None = None,697 on_run_created: Callable[[RunCreateMetadata], None] | None = None,698 durability: Durability | None = None,699 ) -> builtins.list[dict] | dict[str, Any]:700 """Create a run, wait until it finishes and return the final state.701 702 Args:703 thread_id: the thread ID to create the run on.704 If `None` will create a stateless run.705 assistant_id: The assistant ID or graph name to run.706 If using graph name, will default to first assistant created from that graph.707 input: The input to the graph.708 command: The command to execute.709 metadata: Metadata to assign to the run.710 config: The configuration for the assistant.711 context: Static context to add to the assistant.712 !!! version-added "Added in version 0.6.0"713 checkpoint: The checkpoint to resume from.714 checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).715 interrupt_before: Nodes to interrupt immediately before they get executed.716 interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.717 webhook: Webhook to call after LangGraph API call is done.718 on_disconnect: The disconnect mode to use.719 Must be one of 'cancel' or 'continue'.720 on_completion: Whether to delete or keep the thread created for a stateless run.721 Must be one of 'delete' or 'keep'.722 multitask_strategy: Multitask strategy to use.723 Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.724 if_not_exists: How to handle missing thread. Defaults to 'reject'.725 Must be either 'reject' (raise error if missing), or 'create' (create new thread).726 after_seconds: The number of seconds to wait before starting the run.727 Use to schedule future runs.728 langsmith_tracing: LangSmith tracing configuration. Allows routing traces729 to a specific project or associating with a dataset example.730 raise_error: Whether to raise an error if the run fails.731 headers: Optional custom headers to include with the request.732 on_run_created: Optional callback to call when a run is created.733 durability: The durability to use for the run. Values are "sync", "async", or "exit".734 "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True735 "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False736 "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps737 738 Returns:739 The output of the `Run`.740 741 ???+ example "Example Usage"742 743 ```python744 745 final_state_of_run = client.runs.wait(746 thread_id=None,747 assistant_id="agent",748 input={"messages": [{"role": "user", "content": "how are you?"}]},749 metadata={"name":"my_run"},750 context={"model_name": "anthropic"},751 interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],752 interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],753 webhook="https://my.fake.webhook.com",754 multitask_strategy="interrupt"755 )756 print(final_state_of_run)757 ```758 759 ```shell760 761 -------------------------------------------------------------------------------------------------------------------------------------------762 763 {764 'messages': [765 {766 'content': 'how are you?',767 'additional_kwargs': {},768 'response_metadata': {},769 'type': 'human',770 'name': None,771 'id': 'f51a862c-62fe-4866-863b-b0863e8ad78a',772 'example': False773 },774 {775 'content': "I'm doing well, thanks for asking! I'm an AI assistant created by Anthropic to be helpful, honest, and harmless.",776 'additional_kwargs': {},777 'response_metadata': {},778 'type': 'ai',779 'name': None,780 'id': 'run-bf1cd3c6-768f-4c16-b62d-ba6f17ad8b36',781 'example': False,782 'tool_calls': [],783 'invalid_tool_calls': [],784 'usage_metadata': None785 }786 ]787 }788 ```789 790 """791 if checkpoint_during is not None:792 warnings.warn(793 "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",794 DeprecationWarning,795 stacklevel=2,796 )797 payload = {798 "input": input,799 "command": (800 {k: v for k, v in command.items() if v is not None} if command else None801 ),802 "config": config,803 "context": context,804 "metadata": metadata,805 "assistant_id": assistant_id,806 "interrupt_before": interrupt_before,807 "interrupt_after": interrupt_after,808 "webhook": webhook,809 "checkpoint": checkpoint,810 "checkpoint_id": checkpoint_id,811 "multitask_strategy": multitask_strategy,812 "if_not_exists": if_not_exists,813 "on_disconnect": on_disconnect,814 "checkpoint_during": checkpoint_during,815 "on_completion": on_completion,816 "after_seconds": after_seconds,817 "raise_error": raise_error,818 "durability": durability,819 "langsmith_tracer": langsmith_tracing,820 }821 822 def on_response(res: httpx.Response):823 """Callback function to handle the response."""824 if on_run_created and (metadata := _get_run_metadata_from_response(res)):825 on_run_created(metadata)826 827 endpoint = (828 f"/threads/{thread_id}/runs/wait" if thread_id is not None else "/runs/wait"829 )830 return self.http.request_reconnect(831 endpoint,832 "POST",833 json={k: v for k, v in payload.items() if v is not None},834 params=params,835 headers=headers,836 on_response=on_response if on_run_created else None,837 )838 839 def list(840 self,841 thread_id: str,842 *,843 limit: int = 10,844 offset: int = 0,845 status: RunStatus | None = None,846 select: builtins.list[RunSelectField] | None = None,847 headers: Mapping[str, str] | None = None,848 params: QueryParamTypes | None = None,849 ) -> builtins.list[Run]:850 """List runs.851 852 Args:853 thread_id: The thread ID to list runs for.854 limit: The maximum number of results to return.855 offset: The number of results to skip.856 headers: Optional custom headers to include with the request.857 params: Optional query parameters to include with the request.858 859 Returns:860 The runs for the thread.861 862 ???+ example "Example Usage"863 864 ```python865 client = get_sync_client(url="http://localhost:2024")866 client.runs.list(867 thread_id="thread_id",868 limit=5,869 offset=5,870 )871 ```872 873 """874 query_params: dict[str, Any] = {"limit": limit, "offset": offset}875 if status is not None:876 query_params["status"] = status877 if select:878 query_params["select"] = select879 if params:880 query_params.update(params)881 return self.http.get(882 f"/threads/{thread_id}/runs", params=query_params, headers=headers883 )884 885 def get(886 self,887 thread_id: str,888 run_id: str,889 *,890 headers: Mapping[str, str] | None = None,891 params: QueryParamTypes | None = None,892 ) -> Run:893 """Get a run.894 895 Args:896 thread_id: The thread ID to get.897 run_id: The run ID to get.898 headers: Optional custom headers to include with the request.899 900 Returns:901 `Run` object.902 903 ???+ example "Example Usage"904 905 ```python906 907 run = client.runs.get(908 thread_id="thread_id_to_delete",909 run_id="run_id_to_delete",910 )911 ```912 """913 914 return self.http.get(915 f"/threads/{thread_id}/runs/{run_id}", headers=headers, params=params916 )917 918 def cancel(919 self,920 thread_id: str,921 run_id: str,922 *,923 wait: bool = False,924 action: CancelAction = "interrupt",925 headers: Mapping[str, str] | None = None,926 params: QueryParamTypes | None = None,927 ) -> None:928 """Get a run.929 930 Args:931 thread_id: The thread ID to cancel.932 run_id: The run ID to cancel.933 wait: Whether to wait until run has completed.934 action: Action to take when cancelling the run. Possible values935 are `interrupt` or `rollback`. Default is `interrupt`.936 headers: Optional custom headers to include with the request.937 params: Optional query parameters to include with the request.938 939 Returns:940 `None`941 942 ???+ example "Example Usage"943 944 ```python945 client = get_sync_client(url="http://localhost:2024")946 client.runs.cancel(947 thread_id="thread_id_to_cancel",948 run_id="run_id_to_cancel",949 wait=True,950 action="interrupt"951 )952 ```953 954 """955 query_params = {956 "wait": 1 if wait else 0,957 "action": action,958 }959 if params:960 query_params.update(params)961 if wait:962 return self.http.request_reconnect(963 f"/threads/{thread_id}/runs/{run_id}/cancel",964 "POST",965 json=None,966 params=query_params,967 headers=headers,968 )969 return self.http.post(970 f"/threads/{thread_id}/runs/{run_id}/cancel",971 json=None,972 params=query_params,973 headers=headers,974 )975 976 def cancel_many(977 self,978 *,979 thread_id: str | None = None,980 run_ids: Sequence[str] | None = None,981 status: BulkCancelRunsStatus | None = None,982 action: CancelAction = "interrupt",983 headers: Mapping[str, str] | None = None,984 params: QueryParamTypes | None = None,985 ) -> None:986 """Cancel one or more runs.987 988 Can cancel runs by thread ID and run IDs, or by status filter.989 990 Args:991 thread_id: The ID of the thread containing runs to cancel.992 run_ids: List of run IDs to cancel.993 status: Filter runs by status to cancel. Must be one of994 `"pending"`, `"running"`, or `"all"`.995 action: Action to take when cancelling the run. Possible values996 are `"interrupt"` or `"rollback"`. Default is `"interrupt"`.997 headers: Optional custom headers to include with the request.998 params: Optional query parameters to include with the request.999 1000 Returns:1001 `None`1002 1003 ???+ example "Example Usage"1004 1005 ```python1006 client = get_sync_client(url="http://localhost:2024")1007 # Cancel all pending runs1008 client.runs.cancel_many(status="pending")1009 # Cancel specific runs on a thread1010 client.runs.cancel_many(1011 thread_id="my_thread_id",1012 run_ids=["run_1", "run_2"],1013 action="rollback",1014 )1015 ```1016 1017 """1018 payload: dict[str, Any] = {}1019 if thread_id:1020 payload["thread_id"] = thread_id1021 if run_ids:1022 payload["run_ids"] = run_ids1023 if status:1024 payload["status"] = status1025 query_params: dict[str, Any] = {"action": action}1026 if params:1027 query_params.update(params)1028 self.http.post(1029 "/runs/cancel",1030 json=payload,1031 headers=headers,1032 params=query_params,1033 )1034 1035 def join(1036 self,1037 thread_id: str,1038 run_id: str,1039 *,1040 headers: Mapping[str, str] | None = None,1041 params: QueryParamTypes | None = None,1042 ) -> dict:1043 """Block until a run is done. Returns the final state of the thread.1044 1045 Args:1046 thread_id: The thread ID to join.1047 run_id: The run ID to join.1048 headers: Optional custom headers to include with the request.1049 params: Optional query parameters to include with the request.1050 1051 Returns:1052 `None`1053 1054 ???+ example "Example Usage"1055 1056 ```python1057 client = get_sync_client(url="http://localhost:2024")1058 client.runs.join(1059 thread_id="thread_id_to_join",1060 run_id="run_id_to_join"1061 )1062 ```1063 1064 """1065 return self.http.request_reconnect(1066 f"/threads/{thread_id}/runs/{run_id}/join",1067 "GET",1068 headers=headers,1069 params=params,1070 )1071 1072 def join_stream(1073 self,1074 thread_id: str,1075 run_id: str,1076 *,1077 cancel_on_disconnect: bool = False,1078 stream_mode: StreamMode | Sequence[StreamMode] | None = None,1079 headers: Mapping[str, str] | None = None,1080 params: QueryParamTypes | None = None,1081 last_event_id: str | None = None,1082 ) -> Iterator[StreamPart]:1083 """Stream output from a run in real-time, until the run is done.1084 Output is not buffered, so any output produced before this call will1085 not be received here.1086 1087 Args:1088 thread_id: The thread ID to join.1089 run_id: The run ID to join.1090 stream_mode: The stream mode(s) to use. Must be a subset of the stream modes passed1091 when creating the run. Background runs default to having the union of all1092 stream modes.1093 cancel_on_disconnect: Whether to cancel the run when the stream is disconnected.1094 headers: Optional custom headers to include with the request.1095 params: Optional query parameters to include with the request.1096 last_event_id: The last event ID to use for the stream.1097 1098 Returns:1099 `None`1100 1101 ???+ example "Example Usage"1102 1103 ```python1104 client = get_sync_client(url="http://localhost:2024")1105 client.runs.join_stream(1106 thread_id="thread_id_to_join",1107 run_id="run_id_to_join",1108 stream_mode=["values", "debug"]1109 )1110 ```1111 1112 """1113 query_params = {1114 "stream_mode": stream_mode,1115 "cancel_on_disconnect": cancel_on_disconnect,1116 }1117 if params:1118 query_params.update(params)1119 return self.http.stream(1120 f"/threads/{thread_id}/runs/{run_id}/stream",1121 "GET",1122 params=query_params,1123 headers={1124 **({"Last-Event-ID": last_event_id} if last_event_id else {}),1125 **(headers or {}),1126 }1127 or None,1128 )1129 1130 def delete(1131 self,1132 thread_id: str,1133 run_id: str,1134 *,1135 headers: Mapping[str, str] | None = None,1136 params: QueryParamTypes | None = None,1137 ) -> None:1138 """Delete a run.1139 1140 Args:1141 thread_id: The thread ID to delete.1142 run_id: The run ID to delete.1143 headers: Optional custom headers to include with the request.1144 params: Optional query parameters to include with the request.1145 1146 Returns:1147 `None`1148 1149 ???+ example "Example Usage"1150 1151 ```python1152 client = get_sync_client(url="http://localhost:2024")1153 client.runs.delete(1154 thread_id="thread_id_to_delete",1155 run_id="run_id_to_delete"1156 )1157 ```1158 1159 """1160 self.http.delete(1161 f"/threads/{thread_id}/runs/{run_id}", headers=headers, params=params1162 )1163 