mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-09-06 09:47:51 +02:00
Compare commits
36
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9681ecef12 | ||
|
|
81bf17b231 | ||
|
|
11738d83db | ||
|
|
1983e60971 | ||
|
|
d07198a53c | ||
|
|
c0a13bb22d | ||
|
|
7bfee11ad5 | ||
|
|
eae02c622d | ||
|
|
bbd6e2da0b | ||
|
|
b7bcf0083c | ||
|
|
7815573c76 | ||
|
|
0cf2f960f6 | ||
|
|
bd1f031ac0 | ||
|
|
d80317dd6c | ||
|
|
394c109737 | ||
|
|
0c1a062f88 | ||
|
|
4b29b60d13 | ||
|
|
59738d1550 | ||
|
|
2efb0073a5 | ||
|
|
408e964501 | ||
|
|
11ee185999 | ||
|
|
d5f4b2aa96 | ||
|
|
5a77be5e8b | ||
|
|
bdb8a9c7a4 | ||
|
|
38031739e5 | ||
|
|
f09cfe8ffc | ||
|
|
837212b969 | ||
|
|
70918557ca | ||
|
|
1e44bda48f | ||
|
|
644815f9e5 | ||
|
|
7d6b5790ba | ||
|
|
d56666f7fb | ||
|
|
6a2822d3c7 | ||
|
|
fde3068970 | ||
|
|
f55e77274d | ||
|
|
a90ab44358 |
@@ -31,7 +31,7 @@ jobs:
|
||||
sdk_py: ${{ steps.filter.outputs.sdk_py || 'true' }}
|
||||
steps:
|
||||
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
|
||||
- uses: dorny/paths-filter@7b450fff21473bca461d4b92ce414b9d0420d706 # v4
|
||||
- uses: dorny/paths-filter@ceb8a2b8f2d89434be7ff52d3de7ec3738c5cc9d # v4
|
||||
if: github.event_name != 'workflow_dispatch'
|
||||
id: filter
|
||||
with:
|
||||
|
||||
@@ -76,7 +76,12 @@ __pypackages__/
|
||||
|
||||
# Environments
|
||||
.env
|
||||
.env.*
|
||||
.envrc
|
||||
*.crt
|
||||
*.key
|
||||
*.pem
|
||||
credentials.json
|
||||
.venv
|
||||
.venvs
|
||||
env/
|
||||
@@ -98,6 +103,7 @@ dmypy.json
|
||||
|
||||
.vercel
|
||||
.turbo
|
||||
node_modules/
|
||||
.editorconfig
|
||||
.scratch
|
||||
.worktrees/
|
||||
|
||||
@@ -2,6 +2,14 @@
|
||||
|
||||
This repository is a monorepo. Each library lives in a subdirectory under `libs/`.
|
||||
|
||||
<corridor>
|
||||
|
||||
## Corridor security analysis
|
||||
|
||||
When Corridor's `analyzePlan` tool is available, create a plan and use the tool to analyze it before generating or modifying code. Apply the resulting security guidance before writing code.
|
||||
|
||||
</corridor>
|
||||
|
||||
When you modify code in any library, run the following commands in that library's directory before creating a pull request:
|
||||
|
||||
- `make format` – run code formatters
|
||||
|
||||
@@ -2,6 +2,14 @@
|
||||
|
||||
This repository is a monorepo. Each library lives in a subdirectory under `libs/`.
|
||||
|
||||
<corridor>
|
||||
|
||||
## Corridor security analysis
|
||||
|
||||
Before generating or modifying code, create a plan and use Corridor's `analyzePlan` tool to analyze it. Apply the resulting security guidance before writing code.
|
||||
|
||||
</corridor>
|
||||
|
||||
When you modify code in any library, run the following commands in that library's directory before creating a pull request:
|
||||
|
||||
- `make format` – run code formatters
|
||||
|
||||
+73
-13
@@ -6,9 +6,12 @@ import traceback
|
||||
from collections.abc import Callable
|
||||
from uuid import uuid4
|
||||
|
||||
from langgraph.checkpoint.base import BaseCheckpointSaver
|
||||
from langgraph.checkpoint.base import BaseCheckpointSaver, Checkpoint
|
||||
from langgraph.checkpoint.base.id import uuid6
|
||||
from langgraph.checkpoint.serde.types import _DeltaSnapshot
|
||||
|
||||
from langgraph.checkpoint.conformance.spec._delta_fixtures import build_delta_chain
|
||||
from langgraph.checkpoint.conformance.test_utils import generate_metadata
|
||||
|
||||
|
||||
async def test_history_returns_writes_oldest_first(
|
||||
@@ -48,8 +51,6 @@ async def test_history_seed_is_nearest_snapshot(
|
||||
result = await saver.aget_delta_channel_history(config=head, channels=["ch"])
|
||||
assert "seed" in result["ch"], "Expected seed from snapshot at step 3"
|
||||
seed = result["ch"]["seed"]
|
||||
from langgraph.checkpoint.serde.types import _DeltaSnapshot
|
||||
|
||||
actual_value = seed.value if isinstance(seed, _DeltaSnapshot) else seed
|
||||
assert actual_value == 3, f"Expected seed value 3 (step 3), got {actual_value}"
|
||||
writes = result["ch"]["writes"]
|
||||
@@ -81,11 +82,6 @@ async def test_history_multi_channel(
|
||||
tid = str(uuid4())
|
||||
configs: list = []
|
||||
parent_cfg = None
|
||||
from langgraph.checkpoint.base import Checkpoint
|
||||
from langgraph.checkpoint.base.id import uuid6
|
||||
from langgraph.checkpoint.serde.types import _DeltaSnapshot
|
||||
|
||||
from langgraph.checkpoint.conformance.test_utils import generate_metadata
|
||||
|
||||
for step in range(5):
|
||||
config = {"configurable": {"thread_id": tid, "checkpoint_ns": ""}}
|
||||
@@ -161,11 +157,6 @@ async def test_history_migration_plain_value_as_seed(
|
||||
channel_values[ch] (not a _DeltaSnapshot). The walk should treat it as the
|
||||
seed and terminate there.
|
||||
"""
|
||||
from langgraph.checkpoint.base import Checkpoint
|
||||
from langgraph.checkpoint.base.id import uuid6
|
||||
|
||||
from langgraph.checkpoint.conformance.test_utils import generate_metadata
|
||||
|
||||
tid = str(uuid4())
|
||||
configs: list = []
|
||||
parent_cfg = None
|
||||
@@ -208,6 +199,74 @@ async def test_history_migration_plain_value_as_seed(
|
||||
assert values == [2], f"Expected [2], got {values}"
|
||||
|
||||
|
||||
async def test_history_seed_ancestor_own_writes_are_replayed(
|
||||
saver: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Writes stored AT the seed ancestor must be included in `writes`.
|
||||
|
||||
A stored value is the state ENTERING its checkpoint; the writes stored
|
||||
under that same checkpoint are what produced its child and are therefore
|
||||
NOT subsumed by it. Only writes at ancestors OLDER than the seed are
|
||||
subsumed, and the walk terminates before reaching them.
|
||||
|
||||
This holds for plain-value seeds (migration from a pre-delta channel type)
|
||||
exactly as it does for `_DeltaSnapshot` seeds. Skipping the seed
|
||||
ancestor's own writes silently drops the first post-migration write.
|
||||
"""
|
||||
tid = str(uuid4())
|
||||
configs: list = []
|
||||
parent_cfg = None
|
||||
|
||||
# Each step's write is labelled by the role it plays, so the assertion
|
||||
# below reads directly rather than by step index.
|
||||
writes_by_step = {
|
||||
0: "older-than-seed", # subsumed by the value stored at step 1
|
||||
1: "at-seed", # the seed's own write, produced step 2
|
||||
2: "after-seed", # delta-era write on the path to the head
|
||||
3: "pending-at-head", # pending for the next step, never replayed
|
||||
}
|
||||
# Steps 0 and 1 store a plain value; 1 is the nearest, so it is the seed.
|
||||
values_by_step = {0: [10], 1: [10, 20]}
|
||||
|
||||
for step in range(4):
|
||||
config = {"configurable": {"thread_id": tid, "checkpoint_ns": ""}}
|
||||
if parent_cfg:
|
||||
config["configurable"]["checkpoint_id"] = parent_cfg["configurable"][
|
||||
"checkpoint_id"
|
||||
]
|
||||
cv: dict = {}
|
||||
cvs: dict = {}
|
||||
if step in values_by_step:
|
||||
cv["ch"] = values_by_step[step]
|
||||
cvs["ch"] = step + 1
|
||||
cp = Checkpoint(
|
||||
v=1,
|
||||
id=str(uuid6(clock_seq=-1)),
|
||||
ts="",
|
||||
channel_values=cv,
|
||||
channel_versions=cvs,
|
||||
versions_seen={},
|
||||
updated_channels=None,
|
||||
)
|
||||
parent_cfg = await saver.aput(config, cp, generate_metadata(step=step), cvs)
|
||||
configs.append(parent_cfg)
|
||||
await saver.aput_writes(
|
||||
parent_cfg, [("ch", writes_by_step[step])], str(uuid4())
|
||||
)
|
||||
|
||||
head = configs[-1]
|
||||
result = await saver.aget_delta_channel_history(config=head, channels=["ch"])
|
||||
|
||||
assert "seed" in result["ch"], "Expected seed from plain value at step 1"
|
||||
assert result["ch"]["seed"] == [10, 20], (
|
||||
f"Expected nearest plain value [10, 20], got {result['ch']['seed']}"
|
||||
)
|
||||
values = [w[2] for w in result["ch"]["writes"]]
|
||||
assert values == ["at-seed", "after-seed"], (
|
||||
f'Expected ["at-seed", "after-seed"], got {values}'
|
||||
)
|
||||
|
||||
|
||||
ALL_DELTA_CHANNEL_HISTORY_TESTS = [
|
||||
test_history_returns_writes_oldest_first,
|
||||
test_history_seed_is_nearest_snapshot,
|
||||
@@ -216,6 +275,7 @@ ALL_DELTA_CHANNEL_HISTORY_TESTS = [
|
||||
test_history_empty_channels_returns_empty,
|
||||
test_history_walk_to_root_no_seed,
|
||||
test_history_migration_plain_value_as_seed,
|
||||
test_history_seed_ancestor_own_writes_are_replayed,
|
||||
]
|
||||
|
||||
|
||||
|
||||
Generated
+1048
-601
File diff suppressed because it is too large
Load Diff
@@ -7,7 +7,7 @@ import logging
|
||||
import re
|
||||
import threading
|
||||
from collections import defaultdict
|
||||
from collections.abc import Callable, Iterable, Iterator, Sequence
|
||||
from collections.abc import Callable, Iterable, Iterator, Mapping, Sequence
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime
|
||||
from typing import (
|
||||
@@ -354,7 +354,7 @@ class BasePostgresStore(Generic[C]):
|
||||
(
|
||||
_namespace_to_text(op.namespace),
|
||||
op.key,
|
||||
Jsonb(cast(dict, op.value)),
|
||||
Jsonb(dict(cast(Mapping[str, Any], op.value))),
|
||||
)
|
||||
)
|
||||
if op.ttl is not None:
|
||||
|
||||
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "langgraph-checkpoint-postgres"
|
||||
version = "3.1.1"
|
||||
version = "3.1.2"
|
||||
description = "Library with a Postgres implementation of LangGraph checkpoint saver."
|
||||
authors = []
|
||||
requires-python = ">=3.10"
|
||||
|
||||
Generated
+1161
-665
File diff suppressed because it is too large
Load Diff
@@ -7,7 +7,7 @@ import re
|
||||
import sqlite3
|
||||
import threading
|
||||
from collections import defaultdict
|
||||
from collections.abc import Callable, Iterable, Iterator, Sequence
|
||||
from collections.abc import Callable, Iterable, Iterator, Mapping, Sequence
|
||||
from contextlib import contextmanager
|
||||
from typing import Any, Literal, NamedTuple, cast
|
||||
|
||||
@@ -387,7 +387,7 @@ class BaseSqliteStore:
|
||||
[
|
||||
_namespace_to_text(op.namespace),
|
||||
op.key,
|
||||
orjson.dumps(cast(dict, op.value)),
|
||||
orjson.dumps(dict(cast(Mapping[str, Any], op.value))),
|
||||
expires_at,
|
||||
op.ttl,
|
||||
]
|
||||
|
||||
Generated
+1098
-602
File diff suppressed because it is too large
Load Diff
+11
-18
@@ -39,7 +39,12 @@ You must pass these when invoking the graph as part of the configurable part of
|
||||
|
||||
```python
|
||||
{"configurable": {"thread_id": "1"}} # valid config
|
||||
{"configurable": {"thread_id": "1", "checkpoint_id": "0c62ca34-ac19-445d-bbb0-5b4984975b2a"}} # also valid config
|
||||
{
|
||||
"configurable": {
|
||||
"thread_id": "1",
|
||||
"checkpoint_id": "0c62ca34-ac19-445d-bbb0-5b4984975b2a",
|
||||
}
|
||||
} # also valid config
|
||||
```
|
||||
|
||||
### Serde
|
||||
@@ -79,24 +84,12 @@ checkpoint = {
|
||||
"v": 4,
|
||||
"ts": "2024-07-31T20:14:19.804150+00:00",
|
||||
"id": "1ef4f797-8335-6428-8001-8a1503f9b875",
|
||||
"channel_values": {
|
||||
"my_key": "meow",
|
||||
"node": "node"
|
||||
},
|
||||
"channel_versions": {
|
||||
"__start__": 2,
|
||||
"my_key": 3,
|
||||
"start:node": 3,
|
||||
"node": 3
|
||||
},
|
||||
"channel_values": {"my_key": "meow", "node": "node"},
|
||||
"channel_versions": {"__start__": 2, "my_key": 3, "start:node": 3, "node": 3},
|
||||
"versions_seen": {
|
||||
"__input__": {},
|
||||
"__start__": {
|
||||
"__start__": 1
|
||||
},
|
||||
"node": {
|
||||
"start:node": 2
|
||||
}
|
||||
"__input__": {},
|
||||
"__start__": {"__start__": 1},
|
||||
"node": {"start:node": 2},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -10,6 +10,7 @@ from typing import (
|
||||
NamedTuple,
|
||||
TypedDict,
|
||||
TypeVar,
|
||||
cast,
|
||||
)
|
||||
|
||||
from langchain_core.runnables import RunnableConfig
|
||||
@@ -781,7 +782,7 @@ def get_serializable_checkpoint_metadata(
|
||||
"""Get checkpoint metadata in a backwards-compatible manner."""
|
||||
checkpoint_metadata = get_checkpoint_metadata(config, metadata)
|
||||
if "writes" in checkpoint_metadata:
|
||||
checkpoint_metadata.pop("writes")
|
||||
cast(dict[str, Any], checkpoint_metadata).pop("writes")
|
||||
return checkpoint_metadata
|
||||
|
||||
|
||||
|
||||
@@ -148,17 +148,16 @@ class InMemorySaver(
|
||||
whose stored blob is non-empty. Other channels keep walking until
|
||||
they find their own terminator or hit the root.
|
||||
|
||||
Pre-delta plain-value blobs subsume their ancestor's pending
|
||||
writes (the value already includes them); `_DeltaSnapshot` blobs
|
||||
do not (snapshot is the value AT that ancestor, prior to its own
|
||||
pending writes that produce the child).
|
||||
A blob is the value AT its ancestor, prior to the writes stored
|
||||
under that same ancestor (those writes produce its child, which
|
||||
is on the path to the target). This holds for `_DeltaSnapshot`
|
||||
blobs and for pre-delta plain values alike, so the seed
|
||||
ancestor's own writes are always collected. Writes at ancestors
|
||||
older than the seed are subsumed by the seed value and are never
|
||||
reached — the walk terminates there.
|
||||
"""
|
||||
if not channels:
|
||||
return {}
|
||||
# Imported lazily to avoid a hard checkpoint→serde-types coupling at
|
||||
# module import; only this override needs the runtime check.
|
||||
from langgraph.checkpoint.serde.types import _DeltaSnapshot
|
||||
|
||||
thread_id = config["configurable"]["thread_id"]
|
||||
checkpoint_ns = config["configurable"].get("checkpoint_ns", "")
|
||||
checkpoint_id = config["configurable"].get("checkpoint_id", "")
|
||||
@@ -205,11 +204,6 @@ class InMemorySaver(
|
||||
):
|
||||
if ch not in remaining:
|
||||
continue
|
||||
blob_value = blob_value_by_ch.get(ch)
|
||||
if blob_value is not None and not isinstance(
|
||||
blob_value, _DeltaSnapshot
|
||||
):
|
||||
continue
|
||||
collected_by_ch[ch].append(
|
||||
(tid, ch, self.serde.loads_typed(serialized))
|
||||
)
|
||||
|
||||
@@ -12,7 +12,7 @@ Core types:
|
||||
from __future__ import annotations
|
||||
|
||||
from abc import ABC, abstractmethod
|
||||
from collections.abc import Iterable
|
||||
from collections.abc import Iterable, Mapping
|
||||
from datetime import datetime
|
||||
from typing import (
|
||||
Any,
|
||||
@@ -473,10 +473,10 @@ class PutOp(NamedTuple):
|
||||
the full path would effectively be `"documents/user123/report1"`
|
||||
"""
|
||||
|
||||
value: dict[str, Any] | None
|
||||
value: Mapping[str, Any] | None
|
||||
"""The data to store, or `None` to mark the item for deletion.
|
||||
|
||||
The value must be a dictionary with string keys and JSON-serializable values.
|
||||
The value must be a mapping with string keys and JSON-serializable values.
|
||||
Setting this to `None` signals that the item should be deleted.
|
||||
|
||||
Example:
|
||||
@@ -857,7 +857,7 @@ class BaseStore(ABC):
|
||||
self,
|
||||
namespace: tuple[str, ...],
|
||||
key: str,
|
||||
value: dict[str, Any],
|
||||
value: Mapping[str, Any],
|
||||
index: Literal[False] | list[str] | None = None,
|
||||
*,
|
||||
ttl: float | None | NotProvided = NOT_PROVIDED,
|
||||
@@ -869,7 +869,7 @@ class BaseStore(ABC):
|
||||
Example: `("documents", "user123")`
|
||||
key: Unique identifier within the namespace. Together with namespace forms
|
||||
the complete path to the item.
|
||||
value: Dictionary containing the item's data. Must contain string keys
|
||||
value: Mapping containing the item's data. Must contain string keys
|
||||
and JSON-serializable values.
|
||||
index: Controls how the item's fields are indexed for search:
|
||||
|
||||
@@ -1110,7 +1110,7 @@ class BaseStore(ABC):
|
||||
self,
|
||||
namespace: tuple[str, ...],
|
||||
key: str,
|
||||
value: dict[str, Any],
|
||||
value: Mapping[str, Any],
|
||||
index: Literal[False] | list[str] | None = None,
|
||||
*,
|
||||
ttl: float | None | NotProvided = NOT_PROVIDED,
|
||||
@@ -1122,7 +1122,7 @@ class BaseStore(ABC):
|
||||
Example: `("documents", "user123")`
|
||||
key: Unique identifier within the namespace. Together with namespace forms
|
||||
the complete path to the item.
|
||||
value: Dictionary containing the item's data. Must contain string keys
|
||||
value: Mapping containing the item's data. Must contain string keys
|
||||
and JSON-serializable values.
|
||||
index: Controls how the item's fields are indexed for search:
|
||||
|
||||
|
||||
@@ -5,7 +5,7 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import functools
|
||||
import weakref
|
||||
from collections.abc import Callable, Iterable
|
||||
from collections.abc import Callable, Iterable, Mapping
|
||||
from typing import Any, Literal, TypeVar
|
||||
|
||||
from langgraph.store.base import (
|
||||
@@ -132,7 +132,7 @@ class AsyncBatchedBaseStore(BaseStore):
|
||||
self,
|
||||
namespace: tuple[str, ...],
|
||||
key: str,
|
||||
value: dict[str, Any],
|
||||
value: Mapping[str, Any],
|
||||
index: Literal[False] | list[str] | None = None,
|
||||
*,
|
||||
ttl: float | None | NotProvided = NOT_PROVIDED,
|
||||
@@ -231,7 +231,7 @@ class AsyncBatchedBaseStore(BaseStore):
|
||||
self,
|
||||
namespace: tuple[str, ...],
|
||||
key: str,
|
||||
value: dict[str, Any],
|
||||
value: Mapping[str, Any],
|
||||
index: Literal[False] | list[str] | None = None,
|
||||
*,
|
||||
ttl: float | None | NotProvided = NOT_PROVIDED,
|
||||
|
||||
@@ -11,7 +11,7 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import functools
|
||||
import json
|
||||
from collections.abc import Awaitable, Callable, Sequence
|
||||
from collections.abc import Awaitable, Callable, Mapping, Sequence
|
||||
from typing import Any
|
||||
|
||||
from langchain_core.embeddings import Embeddings
|
||||
@@ -244,6 +244,9 @@ def get_text_at_path(obj: Any, path: str | list[str]) -> list[str]:
|
||||
- Multi-field selection: "{field1,field2}"
|
||||
- Nested paths in multi-field: "{field1,nested.field2}"
|
||||
"""
|
||||
if isinstance(obj, Mapping) and not isinstance(obj, dict):
|
||||
obj = dict(obj)
|
||||
|
||||
if not path or path == "$":
|
||||
return [json.dumps(obj, sort_keys=True, ensure_ascii=False)]
|
||||
|
||||
|
||||
@@ -408,7 +408,7 @@ class InMemoryStore(BaseStore):
|
||||
self._vectors[namespace].pop(key, None)
|
||||
else:
|
||||
self._data[namespace][key] = Item(
|
||||
value=op.value,
|
||||
value=dict(op.value),
|
||||
key=key,
|
||||
namespace=namespace,
|
||||
created_at=datetime.now(timezone.utc),
|
||||
|
||||
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "langgraph-checkpoint"
|
||||
version = "4.1.1"
|
||||
version = "4.2.0"
|
||||
description = "Library with base interfaces for LangGraph checkpoint savers."
|
||||
authors = []
|
||||
requires-python = ">=3.10"
|
||||
|
||||
@@ -577,9 +577,19 @@ class TestPreDeltaBlobTerminator:
|
||||
"""
|
||||
|
||||
def _build_mixed_thread(self) -> tuple[InMemorySaver, str, str, str, str]:
|
||||
"""Three-checkpoint chain: cp1 (pre-delta, blob=[A]), cp2 (delta,
|
||||
write=B), cp3 (delta, write=C). Reconstructing at cp3 must yield
|
||||
seed=[A] + writes=[B, C].
|
||||
"""Four-checkpoint chain spanning the migration boundary:
|
||||
|
||||
* `cp0` — pre-delta ancestor OLDER than the seed. Its write
|
||||
(`OLDER-WRITE`) is already folded into `cp1`'s stored value, so the
|
||||
walk must terminate at `cp1` and never reach it.
|
||||
* `cp1` — pre-delta, blob `["A"]`. That value is the state ENTERING
|
||||
`cp1`; the write stored under `cp1` (`PRE-DELTA-WRITE`) is what
|
||||
produced `cp2` and is NOT subsumed by the blob.
|
||||
* `cp2` — delta-era, no stored value, write `B`.
|
||||
* `cp3` — target, delta-era, write `PENDING-AT-TARGET`.
|
||||
|
||||
Reconstructing at `cp3` must yield seed `["A"]` plus writes
|
||||
`["PRE-DELTA-WRITE", "B"]`.
|
||||
|
||||
Returns `(saver, thread_id, ns, channel, cp3_id)`.
|
||||
"""
|
||||
@@ -587,16 +597,21 @@ class TestPreDeltaBlobTerminator:
|
||||
serde = JsonPlusSerializer()
|
||||
thread_id, ns, channel = "t1", "", "messages"
|
||||
|
||||
v0 = "00000000000000000000000000000000.0"
|
||||
v1 = "00000000000000000000000000000001.0"
|
||||
v2 = "00000000000000000000000000000002.0"
|
||||
v3 = "00000000000000000000000000000003.0"
|
||||
|
||||
# Pre-delta: cp1 stored a real blob for the channel.
|
||||
# Pre-delta: cp0 and cp1 stored real blobs for the channel.
|
||||
saver.blobs[(thread_id, ns, channel, v0)] = serde.dumps_typed([])
|
||||
saver.blobs[(thread_id, ns, channel, v1)] = serde.dumps_typed(["A"])
|
||||
# Delta-era: cp2 and cp3 store "empty"; real writes in checkpoint_writes.
|
||||
saver.blobs[(thread_id, ns, channel, v2)] = ("empty", b"")
|
||||
saver.blobs[(thread_id, ns, channel, v3)] = ("empty", b"")
|
||||
|
||||
cp0 = empty_checkpoint()
|
||||
cp0["id"] = "cp0"
|
||||
cp0["channel_versions"][channel] = v0
|
||||
cp1 = empty_checkpoint()
|
||||
cp1["id"] = "cp1"
|
||||
cp1["channel_versions"][channel] = v1
|
||||
@@ -608,16 +623,24 @@ class TestPreDeltaBlobTerminator:
|
||||
cp3["channel_versions"][channel] = v3
|
||||
|
||||
saver.storage[thread_id][ns] = {
|
||||
"cp1": (serde.dumps_typed(cp1), serde.dumps_typed({}), None),
|
||||
"cp0": (serde.dumps_typed(cp0), serde.dumps_typed({}), None),
|
||||
"cp1": (serde.dumps_typed(cp1), serde.dumps_typed({}), "cp0"),
|
||||
"cp2": (serde.dumps_typed(cp2), serde.dumps_typed({}), "cp1"),
|
||||
"cp3": (serde.dumps_typed(cp3), serde.dumps_typed({}), "cp2"),
|
||||
}
|
||||
# Write under cp1 would be from the pre-delta era and MUST be ignored
|
||||
# (the blob already captures it). We add one and assert it is not
|
||||
# folded into the reconstructed result.
|
||||
saver.writes[(thread_id, ns, "cp1")][("task0", 0)] = (
|
||||
# Write under cp0 is older than the seed — cp1's blob already folded
|
||||
# it in, and the terminator must stop before reaching it.
|
||||
saver.writes[(thread_id, ns, "cp0")][("task0", 0)] = (
|
||||
"task0",
|
||||
channel,
|
||||
serde.dumps_typed("OLDER-WRITE"),
|
||||
"",
|
||||
)
|
||||
# Write under cp1 postdates cp1's blob (it is what produced cp2, which
|
||||
# stores no value of its own) and MUST be replayed.
|
||||
saver.writes[(thread_id, ns, "cp1")][("task1", 0)] = (
|
||||
"task1",
|
||||
channel,
|
||||
serde.dumps_typed("PRE-DELTA-WRITE"),
|
||||
"",
|
||||
)
|
||||
@@ -651,15 +674,19 @@ class TestPreDeltaBlobTerminator:
|
||||
|
||||
# Seed came from the pre-delta blob at cp1.
|
||||
assert result["seed"] == ["A"]
|
||||
# Delta-era writes from cp2 replay through the reducer on top of seed.
|
||||
# cp3 is the target — its own write is pending for the NEXT step and
|
||||
# must be excluded.
|
||||
# The seed ancestor's own write and the delta-era write from cp2 both
|
||||
# replay through the reducer on top of the seed, oldest first. cp3 is
|
||||
# the target — its own write is pending for the NEXT step and must be
|
||||
# excluded.
|
||||
values = [v for _, _, v in result["writes"]]
|
||||
assert values == ["B"]
|
||||
assert values == ["PRE-DELTA-WRITE", "B"]
|
||||
|
||||
def test_pre_delta_blob_terminates_walk_before_older_writes(self) -> None:
|
||||
"""Writes stored at the pre-delta ancestor itself must not be replayed
|
||||
(the blob subsumes them)."""
|
||||
def test_seed_bounds_walk_without_dropping_its_own_writes(self) -> None:
|
||||
"""The seed terminator bounds the walk: writes at ancestors OLDER than
|
||||
the seed are already folded into the seed value and must not be
|
||||
replayed. The seed ancestor's own write is not one of them — it
|
||||
postdates the stored value and produced the next checkpoint.
|
||||
"""
|
||||
saver, thread_id, ns, channel, target = self._build_mixed_thread()
|
||||
config: RunnableConfig = {
|
||||
"configurable": {
|
||||
@@ -674,7 +701,9 @@ class TestPreDeltaBlobTerminator:
|
||||
]
|
||||
|
||||
values = [v for _, _, v in result["writes"]]
|
||||
# The pre-delta write under cp1 must not appear (the blob subsumes it).
|
||||
assert "PRE-DELTA-WRITE" not in values
|
||||
# Older than the seed — subsumed by cp1's blob, so the walk stops first.
|
||||
assert "OLDER-WRITE" not in values
|
||||
# Stored AT the seed ancestor — not subsumed, so it must be replayed.
|
||||
assert "PRE-DELTA-WRITE" in values
|
||||
# And the pending write at the target is never folded in.
|
||||
assert "PENDING-AT-TARGET" not in values
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
import asyncio
|
||||
import json
|
||||
from collections.abc import Iterable
|
||||
from collections import UserDict
|
||||
from collections.abc import Iterable, Mapping
|
||||
from datetime import datetime
|
||||
from types import MappingProxyType
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
@@ -137,6 +139,18 @@ def test_get_text_at_path() -> None:
|
||||
assert get_text_at_path(nested_data, "nested[{invalid}]") == []
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"mapping",
|
||||
[
|
||||
UserDict({"text": "searchable"}),
|
||||
MappingProxyType({"text": "searchable"}),
|
||||
],
|
||||
)
|
||||
def test_get_text_at_path_with_non_dict_mapping(mapping: Mapping[str, str]) -> None:
|
||||
assert get_text_at_path(mapping, "$") == ['{"text": "searchable"}']
|
||||
assert get_text_at_path(mapping, "text") == ["searchable"]
|
||||
|
||||
|
||||
async def test_async_batch_store(mocker: MockerFixture) -> None:
|
||||
abatch = mocker.stub()
|
||||
|
||||
|
||||
Generated
+1392
-703
File diff suppressed because it is too large
Load Diff
@@ -21,8 +21,8 @@
|
||||
"test:all": "yarn test && yarn test:int && yarn lint:langgraph"
|
||||
},
|
||||
"dependencies": {
|
||||
"@langchain/core": "^1.2.4",
|
||||
"@langchain/langgraph": "^1.4.8"
|
||||
"@langchain/core": "^1.2.9",
|
||||
"@langchain/langgraph": "^1.4.13"
|
||||
},
|
||||
"resolutions": {
|
||||
"@langchain/langgraph-checkpoint": "1.0.4"
|
||||
@@ -32,15 +32,15 @@
|
||||
"@eslint/js": "^10.0.1",
|
||||
"@tsconfig/recommended": "^1.0.13",
|
||||
"@types/jest": "^30.0.0",
|
||||
"@typescript-eslint/eslint-plugin": "^8.65.0",
|
||||
"@typescript-eslint/parser": "^8.65.0",
|
||||
"@typescript-eslint/eslint-plugin": "^8.68.0",
|
||||
"@typescript-eslint/parser": "^8.68.0",
|
||||
"dotenv": "^17.4.2",
|
||||
"eslint": "^10.8.0",
|
||||
"eslint": "^10.9.1",
|
||||
"eslint-config-prettier": "^10.1.8",
|
||||
"eslint-plugin-import": "^2.32.0",
|
||||
"eslint-plugin-no-instanceof": "^1.0.1",
|
||||
"eslint-plugin-prettier": "^5.5.6",
|
||||
"jest": "^30.4.2",
|
||||
"jest": "^30.5.0",
|
||||
"prettier": "^3.9.6",
|
||||
"ts-jest": "^29.4.12",
|
||||
"typescript": "^7.0.2"
|
||||
|
||||
+830
-752
File diff suppressed because it is too large
Load Diff
@@ -9,8 +9,8 @@
|
||||
},
|
||||
"dependencies": {
|
||||
"@js-monorepo-example/shared": "*",
|
||||
"@langchain/core": "^1.2.4",
|
||||
"@langchain/langgraph": "^1.4.8"
|
||||
"@langchain/core": "^1.2.9",
|
||||
"@langchain/langgraph": "^1.4.13"
|
||||
},
|
||||
"devDependencies": {
|
||||
"typescript": "^7.0.2"
|
||||
|
||||
@@ -17,18 +17,18 @@
|
||||
"lint": "eslint 'apps/**/*.ts' 'libs/**/*.ts'"
|
||||
},
|
||||
"devDependencies": {
|
||||
"turbo": "^2.10.8",
|
||||
"turbo": "^2.10.12",
|
||||
"typescript": "^7.0.2",
|
||||
"@tsconfig/recommended": "^1.0.13",
|
||||
"@eslint/eslintrc": "^3.3.6",
|
||||
"@eslint/js": "^10.0.1",
|
||||
"eslint": "^10.8.0",
|
||||
"eslint": "^10.9.1",
|
||||
"eslint-config-prettier": "^10.1.8",
|
||||
"eslint-plugin-import": "^2.27.5",
|
||||
"eslint-plugin-no-instanceof": "^1.0.1",
|
||||
"eslint-plugin-prettier": "^5.5.6",
|
||||
"@typescript-eslint/eslint-plugin": "^8.65.0",
|
||||
"@typescript-eslint/parser": "^8.65.0",
|
||||
"@typescript-eslint/eslint-plugin": "^8.68.0",
|
||||
"@typescript-eslint/parser": "^8.68.0",
|
||||
"prettier": "^3.9.6"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -103,10 +103,10 @@
|
||||
resolved "https://registry.yarnpkg.com/@isaacs/cliui/-/cliui-9.0.0.tgz#4d0a3f127058043bf2e7ee169eaf30ed901302f3"
|
||||
integrity sha512-AokJm4tuBHillT+FpMtxQ60n8ObyXBatq7jD2/JA9dxbDDokKQm8KMht5ibGzLVU9IJDIKK4TPKgMHEYMn3lMg==
|
||||
|
||||
"@langchain/core@^1.2.4":
|
||||
version "1.2.4"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/core/-/core-1.2.4.tgz#3916181e00f93b2e3b5801ab8273c07d144d75ce"
|
||||
integrity sha512-GIrJktdsFPx8gM0C3VyikkeYGZAV7iRIzUJuS5tFatiKLExybou0LI0MuxH9kksN38quyfWdQDl20lIEAEVkVA==
|
||||
"@langchain/core@^1.2.9":
|
||||
version "1.2.9"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/core/-/core-1.2.9.tgz#4f5fb27ba07c51ce4fb8e1dc32eb36945737afa1"
|
||||
integrity sha512-conzSEj9Zu1AyXJLXsSbgrtxtxinmI1yGqQ5CIJZSoV5rvv+yvQE/vgBnoySpBQ/bl3YPgj2FL/gbDjWykLSfg==
|
||||
dependencies:
|
||||
"@cfworker/json-schema" "^4.0.2"
|
||||
"@standard-schema/spec" "^1.1.0"
|
||||
@@ -116,28 +116,28 @@
|
||||
p-queue "^6.6.2"
|
||||
zod "^3.25.76 || ^4"
|
||||
|
||||
"@langchain/langgraph-checkpoint@^1.1.3":
|
||||
version "1.1.3"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/langgraph-checkpoint/-/langgraph-checkpoint-1.1.3.tgz#891f2f9a2e96cf2ab0920f9240f9ad1e0a86a4ab"
|
||||
integrity sha512-wgzdQNeEsdw1e+4lvlj0tdq/RYR/k1vPin10g0ymGoehZDDgd9nvIllGXSXN4TFgF9sf5qQP/KTkOcLfeseIhA==
|
||||
"@langchain/langgraph-checkpoint@^1.1.5":
|
||||
version "1.1.5"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/langgraph-checkpoint/-/langgraph-checkpoint-1.1.5.tgz#c793177f9afcce31a1e9922f317f88fec1254282"
|
||||
integrity sha512-BwDwl5VeTOh6CVuiIPgsUgfK51vTJDMSbFcSCUfjJWsl8/DPdK/mbv+ejxJstkSk/BlSPMP4JfXWcN6jD2ea2Q==
|
||||
|
||||
"@langchain/langgraph-sdk@~1.9.26":
|
||||
version "1.9.28"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/langgraph-sdk/-/langgraph-sdk-1.9.28.tgz#1664cc7b6f2c1509ddffecf17d805a26e70a7e1a"
|
||||
integrity sha512-4j3XuM0PvtmAbL8mPfBS99ez3+ytRfgbOpAR/nOeaejTRF3Q9dNw2QnaGLGng8wLPtGLoSj+SYgUOVxy9Bv9vg==
|
||||
"@langchain/langgraph-sdk@~1.10.0":
|
||||
version "1.10.0"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/langgraph-sdk/-/langgraph-sdk-1.10.0.tgz#6ec1364a97caf615983161ae1f00b196897b9cd6"
|
||||
integrity sha512-cPPkh+hMNgeOaGtJRrqs1AjZde45cG2+Ma9Sc10wz2RyvT8SKToCKS+VvkS18SsLajnmq6/FKVmthq6rnUVYOw==
|
||||
dependencies:
|
||||
"@langchain/protocol" "^0.0.18"
|
||||
"@langchain/protocol" "^0.0.19"
|
||||
"@types/json-schema" "^7.0.15"
|
||||
p-queue "^9.0.1"
|
||||
p-retry "^7.1.1"
|
||||
|
||||
"@langchain/langgraph@^1.4.8":
|
||||
version "1.4.8"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/langgraph/-/langgraph-1.4.8.tgz#65ef5635b5554cd32f502a3808b92d8371998a99"
|
||||
integrity sha512-DN1Np1XefdBEbp1qBKlt39cwoL743AAGpR5Ipja0gY2YbWvsoQnOTIrjnj/orSAhaUYsdTKS8VSWdFzsHZo6Ig==
|
||||
"@langchain/langgraph@^1.4.13":
|
||||
version "1.4.13"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/langgraph/-/langgraph-1.4.13.tgz#168ab4f05c212fd5ab4b2816bcad2e963b345101"
|
||||
integrity sha512-LO1ak6jNQ9jR13tm7Ay4Yh2/otrH7LNVUwWTAI7WJigVdW5Fb6LuYSZUzVn4S7sVSiyVFfnrUcDTd8c7eAzPrQ==
|
||||
dependencies:
|
||||
"@langchain/langgraph-checkpoint" "^1.1.3"
|
||||
"@langchain/langgraph-sdk" "~1.9.26"
|
||||
"@langchain/langgraph-checkpoint" "^1.1.5"
|
||||
"@langchain/langgraph-sdk" "~1.10.0"
|
||||
"@langchain/protocol" "^0.0.18"
|
||||
"@standard-schema/spec" "1.1.0"
|
||||
|
||||
@@ -146,6 +146,11 @@
|
||||
resolved "https://registry.yarnpkg.com/@langchain/protocol/-/protocol-0.0.18.tgz#6d96155e7263c958fbce6d4b7241b4725b72fbad"
|
||||
integrity sha512-XW1egQtPfsGI41w2AMZNFZrUIwFSQHTjVMZs0OaTpCAvht/QLoaPN8FQcsysMVypOhupG28J29yOorrc70otBQ==
|
||||
|
||||
"@langchain/protocol@^0.0.19":
|
||||
version "0.0.19"
|
||||
resolved "https://registry.yarnpkg.com/@langchain/protocol/-/protocol-0.0.19.tgz#7fe43fed115dc34b3b4c246e8d5283fbcbfab7b0"
|
||||
integrity sha512-9hKcRrH7cBX6gfutdfXPoft1OCchHe4FEpALoDJMl5Qu+n/YG5ynZmyu8+8cxORlPwHBoKTxggvXz+76M1yX1Q==
|
||||
|
||||
"@pkgr/core@^0.3.6":
|
||||
version "0.3.6"
|
||||
resolved "https://registry.yarnpkg.com/@pkgr/core/-/core-0.3.6.tgz#3569708bd4be4d8870ba32bf1c456dac81600d97"
|
||||
@@ -166,35 +171,35 @@
|
||||
resolved "https://registry.yarnpkg.com/@tsconfig/recommended/-/recommended-1.0.13.tgz#269fce3ad04ca70b93269ff44cca81b950f542da"
|
||||
integrity sha512-sySRuBfMKyKO/j2ZAhR8kSembhjuPEV4Ra3AHtmWLq51+iGaudr45crPSzNC5b7/Ctrh9dfUpBuTlYrH6rM58Q==
|
||||
|
||||
"@turbo/darwin-64@2.10.8":
|
||||
version "2.10.8"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/darwin-64/-/darwin-64-2.10.8.tgz#0ae7cde6b95b82b0a4fe406ffea64bdfcf68a0ee"
|
||||
integrity sha512-po+7rfJfUnFXjWlcoN2RwhErgzCdRtBc1T26vYPcywHlggmCQiQe1uWaE4j+BibI2uY9/2pDoFzMN0rmSaPFOw==
|
||||
"@turbo/darwin-64@2.10.12":
|
||||
version "2.10.12"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/darwin-64/-/darwin-64-2.10.12.tgz#719ac12ae47caf7b4be855d672582cdbfde177ea"
|
||||
integrity sha512-9nKgKoF6ZOUsM+or0OtNf+TTJSfGvDNP7ZFv/ZGWVwOSCkumyctQiTeHwB4UNljHTnC41AqylgbunLDHoccNrA==
|
||||
|
||||
"@turbo/darwin-arm64@2.10.8":
|
||||
version "2.10.8"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/darwin-arm64/-/darwin-arm64-2.10.8.tgz#18a9127c0a0d9c5ad7de49bdb27e7032f8946334"
|
||||
integrity sha512-+zB2btDJ00lnPRuqOvpVvgl4x34k/djZQGZTTCfjn7JgNCl8QFY5Njo5+dqkY1g/+9gbbsnAvWm9CmJg9ebcXA==
|
||||
"@turbo/darwin-arm64@2.10.12":
|
||||
version "2.10.12"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/darwin-arm64/-/darwin-arm64-2.10.12.tgz#ac2d3dde3a2407f8359ca2e3152690280cff2714"
|
||||
integrity sha512-H4Elb1jqTZVeIC9bbcNwjSzemZ6RegoTOVHeuV5Osirt2Z8UguTyisMEkvZjPVZgMeN9J4ERZBFad40tFnkb7w==
|
||||
|
||||
"@turbo/linux-64@2.10.8":
|
||||
version "2.10.8"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/linux-64/-/linux-64-2.10.8.tgz#1677672f5760d272c0a0608686fe8b8b29491fdb"
|
||||
integrity sha512-K1dxqiVisyN7cViVsfQLs6xscQbYuI8aO2nbUhFURDACgEDfZRdP/b4CCxeosBJpcMfhYyiibWqJorCnvz9kKg==
|
||||
"@turbo/linux-64@2.10.12":
|
||||
version "2.10.12"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/linux-64/-/linux-64-2.10.12.tgz#264a0ec88a1f69f93cf57dab20be1b4c50bf226c"
|
||||
integrity sha512-lr7KIotukvjZwEXiFSYAeOH3BWzjFVBbSzTbv0fuGFsNukYyH0+g1hB5ecqnJkgkYU+KHEMG1edOhnjiKON1wQ==
|
||||
|
||||
"@turbo/linux-arm64@2.10.8":
|
||||
version "2.10.8"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/linux-arm64/-/linux-arm64-2.10.8.tgz#dbdb5aef9e53bb4c88357dafc19f79138ec29253"
|
||||
integrity sha512-Gi77ibVnrE1fEmvr+/wBD/yvRqhwp/RQuCp2+//lv1U1wNFFyVg0V7Wj8FG9FXPFAw5QHReo8rxc9+wBSDZjzA==
|
||||
"@turbo/linux-arm64@2.10.12":
|
||||
version "2.10.12"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/linux-arm64/-/linux-arm64-2.10.12.tgz#3ee53d1f930d2bb65708e904971be0e765381e37"
|
||||
integrity sha512-f0pZDTtvzB5SuNwuXBaKbZHUCMCukgc8nMlHEuvLmj91Fzec+MEbr3cAvGNor5htEDqZnO6Lxt9N/GPI/77oGA==
|
||||
|
||||
"@turbo/windows-64@2.10.8":
|
||||
version "2.10.8"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/windows-64/-/windows-64-2.10.8.tgz#451c3b5421c9ce573c764b1efacae5ad4d1312ee"
|
||||
integrity sha512-znnLO1haJPYTHoKMKwlAvlkjRiYbbhBzME6wIGaMd+fwir23U6jVd1ecaTWWi1fbnRVqxMfgDBKseQ/hLKb83g==
|
||||
"@turbo/windows-64@2.10.12":
|
||||
version "2.10.12"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/windows-64/-/windows-64-2.10.12.tgz#ac0a0b7794f9a541e68e930383f48747f7e648ab"
|
||||
integrity sha512-SDOueJRjS/QcykWf2KCRtTLmIl5YMKsLbXkXQGhDwcTXvKXZiS5ih5lBl/gkwZIpYFjqA/rAlfMzlAFcVHNe0g==
|
||||
|
||||
"@turbo/windows-arm64@2.10.8":
|
||||
version "2.10.8"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/windows-arm64/-/windows-arm64-2.10.8.tgz#5661ca9b75713aece06d319493e48a366a4cd05c"
|
||||
integrity sha512-VN30vh3b3Czh2WzYHNTfF1FE0YMZ5aHsLO8dBMGHJewA6792wX6iJR8ZxlzFW6WdOu0gEAKIvlYhfyT81Wkm4Q==
|
||||
"@turbo/windows-arm64@2.10.12":
|
||||
version "2.10.12"
|
||||
resolved "https://registry.yarnpkg.com/@turbo/windows-arm64/-/windows-arm64-2.10.12.tgz#50eb6c20a20d0aab216d2b6aa13cc3100ea44106"
|
||||
integrity sha512-0i0mVUa4kKk+/B3RwEwPMf9CB+T7ul56hn5FFHNA4VUNTOoLBEd6aNf3FaKfCatDNZ6cicCEf6if9QUTVyzzcA==
|
||||
|
||||
"@types/esrecurse@^4.3.1":
|
||||
version "4.3.1"
|
||||
@@ -216,110 +221,110 @@
|
||||
resolved "https://registry.yarnpkg.com/@types/json5/-/json5-0.0.29.tgz#ee28707ae94e11d2b827bcbe5270bcea7f3e71ee"
|
||||
integrity sha512-dRLjCWHYg4oaA77cxO64oO+7JwCwnIzkZPdrrC71jQmQtlhM556pwKo5bUzqvZndkVbeFLIIi+9TC40JNF5hNQ==
|
||||
|
||||
"@typescript-eslint/eslint-plugin@^8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.65.0.tgz#0a58df6fea8c0bf6b396f518077099bc8b762bb5"
|
||||
integrity sha512-IEgob78X12rHpUmtcwFsXhZdVGJtwTVP8FiCLZkR6GlYVrl2PcuB+KhCE5BlVC/eQpQnu8WXRtkHZuPar+gCRA==
|
||||
"@typescript-eslint/eslint-plugin@^8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.68.0.tgz#a8fbdb1cf49aafaf16071b646daad890151bd149"
|
||||
integrity sha512-WASHDpCm6qO5jj9g1a+8NiW5+GCkAyLReR56/4VruYmNgfUmqpxOfZ2Yfb8xGfJPWv5Qi6LSD8sXdces3vbp/Q==
|
||||
dependencies:
|
||||
"@eslint-community/regexpp" "^4.12.2"
|
||||
"@typescript-eslint/scope-manager" "8.65.0"
|
||||
"@typescript-eslint/type-utils" "8.65.0"
|
||||
"@typescript-eslint/utils" "8.65.0"
|
||||
"@typescript-eslint/visitor-keys" "8.65.0"
|
||||
"@typescript-eslint/scope-manager" "8.68.0"
|
||||
"@typescript-eslint/type-utils" "8.68.0"
|
||||
"@typescript-eslint/utils" "8.68.0"
|
||||
"@typescript-eslint/visitor-keys" "8.68.0"
|
||||
ignore "^7.0.5"
|
||||
natural-compare "^1.4.0"
|
||||
ts-api-utils "^2.5.0"
|
||||
|
||||
"@typescript-eslint/parser@^8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/parser/-/parser-8.65.0.tgz#5295c1058c0a1dd746ef28baaf9c0341dbdf03dc"
|
||||
integrity sha512-CZ4nMxWwgu1HEEFNkeaCptra9QCtkmKdgf3sWh1rl1trIhmxLilgTV4cwcbQ4wemnT4sWQN8CaKOmdYx+g2gMA==
|
||||
"@typescript-eslint/parser@^8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/parser/-/parser-8.68.0.tgz#61de31481354c50457bc9621a7ed746779f09ee7"
|
||||
integrity sha512-fHq2VC1kpyYfvEcbiMjOpySY4WS7voEp89yAThrHRX5sm9j2lzYppCb2umFMEed4fWcyeLjHxrz0mpjNBaBxMQ==
|
||||
dependencies:
|
||||
"@typescript-eslint/scope-manager" "8.65.0"
|
||||
"@typescript-eslint/types" "8.65.0"
|
||||
"@typescript-eslint/typescript-estree" "8.65.0"
|
||||
"@typescript-eslint/visitor-keys" "8.65.0"
|
||||
"@typescript-eslint/scope-manager" "8.68.0"
|
||||
"@typescript-eslint/types" "8.68.0"
|
||||
"@typescript-eslint/typescript-estree" "8.68.0"
|
||||
"@typescript-eslint/visitor-keys" "8.68.0"
|
||||
debug "^4.4.3"
|
||||
|
||||
"@typescript-eslint/project-service@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/project-service/-/project-service-8.65.0.tgz#65fbbc9a1591abffaeab5513200f848271cb0aa5"
|
||||
integrity sha512-SxnPhbTsGahizDgbu7oqFH/xVtzIqMd/s+WtnSxNxJZJpLbdT5IPdzg8EZxO3+PoKahXmwJLeNQOpKJb3/bi7Q==
|
||||
"@typescript-eslint/project-service@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/project-service/-/project-service-8.68.0.tgz#ea4b2869f59165c420cd7a4bbebc38039794e8cc"
|
||||
integrity sha512-5GQtWZCXFcFYux955pvoS02WLc49pXNlvIxocKjS0clvwo3in1RdlzVKyiqQH9vE5AKWFLTaUgeQkOrTS+0Qxw==
|
||||
dependencies:
|
||||
"@typescript-eslint/tsconfig-utils" "^8.65.0"
|
||||
"@typescript-eslint/types" "^8.65.0"
|
||||
"@typescript-eslint/tsconfig-utils" "^8.68.0"
|
||||
"@typescript-eslint/types" "^8.68.0"
|
||||
debug "^4.4.3"
|
||||
|
||||
"@typescript-eslint/scope-manager@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/scope-manager/-/scope-manager-8.65.0.tgz#9547202ce7e608e7b6283df585703b980a0ea70d"
|
||||
integrity sha512-Esbl8OSYiVxBokYgWPf7VVWg/BE798wXhimnn9ML9Pt5qoDf8bfQlgjlKXR/k98+AcNzlLKYrpCcrcuZ9DZLgg==
|
||||
"@typescript-eslint/scope-manager@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/scope-manager/-/scope-manager-8.68.0.tgz#e5a13a1159497faeab4e48279bf07576045b1499"
|
||||
integrity sha512-T5eXpcaJNg8bhjHJ8Rjp68Vq/QBteYtTKY8TZqVNPaUbuz0f6jI9t6aDkylwvalpAB9XTTFeFOjrjXAZ3YvmVA==
|
||||
dependencies:
|
||||
"@typescript-eslint/types" "8.65.0"
|
||||
"@typescript-eslint/visitor-keys" "8.65.0"
|
||||
"@typescript-eslint/types" "8.68.0"
|
||||
"@typescript-eslint/visitor-keys" "8.68.0"
|
||||
|
||||
"@typescript-eslint/tsconfig-utils@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.65.0.tgz#36f168fcdbb1295f7446ff0379667f98c3cf1bf3"
|
||||
integrity sha512-j6GzGqCiRdA7Qhur2VVmKZAkBLfnHFQfx4TaJGL9RMveZqCo48jSHHO0DTgizEnGhtWnqmbtCUSrqSkdiY/0Hg==
|
||||
"@typescript-eslint/tsconfig-utils@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.68.0.tgz#594d7a3c5952055b3c431fc563ca7fd1defcce18"
|
||||
integrity sha512-F7zrGQfiJHojPwi8vhxZQC1tWtJzvL74cK/nqri2lk8YUXvYaYwl263xOJ69jDWPUk1hmcdoayFwk9lX09npVw==
|
||||
|
||||
"@typescript-eslint/tsconfig-utils@^8.65.0":
|
||||
version "8.66.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.66.0.tgz#3a89066c507aa30541dc176804685b4b444e1e52"
|
||||
integrity sha512-9D5gLYZG4rOjcoag8MQ/fWI8WqA9wcPDyOGyWtWFhvM1lHRbliqUSPIY5J3zqCU1tvSwzXxnnjhQhz5Ne7mJ4g==
|
||||
"@typescript-eslint/tsconfig-utils@^8.68.0":
|
||||
version "8.69.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.69.0.tgz#d3b0ccc781ab252a90a0b3989b9d1eb85ab59469"
|
||||
integrity sha512-xNqK7YTDZsLniQMV/4rpFR8Z5JlqeRvVjuG1YgF/mdPVH84HSD19L8CczMA0qg2RfwEV231GHH3VnToJDo4MfQ==
|
||||
|
||||
"@typescript-eslint/type-utils@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/type-utils/-/type-utils-8.65.0.tgz#d316d7522d93cff4cd14f305e02f3df2d804f9c1"
|
||||
integrity sha512-YjaZ7PRI5qY7ax2L3PbvX0rRyGtipAReCWs0mhhDBHjH/vl0g0BonaGXrKdKpMbIIsMIwDgbk/xzkBTyAltS5g==
|
||||
"@typescript-eslint/type-utils@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/type-utils/-/type-utils-8.68.0.tgz#8f3e838dbd740909db27053857468cd037b00220"
|
||||
integrity sha512-X77zqoY1EjeWGs/0JNxeaMfp5C5lIz4Tw8y66F1Ne8Faq6g424sBNYM6xBAqElfGZPLpWS+CZAp0DXyKDzWiHg==
|
||||
dependencies:
|
||||
"@typescript-eslint/types" "8.65.0"
|
||||
"@typescript-eslint/typescript-estree" "8.65.0"
|
||||
"@typescript-eslint/utils" "8.65.0"
|
||||
"@typescript-eslint/types" "8.68.0"
|
||||
"@typescript-eslint/typescript-estree" "8.68.0"
|
||||
"@typescript-eslint/utils" "8.68.0"
|
||||
debug "^4.4.3"
|
||||
ts-api-utils "^2.5.0"
|
||||
|
||||
"@typescript-eslint/types@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/types/-/types-8.65.0.tgz#3e86738416a777c8b8925ab46745f48ecf904c9f"
|
||||
integrity sha512-JSSwWNy+H0E/01jJEM+hrX6N0OFDzFzeIhHFSAS01tlVaevpG8cFyYRPhS5yjGOvBUx3sqQHVMjCL1CAZZMxBg==
|
||||
"@typescript-eslint/types@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/types/-/types-8.68.0.tgz#3f9d4e62fbe5728f09403cdc7b4d58af842ac1af"
|
||||
integrity sha512-9RnpsGJjrAllCMefGVVsImJM24YurhC0Q1h4UbvivtvOqXmR/vEJge2OoE++z9m6hyg8T1Q8t5SNT6tHSbrxcg==
|
||||
|
||||
"@typescript-eslint/types@^8.65.0":
|
||||
version "8.66.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/types/-/types-8.66.0.tgz#3cacab94d3b564c1d48c56eb37b89f89a6d48479"
|
||||
integrity sha512-H6gcYaSDOyvL3AD/jHUtUFo2jqGgn/F6nuyuZSu0QTesxL+cP4dQoIMrODRofuJC09g64+WgZ6tE19Y1N2YIFQ==
|
||||
"@typescript-eslint/types@^8.68.0":
|
||||
version "8.69.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/types/-/types-8.69.0.tgz#5d9ad3f707c2e4f70a2db540031104df3e63bcf5"
|
||||
integrity sha512-K3VrubUPhlo9VDBS6QdI8YB5j7ClpqLRdefcz6PFrhnwicehBweqQ9Evhl4l+FYz0HdDmMqIiSX0aldGRYtDCA==
|
||||
|
||||
"@typescript-eslint/typescript-estree@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/typescript-estree/-/typescript-estree-8.65.0.tgz#f1f514808f6aa713e2d678ae8ff592a65e1632af"
|
||||
integrity sha512-JboAE2swaYt4tb1fHhHTABE2K+OLy09XfcTbhnk4Pw96f9dd2e9iYsJ28gBggHlo5z5x1rkyWvcPoTuNTd4oGg==
|
||||
"@typescript-eslint/typescript-estree@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/typescript-estree/-/typescript-estree-8.68.0.tgz#bf4165029825138ac27231a3ff02923ecd977f38"
|
||||
integrity sha512-OKKsD0tYmoNiU5PW2zehO1yO56jYOm1ShYlxon/Z0SJNidAkdVg86eg9ruRuoXf8xfnuWZGbwDsStkoXbZtIIA==
|
||||
dependencies:
|
||||
"@typescript-eslint/project-service" "8.65.0"
|
||||
"@typescript-eslint/tsconfig-utils" "8.65.0"
|
||||
"@typescript-eslint/types" "8.65.0"
|
||||
"@typescript-eslint/visitor-keys" "8.65.0"
|
||||
"@typescript-eslint/project-service" "8.68.0"
|
||||
"@typescript-eslint/tsconfig-utils" "8.68.0"
|
||||
"@typescript-eslint/types" "8.68.0"
|
||||
"@typescript-eslint/visitor-keys" "8.68.0"
|
||||
debug "^4.4.3"
|
||||
minimatch "^10.2.2"
|
||||
semver "^7.7.3"
|
||||
tinyglobby "^0.2.15"
|
||||
ts-api-utils "^2.5.0"
|
||||
|
||||
"@typescript-eslint/utils@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/utils/-/utils-8.65.0.tgz#afedd974a0c8deeef553b509df5800bafd615a72"
|
||||
integrity sha512-gXiwIHsYreboxeJucHKPvgwl7dXt50mF8s1/c00cP/WoVTyWKFdtfhRWwZiXYFU5H2O8vVoSLNrexFZjYS/SGA==
|
||||
"@typescript-eslint/utils@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/utils/-/utils-8.68.0.tgz#00547f2c8de8aca2a3c21752a9711f73206fd36d"
|
||||
integrity sha512-PB5gJMMOg0Q5P1tsgWtEAqQacJXq0qEqRHDX/YJ4FaTMLfZPpHB3gjl2EJuiZyPABxmj4ZQYiY9m1bdAJ5y7tQ==
|
||||
dependencies:
|
||||
"@eslint-community/eslint-utils" "^4.9.1"
|
||||
"@typescript-eslint/scope-manager" "8.65.0"
|
||||
"@typescript-eslint/types" "8.65.0"
|
||||
"@typescript-eslint/typescript-estree" "8.65.0"
|
||||
"@typescript-eslint/scope-manager" "8.68.0"
|
||||
"@typescript-eslint/types" "8.68.0"
|
||||
"@typescript-eslint/typescript-estree" "8.68.0"
|
||||
|
||||
"@typescript-eslint/visitor-keys@8.65.0":
|
||||
version "8.65.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/visitor-keys/-/visitor-keys-8.65.0.tgz#e3704c13cb4a1c22454c1abf28ff4737e15018c6"
|
||||
integrity sha512-8C71BQkGjiMmXtop7pHVJu1l2NNShFdkCyD6a2ezzs5vU/L3LRtb69EtcteFwz0mYMPzIgOw0n6OV4VBUWZd7A==
|
||||
"@typescript-eslint/visitor-keys@8.68.0":
|
||||
version "8.68.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/visitor-keys/-/visitor-keys-8.68.0.tgz#78db3c9bb258a0309d9e2b1b617127c3a8fb1f54"
|
||||
integrity sha512-YR65gGdGvTUAWLldC3xLOvOzamdGzB4A5/N8rehEaHs3Zvoe39BhgY+u0SPch1OvrVTfLcc55wsSgK2NcnTS/A==
|
||||
dependencies:
|
||||
"@typescript-eslint/types" "8.65.0"
|
||||
"@typescript-eslint/types" "8.68.0"
|
||||
eslint-visitor-keys "^5.0.0"
|
||||
|
||||
"@typescript/typescript-aix-ppc64@7.0.2":
|
||||
@@ -881,10 +886,10 @@ eslint-visitor-keys@^5.0.0, eslint-visitor-keys@^5.0.1:
|
||||
resolved "https://registry.yarnpkg.com/eslint-visitor-keys/-/eslint-visitor-keys-5.0.1.tgz#9e3c9489697824d2d4ce3a8ad12628f91e9f59be"
|
||||
integrity sha512-tD40eHxA35h0PEIZNeIjkHoDR4YjjJp34biM0mDvplBe//mB+IHCqHDGV7pxF+7MklTvighcCPPZC7ynWyjdTA==
|
||||
|
||||
eslint@^10.8.0:
|
||||
version "10.8.0"
|
||||
resolved "https://registry.yarnpkg.com/eslint/-/eslint-10.8.0.tgz#e6d19907a3f090a53a022261ba34c5ff6d04908b"
|
||||
integrity sha512-nuKKvN+oIBO0koN7Tm7dlkmnkc21mtt0QJLwAKzjLq14y6lRTdVG36MZHJ8eQHwdJMwZbQNMlPOYedMq/oVJvQ==
|
||||
eslint@^10.9.1:
|
||||
version "10.9.1"
|
||||
resolved "https://registry.yarnpkg.com/eslint/-/eslint-10.9.1.tgz#409da5c41a5536d5a849f8555a18ca7ef1eb963b"
|
||||
integrity sha512-9VaAkDURekixUQJy0oJYl2DcN6oKMfxay7XzaGYAWQwsb6qfKf+x76R2k1L8kb1boc+FyCAaTA9GmiKaaiaF+A==
|
||||
dependencies:
|
||||
"@eslint-community/eslint-utils" "^4.8.0"
|
||||
"@eslint-community/regexpp" "^4.12.2"
|
||||
@@ -1908,17 +1913,17 @@ tsconfig-paths@^3.15.0:
|
||||
minimist "^1.2.6"
|
||||
strip-bom "^3.0.0"
|
||||
|
||||
turbo@^2.10.8:
|
||||
version "2.10.8"
|
||||
resolved "https://registry.yarnpkg.com/turbo/-/turbo-2.10.8.tgz#09457cb7db79710586da455a29bccc05f15b1b72"
|
||||
integrity sha512-9+8YX5QOkGXzZxcIykTHgaooRHGMWO+jfdyRK0o+rN0U7hBIig2MrJ8r/aNzIPDPhdA73SGb0O+tIztaModTMg==
|
||||
turbo@^2.10.12:
|
||||
version "2.10.12"
|
||||
resolved "https://registry.yarnpkg.com/turbo/-/turbo-2.10.12.tgz#22f552bd88182d58365960a02e5f45e628fb965b"
|
||||
integrity sha512-AswgMPnpOoaVZHrrSBejETzEbuIA69OVGwfkHwfrY0A23VjWXBANzgq9+OymWOHAIArB7D1+1z498WY8fGg1Jw==
|
||||
optionalDependencies:
|
||||
"@turbo/darwin-64" "2.10.8"
|
||||
"@turbo/darwin-arm64" "2.10.8"
|
||||
"@turbo/linux-64" "2.10.8"
|
||||
"@turbo/linux-arm64" "2.10.8"
|
||||
"@turbo/windows-64" "2.10.8"
|
||||
"@turbo/windows-arm64" "2.10.8"
|
||||
"@turbo/darwin-64" "2.10.12"
|
||||
"@turbo/darwin-arm64" "2.10.12"
|
||||
"@turbo/linux-64" "2.10.12"
|
||||
"@turbo/linux-arm64" "2.10.12"
|
||||
"@turbo/windows-64" "2.10.12"
|
||||
"@turbo/windows-arm64" "2.10.12"
|
||||
|
||||
type-check@^0.4.0, type-check@~0.4.0:
|
||||
version "0.4.0"
|
||||
|
||||
@@ -970,6 +970,22 @@ def python_config_to_docker_uv_lock(
|
||||
f"{uv_export_project_dir}/uv.lock",
|
||||
)
|
||||
)
|
||||
for package_root in sorted(
|
||||
plan.all_workspace_roots,
|
||||
key=lambda root: root.as_posix(),
|
||||
):
|
||||
if package_root == plan.project_root:
|
||||
continue
|
||||
package_relative_path = pathlib.PurePosixPath(
|
||||
package_root.relative_to(plan.project_root).as_posix()
|
||||
)
|
||||
package_pyproject_path = package_relative_path / "pyproject.toml"
|
||||
docker_plan.add_raw(
|
||||
copy_from_project_root(
|
||||
package_pyproject_path,
|
||||
f"{uv_export_project_dir}/{package_pyproject_path.as_posix()}",
|
||||
)
|
||||
)
|
||||
docker_plan.add_instruction("WORKDIR", uv_export_project_dir)
|
||||
docker_plan.add_instruction(
|
||||
"RUN",
|
||||
|
||||
@@ -1403,6 +1403,19 @@ def test_config_to_docker_uv_lock():
|
||||
"COPY --from=uv-workspace-root uv.lock /tmp/uv_export/project/uv.lock"
|
||||
in docker
|
||||
)
|
||||
workspace_pyprojects = [
|
||||
"apps/agent/pyproject.toml",
|
||||
"libs/extra/pyproject.toml",
|
||||
"libs/shared/pyproject.toml",
|
||||
]
|
||||
export_instruction = "RUN uv export --package agent"
|
||||
for pyproject_path in workspace_pyprojects:
|
||||
copy_instruction = (
|
||||
"COPY --from=uv-workspace-root "
|
||||
f"{pyproject_path} /tmp/uv_export/project/{pyproject_path}"
|
||||
)
|
||||
assert copy_instruction in docker
|
||||
assert docker.index(copy_instruction) < docker.index(export_instruction)
|
||||
assert additional_contexts == {"uv-workspace-root": str(project_root.resolve())}
|
||||
|
||||
assert (
|
||||
|
||||
Generated
+6
-6
@@ -266,20 +266,20 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint"
|
||||
version = "4.0.1"
|
||||
version = "4.2.0"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
{ name = "ormsgpack" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/b1/44/a8df45d1e8b4637e29789fa8bae1db022c953cc7ac80093cfc52e923547e/langgraph_checkpoint-4.0.1.tar.gz", hash = "sha256:b433123735df11ade28829e40ce25b9be614930cd50245ff2af60629234befd9", size = 158135, upload-time = "2026-02-27T21:06:16.092Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/dc/e1/089c4c9e0a2fec7f883f82ae8e6a727138d50074cfeb6644bc2d13b1019b/langgraph_checkpoint-4.2.0.tar.gz", hash = "sha256:51a593b6bee684b0818e5d6e58e28ab340c6db7794575056ce7bd1b746a84ed7", size = 180239, upload-time = "2026-08-07T20:05:03.756Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/65/4c/09a4a0c42f5d2fc38d6c4d67884788eff7fd2cfdf367fdf7033de908b4c0/langgraph_checkpoint-4.0.1-py3-none-any.whl", hash = "sha256:e3adcd7a0e0166f3b48b8cf508ce0ea366e7420b5a73aa81289888727769b034", size = 50453, upload-time = "2026-02-27T21:06:14.293Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/05/71/3b475f09bd57d3a5649792c66353312b4432afd843f301739dfcebd157f0/langgraph_checkpoint-4.2.0-py3-none-any.whl", hash = "sha256:0547fd228935a0b758865de3a3d6d7a2537c308895d0f9ab092ce9151b5da942", size = 56833, upload-time = "2026-08-07T20:05:02.655Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint-postgres"
|
||||
version = "3.0.5"
|
||||
version = "3.1.1"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "langgraph-checkpoint" },
|
||||
@@ -287,9 +287,9 @@ dependencies = [
|
||||
{ name = "psycopg" },
|
||||
{ name = "psycopg-pool" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/95/7a/8f439966643d32111248a225e6cb33a182d07c90de780c4dbfc1e0377832/langgraph_checkpoint_postgres-3.0.5.tar.gz", hash = "sha256:a8fd7278a63f4f849b5cbc7884a15ca8f41e7d5f7467d0a66b31e8c24492f7eb", size = 127856, upload-time = "2026-03-18T21:25:29.785Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/06/92/1e8959f8cd1b56e672fde3227f6fd642be85af6c5fd662d73921074aa39d/langgraph_checkpoint_postgres-3.1.1.tar.gz", hash = "sha256:d320e147ddad8c374cd546df0b52b532dd54d0541dd9fd23fc738cbd5de76f41", size = 150413, upload-time = "2026-07-30T19:15:39.014Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/e8/87/b0f98b33a67204bca9d5619bcd9574222f6b025cf3c125eedcec9a50ecbc/langgraph_checkpoint_postgres-3.0.5-py3-none-any.whl", hash = "sha256:86d7040a88fd70087eaafb72251d796696a0a2d856168f5c11ef620771411552", size = 42907, upload-time = "2026-03-18T21:25:28.75Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/03/32/ba457698a48a0e18d786caa770033067049fbe36d6846f8e50f13b594b51/langgraph_checkpoint_postgres-3.1.1-py3-none-any.whl", hash = "sha256:6e353aecd8150de144fef8e51a49076f58b7d6830d4cf51392b7ad4d79832ba7", size = 50778, upload-time = "2026-07-30T19:15:37.405Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
|
||||
Generated
+1472
-1135
File diff suppressed because it is too large
Load Diff
@@ -3,6 +3,7 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import enum
|
||||
import inspect
|
||||
import logging
|
||||
import sys
|
||||
import warnings
|
||||
from collections.abc import (
|
||||
@@ -63,6 +64,26 @@ try:
|
||||
except ImportError:
|
||||
_StreamingCallbackHandler = None # type: ignore
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _trace_payload(value: Any, transform: Callable[[Any], Any] | None) -> Any:
|
||||
"""Return the payload to record on a run for `value`.
|
||||
|
||||
When `transform` is unset this is a passthrough, so unspecified nodes record exactly
|
||||
as before. When set it always runs (regardless of tracing), but never affects
|
||||
execution: if it raises, the untransformed value is recorded instead.
|
||||
"""
|
||||
if transform is None:
|
||||
return value
|
||||
try:
|
||||
return transform(value)
|
||||
except Exception:
|
||||
logger.exception(
|
||||
"trace input/output processor raised; recording untransformed payload"
|
||||
)
|
||||
return value
|
||||
|
||||
|
||||
def _set_config_context(
|
||||
config: RunnableConfig, run: Any = None
|
||||
@@ -572,6 +593,7 @@ class RunnableSeq(Runnable):
|
||||
*steps: RunnableLike,
|
||||
name: str | None = None,
|
||||
trace_inputs: Callable[[Any], Any] | None = None,
|
||||
trace_outputs: Callable[[Any], Any] | None = None,
|
||||
) -> None:
|
||||
"""Create a new RunnableSeq.
|
||||
|
||||
@@ -597,6 +619,7 @@ class RunnableSeq(Runnable):
|
||||
self.steps = steps_flat
|
||||
self.name = name
|
||||
self.trace_inputs = trace_inputs
|
||||
self.trace_outputs = trace_outputs
|
||||
|
||||
def __or__(
|
||||
self,
|
||||
@@ -658,7 +681,7 @@ class RunnableSeq(Runnable):
|
||||
# start the root run
|
||||
run_manager = callback_manager.on_chain_start(
|
||||
None,
|
||||
self.trace_inputs(input) if self.trace_inputs is not None else input,
|
||||
_trace_payload(input, self.trace_inputs),
|
||||
name=config.get("run_name") or self.get_name(),
|
||||
run_id=config.pop("run_id", None),
|
||||
)
|
||||
@@ -689,7 +712,7 @@ class RunnableSeq(Runnable):
|
||||
run_manager.on_chain_error(e)
|
||||
raise
|
||||
else:
|
||||
run_manager.on_chain_end(input)
|
||||
run_manager.on_chain_end(_trace_payload(input, self.trace_outputs))
|
||||
return input
|
||||
|
||||
async def ainvoke(
|
||||
@@ -705,7 +728,7 @@ class RunnableSeq(Runnable):
|
||||
# start the root run
|
||||
run_manager = await callback_manager.on_chain_start(
|
||||
None,
|
||||
self.trace_inputs(input) if self.trace_inputs is not None else input,
|
||||
_trace_payload(input, self.trace_inputs),
|
||||
name=config.get("run_name") or self.get_name(),
|
||||
run_id=config.pop("run_id", None),
|
||||
)
|
||||
@@ -742,7 +765,7 @@ class RunnableSeq(Runnable):
|
||||
await run_manager.on_chain_error(e)
|
||||
raise
|
||||
else:
|
||||
await run_manager.on_chain_end(input)
|
||||
await run_manager.on_chain_end(_trace_payload(input, self.trace_outputs))
|
||||
return input
|
||||
|
||||
def stream(
|
||||
@@ -758,7 +781,7 @@ class RunnableSeq(Runnable):
|
||||
# start the root run
|
||||
run_manager = callback_manager.on_chain_start(
|
||||
None,
|
||||
self.trace_inputs(input) if self.trace_inputs is not None else input,
|
||||
_trace_payload(input, self.trace_inputs),
|
||||
name=config.get("run_name") or self.get_name(),
|
||||
run_id=config.pop("run_id", None),
|
||||
)
|
||||
@@ -803,7 +826,7 @@ class RunnableSeq(Runnable):
|
||||
run_manager.on_chain_error(e)
|
||||
raise
|
||||
else:
|
||||
run_manager.on_chain_end(output)
|
||||
run_manager.on_chain_end(_trace_payload(output, self.trace_outputs))
|
||||
|
||||
async def astream(
|
||||
self,
|
||||
@@ -818,7 +841,7 @@ class RunnableSeq(Runnable):
|
||||
# start the root run
|
||||
run_manager = await callback_manager.on_chain_start(
|
||||
None,
|
||||
self.trace_inputs(input) if self.trace_inputs is not None else input,
|
||||
_trace_payload(input, self.trace_inputs),
|
||||
name=config.get("run_name") or self.get_name(),
|
||||
run_id=config.pop("run_id", None),
|
||||
)
|
||||
@@ -873,7 +896,9 @@ class RunnableSeq(Runnable):
|
||||
await run_manager.on_chain_error(e)
|
||||
raise
|
||||
else:
|
||||
await run_manager.on_chain_end(output)
|
||||
await run_manager.on_chain_end(
|
||||
_trace_payload(output, self.trace_outputs)
|
||||
)
|
||||
else:
|
||||
try:
|
||||
async with AsyncExitStack() as stack:
|
||||
@@ -903,7 +928,9 @@ class RunnableSeq(Runnable):
|
||||
await run_manager.on_chain_error(e)
|
||||
raise
|
||||
else:
|
||||
await run_manager.on_chain_end(output)
|
||||
await run_manager.on_chain_end(
|
||||
_trace_payload(output, self.trace_outputs)
|
||||
)
|
||||
|
||||
|
||||
def _consume_iter(it: Iterator[Any]) -> Any:
|
||||
|
||||
@@ -9,7 +9,13 @@ from langgraph.store.base import BaseStore
|
||||
|
||||
from langgraph._internal._typing import EMPTY_SEQ
|
||||
from langgraph.runtime import Runtime
|
||||
from langgraph.types import CachePolicy, RetryPolicy, StreamWriter, TimeoutPolicy
|
||||
from langgraph.types import (
|
||||
CachePolicy,
|
||||
RetryPolicy,
|
||||
StreamWriter,
|
||||
TimeoutPolicy,
|
||||
TracePolicy,
|
||||
)
|
||||
from langgraph.typing import ContextT, NodeInputT, NodeInputT_contra
|
||||
|
||||
|
||||
@@ -93,3 +99,5 @@ class StateNodeSpec(Generic[NodeInputT, ContextT]):
|
||||
ends: tuple[str, ...] | dict[str, str] | None = EMPTY_SEQ
|
||||
defer: bool = False
|
||||
timeout: TimeoutPolicy | None = None
|
||||
trace_policy: TracePolicy | None = None
|
||||
"""Optional policy controlling what this node records on its trace run."""
|
||||
|
||||
@@ -85,6 +85,7 @@ from langgraph.types import (
|
||||
RetryPolicy,
|
||||
Send,
|
||||
TimeoutPolicy,
|
||||
TracePolicy,
|
||||
ensure_valid_checkpointer,
|
||||
)
|
||||
from langgraph.typing import ContextT, InputT, NodeInputT, OutputT, StateT
|
||||
@@ -384,6 +385,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
error_handler: StateNode[Any, ContextT] | None = None,
|
||||
destinations: dict[str, str] | tuple[str, ...] | None = None,
|
||||
timeout: float | timedelta | TimeoutPolicy | None = None,
|
||||
trace_policy: TracePolicy | None = None,
|
||||
**kwargs: Unpack[DeprecatedKwargs],
|
||||
) -> Self:
|
||||
"""Add a new node to the `StateGraph`, input schema is inferred as the state schema.
|
||||
@@ -453,6 +455,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
error_handler: StateNode[Any, ContextT] | None = None,
|
||||
destinations: dict[str, str] | tuple[str, ...] | None = None,
|
||||
timeout: float | timedelta | TimeoutPolicy | None = None,
|
||||
trace_policy: TracePolicy | None = None,
|
||||
**kwargs: Unpack[DeprecatedKwargs],
|
||||
) -> Self:
|
||||
"""Add a new node to the `StateGraph` where input schema is specified.
|
||||
@@ -527,6 +530,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
error_handler: StateNode[Any, ContextT] | None = None,
|
||||
destinations: dict[str, str] | tuple[str, ...] | None = None,
|
||||
timeout: float | timedelta | TimeoutPolicy | None = None,
|
||||
trace_policy: TracePolicy | None = None,
|
||||
**kwargs: Unpack[DeprecatedKwargs],
|
||||
) -> Self:
|
||||
"""Add a new node to the `StateGraph`, input schema is inferred as the state schema.
|
||||
@@ -596,6 +600,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
error_handler: StateNode[Any, ContextT] | None = None,
|
||||
destinations: dict[str, str] | tuple[str, ...] | None = None,
|
||||
timeout: float | timedelta | TimeoutPolicy | None = None,
|
||||
trace_policy: TracePolicy | None = None,
|
||||
**kwargs: Unpack[DeprecatedKwargs],
|
||||
) -> Self:
|
||||
"""Add a new node to the `StateGraph`, input schema is specified.
|
||||
@@ -672,6 +677,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
error_handler: StateNode[Any, ContextT] | None = None,
|
||||
destinations: dict[str, str] | tuple[str, ...] | None = None,
|
||||
timeout: float | timedelta | TimeoutPolicy | None = None,
|
||||
trace_policy: TracePolicy | None = None,
|
||||
**kwargs: Unpack[DeprecatedKwargs],
|
||||
) -> Self:
|
||||
"""Add a new node to the `StateGraph`.
|
||||
@@ -691,6 +697,10 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
If a sequence is provided, the first matching policy will be applied.
|
||||
cache_policy: The cache policy for the node.
|
||||
error_handler: Optional node-level error handler callable for this node.
|
||||
trace_policy: Optional policy controlling how this node's run is traced. Its
|
||||
`process_inputs` callable transforms the node's input before it is
|
||||
recorded (e.g. to omit or summarize large message history) without
|
||||
changing the value passed to the node. Does not affect execution.
|
||||
destinations: Destinations that indicate where a node can route to.
|
||||
|
||||
Useful for edgeless graphs with nodes that return `Command` objects.
|
||||
@@ -880,6 +890,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
ends=ends,
|
||||
defer=defer,
|
||||
timeout=timeout,
|
||||
trace_policy=trace_policy,
|
||||
)
|
||||
elif inferred_input_schema is not None:
|
||||
self.nodes[node] = StateNodeSpec(
|
||||
@@ -892,6 +903,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
ends=ends,
|
||||
defer=defer,
|
||||
timeout=timeout,
|
||||
trace_policy=trace_policy,
|
||||
)
|
||||
else:
|
||||
self.nodes[node] = StateNodeSpec[StateT, ContextT](
|
||||
@@ -904,6 +916,7 @@ class StateGraph(Generic[StateT, ContextT, InputT, OutputT]):
|
||||
ends=ends,
|
||||
defer=defer,
|
||||
timeout=timeout,
|
||||
trace_policy=trace_policy,
|
||||
)
|
||||
|
||||
input_schema = input_schema or inferred_input_schema
|
||||
@@ -1530,6 +1543,7 @@ class CompiledStateGraph(
|
||||
error_handler_node=node.error_handler_node,
|
||||
bound=node.runnable, # type: ignore[arg-type]
|
||||
timeout=node.timeout,
|
||||
trace_policy=node.trace_policy,
|
||||
)
|
||||
else:
|
||||
raise RuntimeError
|
||||
|
||||
@@ -16,7 +16,7 @@ from langgraph._internal._timeout import coerce_timeout_policy
|
||||
from langgraph.pregel._utils import find_subgraph_pregel
|
||||
from langgraph.pregel._write import ChannelWrite
|
||||
from langgraph.pregel.protocol import PregelProtocol
|
||||
from langgraph.types import CachePolicy, RetryPolicy, TimeoutPolicy
|
||||
from langgraph.types import CachePolicy, RetryPolicy, TimeoutPolicy, TracePolicy
|
||||
|
||||
READ_TYPE = Callable[[str | Sequence[str], bool], Any | dict[str, Any]]
|
||||
INPUT_CACHE_KEY_TYPE = tuple[Callable[..., Any], tuple[str, ...]]
|
||||
@@ -138,6 +138,9 @@ class PregelNode:
|
||||
metadata: Mapping[str, Any] | None
|
||||
"""Metadata to attach to the node for tracing."""
|
||||
|
||||
trace_policy: TracePolicy | None
|
||||
"""Optional policy controlling what this node records on its trace run."""
|
||||
|
||||
is_error_handler: bool
|
||||
"""Whether this node is registered as an error handler node."""
|
||||
|
||||
@@ -156,6 +159,7 @@ class PregelNode:
|
||||
writers: list[Runnable] | None = None,
|
||||
tags: list[str] | None = None,
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
trace_policy: TracePolicy | None = None,
|
||||
bound: Runnable[Any, Any] | None = None,
|
||||
retry_policy: RetryPolicy | Sequence[RetryPolicy] | None = None,
|
||||
cache_policy: CachePolicy | None = None,
|
||||
@@ -177,6 +181,7 @@ class PregelNode:
|
||||
self.timeout = coerce_timeout_policy(timeout)
|
||||
self.tags = tags
|
||||
self.metadata = metadata
|
||||
self.trace_policy = trace_policy
|
||||
self.is_error_handler = is_error_handler
|
||||
self.error_handler_node = error_handler_node
|
||||
if subgraphs is not None:
|
||||
@@ -222,14 +227,23 @@ class PregelNode:
|
||||
def node(self) -> Runnable[Any, Any] | None:
|
||||
"""Get a runnable that combines `bound` and `writers`."""
|
||||
writers = self.flat_writers
|
||||
trace_inputs = self.trace_policy.process_inputs if self.trace_policy else None
|
||||
trace_outputs = self.trace_policy.process_outputs if self.trace_policy else None
|
||||
if self.bound is DEFAULT_BOUND and not writers:
|
||||
return None
|
||||
elif self.bound is DEFAULT_BOUND and len(writers) == 1:
|
||||
return writers[0]
|
||||
elif self.bound is DEFAULT_BOUND:
|
||||
return RunnableSeq(*writers)
|
||||
return RunnableSeq(
|
||||
*writers, trace_inputs=trace_inputs, trace_outputs=trace_outputs
|
||||
)
|
||||
elif writers:
|
||||
return RunnableSeq(self.bound, *writers)
|
||||
return RunnableSeq(
|
||||
self.bound,
|
||||
*writers,
|
||||
trace_inputs=trace_inputs,
|
||||
trace_outputs=trace_outputs,
|
||||
)
|
||||
else:
|
||||
return self.bound
|
||||
|
||||
|
||||
@@ -1,11 +1,10 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import ast
|
||||
import inspect
|
||||
import dis
|
||||
import re
|
||||
import textwrap
|
||||
from collections.abc import Callable, Sequence
|
||||
from functools import partial
|
||||
from types import CodeType, FunctionType
|
||||
from typing import Any
|
||||
|
||||
from langchain_core.runnables import (
|
||||
@@ -17,7 +16,6 @@ from langchain_core.runnables import (
|
||||
from langchain_core.runnables.base import RunnableBindingBase
|
||||
from langchain_core.runnables.config import run_in_executor
|
||||
from langgraph.checkpoint.base import ChannelVersions
|
||||
from typing_extensions import override
|
||||
|
||||
from langgraph._internal._runnable import RunnableCallable, RunnableSeq
|
||||
from langgraph._internal._timeout import sync_timeout_unsupported
|
||||
@@ -137,155 +135,87 @@ def validate_timeout_supported(runnable: Runnable, *, name: str) -> None:
|
||||
raise sync_timeout_unsupported(name)
|
||||
|
||||
|
||||
# Values treated as dead ends when deciding whether to walk a function's
|
||||
# bytecode. A container can hold a graph, but `find_subgraph_pregel` does not
|
||||
# look inside one, so skipping it costs nothing while that holds. Matched by
|
||||
# exact type, since a subclass of a builtin can carry attributes.
|
||||
_LEAF_TYPES = frozenset(
|
||||
{
|
||||
int,
|
||||
float,
|
||||
complex,
|
||||
bool,
|
||||
str,
|
||||
bytes,
|
||||
bytearray,
|
||||
list,
|
||||
tuple,
|
||||
dict,
|
||||
set,
|
||||
frozenset,
|
||||
type(None),
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def get_function_nonlocals(func: Callable) -> list[Any]:
|
||||
"""Get the nonlocal variables accessed by a function.
|
||||
"""Get the values a function reaches from outside its own scope.
|
||||
|
||||
Args:
|
||||
func: The function to check.
|
||||
|
||||
Returns:
|
||||
List[Any]: The nonlocal variables accessed by the function.
|
||||
Every captured cell value, the globals the function names, and each
|
||||
value along an attribute path it loads. Over-approximates: a value can
|
||||
come back without the function reaching it at runtime.
|
||||
"""
|
||||
try:
|
||||
code = inspect.getsource(func)
|
||||
tree = ast.parse(textwrap.dedent(code))
|
||||
visitor = FunctionNonLocals()
|
||||
visitor.visit(tree)
|
||||
values: list[Any] = []
|
||||
closure = (
|
||||
inspect.getclosurevars(func.__wrapped__)
|
||||
if hasattr(func, "__wrapped__") and callable(func.__wrapped__)
|
||||
else inspect.getclosurevars(func)
|
||||
)
|
||||
candidates = {**closure.globals, **closure.nonlocals}
|
||||
for k, v in candidates.items():
|
||||
if k in visitor.nonlocals:
|
||||
values.append(v)
|
||||
for kk in visitor.nonlocals:
|
||||
if "." in kk and kk.startswith(k):
|
||||
vv = v
|
||||
for part in kk.split(".")[1:]:
|
||||
if vv is None:
|
||||
break
|
||||
else:
|
||||
try:
|
||||
vv = getattr(vv, part)
|
||||
except AttributeError:
|
||||
break
|
||||
else:
|
||||
values.append(vv)
|
||||
except (SyntaxError, TypeError, OSError, SystemError):
|
||||
func = getattr(func, "__func__", func) # bound method -> function
|
||||
wrapped = getattr(func, "__wrapped__", None)
|
||||
if callable(wrapped):
|
||||
func = getattr(wrapped, "__func__", wrapped)
|
||||
if not isinstance(func, FunctionType):
|
||||
return []
|
||||
code = func.__code__
|
||||
|
||||
cells: dict[str, Any] = {}
|
||||
for name, cell in zip(code.co_freevars, func.__closure__ or ()):
|
||||
try:
|
||||
cells[name] = cell.cell_contents
|
||||
except ValueError:
|
||||
continue # empty cell: a recursive def not yet bound
|
||||
|
||||
# Every captured value counts, referenced or not: over-declaring costs an
|
||||
# introspection entry, under-declaring drops the subgraph's checkpoints and
|
||||
# stream events. Checking each cell against the bytecode would cost more and
|
||||
# only trade the cheap error for the expensive one.
|
||||
values: list[Any] = list(cells.values())
|
||||
global_ns = func.__globals__
|
||||
globals_ = {name: global_ns[name] for name in code.co_names if name in global_ns}
|
||||
if all(type(v) in _LEAF_TYPES for v in (*cells.values(), *globals_.values())):
|
||||
return values
|
||||
|
||||
# Nested code objects hold the references made by inner defs, lambdas and
|
||||
# comprehensions, which resolve against the namespaces gathered above.
|
||||
codes = [code]
|
||||
for c in codes:
|
||||
codes.extend(k for k in c.co_consts if isinstance(k, CodeType))
|
||||
value: Any = None
|
||||
for instruction in dis.get_instructions(c):
|
||||
opname = instruction.opname
|
||||
if opname == "LOAD_GLOBAL":
|
||||
value = globals_.get(instruction.argval)
|
||||
elif opname == "LOAD_DEREF":
|
||||
value = cells.get(instruction.argval)
|
||||
elif opname in ("LOAD_ATTR", "LOAD_METHOD"):
|
||||
value = getattr(value, instruction.argval, None)
|
||||
else:
|
||||
value = None # anything else ends the chain: `a, b.c` is not `a.c`
|
||||
continue
|
||||
if value is not None:
|
||||
values.append(value)
|
||||
return values
|
||||
|
||||
|
||||
class FunctionNonLocals(ast.NodeVisitor):
|
||||
"""Get the nonlocal variables accessed of a function."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.nonlocals: set[str] = set()
|
||||
|
||||
@override
|
||||
def visit_FunctionDef(self, node: ast.FunctionDef) -> Any:
|
||||
"""Visit a function definition.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
visitor = NonLocals()
|
||||
visitor.visit(node)
|
||||
self.nonlocals.update(visitor.loads - visitor.stores)
|
||||
|
||||
@override
|
||||
def visit_AsyncFunctionDef(self, node: ast.AsyncFunctionDef) -> Any:
|
||||
"""Visit an async function definition.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
visitor = NonLocals()
|
||||
visitor.visit(node)
|
||||
self.nonlocals.update(visitor.loads - visitor.stores)
|
||||
|
||||
@override
|
||||
def visit_Lambda(self, node: ast.Lambda) -> Any:
|
||||
"""Visit a lambda function.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
visitor = NonLocals()
|
||||
visitor.visit(node)
|
||||
self.nonlocals.update(visitor.loads - visitor.stores)
|
||||
|
||||
|
||||
class NonLocals(ast.NodeVisitor):
|
||||
"""Get nonlocal variables accessed."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.loads: set[str] = set()
|
||||
self.stores: set[str] = set()
|
||||
|
||||
@override
|
||||
def visit_Name(self, node: ast.Name) -> Any:
|
||||
"""Visit a name node.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
if isinstance(node.ctx, ast.Load):
|
||||
self.loads.add(node.id)
|
||||
elif isinstance(node.ctx, ast.Store):
|
||||
self.stores.add(node.id)
|
||||
|
||||
@override
|
||||
def visit_Attribute(self, node: ast.Attribute) -> Any:
|
||||
"""Visit an attribute node.
|
||||
|
||||
Args:
|
||||
node: The node to visit.
|
||||
|
||||
Returns:
|
||||
Any: The result of the visit.
|
||||
"""
|
||||
if isinstance(node.ctx, ast.Load):
|
||||
parent = node.value
|
||||
attr_expr = node.attr
|
||||
while isinstance(parent, ast.Attribute):
|
||||
attr_expr = parent.attr + "." + attr_expr
|
||||
parent = parent.value
|
||||
if isinstance(parent, ast.Name):
|
||||
self.loads.add(parent.id + "." + attr_expr)
|
||||
self.loads.discard(parent.id)
|
||||
elif isinstance(parent, ast.Call):
|
||||
if isinstance(parent.func, ast.Name):
|
||||
self.loads.add(parent.func.id)
|
||||
else:
|
||||
parent = parent.func
|
||||
attr_expr = ""
|
||||
while isinstance(parent, ast.Attribute):
|
||||
if attr_expr:
|
||||
attr_expr = parent.attr + "." + attr_expr
|
||||
else:
|
||||
attr_expr = parent.attr
|
||||
parent = parent.value
|
||||
if isinstance(parent, ast.Name):
|
||||
self.loads.add(parent.id + "." + attr_expr)
|
||||
|
||||
|
||||
def is_xxh3_128_hexdigest(value: str) -> bool:
|
||||
"""Check if the given string matches the format of xxh3_128_hexdigest."""
|
||||
return bool(re.fullmatch(r"[0-9a-f]{32}", value))
|
||||
|
||||
@@ -2,7 +2,7 @@ from __future__ import annotations
|
||||
|
||||
from abc import abstractmethod
|
||||
from collections.abc import AsyncIterator, Callable, Iterator, Sequence
|
||||
from typing import Any, Generic, Literal, cast, overload
|
||||
from typing import Any, Generic, Literal, overload
|
||||
|
||||
from langchain_core.runnables import Runnable, RunnableConfig
|
||||
from langchain_core.runnables.graph import Graph as DrawableGraph
|
||||
@@ -277,12 +277,12 @@ class StreamProtocol:
|
||||
|
||||
modes: set[StreamMode]
|
||||
|
||||
__call__: Callable[[Self, StreamChunk], None]
|
||||
__call__: Callable[[StreamChunk], None]
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
__call__: Callable[[StreamChunk], None],
|
||||
modes: set[StreamMode],
|
||||
) -> None:
|
||||
self.__call__ = cast(Callable[[Self, StreamChunk], None], __call__)
|
||||
self.__call__ = __call__
|
||||
self.modes = modes
|
||||
|
||||
@@ -3,7 +3,7 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
from collections.abc import AsyncIterator, Awaitable, Callable, Iterator, Mapping
|
||||
from types import MappingProxyType, TracebackType
|
||||
from typing import TYPE_CHECKING, Any
|
||||
from typing import TYPE_CHECKING, Any, NoReturn
|
||||
|
||||
from langchain_core._api import beta
|
||||
|
||||
@@ -33,6 +33,26 @@ async def _adrive_until_done(pump: Callable[[], Awaitable[bool]]) -> None:
|
||||
pass
|
||||
|
||||
|
||||
def _raise_missing_projection(run: object, name: str) -> NoReturn:
|
||||
"""Raise after normal attribute lookup fails for a projection.
|
||||
|
||||
Registered native projections are installed directly on the run instance
|
||||
during `__init__`, so `__getattr__` is never called for them. At this point
|
||||
the requested name is necessarily missing; the mux is inspected only to
|
||||
include the valid registered projection names in the error message.
|
||||
|
||||
Read `_mux` directly from `__dict__` because it may not exist yet on a
|
||||
partially initialized instance. Accessing `run._mux` in that case would
|
||||
invoke `__getattr__` again and recurse indefinitely.
|
||||
"""
|
||||
mux = run.__dict__.get("_mux")
|
||||
registered = sorted(mux.native_keys) if mux is not None else []
|
||||
raise AttributeError(
|
||||
f"{type(run).__name__!r} object has no attribute {name!r} "
|
||||
f"(registered projections: {', '.join(registered) or 'none'})"
|
||||
)
|
||||
|
||||
|
||||
@beta(message="The v3 streaming protocol on Pregel is experimental.")
|
||||
class GraphRunStream:
|
||||
"""Sync run stream with caller-driven pumping.
|
||||
@@ -54,15 +74,32 @@ class GraphRunStream:
|
||||
experimental and may change.
|
||||
"""
|
||||
|
||||
# Native projections always registered by `stream_events(version="v3")`.
|
||||
# Attached dynamically by the `setattr` loop in `__init__`; declared here
|
||||
# so type checkers see them. Opt-in native projections (`updates`,
|
||||
# `custom`, `checkpoints`, `debug`, `tasks`) are only present when their
|
||||
# transformer is registered, so they are reached via `extensions[...]`.
|
||||
# Native projections, attached dynamically by the `setattr` loop in
|
||||
# `__init__` and declared here so type checkers see them.
|
||||
#
|
||||
# Always registered by `stream_events(version="v3")`:
|
||||
values: StreamChannel[dict[str, Any]]
|
||||
messages: StreamChannel[ChatModelStream]
|
||||
lifecycle: StreamChannel[LifecyclePayload]
|
||||
subgraphs: StreamChannel[SubgraphRunStream]
|
||||
# Registered on demand via `compile(transformers=...)` or
|
||||
# `stream_events(transformers=...)`; reading one whose transformer was not
|
||||
# registered raises AttributeError. Projections contributed by transformers
|
||||
# outside this package are covered by `__getattr__` instead.
|
||||
updates: StreamChannel[dict[str, Any]]
|
||||
custom: StreamChannel[Any]
|
||||
checkpoints: StreamChannel[dict[str, Any]]
|
||||
debug: StreamChannel[dict[str, Any]]
|
||||
tasks: StreamChannel[dict[str, Any]]
|
||||
|
||||
def __getattr__(self, name: str) -> StreamChannel[Any]:
|
||||
"""Type the projections of transformers declared outside this package.
|
||||
|
||||
Projection names come from a registry, so no annotation here can name
|
||||
them all. The cost is that a misspelling type-checks too, and fails at
|
||||
runtime instead.
|
||||
"""
|
||||
_raise_missing_projection(self, name)
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
@@ -345,15 +382,24 @@ class AsyncGraphRunStream:
|
||||
experimental and may change.
|
||||
"""
|
||||
|
||||
# Native projections always registered by `astream_events(version="v3")`.
|
||||
# Attached dynamically by the `setattr` loop in `__init__`; declared here
|
||||
# so type checkers see them. Opt-in native projections (`updates`,
|
||||
# `custom`, `checkpoints`, `debug`, `tasks`) are only present when their
|
||||
# transformer is registered, so they are reached via `extensions[...]`.
|
||||
# Native projections, attached dynamically by the `setattr` loop in
|
||||
# `__init__` and declared here so type checkers see them.
|
||||
#
|
||||
# Always registered by `astream_events(version="v3")`:
|
||||
values: StreamChannel[dict[str, Any]]
|
||||
messages: StreamChannel[AsyncChatModelStream]
|
||||
lifecycle: StreamChannel[LifecyclePayload]
|
||||
subgraphs: StreamChannel[AsyncSubgraphRunStream]
|
||||
# Registered on demand; see `GraphRunStream`.
|
||||
updates: StreamChannel[dict[str, Any]]
|
||||
custom: StreamChannel[Any]
|
||||
checkpoints: StreamChannel[dict[str, Any]]
|
||||
debug: StreamChannel[dict[str, Any]]
|
||||
tasks: StreamChannel[dict[str, Any]]
|
||||
|
||||
def __getattr__(self, name: str) -> StreamChannel[Any]:
|
||||
"""Type projections declared elsewhere. See `GraphRunStream`."""
|
||||
_raise_missing_projection(self, name)
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
|
||||
@@ -70,6 +70,8 @@ __all__ = (
|
||||
"RetryPolicy",
|
||||
"TimeoutPolicy",
|
||||
"CachePolicy",
|
||||
"TracePolicy",
|
||||
"omit_payload",
|
||||
"Interrupt",
|
||||
"StateUpdate",
|
||||
"PregelTask",
|
||||
@@ -527,6 +529,44 @@ class CachePolicy(Generic[KeyFuncT]):
|
||||
"""Time to live for the cache entry in seconds. If `None`, the entry never expires."""
|
||||
|
||||
|
||||
@dataclass(**_DC_KWARGS)
|
||||
class TracePolicy:
|
||||
"""Configuration for how a node's run is traced.
|
||||
|
||||
Scope: this only transforms what the node's *own* run records. Child runs created
|
||||
by a traced `bound` runnable and the root graph run are not affected. Plain
|
||||
function nodes are traced with `trace=False`, so they have no such child runs.
|
||||
|
||||
Not intended to redact secrets. To redact inputs/outputs across all runs
|
||||
(children included), use the LangSmith client's
|
||||
`hide_inputs`/`hide_outputs`/`anonymizer` instead.
|
||||
|
||||
Each processor receives the node's raw input/output value (not a normalized
|
||||
kwargs dict) and returns the value to record.
|
||||
"""
|
||||
|
||||
process_inputs: Callable[[Any], Any] | None = None
|
||||
"""Optional callable to transform the node's input before it is recorded on the
|
||||
node's trace run. Can be used to omit or summarize large payloads
|
||||
(e.g. message history). Not intended to affect the value passed to the node; avoid
|
||||
mutating arguments in place."""
|
||||
|
||||
process_outputs: Callable[[Any], Any] | None = None
|
||||
"""Optional callable to transform the node's output before it is recorded on the
|
||||
node's trace run. Can be used to omit or summarize large payloads
|
||||
(e.g. message history). Not intended to affect the value returned by the node; avoid
|
||||
mutating arguments in place."""
|
||||
|
||||
|
||||
def omit_payload(_value: Any) -> dict[str, Any]:
|
||||
"""`TracePolicy` helper that records an empty payload, dropping the value entirely.
|
||||
|
||||
Use as `process_inputs` and/or `process_outputs` on a `TracePolicy` to keep a node's
|
||||
span and its timing while omitting its inputs/outputs from the trace.
|
||||
"""
|
||||
return {}
|
||||
|
||||
|
||||
_DEFAULT_INTERRUPT_ID = "placeholder-id"
|
||||
|
||||
|
||||
|
||||
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "langgraph"
|
||||
version = "1.2.10"
|
||||
version = "1.2.11"
|
||||
description = "Building stateful, multi-actor applications with LLMs"
|
||||
authors = []
|
||||
requires-python = ">=3.10"
|
||||
|
||||
@@ -616,3 +616,117 @@ async def test_add_messages_to_delta_migration_preserves_message_history_async()
|
||||
assert [m.id for m in snap.values["messages"]] == ["h1", "a1"], (
|
||||
f"async tip hydration mismatch: got {[m.id for m in snap.values['messages']]}"
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 8. First post-migration write, read back cold (regression for #8384)
|
||||
#
|
||||
# The migration boundary produces a checkpoint that carries BOTH a pre-delta
|
||||
# plain-value blob AND the pending write that produced its (delta-era) child.
|
||||
# That write is not subsumed by the blob — the blob is the value ENTERING that
|
||||
# checkpoint. A saver whose ancestor walk skips the seed checkpoint's own
|
||||
# writes silently drops the first post-migration write.
|
||||
#
|
||||
# The failure is invisible to the live `invoke` return value (computed
|
||||
# in-memory before persistence), so these tests must assert on a COLD read.
|
||||
# It is also invisible at `snapshot_frequency=1`, where every write is its own
|
||||
# snapshot boundary and the walk never terminates on a plain value — hence the
|
||||
# explicit default-frequency coverage.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_first_post_migration_write_survives_cold_read() -> None:
|
||||
"""One non-snapshotting write after migrating a thread to `DeltaChannel`
|
||||
must still be present when the state is read back from the checkpointer.
|
||||
|
||||
Regression for #8384: `invoke` returned the correct value while
|
||||
`get_state` dropped the write permanently.
|
||||
"""
|
||||
|
||||
checkpointer = InMemorySaver()
|
||||
config = {"configurable": {"thread_id": "first-post-migration"}}
|
||||
|
||||
binop = _binop_graph(checkpointer)
|
||||
binop.invoke({"items": ["a"]}, config)
|
||||
|
||||
delta = _delta_graph(checkpointer)
|
||||
live = delta.invoke({"items": ["b"]}, config)
|
||||
assert list(live["items"]) == ["a", "b"], "live invoke lost the write"
|
||||
|
||||
cold = delta.get_state(config)
|
||||
assert list(cold.values["items"]) == ["a", "b"], (
|
||||
"first post-migration write dropped on cold read: "
|
||||
f"got {list(cold.values['items'])}"
|
||||
)
|
||||
|
||||
|
||||
async def test_first_post_migration_write_survives_cold_read_async() -> None:
|
||||
"""Async variant of the #8384 regression."""
|
||||
|
||||
checkpointer = InMemorySaver()
|
||||
config = {"configurable": {"thread_id": "first-post-migration-async"}}
|
||||
|
||||
binop = _binop_graph(checkpointer)
|
||||
await binop.ainvoke({"items": ["a"]}, config)
|
||||
|
||||
delta = _delta_graph(checkpointer)
|
||||
live = await delta.ainvoke({"items": ["b"]}, config)
|
||||
assert list(live["items"]) == ["a", "b"], "live ainvoke lost the write"
|
||||
|
||||
cold = await delta.aget_state(config)
|
||||
assert list(cold.values["items"]) == ["a", "b"], (
|
||||
"first post-migration write dropped on cold read: "
|
||||
f"got {list(cold.values['items'])}"
|
||||
)
|
||||
|
||||
|
||||
def test_post_migration_writes_match_base_saver_fallback() -> None:
|
||||
"""Parity across the migration boundary WITH post-migration writes.
|
||||
|
||||
`test_base_saver_fallback_matches_optimized_override` only reads a
|
||||
pre-migration chain, so the optimized override and the reference walk
|
||||
never disagree there. Driving writes after the migration is what
|
||||
separates them.
|
||||
"""
|
||||
|
||||
def _run(saver: Any, thread: str) -> list[tuple[Any, list]]:
|
||||
config = {"configurable": {"thread_id": thread}}
|
||||
_drive(_binop_graph(saver), config, "u", 2)
|
||||
delta = _delta_graph(saver)
|
||||
_drive(delta, config, "d", 3)
|
||||
return [
|
||||
(s.next, list(s.values.get("items", [])))
|
||||
for s in delta.get_state_history(config)
|
||||
]
|
||||
|
||||
fast = _run(InMemorySaver(), "fast")
|
||||
slow = _run(_ThirdPartyStyleSaver(), "slow")
|
||||
|
||||
assert fast == slow, (
|
||||
"optimized override diverges from the base-saver fallback once "
|
||||
f"post-migration writes exist; fast={fast}, slow={slow}"
|
||||
)
|
||||
# Guard the assertion above against both paths being wrong in the same way.
|
||||
assert fast[0][1] == ["u0", "u1", "d0", "d1", "d2"], (
|
||||
f"unexpected accumulated state: {fast[0][1]}"
|
||||
)
|
||||
|
||||
|
||||
def test_add_messages_migration_keeps_first_post_migration_message() -> None:
|
||||
"""The `add_messages` -> `DeltaChannel` path is the one Deep Agents takes;
|
||||
dropping the first post-migration write loses a real user message.
|
||||
"""
|
||||
|
||||
checkpointer = InMemorySaver()
|
||||
config = {"configurable": {"thread_id": "add-messages-first-write"}}
|
||||
|
||||
pre_graph = _add_messages_graph(checkpointer)
|
||||
pre_graph.invoke({"messages": [HumanMessage(content="hello", id="h1")]}, config)
|
||||
|
||||
delta_graph = _delta_messages_graph(checkpointer)
|
||||
delta_graph.invoke({"messages": [HumanMessage(content="second", id="h2")]}, config)
|
||||
|
||||
ids = [m.id for m in delta_graph.get_state(config).values["messages"]]
|
||||
# h1 is the pre-migration seed, h2 the write that was being dropped; both
|
||||
# have to survive, and in order.
|
||||
assert ids == ["h1", "h2"], f"expected ['h1', 'h2'], got {ids}"
|
||||
|
||||
@@ -6,6 +6,7 @@ Type-narrowing is validated via `assert_type` calls in `_check_type_narrowing`.
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import copy
|
||||
import operator
|
||||
import sys
|
||||
from dataclasses import dataclass
|
||||
@@ -34,8 +35,10 @@ from langgraph.stream import (
|
||||
GraphRunStream,
|
||||
LifecyclePayload,
|
||||
StreamChannel,
|
||||
StreamTransformer,
|
||||
SubgraphRunStream,
|
||||
)
|
||||
from langgraph.stream._types import ProtocolEvent
|
||||
from langgraph.types import (
|
||||
CheckpointPayload,
|
||||
CheckpointStreamPart,
|
||||
@@ -1199,6 +1202,29 @@ def _check_type_narrowing(part: StreamPart[_StateT, _OutputT]) -> None:
|
||||
# type and the always-registered native projections.
|
||||
|
||||
|
||||
class _MarkerTransformer(StreamTransformer):
|
||||
"""Native transformer contributing a key this module doesn't declare.
|
||||
|
||||
Stands in for any transformer defined outside this package — projections
|
||||
whose names `GraphRunStream` can't enumerate, so they resolve through
|
||||
`__getattr__` instead of a class annotation.
|
||||
"""
|
||||
|
||||
_native = True
|
||||
|
||||
def __init__(self, scope: tuple[str, ...] = ()) -> None:
|
||||
super().__init__(scope)
|
||||
self._log: StreamChannel[str] = StreamChannel()
|
||||
|
||||
def init(self) -> dict[str, Any]:
|
||||
return {"marker": self._log}
|
||||
|
||||
def process(self, event: ProtocolEvent) -> bool:
|
||||
if event["method"] == "values":
|
||||
self._log.push("saw_values")
|
||||
return True
|
||||
|
||||
|
||||
def _check_stream_events_v3_typing() -> None:
|
||||
"""Compile-time checks for sync v3 typing — never called at runtime."""
|
||||
graph = _make_simple_graph().compile()
|
||||
@@ -1208,6 +1234,16 @@ def _check_stream_events_v3_typing() -> None:
|
||||
assert_type(run.messages, StreamChannel[ChatModelStream])
|
||||
assert_type(run.lifecycle, StreamChannel[LifecyclePayload])
|
||||
assert_type(run.subgraphs, StreamChannel[SubgraphRunStream])
|
||||
# Opt-in projections from transformers this package ships carry their real
|
||||
# item type even though they are only present once registered.
|
||||
assert_type(run.updates, StreamChannel[dict[str, Any]])
|
||||
assert_type(run.custom, StreamChannel[Any])
|
||||
assert_type(run.checkpoints, StreamChannel[dict[str, Any]])
|
||||
assert_type(run.debug, StreamChannel[dict[str, Any]])
|
||||
assert_type(run.tasks, StreamChannel[dict[str, Any]])
|
||||
# Projections this module can't enumerate resolve through `__getattr__`
|
||||
# as `StreamChannel[Any]` rather than failing with attr-defined.
|
||||
assert_type(run.marker, StreamChannel[Any])
|
||||
|
||||
|
||||
async def _check_astream_events_v3_typing() -> None:
|
||||
@@ -1219,3 +1255,66 @@ async def _check_astream_events_v3_typing() -> None:
|
||||
assert_type(run.messages, StreamChannel[AsyncChatModelStream])
|
||||
assert_type(run.lifecycle, StreamChannel[LifecyclePayload])
|
||||
assert_type(run.subgraphs, StreamChannel[AsyncSubgraphRunStream])
|
||||
assert_type(run.updates, StreamChannel[dict[str, Any]])
|
||||
assert_type(run.custom, StreamChannel[Any])
|
||||
assert_type(run.checkpoints, StreamChannel[dict[str, Any]])
|
||||
assert_type(run.debug, StreamChannel[dict[str, Any]])
|
||||
assert_type(run.tasks, StreamChannel[dict[str, Any]])
|
||||
assert_type(run.marker, StreamChannel[Any])
|
||||
|
||||
|
||||
def test_undeclared_native_projection_is_attached() -> None:
|
||||
"""A native projection this module doesn't declare still works at runtime.
|
||||
|
||||
`__getattr__` is a type-checker fallback only — it must not shadow the
|
||||
`setattr` loop that attaches registered native projections.
|
||||
"""
|
||||
graph = _make_simple_graph().compile()
|
||||
run = graph.stream_events(
|
||||
_SIMPLE_INPUT, version="v3", transformers=[_MarkerTransformer]
|
||||
)
|
||||
|
||||
marker_iter = iter(run.marker)
|
||||
assert run.output is not None
|
||||
assert run.marker is run.extensions["marker"]
|
||||
assert "saw_values" in list(marker_iter)
|
||||
|
||||
|
||||
def test_unregistered_projection_raises_attribute_error() -> None:
|
||||
"""An unregistered projection name still fails at runtime.
|
||||
|
||||
The `__getattr__` fallback exists to satisfy type checkers; it must not
|
||||
make unknown names resolve to anything. The message lists what *is*
|
||||
registered so a typo is diagnosable from the traceback alone.
|
||||
"""
|
||||
graph = _make_simple_graph().compile()
|
||||
run = graph.stream_events(_SIMPLE_INPUT, version="v3")
|
||||
|
||||
# `marker` type-checks via `__getattr__` but was never registered here.
|
||||
with pytest.raises(AttributeError) as exc_info:
|
||||
run.marker
|
||||
|
||||
message = str(exc_info.value)
|
||||
assert "marker" in message
|
||||
# The always-registered natives are listed as the alternatives.
|
||||
assert "messages" in message
|
||||
|
||||
# Registered projections still resolve, and the run is unaffected.
|
||||
assert isinstance(run.messages, StreamChannel)
|
||||
assert run.output is not None
|
||||
|
||||
|
||||
def test_getattr_fallback_does_not_recurse_before_init() -> None:
|
||||
"""`__getattr__` reads the mux from `__dict__`, so it is safe pre-init.
|
||||
|
||||
`self._mux` would re-enter `__getattr__` and overflow the stack when the
|
||||
attribute is missing, which is reachable via `hasattr` / `copy` / pickle
|
||||
probing on a partially constructed instance.
|
||||
"""
|
||||
bare = GraphRunStream.__new__(GraphRunStream)
|
||||
|
||||
with pytest.raises(AttributeError, match="anything"):
|
||||
bare.anything
|
||||
|
||||
assert hasattr(bare, "anything") is False
|
||||
assert copy.copy(bare) is not None
|
||||
|
||||
@@ -0,0 +1,286 @@
|
||||
"""Tests for subgraph auto-detection (`pregel/_utils.py`).
|
||||
|
||||
Detection failing is silent — the graph still runs, only introspection goes
|
||||
quiet — so every shape a node can hold a graph in is pinned here. The expected
|
||||
values are what the source-parsing implementation this replaced produced for
|
||||
the same shapes, except for `sourceless`, whose source it could not read,
|
||||
`empty_closure_cell`, on which it raised, and `unreachable_attribute_chain`,
|
||||
where it reported a graph that dropped code could never invoke.
|
||||
"""
|
||||
|
||||
import functools
|
||||
import operator
|
||||
from typing import Annotated, Any
|
||||
|
||||
import pytest
|
||||
from typing_extensions import TypedDict
|
||||
|
||||
from langgraph.graph import END, START, StateGraph
|
||||
from langgraph.pregel._utils import get_function_nonlocals
|
||||
|
||||
|
||||
class State(TypedDict):
|
||||
log: Annotated[list, operator.add]
|
||||
|
||||
|
||||
def _leaf(tag: str) -> Any:
|
||||
"""Return a compiled graph that reports itself as `tag`."""
|
||||
builder = StateGraph(State)
|
||||
builder.add_node(tag, lambda s: {"log": [tag]})
|
||||
builder.add_edge(START, tag)
|
||||
builder.add_edge(tag, END)
|
||||
compiled = builder.compile()
|
||||
compiled.name = tag
|
||||
return compiled
|
||||
|
||||
|
||||
def _detect(node: Any) -> str | None:
|
||||
"""Return the name of the subgraph detected for `node`, or None."""
|
||||
builder = StateGraph(State)
|
||||
builder.add_node("n", node)
|
||||
builder.add_edge(START, "n")
|
||||
builder.add_edge("n", END)
|
||||
subgraphs = builder.compile().nodes["n"].subgraphs
|
||||
return getattr(subgraphs[0], "name", "?") if subgraphs else None
|
||||
|
||||
|
||||
class _Box:
|
||||
def __init__(self, payload: Any) -> None:
|
||||
self.payload = payload
|
||||
|
||||
|
||||
class _ListSubclass(list):
|
||||
pass
|
||||
|
||||
|
||||
class _MethodHolder:
|
||||
def __init__(self) -> None:
|
||||
self.graph = _leaf("via_self")
|
||||
|
||||
def as_node(self, state: State) -> Any:
|
||||
return self.graph.invoke(state)
|
||||
|
||||
|
||||
MODULE_GRAPH = _leaf("module_global")
|
||||
CHAIN = _Box(_Box(_leaf("attr_chain")))
|
||||
GRAPH_IN_PLAIN_LIST = [_leaf("in_list")]
|
||||
METHOD_HOLDER = _MethodHolder()
|
||||
|
||||
|
||||
def closure_capture() -> Any:
|
||||
sub = _leaf("closure")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return sub.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def module_global() -> Any:
|
||||
def node(state: State) -> Any:
|
||||
return MODULE_GRAPH.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def attribute_chain() -> Any:
|
||||
def node(state: State) -> Any:
|
||||
return CHAIN.payload.payload.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def nested_def_captured_attribute() -> Any:
|
||||
"""A chain on a captured holder, named only inside a nested code object.
|
||||
|
||||
The captured value is the holder, not the graph, so the chain itself has to
|
||||
be recovered from the nested scope.
|
||||
"""
|
||||
holder = _Box(_leaf("nested_captured"))
|
||||
|
||||
def node(state: State) -> Any:
|
||||
def inner() -> Any:
|
||||
return holder.payload.invoke(state)
|
||||
|
||||
return inner()
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def unreachable_branch() -> Any:
|
||||
"""A captured graph referenced only from code the compiler removes."""
|
||||
sub = _leaf("unreachable")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
if False:
|
||||
sub.invoke(state)
|
||||
return {"log": []}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def unreachable_attribute_chain() -> Any:
|
||||
"""A graph named only along an attribute path the compiler dropped.
|
||||
|
||||
The closure keeps `holder`, but the `.payload` load is gone. The source
|
||||
parser reported this one; dropped code cannot invoke anything, so that was
|
||||
a phantom rather than a detection.
|
||||
"""
|
||||
holder = _Box(_leaf("unreachable_attr"))
|
||||
|
||||
def node(state: State) -> Any:
|
||||
if False:
|
||||
holder.payload.invoke(state)
|
||||
return {"log": []}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def wrapper_referencing_nothing() -> Any:
|
||||
"""A wrapper whose own scope holds nothing, so only `__wrapped__` leads on.
|
||||
|
||||
`functools.wraps` would leave the wrapper closing over the inner function;
|
||||
setting the attribute by hand does not.
|
||||
"""
|
||||
sub = _leaf("via_wrapped")
|
||||
|
||||
def inner(state: State) -> Any:
|
||||
return sub.invoke(state)
|
||||
|
||||
def wrapper(state: State) -> Any:
|
||||
return {"log": []}
|
||||
|
||||
wrapper.__wrapped__ = inner
|
||||
return wrapper
|
||||
|
||||
|
||||
def captured_list_subclass() -> Any:
|
||||
"""A `list` subclass is not a leaf: it can carry a graph as an attribute."""
|
||||
holder = _ListSubclass()
|
||||
holder.payload = _leaf("list_subclass")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return holder.payload.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def empty_closure_cell() -> Any:
|
||||
"""An unassigned closure variable leaves a cell that cannot be read."""
|
||||
sub = _leaf("beside_empty_cell")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return unassigned, sub.invoke(state)
|
||||
|
||||
return node
|
||||
unassigned = 1 # never runs, so the cell it creates is never filled
|
||||
|
||||
|
||||
def sourceless() -> Any:
|
||||
"""A node compiled without a source file, which `getsource` could not read."""
|
||||
namespace: dict[str, Any] = {"SOURCELESS": _leaf("sourceless")}
|
||||
exec(
|
||||
compile(
|
||||
"def node(state):\n return SOURCELESS.invoke(state)", "<test>", "exec"
|
||||
),
|
||||
namespace,
|
||||
)
|
||||
return namespace["node"]
|
||||
|
||||
|
||||
async def _async_node(state: State) -> Any:
|
||||
return await MODULE_GRAPH.ainvoke(state)
|
||||
|
||||
|
||||
def async_node() -> Any:
|
||||
return _async_node
|
||||
|
||||
|
||||
def no_subgraph() -> Any:
|
||||
"""Nothing but leaf values in reach, so the bytecode walk is skipped."""
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return {"log": [len("abc") + 1]}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def recombined_names() -> Any:
|
||||
"""Loads `CHAIN.payload` and `local.payload`, never `CHAIN.payload.payload`."""
|
||||
|
||||
def node(state: State) -> Any:
|
||||
local = _Box("not a graph")
|
||||
return {"log": [CHAIN.payload, local.payload]}
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def broken_attribute_chain() -> Any:
|
||||
holder = _Box("a string, so `.payload.missing` cannot resolve")
|
||||
|
||||
def node(state: State) -> Any:
|
||||
return holder.payload.missing.invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def nested_def_global() -> Any:
|
||||
"""A global named only in a nested code object: out of reach, as before."""
|
||||
|
||||
def node(state: State) -> Any:
|
||||
def inner() -> Any:
|
||||
return MODULE_GRAPH.invoke(state)
|
||||
|
||||
return inner()
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def graph_in_plain_list() -> Any:
|
||||
def node(state: State) -> Any:
|
||||
return GRAPH_IN_PLAIN_LIST[0].invoke(state)
|
||||
|
||||
return node
|
||||
|
||||
|
||||
def bound_method_self() -> Any:
|
||||
return METHOD_HOLDER.as_node
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("factory", "expected"),
|
||||
[
|
||||
(closure_capture, "closure"),
|
||||
(module_global, "module_global"),
|
||||
(attribute_chain, "attr_chain"),
|
||||
(nested_def_captured_attribute, "nested_captured"),
|
||||
(unreachable_branch, "unreachable"),
|
||||
(wrapper_referencing_nothing, "via_wrapped"),
|
||||
(captured_list_subclass, "list_subclass"),
|
||||
(empty_closure_cell, "beside_empty_cell"),
|
||||
(sourceless, "sourceless"),
|
||||
(async_node, "module_global"),
|
||||
# Shapes no reference chain reaches: a subscript, an instance attribute
|
||||
# of `self`, a global named only in a nested scope, and an attribute
|
||||
# path the compiler dropped.
|
||||
(no_subgraph, None),
|
||||
(recombined_names, None),
|
||||
(broken_attribute_chain, None),
|
||||
(nested_def_global, None),
|
||||
(graph_in_plain_list, None),
|
||||
(bound_method_self, None),
|
||||
(unreachable_attribute_chain, None),
|
||||
],
|
||||
ids=lambda value: value.__name__ if callable(value) else str(value),
|
||||
)
|
||||
def test_subgraph_detection(factory: Any, expected: str | None) -> None:
|
||||
assert _detect(factory()) == expected
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"candidate",
|
||||
[functools.partial(lambda state, extra: {"log": [extra]}, extra="x"), len],
|
||||
ids=["partial", "builtin"],
|
||||
)
|
||||
def test_callables_without_a_code_object_are_handled(candidate: Any) -> None:
|
||||
assert get_function_nonlocals(candidate) == []
|
||||
@@ -0,0 +1,149 @@
|
||||
"""End-to-end tests for `TracePolicy` input processing on node trace runs."""
|
||||
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
from typing_extensions import TypedDict
|
||||
|
||||
from langgraph.graph import END, START, StateGraph
|
||||
from langgraph.types import TracePolicy
|
||||
from tests.fake_tracer import FakeTracer, Run
|
||||
|
||||
|
||||
class State(TypedDict):
|
||||
value: int
|
||||
|
||||
|
||||
def _node_run(tracer: FakeTracer, name: str) -> Run:
|
||||
return next(r for r in tracer.flattened_runs() if r.name == name)
|
||||
|
||||
|
||||
def _incr(state: State) -> State:
|
||||
return {"value": state["value"] + 1}
|
||||
|
||||
|
||||
def test_trace_policy_transforms_recorded_inputs() -> None:
|
||||
seen: dict[str, Any] = {}
|
||||
|
||||
def process_inputs(inp: Any) -> Any:
|
||||
seen["inputs"] = inp
|
||||
return {"scrubbed_in": True}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("n", _incr, trace_policy=TracePolicy(process_inputs=process_inputs))
|
||||
.add_edge(START, "n")
|
||||
.add_edge("n", END)
|
||||
.compile()
|
||||
)
|
||||
|
||||
tracer = FakeTracer()
|
||||
# the real graph output is unaffected by the trace policy
|
||||
assert graph.invoke({"value": 1}, {"callbacks": [tracer]}) == {"value": 2}
|
||||
|
||||
run = _node_run(tracer, "n")
|
||||
# the recorded input is transformed; the output is recorded as-is
|
||||
assert run.inputs == {"scrubbed_in": True}
|
||||
assert run.outputs == {"value": 2}
|
||||
# process_inputs observed the real, untransformed input
|
||||
assert seen["inputs"] == {"value": 1}
|
||||
|
||||
|
||||
def test_trace_policy_transforms_recorded_outputs() -> None:
|
||||
seen: dict[str, Any] = {}
|
||||
|
||||
def process_outputs(out: Any) -> Any:
|
||||
seen["outputs"] = out
|
||||
return {"scrubbed_out": True}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("n", _incr, trace_policy=TracePolicy(process_outputs=process_outputs))
|
||||
.add_edge(START, "n")
|
||||
.add_edge("n", END)
|
||||
.compile()
|
||||
)
|
||||
|
||||
tracer = FakeTracer()
|
||||
# the real graph output is unaffected by the trace policy
|
||||
assert graph.invoke({"value": 1}, {"callbacks": [tracer]}) == {"value": 2}
|
||||
|
||||
run = _node_run(tracer, "n")
|
||||
# the recorded output is transformed; the input is recorded as-is
|
||||
assert run.inputs == {"value": 1}
|
||||
assert run.outputs == {"scrubbed_out": True}
|
||||
# process_outputs observed the real, untransformed output
|
||||
assert seen["outputs"] == {"value": 2}
|
||||
|
||||
|
||||
def test_trace_policy_none_records_real_payloads() -> None:
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("n", _incr)
|
||||
.add_edge(START, "n")
|
||||
.add_edge("n", END)
|
||||
.compile()
|
||||
)
|
||||
|
||||
tracer = FakeTracer()
|
||||
assert graph.invoke({"value": 1}, {"callbacks": [tracer]}) == {"value": 2}
|
||||
|
||||
run = _node_run(tracer, "n")
|
||||
assert run.inputs == {"value": 1}
|
||||
assert run.outputs == {"value": 2}
|
||||
|
||||
|
||||
def test_trace_policy_processor_error_safe_without_callbacks() -> None:
|
||||
def boom(_inp: Any) -> Any:
|
||||
raise RuntimeError("processor failed")
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("n", _incr, trace_policy=TracePolicy(process_inputs=boom))
|
||||
.add_edge(START, "n")
|
||||
.add_edge("n", END)
|
||||
.compile()
|
||||
)
|
||||
|
||||
# no callbacks: the processor still runs but is fail-open, so execution is unaffected
|
||||
assert graph.invoke({"value": 1}) == {"value": 2}
|
||||
|
||||
|
||||
def test_trace_policy_processor_error_does_not_break_execution() -> None:
|
||||
def boom(_inp: Any) -> Any:
|
||||
raise RuntimeError("processor failed")
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("n", _incr, trace_policy=TracePolicy(process_inputs=boom))
|
||||
.add_edge(START, "n")
|
||||
.add_edge("n", END)
|
||||
.compile()
|
||||
)
|
||||
|
||||
tracer = FakeTracer()
|
||||
# a raising processor must not abort the node; the untransformed input is recorded
|
||||
assert graph.invoke({"value": 1}, {"callbacks": [tracer]}) == {"value": 2}
|
||||
assert _node_run(tracer, "n").inputs == {"value": 1}
|
||||
|
||||
|
||||
@pytest.mark.anyio
|
||||
async def test_trace_policy_transforms_recorded_inputs_async() -> None:
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node(
|
||||
"n",
|
||||
_incr,
|
||||
trace_policy=TracePolicy(process_inputs=lambda _: {"scrubbed_in": True}),
|
||||
)
|
||||
.add_edge(START, "n")
|
||||
.add_edge("n", END)
|
||||
.compile()
|
||||
)
|
||||
|
||||
tracer = FakeTracer()
|
||||
assert await graph.ainvoke({"value": 5}, {"callbacks": [tracer]}) == {"value": 6}
|
||||
|
||||
run = _node_run(tracer, "n")
|
||||
assert run.inputs == {"scrubbed_in": True}
|
||||
assert run.outputs == {"value": 6}
|
||||
Generated
+2053
-1526
File diff suppressed because it is too large
Load Diff
+17
-4
@@ -37,6 +37,7 @@ uv add langchain-anthropic
|
||||
from langchain_anthropic import ChatAnthropic
|
||||
from langgraph.prebuilt import create_react_agent
|
||||
|
||||
|
||||
# Define the tools for the agent to use
|
||||
def search(query: str):
|
||||
"""Call to surf the web."""
|
||||
@@ -45,6 +46,7 @@ def search(query: str):
|
||||
return "It's 60 degrees and foggy."
|
||||
return "It's 90 degrees and sunny."
|
||||
|
||||
|
||||
tools = [search]
|
||||
model = ChatAnthropic(model="claude-3-7-sonnet-latest")
|
||||
|
||||
@@ -65,6 +67,7 @@ app.invoke(
|
||||
from langgraph.prebuilt import ToolNode
|
||||
from langchain_core.messages import AIMessage
|
||||
|
||||
|
||||
def search(query: str):
|
||||
"""Call to surf the web."""
|
||||
# This is a placeholder, but don't tell the LLM that...
|
||||
@@ -72,8 +75,11 @@ def search(query: str):
|
||||
return "It's 60 degrees and foggy."
|
||||
return "It's 90 degrees and sunny."
|
||||
|
||||
|
||||
tool_node = ToolNode([search])
|
||||
tool_calls = [{"name": "search", "args": {"query": "what is the weather in sf"}, "id": "1"}]
|
||||
tool_calls = [
|
||||
{"name": "search", "args": {"query": "what is the weather in sf"}, "id": "1"}
|
||||
]
|
||||
ai_message = AIMessage(content="", tool_calls=tool_calls)
|
||||
# execute tool call
|
||||
tool_node.invoke({"messages": [ai_message]})
|
||||
@@ -98,10 +104,17 @@ class SelectNumber(BaseModel):
|
||||
raise ValueError("Only 37 is allowed")
|
||||
return v
|
||||
|
||||
|
||||
validation_node = ValidationNode([SelectNumber])
|
||||
validation_node.invoke({
|
||||
"messages": [AIMessage("", tool_calls=[{"name": "SelectNumber", "args": {"a": 42}, "id": "1"}])]
|
||||
})
|
||||
validation_node.invoke(
|
||||
{
|
||||
"messages": [
|
||||
AIMessage(
|
||||
"", tool_calls=[{"name": "SelectNumber", "args": {"a": 42}, "id": "1"}]
|
||||
)
|
||||
]
|
||||
}
|
||||
)
|
||||
```
|
||||
|
||||
## Agent Inbox
|
||||
|
||||
Generated
+1002
-698
File diff suppressed because it is too large
Load Diff
@@ -39,6 +39,15 @@
|
||||
- `client.threads.stream()` now accepts `transport="sse"` (default) or
|
||||
`transport="websocket"` in place of the previous transport-agnostic default.
|
||||
|
||||
### Fixed
|
||||
|
||||
- Resource-scoped auth decorators now honor `actions=` and reject empty or
|
||||
invalid action lists. Because unmatched custom-auth paths remain allowed,
|
||||
deployments using action-scoped handlers should configure a global
|
||||
default-deny handler; `langgraph-api` 0.10+ warns about uncovered paths at
|
||||
startup. Resource-specific decorators retain matching `resources=` selectors
|
||||
for backward compatibility; use `@auth.on(resources=...)` for other resources.
|
||||
|
||||
### Notes
|
||||
|
||||
- The v3 streaming surface (`AsyncThreadStream`, `SyncThreadStream`, and all
|
||||
|
||||
@@ -87,7 +87,9 @@ async with client.threads.stream(assistant_id="agent") as thread:
|
||||
|
||||
```python
|
||||
async with client.threads.stream(assistant_id="agent") as thread:
|
||||
await thread.run.start(input={"messages": [{"role": "user", "content": "book a flight"}]})
|
||||
await thread.run.start(
|
||||
input={"messages": [{"role": "user", "content": "book a flight"}]}
|
||||
)
|
||||
|
||||
# Wait for the run to pause at an interrupt node.
|
||||
# thread.interrupted becomes True when input.requested arrives.
|
||||
|
||||
@@ -43,7 +43,9 @@ thread = await client.threads.create()
|
||||
|
||||
# Start a streaming run
|
||||
input = {"messages": [{"role": "human", "content": "what's the weather in la"}]}
|
||||
async for chunk in client.runs.stream(thread['thread_id'], agent['assistant_id'], input=input):
|
||||
async for chunk in client.runs.stream(
|
||||
thread["thread_id"], agent["assistant_id"], input=input
|
||||
):
|
||||
print(chunk)
|
||||
```
|
||||
|
||||
@@ -80,9 +82,9 @@ async with client.threads.stream(
|
||||
messages, tool_calls = await asyncio.gather(get_messages(), get_tool_calls())
|
||||
|
||||
for stream in messages:
|
||||
print(await stream.text) # accumulated text
|
||||
print(await stream.text) # accumulated text
|
||||
|
||||
final = await thread.output # terminal state values
|
||||
final = await thread.output # terminal state values
|
||||
```
|
||||
|
||||
## 📕 Releases & Versioning
|
||||
|
||||
@@ -1,8 +1,15 @@
|
||||
from langgraph_sdk.auth import Auth
|
||||
from langgraph_sdk.client import get_client, get_sync_client
|
||||
from langgraph_sdk.encryption import Encryption
|
||||
from langgraph_sdk.encryption.types import EncryptionContext
|
||||
from langgraph_sdk.encryption.types import DecryptResult, EncryptionContext
|
||||
|
||||
__version__ = "0.4.2"
|
||||
__version__ = "0.4.4"
|
||||
|
||||
__all__ = ["Auth", "Encryption", "EncryptionContext", "get_client", "get_sync_client"]
|
||||
__all__ = [
|
||||
"Auth",
|
||||
"DecryptResult",
|
||||
"Encryption",
|
||||
"EncryptionContext",
|
||||
"get_client",
|
||||
"get_sync_client",
|
||||
]
|
||||
|
||||
@@ -24,7 +24,7 @@ from langchain_core.language_models.chat_model_stream import AsyncChatModelStrea
|
||||
from langchain_protocol import Event, SubscribeParams
|
||||
|
||||
from langgraph_sdk._async.http import HttpClient
|
||||
from langgraph_sdk.schema import QueryParamTypes
|
||||
from langgraph_sdk.schema import LangSmithTracing, QueryParamTypes
|
||||
from langgraph_sdk.stream.controller import _SeenEventIds
|
||||
from langgraph_sdk.stream.decoders import (
|
||||
DataDecoder,
|
||||
@@ -172,6 +172,7 @@ class RunModule:
|
||||
input: Any = None,
|
||||
config: dict[str, Any] | None = None,
|
||||
metadata: dict[str, Any] | None = None,
|
||||
langsmith_tracing: LangSmithTracing | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Send `run.start` to the server. Returns the result (`{"run_id": ...}`)."""
|
||||
params: dict[str, Any] = {"assistant_id": self._owner.assistant_id}
|
||||
@@ -181,6 +182,8 @@ class RunModule:
|
||||
params["config"] = config
|
||||
if metadata is not None:
|
||||
params["metadata"] = metadata
|
||||
if langsmith_tracing is not None:
|
||||
params["langsmith_tracer"] = langsmith_tracing
|
||||
loop = asyncio.get_running_loop()
|
||||
gate: asyncio.Future[None] = loop.create_future()
|
||||
self._owner._run_start_ready = gate
|
||||
|
||||
@@ -23,7 +23,7 @@ from langchain_core.language_models.chat_model_stream import ChatModelStream
|
||||
from langchain_protocol import Event, SubscribeParams
|
||||
|
||||
from langgraph_sdk._sync.http import SyncHttpClient
|
||||
from langgraph_sdk.schema import QueryParamTypes
|
||||
from langgraph_sdk.schema import LangSmithTracing, QueryParamTypes
|
||||
from langgraph_sdk.stream.decoders import (
|
||||
DataDecoder,
|
||||
Decoder,
|
||||
@@ -215,6 +215,7 @@ class SyncRunModule:
|
||||
input: Any = None,
|
||||
config: dict[str, Any] | None = None,
|
||||
metadata: dict[str, Any] | None = None,
|
||||
langsmith_tracing: LangSmithTracing | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Send `run.start` to the server. Returns the result (`{"run_id": ...}`)."""
|
||||
params: dict[str, Any] = {"assistant_id": self._owner.assistant_id}
|
||||
@@ -224,6 +225,8 @@ class SyncRunModule:
|
||||
params["config"] = config
|
||||
if metadata is not None:
|
||||
params["metadata"] = metadata
|
||||
if langsmith_tracing is not None:
|
||||
params["langsmith_tracer"] = langsmith_tracing
|
||||
result = self._owner._send_command("run.start", params)
|
||||
self._owner._run_seen = True
|
||||
controller = self._owner._controller
|
||||
|
||||
@@ -341,9 +341,15 @@ VUpdate = typing.TypeVar("VUpdate", covariant=True)
|
||||
VRead = typing.TypeVar("VRead", covariant=True)
|
||||
VDelete = typing.TypeVar("VDelete", covariant=True)
|
||||
VSearch = typing.TypeVar("VSearch", covariant=True)
|
||||
ResourceActionT = typing.TypeVar("ResourceActionT", bound=str)
|
||||
|
||||
_ResourceAction = typing.Literal["create", "read", "update", "delete", "search"]
|
||||
_ThreadAction = _ResourceAction | typing.Literal["create_run"]
|
||||
|
||||
|
||||
class _ResourceOn(typing.Generic[VCreate, VRead, VUpdate, VDelete, VSearch]):
|
||||
class _ResourceOn(
|
||||
typing.Generic[VCreate, VRead, VUpdate, VDelete, VSearch, ResourceActionT]
|
||||
):
|
||||
"""
|
||||
Generic base class for resource-specific handlers.
|
||||
"""
|
||||
@@ -392,8 +398,8 @@ class _ResourceOn(typing.Generic[VCreate, VRead, VUpdate, VDelete, VSearch]):
|
||||
def __call__(
|
||||
self,
|
||||
*,
|
||||
resources: str | Sequence[str],
|
||||
actions: str | Sequence[str] | None = None,
|
||||
resources: str | Sequence[str] | None = None,
|
||||
actions: ResourceActionT | Sequence[ResourceActionT] | None = None,
|
||||
) -> Callable[
|
||||
[_ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch]],
|
||||
_ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch],
|
||||
@@ -408,7 +414,7 @@ class _ResourceOn(typing.Generic[VCreate, VRead, VUpdate, VDelete, VSearch]):
|
||||
) = None,
|
||||
*,
|
||||
resources: str | Sequence[str] | None = None,
|
||||
actions: str | Sequence[str] | None = None,
|
||||
actions: ResourceActionT | Sequence[ResourceActionT] | None = None,
|
||||
) -> (
|
||||
_ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch]
|
||||
| Callable[
|
||||
@@ -416,24 +422,66 @@ class _ResourceOn(typing.Generic[VCreate, VRead, VUpdate, VDelete, VSearch]):
|
||||
_ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch],
|
||||
]
|
||||
):
|
||||
if fn is not None:
|
||||
_validate_handler(fn)
|
||||
return typing.cast(
|
||||
"_ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch]",
|
||||
_register_handler(self.auth, self.resource, "*", fn),
|
||||
)
|
||||
|
||||
def decorator(
|
||||
handler: _ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch],
|
||||
) -> _ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch]:
|
||||
_validate_handler(handler)
|
||||
return typing.cast(
|
||||
"_ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch]",
|
||||
_register_handler(self.auth, self.resource, "*", handler),
|
||||
if resources is None:
|
||||
resource_list = [self.resource]
|
||||
elif isinstance(resources, str):
|
||||
resource_list = [resources]
|
||||
elif isinstance(resources, Sequence):
|
||||
resource_list = list(resources)
|
||||
else:
|
||||
raise TypeError("resources must be a string or sequence of strings")
|
||||
if resource_list != [self.resource]:
|
||||
raise ValueError(
|
||||
f"Resource-specific decorator for {self.resource!r} cannot "
|
||||
f"register handlers for {resource_list!r}. Use @auth.on(...) "
|
||||
"for other or multiple resources."
|
||||
)
|
||||
if actions is None:
|
||||
action_list = ["*"]
|
||||
elif isinstance(actions, str):
|
||||
action_list = [actions]
|
||||
elif isinstance(actions, Sequence):
|
||||
action_list = list(actions)
|
||||
else:
|
||||
raise TypeError("actions must be a string or sequence of strings")
|
||||
if not action_list:
|
||||
raise ValueError("actions must not be empty")
|
||||
if not all(isinstance(action, str) for action in action_list):
|
||||
raise TypeError("actions must be a string or sequence of strings")
|
||||
valid_actions = {
|
||||
value.action
|
||||
for value in vars(self).values()
|
||||
if isinstance(value, _ResourceActionOn)
|
||||
}
|
||||
invalid_actions = (
|
||||
sorted(set(action_list) - valid_actions) if actions is not None else []
|
||||
)
|
||||
if invalid_actions:
|
||||
raise ValueError(
|
||||
f"Invalid action(s) for {self.resource}: {', '.join(invalid_actions)}"
|
||||
)
|
||||
if len(action_list) != len(set(action_list)):
|
||||
raise ValueError("actions must not contain duplicates")
|
||||
for action in action_list:
|
||||
if (self.resource, action) in self.auth._handlers:
|
||||
raise ValueError(
|
||||
f"types.Handler already set for {self.resource}, {action}."
|
||||
)
|
||||
for action in action_list:
|
||||
_register_handler(self.auth, self.resource, action, handler)
|
||||
return handler
|
||||
|
||||
# Accept keyword-only parameters for future filtering behavior; referenced to satisfy linters.
|
||||
_ = resources, actions
|
||||
if fn is not None:
|
||||
return decorator(
|
||||
typing.cast(
|
||||
"_ActionHandler[VCreate | VUpdate | VRead | VDelete | VSearch]",
|
||||
fn,
|
||||
)
|
||||
)
|
||||
return decorator
|
||||
|
||||
|
||||
@@ -444,6 +492,7 @@ class _AssistantsOn(
|
||||
types.AssistantsUpdate,
|
||||
types.AssistantsDelete,
|
||||
types.AssistantsSearch,
|
||||
_ResourceAction,
|
||||
]
|
||||
):
|
||||
value = (
|
||||
@@ -467,6 +516,7 @@ class _ThreadsOn(
|
||||
types.ThreadsUpdate,
|
||||
types.ThreadsDelete,
|
||||
types.ThreadsSearch,
|
||||
_ThreadAction,
|
||||
]
|
||||
):
|
||||
value = (
|
||||
@@ -502,6 +552,7 @@ class _CronsOn(
|
||||
types.CronsUpdate,
|
||||
types.CronsDelete,
|
||||
types.CronsSearch,
|
||||
_ResourceAction,
|
||||
]
|
||||
):
|
||||
value = type[
|
||||
|
||||
@@ -18,6 +18,9 @@ import warnings
|
||||
|
||||
from langgraph_sdk.encryption import types
|
||||
|
||||
_BlobDecryptorT = typing.TypeVar("_BlobDecryptorT", bound=types.BlobDecryptor)
|
||||
_JsonDecryptorT = typing.TypeVar("_JsonDecryptorT", bound=types.JsonDecryptor)
|
||||
|
||||
|
||||
class LangGraphBetaWarning(UserWarning):
|
||||
"""Warning for beta features in LangGraph SDK."""
|
||||
@@ -141,7 +144,7 @@ class _DecryptDecorators:
|
||||
def __init__(self, parent: Encryption):
|
||||
self._parent = parent
|
||||
|
||||
def blob(self, fn: types.BlobDecryptor) -> types.BlobDecryptor:
|
||||
def blob(self, fn: _BlobDecryptorT) -> _BlobDecryptorT:
|
||||
"""Register a blob decryption handler.
|
||||
|
||||
The handler will be called to decrypt opaque data like checkpoint blobs.
|
||||
@@ -149,7 +152,9 @@ class _DecryptDecorators:
|
||||
Example:
|
||||
```python
|
||||
@encryption.decrypt.blob
|
||||
async def decrypt_blob(ctx: EncryptionContext, blob: bytes) -> bytes:
|
||||
async def decrypt_blob(
|
||||
ctx: EncryptionContext, blob: bytes
|
||||
) -> bytes | DecryptResult[bytes]:
|
||||
# Decrypt the blob using your encryption service
|
||||
return decrypted_blob
|
||||
```
|
||||
@@ -170,13 +175,15 @@ class _DecryptDecorators:
|
||||
self._parent._blob_decryptor = fn
|
||||
return fn
|
||||
|
||||
def json(self, fn: types.JsonDecryptor) -> types.JsonDecryptor:
|
||||
def json(self, fn: _JsonDecryptorT) -> _JsonDecryptorT:
|
||||
"""Register the JSON decryption handler.
|
||||
|
||||
Example:
|
||||
```python
|
||||
@encryption.decrypt.json
|
||||
async def decrypt_json(ctx: EncryptionContext, data: dict) -> dict:
|
||||
async def decrypt_json(
|
||||
ctx: EncryptionContext, data: dict
|
||||
) -> dict | DecryptResult[dict]:
|
||||
# Decrypt the data
|
||||
return decrypt_data(data)
|
||||
```
|
||||
@@ -369,7 +376,7 @@ class Encryption:
|
||||
"""Reference to encryption type definitions.
|
||||
|
||||
Provides access to all type definitions used in the encryption system,
|
||||
including EncryptionContext, BlobEncryptor, BlobDecryptor,
|
||||
including EncryptionContext, DecryptResult, BlobEncryptor, BlobDecryptor,
|
||||
JsonEncryptor, and JsonDecryptor.
|
||||
"""
|
||||
|
||||
|
||||
@@ -9,10 +9,30 @@ from __future__ import annotations
|
||||
|
||||
import typing
|
||||
from collections.abc import Awaitable, Callable
|
||||
from dataclasses import dataclass
|
||||
|
||||
Json = dict[str, typing.Any]
|
||||
"""JSON-serializable dictionary type for structured data encryption."""
|
||||
|
||||
T = typing.TypeVar("T")
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class DecryptResult(typing.Generic[T]):
|
||||
"""Decrypted data and optional replacement ciphertext.
|
||||
|
||||
Return this from a decrypt handler when encrypted data should be replaced,
|
||||
such as after rotating its encryption key. Returning plaintext directly
|
||||
remains supported when no replacement is needed.
|
||||
|
||||
Attributes:
|
||||
plaintext: Decrypted data returned to the caller
|
||||
replacement: New encrypted data to persist in place of the input
|
||||
"""
|
||||
|
||||
plaintext: T
|
||||
replacement: T | None = None
|
||||
|
||||
|
||||
class EncryptionContext:
|
||||
"""Context passed to encryption/decryption handlers.
|
||||
@@ -57,7 +77,9 @@ Returns:
|
||||
Awaitable that resolves to encrypted bytes
|
||||
"""
|
||||
|
||||
BlobDecryptor = Callable[[EncryptionContext, bytes], Awaitable[bytes]]
|
||||
BlobDecryptor = Callable[
|
||||
[EncryptionContext, bytes], Awaitable[bytes | DecryptResult[bytes]]
|
||||
]
|
||||
"""Handler for decrypting opaque blob data like checkpoints.
|
||||
|
||||
Note: Must be an async function. Decryption typically involves I/O operations
|
||||
@@ -68,7 +90,8 @@ Args:
|
||||
blob: The encrypted bytes to decrypt
|
||||
|
||||
Returns:
|
||||
Awaitable that resolves to decrypted bytes
|
||||
Awaitable that resolves to decrypted bytes, or a DecryptResult containing
|
||||
decrypted bytes and replacement ciphertext
|
||||
"""
|
||||
|
||||
JsonEncryptor = Callable[[EncryptionContext, Json], Awaitable[Json]]
|
||||
@@ -101,7 +124,9 @@ Returns:
|
||||
Awaitable that resolves to encrypted JSON dictionary
|
||||
"""
|
||||
|
||||
JsonDecryptor = Callable[[EncryptionContext, Json], Awaitable[Json]]
|
||||
JsonDecryptor = Callable[
|
||||
[EncryptionContext, Json], Awaitable[Json | DecryptResult[Json]]
|
||||
]
|
||||
"""Handler for decrypting structured JSON data.
|
||||
|
||||
Note: Must be an async function. Decryption typically involves I/O operations
|
||||
@@ -115,7 +140,8 @@ Args:
|
||||
data: The encrypted JSON dictionary
|
||||
|
||||
Returns:
|
||||
Awaitable that resolves to decrypted JSON dictionary
|
||||
Awaitable that resolves to a decrypted JSON dictionary, or a DecryptResult
|
||||
containing decrypted JSON and replacement ciphertext
|
||||
"""
|
||||
|
||||
if typing.TYPE_CHECKING:
|
||||
|
||||
@@ -963,7 +963,7 @@ class _BaseModelLike(Protocol):
|
||||
) -> dict[str, Any]: ...
|
||||
|
||||
|
||||
_JSONLike: TypeAlias = None | str | int | float | bool
|
||||
_JSONLike: TypeAlias = str | int | float | bool | None
|
||||
_JSONMap: TypeAlias = Mapping[
|
||||
str, Union[_JSONLike, list[_JSONLike], "_JSONMap", list["_JSONMap"]]
|
||||
]
|
||||
|
||||
@@ -36,7 +36,7 @@ test = [
|
||||
"pytest-watch",
|
||||
]
|
||||
lint = [
|
||||
"ruff==0.15.20",
|
||||
"ruff==0.16.5",
|
||||
"codespell",
|
||||
"ty",
|
||||
"starlette",
|
||||
|
||||
@@ -426,11 +426,17 @@ def test_sync_run_start_sends_command():
|
||||
with httpx.Client(transport=fake.transport, base_url="http://test") as raw:
|
||||
threads = SyncThreadsClient(SyncHttpClient(raw))
|
||||
with threads.stream(thread_id="t-1", assistant_id="agent") as thread:
|
||||
result = thread.run.start(input={"x": 1})
|
||||
result = thread.run.start(
|
||||
input={"x": 1},
|
||||
langsmith_tracing={"project_name": "replica-project"},
|
||||
)
|
||||
|
||||
assert result == {"run_id": "run-1"}
|
||||
assert fake.received_commands[0]["method"] == "run.start"
|
||||
assert fake.received_commands[0]["params"]["assistant_id"] == "agent"
|
||||
assert fake.received_commands[0]["params"]["langsmith_tracer"] == {
|
||||
"project_name": "replica-project"
|
||||
}
|
||||
|
||||
|
||||
def test_sync_events_iterates_raw_events():
|
||||
@@ -584,10 +590,10 @@ def test_v3_streaming_sync_surface_smoke():
|
||||
assert results["values"] == fake.state["values"]
|
||||
messages_result = results["messages"]
|
||||
assert isinstance(messages_result, list)
|
||||
assert [str(m.text) for m in messages_result] == ["hi"] # ty: ignore[unresolved-attribute]
|
||||
assert [str(m.text) for m in messages_result] == ["hi"]
|
||||
tools_result = results["tools"]
|
||||
assert isinstance(tools_result, list)
|
||||
assert tools_result[0].name == "search" # ty: ignore[unresolved-attribute]
|
||||
assert tools_result[0].name == "search"
|
||||
assert results["progress"] == [{"name": "progress", "step": 1}]
|
||||
assert final == {"final": True}
|
||||
|
||||
|
||||
@@ -287,7 +287,7 @@ async def test_command_ids_are_monotonic():
|
||||
assert [c["id"] for c in fake.received_commands] == [1, 2]
|
||||
|
||||
|
||||
async def test_run_start_forwards_config_and_metadata():
|
||||
async def test_run_start_forwards_config_metadata_and_langsmith_tracing():
|
||||
fake = FakeServer()
|
||||
transport = httpx.ASGITransport(app=fake.app)
|
||||
async with httpx.AsyncClient(transport=transport, base_url="http://test") as raw:
|
||||
@@ -297,10 +297,18 @@ async def test_run_start_forwards_config_and_metadata():
|
||||
input={"x": 1},
|
||||
config={"recursion_limit": 5},
|
||||
metadata={"trace": "abc"},
|
||||
langsmith_tracing={
|
||||
"project_name": "replica-project",
|
||||
"example_id": "example-1",
|
||||
},
|
||||
)
|
||||
params = fake.received_commands[0]["params"]
|
||||
assert params["config"] == {"recursion_limit": 5}
|
||||
assert params["metadata"] == {"trace": "abc"}
|
||||
assert params["langsmith_tracer"] == {
|
||||
"project_name": "replica-project",
|
||||
"example_id": "example-1",
|
||||
}
|
||||
|
||||
|
||||
async def test_run_start_raises_outside_context_manager():
|
||||
|
||||
@@ -0,0 +1,132 @@
|
||||
import pytest
|
||||
|
||||
from langgraph_sdk import Auth
|
||||
|
||||
|
||||
def test_handler_multiple_resources_and_actions() -> None:
|
||||
auth = Auth()
|
||||
|
||||
@auth.on(resources=["threads", "assistants"], actions=["read", "search"])
|
||||
async def allow_reads(ctx, value):
|
||||
del value
|
||||
return {"owner": ctx.user.identity}
|
||||
|
||||
assert auth._handlers == {
|
||||
("threads", "read"): [allow_reads],
|
||||
("threads", "search"): [allow_reads],
|
||||
("assistants", "read"): [allow_reads],
|
||||
("assistants", "search"): [allow_reads],
|
||||
}
|
||||
|
||||
|
||||
def test_resource_handler_actions_are_scoped() -> None:
|
||||
auth = Auth()
|
||||
|
||||
@auth.on
|
||||
async def deny_all(ctx, value):
|
||||
del ctx, value
|
||||
return False
|
||||
|
||||
@auth.on.threads(actions=["create", "search"])
|
||||
async def handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
@auth.on.threads(actions="create_run")
|
||||
async def run_handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
assert auth._handlers == {
|
||||
("threads", "create"): [handler],
|
||||
("threads", "search"): [handler],
|
||||
("threads", "create_run"): [run_handler],
|
||||
}
|
||||
assert auth._global_handlers == [deny_all]
|
||||
|
||||
|
||||
def test_resource_handler_preserves_wildcard() -> None:
|
||||
auth = Auth()
|
||||
|
||||
@auth.on.threads
|
||||
async def handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
assert auth._handlers == {("threads", "*"): [handler]}
|
||||
|
||||
|
||||
def test_resource_handler_preserves_wildcard_with_parentheses() -> None:
|
||||
auth = Auth()
|
||||
|
||||
@auth.on.threads()
|
||||
async def handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
assert auth._handlers == {("threads", "*"): [handler]}
|
||||
|
||||
|
||||
def test_resource_handler_accepts_matching_resource() -> None:
|
||||
auth = Auth()
|
||||
|
||||
@auth.on.threads(resources=["threads"], actions="read")
|
||||
async def handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
assert auth._handlers == {("threads", "read"): [handler]}
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"resources", [["assistants"], ["threads", "assistants"], [], [1]]
|
||||
)
|
||||
def test_resource_handler_rejects_nonmatching_resources(resources) -> None:
|
||||
auth = Auth()
|
||||
|
||||
async def handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
with pytest.raises(ValueError, match=r"Use @auth\.on"):
|
||||
auth.on.threads(resources=resources)(handler)
|
||||
assert auth._handlers == {}
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("resource", "actions", "error"),
|
||||
[
|
||||
("threads", [], ValueError),
|
||||
("threads", ["reed"], ValueError),
|
||||
("threads", ["create", "create"], ValueError),
|
||||
("threads", {"create": True}, TypeError),
|
||||
("crons", ["create_run"], ValueError),
|
||||
],
|
||||
)
|
||||
def test_resource_handler_rejects_invalid_actions(resource, actions, error) -> None:
|
||||
auth = Auth()
|
||||
|
||||
async def handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
with pytest.raises(error):
|
||||
getattr(auth.on, resource)(actions=actions)(handler)
|
||||
assert auth._handlers == {}
|
||||
|
||||
|
||||
def test_resource_handler_registration_is_atomic() -> None:
|
||||
auth = Auth()
|
||||
|
||||
@auth.on.threads.read
|
||||
async def read_handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
async def handler(ctx, value):
|
||||
del ctx, value
|
||||
return None
|
||||
|
||||
with pytest.raises(ValueError, match="already set"):
|
||||
auth.on.threads(actions=["create", "read"])(handler)
|
||||
assert auth._handlers == {("threads", "read"): [read_handler]}
|
||||
@@ -1,8 +1,40 @@
|
||||
from collections.abc import Awaitable, Callable
|
||||
|
||||
import pytest
|
||||
|
||||
from langgraph_sdk import DecryptResult, EncryptionContext
|
||||
from langgraph_sdk.encryption import DuplicateHandlerError, Encryption
|
||||
|
||||
|
||||
def test_decrypt_result():
|
||||
result = DecryptResult(plaintext=b"plain", replacement=b"rotated")
|
||||
|
||||
assert result.plaintext == b"plain"
|
||||
assert result.replacement == b"rotated"
|
||||
assert DecryptResult(plaintext={"plain": True}).replacement is None
|
||||
|
||||
|
||||
def test_decrypt_decorators_preserve_return_types():
|
||||
encryption = Encryption()
|
||||
|
||||
@encryption.decrypt.blob
|
||||
async def blob_dec(_ctx: EncryptionContext, data: bytes) -> bytes:
|
||||
return data
|
||||
|
||||
@encryption.decrypt.json
|
||||
async def json_dec(
|
||||
_ctx: EncryptionContext, data: dict[str, object]
|
||||
) -> dict[str, object]:
|
||||
return data
|
||||
|
||||
blob_handler: Callable[[EncryptionContext, bytes], Awaitable[bytes]] = blob_dec
|
||||
json_handler: Callable[
|
||||
[EncryptionContext, dict[str, object]], Awaitable[dict[str, object]]
|
||||
] = json_dec
|
||||
assert blob_handler is blob_dec
|
||||
assert json_handler is json_dec
|
||||
|
||||
|
||||
class TestHandlerValidation:
|
||||
"""Test duplicate handler and signature validation."""
|
||||
|
||||
|
||||
Generated
+940
-661
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user