mirror of
https://github.com/langchain-ai/langgraph.git
synced 2026-10-02 22:45:07 +02:00
Compare commits
19
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8563948e70 | ||
|
|
b0f14649e0 | ||
|
|
ea20432b9b | ||
|
|
e2efab8061 | ||
|
|
9babffa054 | ||
|
|
73cebea3c2 | ||
|
|
b73b2d19eb | ||
|
|
ca26805b5f | ||
|
|
5ac837d7cd | ||
|
|
c4f5861166 | ||
|
|
172238b2d5 | ||
|
|
095da17833 | ||
|
|
e931c68669 | ||
|
|
666c224c2d | ||
|
|
21a6f41e0a | ||
|
|
acdf85aba6 | ||
|
|
9b18243fa6 | ||
|
|
0225b998af | ||
|
|
bec122d4a2 |
@@ -31,6 +31,8 @@ jobs:
|
||||
workdir: libs/cli/examples/graphs_reqs_b
|
||||
tag: langgraph-test-d
|
||||
name: "CLI integration test"
|
||||
env:
|
||||
HAS_LANGSMITH_API_KEY: ${{ secrets.LANGSMITH_API_KEY != '' }}
|
||||
defaults:
|
||||
run:
|
||||
working-directory: libs/cli
|
||||
@@ -58,7 +60,7 @@ jobs:
|
||||
run: |
|
||||
langgraph build -t ${{ matrix.example.tag }}
|
||||
- name: Test service ${{ matrix.example.name }}
|
||||
if: ${{ steps.changed-files.outputs.all && secrets.LANGSMITH_API_KEY != '' }}
|
||||
if: ${{ steps.changed-files.outputs.all && env.HAS_LANGSMITH_API_KEY == 'true' }}
|
||||
working-directory: ${{ matrix.example.workdir }}
|
||||
env:
|
||||
LANGSMITH_API_KEY: ${{ secrets.LANGSMITH_API_KEY }}
|
||||
@@ -89,7 +91,7 @@ jobs:
|
||||
run: |
|
||||
langgraph build -t langgraph-test-g -c apps/agent/langgraph.json
|
||||
- name: Test Python monorepo service
|
||||
if: ${{ steps.changed-files.outputs.all && matrix.example.name == 'A' && secrets.LANGSMITH_API_KEY != '' }}
|
||||
if: ${{ steps.changed-files.outputs.all && matrix.example.name == 'A' && env.HAS_LANGSMITH_API_KEY == 'true' }}
|
||||
working-directory: libs/cli/python-monorepo-example
|
||||
env:
|
||||
LANGSMITH_API_KEY: ${{ secrets.LANGSMITH_API_KEY }}
|
||||
@@ -104,7 +106,7 @@ jobs:
|
||||
run: |
|
||||
langgraph build -t langgraph-test-h
|
||||
- name: Test prerelease reqs service
|
||||
if: ${{ steps.changed-files.outputs.all && matrix.example.name == 'A' && secrets.LANGSMITH_API_KEY != '' }}
|
||||
if: ${{ steps.changed-files.outputs.all && matrix.example.name == 'A' && env.HAS_LANGSMITH_API_KEY == 'true' }}
|
||||
working-directory: libs/cli/examples/graph_prerelease_reqs
|
||||
env:
|
||||
LANGSMITH_API_KEY: ${{ secrets.LANGSMITH_API_KEY }}
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
name: Deploy Redirects to GitHub Pages
|
||||
|
||||
on:
|
||||
push:
|
||||
branches:
|
||||
- main
|
||||
paths:
|
||||
- 'docs/**'
|
||||
- '.github/workflows/deploy-redirects.yml'
|
||||
workflow_dispatch:
|
||||
|
||||
permissions:
|
||||
contents: read
|
||||
pages: write
|
||||
id-token: write
|
||||
|
||||
concurrency:
|
||||
group: "pages"
|
||||
cancel-in-progress: false
|
||||
|
||||
jobs:
|
||||
deploy:
|
||||
environment:
|
||||
name: github-pages
|
||||
url: ${{ steps.deployment.outputs.page_url }}
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@v4
|
||||
|
||||
- name: Setup Python
|
||||
uses: actions/setup-python@v5
|
||||
with:
|
||||
python-version: '3.11'
|
||||
|
||||
- name: Generate redirect files
|
||||
run: python docs/generate_redirects.py
|
||||
|
||||
- name: Setup Pages
|
||||
uses: actions/configure-pages@v4
|
||||
|
||||
- name: Upload artifact
|
||||
uses: actions/upload-pages-artifact@v3
|
||||
with:
|
||||
path: 'docs/_site'
|
||||
|
||||
- name: Deploy to GitHub Pages
|
||||
id: deployment
|
||||
uses: actions/deploy-pages@v4
|
||||
@@ -0,0 +1 @@
|
||||
_site/
|
||||
@@ -0,0 +1,142 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Generate HTML redirect files from redirects.json.
|
||||
|
||||
Usage:
|
||||
python generate_redirects.py
|
||||
|
||||
This script reads redirects.json and generates individual HTML files
|
||||
for each redirect path. Each HTML file uses meta refresh (0 delay)
|
||||
which is SEO-friendly and treated similarly to 301 redirects by Google.
|
||||
|
||||
To add new redirects, simply edit redirects.json and re-run this script.
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
# Default fallback URL for any path not in the redirect map
|
||||
DEFAULT_REDIRECT = "https://docs.langchain.com/oss/python/langgraph/overview"
|
||||
|
||||
HTML_TEMPLATE = """<!doctype html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="utf-8">
|
||||
<title>Redirecting...</title>
|
||||
<link rel="canonical" href="{url}">
|
||||
<meta name="robots" content="noindex">
|
||||
<script>var anchor=window.location.hash.substr(1);location.href="{url}"+(anchor?"#"+anchor:"")</script>
|
||||
<meta http-equiv="refresh" content="0; url={url}">
|
||||
</head>
|
||||
<body>
|
||||
Redirecting...
|
||||
</body>
|
||||
</html>
|
||||
"""
|
||||
|
||||
ROOT_HTML_TEMPLATE = """<!doctype html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="utf-8">
|
||||
<title>Redirecting to LangGraph Documentation</title>
|
||||
<link rel="canonical" href="{url}">
|
||||
<meta name="robots" content="noindex">
|
||||
<script>var anchor=window.location.hash.substr(1);location.href="{url}"+(anchor?"#"+anchor:"")</script>
|
||||
<meta http-equiv="refresh" content="0; url={url}">
|
||||
</head>
|
||||
<body>
|
||||
<h1>Documentation has moved</h1>
|
||||
<p>The LangGraph documentation has moved to <a href="{url}">docs.langchain.com</a>.</p>
|
||||
<p>Redirecting you now...</p>
|
||||
</body>
|
||||
</html>
|
||||
"""
|
||||
|
||||
CATCHALL_404_TEMPLATE = """<!doctype html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="utf-8">
|
||||
<title>Redirecting to LangGraph Documentation</title>
|
||||
<link rel="canonical" href="{default_url}">
|
||||
<meta name="robots" content="noindex">
|
||||
<script>
|
||||
// Catchall redirect for any unmapped paths
|
||||
window.location.replace("{default_url}");
|
||||
</script>
|
||||
<meta http-equiv="refresh" content="0; url={default_url}">
|
||||
</head>
|
||||
<body>
|
||||
<h1>Documentation has moved</h1>
|
||||
<p>The LangGraph documentation has moved to <a href="{default_url}">docs.langchain.com</a>.</p>
|
||||
<p>Redirecting you now...</p>
|
||||
</body>
|
||||
</html>
|
||||
"""
|
||||
|
||||
|
||||
def generate_redirects():
|
||||
script_dir = Path(__file__).parent
|
||||
output_dir = script_dir / "_site"
|
||||
|
||||
# Load redirects
|
||||
with open(script_dir / "redirects.json") as f:
|
||||
redirects = json.load(f)
|
||||
|
||||
# Clean output directory
|
||||
if output_dir.exists():
|
||||
import shutil
|
||||
shutil.rmtree(output_dir)
|
||||
output_dir.mkdir(parents=True)
|
||||
|
||||
# Generate individual HTML files for each redirect
|
||||
for old_path, new_url in redirects.items():
|
||||
# Remove leading slash and create directory structure
|
||||
path = old_path.lstrip("/")
|
||||
|
||||
# Check if path has a file extension (e.g., .txt, .xml)
|
||||
# If so, create the file directly instead of a directory with index.html
|
||||
path_obj = Path(path)
|
||||
has_extension = path_obj.suffix and len(path_obj.suffix) <= 5
|
||||
|
||||
if not path:
|
||||
html_path = output_dir / "index.html"
|
||||
elif has_extension:
|
||||
# For files with extensions, create the file directly
|
||||
html_path = output_dir / path
|
||||
else:
|
||||
# For directory-style URLs, create index.html inside
|
||||
html_path = output_dir / path / "index.html"
|
||||
|
||||
# Create parent directories
|
||||
html_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# Write the redirect HTML
|
||||
html_path.write_text(HTML_TEMPLATE.format(url=new_url))
|
||||
print(f"Created: {html_path}")
|
||||
|
||||
# Create root index.html
|
||||
root_index = output_dir / "index.html"
|
||||
if not root_index.exists():
|
||||
root_index.write_text(ROOT_HTML_TEMPLATE.format(url=DEFAULT_REDIRECT))
|
||||
print(f"Created: {root_index}")
|
||||
|
||||
# Create 404.html for catchall
|
||||
catchall_404 = output_dir / "404.html"
|
||||
catchall_404.write_text(CATCHALL_404_TEMPLATE.format(default_url=DEFAULT_REDIRECT))
|
||||
print(f"Created: {catchall_404}")
|
||||
|
||||
# Copy static files (like llms.txt) that can't be redirected via HTML
|
||||
static_files = ["llms.txt"]
|
||||
for static_file in static_files:
|
||||
src = script_dir / static_file
|
||||
if src.exists():
|
||||
dst = output_dir / static_file
|
||||
dst.write_text(src.read_text())
|
||||
print(f"Copied: {dst}")
|
||||
|
||||
print(f"\nGenerated {len(redirects)} redirect files in {output_dir}")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
generate_redirects()
|
||||
@@ -0,0 +1,35 @@
|
||||
# LangGraph
|
||||
|
||||
LangGraph documentation has moved to docs.langchain.com.
|
||||
|
||||
## Overview
|
||||
|
||||
- [LangGraph Overview](https://docs.langchain.com/oss/python/langgraph/overview): Introduction to LangGraph, a library for building stateful, multi-actor applications with LLMs.
|
||||
- [Why LangGraph?](https://docs.langchain.com/oss/python/langgraph/why-langgraph): Motivation for LangGraph and its key features.
|
||||
|
||||
## Core Concepts
|
||||
|
||||
- [Graph API](https://docs.langchain.com/oss/python/langgraph/graph-api): Learn how to define state, create nodes, and connect them with edges.
|
||||
- [Streaming](https://docs.langchain.com/oss/python/langgraph/streaming): Stream outputs from your graph for better UX.
|
||||
- [Persistence](https://docs.langchain.com/oss/python/langgraph/persistence): Add memory and checkpointing to your graphs.
|
||||
- [Add Memory](https://docs.langchain.com/oss/python/langgraph/add-memory): Implement short-term and long-term memory.
|
||||
- [Workflows & Agents](https://docs.langchain.com/oss/python/langgraph/workflows-agents): Build agents and workflows with LangGraph.
|
||||
|
||||
## How-To Guides
|
||||
|
||||
- [Use Subgraphs](https://docs.langchain.com/oss/python/langgraph/use-subgraphs): Compose graphs using subgraphs.
|
||||
- [Observability](https://docs.langchain.com/oss/python/langgraph/observability): Add tracing and debugging to your graphs.
|
||||
- [Common Errors](https://docs.langchain.com/oss/python/langgraph/common-errors): Troubleshoot common LangGraph errors.
|
||||
|
||||
## Tutorials
|
||||
|
||||
- [Agentic RAG](https://docs.langchain.com/oss/python/langgraph/agentic-rag): Build an agentic RAG system with LangGraph.
|
||||
- [SQL Agent](https://docs.langchain.com/oss/python/langgraph/sql-agent): Create a SQL agent with LangGraph.
|
||||
|
||||
## Reference
|
||||
|
||||
- [API Reference](https://reference.langchain.com/python/langgraph/): Complete API documentation for LangGraph.
|
||||
|
||||
## LangGraph Platform
|
||||
|
||||
For deploying LangGraph applications in production, see the [LangSmith documentation](https://docs.langchain.com/langsmith/agent-server).
|
||||
@@ -0,0 +1,296 @@
|
||||
{
|
||||
"/how-tos/stream-values": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/stream-updates": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/streaming-content": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/stream-multiple": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/streaming-tokens-without-langchain": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/streaming-from-final-node": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/streaming-events-from-within-tools-without-langchain": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/state-reducers": "https://docs.langchain.com/oss/python/langgraph/graph-api#define-and-update-state",
|
||||
"/how-tos/sequence": "https://docs.langchain.com/oss/python/langgraph/graph-api#create-a-sequence-of-steps",
|
||||
"/how-tos/branching": "https://docs.langchain.com/oss/python/langgraph/graph-api#create-branches",
|
||||
"/how-tos/recursion-limit": "https://docs.langchain.com/oss/python/langgraph/graph-api#create-and-control-loops",
|
||||
"/how-tos/visualization": "https://docs.langchain.com/oss/python/langgraph/graph-api#visualize-your-graph",
|
||||
"/how-tos/input_output_schema": "https://docs.langchain.com/oss/python/langgraph/graph-api#define-input-and-output-schemas",
|
||||
"/how-tos/pass_private_state": "https://docs.langchain.com/oss/python/langgraph/graph-api#pass-private-state-between-nodes",
|
||||
"/how-tos/state-model": "https://docs.langchain.com/oss/python/langgraph/graph-api#use-pydantic-models-for-graph-state",
|
||||
"/how-tos/map-reduce": "https://docs.langchain.com/oss/python/langgraph/graph-api#map-reduce-and-the-send-api",
|
||||
"/how-tos/command": "https://docs.langchain.com/oss/python/langgraph/graph-api#combine-control-flow-and-state-updates-with-command",
|
||||
"/how-tos/configuration": "https://docs.langchain.com/oss/python/langgraph/graph-api#add-runtime-configuration",
|
||||
"/how-tos/node-retries": "https://docs.langchain.com/oss/python/langgraph/graph-api#add-retry-policies",
|
||||
"/how-tos/return-when-recursion-limit-hits": "https://docs.langchain.com/oss/python/langgraph/graph-api#impose-a-recursion-limit",
|
||||
"/how-tos/async": "https://docs.langchain.com/oss/python/langgraph/graph-api#async",
|
||||
"/how-tos/memory/manage-conversation-history": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/memory/delete-messages": "https://docs.langchain.com/oss/python/langgraph/add-memory#delete-messages",
|
||||
"/how-tos/memory/add-summary-conversation-history": "https://docs.langchain.com/oss/python/langgraph/add-memory#summarize-messages",
|
||||
"/how-tos/memory": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/agents/memory": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/subgraph-transform-state": "https://docs.langchain.com/oss/python/langgraph/use-subgraphs#different-state-schemas",
|
||||
"/how-tos/subgraphs-manage-state": "https://docs.langchain.com/oss/python/langgraph/use-subgraphs#add-persistence",
|
||||
"/how-tos/persistence_postgres": "https://docs.langchain.com/oss/python/langgraph/add-memory#use-in-production",
|
||||
"/how-tos/persistence_mongodb": "https://docs.langchain.com/oss/python/langgraph/add-memory#use-in-production",
|
||||
"/how-tos/persistence_redis": "https://docs.langchain.com/oss/python/langgraph/add-memory#use-in-production",
|
||||
"/how-tos/subgraph-persistence": "https://docs.langchain.com/oss/python/langgraph/add-memory#use-with-subgraphs",
|
||||
"/how-tos/cross-thread-persistence": "https://docs.langchain.com/oss/python/langgraph/add-memory#add-long-term-memory",
|
||||
"/cloud/how-tos/copy_threads": "https://docs.langchain.com/langsmith/use-threads",
|
||||
"/cloud/how-tos/check-thread-status": "https://docs.langchain.com/langsmith/use-threads",
|
||||
"/cloud/concepts/threads": "https://docs.langchain.com/oss/python/langgraph/persistence#threads",
|
||||
"/how-tos/persistence": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/tool-calling-errors": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/how-tos/pass-config-to-tools": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/how-tos/pass-run-time-values-to-tools": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/how-tos/update-state-from-tools": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/agents/tools": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/how-tos/agent-handoffs": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/how-tos/multi-agent-network": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/how-tos/multi-agent-multi-turn-convo": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/cloud/index": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/cloud/how-tos/index": "https://docs.langchain.com/langsmith/home",
|
||||
"/cloud/concepts/api": "https://docs.langchain.com/langsmith/agent-server",
|
||||
"/cloud/concepts/cloud": "https://docs.langchain.com/langsmith/cloud",
|
||||
"/cloud/faq/studio": "https://docs.langchain.com/langsmith/studio",
|
||||
"/cloud/how-tos/human_in_the_loop_edit_state": "https://docs.langchain.com/langsmith/add-human-in-the-loop",
|
||||
"/cloud/how-tos/human_in_the_loop_user_input": "https://docs.langchain.com/langsmith/add-human-in-the-loop",
|
||||
"/concepts/platform_architecture": "https://docs.langchain.com/langsmith/cloud#architecture",
|
||||
"/cloud/how-tos/stream_values": "https://docs.langchain.com/langsmith/streaming",
|
||||
"/cloud/how-tos/stream_updates": "https://docs.langchain.com/langsmith/streaming",
|
||||
"/cloud/how-tos/stream_messages": "https://docs.langchain.com/langsmith/streaming",
|
||||
"/cloud/how-tos/stream_events": "https://docs.langchain.com/langsmith/streaming",
|
||||
"/cloud/how-tos/stream_debug": "https://docs.langchain.com/langsmith/streaming",
|
||||
"/cloud/how-tos/stream_multiple": "https://docs.langchain.com/langsmith/streaming",
|
||||
"/cloud/concepts/streaming": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/agents/streaming": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/create-react-agent": "https://docs.langchain.com/oss/python/langchain/agents#basic-configuration",
|
||||
"/how-tos/create-react-agent-memory": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/create-react-agent-system-prompt": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/create-react-agent-structured-output": "https://docs.langchain.com/oss/python/langchain/agents#structured-output",
|
||||
"/prebuilt": "https://docs.langchain.com/oss/python/langchain/agents",
|
||||
"/reference/prebuilt": "https://reference.langchain.com/python/langgraph/agents/",
|
||||
"/concepts/high_level": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/concepts/index": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/concepts/v0-human-in-the-loop": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/how-tos/index": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/introduction": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/agents/deployment": "https://docs.langchain.com/oss/python/langgraph/local-server",
|
||||
"/how-tos/deploy-self-hosted": "https://docs.langchain.com/langsmith/platform-setup",
|
||||
"/concepts/self_hosted": "https://docs.langchain.com/langsmith/platform-setup",
|
||||
"/tutorials/deployment": "https://docs.langchain.com/langsmith/deployments",
|
||||
"/cloud/how-tos/assistant_versioning": "https://docs.langchain.com/langsmith/configuration-cloud",
|
||||
"/cloud/concepts/runs": "https://docs.langchain.com/langsmith/assistants#execution",
|
||||
"/how-tos/wait-user-input-functional": "https://docs.langchain.com/oss/python/langgraph/functional-api",
|
||||
"/how-tos/review-tool-calls-functional": "https://docs.langchain.com/oss/python/langgraph/functional-api",
|
||||
"/how-tos/create-react-agent-hitl": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/agents/human-in-the-loop": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/how-tos/human_in_the_loop/dynamic_breakpoints": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/concepts/breakpoints": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/how-tos/human_in_the_loop/breakpoints": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/cloud/how-tos/human_in_the_loop_breakpoint": "https://docs.langchain.com/langsmith/add-human-in-the-loop",
|
||||
"/how-tos/human_in_the_loop/edit-graph-state": "https://docs.langchain.com/oss/python/langgraph/use-time-travel",
|
||||
"/examples/index": "https://docs.langchain.com/oss/python/langgraph/case-studies",
|
||||
"/guides/index": "https://docs.langchain.com/oss/python/langchain/overview",
|
||||
"/tutorials/index": "https://docs.langchain.com/oss/python/learn",
|
||||
"/llms-txt-overview": "https://docs.langchain.com/llms.txt",
|
||||
"/tutorials/rag/langgraph_adaptive_rag": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/tutorials/multi_agent/multi-agent-collaboration": "https://docs.langchain.com/oss/python/langchain/multi-agent",
|
||||
"/how-tos/create-react-agent-manage-message-history": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/many-tools": "https://docs.langchain.com/oss/python/langchain/tools",
|
||||
"/tutorials/customer-support/customer-support": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/how-tos/react-agent-structured-output": "https://docs.langchain.com/oss/python/langchain/agents#structured-output",
|
||||
"/tutorials/code_assistant/langgraph_code_assistant": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/tutorials/multi_agent/hierarchical_agent_teams": "https://docs.langchain.com/oss/python/langchain/supervisor",
|
||||
"/tutorials/auth/getting_started": "https://docs.langchain.com/langsmith/auth",
|
||||
"/tutorials/auth/resource_auth": "https://docs.langchain.com/langsmith/resource-auth",
|
||||
"/tutorials/auth/add_auth_server": "https://docs.langchain.com/langsmith/add-auth-server",
|
||||
"/how-tos/use-remote-graph": "https://docs.langchain.com/langsmith/use-remote-graph",
|
||||
"/how-tos/autogen-integration": "https://docs.langchain.com/langsmith/autogen-integration",
|
||||
"/how-tos/human_in_the_loop/wait-user-input": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/cloud/how-tos/use_stream_react": "https://docs.langchain.com/langsmith/use-stream-react",
|
||||
"/cloud/how-tos/generative_ui_react": "https://docs.langchain.com/langsmith/generative-ui-react",
|
||||
"/concepts/langgraph_platform": "https://docs.langchain.com/langsmith/home",
|
||||
"/concepts/langgraph_components": "https://docs.langchain.com/langsmith/components",
|
||||
"/concepts/langgraph_server": "https://docs.langchain.com/langsmith/agent-server",
|
||||
"/concepts/langgraph_data_plane": "https://docs.langchain.com/langsmith/data-plane",
|
||||
"/concepts/langgraph_control_plane": "https://docs.langchain.com/langsmith/control-plane",
|
||||
"/concepts/langgraph_cli": "https://docs.langchain.com/langsmith/cli",
|
||||
"/concepts/langgraph_studio": "https://docs.langchain.com/langsmith/studio",
|
||||
"/cloud/how-tos/studio/quick_start": "https://docs.langchain.com/langsmith/quick-start-studio",
|
||||
"/cloud/how-tos/invoke_studio": "https://docs.langchain.com/langsmith/use-studio",
|
||||
"/cloud/how-tos/studio/manage_assistants": "https://docs.langchain.com/langsmith/use-studio",
|
||||
"/cloud/how-tos/threads_studio": "https://docs.langchain.com/langsmith/use-threads",
|
||||
"/cloud/how-tos/iterate_graph_studio": "https://docs.langchain.com/langsmith/use-studio",
|
||||
"/cloud/how-tos/studio/run_evals": "https://docs.langchain.com/langsmith/observability",
|
||||
"/cloud/how-tos/clone_traces_studio": "https://docs.langchain.com/langsmith/observability",
|
||||
"/cloud/how-tos/datasets_studio": "https://docs.langchain.com/langsmith/use-studio",
|
||||
"/concepts/sdk": "https://docs.langchain.com/langsmith/sdk",
|
||||
"/concepts/plans": "https://docs.langchain.com/langsmith/home",
|
||||
"/concepts/application_structure": "https://docs.langchain.com/langsmith/application-structure",
|
||||
"/concepts/scalability_and_resilience": "https://docs.langchain.com/langsmith/scalability-and-resilience",
|
||||
"/concepts/auth": "https://docs.langchain.com/langsmith/auth",
|
||||
"/how-tos/auth/custom_auth": "https://docs.langchain.com/langsmith/custom-auth",
|
||||
"/how-tos/auth/openapi_security": "https://docs.langchain.com/langsmith/openapi-security",
|
||||
"/concepts/assistants": "https://docs.langchain.com/langsmith/assistants",
|
||||
"/cloud/how-tos/configuration_cloud": "https://docs.langchain.com/langsmith/configuration-cloud",
|
||||
"/cloud/how-tos/use_threads": "https://docs.langchain.com/langsmith/use-threads",
|
||||
"/cloud/how-tos/background_run": "https://docs.langchain.com/langsmith/background-run",
|
||||
"/cloud/how-tos/same-thread": "https://docs.langchain.com/langsmith/same-thread",
|
||||
"/cloud/how-tos/stateless_runs": "https://docs.langchain.com/langsmith/stateless-runs",
|
||||
"/cloud/how-tos/configurable_headers": "https://docs.langchain.com/langsmith/configurable-headers",
|
||||
"/concepts/double_texting": "https://docs.langchain.com/langsmith/double-texting",
|
||||
"/cloud/how-tos/interrupt_concurrent": "https://docs.langchain.com/langsmith/interrupt-concurrent",
|
||||
"/cloud/how-tos/rollback_concurrent": "https://docs.langchain.com/langsmith/rollback-concurrent",
|
||||
"/cloud/how-tos/reject_concurrent": "https://docs.langchain.com/langsmith/reject-concurrent",
|
||||
"/cloud/how-tos/enqueue_concurrent": "https://docs.langchain.com/langsmith/enqueue-concurrent",
|
||||
"/cloud/concepts/webhooks": "https://docs.langchain.com/langsmith/use-webhooks",
|
||||
"/cloud/how-tos/webhooks": "https://docs.langchain.com/langsmith/use-webhooks",
|
||||
"/cloud/concepts/cron_jobs": "https://docs.langchain.com/langsmith/cron-jobs",
|
||||
"/cloud/how-tos/cron_jobs": "https://docs.langchain.com/langsmith/cron-jobs",
|
||||
"/how-tos/http/custom_lifespan": "https://docs.langchain.com/langsmith/custom-lifespan",
|
||||
"/how-tos/http/custom_middleware": "https://docs.langchain.com/langsmith/custom-middleware",
|
||||
"/how-tos/http/custom_routes": "https://docs.langchain.com/langsmith/custom-routes",
|
||||
"/cloud/concepts/data_storage_and_privacy": "https://docs.langchain.com/langsmith/data-storage-and-privacy",
|
||||
"/cloud/deployment/semantic_search": "https://docs.langchain.com/langsmith/semantic-search",
|
||||
"/how-tos/ttl/configure_ttl": "https://docs.langchain.com/langsmith/configure-ttl",
|
||||
"/concepts/deployment_options": "https://docs.langchain.com/langsmith/deployments",
|
||||
"/cloud/quick_start": "https://docs.langchain.com/langsmith/deployment-quickstart",
|
||||
"/cloud/deployment/setup": "https://docs.langchain.com/langsmith/setup-app-requirements-txt",
|
||||
"/cloud/deployment/setup_pyproject": "https://docs.langchain.com/langsmith/setup-pyproject",
|
||||
"/cloud/deployment/setup_javascript": "https://docs.langchain.com/langsmith/setup-javascript",
|
||||
"/cloud/deployment/custom_docker": "https://docs.langchain.com/langsmith/custom-docker",
|
||||
"/cloud/deployment/graph_rebuild": "https://docs.langchain.com/langsmith/graph-rebuild",
|
||||
"/concepts/langgraph_cloud": "https://docs.langchain.com/langsmith/cloud",
|
||||
"/concepts/langgraph_self_hosted_data_plane": "https://docs.langchain.com/langsmith/platform-setup",
|
||||
"/concepts/langgraph_self_hosted_control_plane": "https://docs.langchain.com/langsmith/platform-setup",
|
||||
"/concepts/langgraph_standalone_container": "https://docs.langchain.com/langsmith/docker",
|
||||
"/cloud/deployment/cloud": "https://docs.langchain.com/langsmith/cloud",
|
||||
"/cloud/deployment/self_hosted_data_plane": "https://docs.langchain.com/langsmith/platform-setup",
|
||||
"/cloud/deployment/self_hosted_control_plane": "https://docs.langchain.com/langsmith/platform-setup",
|
||||
"/cloud/deployment/standalone_container": "https://docs.langchain.com/langsmith/docker",
|
||||
"/concepts/server-mcp": "https://docs.langchain.com/langsmith/server-mcp",
|
||||
"/cloud/how-tos/human_in_the_loop_time_travel": "https://docs.langchain.com/langsmith/human-in-the-loop-time-travel",
|
||||
"/cloud/how-tos/add-human-in-the-loop": "https://docs.langchain.com/langsmith/add-human-in-the-loop",
|
||||
"/cloud/deployment/egress": "https://docs.langchain.com/langsmith/env-var",
|
||||
"/cloud/how-tos/streaming": "https://docs.langchain.com/langsmith/streaming",
|
||||
"/cloud/reference/api/api_ref": "https://docs.langchain.com/langsmith/server-api-ref",
|
||||
"/cloud/reference/langgraph_server_changelog": "https://docs.langchain.com/langsmith/agent-server-changelog",
|
||||
"/cloud/reference/api/api_ref_control_plane": "https://docs.langchain.com/langsmith/api-ref-control-plane",
|
||||
"/cloud/reference/cli": "https://docs.langchain.com/langsmith/cli",
|
||||
"/cloud/reference/env_var": "https://docs.langchain.com/langsmith/env-var",
|
||||
"/troubleshooting/studio": "https://docs.langchain.com/langsmith/troubleshooting-studio",
|
||||
"/index": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/agents/agents": "https://docs.langchain.com/oss/python/langchain/agents",
|
||||
"/concepts/why-langgraph": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/get-started/1-build-basic-chatbot": "https://docs.langchain.com/oss/python/langgraph/quickstart",
|
||||
"/tutorials/get-started/2-add-tools": "https://docs.langchain.com/oss/python/langgraph/quickstart",
|
||||
"/tutorials/get-started/3-add-memory": "https://docs.langchain.com/oss/python/langgraph/quickstart",
|
||||
"/tutorials/get-started/4-human-in-the-loop": "https://docs.langchain.com/oss/python/langgraph/quickstart",
|
||||
"/tutorials/get-started/5-customize-state": "https://docs.langchain.com/oss/python/langgraph/quickstart",
|
||||
"/tutorials/get-started/6-time-travel": "https://docs.langchain.com/oss/python/langgraph/quickstart",
|
||||
"/tutorials/langsmith/local-server": "https://docs.langchain.com/oss/python/langgraph/local-server",
|
||||
"/tutorials/workflows": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/tutorials/plan-and-execute/plan-and-execute": "https://docs.langchain.com/oss/python/langchain/middleware/built-in#to-do-list",
|
||||
"/tutorials/langgraph-platform/local-server/local-server": "https://docs.langchain.com/langsmith/local-server",
|
||||
"/concepts/agentic_concepts": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/agents/overview": "https://docs.langchain.com/oss/python/langchain/agents",
|
||||
"/agents/run_agents": "https://docs.langchain.com/oss/python/langgraph/quickstart",
|
||||
"/concepts/low_level": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/how-tos/graph-api": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/how-tos/react-agent-from-scratch": "https://docs.langchain.com/oss/python/langchain/quickstart",
|
||||
"/concepts/functional_api": "https://docs.langchain.com/oss/python/langgraph/functional-api",
|
||||
"/how-tos/use-functional-api": "https://docs.langchain.com/oss/python/langgraph/functional-api",
|
||||
"/concepts/pregel": "https://docs.langchain.com/oss/python/langgraph/pregel",
|
||||
"/concepts/streaming": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/streaming": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/concepts/persistence": "https://docs.langchain.com/oss/python/langgraph/persistence",
|
||||
"/concepts/durable_execution": "https://docs.langchain.com/oss/python/langgraph/durable-execution",
|
||||
"/concepts/memory": "https://docs.langchain.com/oss/python/langgraph/memory",
|
||||
"/how-tos/memory/add-memory": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/agents/context": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/agents/models": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/concepts/tools": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/how-tos/tool-calling": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/concepts/human_in_the_loop": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/how-tos/human_in_the_loop/add-human-in-the-loop": "https://docs.langchain.com/oss/python/langgraph/interrupts",
|
||||
"/concepts/time-travel": "https://docs.langchain.com/oss/python/langgraph/persistence",
|
||||
"/how-tos/human_in_the_loop/time-travel": "https://docs.langchain.com/oss/python/langgraph/use-time-travel",
|
||||
"/concepts/subgraphs": "https://docs.langchain.com/oss/python/langgraph/use-subgraphs",
|
||||
"/how-tos/subgraph": "https://docs.langchain.com/oss/python/langgraph/use-subgraphs",
|
||||
"/concepts/multi_agent": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/agents/multi-agent": "https://docs.langchain.com/oss/python/langchain/multi-agent",
|
||||
"/how-tos/multi_agent": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/concepts/mcp": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/agents/mcp": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/concepts/tracing": "https://docs.langchain.com/oss/python/langgraph/observability",
|
||||
"/how-tos/enable-tracing": "https://docs.langchain.com/oss/python/langgraph/observability",
|
||||
"/agents/evals": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/concepts/template_applications": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/rag/langgraph_agentic_rag": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/tutorials/multi_agent/agent_supervisor": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/tutorials/sql/sql-agent": "https://docs.langchain.com/oss/python/langgraph/sql-agent",
|
||||
"/agents/ui": "https://docs.langchain.com/oss/python/langgraph/ui",
|
||||
"/how-tos/run-id-langsmith": "https://docs.langchain.com/oss/python/langgraph/observability",
|
||||
"/troubleshooting/errors/index": "https://docs.langchain.com/oss/python/langgraph/common-errors",
|
||||
"/troubleshooting/errors/INVALID_CHAT_HISTORY": "https://docs.langchain.com/oss/python/langgraph/INVALID_CHAT_HISTORY",
|
||||
"/troubleshooting/errors/INVALID_LICENSE": "https://docs.langchain.com/oss/python/langgraph/common-errors",
|
||||
"/adopters": "https://docs.langchain.com/oss/python/langgraph/case-studies",
|
||||
"/concepts/faq": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/agents/prebuilt": "https://docs.langchain.com/oss/python/langchain/agents",
|
||||
"/reference/index": "https://reference.langchain.com/python/langgraph/",
|
||||
"/reference/graphs": "https://reference.langchain.com/python/langgraph/graphs/",
|
||||
"/reference/func": "https://reference.langchain.com/python/langgraph/func/",
|
||||
"/reference/pregel": "https://reference.langchain.com/python/langgraph/pregel/",
|
||||
"/reference/checkpoints": "https://reference.langchain.com/python/langgraph/checkpoints/",
|
||||
"/reference/store": "https://reference.langchain.com/python/langgraph/store/",
|
||||
"/reference/cache": "https://reference.langchain.com/python/langgraph/cache/",
|
||||
"/reference/types": "https://reference.langchain.com/python/langgraph/types/",
|
||||
"/reference/runtime": "https://reference.langchain.com/python/langgraph/runtime/",
|
||||
"/reference/config": "https://reference.langchain.com/python/langgraph/config/",
|
||||
"/reference/errors": "https://reference.langchain.com/python/langgraph/errors/",
|
||||
"/reference/constants": "https://reference.langchain.com/python/langgraph/constants/",
|
||||
"/reference/channels": "https://reference.langchain.com/python/langgraph/channels/",
|
||||
"/reference/agents": "https://reference.langchain.com/python/langgraph/agents/",
|
||||
"/reference/supervisor": "https://reference.langchain.com/python/langgraph/supervisor/",
|
||||
"/reference/swarm": "https://reference.langchain.com/python/langgraph/swarm/",
|
||||
"/reference/mcp": "https://reference.langchain.com/python/langgraph/mcp/",
|
||||
"/cloud/reference/sdk/python_sdk_ref": "https://reference.langchain.com/python/langsmith/deployment/sdk/",
|
||||
"/reference/remote_graph": "https://reference.langchain.com/python/langsmith/deployment/remote_graph/",
|
||||
"/additional-resources/index": "https://docs.langchain.com/oss/python/langchain/overview",
|
||||
"/cloud/reference/sdk/js_ts_sdk_ref": "https://reference.langchain.com/javascript/modules/langsmith.html",
|
||||
"/snippets/chat_model_tabs": "https://docs.langchain.com/oss/python/langchain/overview",
|
||||
"/troubleshooting/errors/GRAPH_RECURSION_LIMIT": "https://docs.langchain.com/oss/python/langgraph/GRAPH_RECURSION_LIMIT",
|
||||
"/troubleshooting/errors/INVALID_CONCURRENT_GRAPH_UPDATE": "https://docs.langchain.com/oss/python/langgraph/INVALID_CONCURRENT_GRAPH_UPDATE",
|
||||
"/troubleshooting/errors/INVALID_GRAPH_NODE_RETURN_VALUE": "https://docs.langchain.com/oss/python/langgraph/INVALID_GRAPH_NODE_RETURN_VALUE",
|
||||
"/troubleshooting/errors/MULTIPLE_SUBGRAPHS": "https://docs.langchain.com/oss/python/langgraph/MULTIPLE_SUBGRAPHS",
|
||||
"/tutorials/rag/langgraph_self_rag": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/additional-resources": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/examples": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/guides": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/how-tos/autogen-integration-functional": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/how-tos/cross-thread-persistence-functional": "https://docs.langchain.com/oss/python/langgraph/add-memory#add-long-term-memory",
|
||||
"/how-tos/disable-streaming": "https://docs.langchain.com/oss/python/langgraph/streaming",
|
||||
"/how-tos/memory/semantic-search": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/multi-agent-multi-turn-convo-functional": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/how-tos/multi-agent-network-functional": "https://docs.langchain.com/oss/python/langgraph/graph-api",
|
||||
"/how-tos/persistence-functional": "https://docs.langchain.com/oss/python/langgraph/add-memory",
|
||||
"/how-tos/react-agent-from-scratch-functional": "https://docs.langchain.com/oss/python/langgraph/workflows-agents",
|
||||
"/reference": "https://reference.langchain.com/python/langgraph/",
|
||||
"/troubleshooting/errors": "https://docs.langchain.com/oss/python/langgraph/common-errors",
|
||||
"/tutorials/chatbot-simulation-evaluation/agent-simulation-evaluation": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/chatbot-simulation-evaluation/langsmith-agent-simulation-evaluation": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/chatbots/information-gather-prompting": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/extraction/retries": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/langgraph-platform/local-server": "https://docs.langchain.com/langsmith/agent-server",
|
||||
"/tutorials/lats/lats": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/llm-compiler/LLMCompiler": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/rag/langgraph_adaptive_rag_local": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/tutorials/rag/langgraph_crag": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/tutorials/rag/langgraph_crag_local": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/tutorials/rag/langgraph_self_rag_local": "https://docs.langchain.com/oss/python/langgraph/agentic-rag",
|
||||
"/tutorials/reflection/reflection": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/reflexion/reflexion": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/rewoo/rewoo": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/self-discover/self-discover": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/tnt-llm/tnt-llm": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/tot/tot": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/usaco/usaco": "https://docs.langchain.com/oss/python/langgraph/overview",
|
||||
"/tutorials/web-navigation/web_voyager": "https://docs.langchain.com/oss/python/langgraph/overview"
|
||||
}
|
||||
@@ -57,3 +57,9 @@ lint.select = [
|
||||
]
|
||||
lint.ignore = ["E501", "B008"]
|
||||
target-version = "py310"
|
||||
|
||||
[[tool.uv.index]]
|
||||
name = "testpypi"
|
||||
url = "https://test.pypi.org/simple/"
|
||||
publish-url = "https://test.pypi.org/legacy/"
|
||||
explicit = true
|
||||
|
||||
@@ -390,6 +390,142 @@ class PostgresSaver(BasePostgresSaver):
|
||||
(str(thread_id),),
|
||||
)
|
||||
|
||||
def delete_for_runs(self, run_ids: Sequence[str]) -> None:
|
||||
"""Delete all checkpoints and writes for the given run IDs.
|
||||
|
||||
Args:
|
||||
run_ids: The run IDs whose checkpoints should be deleted.
|
||||
"""
|
||||
if not run_ids:
|
||||
return
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE (thread_id, checkpoint_ns, checkpoint_id) IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)""",
|
||||
(list(run_ids),),
|
||||
)
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE (thread_id, checkpoint_ns) IN (
|
||||
SELECT DISTINCT thread_id, checkpoint_ns
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c2
|
||||
WHERE c2.thread_id = checkpoint_blobs.thread_id
|
||||
AND c2.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
AND (c2.metadata->>'run_id' IS NULL OR c2.metadata->>'run_id' != ALL(%s))
|
||||
AND c2.checkpoint->'channel_versions' ? checkpoint_blobs.channel
|
||||
)""",
|
||||
(list(run_ids), list(run_ids)),
|
||||
)
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoints WHERE metadata->>'run_id' = ANY(%s)",
|
||||
(list(run_ids),),
|
||||
)
|
||||
|
||||
def copy_thread(self, source_thread_id: str, target_thread_id: str) -> None:
|
||||
"""Copy all checkpoints and writes from source thread to target thread.
|
||||
|
||||
Args:
|
||||
source_thread_id: The thread ID to copy from.
|
||||
target_thread_id: The thread ID to copy to.
|
||||
"""
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"""INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata
|
||||
FROM checkpoints
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
cur.execute(
|
||||
"""INSERT INTO checkpoint_blobs (thread_id, checkpoint_ns, channel, version, type, blob)
|
||||
SELECT %s, checkpoint_ns, channel, version, type, blob
|
||||
FROM checkpoint_blobs
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, channel, version) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
cur.execute(
|
||||
"""INSERT INTO checkpoint_writes (thread_id, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob
|
||||
FROM checkpoint_writes
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
|
||||
def prune(
|
||||
self,
|
||||
thread_ids: Sequence[str],
|
||||
*,
|
||||
strategy: str = "keep_latest",
|
||||
) -> None:
|
||||
"""Prune checkpoints for given threads.
|
||||
|
||||
Args:
|
||||
thread_ids: The thread IDs to prune.
|
||||
strategy: `"keep_latest"` keeps only the latest checkpoint per
|
||||
thread+namespace. `"delete"` removes everything.
|
||||
"""
|
||||
if not thread_ids:
|
||||
return
|
||||
if strategy == "delete":
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoints WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoint_blobs WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
cur.execute(
|
||||
"DELETE FROM checkpoint_writes WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
elif strategy == "keep_latest":
|
||||
with self._cursor(pipeline=True) as cur:
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT DISTINCT ON (thread_id, checkpoint_ns)
|
||||
thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
ORDER BY thread_id, checkpoint_ns, checkpoint_id DESC
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c
|
||||
WHERE c.thread_id = checkpoint_blobs.thread_id
|
||||
AND c.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
)""",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
|
||||
@contextmanager
|
||||
def _cursor(self, *, pipeline: bool = False) -> Iterator[Cursor[DictRow]]:
|
||||
"""Create a database cursor as a context manager.
|
||||
|
||||
@@ -349,6 +349,152 @@ class AsyncPostgresSaver(BasePostgresSaver):
|
||||
(str(thread_id),),
|
||||
)
|
||||
|
||||
async def adelete_for_runs(self, run_ids: Sequence[str]) -> None:
|
||||
"""Delete all checkpoints and writes for the given run IDs.
|
||||
|
||||
Args:
|
||||
run_ids: The run IDs whose checkpoints should be deleted.
|
||||
"""
|
||||
if not run_ids:
|
||||
return
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
# Delete writes associated with checkpoints that have matching run_ids
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE (thread_id, checkpoint_ns, checkpoint_id) IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)""",
|
||||
(list(run_ids),),
|
||||
)
|
||||
# Delete blobs associated with checkpoints that have matching run_ids
|
||||
# We need to delete blobs for channels referenced by these checkpoints
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE (thread_id, checkpoint_ns) IN (
|
||||
SELECT DISTINCT thread_id, checkpoint_ns
|
||||
FROM checkpoints
|
||||
WHERE metadata->>'run_id' = ANY(%s)
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c2
|
||||
WHERE c2.thread_id = checkpoint_blobs.thread_id
|
||||
AND c2.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
AND (c2.metadata->>'run_id' IS NULL OR c2.metadata->>'run_id' != ALL(%s))
|
||||
AND c2.checkpoint->'channel_versions' ? checkpoint_blobs.channel
|
||||
)""",
|
||||
(list(run_ids), list(run_ids)),
|
||||
)
|
||||
# Delete the checkpoints themselves
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoints WHERE metadata->>'run_id' = ANY(%s)",
|
||||
(list(run_ids),),
|
||||
)
|
||||
|
||||
async def acopy_thread(self, source_thread_id: str, target_thread_id: str) -> None:
|
||||
"""Copy all checkpoints and writes from source thread to target thread.
|
||||
|
||||
Args:
|
||||
source_thread_id: The thread ID to copy from.
|
||||
target_thread_id: The thread ID to copy to.
|
||||
"""
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
# Copy checkpoints
|
||||
await cur.execute(
|
||||
"""INSERT INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, parent_checkpoint_id, checkpoint, metadata
|
||||
FROM checkpoints
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
# Copy blobs
|
||||
await cur.execute(
|
||||
"""INSERT INTO checkpoint_blobs (thread_id, checkpoint_ns, channel, version, type, blob)
|
||||
SELECT %s, checkpoint_ns, channel, version, type, blob
|
||||
FROM checkpoint_blobs
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, channel, version) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
# Copy writes
|
||||
await cur.execute(
|
||||
"""INSERT INTO checkpoint_writes (thread_id, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob)
|
||||
SELECT %s, checkpoint_ns, checkpoint_id, task_id, task_path, idx, channel, type, blob
|
||||
FROM checkpoint_writes
|
||||
WHERE thread_id = %s
|
||||
ON CONFLICT (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) DO NOTHING""",
|
||||
(target_thread_id, source_thread_id),
|
||||
)
|
||||
|
||||
async def aprune(
|
||||
self,
|
||||
thread_ids: Sequence[str],
|
||||
*,
|
||||
strategy: str = "keep_latest",
|
||||
) -> None:
|
||||
"""Prune checkpoints for given threads.
|
||||
|
||||
Args:
|
||||
thread_ids: The thread IDs to prune.
|
||||
strategy: `"keep_latest"` keeps only the latest checkpoint per
|
||||
thread+namespace. `"delete"` removes everything.
|
||||
"""
|
||||
if not thread_ids:
|
||||
return
|
||||
if strategy == "delete":
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoints WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoint_blobs WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
await cur.execute(
|
||||
"DELETE FROM checkpoint_writes WHERE thread_id = ANY(%s)",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
elif strategy == "keep_latest":
|
||||
async with self._cursor(pipeline=True) as cur:
|
||||
# Delete non-latest checkpoints
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT DISTINCT ON (thread_id, checkpoint_ns)
|
||||
thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
ORDER BY thread_id, checkpoint_ns, checkpoint_id DESC
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
# Delete writes for removed checkpoints (keep only writes for remaining checkpoints)
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_writes
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND (thread_id, checkpoint_ns, checkpoint_id) NOT IN (
|
||||
SELECT thread_id, checkpoint_ns, checkpoint_id
|
||||
FROM checkpoints
|
||||
WHERE thread_id = ANY(%s)
|
||||
)""",
|
||||
(list(thread_ids), list(thread_ids)),
|
||||
)
|
||||
# Clean up orphaned blobs
|
||||
await cur.execute(
|
||||
"""DELETE FROM checkpoint_blobs
|
||||
WHERE thread_id = ANY(%s)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM checkpoints c
|
||||
WHERE c.thread_id = checkpoint_blobs.thread_id
|
||||
AND c.checkpoint_ns = checkpoint_blobs.checkpoint_ns
|
||||
)""",
|
||||
(list(thread_ids),),
|
||||
)
|
||||
|
||||
@asynccontextmanager
|
||||
async def _cursor(
|
||||
self, *, pipeline: bool = False
|
||||
@@ -578,5 +724,43 @@ class AsyncPostgresSaver(BasePostgresSaver):
|
||||
self.adelete_thread(thread_id), self.loop
|
||||
).result()
|
||||
|
||||
def delete_for_runs(self, run_ids: Sequence[str]) -> None:
|
||||
"""Delete all checkpoints and writes for the given run IDs.
|
||||
|
||||
Args:
|
||||
run_ids: The run IDs whose checkpoints should be deleted.
|
||||
"""
|
||||
return asyncio.run_coroutine_threadsafe(
|
||||
self.adelete_for_runs(run_ids), self.loop
|
||||
).result()
|
||||
|
||||
def copy_thread(self, source_thread_id: str, target_thread_id: str) -> None:
|
||||
"""Copy all checkpoints and writes from source thread to target thread.
|
||||
|
||||
Args:
|
||||
source_thread_id: The thread ID to copy from.
|
||||
target_thread_id: The thread ID to copy to.
|
||||
"""
|
||||
return asyncio.run_coroutine_threadsafe(
|
||||
self.acopy_thread(source_thread_id, target_thread_id), self.loop
|
||||
).result()
|
||||
|
||||
def prune(
|
||||
self,
|
||||
thread_ids: Sequence[str],
|
||||
*,
|
||||
strategy: str = "keep_latest",
|
||||
) -> None:
|
||||
"""Prune checkpoints for given threads.
|
||||
|
||||
Args:
|
||||
thread_ids: The thread IDs to prune.
|
||||
strategy: `"keep_latest"` keeps only the latest checkpoint per
|
||||
thread+namespace. `"delete"` removes everything.
|
||||
"""
|
||||
return asyncio.run_coroutine_threadsafe(
|
||||
self.aprune(thread_ids, strategy=strategy), self.loop
|
||||
).result()
|
||||
|
||||
|
||||
__all__ = ["AsyncPostgresSaver", "AsyncShallowPostgresSaver", "Conn"]
|
||||
|
||||
@@ -13,7 +13,7 @@ license = "MIT"
|
||||
license-files = ['LICENSE']
|
||||
dependencies = [
|
||||
"langgraph-checkpoint>=2.1.2,<5.0.0",
|
||||
"orjson>=3.10.1",
|
||||
"orjson>=3.11.5",
|
||||
"psycopg>=3.2.0",
|
||||
"psycopg-pool>=3.2.0",
|
||||
]
|
||||
@@ -32,6 +32,7 @@ test = [
|
||||
"pytest-mock",
|
||||
"psycopg[binary]",
|
||||
"langgraph-checkpoint",
|
||||
"langgraph-checkpoint-conformance",
|
||||
"pytest-watcher",
|
||||
]
|
||||
lint = [
|
||||
@@ -49,6 +50,7 @@ default-groups = ['dev']
|
||||
|
||||
[tool.uv.sources]
|
||||
langgraph-checkpoint = { path = "../checkpoint", editable = true }
|
||||
langgraph-checkpoint-conformance = { path = "../checkpoint-conformance", editable = true }
|
||||
|
||||
[tool.hatch.build.targets.wheel]
|
||||
include = ["langgraph"]
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
"""Conformance tests for AsyncPostgresSaver."""
|
||||
# mypy: disable-error-code="import-untyped"
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import AsyncGenerator
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
from langgraph.checkpoint.conformance import checkpointer_test, validate
|
||||
from langgraph.checkpoint.conformance.report import ProgressCallbacks
|
||||
from psycopg import AsyncConnection
|
||||
from psycopg.rows import dict_row
|
||||
|
||||
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
|
||||
from tests.conftest import DEFAULT_POSTGRES_URI
|
||||
|
||||
|
||||
async def pg_lifespan() -> AsyncGenerator[None, None]:
|
||||
"""No-op lifespan; databases are created per-checkpointer instance."""
|
||||
yield
|
||||
|
||||
|
||||
@checkpointer_test(name="AsyncPostgresSaver", lifespan=pg_lifespan)
|
||||
async def postgres_checkpointer() -> AsyncGenerator[AsyncPostgresSaver, None]:
|
||||
database = f"test_{uuid4().hex[:16]}"
|
||||
async with await AsyncConnection.connect(
|
||||
DEFAULT_POSTGRES_URI, autocommit=True
|
||||
) as conn:
|
||||
await conn.execute(f"CREATE DATABASE {database}")
|
||||
try:
|
||||
async with await AsyncConnection.connect(
|
||||
DEFAULT_POSTGRES_URI + database,
|
||||
autocommit=True,
|
||||
prepare_threshold=0,
|
||||
row_factory=dict_row,
|
||||
) as conn:
|
||||
saver = AsyncPostgresSaver(conn)
|
||||
await saver.setup()
|
||||
yield saver
|
||||
finally:
|
||||
async with await AsyncConnection.connect(
|
||||
DEFAULT_POSTGRES_URI, autocommit=True
|
||||
) as conn:
|
||||
await conn.execute(f"DROP DATABASE {database}")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_full_conformance() -> None:
|
||||
"""AsyncPostgresSaver passes ALL conformance tests."""
|
||||
report = await validate(
|
||||
postgres_checkpointer,
|
||||
progress=ProgressCallbacks.verbose(),
|
||||
)
|
||||
report.print_report()
|
||||
assert report.passed_all(), f"Conformance failed: {report.to_dict()}"
|
||||
Generated
+32
-1
@@ -304,6 +304,33 @@ test = [
|
||||
{ name = "redis" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint-conformance"
|
||||
version = "0.0.1"
|
||||
source = { editable = "../checkpoint-conformance" }
|
||||
dependencies = [
|
||||
{ name = "langgraph-checkpoint" },
|
||||
]
|
||||
|
||||
[package.metadata]
|
||||
requires-dist = [{ name = "langgraph-checkpoint", specifier = ">=2.0.0" }]
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
dev = [
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
{ name = "ruff" },
|
||||
{ name = "ty" },
|
||||
]
|
||||
lint = [
|
||||
{ name = "ruff" },
|
||||
{ name = "ty" },
|
||||
]
|
||||
test = [
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-checkpoint-postgres"
|
||||
version = "3.0.4"
|
||||
@@ -320,6 +347,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extra = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -336,6 +364,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance" },
|
||||
{ name = "psycopg", extra = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
@@ -346,7 +375,7 @@ test = [
|
||||
[package.metadata]
|
||||
requires-dist = [
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "orjson", specifier = ">=3.10.1" },
|
||||
{ name = "orjson", specifier = ">=3.11.5" },
|
||||
{ name = "psycopg", specifier = ">=3.2.0" },
|
||||
{ name = "psycopg-pool", specifier = ">=3.2.0" },
|
||||
]
|
||||
@@ -356,6 +385,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -372,6 +402,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
|
||||
Generated
+6
-6
@@ -267,7 +267,7 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "langchain-core"
|
||||
version = "1.2.12"
|
||||
version = "1.2.13"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "jsonpatch" },
|
||||
@@ -279,9 +279,9 @@ dependencies = [
|
||||
{ name = "typing-extensions" },
|
||||
{ name = "uuid-utils" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/2a/1d/08e935d1532fcc90981f6e5bb6825914c9227ea7a962c62b1e18619b49e7/langchain_core-1.2.12.tar.gz", hash = "sha256:4d7fa6643d7ab06fc1905a9b7dcbe96a6f3c181046b56edf9c0c17ecd412d9e9", size = 831329, upload-time = "2026-02-12T20:53:15.01Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/fb/bb/c501ca60556c11ac80d1454bdcac63cb33583ce4e64fc4535ad5a7d5c6ba/langchain_core-1.2.13.tar.gz", hash = "sha256:d2773d0d0130a356378db9a858cfeef64c3d64bc03722f1d4d6c40eb46fdf01b", size = 831612, upload-time = "2026-02-15T07:45:57.014Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/8c/a5/678ab0e5cc57794f20ae5ed12c1442506ef1108c9434f950aebc6044e5a3/langchain_core-1.2.12-py3-none-any.whl", hash = "sha256:66ca17a2a9cb007ab29021968e6adfcf4228067151dc2bd6ebfff265ffaf92f5", size = 500132, upload-time = "2026-02-12T20:53:13.806Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/12/ab/60fd69e5d55f67d422baefddaaca523c42cd7510ab6aeb17db6ae57fb107/langchain_core-1.2.13-py3-none-any.whl", hash = "sha256:b31823e28d3eff1e237096d0bd3bf80c6f9624eb471a9496dbfbd427779f8d82", size = 500485, upload-time = "2026-02-15T07:45:55.422Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
@@ -1198,14 +1198,14 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "redis"
|
||||
version = "7.1.1"
|
||||
version = "7.2.0"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "async-timeout", marker = "python_full_version < '3.11.3'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/f7/80/2971931d27651affa88a44c0ad7b8c4a19dc29c998abb20b23868d319b59/redis-7.1.1.tar.gz", hash = "sha256:a2814b2bda15b39dad11391cc48edac4697214a8a5a4bd10abe936ab4892eb43", size = 4800064, upload-time = "2026-02-09T18:39:40.292Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/9f/32/6fac13a11e73e1bc67a2ae821a72bfe4c2d8c4c48f0267e4a952be0f1bae/redis-7.2.0.tar.gz", hash = "sha256:4dd5bf4bd4ae80510267f14185a15cba2a38666b941aff68cccf0256b51c1f26", size = 4901247, upload-time = "2026-02-16T17:16:22.797Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/29/55/1de1d812ba1481fa4b37fb03b4eec0fcb71b6a0d44c04ea3482eb017600f/redis-7.1.1-py3-none-any.whl", hash = "sha256:f77817f16071c2950492c67d40b771fa493eb3fccc630a424a10976dbb794b7a", size = 356057, upload-time = "2026-02-09T18:39:38.602Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/86/cf/f6180b67f99688d83e15c84c5beda831d1d341e95872d224f87ccafafe61/redis-7.2.0-py3-none-any.whl", hash = "sha256:01f591f8598e483f1842d429e8ae3a820804566f1c73dca1b80e23af9fba0497", size = 394898, upload-time = "2026-02-16T17:16:20.693Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
|
||||
@@ -29,8 +29,8 @@
|
||||
"@eslint/js": "^10.0.1",
|
||||
"@tsconfig/recommended": "^1.0.13",
|
||||
"@types/jest": "^30.0.0",
|
||||
"@typescript-eslint/eslint-plugin": "^8.55.0",
|
||||
"@typescript-eslint/parser": "^8.55.0",
|
||||
"@typescript-eslint/eslint-plugin": "^8.56.0",
|
||||
"@typescript-eslint/parser": "^8.56.0",
|
||||
"dotenv": "^17.3.1",
|
||||
"eslint": "^10.0.0",
|
||||
"eslint-config-prettier": "^10.1.8",
|
||||
|
||||
@@ -1079,101 +1079,101 @@
|
||||
dependencies:
|
||||
"@types/yargs-parser" "*"
|
||||
|
||||
"@typescript-eslint/eslint-plugin@^8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.55.0.tgz#086d2ef661507b561f7b17f62d3179d692a0765f"
|
||||
integrity sha512-1y/MVSz0NglV1ijHC8OT49mPJ4qhPYjiK08YUQVbIOyu+5k862LKUHFkpKHWu//zmr7hDR2rhwUm6gnCGNmGBQ==
|
||||
"@typescript-eslint/eslint-plugin@^8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.56.0.tgz#5aec3db807a6b8437ea5d5ebf7bd16b4119aba8d"
|
||||
integrity sha512-lRyPDLzNCuae71A3t9NEINBiTn7swyOhvUj3MyUOxb8x6g6vPEFoOU+ZRmGMusNC3X3YMhqMIX7i8ShqhT74Pw==
|
||||
dependencies:
|
||||
"@eslint-community/regexpp" "^4.12.2"
|
||||
"@typescript-eslint/scope-manager" "8.55.0"
|
||||
"@typescript-eslint/type-utils" "8.55.0"
|
||||
"@typescript-eslint/utils" "8.55.0"
|
||||
"@typescript-eslint/visitor-keys" "8.55.0"
|
||||
"@typescript-eslint/scope-manager" "8.56.0"
|
||||
"@typescript-eslint/type-utils" "8.56.0"
|
||||
"@typescript-eslint/utils" "8.56.0"
|
||||
"@typescript-eslint/visitor-keys" "8.56.0"
|
||||
ignore "^7.0.5"
|
||||
natural-compare "^1.4.0"
|
||||
ts-api-utils "^2.4.0"
|
||||
|
||||
"@typescript-eslint/parser@^8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/parser/-/parser-8.55.0.tgz#6eace4e9e95f178d3447ed1f17f3d6a5dfdb345c"
|
||||
integrity sha512-4z2nCSBfVIMnbuu8uinj+f0o4qOeggYJLbjpPHka3KH1om7e+H9yLKTYgksTaHcGco+NClhhY2vyO3HsMH1RGw==
|
||||
"@typescript-eslint/parser@^8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/parser/-/parser-8.56.0.tgz#8ecff1678b8b1a742d29c446ccf5eeea7f971d72"
|
||||
integrity sha512-IgSWvLobTDOjnaxAfDTIHaECbkNlAlKv2j5SjpB2v7QHKv1FIfjwMy8FsDbVfDX/KjmCmYICcw7uGaXLhtsLNg==
|
||||
dependencies:
|
||||
"@typescript-eslint/scope-manager" "8.55.0"
|
||||
"@typescript-eslint/types" "8.55.0"
|
||||
"@typescript-eslint/typescript-estree" "8.55.0"
|
||||
"@typescript-eslint/visitor-keys" "8.55.0"
|
||||
"@typescript-eslint/scope-manager" "8.56.0"
|
||||
"@typescript-eslint/types" "8.56.0"
|
||||
"@typescript-eslint/typescript-estree" "8.56.0"
|
||||
"@typescript-eslint/visitor-keys" "8.56.0"
|
||||
debug "^4.4.3"
|
||||
|
||||
"@typescript-eslint/project-service@8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/project-service/-/project-service-8.55.0.tgz#b8a71c06a625bdad481c24d5614b68e252f3ae9b"
|
||||
integrity sha512-zRcVVPFUYWa3kNnjaZGXSu3xkKV1zXy8M4nO/pElzQhFweb7PPtluDLQtKArEOGmjXoRjnUZ29NjOiF0eCDkcQ==
|
||||
"@typescript-eslint/project-service@8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/project-service/-/project-service-8.56.0.tgz#bb8562fecd8f7922e676fc6a1189c20dd7991d73"
|
||||
integrity sha512-M3rnyL1vIQOMeWxTWIW096/TtVP+8W3p/XnaFflhmcFp+U4zlxUxWj4XwNs6HbDeTtN4yun0GNTTDBw/SvufKg==
|
||||
dependencies:
|
||||
"@typescript-eslint/tsconfig-utils" "^8.55.0"
|
||||
"@typescript-eslint/types" "^8.55.0"
|
||||
"@typescript-eslint/tsconfig-utils" "^8.56.0"
|
||||
"@typescript-eslint/types" "^8.56.0"
|
||||
debug "^4.4.3"
|
||||
|
||||
"@typescript-eslint/scope-manager@8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/scope-manager/-/scope-manager-8.55.0.tgz#8a0752c31c788651840dc98f840b0c2ebe143b8c"
|
||||
integrity sha512-fVu5Omrd3jeqeQLiB9f1YsuK/iHFOwb04bCtY4BSCLgjNbOD33ZdV6KyEqplHr+IlpgT0QTZ/iJ+wT7hvTx49Q==
|
||||
"@typescript-eslint/scope-manager@8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/scope-manager/-/scope-manager-8.56.0.tgz#604030a4c6433df3728effdd441d47f45a86edb4"
|
||||
integrity sha512-7UiO/XwMHquH+ZzfVCfUNkIXlp/yQjjnlYUyYz7pfvlK3/EyyN6BK+emDmGNyQLBtLGaYrTAI6KOw8tFucWL2w==
|
||||
dependencies:
|
||||
"@typescript-eslint/types" "8.55.0"
|
||||
"@typescript-eslint/visitor-keys" "8.55.0"
|
||||
"@typescript-eslint/types" "8.56.0"
|
||||
"@typescript-eslint/visitor-keys" "8.56.0"
|
||||
|
||||
"@typescript-eslint/tsconfig-utils@8.55.0", "@typescript-eslint/tsconfig-utils@^8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.55.0.tgz#62f1d005419985e09d37a040b2f1450e4e805afa"
|
||||
integrity sha512-1R9cXqY7RQd7WuqSN47PK9EDpgFUK3VqdmbYrvWJZYDd0cavROGn+74ktWBlmJ13NXUQKlZ/iAEQHI/V0kKe0Q==
|
||||
"@typescript-eslint/tsconfig-utils@8.56.0", "@typescript-eslint/tsconfig-utils@^8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.56.0.tgz#2538ce83cbc376e685487960cbb24b65fe2abc4e"
|
||||
integrity sha512-bSJoIIt4o3lKXD3xmDh9chZcjCz5Lk8xS7Rxn+6l5/pKrDpkCwtQNQQwZ2qRPk7TkUYhrq3WPIHXOXlbXP0itg==
|
||||
|
||||
"@typescript-eslint/type-utils@8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/type-utils/-/type-utils-8.55.0.tgz#195d854b3e56308ce475fdea2165313bb1190200"
|
||||
integrity sha512-x1iH2unH4qAt6I37I2CGlsNs+B9WGxurP2uyZLRz6UJoZWDBx9cJL1xVN/FiOmHEONEg6RIufdvyT0TEYIgC5g==
|
||||
"@typescript-eslint/type-utils@8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/type-utils/-/type-utils-8.56.0.tgz#72b4edc1fc73988998f1632b3ec99c2a66eaac6e"
|
||||
integrity sha512-qX2L3HWOU2nuDs6GzglBeuFXviDODreS58tLY/BALPC7iu3Fa+J7EOTwnX9PdNBxUI7Uh0ntP0YWGnxCkXzmfA==
|
||||
dependencies:
|
||||
"@typescript-eslint/types" "8.55.0"
|
||||
"@typescript-eslint/typescript-estree" "8.55.0"
|
||||
"@typescript-eslint/utils" "8.55.0"
|
||||
"@typescript-eslint/types" "8.56.0"
|
||||
"@typescript-eslint/typescript-estree" "8.56.0"
|
||||
"@typescript-eslint/utils" "8.56.0"
|
||||
debug "^4.4.3"
|
||||
ts-api-utils "^2.4.0"
|
||||
|
||||
"@typescript-eslint/types@8.55.0", "@typescript-eslint/types@^8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/types/-/types-8.55.0.tgz#8449c5a7adac61184cac92dbf6315733569708c2"
|
||||
integrity sha512-ujT0Je8GI5BJWi+/mMoR0wxwVEQaxM+pi30xuMiJETlX80OPovb2p9E8ss87gnSVtYXtJoU9U1Cowcr6w2FE0w==
|
||||
"@typescript-eslint/types@8.56.0", "@typescript-eslint/types@^8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/types/-/types-8.56.0.tgz#a2444011b9a98ca13d70411d2cbfed5443b3526a"
|
||||
integrity sha512-DBsLPs3GsWhX5HylbP9HNG15U0bnwut55Lx12bHB9MpXxQ+R5GC8MwQe+N1UFXxAeQDvEsEDY6ZYwX03K7Z6HQ==
|
||||
|
||||
"@typescript-eslint/typescript-estree@8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/typescript-estree/-/typescript-estree-8.55.0.tgz#c83ac92c11ce79bedd984937c7780a65e7f7b2e3"
|
||||
integrity sha512-EwrH67bSWdx/3aRQhCoxDaHM+CrZjotc2UCCpEDVqfCE+7OjKAGWNY2HsCSTEVvWH2clYQK8pdeLp42EVs+xQw==
|
||||
"@typescript-eslint/typescript-estree@8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/typescript-estree/-/typescript-estree-8.56.0.tgz#fadbc74c14c5bac947db04980ff58bb178701c2e"
|
||||
integrity sha512-ex1nTUMWrseMltXUHmR2GAQ4d+WjkZCT4f+4bVsps8QEdh0vlBsaCokKTPlnqBFqqGaxilDNJG7b8dolW2m43Q==
|
||||
dependencies:
|
||||
"@typescript-eslint/project-service" "8.55.0"
|
||||
"@typescript-eslint/tsconfig-utils" "8.55.0"
|
||||
"@typescript-eslint/types" "8.55.0"
|
||||
"@typescript-eslint/visitor-keys" "8.55.0"
|
||||
"@typescript-eslint/project-service" "8.56.0"
|
||||
"@typescript-eslint/tsconfig-utils" "8.56.0"
|
||||
"@typescript-eslint/types" "8.56.0"
|
||||
"@typescript-eslint/visitor-keys" "8.56.0"
|
||||
debug "^4.4.3"
|
||||
minimatch "^9.0.5"
|
||||
semver "^7.7.3"
|
||||
tinyglobby "^0.2.15"
|
||||
ts-api-utils "^2.4.0"
|
||||
|
||||
"@typescript-eslint/utils@8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/utils/-/utils-8.55.0.tgz#c1744d94a3901deb01f58b09d3478d811f96d619"
|
||||
integrity sha512-BqZEsnPGdYpgyEIkDC1BadNY8oMwckftxBT+C8W0g1iKPdeqKZBtTfnvcq0nf60u7MkjFO8RBvpRGZBPw4L2ow==
|
||||
"@typescript-eslint/utils@8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/utils/-/utils-8.56.0.tgz#063ce6f702ec603de1b83ee795ed5e877d6f7841"
|
||||
integrity sha512-RZ3Qsmi2nFGsS+n+kjLAYDPVlrzf7UhTffrDIKr+h2yzAlYP/y5ZulU0yeDEPItos2Ph46JAL5P/On3pe7kDIQ==
|
||||
dependencies:
|
||||
"@eslint-community/eslint-utils" "^4.9.1"
|
||||
"@typescript-eslint/scope-manager" "8.55.0"
|
||||
"@typescript-eslint/types" "8.55.0"
|
||||
"@typescript-eslint/typescript-estree" "8.55.0"
|
||||
"@typescript-eslint/scope-manager" "8.56.0"
|
||||
"@typescript-eslint/types" "8.56.0"
|
||||
"@typescript-eslint/typescript-estree" "8.56.0"
|
||||
|
||||
"@typescript-eslint/visitor-keys@8.55.0":
|
||||
version "8.55.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/visitor-keys/-/visitor-keys-8.55.0.tgz#3d9a40fd4e3705c63d8fae3af58988add3ed464d"
|
||||
integrity sha512-AxNRwEie8Nn4eFS1FzDMJWIISMGoXMb037sgCBJ3UR6o0fQTzr2tqN9WT+DkWJPhIdQCfV7T6D387566VtnCJA==
|
||||
"@typescript-eslint/visitor-keys@8.56.0":
|
||||
version "8.56.0"
|
||||
resolved "https://registry.yarnpkg.com/@typescript-eslint/visitor-keys/-/visitor-keys-8.56.0.tgz#7d6592ab001827d3ce052155edf7ecad19688d7d"
|
||||
integrity sha512-q+SL+b+05Ud6LbEE35qe4A99P+htKTKVbyiNEe45eCbJFyh/HVK9QXwlrbz+Q4L8SOW4roxSVwXYj4DMBT7Ieg==
|
||||
dependencies:
|
||||
"@typescript-eslint/types" "8.55.0"
|
||||
eslint-visitor-keys "^4.2.1"
|
||||
"@typescript-eslint/types" "8.56.0"
|
||||
eslint-visitor-keys "^5.0.0"
|
||||
|
||||
"@ungap/structured-clone@^1.3.0":
|
||||
version "1.3.0"
|
||||
@@ -2241,7 +2241,7 @@ eslint-visitor-keys@^3.4.3:
|
||||
resolved "https://registry.yarnpkg.com/eslint-visitor-keys/-/eslint-visitor-keys-3.4.3.tgz#0cd72fe8550e3c2eae156a96a4dddcd1c8ac5800"
|
||||
integrity sha512-wpc+LXeiyiisxPlEkUzU6svyS1frIO3Mgxj1fdy7Pm8Ygzguax2N3Fa/D/ag1WqbOprdI+uY6wMUl8/a2G+iag==
|
||||
|
||||
eslint-visitor-keys@^4.0.0, eslint-visitor-keys@^4.2.1:
|
||||
eslint-visitor-keys@^4.0.0:
|
||||
version "4.2.1"
|
||||
resolved "https://registry.yarnpkg.com/eslint-visitor-keys/-/eslint-visitor-keys-4.2.1.tgz#4cfea60fe7dd0ad8e816e1ed026c1d5251b512c1"
|
||||
integrity sha512-Uhdk5sfqcee/9H/rCOJikYz67o0a2Tw2hGRPOG2Y1R2dg7brRe1uG0yaNQDHu+TO/uQPF/5eCapvYSmHUjt7JQ==
|
||||
@@ -3533,9 +3533,9 @@ js-tokens@^4.0.0:
|
||||
integrity sha512-RdJUflcE3cUzKiMqQgsCu06FPu9UdIJO0beYbPhHN4k6apgJtifcoCtT9bcxOpYBtpD2kCM6Sbzg4CausW/PKQ==
|
||||
|
||||
js-yaml@^3.13.1:
|
||||
version "3.14.1"
|
||||
resolved "https://registry.yarnpkg.com/js-yaml/-/js-yaml-3.14.1.tgz#dae812fdb3825fa306609a8717383c50c36a0537"
|
||||
integrity sha512-okMH7OXXJ7YrN9Ok3/SXrnu4iX9yOk+25nqX4imS2npuvTYDmo/QEZoqwZkYaIDk3jVvBOTOIEgEhaLOynBS9g==
|
||||
version "3.14.2"
|
||||
resolved "https://registry.yarnpkg.com/js-yaml/-/js-yaml-3.14.2.tgz#77485ce1dd7f33c061fd1b16ecea23b55fcb04b0"
|
||||
integrity sha512-PMSmkqxr106Xa156c2M265Z+FTrPl+oxd/rgOQy2tijQeK5TxQ43psO1ZCwhVOSdnn+RzkzlRz/eY4BgJBYVpg==
|
||||
dependencies:
|
||||
argparse "^1.0.7"
|
||||
esprima "^4.0.0"
|
||||
|
||||
@@ -678,6 +678,51 @@ def _update_encryption_path(
|
||||
)
|
||||
|
||||
|
||||
def _update_checkpointer_path(
|
||||
config_path: pathlib.Path, config: Config, local_deps: LocalDeps
|
||||
) -> None:
|
||||
"""Update checkpointer.path to use Docker container paths."""
|
||||
checkpointer_conf = config.get("checkpointer")
|
||||
if not checkpointer_conf or not isinstance(checkpointer_conf, dict):
|
||||
return
|
||||
if not (path_str := checkpointer_conf.get("path")):
|
||||
return
|
||||
|
||||
module_str, sep, attr_str = path_str.partition(":")
|
||||
if not sep or not module_str.startswith("."):
|
||||
return # Already validated or absolute path
|
||||
|
||||
resolved = config_path.parent / module_str
|
||||
if not resolved.exists():
|
||||
raise FileNotFoundError(
|
||||
f"Checkpointer file not found: {resolved} (from {path_str})"
|
||||
)
|
||||
if not resolved.is_file():
|
||||
raise IsADirectoryError(f"Checkpointer path must be a file: {resolved}")
|
||||
|
||||
# Check faux packages first (higher priority)
|
||||
for faux_path, (_, destpath) in local_deps.faux_pkgs.items():
|
||||
if resolved.is_relative_to(faux_path):
|
||||
new_path = f"{destpath}/{resolved.relative_to(faux_path)}:{attr_str}"
|
||||
checkpointer_conf["path"] = new_path
|
||||
return
|
||||
|
||||
# Check real packages
|
||||
for real_path in local_deps.real_pkgs:
|
||||
if resolved.is_relative_to(real_path):
|
||||
new_path = (
|
||||
f"/deps/{real_path.name}/{resolved.relative_to(real_path)}:{attr_str}"
|
||||
)
|
||||
checkpointer_conf["path"] = new_path
|
||||
return
|
||||
|
||||
raise ValueError(
|
||||
f"Checkpointer file '{resolved}' not covered by dependencies.\n"
|
||||
"Add its parent directory to the 'dependencies' array in your config.\n"
|
||||
f"Current dependencies: {config['dependencies']}"
|
||||
)
|
||||
|
||||
|
||||
def _update_http_app_path(
|
||||
config_path: pathlib.Path, config: Config, local_deps: LocalDeps
|
||||
) -> None:
|
||||
@@ -877,6 +922,8 @@ def python_config_to_docker(
|
||||
_update_auth_path(config_path, config, local_deps)
|
||||
# Rewrite encryption path, so it points to the correct location in the Docker container
|
||||
_update_encryption_path(config_path, config, local_deps)
|
||||
# Rewrite checkpointer path, so it points to the correct location in the Docker container
|
||||
_update_checkpointer_path(config_path, config, local_deps)
|
||||
# Rewrite HTTP app path, so it points to the correct location in the Docker container
|
||||
_update_http_app_path(config_path, config, local_deps)
|
||||
|
||||
|
||||
@@ -167,6 +167,26 @@ class CheckpointerConfig(TypedDict, total=False):
|
||||
If omitted, no checkpointer is set up (the object store will still be present, however).
|
||||
"""
|
||||
|
||||
path: str
|
||||
"""Import path to an async context manager that yields a `BaseCheckpointSaver`
|
||||
instance.
|
||||
|
||||
The referenced object should be an `@asynccontextmanager`-decorated function
|
||||
so that the server can properly manage the checkpointer's lifecycle (e.g.
|
||||
opening and closing connections).
|
||||
|
||||
Examples:
|
||||
- "./my_checkpointer.py:create_checkpointer"
|
||||
- "my_package.checkpointer:create_checkpointer"
|
||||
|
||||
When provided, this replaces the default checkpointer.
|
||||
|
||||
You can use the `langgraph-checkpoint-conformance` package
|
||||
(https://pypi.org/project/langgraph-checkpoint-conformance/) to run simple
|
||||
conformance tests against your custom checkpointer and catch
|
||||
incompatibilities early.
|
||||
"""
|
||||
|
||||
ttl: ThreadTTLConfig | None
|
||||
"""Optional. Defines the TTL (time-to-live) behavior configuration.
|
||||
|
||||
|
||||
@@ -542,6 +542,10 @@
|
||||
"description": "Configuration for the built-in checkpointer, which handles checkpointing of state.\n\nIf omitted, no checkpointer is set up (the object store will still be present, however).",
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"path": {
|
||||
"type": "string",
|
||||
"description": "Import path to an async context manager that yields a `BaseCheckpointSaver`\ninstance.\n\nThe referenced object should be an `@asynccontextmanager`-decorated function\nso that the server can properly manage the checkpointer's lifecycle (e.g.\nopening and closing connections).\n"
|
||||
},
|
||||
"serde": {
|
||||
"anyOf": [
|
||||
{
|
||||
|
||||
@@ -542,6 +542,10 @@
|
||||
"description": "Configuration for the built-in checkpointer, which handles checkpointing of state.\n\nIf omitted, no checkpointer is set up (the object store will still be present, however).",
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"path": {
|
||||
"type": "string",
|
||||
"description": "Import path to an async context manager that yields a `BaseCheckpointSaver`\ninstance.\n\nThe referenced object should be an `@asynccontextmanager`-decorated function\nso that the server can properly manage the checkpointer's lifecycle (e.g.\nopening and closing connections).\n"
|
||||
},
|
||||
"serde": {
|
||||
"anyOf": [
|
||||
{
|
||||
|
||||
Generated
+3
-3
@@ -1071,15 +1071,15 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-sdk"
|
||||
version = "0.3.5"
|
||||
version = "0.3.6"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "httpx", marker = "python_full_version >= '3.11'" },
|
||||
{ name = "orjson", marker = "python_full_version >= '3.11'" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/60/2b/2dae368ac76e315197f07ab58077aadf20833c226fbfd450d71745850314/langgraph_sdk-0.3.5.tar.gz", hash = "sha256:64669e9885a908578eed921ef9a8e52b8d0cd38db1e3e5d6d299d4e6f8830ac0", size = 177470, upload-time = "2026-02-10T16:56:09.18Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/3e/ec/477fa8b408f948b145d90fd935c0a9f37945fa5ec1dfabfc71e7cafba6d8/langgraph_sdk-0.3.6.tar.gz", hash = "sha256:7650f607f89c1586db5bee391b1a8754cbe1fc83b721ff2f1450f8906e790bd7", size = 182666, upload-time = "2026-02-14T19:46:03.752Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/84/d5/a14d957c515ba7a9713bf0f03f2b9277979c403bc50f829bdfd54ae7dc9e/langgraph_sdk-0.3.5-py3-none-any.whl", hash = "sha256:bcfa1dcbddadb604076ce46f5e08969538735e5ac47fa863d4fac5a512dab5c9", size = 70851, upload-time = "2026-02-10T16:56:07.983Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/d8/61/12508e12652edd1874327271a5a8834c728a605f53a1a1c945f13ab69664/langgraph_sdk-0.3.6-py3-none-any.whl", hash = "sha256:7df2fd552ad7262d0baf8e1f849dce1d62186e76dcdd36db9dc5bdfa5c3fc20f", size = 88277, upload-time = "2026-02-14T19:46:02.48Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
|
||||
@@ -728,7 +728,6 @@ async def _acall_impl(
|
||||
)
|
||||
else:
|
||||
fut.set_result(None)
|
||||
futures()[fut] = next_task # type: ignore[index]
|
||||
else:
|
||||
# schedule the next task
|
||||
fut = cast(
|
||||
|
||||
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "langgraph"
|
||||
version = "1.0.8"
|
||||
version = "1.0.9"
|
||||
description = "Building stateful, multi-actor applications with LLMs"
|
||||
authors = []
|
||||
requires-python = ">=3.10"
|
||||
@@ -27,7 +27,7 @@ dependencies = [
|
||||
"langchain-core>=0.1",
|
||||
"langgraph-checkpoint>=2.1.0,<5.0.0",
|
||||
"langgraph-sdk>=0.3.0,<0.4.0",
|
||||
"langgraph-prebuilt>=1.0.7,<1.1.0",
|
||||
"langgraph-prebuilt>=1.0.8,<1.1.0",
|
||||
"xxhash>=3.5.0",
|
||||
"pydantic>=2.7.4",
|
||||
]
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -5800,6 +5800,284 @@ def test_multiple_interrupts_functional_cache(
|
||||
assert counter == 6
|
||||
|
||||
|
||||
def test_task_before_interrupt_resume(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly when a @task runs
|
||||
before interrupt-producing tasks in an @entrypoint.
|
||||
|
||||
The @task wrapper on both setup and ask is essential to reproduce the bug:
|
||||
- @task on setup triggers a mid-step put_writes (creating a new pending_writes list)
|
||||
- @task on ask means interrupt() runs in a child scratchpad that must
|
||||
delegate to the parent for null resume consumption tracking
|
||||
"""
|
||||
|
||||
@entrypoint(checkpointer=sync_checkpointer)
|
||||
def workflow(number_of_topics: int) -> dict:
|
||||
@task
|
||||
def setup() -> int:
|
||||
return number_of_topics
|
||||
|
||||
@task
|
||||
def ask(question: str) -> str:
|
||||
return interrupt(question)
|
||||
|
||||
n = setup().result()
|
||||
|
||||
answers = []
|
||||
for i in range(n):
|
||||
q = f"Whats the answer for topic {i + 1}?"
|
||||
answers.append(ask(q).result())
|
||||
|
||||
return {"answers": answers}
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - should get first interrupt
|
||||
result = workflow.invoke(2, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 1?"
|
||||
|
||||
# Resume with answer for topic 1 - should get second interrupt
|
||||
result = workflow.invoke(Command(resume="answer1"), config=config)
|
||||
assert "__interrupt__" in result, f"Expected interrupt for topic 2, got: {result}"
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 2?"
|
||||
|
||||
# Resume with answer for topic 2 - should get final result
|
||||
result = workflow.invoke(Command(resume="answer2"), config=config)
|
||||
assert result == {"answers": ["answer1", "answer2"]}
|
||||
|
||||
|
||||
def test_multiple_tasks_before_interrupt_resume(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly when multiple @tasks
|
||||
run before an interrupt-producing task in an @entrypoint."""
|
||||
|
||||
@entrypoint(checkpointer=sync_checkpointer)
|
||||
def workflow(inputs: dict) -> dict:
|
||||
@task
|
||||
def step_a(x: int) -> int:
|
||||
return x + 1
|
||||
|
||||
@task
|
||||
def step_b(x: int) -> int:
|
||||
return x * 2
|
||||
|
||||
@task
|
||||
def ask(question: str) -> str:
|
||||
return interrupt(question)
|
||||
|
||||
a = step_a(inputs["x"]).result()
|
||||
b = step_b(a).result()
|
||||
|
||||
answer = ask(f"Result so far is {b}. What next?").result()
|
||||
|
||||
return {"computed": b, "answer": answer}
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - should get interrupt
|
||||
result = workflow.invoke({"x": 5}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "Result so far is 12. What next?"
|
||||
|
||||
# Resume
|
||||
result = workflow.invoke(Command(resume="continue"), config=config)
|
||||
assert result == {"computed": 12, "answer": "continue"}
|
||||
|
||||
|
||||
def test_no_redundant_put_writes_for_cached_task(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Cached @tasks on resume must not trigger redundant put_writes."""
|
||||
from unittest.mock import patch
|
||||
|
||||
from langgraph.pregel._loop import PregelLoop
|
||||
|
||||
@task
|
||||
def setup(x: int) -> int:
|
||||
return x
|
||||
|
||||
@task
|
||||
def ask(question: str) -> str:
|
||||
return interrupt(question)
|
||||
|
||||
@entrypoint(checkpointer=sync_checkpointer)
|
||||
def workflow(x: int) -> dict:
|
||||
n = setup(x).result()
|
||||
answer = ask(f"q{n}").result()
|
||||
return {"answer": answer}
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
result = workflow.invoke(1, config=config)
|
||||
assert "__interrupt__" in result
|
||||
|
||||
put_writes_task_ids: list[str] = []
|
||||
orig = PregelLoop.put_writes
|
||||
|
||||
def spy(self, task_id, writes):
|
||||
put_writes_task_ids.append(task_id)
|
||||
return orig(self, task_id, writes)
|
||||
|
||||
with patch.object(PregelLoop, "put_writes", spy):
|
||||
result = workflow.invoke(Command(resume="ans"), config=config)
|
||||
|
||||
assert result == {"answer": "ans"}
|
||||
# Count unique non-null task IDs that got put_writes.
|
||||
# Should be exactly 2: the ask task and the entrypoint task.
|
||||
# If 3, the cached setup task is being redundantly re-committed.
|
||||
non_null = set(tid for tid in put_writes_task_ids if not tid.startswith("00000000"))
|
||||
assert len(non_null) == 2, (
|
||||
f"Expected 2 task IDs in put_writes (ask + entrypoint), got {len(non_null)}"
|
||||
)
|
||||
|
||||
|
||||
def test_node_before_interrupt_resume_graph_api(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly in a StateGraph when a
|
||||
node runs before a node that calls interrupt(). This is the graph-API
|
||||
analog of test_task_before_interrupt_resume (entrypoint API)."""
|
||||
|
||||
class State(TypedDict):
|
||||
topics: list[str]
|
||||
answers: Annotated[list[str], operator.add]
|
||||
|
||||
def setup(state: State) -> dict:
|
||||
return {"topics": [f"topic {i + 1}" for i in range(len(state["topics"]))]}
|
||||
|
||||
def ask(state: State) -> dict:
|
||||
answers = []
|
||||
for topic in state["topics"]:
|
||||
answer = interrupt(f"Whats the answer for {topic}?")
|
||||
answers.append(answer)
|
||||
return {"answers": answers}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("setup", setup)
|
||||
.add_node("ask", ask)
|
||||
.add_edge(START, "setup")
|
||||
.add_edge("setup", "ask")
|
||||
.add_edge("ask", END)
|
||||
.compile(checkpointer=sync_checkpointer)
|
||||
)
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - setup runs, then ask interrupts on the first topic
|
||||
result = graph.invoke({"topics": ["a", "b"], "answers": []}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 1?"
|
||||
|
||||
# Resume with answer for topic 1 - should get second interrupt
|
||||
result = graph.invoke(Command(resume="answer1"), config=config)
|
||||
assert "__interrupt__" in result, f"Expected interrupt for topic 2, got: {result}"
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 2?"
|
||||
|
||||
# Resume with answer for topic 2 - should complete
|
||||
result = graph.invoke(Command(resume="answer2"), config=config)
|
||||
assert result == {
|
||||
"topics": ["topic 1", "topic 2"],
|
||||
"answers": ["answer1", "answer2"],
|
||||
}
|
||||
|
||||
|
||||
def test_multiple_nodes_before_interrupt_resume_graph_api(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly in a StateGraph when
|
||||
multiple nodes run before a node that calls interrupt(). This is the
|
||||
graph-API analog of test_multiple_tasks_before_interrupt_resume."""
|
||||
|
||||
class State(TypedDict):
|
||||
value: int
|
||||
answer: str
|
||||
|
||||
def step_a(state: State) -> dict:
|
||||
return {"value": state["value"] + 1}
|
||||
|
||||
def step_b(state: State) -> dict:
|
||||
return {"value": state["value"] * 2}
|
||||
|
||||
def ask(state: State) -> dict:
|
||||
answer = interrupt(f"Result so far is {state['value']}. What next?")
|
||||
return {"answer": answer}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("step_a", step_a)
|
||||
.add_node("step_b", step_b)
|
||||
.add_node("ask", ask)
|
||||
.add_edge(START, "step_a")
|
||||
.add_edge("step_a", "step_b")
|
||||
.add_edge("step_b", "ask")
|
||||
.add_edge("ask", END)
|
||||
.compile(checkpointer=sync_checkpointer)
|
||||
)
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - step_a and step_b run, then ask interrupts
|
||||
result = graph.invoke({"value": 5, "answer": ""}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "Result so far is 12. What next?"
|
||||
|
||||
# Resume - should complete
|
||||
result = graph.invoke(Command(resume="continue"), config=config)
|
||||
assert result == {"value": 12, "answer": "continue"}
|
||||
|
||||
|
||||
def test_node_before_multiple_interrupt_cycles_graph_api(
|
||||
sync_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that a node running before an interrupt node does not interfere
|
||||
with multiple interrupt/resume cycles in a StateGraph."""
|
||||
|
||||
class State(TypedDict):
|
||||
count: int
|
||||
data: str
|
||||
|
||||
def prepare(state: State) -> dict:
|
||||
return {"count": state["count"] + 10}
|
||||
|
||||
def multi_interrupt(state: State) -> dict:
|
||||
first = interrupt("First question?")
|
||||
second = interrupt("Second question?")
|
||||
return {"data": f"{first},{second}"}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("prepare", prepare)
|
||||
.add_node("multi_interrupt", multi_interrupt)
|
||||
.add_edge(START, "prepare")
|
||||
.add_edge("prepare", "multi_interrupt")
|
||||
.add_edge("multi_interrupt", END)
|
||||
.compile(checkpointer=sync_checkpointer)
|
||||
)
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - prepare runs, multi_interrupt hits first interrupt
|
||||
result = graph.invoke({"count": 0, "data": ""}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "First question?"
|
||||
|
||||
# Resume first interrupt - hits second interrupt
|
||||
result = graph.invoke(Command(resume="first_answer"), config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "Second question?"
|
||||
|
||||
# Resume second interrupt - completes
|
||||
result = graph.invoke(Command(resume="second_answer"), config=config)
|
||||
assert result == {"count": 10, "data": "first_answer,second_answer"}
|
||||
|
||||
|
||||
def test_double_interrupt_subgraph(sync_checkpointer: BaseCheckpointSaver) -> None:
|
||||
class AgentState(TypedDict):
|
||||
input: str
|
||||
|
||||
@@ -7920,6 +7920,290 @@ async def test_interrupts_in_tasks_surfaced_once(
|
||||
assert result[1] == "Added Will!"
|
||||
|
||||
|
||||
@NEEDS_CONTEXTVARS
|
||||
async def test_task_before_interrupt_resume(
|
||||
async_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly when a @task runs
|
||||
before interrupt-producing tasks in an @entrypoint.
|
||||
|
||||
The @task wrapper on both setup and ask is essential to reproduce the bug:
|
||||
- @task on setup triggers a mid-step put_writes (creating a new pending_writes list)
|
||||
- @task on ask means interrupt() runs in a child scratchpad that must
|
||||
delegate to the parent for null resume consumption tracking
|
||||
"""
|
||||
|
||||
@entrypoint(checkpointer=async_checkpointer)
|
||||
async def workflow(number_of_topics: int) -> dict:
|
||||
@task
|
||||
async def setup() -> int:
|
||||
return number_of_topics
|
||||
|
||||
@task
|
||||
async def ask(question: str) -> str:
|
||||
return interrupt(question)
|
||||
|
||||
n = await setup()
|
||||
|
||||
answers = []
|
||||
for i in range(n):
|
||||
q = f"Whats the answer for topic {i + 1}?"
|
||||
answers.append(await ask(q))
|
||||
|
||||
return {"answers": answers}
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - should get first interrupt
|
||||
result = await workflow.ainvoke(2, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 1?"
|
||||
|
||||
# Resume with answer for topic 1 - should get second interrupt
|
||||
result = await workflow.ainvoke(Command(resume="answer1"), config=config)
|
||||
assert "__interrupt__" in result, f"Expected interrupt for topic 2, got: {result}"
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 2?"
|
||||
|
||||
# Resume with answer for topic 2 - should get final result
|
||||
result = await workflow.ainvoke(Command(resume="answer2"), config=config)
|
||||
assert result == {"answers": ["answer1", "answer2"]}
|
||||
|
||||
|
||||
@NEEDS_CONTEXTVARS
|
||||
async def test_multiple_tasks_before_interrupt_resume(
|
||||
async_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly when multiple @tasks
|
||||
run before an interrupt-producing task in an @entrypoint."""
|
||||
|
||||
@entrypoint(checkpointer=async_checkpointer)
|
||||
async def workflow(inputs: dict) -> dict:
|
||||
@task
|
||||
async def step_a(x: int) -> int:
|
||||
return x + 1
|
||||
|
||||
@task
|
||||
async def step_b(x: int) -> int:
|
||||
return x * 2
|
||||
|
||||
@task
|
||||
async def ask(question: str) -> str:
|
||||
return interrupt(question)
|
||||
|
||||
a = await step_a(inputs["x"])
|
||||
b = await step_b(a)
|
||||
|
||||
answer = await ask(f"Result so far is {b}. What next?")
|
||||
|
||||
return {"computed": b, "answer": answer}
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - should get interrupt
|
||||
result = await workflow.ainvoke({"x": 5}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "Result so far is 12. What next?"
|
||||
|
||||
# Resume
|
||||
result = await workflow.ainvoke(Command(resume="continue"), config=config)
|
||||
assert result == {"computed": 12, "answer": "continue"}
|
||||
|
||||
|
||||
@NEEDS_CONTEXTVARS
|
||||
async def test_no_redundant_put_writes_for_cached_task(
|
||||
async_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Cached @tasks on resume must not trigger redundant put_writes."""
|
||||
from unittest.mock import patch
|
||||
|
||||
from langgraph.pregel._loop import PregelLoop
|
||||
|
||||
@task
|
||||
async def setup(x: int) -> int:
|
||||
return x
|
||||
|
||||
@task
|
||||
async def ask(question: str) -> str:
|
||||
return interrupt(question)
|
||||
|
||||
@entrypoint(checkpointer=async_checkpointer)
|
||||
async def workflow(x: int) -> dict:
|
||||
n = await setup(x)
|
||||
answer = await ask(f"q{n}")
|
||||
return {"answer": answer}
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
result = await workflow.ainvoke(1, config=config)
|
||||
assert "__interrupt__" in result
|
||||
|
||||
put_writes_task_ids: list[str] = []
|
||||
orig = PregelLoop.put_writes
|
||||
|
||||
def spy(self, task_id, writes):
|
||||
put_writes_task_ids.append(task_id)
|
||||
return orig(self, task_id, writes)
|
||||
|
||||
with patch.object(PregelLoop, "put_writes", spy):
|
||||
result = await workflow.ainvoke(Command(resume="ans"), config=config)
|
||||
|
||||
assert result == {"answer": "ans"}
|
||||
# Count unique non-null task IDs that got put_writes.
|
||||
# Should be exactly 2: the ask task and the entrypoint task.
|
||||
# If 3, the cached setup task is being redundantly re-committed.
|
||||
non_null = set(tid for tid in put_writes_task_ids if not tid.startswith("00000000"))
|
||||
assert len(non_null) == 2, (
|
||||
f"Expected 2 task IDs in put_writes (ask + entrypoint), got {len(non_null)}"
|
||||
)
|
||||
|
||||
|
||||
@NEEDS_CONTEXTVARS
|
||||
async def test_node_before_interrupt_resume_graph_api(
|
||||
async_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly in a StateGraph when a
|
||||
node runs before a node that calls interrupt(). This is the graph-API
|
||||
analog of test_task_before_interrupt_resume (entrypoint API)."""
|
||||
|
||||
class State(TypedDict):
|
||||
topics: list[str]
|
||||
answers: Annotated[list[str], operator.add]
|
||||
|
||||
def setup(state: State) -> dict:
|
||||
return {"topics": [f"topic {i + 1}" for i in range(len(state["topics"]))]}
|
||||
|
||||
def ask(state: State) -> dict:
|
||||
answers = []
|
||||
for topic in state["topics"]:
|
||||
answer = interrupt(f"Whats the answer for {topic}?")
|
||||
answers.append(answer)
|
||||
return {"answers": answers}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("setup", setup)
|
||||
.add_node("ask", ask)
|
||||
.add_edge(START, "setup")
|
||||
.add_edge("setup", "ask")
|
||||
.add_edge("ask", END)
|
||||
.compile(checkpointer=async_checkpointer)
|
||||
)
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - setup runs, then ask interrupts on the first topic
|
||||
result = await graph.ainvoke({"topics": ["a", "b"], "answers": []}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 1?"
|
||||
|
||||
# Resume with answer for topic 1 - should get second interrupt
|
||||
result = await graph.ainvoke(Command(resume="answer1"), config=config)
|
||||
assert "__interrupt__" in result, f"Expected interrupt for topic 2, got: {result}"
|
||||
assert len(result["__interrupt__"]) == 1
|
||||
assert result["__interrupt__"][0].value == "Whats the answer for topic 2?"
|
||||
|
||||
# Resume with answer for topic 2 - should complete
|
||||
result = await graph.ainvoke(Command(resume="answer2"), config=config)
|
||||
assert result == {
|
||||
"topics": ["topic 1", "topic 2"],
|
||||
"answers": ["answer1", "answer2"],
|
||||
}
|
||||
|
||||
|
||||
@NEEDS_CONTEXTVARS
|
||||
async def test_multiple_nodes_before_interrupt_resume_graph_api(
|
||||
async_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that Command(resume=value) works correctly in a StateGraph when
|
||||
multiple nodes run before a node that calls interrupt(). This is the
|
||||
graph-API analog of test_multiple_tasks_before_interrupt_resume."""
|
||||
|
||||
class State(TypedDict):
|
||||
value: int
|
||||
answer: str
|
||||
|
||||
def step_a(state: State) -> dict:
|
||||
return {"value": state["value"] + 1}
|
||||
|
||||
def step_b(state: State) -> dict:
|
||||
return {"value": state["value"] * 2}
|
||||
|
||||
def ask(state: State) -> dict:
|
||||
answer = interrupt(f"Result so far is {state['value']}. What next?")
|
||||
return {"answer": answer}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("step_a", step_a)
|
||||
.add_node("step_b", step_b)
|
||||
.add_node("ask", ask)
|
||||
.add_edge(START, "step_a")
|
||||
.add_edge("step_a", "step_b")
|
||||
.add_edge("step_b", "ask")
|
||||
.add_edge("ask", END)
|
||||
.compile(checkpointer=async_checkpointer)
|
||||
)
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - step_a and step_b run, then ask interrupts
|
||||
result = await graph.ainvoke({"value": 5, "answer": ""}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "Result so far is 12. What next?"
|
||||
|
||||
# Resume - should complete
|
||||
result = await graph.ainvoke(Command(resume="continue"), config=config)
|
||||
assert result == {"value": 12, "answer": "continue"}
|
||||
|
||||
|
||||
@NEEDS_CONTEXTVARS
|
||||
async def test_node_before_multiple_interrupt_cycles_graph_api(
|
||||
async_checkpointer: BaseCheckpointSaver,
|
||||
) -> None:
|
||||
"""Test that a node running before an interrupt node does not interfere
|
||||
with multiple interrupt/resume cycles in a StateGraph."""
|
||||
|
||||
class State(TypedDict):
|
||||
count: int
|
||||
data: str
|
||||
|
||||
def prepare(state: State) -> dict:
|
||||
return {"count": state["count"] + 10}
|
||||
|
||||
def multi_interrupt(state: State) -> dict:
|
||||
first = interrupt("First question?")
|
||||
second = interrupt("Second question?")
|
||||
return {"data": f"{first},{second}"}
|
||||
|
||||
graph = (
|
||||
StateGraph(State)
|
||||
.add_node("prepare", prepare)
|
||||
.add_node("multi_interrupt", multi_interrupt)
|
||||
.add_edge(START, "prepare")
|
||||
.add_edge("prepare", "multi_interrupt")
|
||||
.add_edge("multi_interrupt", END)
|
||||
.compile(checkpointer=async_checkpointer)
|
||||
)
|
||||
|
||||
config = {"configurable": {"thread_id": "1"}}
|
||||
|
||||
# First invocation - prepare runs, multi_interrupt hits first interrupt
|
||||
result = await graph.ainvoke({"count": 0, "data": ""}, config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "First question?"
|
||||
|
||||
# Resume first interrupt - hits second interrupt
|
||||
result = await graph.ainvoke(Command(resume="first_answer"), config=config)
|
||||
assert "__interrupt__" in result
|
||||
assert result["__interrupt__"][0].value == "Second question?"
|
||||
|
||||
# Resume second interrupt - completes
|
||||
result = await graph.ainvoke(Command(resume="second_answer"), config=config)
|
||||
assert result == {"count": 10, "data": "first_answer,second_answer"}
|
||||
|
||||
|
||||
async def test_pregel_loop_refcount():
|
||||
gc.collect()
|
||||
try:
|
||||
|
||||
Generated
+6
-4
@@ -1367,7 +1367,7 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph"
|
||||
version = "1.0.8"
|
||||
version = "1.0.9"
|
||||
source = { editable = "." }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
@@ -1607,7 +1607,7 @@ dependencies = [
|
||||
[package.metadata]
|
||||
requires-dist = [
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "orjson", specifier = ">=3.10.1" },
|
||||
{ name = "orjson", specifier = ">=3.11.5" },
|
||||
{ name = "psycopg", specifier = ">=3.2.0" },
|
||||
{ name = "psycopg-pool", specifier = ">=3.2.0" },
|
||||
]
|
||||
@@ -1617,6 +1617,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -1633,6 +1634,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
@@ -1735,7 +1737,7 @@ test = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-prebuilt"
|
||||
version = "1.0.7"
|
||||
version = "1.0.8"
|
||||
source = { editable = "../prebuilt" }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
@@ -1812,7 +1814,7 @@ dependencies = [
|
||||
[package.metadata]
|
||||
requires-dist = [
|
||||
{ name = "httpx", specifier = ">=0.25.2" },
|
||||
{ name = "orjson", specifier = ">=3.10.1" },
|
||||
{ name = "orjson", specifier = ">=3.11.5" },
|
||||
]
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
|
||||
@@ -922,7 +922,7 @@ class ToolNode(RunnableCallable):
|
||||
raise TypeError(msg)
|
||||
|
||||
# Inject state, store, and runtime right before invocation
|
||||
injected_call = self._inject_tool_args(call, request.runtime)
|
||||
injected_call = self._inject_tool_args(call, request.runtime, tool)
|
||||
call_args = {**injected_call, "type": "tool_call"}
|
||||
|
||||
try:
|
||||
@@ -1075,7 +1075,7 @@ class ToolNode(RunnableCallable):
|
||||
raise TypeError(msg)
|
||||
|
||||
# Inject state, store, and runtime right before invocation
|
||||
injected_call = self._inject_tool_args(call, request.runtime)
|
||||
injected_call = self._inject_tool_args(call, request.runtime, tool)
|
||||
call_args = {**injected_call, "type": "tool_call"}
|
||||
|
||||
try:
|
||||
@@ -1281,6 +1281,7 @@ class ToolNode(RunnableCallable):
|
||||
self,
|
||||
tool_call: ToolCall,
|
||||
tool_runtime: ToolRuntime,
|
||||
tool: BaseTool | None = None,
|
||||
) -> ToolCall:
|
||||
"""Inject graph state, store, and runtime into tool call arguments.
|
||||
|
||||
@@ -1299,6 +1300,9 @@ class ToolNode(RunnableCallable):
|
||||
Must contain 'name', 'args', 'id', and 'type' fields.
|
||||
tool_runtime: The ToolRuntime instance containing all runtime context
|
||||
(state, config, store, context, stream_writer) to inject into tools.
|
||||
tool: Optional tool instance. When provided, allows injection for
|
||||
dynamically registered tools that are not in self.tools_by_name
|
||||
(e.g., tools added via middleware's wrap_tool_call).
|
||||
|
||||
Returns:
|
||||
A new ToolCall dictionary with the same structure as the input but with
|
||||
@@ -1312,10 +1316,12 @@ class ToolNode(RunnableCallable):
|
||||
This method is called automatically during tool execution. It should not
|
||||
be called from outside the `ToolNode`.
|
||||
"""
|
||||
if tool_call["name"] not in self.tools_by_name:
|
||||
return tool_call
|
||||
|
||||
injected = self._injected_args.get(tool_call["name"])
|
||||
if not injected and tool is not None:
|
||||
# For dynamically registered tools (e.g., added via middleware's
|
||||
# wrap_tool_call), compute injected args on-the-fly since they
|
||||
# were not present during ToolNode initialization.
|
||||
injected = _get_all_injected_args(tool)
|
||||
if not injected:
|
||||
return tool_call
|
||||
|
||||
|
||||
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "langgraph-prebuilt"
|
||||
version = "1.0.7"
|
||||
version = "1.0.8"
|
||||
description = "Library with high-level APIs for creating and executing LangGraph agents and tools."
|
||||
authors = []
|
||||
requires-python = ">=3.10"
|
||||
|
||||
@@ -1902,3 +1902,109 @@ async def test_tool_node_tool_runtime_generic() -> None:
|
||||
assert tool_message.type == "tool"
|
||||
assert tool_message.content == "test_info"
|
||||
assert tool_message.tool_call_id == "call_1"
|
||||
|
||||
|
||||
def test_tool_node_inject_runtime_dynamic_tool_via_wrap_tool_call() -> None:
|
||||
"""Test that ToolRuntime is injected for dynamically registered tools.
|
||||
|
||||
Regression test for https://github.com/langchain-ai/langchain/issues/35305.
|
||||
When a tool is dynamically provided via wrap_tool_call (not registered at
|
||||
ToolNode init time), ToolRuntime should still be injected into the tool.
|
||||
"""
|
||||
|
||||
@dec_tool
|
||||
def static_tool(x: int) -> str:
|
||||
"""A static tool registered at init."""
|
||||
return f"static: {x}"
|
||||
|
||||
@dec_tool
|
||||
def dynamic_tool_with_runtime(x: int, runtime: ToolRuntime) -> str:
|
||||
"""A dynamic tool that needs ToolRuntime injection."""
|
||||
return f"dynamic: x={x}, tool_call_id={runtime.tool_call_id}"
|
||||
|
||||
def wrap_tool_call(request, execute):
|
||||
"""Middleware that swaps in a dynamic tool."""
|
||||
if request.tool_call["name"] == "dynamic_tool_with_runtime":
|
||||
# Override tool to the dynamic one (not registered at init)
|
||||
new_request = request.override(tool=dynamic_tool_with_runtime)
|
||||
return execute(new_request)
|
||||
return execute(request)
|
||||
|
||||
# ToolNode only knows about static_tool at init time
|
||||
tool_node = ToolNode(
|
||||
[static_tool],
|
||||
wrap_tool_call=wrap_tool_call,
|
||||
)
|
||||
|
||||
# Verify the dynamic tool is NOT in the tool node's registered tools
|
||||
assert "dynamic_tool_with_runtime" not in tool_node.tools_by_name
|
||||
|
||||
# Call the dynamic tool
|
||||
tool_call = {
|
||||
"name": "dynamic_tool_with_runtime",
|
||||
"args": {"x": 42},
|
||||
"id": "call_dynamic_1",
|
||||
"type": "tool_call",
|
||||
}
|
||||
msg = AIMessage("", tool_calls=[tool_call])
|
||||
result = tool_node.invoke(
|
||||
{"messages": [msg]},
|
||||
config=_create_config_with_runtime(),
|
||||
)
|
||||
|
||||
# ToolRuntime should be injected and the tool should execute successfully
|
||||
tool_message = result["messages"][-1]
|
||||
assert tool_message.content == "dynamic: x=42, tool_call_id=call_dynamic_1"
|
||||
assert tool_message.tool_call_id == "call_dynamic_1"
|
||||
|
||||
|
||||
async def test_tool_node_inject_runtime_dynamic_tool_via_wrap_tool_call_async() -> None:
|
||||
"""Test that ToolRuntime is injected for dynamically registered tools (async).
|
||||
|
||||
Async version of the regression test for
|
||||
https://github.com/langchain-ai/langchain/issues/35305.
|
||||
"""
|
||||
|
||||
@dec_tool
|
||||
def static_tool(x: int) -> str:
|
||||
"""A static tool registered at init."""
|
||||
return f"static: {x}"
|
||||
|
||||
@dec_tool
|
||||
async def dynamic_tool_with_runtime(x: int, runtime: ToolRuntime) -> str:
|
||||
"""A dynamic async tool that needs ToolRuntime injection."""
|
||||
return f"dynamic: x={x}, tool_call_id={runtime.tool_call_id}"
|
||||
|
||||
async def awrap_tool_call(request, execute):
|
||||
"""Async middleware that swaps in a dynamic tool."""
|
||||
if request.tool_call["name"] == "dynamic_tool_with_runtime":
|
||||
new_request = request.override(tool=dynamic_tool_with_runtime)
|
||||
return await execute(new_request)
|
||||
return await execute(request)
|
||||
|
||||
# ToolNode only knows about static_tool at init time
|
||||
tool_node = ToolNode(
|
||||
[static_tool],
|
||||
awrap_tool_call=awrap_tool_call,
|
||||
)
|
||||
|
||||
# Verify the dynamic tool is NOT in the tool node's registered tools
|
||||
assert "dynamic_tool_with_runtime" not in tool_node.tools_by_name
|
||||
|
||||
# Call the dynamic tool
|
||||
tool_call = {
|
||||
"name": "dynamic_tool_with_runtime",
|
||||
"args": {"x": 42},
|
||||
"id": "call_dynamic_2",
|
||||
"type": "tool_call",
|
||||
}
|
||||
msg = AIMessage("", tool_calls=[tool_call])
|
||||
result = await tool_node.ainvoke(
|
||||
{"messages": [msg]},
|
||||
config=_create_config_with_runtime(),
|
||||
)
|
||||
|
||||
# ToolRuntime should be injected and the tool should execute successfully
|
||||
tool_message = result["messages"][-1]
|
||||
assert tool_message.content == "dynamic: x=42, tool_call_id=call_dynamic_2"
|
||||
assert tool_message.tool_call_id == "call_dynamic_2"
|
||||
|
||||
Generated
+9
-7
@@ -249,7 +249,7 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "langchain-core"
|
||||
version = "1.2.12"
|
||||
version = "1.2.13"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "jsonpatch" },
|
||||
@@ -261,14 +261,14 @@ dependencies = [
|
||||
{ name = "typing-extensions" },
|
||||
{ name = "uuid-utils" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/2a/1d/08e935d1532fcc90981f6e5bb6825914c9227ea7a962c62b1e18619b49e7/langchain_core-1.2.12.tar.gz", hash = "sha256:4d7fa6643d7ab06fc1905a9b7dcbe96a6f3c181046b56edf9c0c17ecd412d9e9", size = 831329, upload-time = "2026-02-12T20:53:15.01Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/fb/bb/c501ca60556c11ac80d1454bdcac63cb33583ce4e64fc4535ad5a7d5c6ba/langchain_core-1.2.13.tar.gz", hash = "sha256:d2773d0d0130a356378db9a858cfeef64c3d64bc03722f1d4d6c40eb46fdf01b", size = 831612, upload-time = "2026-02-15T07:45:57.014Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/8c/a5/678ab0e5cc57794f20ae5ed12c1442506ef1108c9434f950aebc6044e5a3/langchain_core-1.2.12-py3-none-any.whl", hash = "sha256:66ca17a2a9cb007ab29021968e6adfcf4228067151dc2bd6ebfff265ffaf92f5", size = 500132, upload-time = "2026-02-12T20:53:13.806Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/12/ab/60fd69e5d55f67d422baefddaaca523c42cd7510ab6aeb17db6ae57fb107/langchain_core-1.2.13-py3-none-any.whl", hash = "sha256:b31823e28d3eff1e237096d0bd3bf80c6f9624eb471a9496dbfbd427779f8d82", size = 500485, upload-time = "2026-02-15T07:45:55.422Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "langgraph"
|
||||
version = "1.0.8"
|
||||
version = "1.0.9"
|
||||
source = { editable = "../langgraph" }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
@@ -411,7 +411,7 @@ dependencies = [
|
||||
[package.metadata]
|
||||
requires-dist = [
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "orjson", specifier = ">=3.10.1" },
|
||||
{ name = "orjson", specifier = ">=3.11.5" },
|
||||
{ name = "psycopg", specifier = ">=3.2.0" },
|
||||
{ name = "psycopg-pool", specifier = ">=3.2.0" },
|
||||
]
|
||||
@@ -421,6 +421,7 @@ dev = [
|
||||
{ name = "anyio" },
|
||||
{ name = "codespell" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "mypy" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
@@ -437,6 +438,7 @@ lint = [
|
||||
test = [
|
||||
{ name = "anyio" },
|
||||
{ name = "langgraph-checkpoint", editable = "../checkpoint" },
|
||||
{ name = "langgraph-checkpoint-conformance", editable = "../checkpoint-conformance" },
|
||||
{ name = "psycopg", extras = ["binary"] },
|
||||
{ name = "pytest" },
|
||||
{ name = "pytest-asyncio" },
|
||||
@@ -489,7 +491,7 @@ test = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-prebuilt"
|
||||
version = "1.0.7"
|
||||
version = "1.0.8"
|
||||
source = { editable = "." }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
@@ -585,7 +587,7 @@ dependencies = [
|
||||
[package.metadata]
|
||||
requires-dist = [
|
||||
{ name = "httpx", specifier = ">=0.25.2" },
|
||||
{ name = "orjson", specifier = ">=3.10.1" },
|
||||
{ name = "orjson", specifier = ">=3.11.5" },
|
||||
]
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
|
||||
@@ -3,6 +3,6 @@ from langgraph_sdk.client import get_client, get_sync_client
|
||||
from langgraph_sdk.encryption import Encryption
|
||||
from langgraph_sdk.encryption.types import EncryptionContext
|
||||
|
||||
__version__ = "0.3.6"
|
||||
__version__ = "0.3.8"
|
||||
|
||||
__all__ = ["Auth", "Encryption", "EncryptionContext", "get_client", "get_sync_client"]
|
||||
|
||||
@@ -290,7 +290,7 @@ class AssistantsClient:
|
||||
"""
|
||||
get_params = {"recurse": recurse}
|
||||
if params:
|
||||
get_params = {**get_params, **params}
|
||||
get_params = {**get_params, **dict(params)}
|
||||
if namespace is not None:
|
||||
return await self.http.get(
|
||||
f"/assistants/{assistant_id}/subgraphs/{namespace}",
|
||||
@@ -425,9 +425,9 @@ class AssistantsClient:
|
||||
payload: dict[str, Any] = {}
|
||||
if graph_id:
|
||||
payload["graph_id"] = graph_id
|
||||
if config:
|
||||
if config is not None:
|
||||
payload["config"] = config
|
||||
if context:
|
||||
if context is not None:
|
||||
payload["context"] = context
|
||||
if metadata:
|
||||
payload["metadata"] = metadata
|
||||
|
||||
@@ -110,7 +110,7 @@ def get_client(
|
||||
if url is None:
|
||||
url = "http://api"
|
||||
if os.environ.get("__LANGGRAPH_DEFER_LOOPBACK_TRANSPORT") == "true":
|
||||
transport = get_asgi_transport()(app=None, root_path="/noauth")
|
||||
transport = get_asgi_transport()(app=None, root_path="/noauth") # type: ignore[invalid-argument-type]
|
||||
_registered_transports.append(transport)
|
||||
else:
|
||||
try:
|
||||
@@ -122,7 +122,7 @@ def get_client(
|
||||
"Failed to connect to in-process LangGraph server. Deferring configuration.",
|
||||
exc_info=True,
|
||||
)
|
||||
transport = get_asgi_transport()(app=None, root_path="/noauth")
|
||||
transport = get_asgi_transport()(app=None, root_path="/noauth") # type: ignore[invalid-argument-type]
|
||||
_registered_transports.append(transport)
|
||||
|
||||
if transport is None:
|
||||
|
||||
@@ -2,7 +2,8 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
import warnings
|
||||
from collections.abc import Mapping, Sequence
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
@@ -14,11 +15,13 @@ from langgraph_sdk.schema import (
|
||||
Cron,
|
||||
CronSelectField,
|
||||
CronSortBy,
|
||||
Durability,
|
||||
Input,
|
||||
OnCompletionBehavior,
|
||||
QueryParamTypes,
|
||||
Run,
|
||||
SortOrder,
|
||||
StreamMode,
|
||||
)
|
||||
|
||||
|
||||
@@ -60,13 +63,17 @@ class CronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -83,7 +90,7 @@ class CronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
@@ -93,6 +100,13 @@ class CronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled or not.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -118,6 +132,13 @@ class CronClient:
|
||||
)
|
||||
```
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -131,6 +152,10 @@ class CronClient:
|
||||
"webhook": webhook,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
if multitask_strategy:
|
||||
payload["multitask_strategy"] = multitask_strategy
|
||||
@@ -151,7 +176,7 @@ class CronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
@@ -159,6 +184,10 @@ class CronClient:
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -174,7 +203,7 @@ class CronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
webhook: Webhook to call after LangGraph API call is done.
|
||||
@@ -186,6 +215,13 @@ class CronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled or not.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -211,6 +247,13 @@ class CronClient:
|
||||
```
|
||||
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -225,6 +268,10 @@ class CronClient:
|
||||
"on_run_completed": on_run_completed,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
if multitask_strategy:
|
||||
payload["multitask_strategy"] = multitask_strategy
|
||||
@@ -277,6 +324,10 @@ class CronClient:
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
on_run_completed: OnCompletionBehavior | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Cron:
|
||||
@@ -299,6 +350,10 @@ class CronClient:
|
||||
after execution. 'keep' creates a new thread for each execution but does not
|
||||
clean them up.
|
||||
enabled: Enable or disable the cron job.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -329,6 +384,10 @@ class CronClient:
|
||||
"interrupt_after": interrupt_after,
|
||||
"on_run_completed": on_run_completed,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return await self.http.patch(
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import builtins
|
||||
import warnings
|
||||
from collections.abc import AsyncIterator, Callable, Mapping, Sequence
|
||||
from typing import Any, overload
|
||||
@@ -507,11 +508,11 @@ class RunsClient:
|
||||
|
||||
async def create_batch(
|
||||
self,
|
||||
payloads: list[RunCreate],
|
||||
payloads: builtins.list[RunCreate],
|
||||
*,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> list[Run]:
|
||||
) -> builtins.list[Run]:
|
||||
"""Create a batch of stateless background runs."""
|
||||
|
||||
def filter_payload(payload: RunCreate):
|
||||
@@ -547,7 +548,7 @@ class RunsClient:
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
on_run_created: Callable[[RunCreateMetadata], None] | None = None,
|
||||
) -> list[dict] | dict[str, Any]: ...
|
||||
) -> builtins.list[dict] | dict[str, Any]: ...
|
||||
|
||||
@overload
|
||||
async def wait(
|
||||
@@ -572,7 +573,7 @@ class RunsClient:
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
on_run_created: Callable[[RunCreateMetadata], None] | None = None,
|
||||
) -> list[dict] | dict[str, Any]: ...
|
||||
) -> builtins.list[dict] | dict[str, Any]: ...
|
||||
|
||||
async def wait(
|
||||
self,
|
||||
@@ -600,7 +601,7 @@ class RunsClient:
|
||||
params: QueryParamTypes | None = None,
|
||||
on_run_created: Callable[[RunCreateMetadata], None] | None = None,
|
||||
durability: Durability | None = None,
|
||||
) -> list[dict] | dict[str, Any]:
|
||||
) -> builtins.list[dict] | dict[str, Any]:
|
||||
"""Create a run, wait until it finishes and return the final state.
|
||||
|
||||
Args:
|
||||
@@ -751,10 +752,10 @@ class RunsClient:
|
||||
limit: int = 10,
|
||||
offset: int = 0,
|
||||
status: RunStatus | None = None,
|
||||
select: list[RunSelectField] | None = None,
|
||||
select: builtins.list[RunSelectField] | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> list[Run]:
|
||||
) -> builtins.list[Run]:
|
||||
"""List runs.
|
||||
|
||||
Args:
|
||||
|
||||
@@ -138,7 +138,7 @@ class StoreClient:
|
||||
if refresh_ttl is not None:
|
||||
get_params["refresh_ttl"] = refresh_ttl
|
||||
if params:
|
||||
get_params = {**get_params, **params}
|
||||
get_params = {**get_params, **dict(params)}
|
||||
return await self.http.get("/store/items", params=get_params, headers=headers)
|
||||
|
||||
async def delete_item(
|
||||
|
||||
@@ -543,7 +543,7 @@ class ThreadsClient:
|
||||
elif checkpoint_id:
|
||||
get_params = {"subgraphs": subgraphs}
|
||||
if params:
|
||||
get_params = {**get_params, **params}
|
||||
get_params = {**get_params, **dict(params)}
|
||||
return await self.http.get(
|
||||
f"/threads/{thread_id}/state/{checkpoint_id}",
|
||||
params=get_params,
|
||||
@@ -552,7 +552,7 @@ class ThreadsClient:
|
||||
else:
|
||||
get_params = {"subgraphs": subgraphs}
|
||||
if params:
|
||||
get_params = {**get_params, **params}
|
||||
get_params = {**get_params, **dict(params)}
|
||||
return await self.http.get(
|
||||
f"/threads/{thread_id}/state",
|
||||
params=get_params,
|
||||
|
||||
@@ -294,7 +294,7 @@ class SyncAssistantsClient:
|
||||
"""
|
||||
get_params = {"recurse": recurse}
|
||||
if params:
|
||||
get_params = {**get_params, **params}
|
||||
get_params = {**get_params, **dict(params)}
|
||||
if namespace is not None:
|
||||
return self.http.get(
|
||||
f"/assistants/{assistant_id}/subgraphs/{namespace}",
|
||||
@@ -427,9 +427,9 @@ class SyncAssistantsClient:
|
||||
payload: dict[str, Any] = {}
|
||||
if graph_id:
|
||||
payload["graph_id"] = graph_id
|
||||
if config:
|
||||
if config is not None:
|
||||
payload["config"] = config
|
||||
if context:
|
||||
if context is not None:
|
||||
payload["context"] = context
|
||||
if metadata:
|
||||
payload["metadata"] = metadata
|
||||
|
||||
@@ -2,7 +2,8 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
import warnings
|
||||
from collections.abc import Mapping, Sequence
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
@@ -14,11 +15,13 @@ from langgraph_sdk.schema import (
|
||||
Cron,
|
||||
CronSelectField,
|
||||
CronSortBy,
|
||||
Durability,
|
||||
Input,
|
||||
OnCompletionBehavior,
|
||||
QueryParamTypes,
|
||||
Run,
|
||||
SortOrder,
|
||||
StreamMode,
|
||||
)
|
||||
|
||||
|
||||
@@ -54,13 +57,17 @@ class SyncCronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -77,7 +84,7 @@ class SyncCronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
webhook: Webhook to call after LangGraph API call is done.
|
||||
@@ -85,6 +92,13 @@ class SyncCronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled. By default, it is considered enabled.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
|
||||
Returns:
|
||||
@@ -109,6 +123,13 @@ class SyncCronClient:
|
||||
)
|
||||
```
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -123,6 +144,10 @@ class SyncCronClient:
|
||||
"multitask_strategy": multitask_strategy,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return self.http.post(
|
||||
@@ -141,7 +166,7 @@ class SyncCronClient:
|
||||
metadata: Mapping[str, Any] | None = None,
|
||||
config: Config | None = None,
|
||||
context: Context | None = None,
|
||||
checkpoint_during: bool | None = None,
|
||||
checkpoint_during: bool | None = None, # deprecated
|
||||
interrupt_before: All | list[str] | None = None,
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
webhook: str | None = None,
|
||||
@@ -149,6 +174,10 @@ class SyncCronClient:
|
||||
multitask_strategy: str | None = None,
|
||||
end_time: datetime | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Run:
|
||||
@@ -164,7 +193,7 @@ class SyncCronClient:
|
||||
config: The configuration for the assistant.
|
||||
context: Static context to add to the assistant.
|
||||
!!! version-added "Added in version 0.6.0"
|
||||
checkpoint_during: Whether to checkpoint during the run (or only at the end/interruption).
|
||||
checkpoint_during: (deprecated) Whether to checkpoint during the run (or only at the end/interruption).
|
||||
interrupt_before: Nodes to interrupt immediately before they get executed.
|
||||
interrupt_after: Nodes to Nodes to interrupt immediately after they get executed.
|
||||
webhook: Webhook to call after LangGraph API call is done.
|
||||
@@ -176,6 +205,13 @@ class SyncCronClient:
|
||||
Must be one of 'reject', 'interrupt', 'rollback', or 'enqueue'.
|
||||
end_time: The time to stop running the cron job. If not provided, the cron job will run indefinitely.
|
||||
enabled: Whether the cron job is enabled. By default, it is considered enabled.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
"async" means checkpoints are persisted async while next graph step executes, replaces checkpoint_during=True
|
||||
"sync" means checkpoints are persisted sync after graph step executes, replaces checkpoint_during=False
|
||||
"exit" means checkpoints are only persisted when the run exits, does not save intermediate steps
|
||||
headers: Optional custom headers to include with the request.
|
||||
|
||||
Returns:
|
||||
@@ -201,6 +237,13 @@ class SyncCronClient:
|
||||
```
|
||||
|
||||
"""
|
||||
if checkpoint_during is not None:
|
||||
warnings.warn(
|
||||
"`checkpoint_during` is deprecated and will be removed in a future version. Use `durability` instead.",
|
||||
DeprecationWarning,
|
||||
stacklevel=2,
|
||||
)
|
||||
|
||||
payload = {
|
||||
"schedule": schedule,
|
||||
"input": input,
|
||||
@@ -216,6 +259,10 @@ class SyncCronClient:
|
||||
"multitask_strategy": multitask_strategy,
|
||||
"end_time": end_time.isoformat() if end_time else None,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return self.http.post(
|
||||
@@ -266,6 +313,10 @@ class SyncCronClient:
|
||||
interrupt_after: All | list[str] | None = None,
|
||||
on_run_completed: OnCompletionBehavior | None = None,
|
||||
enabled: bool | None = None,
|
||||
stream_mode: StreamMode | Sequence[StreamMode] | None = None,
|
||||
stream_subgraphs: bool | None = None,
|
||||
stream_resumable: bool | None = None,
|
||||
durability: Durability | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> Cron:
|
||||
@@ -288,6 +339,10 @@ class SyncCronClient:
|
||||
after execution. 'keep' creates a new thread for each execution but does not
|
||||
clean them up.
|
||||
enabled: Enable or disable the cron job.
|
||||
stream_mode: The stream mode(s) to use.
|
||||
stream_subgraphs: Whether to stream output from subgraphs.
|
||||
stream_resumable: Whether to persist the stream chunks in order to resume the stream later.
|
||||
durability: Durability level for the run. Must be one of 'sync', 'async', or 'exit'.
|
||||
headers: Optional custom headers to include with the request.
|
||||
params: Optional query parameters to include with the request.
|
||||
|
||||
@@ -318,6 +373,10 @@ class SyncCronClient:
|
||||
"interrupt_after": interrupt_after,
|
||||
"on_run_completed": on_run_completed,
|
||||
"enabled": enabled,
|
||||
"stream_mode": stream_mode,
|
||||
"stream_subgraphs": stream_subgraphs,
|
||||
"stream_resumable": stream_resumable,
|
||||
"durability": durability,
|
||||
}
|
||||
payload = {k: v for k, v in payload.items() if v is not None}
|
||||
return self.http.patch(
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import builtins
|
||||
import warnings
|
||||
from collections.abc import Callable, Iterator, Mapping, Sequence
|
||||
from typing import Any, overload
|
||||
@@ -503,11 +504,11 @@ class SyncRunsClient:
|
||||
|
||||
def create_batch(
|
||||
self,
|
||||
payloads: list[RunCreate],
|
||||
payloads: builtins.list[RunCreate],
|
||||
*,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> list[Run]:
|
||||
) -> builtins.list[Run]:
|
||||
"""Create a batch of stateless background runs."""
|
||||
|
||||
def filter_payload(payload: RunCreate):
|
||||
@@ -543,7 +544,7 @@ class SyncRunsClient:
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
on_run_created: Callable[[RunCreateMetadata], None] | None = None,
|
||||
) -> list[dict] | dict[str, Any]: ...
|
||||
) -> builtins.list[dict] | dict[str, Any]: ...
|
||||
|
||||
@overload
|
||||
def wait(
|
||||
@@ -568,7 +569,7 @@ class SyncRunsClient:
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
on_run_created: Callable[[RunCreateMetadata], None] | None = None,
|
||||
) -> list[dict] | dict[str, Any]: ...
|
||||
) -> builtins.list[dict] | dict[str, Any]: ...
|
||||
|
||||
def wait(
|
||||
self,
|
||||
@@ -596,7 +597,7 @@ class SyncRunsClient:
|
||||
params: QueryParamTypes | None = None,
|
||||
on_run_created: Callable[[RunCreateMetadata], None] | None = None,
|
||||
durability: Durability | None = None,
|
||||
) -> list[dict] | dict[str, Any]:
|
||||
) -> builtins.list[dict] | dict[str, Any]:
|
||||
"""Create a run, wait until it finishes and return the final state.
|
||||
|
||||
Args:
|
||||
@@ -740,10 +741,10 @@ class SyncRunsClient:
|
||||
limit: int = 10,
|
||||
offset: int = 0,
|
||||
status: RunStatus | None = None,
|
||||
select: list[RunSelectField] | None = None,
|
||||
select: builtins.list[RunSelectField] | None = None,
|
||||
headers: Mapping[str, str] | None = None,
|
||||
params: QueryParamTypes | None = None,
|
||||
) -> list[Run]:
|
||||
) -> builtins.list[Run]:
|
||||
"""List runs.
|
||||
|
||||
Args:
|
||||
|
||||
@@ -530,7 +530,7 @@ class SyncThreadsClient:
|
||||
elif checkpoint_id:
|
||||
get_params = {"subgraphs": subgraphs}
|
||||
if params:
|
||||
get_params = {**get_params, **params}
|
||||
get_params = {**get_params, **dict(params)}
|
||||
return self.http.get(
|
||||
f"/threads/{thread_id}/state/{checkpoint_id}",
|
||||
params=get_params,
|
||||
@@ -539,7 +539,7 @@ class SyncThreadsClient:
|
||||
else:
|
||||
get_params = {"subgraphs": subgraphs}
|
||||
if params:
|
||||
get_params = {**get_params, **params}
|
||||
get_params = {**get_params, **dict(params)}
|
||||
return self.http.get(
|
||||
f"/threads/{thread_id}/state",
|
||||
params=get_params,
|
||||
|
||||
@@ -72,8 +72,13 @@ class Auth:
|
||||
assert params.get("metadata", {}).get("owner") == "allowed_user"
|
||||
|
||||
@auth.on.store
|
||||
async def authorize_store(ctx: Auth.types.AuthContext, value: Auth.types.on):
|
||||
assert ctx.user.identity in value["namespace"], "Not authorized"
|
||||
async def authorize_store(ctx: Auth.types.AuthContext, value: Auth.types.on.store.value):
|
||||
# Automatically scope all store operations to the user's namespace.
|
||||
namespace = tuple(value["namespace"]) if value.get("namespace") else ()
|
||||
assert isinstance(namespace, tuple)
|
||||
if not namespace or namespace[0] != ctx.user.identity:
|
||||
namespace = (ctx.user.identity, *namespace)
|
||||
value["namespace"] = namespace
|
||||
```
|
||||
|
||||
???+ note "Request Processing Flow"
|
||||
@@ -170,13 +175,32 @@ class Auth:
|
||||
```
|
||||
|
||||
Auth for the `store` resource is a bit different since its structure is developer defined.
|
||||
You typically want to enforce user creds in the namespace.
|
||||
You typically want to scope store operations by rewriting the namespace to include the user's identity.
|
||||
The `value` dict is mutable — changes to `value["namespace"]` are used by the server for the actual operation.
|
||||
|
||||
```python
|
||||
@auth.on.store
|
||||
async def check_store_access(ctx: AuthContext, value: Auth.types.on) -> bool:
|
||||
# Assuming you structure your store like (store.aput((user_id, application_context), key, value))
|
||||
assert value["namespace"][0] == ctx.user.identity
|
||||
async def authorize_store(ctx: AuthContext, value: Auth.types.on.store.value):
|
||||
# Automatically scope all store operations to the user's namespace.
|
||||
namespace = tuple(value["namespace"]) if value.get("namespace") else ()
|
||||
assert isinstance(namespace, tuple)
|
||||
if not namespace or namespace[0] != ctx.user.identity:
|
||||
namespace = (ctx.user.identity, *namespace)
|
||||
value["namespace"] = namespace
|
||||
```
|
||||
|
||||
You can also register handlers for specific store actions:
|
||||
|
||||
```python
|
||||
@auth.on.store.put
|
||||
async def on_put(ctx: AuthContext, value: Auth.types.on.store.put.value):
|
||||
# value has typed fields: namespace, key, value, index
|
||||
...
|
||||
|
||||
@auth.on.store.get
|
||||
async def on_get(ctx: AuthContext, value: Auth.types.on.store.get.value):
|
||||
# value has typed fields: namespace, key
|
||||
...
|
||||
```
|
||||
"""
|
||||
# These are accessed by the API. Changes to their names or types is
|
||||
@@ -483,9 +507,85 @@ class _CronsOn(
|
||||
Search = types.CronsSearch
|
||||
|
||||
|
||||
class _StoreActionOn(typing.Generic[T]):
|
||||
"""Decorator for registering a handler for a specific store action."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
auth: Auth,
|
||||
action: typing.Literal["put", "get", "search", "delete", "list_namespaces"],
|
||||
value: type[T],
|
||||
) -> None:
|
||||
self.auth = auth
|
||||
self.action = action
|
||||
self.value = value
|
||||
|
||||
def __call__(self, fn: _ActionHandler[T]) -> _ActionHandler[T]:
|
||||
_validate_handler(fn)
|
||||
_register_handler(self.auth, "store", self.action, fn)
|
||||
return fn
|
||||
|
||||
|
||||
class _StoreOn:
|
||||
def __init__(self, auth: Auth) -> None:
|
||||
self._auth = auth
|
||||
self.put = _StoreActionOn(auth, "put", types.StorePut)
|
||||
"""Register a handler for store put operations.
|
||||
|
||||
???+ example "Example"
|
||||
```python
|
||||
@auth.on.store.put
|
||||
async def on_store_put(ctx: Auth.types.AuthContext, value: Auth.types.on.store.put.value):
|
||||
# Scope puts to user's namespace
|
||||
...
|
||||
```
|
||||
"""
|
||||
self.get = _StoreActionOn(auth, "get", types.StoreGet)
|
||||
"""Register a handler for store get operations.
|
||||
|
||||
???+ example "Example"
|
||||
```python
|
||||
@auth.on.store.get
|
||||
async def on_store_get(ctx: Auth.types.AuthContext, value: Auth.types.on.store.get.value):
|
||||
# Scope gets to user's namespace
|
||||
...
|
||||
```
|
||||
"""
|
||||
self.search = _StoreActionOn(auth, "search", types.StoreSearch)
|
||||
"""Register a handler for store search operations.
|
||||
|
||||
???+ example "Example"
|
||||
```python
|
||||
@auth.on.store.search
|
||||
async def on_store_search(ctx: Auth.types.AuthContext, value: Auth.types.on.store.search.value):
|
||||
# Scope searches to user's namespace
|
||||
...
|
||||
```
|
||||
"""
|
||||
self.delete = _StoreActionOn(auth, "delete", types.StoreDelete)
|
||||
"""Register a handler for store delete operations.
|
||||
|
||||
???+ example "Example"
|
||||
```python
|
||||
@auth.on.store.delete
|
||||
async def on_store_delete(ctx: Auth.types.AuthContext, value: Auth.types.on.store.delete.value):
|
||||
# Scope deletes to user's namespace
|
||||
...
|
||||
```
|
||||
"""
|
||||
self.list_namespaces = _StoreActionOn(
|
||||
auth, "list_namespaces", types.StoreListNamespaces
|
||||
)
|
||||
"""Register a handler for store list_namespaces operations.
|
||||
|
||||
???+ example "Example"
|
||||
```python
|
||||
@auth.on.store.list_namespaces
|
||||
async def on_list_ns(ctx: Auth.types.AuthContext, value: Auth.types.on.store.list_namespaces.value):
|
||||
# Scope namespace listing to user's prefix
|
||||
...
|
||||
```
|
||||
"""
|
||||
|
||||
@typing.overload
|
||||
def __call__(
|
||||
|
||||
@@ -402,7 +402,7 @@ class AuthContext(BaseAuthContext):
|
||||
"list_namespaces",
|
||||
]
|
||||
"""The action being performed on the resource.
|
||||
|
||||
|
||||
Most resources support the following actions:
|
||||
- create: Create a new resource
|
||||
- read: Read information about a resource
|
||||
@@ -411,8 +411,10 @@ class AuthContext(BaseAuthContext):
|
||||
- search: Search for resources
|
||||
|
||||
The store supports the following actions:
|
||||
- put: Add or update a document in the store
|
||||
- get: Get a document from the store
|
||||
- put: Add or update an item in the store
|
||||
- get: Get an item from the store
|
||||
- search: Search for items within a namespace prefix
|
||||
- delete: Delete an item from the store
|
||||
- list_namespaces: List the namespaces in the store
|
||||
"""
|
||||
|
||||
@@ -851,20 +853,34 @@ class CronsSearch(typing.TypedDict, total=False):
|
||||
|
||||
|
||||
class StoreGet(typing.TypedDict):
|
||||
"""Operation to retrieve a specific item by its namespace and key."""
|
||||
"""Operation to retrieve a specific item by its namespace and key.
|
||||
|
||||
This dict is mutable — auth handlers can modify `namespace` to enforce
|
||||
access scoping (e.g., prepending the user's identity).
|
||||
"""
|
||||
|
||||
namespace: tuple[str, ...]
|
||||
"""Hierarchical path that uniquely identifies the item's location."""
|
||||
"""Hierarchical path that uniquely identifies the item's location.
|
||||
|
||||
Auth handlers can modify this to enforce per-user scoping.
|
||||
"""
|
||||
|
||||
key: str
|
||||
"""Unique identifier for the item within its specific namespace."""
|
||||
|
||||
|
||||
class StoreSearch(typing.TypedDict):
|
||||
"""Operation to search for items within a specified namespace hierarchy."""
|
||||
"""Operation to search for items within a specified namespace hierarchy.
|
||||
|
||||
This dict is mutable — auth handlers can modify `namespace` to enforce
|
||||
access scoping (e.g., prepending the user's identity).
|
||||
"""
|
||||
|
||||
namespace: tuple[str, ...]
|
||||
"""Prefix filter for defining the search scope."""
|
||||
"""Prefix filter for defining the search scope.
|
||||
|
||||
Auth handlers can modify this to enforce per-user scoping.
|
||||
"""
|
||||
|
||||
filter: dict[str, typing.Any] | None
|
||||
"""Key-value pairs for filtering results based on exact matches or comparison operators."""
|
||||
@@ -876,14 +892,22 @@ class StoreSearch(typing.TypedDict):
|
||||
"""Number of matching items to skip for pagination."""
|
||||
|
||||
query: str | None
|
||||
"""Naturalj language search query for semantic search capabilities."""
|
||||
"""Natural language search query for semantic search capabilities."""
|
||||
|
||||
|
||||
class StoreListNamespaces(typing.TypedDict):
|
||||
"""Operation to list and filter namespaces in the store."""
|
||||
"""Operation to list and filter namespaces in the store.
|
||||
|
||||
This dict is mutable — auth handlers can modify `namespace` (the prefix)
|
||||
to enforce access scoping (e.g., prepending the user's identity).
|
||||
"""
|
||||
|
||||
namespace: tuple[str, ...] | None
|
||||
"""Prefix filter namespaces."""
|
||||
"""Prefix filter for namespaces. Can be `None` if no prefix was provided.
|
||||
|
||||
Auth handlers can modify this to enforce per-user scoping. When `None`,
|
||||
handlers should set it to `(user_id,)` to scope listing to the user's namespaces.
|
||||
"""
|
||||
|
||||
suffix: tuple[str, ...] | None
|
||||
"""Optional conditions for filtering namespaces."""
|
||||
@@ -903,10 +927,17 @@ class StoreListNamespaces(typing.TypedDict):
|
||||
|
||||
|
||||
class StorePut(typing.TypedDict):
|
||||
"""Operation to store, update, or delete an item in the store."""
|
||||
"""Operation to store, update, or delete an item in the store.
|
||||
|
||||
This dict is mutable — auth handlers can modify `namespace` to enforce
|
||||
access scoping (e.g., prepending the user's identity).
|
||||
"""
|
||||
|
||||
namespace: tuple[str, ...]
|
||||
"""Hierarchical path that identifies the location of the item."""
|
||||
"""Hierarchical path that identifies the location of the item.
|
||||
|
||||
Auth handlers can modify this to enforce per-user scoping.
|
||||
"""
|
||||
|
||||
key: str
|
||||
"""Unique identifier for the item within its namespace."""
|
||||
@@ -919,10 +950,17 @@ class StorePut(typing.TypedDict):
|
||||
|
||||
|
||||
class StoreDelete(typing.TypedDict):
|
||||
"""Operation to delete an item from the store."""
|
||||
"""Operation to delete an item from the store.
|
||||
|
||||
This dict is mutable — auth handlers can modify `namespace` to enforce
|
||||
access scoping (e.g., prepending the user's identity).
|
||||
"""
|
||||
|
||||
namespace: tuple[str, ...]
|
||||
"""Hierarchical path that uniquely identifies the item's location."""
|
||||
"""Hierarchical path that uniquely identifies the item's location.
|
||||
|
||||
Auth handlers can modify this to enforce per-user scoping.
|
||||
"""
|
||||
|
||||
key: str
|
||||
"""Unique identifier for the item within its specific namespace."""
|
||||
|
||||
@@ -146,7 +146,9 @@ AssistantSortBy = Literal[
|
||||
The field to sort by.
|
||||
"""
|
||||
|
||||
ThreadSortBy = Literal["thread_id", "status", "created_at", "updated_at"]
|
||||
ThreadSortBy = Literal[
|
||||
"thread_id", "status", "created_at", "updated_at", "state_updated_at"
|
||||
]
|
||||
"""
|
||||
The field to sort by.
|
||||
"""
|
||||
@@ -422,6 +424,14 @@ class CronUpdate(TypedDict, total=False):
|
||||
"""What to do with the thread after the run completes."""
|
||||
enabled: bool
|
||||
"""Enable or disable the cron job."""
|
||||
stream_mode: StreamMode | list[StreamMode]
|
||||
"""The stream mode(s) to use."""
|
||||
stream_subgraphs: bool
|
||||
"""Whether to stream output from subgraphs."""
|
||||
stream_resumable: bool
|
||||
"""Whether to persist the stream chunks in order to resume the stream later."""
|
||||
durability: Durability
|
||||
"""Durability level for the run. Must be one of 'sync', 'async', or 'exit'."""
|
||||
|
||||
|
||||
# Select field aliases for client-side typing of `select` parameters.
|
||||
|
||||
@@ -11,7 +11,7 @@ requires-python = ">=3.10"
|
||||
readme = "README.md"
|
||||
license = "MIT"
|
||||
license-files = ['LICENSE']
|
||||
dependencies = ["httpx>=0.25.2", "orjson>=3.10.1"]
|
||||
dependencies = ["httpx>=0.25.2", "orjson>=3.11.5"]
|
||||
|
||||
[tool.hatch.version]
|
||||
path = "langgraph_sdk/__init__.py"
|
||||
|
||||
@@ -19,7 +19,7 @@ class AsyncListByteStream(httpx.AsyncByteStream):
|
||||
self._chunks = list(chunks)
|
||||
self._exc = exc
|
||||
|
||||
async def __aiter__(self): # type: ignore[override]
|
||||
async def __aiter__(self):
|
||||
for chunk in self._chunks:
|
||||
yield chunk
|
||||
if self._exc is not None:
|
||||
@@ -34,7 +34,7 @@ class ListByteStream(httpx.ByteStream):
|
||||
self._chunks = list(chunks)
|
||||
self._exc = exc
|
||||
|
||||
def __iter__(self): # type: ignore[override]
|
||||
def __iter__(self):
|
||||
yield from self._chunks
|
||||
if self._exc is not None:
|
||||
raise self._exc
|
||||
|
||||
@@ -65,7 +65,7 @@ def test_raise_for_status_typed_maps_exceptions_and_sets_status_code(
|
||||
with pytest.raises(exc_type) as ei:
|
||||
_raise_for_status_typed(r)
|
||||
|
||||
err = cast("APIStatusError", ei.value)
|
||||
err = ei.value
|
||||
assert err.status_code == status
|
||||
# response attribute should be present and match
|
||||
assert err.response.status_code == status
|
||||
@@ -113,7 +113,7 @@ def test_error_message_in_str_and_args() -> None:
|
||||
r = make_response(422, json_body={"message": "Validation failed"})
|
||||
with pytest.raises(UnprocessableEntityError) as ei:
|
||||
_raise_for_status_typed(r)
|
||||
err = cast("UnprocessableEntityError", ei.value)
|
||||
err = ei.value
|
||||
assert str(err) == "Validation failed"
|
||||
assert err.args == ("Validation failed",)
|
||||
assert err.message == "Validation failed"
|
||||
|
||||
Generated
+3
-3
@@ -265,7 +265,7 @@ wheels = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph"
|
||||
version = "1.0.8"
|
||||
version = "1.0.9"
|
||||
source = { editable = "../langgraph" }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
@@ -396,7 +396,7 @@ test = [
|
||||
|
||||
[[package]]
|
||||
name = "langgraph-prebuilt"
|
||||
version = "1.0.7"
|
||||
version = "1.0.8"
|
||||
source = { editable = "../prebuilt" }
|
||||
dependencies = [
|
||||
{ name = "langchain-core" },
|
||||
@@ -484,7 +484,7 @@ test = [
|
||||
[package.metadata]
|
||||
requires-dist = [
|
||||
{ name = "httpx", specifier = ">=0.25.2" },
|
||||
{ name = "orjson", specifier = ">=3.10.1" },
|
||||
{ name = "orjson", specifier = ">=3.11.5" },
|
||||
]
|
||||
|
||||
[package.metadata.requires-dev]
|
||||
|
||||
Reference in New Issue
Block a user