Compare commits

...
Author SHA1 Message Date
11ee185999 fix(checkpoint): widen Store put value type to Mapping[str, Any] (#8617)
TypedDict values don't structurally satisfy dict[str, Any] since dict
implies full mutability. Mapping[str, Any] accepts both plain dicts and
TypedDicts while still requiring string keys, matching what put()
actually needs from callers.

Fixes #8616

Verified by running lint/type/test locally across checkpoint,
checkpoint-sqlite, checkpoint-postgres, prebuilt, sdk-py, and a scoped
langgraph subset.

LinkedIn: https://linkedin.com/in/lisandro-navarra

---------

Co-authored-by: Mason Daugherty <github@mdrxy.com>
2026-08-28 08:05:04 -05:00
Sreekara YachamaneniandGitHub d5f4b2aa96 release(sdk-py): 0.4.4 (#8738)
Bumps the Python SDK version from 0.4.3 to 0.4.4.
2026-08-27 17:14:54 -04:00
Sreekara YachamaneniandGitHub 5a77be5e8b Merge commit from fork
* authz fix for custom auth

* add back resource param

* auth on multiple resources
2026-08-27 13:28:59 -07:00
Mason DaughertyGitHubopen-swe[bot] <open-swe@users.noreply.github.com>
bdb8a9c7a4 feat: route LangSmith traces from thread streams (#8723)
## Description
Expose the existing `langsmith_tracing` option on Python sync and async
thread-stream run starts and forward it through the protocol.

## Release Note
Python thread streams can route traces to an additional LangSmith
project per run.

## Test Plan
- [x] Verify sync and async run-start payloads include tracing settings

## Related PRs
- langchain-ai/agent-protocol#95
- langchain-ai/langgraphjs#2745
- langchain-ai/langgraph-api#4033

Made by [Open
SWE](https://openswe.vercel.app/agents/f9e34294-b9c3-52f0-815a-0102188e1181)

Co-authored-by: open-swe[bot] <open-swe@users.noreply.github.com>
2026-08-26 16:57:39 -04:00
16 changed files with 268 additions and 38 deletions
+1
View File
@@ -76,6 +76,7 @@ __pypackages__/
# Environments
.env
.env.*
.envrc
*.crt
*.key
@@ -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:
@@ -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,
]
@@ -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),
+15 -1
View File
@@ -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()
+9
View File
@@ -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
+1 -1
View File
@@ -3,7 +3,7 @@ from langgraph_sdk.client import get_client, get_sync_client
from langgraph_sdk.encryption import Encryption
from langgraph_sdk.encryption.types import DecryptResult, EncryptionContext
__version__ = "0.4.3"
__version__ = "0.4.4"
__all__ = [
"Auth",
+4 -1
View File
@@ -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
+4 -1
View File
@@ -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
+67 -16
View File
@@ -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[
@@ -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():
@@ -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():
+132
View File
@@ -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]}