Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
runs.py1163 linesDownload Raw Back to _sync
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 
codekingpro/portable-devtools · Team Ai