Compare commits

..
Author SHA1 Message Date
Sydney Runkle 8086a20865 add test 2025-08-29 17:04:25 -04:00
Sydney Runkle fcfb9dd3a7 asyncio escape hatch 2025-08-29 17:02:03 -04:00
6 changed files with 43 additions and 337 deletions
+1 -264
View File
@@ -29,10 +29,6 @@
"name": "Store",
"description": "Store is an API for managing persistent key-value store (long-term memory) that is available from any thread."
},
{
"name": "A2A",
"description": "Agent-to-Agent Protocol related endpoints for exposing assistants as A2A-compliant agents."
},
{
"name": "MCP",
"description": "Model Context Protocol related endpoints for exposing an agent as an MCP server."
@@ -1554,29 +1550,6 @@
},
"name": "Last-Event-ID",
"in": "header"
},
{
"required": false,
"schema": {
"anyOf": [
{
"type": "string",
"enum": ["lifecycle", "run_modes", "state_update"]
},
{
"type": "array",
"items": {
"type": "string",
"enum": ["lifecycle", "run_modes", "state_update"]
}
}
],
"default": ["run_modes"],
"title": "Stream Modes",
"description": "Stream modes to control which events are returned. 'lifecycle' returns only run start/end events, 'run_modes' returns all run events (default behavior), 'state_update' returns only state update events."
},
"name": "stream_modes",
"in": "query"
}
],
"responses": {
@@ -3186,195 +3159,6 @@
}
}
},
"/a2a/{assistant_id}": {
"post": {
"operationId": "post_a2a",
"summary": "A2A Post",
"description": "Communicate with an assistant using the Agent-to-Agent Protocol.\nSends a JSON-RPC 2.0 message to the assistant.\n\n- **Request**: Provide an object with `jsonrpc`, `id`, `method`, and optional `params`.\n- **Response**: Returns a JSON-RPC response with task information or error.\n\n**Supported Methods:**\n- `message/send`: Send a message to the assistant\n- `tasks/get`: Get the status and result of a task\n\n**Notes:**\n- Supports threaded conversations via thread context\n- Messages can contain text and data parts\n- Tasks run asynchronously and return completion status\n",
"parameters": [
{
"name": "assistant_id",
"in": "path",
"required": true,
"schema": {
"type": "string",
"format": "uuid"
},
"description": "The ID of the assistant to communicate with"
},
{
"name": "Accept",
"in": "header",
"required": true,
"schema": {
"type": "string",
"enum": ["application/json"]
},
"description": "Must be application/json"
}
],
"requestBody": {
"required": true,
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"jsonrpc": {
"type": "string",
"enum": ["2.0"],
"description": "JSON-RPC version"
},
"id": {
"type": "string",
"description": "Request identifier"
},
"method": {
"type": "string",
"enum": ["message/send", "tasks/get"],
"description": "The method to invoke"
},
"params": {
"type": "object",
"description": "Method parameters",
"oneOf": [
{
"title": "Message Send Parameters",
"properties": {
"message": {
"type": "object",
"properties": {
"role": {
"type": "string",
"enum": ["user", "assistant"],
"description": "Message role"
},
"parts": {
"type": "array",
"items": {
"oneOf": [
{
"title": "Text Part",
"type": "object",
"properties": {
"kind": {
"type": "string",
"enum": ["text"]
},
"text": {
"type": "string"
}
},
"required": ["kind", "text"]
},
{
"title": "Data Part",
"type": "object",
"properties": {
"kind": {
"type": "string",
"enum": ["data"]
},
"data": {
"type": "object"
}
},
"required": ["kind", "data"]
}
]
},
"description": "Message parts"
},
"messageId": {
"type": "string",
"description": "Unique message identifier"
}
},
"required": ["role", "parts", "messageId"]
},
"thread": {
"type": "object",
"properties": {
"threadId": {
"type": "string",
"description": "Thread identifier for conversation context"
}
},
"description": "Optional thread context"
}
},
"required": ["message"]
},
{
"title": "Task Get Parameters",
"properties": {
"taskId": {
"type": "string",
"description": "Task identifier to retrieve"
}
},
"required": ["taskId"]
}
]
}
},
"required": ["jsonrpc", "id", "method"]
}
}
}
},
"responses": {
"200": {
"description": "JSON-RPC response",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"jsonrpc": {
"type": "string",
"enum": ["2.0"]
},
"id": {
"type": "string"
},
"result": {
"type": "object",
"description": "Success result containing task information or task details"
},
"error": {
"type": "object",
"properties": {
"code": {
"type": "integer"
},
"message": {
"type": "string"
}
},
"description": "Error information if request failed"
}
},
"required": ["jsonrpc", "id"]
}
}
}
},
"400": {
"description": "Bad request - invalid JSON-RPC or missing Accept header"
},
"404": {
"description": "Assistant not found"
},
"500": {
"description": "Internal server error"
}
},
"tags": [
"A2A"
]
}
},
"/mcp/": {
"post": {
"operationId": "post_mcp",
@@ -4629,17 +4413,6 @@
"title": "Checkpoint During",
"description": "Whether to checkpoint during the run.",
"default": false
},
"durability": {
"type": "string",
"enum": [
"sync",
"async",
"exit"
],
"title": "Durability",
"description": "Durability level for the run. Must be one of 'sync', 'async', or 'exit'.",
"default": "async"
}
},
"type": "object",
@@ -4876,17 +4649,6 @@
"title": "Checkpoint During",
"description": "Whether to checkpoint during the run.",
"default": false
},
"durability": {
"type": "string",
"enum": [
"sync",
"async",
"exit"
],
"title": "Durability",
"description": "Durability level for the run. Must be one of 'sync', 'async', or 'exit'.",
"default": "async"
}
},
"type": "object",
@@ -5015,12 +4777,6 @@
},
"ThreadSearchRequest": {
"properties": {
"ids": {
"type": "array",
"items": {"type": "string", "format": "uuid"},
"title": "Ids",
"description": "List of thread IDs to include. Others are excluded."
},
"metadata": {
"type": "object",
"title": "Metadata",
@@ -5261,30 +5017,11 @@
"type": "object",
"title": "Metadata",
"description": "Metadata to merge with existing thread metadata."
},
"ttl": {
"type": "object",
"title": "TTL",
"description": "The time-to-live for the thread.",
"properties": {
"strategy": {
"type": "string",
"enum": [
"delete"
],
"description": "The TTL strategy. 'delete' removes the entire thread.",
"default": "delete"
},
"ttl": {
"type": "number",
"description": "The time-to-live in minutes from now until thread should be swept."
}
}
}
},
"type": "object",
"title": "ThreadPatch",
"description": "Payload for updating a thread."
"description": "Payload for creating a thread."
},
"ThreadStateCheckpointRequest": {
"properties": {
@@ -443,7 +443,12 @@ class ToolNode(RunnableCallable):
return invalid_tool_message
try:
call_args = {**call, **{"type": "tool_call"}}
response = self.tools_by_name[call["name"]].invoke(call_args, config)
tool = self.tools_by_name[call["name"]]
try:
response = tool.invoke(call_args, config)
except NotImplementedError:
response = asyncio.run(tool.ainvoke(call_args, config))
# GraphInterrupt is a special exception that will always be raised.
# It can be triggered in the following scenarios,
+33
View File
@@ -1156,3 +1156,36 @@ async def test_tool_node_command_remove_all_messages():
command = result[0]
assert isinstance(command, Command)
assert command.update == {"messages": [RemoveMessage(id=REMOVE_ALL_MESSAGES)]}
async def test_async_tool_called_syncly() -> None:
"""Confirm that async tools can be called synchronously."""
@dec_tool
async def async_tool():
"""An async tool."""
return "async tool"
tool_node = ToolNode([async_tool])
result = tool_node.invoke(
{
"messages": [
AIMessage(
content="",
tool_calls=[
{
"name": "async_tool",
"args": {},
"id": "1",
"type": "tool_call",
}
],
)
]
}
)
assert result == {
"messages": [
ToolMessage(content="async tool", name="async_tool", tool_call_id="1")
]
}
+1 -1
View File
@@ -1,6 +1,6 @@
from langgraph_sdk.auth import Auth
from langgraph_sdk.client import get_client, get_sync_client
__version__ = "0.2.6"
__version__ = "0.2.4"
__all__ = ["Auth", "get_client", "get_sync_client"]
-20
View File
@@ -400,20 +400,6 @@ class AuthContext(BaseAuthContext):
"""
class ThreadTTL(typing.TypedDict, total=False):
"""Time-to-live configuration for a thread.
Matches the OpenAPI schema where TTL is represented as an object with
an optional strategy and a time value in minutes.
"""
strategy: typing.Literal["delete"]
"""TTL strategy. Currently only 'delete' is supported."""
ttl: int
"""Time-to-live in minutes from now until the thread should be swept."""
class ThreadsCreate(typing.TypedDict, total=False):
"""Parameters for creating a new thread.
@@ -436,9 +422,6 @@ class ThreadsCreate(typing.TypedDict, total=False):
if_exists: OnConflictBehavior
"""Behavior when a thread with the same ID already exists."""
ttl: ThreadTTL
"""Optional TTL configuration for the thread."""
class ThreadsRead(typing.TypedDict, total=False):
"""Parameters for reading thread state or run information.
@@ -506,9 +489,6 @@ class ThreadsSearch(typing.TypedDict, total=False):
offset: int
"""Offset for pagination."""
ids: Sequence[UUID] | None
"""typing.Optional list of thread IDs to filter by."""
thread_id: UUID | None
"""typing.Optional thread ID to filter by."""
+2 -51
View File
@@ -1179,7 +1179,6 @@ class ThreadsClient:
if_exists: OnConflictBehavior | None = None,
supersteps: Sequence[dict[str, Sequence[dict[str, Any]]]] | None = None,
graph_id: str | None = None,
ttl: int | Mapping[str, Any] | None = None,
headers: Mapping[str, str] | None = None,
params: QueryParamTypes | None = None,
) -> Thread:
@@ -1194,9 +1193,6 @@ class ThreadsClient:
supersteps: Apply a list of supersteps when creating a thread, each containing a sequence of updates.
Each update has `values` or `command` and `as_node`. Used for copying a thread between deployments.
graph_id: Optional graph ID to associate with the thread.
ttl: Optional time-to-live in minutes for the thread. You can pass an
integer (minutes) or a mapping with keys `ttl` and optional
`strategy` (defaults to "delete").
headers: Optional custom headers to include with the request.
params: Optional query parameters to include with the request.
@@ -1238,11 +1234,6 @@ class ThreadsClient:
}
for s in supersteps
]
if ttl is not None:
if isinstance(ttl, (int, float)):
payload["ttl"] = {"ttl": ttl, "strategy": "delete"}
else:
payload["ttl"] = ttl
return await self.http.post(
"/threads", json=payload, headers=headers, params=params
@@ -1253,7 +1244,6 @@ class ThreadsClient:
thread_id: str,
*,
metadata: Mapping[str, Any],
ttl: int | Mapping[str, Any] | None = None,
headers: Mapping[str, str] | None = None,
params: QueryParamTypes | None = None,
) -> Thread:
@@ -1262,9 +1252,6 @@ class ThreadsClient:
Args:
thread_id: ID of thread to update.
metadata: Metadata to merge with existing thread metadata.
ttl: Optional time-to-live in minutes for the thread. You can pass an
integer (minutes) or a mapping with keys `ttl` and optional
`strategy` (defaults to "delete").
headers: Optional custom headers to include with the request.
params: Optional query parameters to include with the request.
@@ -1278,19 +1265,12 @@ class ThreadsClient:
thread = await client.threads.update(
thread_id="my-thread-id",
metadata={"number":1},
ttl=43_200,
)
```
""" # noqa: E501
payload: dict[str, Any] = {"metadata": metadata}
if ttl is not None:
if isinstance(ttl, (int, float)):
payload["ttl"] = {"ttl": ttl, "strategy": "delete"}
else:
payload["ttl"] = ttl
return await self.http.patch(
f"/threads/{thread_id}",
json=payload,
json={"metadata": metadata},
headers=headers,
params=params,
)
@@ -1329,7 +1309,6 @@ class ThreadsClient:
*,
metadata: Json = None,
values: Json = None,
ids: Sequence[str] | None = None,
status: ThreadStatus | None = None,
limit: int = 10,
offset: int = 0,
@@ -1344,7 +1323,6 @@ class ThreadsClient:
Args:
metadata: Thread metadata to filter on.
values: State values to filter on.
ids: List of thread IDs to filter by.
status: Thread status to filter on.
Must be one of 'idle', 'busy', 'interrupted' or 'error'.
limit: Limit on number of threads to return.
@@ -1378,8 +1356,6 @@ class ThreadsClient:
payload["metadata"] = metadata
if values:
payload["values"] = values
if ids:
payload["ids"] = ids
if status:
payload["status"] = status
if sort_by:
@@ -4352,7 +4328,6 @@ class SyncThreadsClient:
if_exists: OnConflictBehavior | None = None,
supersteps: Sequence[dict[str, Sequence[dict[str, Any]]]] | None = None,
graph_id: str | None = None,
ttl: int | Mapping[str, Any] | None = None,
headers: Mapping[str, str] | None = None,
params: QueryParamTypes | None = None,
) -> Thread:
@@ -4367,9 +4342,6 @@ class SyncThreadsClient:
supersteps: Apply a list of supersteps when creating a thread, each containing a sequence of updates.
Each update has `values` or `command` and `as_node`. Used for copying a thread between deployments.
graph_id: Optional graph ID to associate with the thread.
ttl: Optional time-to-live in minutes for the thread. You can pass an
integer (minutes) or a mapping with keys `ttl` and optional
`strategy` (defaults to "delete").
headers: Optional custom headers to include with the request.
Returns:
@@ -4411,11 +4383,6 @@ class SyncThreadsClient:
}
for s in supersteps
]
if ttl is not None:
if isinstance(ttl, (int, float)):
payload["ttl"] = {"ttl": ttl, "strategy": "delete"}
else:
payload["ttl"] = ttl
return self.http.post("/threads", json=payload, headers=headers, params=params)
@@ -4424,7 +4391,6 @@ class SyncThreadsClient:
thread_id: str,
*,
metadata: Mapping[str, Any],
ttl: int | Mapping[str, Any] | None = None,
headers: Mapping[str, str] | None = None,
params: QueryParamTypes | None = None,
) -> Thread:
@@ -4433,11 +4399,7 @@ class SyncThreadsClient:
Args:
thread_id: ID of thread to update.
metadata: Metadata to merge with existing thread metadata.
ttl: Optional time-to-live in minutes for the thread. You can pass an
integer (minutes) or a mapping with keys `ttl` and optional
`strategy` (defaults to "delete").
headers: Optional custom headers to include with the request.
params: Optional query parameters to include with the request.
Returns:
Thread: The created thread.
@@ -4449,19 +4411,12 @@ class SyncThreadsClient:
thread = client.threads.update(
thread_id="my-thread-id",
metadata={"number":1},
ttl=43_200,
)
```
""" # noqa: E501
payload: dict[str, Any] = {"metadata": metadata}
if ttl is not None:
if isinstance(ttl, (int, float)):
payload["ttl"] = {"ttl": ttl, "strategy": "delete"}
else:
payload["ttl"] = ttl
return self.http.patch(
f"/threads/{thread_id}",
json=payload,
json={"metadata": metadata},
headers=headers,
params=params,
)
@@ -4499,7 +4454,6 @@ class SyncThreadsClient:
*,
metadata: Json = None,
values: Json = None,
ids: Sequence[str] | None = None,
status: ThreadStatus | None = None,
limit: int = 10,
offset: int = 0,
@@ -4514,7 +4468,6 @@ class SyncThreadsClient:
Args:
metadata: Thread metadata to filter on.
values: State values to filter on.
ids: List of thread IDs to filter by.
status: Thread status to filter on.
Must be one of 'idle', 'busy', 'interrupted' or 'error'.
limit: Limit on number of threads to return.
@@ -4544,8 +4497,6 @@ class SyncThreadsClient:
payload["metadata"] = metadata
if values:
payload["values"] = values
if ids:
payload["ids"] = ids
if status:
payload["status"] = status
if sort_by: