codekingpro/portable-devtools
114k
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 