mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-10-03 06:55:13 +02:00
Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8563948e70 | ||
|
|
b0f14649e0 | ||
|
|
ea20432b9b | ||
|
|
e2efab8061 | ||
|
|
9babffa054 |
@@ -390,6 +390,142 @@ class PostgresSaver(BasePostgresSaver):
|
||||
(str(thread_id),),
|
||||
)
|
||||
|
||||
def delete_for_runs(self, run_ids: Sequence[str]) -> None:
|
||||
"""Delete all checkpoints and writes for the given run IDs.
|
||||
|
||||
Args:
|
||||
run_ids: The run IDs whose checkpoints should be deleted.
|
||||
"""
|
||||
if not run_ids:
|
||||
return
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE (thread_id, checkpoint_ns, checkpoint_id) IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)""",
|
||||
(list(run_ids),),
|
||||
)
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE (thread_id, checkpoint_ns) IN (
|
||||
SELECT DISTINCT thread_id, checkpoint_ns
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c2
|
||||
WHERE c2.thread_id = checkpoint_blobs.thread_id
|
||||
AND c2.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
AND (c2.metadata->>'run_id' IS NULL OR c2.metadata->>'run_id' != ALL(%s))
|
||||
AND c2.checkpoint->'channel_versions' ? checkpoint_blobs.channel
|
||||
)""",
|
||||
(list(run_ids), list(run_ids)),
|
||||
)
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoints WHERE metadata->>'run_id' = ANY(%s)",
|
||||
(list(run_ids),),
|
||||
)
|
||||
|
||||
def copy_thread(self, source_thread_id: str, target_thread_id: str) -> None:
|
||||
"""Copy all checkpoints and writes from source thread to target thread.
|
||||
|
||||
Args:
|
||||
source_thread_id: The thread ID to copy from.
|
||||
target_thread_id: The thread ID to copy to.
|
||||
"""
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"""INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata
|
||||
FROM checkpoints
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
cur.execute(
|
||||
"""INSERT INTO checkpoint_blobs (thread_id, checkpoint_ns, channel, version, type, blob)
|
||||
SELECT %s, checkpoint_ns, channel, version, type, blob
|
||||
FROM checkpoint_blobs
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, channel, version) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
cur.execute(
|
||||
"""INSERT INTO checkpoint_writes (thread_id, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob
|
||||
FROM checkpoint_writes
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
|
||||
def prune(
|
||||
self,
|
||||
thread_ids: Sequence[str],
|
||||
*,
|
||||
strategy: str = "keep_latest",
|
||||
) -> None:
|
||||
"""Prune checkpoints for given threads.
|
||||
|
||||
Args:
|
||||
thread_ids: The thread IDs to prune.
|
||||
strategy: `"keep_latest"` keeps only the latest checkpoint per
|
||||
thread+namespace. `"delete"` removes everything.
|
||||
"""
|
||||
if not thread_ids:
|
||||
return
|
||||
if strategy == "delete":
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoints WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoint_blobs WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoint_writes WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
elif strategy == "keep_latest":
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT DISTINCT ON (thread_id, checkpoint_ns)
|
||||
thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
ORDER BY thread_id, checkpoint_ns, checkpoint_id DESC
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c
|
||||
WHERE c.thread_id = checkpoint_blobs.thread_id
|
||||
AND c.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
)""",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
|
||||
@contextmanager
|
||||
def _cursor(self, *, pipeline: bool = False) -> Iterator[Cursor[DictRow]]:
|
||||
"""Create a database cursor as a context manager.
|
||||
|
||||
@@ -349,6 +349,152 @@ class AsyncPostgresSaver(BasePostgresSaver):
|
||||
(str(thread_id),),
|
||||
)
|
||||
|
||||
async def adelete_for_runs(self, run_ids: Sequence[str]) -> None:
|
||||
"""Delete all checkpoints and writes for the given run IDs.
|
||||
|
||||
Args:
|
||||
run_ids: The run IDs whose checkpoints should be deleted.
|
||||
"""
|
||||
if not run_ids:
|
||||
return
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
# Delete writes associated with checkpoints that have matching run_ids
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE (thread_id, checkpoint_ns, checkpoint_id) IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)""",
|
||||
(list(run_ids),),
|
||||
)
|
||||
# Delete blobs associated with checkpoints that have matching run_ids
|
||||
# We need to delete blobs for channels referenced by these checkpoints
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE (thread_id, checkpoint_ns) IN (
|
||||
SELECT DISTINCT thread_id, checkpoint_ns
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c2
|
||||
WHERE c2.thread_id = checkpoint_blobs.thread_id
|
||||
AND c2.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
AND (c2.metadata->>'run_id' IS NULL OR c2.metadata->>'run_id' != ALL(%s))
|
||||
AND c2.checkpoint->'channel_versions' ? checkpoint_blobs.channel
|
||||
)""",
|
||||
(list(run_ids), list(run_ids)),
|
||||
)
|
||||
# Delete the checkpoints themselves
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoints WHERE metadata->>'run_id' = ANY(%s)",
|
||||
(list(run_ids),),
|
||||
)
|
||||
|
||||
async def acopy_thread(self, source_thread_id: str, target_thread_id: str) -> None:
|
||||
"""Copy all checkpoints and writes from source thread to target thread.
|
||||
|
||||
Args:
|
||||
source_thread_id: The thread ID to copy from.
|
||||
target_thread_id: The thread ID to copy to.
|
||||
"""
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
# Copy checkpoints
|
||||
await cur.execute(
|
||||
"""INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata
|
||||
FROM checkpoints
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
# Copy blobs
|
||||
await cur.execute(
|
||||
"""INSERT INTO checkpoint_blobs (thread_id, checkpoint_ns, channel, version, type, blob)
|
||||
SELECT %s, checkpoint_ns, channel, version, type, blob
|
||||
FROM checkpoint_blobs
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, channel, version) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
# Copy writes
|
||||
await cur.execute(
|
||||
"""INSERT INTO checkpoint_writes (thread_id, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob
|
||||
FROM checkpoint_writes
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
|
||||
async def aprune(
|
||||
self,
|
||||
thread_ids: Sequence[str],
|
||||
*,
|
||||
strategy: str = "keep_latest",
|
||||
) -> None:
|
||||
"""Prune checkpoints for given threads.
|
||||
|
||||
Args:
|
||||
thread_ids: The thread IDs to prune.
|
||||
strategy: `"keep_latest"` keeps only the latest checkpoint per
|
||||
thread+namespace. `"delete"` removes everything.
|
||||
"""
|
||||
if not thread_ids:
|
||||
return
|
||||
if strategy == "delete":
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoints WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoint_blobs WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoint_writes WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
elif strategy == "keep_latest":
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
# Delete non-latest checkpoints
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT DISTINCT ON (thread_id, checkpoint_ns)
|
||||
thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
ORDER BY thread_id, checkpoint_ns, checkpoint_id DESC
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
# Delete writes for removed checkpoints (keep only writes for remaining checkpoints)
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
# Clean up orphaned blobs
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c
|
||||
WHERE c.thread_id = checkpoint_blobs.thread_id
|
||||
AND c.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
)""",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
|
||||
@asynccontextmanager
|
||||
async def _cursor(
|
||||
self, *, pipeline: bool = False
|
||||
@@ -578,5 +724,43 @@ class AsyncPostgresSaver(BasePostgresSaver):
|
||||
self.adelete_thread(thread_id), self.loop
|
||||
).result()
|
||||
|
||||
def delete_for_runs(self, run_ids: Sequence[str]) -> None:
|
||||
"""Delete all checkpoints and writes for the given run IDs.
|
||||
|
||||
Args:
|
||||
run_ids: The run IDs whose checkpoints should be deleted.
|
||||
"""
|
||||
return asyncio.run_coroutine_threadsafe(
|
||||
self.adelete_for_runs(run_ids), self.loop
|
||||
).result()
|
||||
|
||||
def copy_thread(self, source_thread_id: str, target_thread_id: str) -> None:
|
||||
"""Copy all checkpoints and writes from source thread to target thread.
|
||||
|
||||
Args:
|
||||
source_thread_id: The thread ID to copy from.
|
||||
target_thread_id: The thread ID to copy to.
|
||||
"""
|
||||
return asyncio.run_coroutine_threadsafe(
|
||||
self.acopy_thread(source_thread_id, target_thread_id), self.loop
|
||||
).result()
|
||||
|
||||
def prune(
|
||||
self,
|
||||
thread_ids: Sequence[str],
|
||||
*,
|
||||
strategy: str = "keep_latest",
|
||||
) -> None:
|
||||
"""Prune checkpoints for given threads.
|
||||
|
||||
Args:
|
||||
thread_ids: The thread IDs to prune.
|
||||
strategy: `"keep_latest"` keeps only the latest checkpoint per
|
||||
thread+namespace. `"delete"` removes everything.
|
||||
"""
|
||||
return asyncio.run_coroutine_threadsafe(
|
||||
self.aprune(thread_ids, strategy=strategy), self.loop
|
||||
).result()
|
||||
|
||||
|
||||
__all__ = ["AsyncPostgresSaver", "AsyncShallowPostgresSaver", "Conn"]
|
||||
|
||||
@@ -32,6 +32,7 @@ test = [
|
||||
"pytest-mock",
|
||||
"psycopg[binary]",
|
||||
"langgraph-checkpoint",
|
||||
"langgraph-checkpoint-conformance",
|
||||
"pytest-watcher",
|
||||
]
|
||||
lint = [
|
||||
@@ -49,6 +50,7 @@ default-groups = ['dev']
|
||||
|
||||
[tool.uv.sources]
|
||||
langgraph-checkpoint = { path = "../checkpoint", editable = true }
|
||||
langgraph-checkpoint-conformance = { path = "../checkpoint-conformance", editable = true }
|
||||
|
||||
[tool.hatch.build.targets.wheel]
|
||||
include = ["langgraph"]
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
"""Conformance tests for AsyncPostgresSaver."""
|
||||
# mypy: disable-error-code="import-untyped"
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import AsyncGenerator
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
from langgraph.checkpoint.conformance import checkpointer_test, validate
|
||||
from langgraph.checkpoint.conformance.report import ProgressCallbacks
|
||||
from psycopg import AsyncConnection
|
||||
from psycopg.rows import dict_row
|
||||
|
||||
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
|
||||
from tests.conftest import DEFAULT_POSTGRES_URI
|
||||
|
||||
|
||||
async def pg_lifespan() -> AsyncGenerator[None, None]:
|
||||
"""No-op lifespan; databases are created per-checkpointer instance."""
|
||||
yield
|
||||
|
||||
|
||||
@checkpointer_test(name="AsyncPostgresSaver", lifespan=pg_lifespan)
|
||||
async def postgres_checkpointer() -> AsyncGenerator[AsyncPostgresSaver, None]:
|
||||
database = f"test_{uuid4().hex[:16]}"
|
||||
async with await AsyncConnection.connect(
|
||||
DEFAULT_POSTGRES_URI, autocommit=True
|
||||
) as conn:
|
||||
await conn.execute(f"CREATE DATABASE {database}")
|
||||
try:
|
||||
async with await AsyncConnection.connect(
|
||||
DEFAULT_POSTGRES_URI + database,
|
||||
autocommit=True,
|
||||
prepare_threshold=0,
|
||||
row_factory=dict_row,
|
||||
) as conn:
|
||||
saver = AsyncPostgresSaver(conn)
|
||||
await saver.setup()
|
||||
yield saver
|
||||
finally:
|
||||
async with await AsyncConnection.connect(
|
||||
DEFAULT_POSTGRES_URI, autocommit=True
|
||||
) as conn:
|
||||
await conn.execute(f"DROP DATABASE {database}")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_full_conformance() -> None:
|
||||
"""AsyncPostgresSaver passes ALL conformance tests."""
|
||||
report = await validate(
|
||||
postgres_checkpointer,
|
||||
progress=ProgressCallbacks.verbose(),
|
||||
)
|
||||
report.print_report()
|
||||
assert report.passed_all(), f"Conformance failed: {report.to_dict()}"
|
||||
Generated
+31
@@ -304,6 +304,33 @@ test = [
|
||||
{ name = "redis" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint-conformance"
|
||||
version = "0.0.1"
|
||||
source = { editable = "../checkpoint-conformance" }
|
||||
dependencies = [
|
||||
{ name = "langgraph-checkpoint" },
|
||||
]
|
||||
|
||||
[package.metadata]
|
||||
requires-dist = [{ name = "langgraph-checkpoint", specifier = ">=2.0.0" }]
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
dev = [
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
{ name = "ruff" },
|
||||
{ name = "ty" },
|
||||
]
|
||||
lint = [
|
||||
{ name = "ruff" },
|
||||
{ name = "ty" },
|
||||
]
|
||||
test = [
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint-postgres"
|
||||
version = "3.0.4"
|
||||
@@ -320,6 +347,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extra = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -336,6 +364,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance" },
|
||||
{ name = "psycopg", extra = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
@@ -356,6 +385,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -372,6 +402,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
|
||||
@@ -3533,9 +3533,9 @@ js-tokens@^4.0.0:
|
||||
integrity sha512-RdJUflcE3cUzKiMqQgsCu06FPu9UdIJO0beYbPhHN4k6apgJtifcoCtT9bcxOpYBtpD2kCM6Sbzg4CausW/PKQ==
|
||||
|
||||
js-yaml@^3.13.1:
|
||||
version "3.14.1"
|
||||
resolved "https://registry.yarnpkg.com/js-yaml/-/js-yaml-3.14.1.tgz#dae812fdb3825fa306609a8717383c50c36a0537"
|
||||
integrity sha512-okMH7OXXJ7YrN9Ok3/SXrnu4iX9yOk+25nqX4imS2npuvTYDmo/QEZoqwZkYaIDk3jVvBOTOIEgEhaLOynBS9g==
|
||||
version "3.14.2"
|
||||
resolved "https://registry.yarnpkg.com/js-yaml/-/js-yaml-3.14.2.tgz#77485ce1dd7f33c061fd1b16ecea23b55fcb04b0"
|
||||
integrity sha512-PMSmkqxr106Xa156c2M265Z+FTrPl+oxd/rgOQy2tijQeK5TxQ43psO1ZCwhVOSdnn+RzkzlRz/eY4BgJBYVpg==
|
||||
dependencies:
|
||||
argparse "^1.0.7"
|
||||
esprima "^4.0.0"
|
||||
|
||||
@@ -678,6 +678,51 @@ def _update_encryption_path(
|
||||
)
|
||||
|
||||
|
||||
def _update_checkpointer_path(
|
||||
config_path: pathlib.Path, config: Config, local_deps: LocalDeps
|
||||
) -> None:
|
||||
"""Update checkpointer.path to use Docker container paths."""
|
||||
checkpointer_conf = config.get("checkpointer")
|
||||
if not checkpointer_conf or not isinstance(checkpointer_conf, dict):
|
||||
return
|
||||
if not (path_str := checkpointer_conf.get("path")):
|
||||
return
|
||||
|
||||
module_str, sep, attr_str = path_str.partition(":")
|
||||
if not sep or not module_str.startswith("."):
|
||||
return # Already validated or absolute path
|
||||
|
||||
resolved = config_path.parent / module_str
|
||||
if not resolved.exists():
|
||||
raise FileNotFoundError(
|
||||
f"Checkpointer file not found: {resolved} (from {path_str})"
|
||||
)
|
||||
if not resolved.is_file():
|
||||
raise IsADirectoryError(f"Checkpointer path must be a file: {resolved}")
|
||||
|
||||
# Check faux packages first (higher priority)
|
||||
for faux_path, (_, destpath) in local_deps.faux_pkgs.items():
|
||||
if resolved.is_relative_to(faux_path):
|
||||
new_path = f"{destpath}/{resolved.relative_to(faux_path)}:{attr_str}"
|
||||
checkpointer_conf["path"] = new_path
|
||||
return
|
||||
|
||||
# Check real packages
|
||||
for real_path in local_deps.real_pkgs:
|
||||
if resolved.is_relative_to(real_path):
|
||||
new_path = (
|
||||
f"/deps/{real_path.name}/{resolved.relative_to(real_path)}:{attr_str}"
|
||||
)
|
||||
checkpointer_conf["path"] = new_path
|
||||
return
|
||||
|
||||
raise ValueError(
|
||||
f"Checkpointer file '{resolved}' not covered by dependencies.\n"
|
||||
"Add its parent directory to the 'dependencies' array in your config.\n"
|
||||
f"Current dependencies: {config['dependencies']}"
|
||||
)
|
||||
|
||||
|
||||
def _update_http_app_path(
|
||||
config_path: pathlib.Path, config: Config, local_deps: LocalDeps
|
||||
) -> None:
|
||||
@@ -877,6 +922,8 @@ def python_config_to_docker(
|
||||
_update_auth_path(config_path, config, local_deps)
|
||||
# Rewrite encryption path, so it points to the correct location in the Docker container
|
||||
_update_encryption_path(config_path, config, local_deps)
|
||||
# Rewrite checkpointer path, so it points to the correct location in the Docker container
|
||||
_update_checkpointer_path(config_path, config, local_deps)
|
||||
# Rewrite HTTP app path, so it points to the correct location in the Docker container
|
||||
_update_http_app_path(config_path, config, local_deps)
|
||||
|
||||
|
||||
@@ -167,6 +167,26 @@ class CheckpointerConfig(TypedDict, total=False):
|
||||
If omitted, no checkpointer is set up (the object store will still be present, however).
|
||||
"""
|
||||
|
||||
path: str
|
||||
"""Import path to an async context manager that yields a `BaseCheckpointSaver`
|
||||
instance.
|
||||
|
||||
The referenced object should be an `@asynccontextmanager`-decorated function
|
||||
so that the server can properly manage the checkpointer's lifecycle (e.g.
|
||||
opening and closing connections).
|
||||
|
||||
Examples:
|
||||
- "./my_checkpointer.py:create_checkpointer"
|
||||
- "my_package.checkpointer:create_checkpointer"
|
||||
|
||||
When provided, this replaces the default checkpointer.
|
||||
|
||||
You can use the `langgraph-checkpoint-conformance` package
|
||||
(https://pypi.org/project/langgraph-checkpoint-conformance/) to run simple
|
||||
conformance tests against your custom checkpointer and catch
|
||||
incompatibilities early.
|
||||
"""
|
||||
|
||||
ttl: ThreadTTLConfig | None
|
||||
"""Optional. Defines the TTL (time-to-live) behavior configuration.
|
||||
|
||||
|
||||
@@ -542,6 +542,10 @@
|
||||
"description": "Configuration for the built-in checkpointer, which handles checkpointing of state.\n\nIf omitted, no checkpointer is set up (the object store will still be present, however).",
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"path": {
|
||||
"type": "string",
|
||||
"description": "Import path to an async context manager that yields a `BaseCheckpointSaver`\ninstance.\n\nThe referenced object should be an `@asynccontextmanager`-decorated function\nso that the server can properly manage the checkpointer's lifecycle (e.g.\nopening and closing connections).\n"
|
||||
},
|
||||
"serde": {
|
||||
"anyOf": [
|
||||
{
|
||||
|
||||
@@ -542,6 +542,10 @@
|
||||
"description": "Configuration for the built-in checkpointer, which handles checkpointing of state.\n\nIf omitted, no checkpointer is set up (the object store will still be present, however).",
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"path": {
|
||||
"type": "string",
|
||||
"description": "Import path to an async context manager that yields a `BaseCheckpointSaver`\ninstance.\n\nThe referenced object should be an `@asynccontextmanager`-decorated function\nso that the server can properly manage the checkpointer's lifecycle (e.g.\nopening and closing connections).\n"
|
||||
},
|
||||
"serde": {
|
||||
"anyOf": [
|
||||
{
|
||||
|
||||
Generated
+2
@@ -1617,6 +1617,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -1633,6 +1634,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
|
||||
Generated
+2
@@ -421,6 +421,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -437,6 +438,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
|
||||
@@ -3,6 +3,6 @@ from langgraph_sdk.client import get_client, get_sync_client
|
||||
from langgraph_sdk.encryption import Encryption
|
||||
from langgraph_sdk.encryption.types import EncryptionContext
|
||||
|
||||
__version__ = "0.3.7"
|
||||
__version__ = "0.3.8"
|
||||
|
||||
__all__ = ["Auth", "Encryption", "EncryptionContext", "get_client", "get_sync_client"]
|
||||
|
||||
@@ -2,7 +2,8 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
import warnings
|
||||
from collections.abc import Mapping, Sequence
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
@@ -14,11 +15,13 @@ from langgraph_sdk.schema import (
|
||||
Cron,
|
||||
CronSelectField,
|
||||
CronSortBy,
|
||||
Durability,
|
||||
Input,
|
||||
OnCompletionBehavior,
|
||||
QueryParamTypes,
|
||||
Run,
|
||||
SortOrder,
|
||||
StreamMode,
|
||||
)
|
||||
|
||||
|
||||
@@ -60,13 +63,17 @@ class CronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -83,7 +90,7 @@ class CronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
@@ -93,6 +100,13 @@ class CronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled or not.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -118,6 +132,13 @@ class CronClient:
|
||||
)
|
||||
```
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -131,6 +152,10 @@ class CronClient:
|
||||
"webhook": webhook,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
if multitask_strategy:
|
||||
payload["multitask_strategy"] = multitask_strategy
|
||||
@@ -151,7 +176,7 @@ class CronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
@@ -159,6 +184,10 @@ class CronClient:
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -174,7 +203,7 @@ class CronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
webhook: Webhook to call after LangGraph API call is done.
|
||||
@@ -186,6 +215,13 @@ class CronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled or not.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -211,6 +247,13 @@ class CronClient:
|
||||
```
|
||||
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -225,6 +268,10 @@ class CronClient:
|
||||
"on_run_completed": on_run_completed,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
if multitask_strategy:
|
||||
payload["multitask_strategy"] = multitask_strategy
|
||||
@@ -277,6 +324,10 @@ class CronClient:
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
on_run_completed: OnCompletionBehavior | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Cron:
|
||||
@@ -299,6 +350,10 @@ class CronClient:
|
||||
after execution. 'keep' creates a new thread for each execution but does not
|
||||
clean them up.
|
||||
enabled: Enable or disable the cron job.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -329,6 +384,10 @@ class CronClient:
|
||||
"interrupt_after": interrupt_after,
|
||||
"on_run_completed": on_run_completed,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return await self.http.patch(
|
||||
|
||||
@@ -2,7 +2,8 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
import warnings
|
||||
from collections.abc import Mapping, Sequence
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
@@ -14,11 +15,13 @@ from langgraph_sdk.schema import (
|
||||
Cron,
|
||||
CronSelectField,
|
||||
CronSortBy,
|
||||
Durability,
|
||||
Input,
|
||||
OnCompletionBehavior,
|
||||
QueryParamTypes,
|
||||
Run,
|
||||
SortOrder,
|
||||
StreamMode,
|
||||
)
|
||||
|
||||
|
||||
@@ -54,13 +57,17 @@ class SyncCronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -77,7 +84,7 @@ class SyncCronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
webhook: Webhook to call after LangGraph API call is done.
|
||||
@@ -85,6 +92,13 @@ class SyncCronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled. By default, it is considered enabled.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
|
||||
Returns:
|
||||
@@ -109,6 +123,13 @@ class SyncCronClient:
|
||||
)
|
||||
```
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -123,6 +144,10 @@ class SyncCronClient:
|
||||
"multitask_strategy": multitask_strategy,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return self.http.post(
|
||||
@@ -141,7 +166,7 @@ class SyncCronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
@@ -149,6 +174,10 @@ class SyncCronClient:
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -164,7 +193,7 @@ class SyncCronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
webhook: Webhook to call after LangGraph API call is done.
|
||||
@@ -176,6 +205,13 @@ class SyncCronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled. By default, it is considered enabled.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
|
||||
Returns:
|
||||
@@ -201,6 +237,13 @@ class SyncCronClient:
|
||||
```
|
||||
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -216,6 +259,10 @@ class SyncCronClient:
|
||||
"multitask_strategy": multitask_strategy,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return self.http.post(
|
||||
@@ -266,6 +313,10 @@ class SyncCronClient:
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
on_run_completed: OnCompletionBehavior | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Cron:
|
||||
@@ -288,6 +339,10 @@ class SyncCronClient:
|
||||
after execution. 'keep' creates a new thread for each execution but does not
|
||||
clean them up.
|
||||
enabled: Enable or disable the cron job.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -318,6 +373,10 @@ class SyncCronClient:
|
||||
"interrupt_after": interrupt_after,
|
||||
"on_run_completed": on_run_completed,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return self.http.patch(
|
||||
|
||||
@@ -424,6 +424,14 @@ class CronUpdate(TypedDict, total=False):
|
||||
"""What to do with the thread after the run completes."""
|
||||
enabled: bool
|
||||
"""Enable or disable the cron job."""
|
||||
stream_mode: StreamMode | list[StreamMode]
|
||||
"""The stream mode(s) to use."""
|
||||
stream_subgraphs: bool
|
||||
"""Whether to stream output from subgraphs."""
|
||||
stream_resumable: bool
|
||||
"""Whether to persist the stream chunks in order to resume the stream later."""
|
||||
durability: Durability
|
||||
"""Durability level for the run. Must be one of 'sync', 'async', or 'exit'."""
|
||||
|
||||
|
||||
# Select field aliases for client-side typing of `select` parameters.
|
||||
|
||||
Reference in New Issue
Block a user