Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
cron.py522 linesDownload Raw Back to _async
1"""Async client for managing recurrent runs (cron jobs) in LangGraph."""2 3from __future__ import annotations4 5import warnings6from collections.abc import Mapping, Sequence7from datetime import datetime, tzinfo8from typing import Any9 10from langgraph_sdk._async.http import HttpClient11from langgraph_sdk._shared.utilities import _resolve_timezone12from langgraph_sdk.schema import (13    All,14    Config,15    Context,16    Cron,17    CronSelectField,18    CronSortBy,19    Durability,20    Input,21    OnCompletionBehavior,22    QueryParamTypes,23    Run,24    SortOrder,25    StreamMode,26)27 28 29class CronClient:30    """Client for managing recurrent runs (cron jobs) in LangGraph.31 32    A run is a single invocation of an assistant with optional input, config, and context.33    This client allows scheduling recurring runs to occur automatically.34 35    ???+ example "Example Usage"36 37        ```python38        client = get_client(url="http://localhost:2024"))39        cron_job = await client.crons.create_for_thread(40            thread_id="thread_123",41            assistant_id="asst_456",42            schedule="0 9 * * *",43            input={"message": "Daily update"}44        )45        ```46 47    !!! note "Feature Availability"48 49        The crons client functionality is not supported on all licenses.50        Please check the relevant license documentation for the most up-to-date51        details on feature availability.52    """53 54    def __init__(self, http_client: HttpClient) -> None:55        self.http = http_client56 57    async def create_for_thread(58        self,59        thread_id: str,60        assistant_id: str,61        *,62        schedule: str,63        input: Input | None = None,64        metadata: Mapping[str, Any] | None = None,65        config: Config | None = None,66        context: Context | None = None,67        checkpoint_during: bool | None = None,  # deprecated68        interrupt_before: All | list[str] | None = None,69        interrupt_after: All | list[str] | None = None,70        webhook: str | None = None,71        multitask_strategy: str | None = None,72        end_time: datetime | None = None,73        enabled: bool | None = None,74        timezone: str | tzinfo | None = None,75        stream_mode: StreamMode | Sequence[StreamMode] | None = None,76        stream_subgraphs: bool | None = None,77        stream_resumable: bool | None = None,78        durability: Durability | None = None,79        headers: Mapping[str, str] | None = None,80        params: QueryParamTypes | None = None,81    ) -> Run:82        """Create a cron job for a thread.83 84        Args:85            thread_id: the thread ID to run the cron job on.86            assistant_id: The assistant ID or graph name to use for the cron job.87                If using graph name, will default to first assistant created from that graph.88            schedule: The cron schedule to execute this job on.89                Schedules are interpreted in UTC unless a timezone is specified.90            input: The input to the graph.91            metadata: Metadata to assign to the cron job runs.92            config: The configuration for the assistant.93            context: Static context to add to the assistant.94                !!! version-added "Added in version 0.6.0"95            checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).96            interrupt_before: Nodes to interrupt immediately before they get executed.97 98            interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.99 100            webhook: Webhook to call after LangGraph API call is done.101            multitask_strategy: Multitask strategy to use.102                Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.103            end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.104            enabled: Whether the cron job is enabled or not.105            timezone: IANA timezone for the cron schedule. Accepts a string (e.g. 'America/New_York') or a ``datetime.tzinfo`` instance (e.g. ``ZoneInfo("America/New_York")``).106            stream_mode: The stream mode(s) to use.107            stream_subgraphs: Whether to stream output from subgraphs.108            stream_resumable: Whether to persist the stream chunks in order to resume the stream later.109            durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.110                "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True111                "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False112                "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps113            headers: Optional custom headers to include with the request.114            params: Optional query parameters to include with the request.115 116        Returns:117            The cron run.118 119        ???+ example "Example Usage"120 121            ```python122            client = get_client(url="http://localhost:2024")123            cron_run = await client.crons.create_for_thread(124                thread_id="my-thread-id",125                assistant_id="agent",126                schedule="27 15 * * *",127                input={"messages": [{"role": "user", "content": "hello!"}]},128                metadata={"name":"my_run"},129                context={"model_name": "openai"},130                interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],131                interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],132                webhook="https://my.fake.webhook.com",133                multitask_strategy="interrupt",134                enabled=True,135            )136            ```137        """138        if checkpoint_during is not None:139            warnings.warn(140                "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",141                DeprecationWarning,142                stacklevel=2,143            )144 145        payload = {146            "schedule": schedule,147            "input": input,148            "config": config,149            "metadata": metadata,150            "context": context,151            "assistant_id": assistant_id,152            "checkpoint_during": checkpoint_during,153            "interrupt_before": interrupt_before,154            "interrupt_after": interrupt_after,155            "webhook": webhook,156            "end_time": end_time.isoformat() if end_time else None,157            "enabled": enabled,158            "timezone": _resolve_timezone(timezone),159            "stream_mode": stream_mode,160            "stream_subgraphs": stream_subgraphs,161            "stream_resumable": stream_resumable,162            "durability": durability,163        }164        if multitask_strategy:165            payload["multitask_strategy"] = multitask_strategy166        payload = {k: v for k, v in payload.items() if v is not None}167        return await self.http.post(168            f"/threads/{thread_id}/runs/crons",169            json=payload,170            headers=headers,171            params=params,172        )173 174    async def create(175        self,176        assistant_id: str,177        *,178        schedule: str,179        input: Input | None = None,180        metadata: Mapping[str, Any] | None = None,181        config: Config | None = None,182        context: Context | None = None,183        checkpoint_during: bool | None = None,  # deprecated184        interrupt_before: All | list[str] | None = None,185        interrupt_after: All | list[str] | None = None,186        webhook: str | None = None,187        on_run_completed: OnCompletionBehavior | None = None,188        multitask_strategy: str | None = None,189        end_time: datetime | None = None,190        enabled: bool | None = None,191        timezone: str | tzinfo | None = None,192        stream_mode: StreamMode | Sequence[StreamMode] | None = None,193        stream_subgraphs: bool | None = None,194        stream_resumable: bool | None = None,195        durability: Durability | None = None,196        headers: Mapping[str, str] | None = None,197        params: QueryParamTypes | None = None,198    ) -> Run:199        """Create a cron run.200 201        Args:202            assistant_id: The assistant ID or graph name to use for the cron job.203                If using graph name, will default to first assistant created from that graph.204            schedule: The cron schedule to execute this job on.205                Schedules are interpreted in UTC unless a timezone is specified.206            input: The input to the graph.207            metadata: Metadata to assign to the cron job runs.208            config: The configuration for the assistant.209            context: Static context to add to the assistant.210                !!! version-added "Added in version 0.6.0"211            checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).212            interrupt_before: Nodes to interrupt immediately before they get executed.213            interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.214            webhook: Webhook to call after LangGraph API call is done.215            on_run_completed: What to do with the thread after the run completes.216                Must be one of 'delete' (default) or 'keep'. 'delete' removes the thread217                after execution. 'keep' creates a new thread for each execution but does not218                clean them up. Clients are responsible for cleaning up kept threads.219            multitask_strategy: Multitask strategy to use.220                Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.221            end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.222            enabled: Whether the cron job is enabled or not.223            timezone: IANA timezone for the cron schedule. Accepts a string (e.g. 'America/New_York') or a ``datetime.tzinfo`` instance (e.g. ``ZoneInfo("America/New_York")``).224            stream_mode: The stream mode(s) to use.225            stream_subgraphs: Whether to stream output from subgraphs.226            stream_resumable: Whether to persist the stream chunks in order to resume the stream later.227            durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.228                "async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True229                "sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False230                "exit" means checkpoints are only persisted when the run exits, does not save intermediate steps231            headers: Optional custom headers to include with the request.232            params: Optional query parameters to include with the request.233 234        Returns:235            The cron run.236 237        ???+ example "Example Usage"238 239            ```python240            client = get_client(url="http://localhost:2024")241            cron_run = client.crons.create(242                assistant_id="agent",243                schedule="27 15 * * *",244                input={"messages": [{"role": "user", "content": "hello!"}]},245                metadata={"name":"my_run"},246                context={"model_name": "openai"},247                interrupt_before=["node_to_stop_before_1","node_to_stop_before_2"],248                interrupt_after=["node_to_stop_after_1","node_to_stop_after_2"],249                webhook="https://my.fake.webhook.com",250                multitask_strategy="interrupt",251                enabled=True,252            )253            ```254 255        """256        if checkpoint_during is not None:257            warnings.warn(258                "`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",259                DeprecationWarning,260                stacklevel=2,261            )262 263        payload = {264            "schedule": schedule,265            "input": input,266            "config": config,267            "metadata": metadata,268            "context": context,269            "assistant_id": assistant_id,270            "checkpoint_during": checkpoint_during,271            "interrupt_before": interrupt_before,272            "interrupt_after": interrupt_after,273            "webhook": webhook,274            "on_run_completed": on_run_completed,275            "end_time": end_time.isoformat() if end_time else None,276            "enabled": enabled,277            "timezone": _resolve_timezone(timezone),278            "stream_mode": stream_mode,279            "stream_subgraphs": stream_subgraphs,280            "stream_resumable": stream_resumable,281            "durability": durability,282        }283        if multitask_strategy:284            payload["multitask_strategy"] = multitask_strategy285        payload = {k: v for k, v in payload.items() if v is not None}286        return await self.http.post(287            "/runs/crons", json=payload, headers=headers, params=params288        )289 290    async def delete(291        self,292        cron_id: str,293        *,294        headers: Mapping[str, str] | None = None,295        params: QueryParamTypes | None = None,296    ) -> None:297        """Delete a cron.298 299        Args:300            cron_id: The cron ID to delete.301            headers: Optional custom headers to include with the request.302            params: Optional query parameters to include with the request.303 304        Returns:305            `None`306 307        ???+ example "Example Usage"308 309            ```python310            client = get_client(url="http://localhost:2024")311            await client.crons.delete(312                cron_id="cron_to_delete"313            )314            ```315 316        """317        await self.http.delete(f"/runs/crons/{cron_id}", headers=headers, params=params)318 319    async def update(320        self,321        cron_id: str,322        *,323        schedule: str | None = None,324        end_time: datetime | None = None,325        input: Input | None = None,326        metadata: Mapping[str, Any] | None = None,327        config: Config | None = None,328        context: Context | None = None,329        webhook: str | None = None,330        interrupt_before: All | list[str] | None = None,331        interrupt_after: All | list[str] | None = None,332        on_run_completed: OnCompletionBehavior | None = None,333        enabled: bool | None = None,334        timezone: str | tzinfo | None = None,335        stream_mode: StreamMode | Sequence[StreamMode] | None = None,336        stream_subgraphs: bool | None = None,337        stream_resumable: bool | None = None,338        durability: Durability | None = None,339        headers: Mapping[str, str] | None = None,340        params: QueryParamTypes | None = None,341    ) -> Cron:342        """Update a cron job by ID.343 344        Args:345            cron_id: The cron ID to update.346            schedule: The cron schedule to execute this job on.347                Schedules are interpreted in UTC unless a timezone is specified.348            end_time: The end date to stop running the cron.349            input: The input to the graph.350            metadata: Metadata to assign to the cron job runs.351            config: The configuration for the assistant.352            context: Static context added to the assistant.353            webhook: Webhook to call after LangGraph API call is done.354            interrupt_before: Nodes to interrupt immediately before they get executed.355            interrupt_after: Nodes to interrupt immediately after they get executed.356            on_run_completed: What to do with the thread after the run completes.357                Must be one of 'delete' or 'keep'. 'delete' removes the thread358                after execution. 'keep' creates a new thread for each execution but does not359                clean them up.360            enabled: Enable or disable the cron job.361            timezone: IANA timezone for the cron schedule. Accepts a string (e.g. 'America/New_York') or a ``datetime.tzinfo`` instance (e.g. ``ZoneInfo("America/New_York")``).362            stream_mode: The stream mode(s) to use.363            stream_subgraphs: Whether to stream output from subgraphs.364            stream_resumable: Whether to persist the stream chunks in order to resume the stream later.365            durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.366            headers: Optional custom headers to include with the request.367            params: Optional query parameters to include with the request.368 369        Returns:370            The updated cron job.371 372        ???+ example "Example Usage"373 374            ```python375            client = get_client(url="http://localhost:2024")376            updated_cron = await client.crons.update(377                cron_id="1ef3cefa-4c09-6926-96d0-3dc97fd5e39b",378                schedule="0 10 * * *",379                enabled=False,380            )381            ```382 383        """384        payload = {385            "schedule": schedule,386            "end_time": end_time.isoformat() if end_time else None,387            "input": input,388            "metadata": metadata,389            "config": config,390            "context": context,391            "webhook": webhook,392            "interrupt_before": interrupt_before,393            "interrupt_after": interrupt_after,394            "on_run_completed": on_run_completed,395            "enabled": enabled,396            "timezone": _resolve_timezone(timezone),397            "stream_mode": stream_mode,398            "stream_subgraphs": stream_subgraphs,399            "stream_resumable": stream_resumable,400            "durability": durability,401        }402        payload = {k: v for k, v in payload.items() if v is not None}403        return await self.http.patch(404            f"/runs/crons/{cron_id}",405            json=payload,406            headers=headers,407            params=params,408        )409 410    async def search(411        self,412        *,413        assistant_id: str | None = None,414        thread_id: str | None = None,415        enabled: bool | None = None,416        limit: int = 10,417        offset: int = 0,418        sort_by: CronSortBy | None = None,419        sort_order: SortOrder | None = None,420        select: list[CronSelectField] | None = None,421        headers: Mapping[str, str] | None = None,422        params: QueryParamTypes | None = None,423    ) -> list[Cron]:424        """Get a list of cron jobs.425 426        Args:427            assistant_id: The assistant ID or graph name to search for.428            thread_id: the thread ID to search for.429            enabled: The enabled status to search for.430            limit: The maximum number of results to return.431            offset: The number of results to skip.432            headers: Optional custom headers to include with the request.433            params: Optional query parameters to include with the request.434 435        Returns:436            The list of cron jobs returned by the search,437 438        ???+ example "Example Usage"439 440            ```python441            client = get_client(url="http://localhost:2024")442            cron_jobs = await client.crons.search(443                assistant_id="my_assistant_id",444                thread_id="my_thread_id",445                enabled=True,446                limit=5,447                offset=5,448            )449            print(cron_jobs)450            ```451            ```shell452 453            ----------------------------------------------------------454 455            [456                {457                    'cron_id': '1ef3cefa-4c09-6926-96d0-3dc97fd5e39b',458                    'assistant_id': 'my_assistant_id',459                    'thread_id': 'my_thread_id',460                    'user_id': None,461                    'payload':462                        {463                            'input': {'start_time': ''},464                            'schedule': '4 * * * *',465                            'assistant_id': 'my_assistant_id'466                        },467                    'schedule': '4 * * * *',468                    'next_run_date': '2024-07-25T17:04:00+00:00',469                    'end_time': None,470                    'created_at': '2024-07-08T06:02:23.073257+00:00',471                    'updated_at': '2024-07-08T06:02:23.073257+00:00'472                }473            ]474            ```475 476        """477        payload: dict[str, Any] = {478            "assistant_id": assistant_id,479            "thread_id": thread_id,480            "enabled": enabled,481            "limit": limit,482            "offset": offset,483        }484        if sort_by:485            payload["sort_by"] = sort_by486        if sort_order:487            payload["sort_order"] = sort_order488        if select:489            payload["select"] = select490        payload = {k: v for k, v in payload.items() if v is not None}491        return await self.http.post(492            "/runs/crons/search", json=payload, headers=headers, params=params493        )494 495    async def count(496        self,497        *,498        assistant_id: str | None = None,499        thread_id: str | None = None,500        headers: Mapping[str, str] | None = None,501        params: QueryParamTypes | None = None,502    ) -> int:503        """Count cron jobs matching filters.504 505        Args:506            assistant_id: Assistant ID to filter by.507            thread_id: Thread ID to filter by.508            headers: Optional custom headers to include with the request.509            params: Optional query parameters to include with the request.510 511        Returns:512            int: Number of crons matching the criteria.513        """514        payload: dict[str, Any] = {}515        if assistant_id:516            payload["assistant_id"] = assistant_id517        if thread_id:518            payload["thread_id"] = thread_id519        return await self.http.post(520            "/runs/crons/count", json=payload, headers=headers, params=params521        )522 
codekingpro/portable-devtools · Team Ai