Team Ai
Datasetpublic

codekingpro/portable-devtools

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