Compare commits

..
Author SHA1 Message Date
Eugene Yurtsev 0b8634b4c6 x 2025-06-25 10:25:06 -04:00
Eugene Yurtsev bc3ef7f913 x 2025-06-24 14:47:20 -04:00
Eugene Yurtsev 0dda1b4b1e x 2025-06-24 12:35:31 -04:00
Eugene Yurtsev d22c2c4dac x 2025-06-24 11:17:21 -04:00
Eugene Yurtsev db03dccb2b x 2025-06-24 09:56:19 -04:00
Eugene Yurtsev c79c9ea733 x 2025-06-20 17:06:13 -04:00
Eugene Yurtsev b115e1dcde tools 2025-06-20 17:04:41 -04:00
Eugene Yurtsev 3b59213311 fix tools 2025-06-20 16:45:24 -04:00
Eugene Yurtsev 0c0e5a299d Replace notebook with markdown file 2025-06-20 16:33:42 -04:00
Lauren Hirata Singh 69d4c37d25 Consolidate assistant conceptual guides 2025-06-18 19:05:34 -04:00
Lauren Hirata Singh 0e7554a1a1 Update navigation 2025-06-18 10:58:18 -04:00
Lauren Hirata Singh 596c60a65c Move LGP to platform section 2025-06-18 10:58:12 -04:00
Lauren Hirata Singh e45797ce19 Fix titles based on feedback 2025-06-18 10:49:55 -04:00
Lauren Hirata Singh e746b54a57 Change titles 2025-06-17 16:30:35 -04:00
Lauren Hirata SinghandGitHub c88e22ffa7 Merge branch 'main' into get-started 2025-06-17 16:16:36 -04:00
Lauren Hirata Singh dfdeb6a6f1 Edit stream modes 2025-06-17 15:41:57 -04:00
Lauren Hirata Singh ba08acb71c Fix broken links 2025-06-17 15:21:50 -04:00
Lauren Hirata Singh dd64636ca8 Remove agents/streaming 2025-06-17 15:14:48 -04:00
Lauren Hirata Singh 664475887d Consolidate streaming 2025-06-17 15:07:07 -04:00
Lauren Hirata Singh 3d3a2bfacd edits 2025-06-16 21:07:12 -04:00
Lauren Hirata Singh bbe90e04ca Remove cookie consent popup 2025-06-16 16:35:27 -04:00
Lauren Hirata Singh 905fcb3d02 Fix broken links 2025-06-16 16:25:31 -04:00
Lauren Hirata Singh 543d7d85af Organize existing content differently 2025-06-16 16:17:52 -04:00
70 changed files with 2310 additions and 2695 deletions
+122 -119
View File
@@ -1,154 +1,157 @@
"""Translate Python markdown to TypeScript and/or consolidate Python-JS markdown into a single document."""
"""Add typescript translation to a given markdown file."""
import argparse
import re
import requests
from langchain_anthropic import ChatAnthropic
# Load reference TypeScript snippets
URL = "https://gist.githubusercontent.com/eyurtsev/e7486731415463a9bc5b4682358859c8/raw/b5a5fda9c7e3387cfcb781f25082814d43675d50/gistfile1.txt"
response = requests.get(URL)
response.raise_for_status()
reference_snippets = response.text
# Initialize model
model = ChatAnthropic(model="claude-sonnet-4-0", max_tokens=64_000)
TRANSLATION_PROMPT = (
"You are a helpful assistant that translates Python-based technical "
"documentation written in Markdown to equivalent TypeScript-based documentation. "
"The input is a Markdown file written in mkdocs format. It contains "
"Python code snippets embedded in prose. "
"Your task is to rewrite the content by translating the Python code to "
"idiomatic TypeScript, using the provided TypeScript reference snippets "
"to ensure accurate and consistent usage (e.g., correct imports, function "
"names, and patterns). "
"Remove the original Python code and replace it with the corresponding "
"TypeScript version. "
"Do not alter the surrounding prose unless a change is necessary to "
"reflect differences between Python and TypeScript. "
"Preserve the structure and formatting of the original Markdown document. "
"Do not make stylistic or structural changes unless they directly support "
"the translation. "
"Use the reference TypeScript snippets as guidance whenever possible to "
"maintain alignment with existing conventions.\n\n"
f"Here are the reference TypeScript snippets:\n\n{reference_snippets}\n\n"
)
CONSOLIDATION_PROMPT = (
"You are a helpful assistant that consolidates parallel Python and JavaScript (TypeScript) technical documentation "
"written in Markdown into a single unified Markdown document. "
"The input consists of two documents: the first is for Python users, and the second is for JavaScript/TypeScript users. "
"Your task is to merge these into one Markdown file using language-specific fenced blocks to separate the content where needed. "
"Use the following syntax to distinguish content for each language:\n\n"
":::python\n"
"# Python-specific content\n"
":::\n\n"
":::js\n"
"# JavaScript/TypeScript-specific content\n"
":::\n\n"
"Follow these consolidation rules:\n"
"- When content (prose or code) is the same or nearly identical in both versions, include it only once—outside of any fenced block.\n"
"- When content differs between the Python and JS versions, wrap each version in its corresponding fenced block.\n"
"- Prefer **paragraph-level separation** of language-specific content. Do not combine Python and JS snippets or terminology in the same sentence or paragraph using conditional phrases.\n"
" For example, avoid inline constructs like:\n"
" `The :::python add_messages ::: :::js reducer ::: function...`\n"
" Instead, write two distinct paragraphs:\n\n"
" :::python\n"
" The `add_messages` function in our `State` will append the LLM's response messages to whatever messages are already in the state.\n"
" ::: \n\n"
" :::js\n"
" The `reducer` function in our `StateAnnotation` will append the LLM's response messages to whatever messages are already in the state.\n"
" :::\n\n"
"- Preserve the overall structure, ordering, and formatting of the original Markdown documents.\n"
"- Do not rephrase or unify content unless it is logically and semantically identical.\n"
"- Use the fenced blocks for both prose and code as needed, and ensure output is clean, readable Markdown suitable for tools that parse these directives.\n"
"Your goal is to produce a cleanly merged documentation file that serves both Python and JavaScript users without redundancy, while maximizing clarity and separation of language-specific details."
)
model = ChatAnthropic(model="claude-3-5-sonnet-latest")
def translate_python_to_ts(markdown_content: str) -> str:
response = model.invoke(
def _get_tqdm():
try:
from tqdm import tqdm
except ImportError:
# If not available return a simple identity function
def tqdm(iterable, *args, **kwargs):
return iterable
return tqdm
_tqdm = _get_tqdm()
opening_pattern = re.compile(r"^\s*```python(?:\s+.*)?\s*$")
closing_pattern = re.compile(r"^\s*```\s*$")
def extract_python_snippets(markdown: str) -> list[str]:
"""
Extract all python code blocks (including their fence lines) from the markdown content.
A python block is defined as any block that starts with a line containing an opening fence
with '```python' (optionally with extra parameters) and ends with a closing fence '```'.
"""
snippets = []
inside_block = False
current_snippet = []
for line in markdown.splitlines(keepends=True):
if not inside_block:
if opening_pattern.match(line):
inside_block = True
current_snippet = [line]
else:
current_snippet.append(line)
if closing_pattern.match(line):
inside_block = False
snippets.append("".join(current_snippet))
current_snippet = []
return snippets
def translate_snippet(python_snippet: str) -> str:
"""Translate a python code block into a TypeScript code block using Langchain.
The response is expected to be a properly fenced TypeScript code block (i.e.
starting with ```typescript and ending with ```).
"""
ai_message = model.invoke(
[
{
"role": "system",
"content": TRANSLATION_PROMPT,
"cache_control": {"type": "ephemeral"},
"content": (
f"You have access to the following up-to-date example TypeScript code "
f"snippets that show examples of building with langgraph "
f"and langchain:\n\n{reference_snippets}\n\n"
"Use this context to translate the following Python code to equivalent "
"TypeScript. Ensure that your output is a valid fenced TypeScript "
"code block (i.e. starts with ```typescript and ends with ```)."
),
},
{"role": "user", "content": markdown_content},
]
)
return response.content
def consolidate_python_and_ts(combined_content: str) -> str:
response = model.invoke(
[
{
"role": "system",
"content": CONSOLIDATION_PROMPT,
"cache_control": {"type": "ephemeral"},
"role": "user",
"content": f"Translate this Python snippet to TypeScript:\n\n{python_snippet}",
},
{"role": "user", "content": combined_content},
]
)
return response.content
# Use a regular expression to search for a TypeScript code block in the response.
pattern = r"```typescript\s*(.*?)\s*```"
match = re.search(pattern, ai_message.content, re.DOTALL)
if match:
# Reconstruct the code block with proper fences.
typescript_code = match.group(1).strip()
return f"```typescript\n{typescript_code}\n```"
else:
raise ValueError("No TypeScript code block found in the model's response.")
def main(file_path: str, translate_only: bool, consolidate_only: bool) -> None:
with open(file_path, "r", encoding="utf-8") as f:
def insert_translations_into_markdown(
markdown: str, typescript_snippets: list[str]
) -> str:
"""Walks through the original markdown content and, after each
Python snippet block, inserts the corresponding translated TypeScript snippet.
It assumes that the ordering of the Python snippets
(from extract_python_snippets) matches the order they appear in the markdown.
"""
output_lines = []
lines = markdown.splitlines(keepends=True)
inside_block = False
snippet_index = 0
for line in lines:
output_lines.append(line)
if not inside_block and opening_pattern.match(line):
# We've encountered the start of a python code block.
inside_block = True
elif inside_block:
if closing_pattern.match(line):
# End of a python snippet block.
inside_block = False
if snippet_index < len(typescript_snippets):
# Insert an extra newline for clarity, then the translated TypeScript snippet.
output_lines.append("\n")
output_lines.append(typescript_snippets[snippet_index])
output_lines.append("\n")
snippet_index += 1
return "".join(output_lines)
def main(file_path: str) -> None:
# Read the markdown file.
with open(file_path, "r") as f:
markdown_content = f.read()
if translate_only:
translated = translate_python_to_ts(markdown_content)
output_path = file_path.replace(".md", ".translated.md")
with open(output_path, "w", encoding="utf-8") as f:
f.write(translated)
print(f"Translated JS/TS version written to: {output_path}")
# 1. Extract all Python snippets.
python_snippets = extract_python_snippets(markdown_content)[:1]
elif consolidate_only:
consolidated = consolidate_python_and_ts(markdown_content)
with open(file_path, "w", encoding="utf-8") as f:
f.write(consolidated)
print(f"Consolidated content written to: {file_path}")
# 2. Translate each Python snippet to TypeScript.
typescript_snippets = []
# Replace with .batch() for faster translation
for python_snippet in _tqdm(python_snippets):
ts_snippet = translate_snippet(python_snippet)
typescript_snippets.append(ts_snippet)
else:
# Default behavior: translate first, then consolidate both
translated = translate_python_to_ts(markdown_content)
combined = f"{markdown_content.strip()}\n\n\n{translated.strip()}"
consolidated = consolidate_python_and_ts(combined)
with open(file_path, "w", encoding="utf-8") as f:
f.write(consolidated)
print(f"Translated and consolidated content written to: {file_path}")
# 3. Insert the TypeScript translations after their respective Python snippets.
updated_markdown = insert_translations_into_markdown(
markdown_content, typescript_snippets
)
# Overwrite the original markdown file with the updated content.
with open(file_path, "w") as f:
f.write(updated_markdown)
if __name__ == "__main__":
parser = argparse.ArgumentParser(
description=(
"Translate Python markdown to TypeScript and/or consolidate "
"Python-JS markdown into one file."
)
description="Translate Python snippets in a markdown file to TypeScript and insert them after each Python snippet."
)
parser.add_argument("file_path", type=str, help="Path to the markdown file.")
parser.add_argument(
"--translate-only",
action="store_true",
help="Only generate the JS translation.",
)
parser.add_argument(
"--consolidate-only",
action="store_true",
help="Only consolidate pre-paired Python and JS content.",
)
args = parser.parse_args()
if args.translate_only and args.consolidate_only:
raise ValueError(
"Cannot use both --translate-only and --consolidate-only at the same time."
)
main(
args.file_path,
translate_only=args.translate_only,
consolidate_only=args.consolidate_only,
)
main(args.file_path)
+6 -10
View File
@@ -3,21 +3,19 @@
import asyncio
import glob
import os
import re
from typing import TypedDict, List, Optional
import pydantic
import re
from pydantic import BaseModel, Field
from langchain_core.rate_limiters import InMemoryRateLimiter
import yaml
from langchain.chat_models import init_chat_model
from langchain_core.rate_limiters import InMemoryRateLimiter
from mkdocs.structure.files import File
from mkdocs.structure.pages import Page
from pydantic import BaseModel, Field
from yaml import SafeLoader
from _scripts.notebook_hooks import (
_on_page_markdown_with_config,
_apply_conditional_rendering,
)
from _scripts.notebook_hooks import _on_page_markdown_with_config
HERE = os.path.dirname(os.path.abspath(__file__))
# Get source directory (parent of HERE / docs)
@@ -213,9 +211,7 @@ async def process_nav_items(nav_items: list[NavItem]) -> list[NavItem]:
# Remove any items that start with http:// or https:// looking only for
# local file at this stages.
nav_items = [
item
for item in nav_items
if not item["url"].startswith(("http://", "https://"))
item for item in nav_items if not item["url"].startswith(("http://", "https://"))
]
# Process items in parallel
tasks = [process_single_item(item) for item in nav_items]
-5
View File
@@ -1,5 +0,0 @@
JS_LINK_MAP = {
"langgraph.types.interrupt": "https://langchain-ai.github.io/langgraphjs/reference/functions/langgraph.interrupt-2.html",
"create_react_agent": "https://langchain-ai.github.io/langgraphjs/reference/functions/langgraph_prebuilt.createReactAgent.html",
"langgraph.types.Command": "https://langchain-ai.github.io/langgraphjs/reference/classes/langgraph.Command.html",
}
+6 -72
View File
@@ -16,7 +16,6 @@ from mkdocs.structure.pages import Page
from _scripts.generate_api_reference_links import update_markdown_with_imports
from _scripts.notebook_convert import convert_notebook
from _scripts.link_map import JS_LINK_MAP
logger = logging.getLogger(__name__)
logging.basicConfig()
@@ -62,6 +61,7 @@ REDIRECT_MAP = {
"how-tos/subgraph-persistence.ipynb": "how-tos/persistence.ipynb#use-with-subgraphs",
"how-tos/cross-thread-persistence.ipynb": "how-tos/persistence.ipynb#add-long-term-memory",
"cloud/how-tos/copy_threads": "cloud/how-tos/use_threads",
"cloud/concepts/threads.md": "concepts/persistence.md#threads",
# tool calling how-tos
"how-tos/tool-calling-errors.ipynb": "how-tos/tool-calling.ipynb#handle-errors",
"how-tos/pass-config-to-tools.ipynb": "how-tos/tool-calling.ipynb#access-config",
@@ -87,7 +87,9 @@ REDIRECT_MAP = {
"cloud/how-tos/stream_events.md": "cloud/how-tos/streaming.md#stream-events",
"cloud/how-tos/stream_debug.md": "cloud/how-tos/streaming.md#debug",
"cloud/how-tos/stream_multiple.md": "cloud/how-tos/streaming.md#stream-multiple-modes",
# prebuilt redirects
"cloud/concepts/streaming.md": "concepts/streaming.md",
"agents/streaming.md": "how-tos/streaming.md",
# prebuit redirects
"how-tos/create-react-agent.ipynb": "agents/agents.md#basic-configuration",
"how-tos/create-react-agent-memory.ipynb": "agents/memory.md",
"how-tos/create-react-agent-system-prompt.ipynb": "agents/context.md#prompts",
@@ -108,8 +110,10 @@ REDIRECT_MAP = {
# deployment redirects
"how-tos/deploy-self-hosted.md": "cloud/deployment/self_hosted_data_plane.md",
"concepts/self_hosted.md": "concepts/langgraph_self_hosted_data_plane.md",
"tutorials/deployment.md": "concepts/deployment_options.md",
# assistant redirects
"cloud/how-tos/assistant_versioning.md": "cloud/how-tos/configuration_cloud.md",
"cloud/concepts/runs.md": "concepts/assistants.md#execution",
}
@@ -159,62 +163,6 @@ def _add_path_to_code_blocks(markdown: str, page: Page) -> str:
return code_block_pattern.sub(replace_code_block_header, markdown)
def _resolve_cross_references(md_text: str, link_map: dict[str, str]) -> str:
"""Replace [title][identifier] with [title](url) using language-specific link_map.
Args:
md_text: The markdown text to process.
link_map: mapping of identifier to URL.
Returns:
The processed markdown text with cross-references resolved.
"""
# Pattern to match [title][identifier]
pattern = re.compile(r"\[([^\]]+)\]\[([^\]]+)\]")
def replace_reference(match: re.Match) -> str:
"""Replace the matched reference with the corresponding URL."""
title, identifier = match.group(1), match.group(2)
url = link_map.get(identifier)
if url:
return f"[{title}]({url})"
else:
# Leave it unchanged if not found
return match.group(0)
return pattern.sub(replace_reference, md_text)
def _apply_conditional_rendering(md_text: str, target_language: str) -> str:
if target_language not in {"python", "js"}:
raise ValueError("target_language must be 'python' or 'js'")
pattern = re.compile(
r"(?P<indent>[ \t]*):::(?P<language>\w+)\s*\n"
r"(?P<content>((?:.*\n)*?))" # Capture the content inside the block
r"(?P=indent):::" # Match closing with the same indentation
)
def replace_conditional_blocks(match: re.Match) -> str:
"""Keep active conditionals."""
language = match.group("language")
content = match.group("content")
if language not in {"python", "js"}:
# If the language is not supported, return the original block
return match.group(0)
if language == target_language:
return content
# If the language does not match, return an empty string
return ""
processed = pattern.sub(replace_conditional_blocks, md_text)
return processed
def _highlight_code_blocks(markdown: str) -> str:
"""Find code blocks with highlight comments and add hl_lines attribute.
@@ -314,20 +262,6 @@ def _on_page_markdown_with_config(
# Apply highlight comments to code blocks
markdown = _highlight_code_blocks(markdown)
# Apply conditional rendering for code blocks
target_language = kwargs.get("target_language", "python")
markdown = _apply_conditional_rendering(markdown, target_language)
if target_language == "js":
markdown = _resolve_cross_references(markdown, JS_LINK_MAP)
elif target_language == "python":
# Via a dedicated plugin
pass
else:
raise ValueError(
f"Unsupported target language: {target_language}. "
"Supported languages are 'python' and 'js'."
)
# Add file path as an attribute to code blocks that are executable.
# This file path is used to associate fixtures with the executable code
# which can be used in CI to test the docs without making network requests.
+1 -1
View File
@@ -89,4 +89,4 @@ LangGraph Studio Web is a specialized UI that you can connect to LangGraph API s
## Deployment
Once your LangGraph app is running locally, you can deploy it using LangGraph Platform. Refer to the [deployment options guide](../tutorials/deployment.md) for detailed instructions on all supported deployment models.
Once your LangGraph app is running locally, you can deploy it using LangGraph Platform. Refer to the [deployment options guide](../concepts/deployment_options.md) for detailed instructions on all supported deployment models.
+2 -2
View File
@@ -29,10 +29,10 @@ LangGraph includes several capabilities essential for building robust, productio
- [**Memory integration**](./memory.md): Native support for *short-term* (session-based) and *long-term* (persistent across sessions) memory, enabling stateful behaviors in chatbots and assistants.
- [**Human-in-the-loop control**](./human-in-the-loop.md): Execution can pause *indefinitely* to await human feedback—unlike websocket-based solutions limited to real-time interaction. This enables asynchronous approval, correction, or intervention at any point in the workflow.
- [**Streaming support**](./streaming.md): Real-time streaming of agent state, model tokens, tool outputs, or combined streams.
- [**Streaming support**](../how-tos/streaming.md): Real-time streaming of agent state, model tokens, tool outputs, or combined streams.
- [**Deployment tooling**](./deployment.md): Includes infrastructure-free deployment tools. [**LangGraph Platform**](https://langchain-ai.github.io/langgraph/concepts/langgraph_platform/) supports testing, debugging, and deployment.
- **[Studio](https://langchain-ai.github.io/langgraph/concepts/langgraph_studio/)**: A visual IDE for inspecting and debugging workflows.
- Supports multiple [**deployment options**](https://langchain-ai.github.io/langgraph/tutorials/deployment/) for production.
- Supports multiple [**deployment options**](https://langchain-ai.github.io/langgraph/concepts/deployment_options.md) for production.
## High-level building blocks
+1 -1
View File
@@ -109,7 +109,7 @@ Streaming is available in both sync and async modes:
!!! tip
For full details, see the [streaming guide](./streaming.md).
For full details, see the [streaming guide](../how-tos/streaming.md).
## Max iterations
-223
View File
@@ -1,223 +0,0 @@
---
search:
boost: 2
tags:
- agent
hide:
- tags
---
# Streaming
Streaming is key to building responsive applications. There are a few types of data you’ll want to stream:
1. [**Agent progress**](#agent-progress) — get updates after each node in the agent graph is executed.
2. [**LLM tokens**](#llm-tokens) — stream tokens as they are generated by the language model.
3. [**Custom updates**](#tool-updates) — emit custom data from tools during execution (e.g., "Fetched 10/100 records")
You can stream [more than one type of data](#stream-multiple-modes) at a time.
<figure markdown="1">
![image](./assets/fast_parrot.png){: style="max-height:300px"}
<figcaption>
Waiting is for pigeons.
</figcaption>
</figure>
## Agent progress
To stream agent progress, use the [`stream()`][langgraph.graph.state.CompiledStateGraph.stream] or [`astream()`][langgraph.graph.state.CompiledStateGraph.astream] methods with [`stream_mode="updates"`](https://langchain-ai.github.io/langgraph/how-tos/streaming/#updates). This emits an event after every agent step.
For example, if you have an agent that calls a tool once, you should see the following updates:
* **LLM node**: AI message with tool call requests
* **Tool node**: Tool message with execution result
* **LLM node**: Final AI response
=== "Sync"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
for chunk in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="updates"
):
print(chunk)
print("\n")
```
=== "Async"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
async for chunk in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="updates"
):
print(chunk)
print("\n")
```
## LLM tokens
To stream tokens as they are produced by the LLM, use `stream_mode="messages"`:
=== "Sync"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
for token, metadata in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="messages"
):
print("Token", token)
print("Metadata", metadata)
print("\n")
```
=== "Async"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
async for token, metadata in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="messages"
):
print("Token", token)
print("Metadata", metadata)
print("\n")
```
## Tool updates
To stream updates from tools as they are executed, you can use [get_stream_writer][langgraph.config.get_stream_writer].
=== "Sync"
```python
# highlight-next-line
from langgraph.config import get_stream_writer
def get_weather(city: str) -> str:
"""Get weather for a given city."""
# highlight-next-line
writer = get_stream_writer()
# stream any arbitrary data
# highlight-next-line
writer(f"Looking up data for city: {city}")
return f"It's always sunny in {city}!"
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
for chunk in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="custom"
):
print(chunk)
print("\n")
```
=== "Async"
```python
# highlight-next-line
from langgraph.config import get_stream_writer
def get_weather(city: str) -> str:
"""Get weather for a given city."""
# highlight-next-line
writer = get_stream_writer()
# stream any arbitrary data
# highlight-next-line
writer(f"Looking up data for city: {city}")
return f"It's always sunny in {city}!"
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
async for chunk in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="custom"
):
print(chunk)
print("\n")
```
!!! Note
If you add `get_stream_writer` inside your tool, you won't be able to invoke the tool outside of a LangGraph execution context.
## Stream multiple modes
You can specify multiple streaming modes by passing stream mode as a list: `stream_mode=["updates", "messages", "custom"]`:
=== "Sync"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
for stream_mode, chunk in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode=["updates", "messages", "custom"]
):
print(chunk)
print("\n")
```
=== "Async"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
async for stream_mode, chunk in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode=["updates", "messages", "custom"]
):
print(chunk)
print("\n")
```
## Disable streaming
In some applications you might need to disable streaming of individual tokens for a given model. This is useful in [multi-agent](./multi-agent.md) systems to control which agents stream their output.
See the [Models](./models.md#disable-streaming) guide to learn how to disable streaming.
## Additional resources
* [Streaming in LangGraph](https://langchain-ai.github.io/langgraph/how-tos/streaming)
+1 -310
View File
@@ -1,310 +1 @@
---
search:
boost: 2
tags:
- agent
hide:
- tags
---
# Tools
[Tools](https://python.langchain.com/docs/concepts/tools/) are a way to encapsulate a function and its input schema in a way that can be passed to a chat model that supports tool calling. This allows the model to request the execution of this function with specific inputs.
You can either [define your own tools](#define-simple-tools) or use [prebuilt integrations](#prebuilt-tools) that LangChain provides.
## Define simple tools
You can pass a vanilla function to `create_react_agent` to use as a tool:
```python
from langgraph.prebuilt import create_react_agent
def multiply(a: int, b: int) -> int:
"""Multiply two numbers."""
return a * b
create_react_agent(
model="anthropic:claude-3-7-sonnet",
tools=[multiply]
)
```
`create_react_agent` automatically converts vanilla functions to [LangChain tools](https://python.langchain.com/docs/concepts/tools/#tool-interface).
## Customize tools
For more control over tool behavior, use the `@tool` decorator:
```python
# highlight-next-line
from langchain_core.tools import tool
# highlight-next-line
@tool("multiply_tool", parse_docstring=True)
def multiply(a: int, b: int) -> int:
"""Multiply two numbers.
Args:
a: First operand
b: Second operand
"""
return a * b
```
You can also define a custom input schema using Pydantic:
```python
from pydantic import BaseModel, Field
class MultiplyInputSchema(BaseModel):
"""Multiply two numbers"""
a: int = Field(description="First operand")
b: int = Field(description="Second operand")
# highlight-next-line
@tool("multiply_tool", args_schema=MultiplyInputSchema)
def multiply(a: int, b: int) -> int:
return a * b
```
For additional customization, refer to the [custom tools guide](https://python.langchain.com/docs/how_to/custom_tools/).
## Hide arguments from the model
Some tools require runtime-only arguments (e.g., user ID or session context) that should not be controllable by the model.
You can put these arguments in the `state` or `config` of the agent, and access
this information inside the tool:
```python
from langgraph.prebuilt import InjectedState
from langgraph.prebuilt.chat_agent_executor import AgentState
from langchain_core.runnables import RunnableConfig
def my_tool(
# This will be populated by an LLM
tool_arg: str,
# access information that's dynamically updated inside the agent
# highlight-next-line
state: Annotated[AgentState, InjectedState],
# access static data that is passed at agent invocation
# highlight-next-line
config: RunnableConfig,
) -> str:
"""My tool."""
do_something_with_state(state["messages"])
do_something_with_config(config)
...
```
## Disable parallel tool calling
Some model providers support executing multiple tools in parallel, but
allow users to disable this feature.
For supported providers, you can disable parallel tool calling by setting `parallel_tool_calls=False` via the `model.bind_tools()` method:
```python
from langchain.chat_models import init_chat_model
def add(a: int, b: int) -> int:
"""Add two numbers"""
return a + b
def multiply(a: int, b: int) -> int:
"""Multiply two numbers."""
return a * b
model = init_chat_model("anthropic:claude-3-5-sonnet-latest", temperature=0)
tools = [add, multiply]
agent = create_react_agent(
# disable parallel tool calls
# highlight-next-line
model=model.bind_tools(tools, parallel_tool_calls=False),
tools=tools
)
agent.invoke(
{"messages": [{"role": "user", "content": "what's 3 + 5 and 4 * 7?"}]}
)
```
## Return tool results directly
Use `return_direct=True` to return tool results immediately and stop the agent loop:
```python
from langchain_core.tools import tool
# highlight-next-line
@tool(return_direct=True)
def add(a: int, b: int) -> int:
"""Add two numbers"""
return a + b
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[add]
)
agent.invoke(
{"messages": [{"role": "user", "content": "what's 3 + 5?"}]}
)
```
## Force tool use
To force the agent to use specific tools, you can set the `tool_choice` option in `model.bind_tools()`:
```python
from langchain_core.tools import tool
# highlight-next-line
@tool(return_direct=True)
def greet(user_name: str) -> int:
"""Greet user."""
return f"Hello {user_name}!"
tools = [greet]
agent = create_react_agent(
# highlight-next-line
model=model.bind_tools(tools, tool_choice={"type": "tool", "name": "greet"}),
tools=tools
)
agent.invoke(
{"messages": [{"role": "user", "content": "Hi, I am Bob"}]}
)
```
!!! Warning "Avoid infinite loops"
Forcing tool usage without stopping conditions can create infinite loops. Use one of the following safeguards:
- Mark the tool with [`return_direct=True`](#return-tool-results-directly) to end the loop after execution.
- Set [`recursion_limit`](../concepts/low_level.md#recursion-limit) to restrict the number of execution steps.
## Handle tool errors
By default, the agent will catch all exceptions raised during tool calls and will pass those as tool messages to the LLM. To control how the errors are handled, you can use the prebuilt [`ToolNode`][langgraph.prebuilt.tool_node.ToolNode] — the node that executes tools inside `create_react_agent` — via its `handle_tool_errors` parameter:
=== "Enable error handling (default)"
```python
from langgraph.prebuilt import create_react_agent
def multiply(a: int, b: int) -> int:
"""Multiply two numbers."""
if a == 42:
raise ValueError("The ultimate error")
return a * b
# Run with error handling (default)
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[multiply]
)
agent.invoke(
{"messages": [{"role": "user", "content": "what's 42 x 7?"}]}
)
```
=== "Disable error handling"
```python
from langgraph.prebuilt import create_react_agent, ToolNode
def multiply(a: int, b: int) -> int:
"""Multiply two numbers."""
if a == 42:
raise ValueError("The ultimate error")
return a * b
# highlight-next-line
tool_node = ToolNode(
[multiply],
# highlight-next-line
handle_tool_errors=False # (1)!
)
agent_no_error_handling = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=tool_node
)
agent_no_error_handling.invoke(
{"messages": [{"role": "user", "content": "what's 42 x 7?"}]}
)
```
1. This disables error handling (enabled by default). See all available strategies in the [API reference][langgraph.prebuilt.tool_node.ToolNode].
=== "Custom error handling"
```python
from langgraph.prebuilt import create_react_agent, ToolNode
def multiply(a: int, b: int) -> int:
"""Multiply two numbers."""
if a == 42:
raise ValueError("The ultimate error")
return a * b
# highlight-next-line
tool_node = ToolNode(
[multiply],
# highlight-next-line
handle_tool_errors=(
"Can't use 42 as a first operand, you must switch operands!" # (1)!
)
)
agent_custom_error_handling = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=tool_node
)
agent_custom_error_handling.invoke(
{"messages": [{"role": "user", "content": "what's 42 x 7?"}]}
)
```
1. This provides a custom message to send to the LLM in case of an exception. See all available strategies in the [API reference][langgraph.prebuilt.tool_node.ToolNode].
See [API reference][langgraph.prebuilt.tool_node.ToolNode] for more information on different tool error handling options.
## Working with memory
LangGraph allows access to short-term and long-term memory from tools. See [Memory](./memory.md) guide for more information on:
* how to [read](./memory.md#read-short-term) from and [write](./memory.md#write-short-term) to **short-term** memory
* how to [read](./memory.md#read-long-term) from and [write](./memory.md#write-long-term) to **long-term** memory
## Prebuilt tools
You can use prebuilt tools from model providers by passing a dictionary with tool specs to the `tools` parameter of `create_react_agent`. For example, to use the `web_search_preview` tool from OpenAI:
```python
from langgraph.prebuilt import create_react_agent
agent = create_react_agent(
model="openai:gpt-4o-mini",
tools=[{"type": "web_search_preview"}]
)
response = agent.invoke(
{"messages": ["What was a positive news story from today?"]}
)
```
Additionally, LangChain supports a wide range of prebuilt tool integrations for interacting with APIs, databases, file systems, web data, and more. These tools extend the functionality of agents and enable rapid development.
You can browse the full list of available integrations in the [LangChain integrations directory](https://python.langchain.com/docs/integrations/tools/).
Some commonly used tool categories include:
- **Search**: Bing, SerpAPI, Tavily
- **Code interpreters**: Python REPL, Node.js REPL
- **Databases**: SQL, MongoDB, Redis
- **Web data**: Web scraping and browsing
- **APIs**: OpenWeatherMap, NewsAPI, and others
These integrations can be configured and added to your agents using the same `tools` parameter shown in the examples above.
delete me
-5
View File
@@ -1,5 +0,0 @@
# Runs
A run is an invocation of an [assistant](../../concepts/assistants.md). Each run may have its own input, configuration, and metadata, which may affect execution and output of the underlying graph. A run can optionally be executed on a [thread](./threads.md).
The LangGraph Platform API provides several endpoints for creating and managing runs. See the [API reference](../../cloud/reference/api/api_ref.html#tag/thread-runs/) for more details.
-138
View File
@@ -1,138 +0,0 @@
# Streaming
Streaming is critical for making LLM applications feel responsive to end users.
When creating a streaming run, the **streaming mode** determines what kinds of data are streamed back to the API client.
## Supported streaming modes
LangGraph Platform supports the following streaming modes:
| Mode | Description | LangGraph Library Method |
|----------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------|-------------------------------------------------------------------------|
| **`values`** | Stream the full graph state after each [super-step](https://langchain-ai.github.io/langgraph/concepts/low_level/#graphs). [Guide](../how-tos/streaming.md#stream-graph-state) | `.stream()` / `.astream()` with `stream_mode="values"` |
| **`updates`** | Stream only the updates to the graph state after each node. [Guide](../how-tos/streaming.md#stream-graph-state) | `.stream()` / `.astream()` with `stream_mode="updates"` |
| **`messages-tuple`** | Stream LLM tokens for any messages generated inside the graph (useful for chat apps). [Guide](../how-tos/streaming.md#messages) | `.stream()` / `.astream()` with `stream_mode="messages"` |
| **`debug`** | Stream debug information throughout graph execution. [Guide](../how-tos/streaming.md#debug) | `.stream()` / `.astream()` with `stream_mode="debug"` |
| **`custom`** | Stream custom data. [Guide](../../how-tos/streaming.md#stream-custom-data) | `.stream()` / `.astream()` with `stream_mode="custom"` |
| **`events`** | Stream all events (including the state of the graph); mainly useful when migrating large LCEL apps. [Guide](../how-tos/streaming.md#stream-events) | `.astream_events()` |
✅ You can also **combine multiple modes** at the same time. See the [how-to guide](../how-tos/streaming.md#stream-multiple-modes) for configuration details.
## Stateless runs
If you don't want to **persist the outputs** of a streaming run in the [checkpointer](../../concepts/persistence.md) DB, you can create a stateless run without creating a thread:
=== "Python"
```python
from langgraph_sdk import get_client
client = get_client(url=<DEPLOYMENT_URL>, api_key=<API_KEY>)
async for chunk in client.runs.stream(
# highlight-next-line
None, # (1)!
assistant_id,
input=inputs,
stream_mode="updates"
):
print(chunk.data)
```
1. We are passing `None` instead of a `thread_id` UUID.
=== "JavaScript"
```js
import { Client } from "@langchain/langgraph-sdk";
const client = new Client({ apiUrl: <DEPLOYMENT_URL>, apiKey: <API_KEY> });
// create a streaming run
// highlight-next-line
const streamResponse = client.runs.stream(
// highlight-next-line
null, // (1)!
assistantID,
{
input,
streamMode: "updates"
}
);
for await (const chunk of streamResponse) {
console.log(chunk.data);
}
```
1. We are passing `None` instead of a `thread_id` UUID.
=== "cURL"
```bash
curl --request POST \
--url <DEPLOYMENT_URL>/runs/stream \
--header 'Content-Type: application/json' \
--header 'x-api-key: <API_KEY>'
--data "{
\"assistant_id\": \"agent\",
\"input\": <inputs>,
\"stream_mode\": \"updates\"
}"
```
## Join and stream
LangGraph Platform allows you to join an active [background run](../how-tos/background_run.md) and stream outputs from it. To do so, you can use [LangGraph SDK's](https://langchain-ai.github.io/langgraph/cloud/reference/sdk/python_sdk_ref/) `client.runs.join_stream` method:
=== "Python"
```python
from langgraph_sdk import get_client
client = get_client(url=<DEPLOYMENT_URL>, api_key=<API_KEY>)
# highlight-next-line
async for chunk in client.runs.join_stream(
thread_id,
# highlight-next-line
run_id, # (1)!
):
print(chunk)
```
1. This is the `run_id` of an existing run you want to join.
=== "JavaScript"
```js
import { Client } from "@langchain/langgraph-sdk";
const client = new Client({ apiUrl: <DEPLOYMENT_URL>, apiKey: <API_KEY> });
// highlight-next-line
const streamResponse = client.runs.joinStream(
threadID,
// highlight-next-line
runId // (1)!
);
for await (const chunk of streamResponse) {
console.log(chunk);
}
```
1. This is the `run_id` of an existing run you want to join.
=== "cURL"
```bash
curl --request GET \
--url <DEPLOYMENT_URL>/threads/<THREAD_ID>/runs/<RUN_ID>/stream \
--header 'Content-Type: application/json' \
--header 'x-api-key: <API_KEY>'
```
!!! warning "Outputs not buffered"
When you use `.join_stream`, output is not buffered, so any output produced before joining will not be received.
## API Reference
For API usage and implementation, refer to the [API reference](../reference/api/api_ref.html#tag/thread-runs/POST/threads/{thread_id}/runs/stream).
@@ -212,7 +212,6 @@ We have now created an assistant called "Open AI Assistant" that has `model_name
Output:
```
Receiving event of type: metadata
{'run_id': '1ef6746e-5893-67b1-978a-0f1cd4060e16'}
@@ -220,7 +219,6 @@ Output:
Receiving event of type: updates
{'agent': {'messages': [{'content': 'I was created by OpenAI, a research organization focused on developing and advancing artificial intelligence technology.', 'additional_kwargs': {}, 'response_metadata': {'finish_reason': 'stop', 'model_name': 'gpt-4o-2024-05-13', 'system_fingerprint': 'fp_157b3831f5'}, 'type': 'ai', 'name': None, 'id': 'run-e1a6b25c-8416-41f2-9981-f9cfe043f414', 'example': False, 'tool_calls': [], 'invalid_tool_calls': [], 'usage_metadata': None}]}}
```
### LangGraph Platform UI
@@ -233,11 +231,9 @@ Inside your deployment, select the "Assistants" tab. For the assistant you would
To edit the assistant, use the `update` method. This will create a new version of the assistant with the provided edits. See the [Python](https://langchain-ai.github.io/langgraph/cloud/reference/sdk/python_sdk_ref/#langgraph_sdk.client.AssistantsClient.update) and [JS](https://langchain-ai.github.io/langgraph/cloud/reference/sdk/js_ts_sdk_ref/#update) SDK reference docs for more information.
!!! note "Note"
You must pass in the ENTIRE config (and metadata if you are using it). The update endpoint creates new versions completely from scratch and does not rely on previous versions.
You must pass in the ENTIRE config (and metadata if you are using it). The update endpoint creates new versions completely from scratch and does not rely on previous versions.
For example, to update your assistant's system prompt:
=== "Python"
```python
@@ -247,7 +247,5 @@ Verify that the original, interrupted run was interrupted
Output:
```
'interrupted'
```
+1 -1
View File
@@ -33,7 +33,7 @@ For more information on breakpoints see [here](../../concepts/breakpoints.md).
### Submit run
To submit the run with the specified input and run settings, click the "Submit" button. This will add a [run](../concepts/runs.md) to the existing selected [thread](../concepts/threads.md). If no thread is currently selected, a new one will be created.
To submit the run with the specified input and run settings, click the "Submit" button. This will add a [run](../concepts/runs.md) to the existing selected [thread](../../concepts/persistence.md#threads). If no thread is currently selected, a new one will be created.
To cancel the ongoing run, click the "Cancel" button.
+125 -3
View File
@@ -1,8 +1,12 @@
# Stream outputs
# Streaming API
## Streaming API
[LangGraph SDK](https://langchain-ai.github.io/langgraph/cloud/reference/sdk/python_sdk_ref/) allows you to [stream outputs](../../concepts/streaming.md) from the LangGraph API server.
[LangGraph SDK](https://langchain-ai.github.io/langgraph/cloud/reference/sdk/python_sdk_ref/) allows you to stream outputs from the LangGraph API server.
!!! note
LangGraph SDK and LangGraph Server are a part of [LangGraph Platform](../../concepts/langgraph_platform.md).
## Basic usage
Basic usage example:
@@ -833,3 +837,121 @@ To stream all events, including the state of the graph:
\"stream_mode\": \"events\"
}"
```
## Stateless runs
If you don't want to **persist the outputs** of a streaming run in the [checkpointer](../../concepts/persistence.md) DB, you can create a stateless run without creating a thread:
=== "Python"
```python
from langgraph_sdk import get_client
client = get_client(url=<DEPLOYMENT_URL>, api_key=<API_KEY>)
async for chunk in client.runs.stream(
# highlight-next-line
None, # (1)!
assistant_id,
input=inputs,
stream_mode="updates"
):
print(chunk.data)
```
1. We are passing `None` instead of a `thread_id` UUID.
=== "JavaScript"
```js
import { Client } from "@langchain/langgraph-sdk";
const client = new Client({ apiUrl: <DEPLOYMENT_URL>, apiKey: <API_KEY> });
// create a streaming run
// highlight-next-line
const streamResponse = client.runs.stream(
// highlight-next-line
null, // (1)!
assistantID,
{
input,
streamMode: "updates"
}
);
for await (const chunk of streamResponse) {
console.log(chunk.data);
}
```
1. We are passing `None` instead of a `thread_id` UUID.
=== "cURL"
```bash
curl --request POST \
--url <DEPLOYMENT_URL>/runs/stream \
--header 'Content-Type: application/json' \
--header 'x-api-key: <API_KEY>'
--data "{
\"assistant_id\": \"agent\",
\"input\": <inputs>,
\"stream_mode\": \"updates\"
}"
```
## Join and stream
LangGraph Platform allows you to join an active [background run](../how-tos/background_run.md) and stream outputs from it. To do so, you can use [LangGraph SDK's](https://langchain-ai.github.io/langgraph/cloud/reference/sdk/python_sdk_ref/) `client.runs.join_stream` method:
=== "Python"
```python
from langgraph_sdk import get_client
client = get_client(url=<DEPLOYMENT_URL>, api_key=<API_KEY>)
# highlight-next-line
async for chunk in client.runs.join_stream(
thread_id,
# highlight-next-line
run_id, # (1)!
):
print(chunk)
```
1. This is the `run_id` of an existing run you want to join.
=== "JavaScript"
```js
import { Client } from "@langchain/langgraph-sdk";
const client = new Client({ apiUrl: <DEPLOYMENT_URL>, apiKey: <API_KEY> });
// highlight-next-line
const streamResponse = client.runs.joinStream(
threadID,
// highlight-next-line
runId // (1)!
);
for await (const chunk of streamResponse) {
console.log(chunk);
}
```
1. This is the `run_id` of an existing run you want to join.
=== "cURL"
```bash
curl --request GET \
--url <DEPLOYMENT_URL>/threads/<THREAD_ID>/runs/<RUN_ID>/stream \
--header 'Content-Type: application/json' \
--header 'x-api-key: <API_KEY>'
```
!!! warning "Outputs not buffered"
When you use `.join_stream`, output is not buffered, so any output produced before joining will not be received.
## API Reference
For API usage and implementation, refer to the [API reference](../reference/api/api_ref.html#tag/thread-runs/POST/threads/{thread_id}/runs/stream).
+8 -15
View File
@@ -13,7 +13,7 @@ LangGraph Studio is accessed from the LangSmith UI, within the LangGraph Platfor
For applications that are [deployed](../../quick_start.md) on LangGraph Platform, you can access Studio as part of that deployment. To do so, navigate to the deployment in LangGraph Platform within the LangSmith UI and click the "LangGraph Studio" button.
This will load the Studio UI connected to your live deployment, allowing you to create, read, and update the [threads](../../concepts/threads.md), [assistants](../../../concepts/assistants.md), and [memory](../../../concepts//memory.md) in that deployment.
This will load the Studio UI connected to your live deployment, allowing you to create, read, and update the [threads](../../../concepts/persistence.md#threads), [assistants](../../../concepts/assistants.md), and [memory](../../../concepts//memory.md) in that deployment.
## Local development server
@@ -73,11 +73,9 @@ langgraph dev --debug-port 5678
Then attach your preferred debugger:
=== "VS Code"
Add this configuration to `launch.json`:
```json
{
Add this configuration to `launch.json`:
`json
{
"name": "Attach to LangGraph",
"type": "debugpy",
"request": "attach",
@@ -85,16 +83,11 @@ Then attach your preferred debugger:
"host": "0.0.0.0",
"port": 5678
}
}
```
}
`
Specify the port number you chose in the previous step.
=== "PyCharm"
1. Go to Run → Edit Configurations
2. Click + and select "Python Debug Server"
3. Set IDE host name: `localhost`
4. Set port: `5678` (or the port number you chose in the previous step)
5. Click "OK" and start debugging
=== "PyCharm" 1. Go to Run → Edit Configurations 2. Click + and select "Python Debug Server" 3. Set IDE host name: `localhost` 4. Set port: `5678` (or the port number you chose in the previous step) 5. Click "OK" and start debugging
## Troubleshooting
+1 -5
View File
@@ -1,10 +1,6 @@
# Manage threads
!!! info "Prerequisites"
- [Threads Overview](../concepts/threads.md)
Studio allows you to view threads from the server and edit their state.
Studio allows you to view [threads](../../concepts/persistence.md#threads) from the server and edit their state.
## View threads
+1 -5
View File
@@ -1,10 +1,6 @@
# Use threads
!!! info "Prerequisites"
- [Threads Overview](../concepts/threads.md)
In this guide, we will show how to create, view, and inspect threads.
In this guide, we will show how to create, view, and inspect [threads](../../concepts/persistence.md#threads).
## Create a thread
+68 -74
View File
@@ -8,15 +8,15 @@ Currently, the SDK does not provide built-in support for defining webhook endpoi
The following API endpoints accept a `webhook` parameter:
| Operation | HTTP Method | Endpoint |
|----------------------|-------------|-----------------------------------|
| Create Run | `POST` | `/thread/{thread_id}/runs` |
| Create Thread Cron | `POST` | `/thread/{thread_id}/runs/crons` |
| Stream Run | `POST` | `/thread/{thread_id}/runs/stream` |
| Wait Run | `POST` | `/thread/{thread_id}/runs/wait` |
| Create Cron | `POST` | `/runs/crons` |
| Stream Run Stateless | `POST` | `/runs/stream` |
| Wait Run Stateless | `POST` | `/runs/wait` |
| Operation | HTTP Method | Endpoint |
|-----------|------------|----------|
| Create Run | `POST` | `/thread/{thread_id}/runs` |
| Create Thread Cron | `POST` | `/thread/{thread_id}/runs/crons` |
| Stream Run | `POST` | `/thread/{thread_id}/runs/stream` |
| Wait Run | `POST` | `/thread/{thread_id}/runs/wait` |
| Create Cron | `POST` | `/runs/crons` |
| Stream Run Stateless | `POST` | `/runs/stream` |
| Wait Run Stateless | `POST` | `/runs/wait` |
In this guide, we’ll show how to trigger a webhook after streaming a run.
@@ -25,39 +25,36 @@ In this guide, we’ll show how to trigger a webhook after streaming a run.
Before making API calls, set up your assistant and thread.
=== "Python"
```python
from langgraph_sdk import get_client
```python
from langgraph_sdk import get_client
client = get_client(url=<DEPLOYMENT_URL>)
assistant_id = "agent"
thread = await client.threads.create()
print(thread)
```
client = get_client(url=<DEPLOYMENT_URL>)
assistant_id = "agent"
thread = await client.threads.create()
print(thread)
```
=== "JavaScript"
```js
import { Client } from "@langchain/langgraph-sdk";
```js
import { Client } from "@langchain/langgraph-sdk";
const client = new Client({ apiUrl: <DEPLOYMENT_URL> });
const assistantID = "agent";
const thread = await client.threads.create();
console.log(thread);
```
const client = new Client({ apiUrl: <DEPLOYMENT_URL> });
const assistantID = "agent";
const thread = await client.threads.create();
console.log(thread);
```
=== "CURL"
```bash
curl --request POST \
--url <DEPLOYMENT_URL>/assistants/search \
--header 'Content-Type: application/json' \
--data '{ "limit": 10, "offset": 0 }' | jq -c 'map(select(.config == null or .config == {})) | .[0]' && \
curl --request POST \
--url <DEPLOYMENT_URL>/threads \
--header 'Content-Type: application/json' \
--data '{}'
```
```bash
curl --request POST \
--url <DEPLOYMENT_URL>/assistants/search \
--header 'Content-Type: application/json' \
--data '{ "limit": 10, "offset": 0 }' | jq -c 'map(select(.config == null or .config == {})) | .[0]' && \
curl --request POST \
--url <DEPLOYMENT_URL>/threads \
--header 'Content-Type: application/json' \
--data '{}'
```
Example response:
@@ -80,51 +77,48 @@ To use a webhook, specify the `webhook` parameter in your API request. When the
For example, if your server listens for webhook events at `https://my-server.app/my-webhook-endpoint`, include this in your request:
=== "Python"
```python
input = { "messages": [{ "role": "user", "content": "Hello!" }] }
```python
input = { "messages": [{ "role": "user", "content": "Hello!" }] }
async for chunk in client.runs.stream(
thread_id=thread["thread_id"],
assistant_id=assistant_id,
input=input,
stream_mode="events",
webhook="https://my-server.app/my-webhook-endpoint"
):
pass
```
async for chunk in client.runs.stream(
thread_id=thread["thread_id"],
assistant_id=assistant_id,
input=input,
stream_mode="events",
webhook="https://my-server.app/my-webhook-endpoint"
):
pass
```
=== "JavaScript"
```js
const input = { messages: [{ role: "human", content: "Hello!" }] };
```js
const input = { messages: [{ role: "human", content: "Hello!" }] };
const streamResponse = client.runs.stream(
thread["thread_id"],
assistantID,
{
input: input,
webhook: "https://my-server.app/my-webhook-endpoint"
}
);
const streamResponse = client.runs.stream(
thread["thread_id"],
assistantID,
{
input: input,
webhook: "https://my-server.app/my-webhook-endpoint"
}
);
for await (const chunk of streamResponse) {
// Handle stream output
}
```
for await (const chunk of streamResponse) {
// Handle stream output
}
```
=== "CURL"
```bash
curl --request POST \
--url <DEPLOYMENT_URL>/threads/<THREAD_ID>/runs/stream \
--header 'Content-Type: application/json' \
--data '{
"assistant_id": <ASSISTANT_ID>,
"input": {"messages": [{"role": "user", "content": "Hello!"}]},
"webhook": "https://my-server.app/my-webhook-endpoint"
}'
```
```bash
curl --request POST \
--url <DEPLOYMENT_URL>/threads/<THREAD_ID>/runs/stream \
--header 'Content-Type: application/json' \
--data '{
"assistant_id": <ASSISTANT_ID>,
"input": {"messages": [{"role": "user", "content": "Hello!"}]},
"webhook": "https://my-server.app/my-webhook-endpoint"
}'
```
## Webhook payload
+6 -7
View File
@@ -50,10 +50,9 @@ The LangGraph CLI requires a JSON configuration file that follows this [schema](
| <span style="white-space: nowrap;">`python_version`</span> | `3.11`, `3.12`, or `3.13`. Defaults to `3.11`. |
| <span style="white-space: nowrap;">`node_version`</span> | Specify `node_version: 20` to use LangGraph.js. |
| <span style="white-space: nowrap;">`pip_config_file`</span> | Path to `pip` config file. |
| <span style="white-space: nowrap;">`pip_installer`</span> | _(Added in v0.3)_ Optional. Python package installer selector. It can be set to `"auto"`, `"pip"`, or `"uv"`. From version&nbsp;0.3 onward the default strategy is to run `uv pip`, which typically delivers faster builds while remaining a drop-in replacement. In the uncommon situation where `uv` cannot handle your dependency graph or the structure of your `pyproject.toml`, specify `"pip"` here to revert to the earlier behaviour. |
| <span style="white-space: nowrap;">`dockerfile_lines`</span> | Array of additional lines to add to Dockerfile following the import from parent image. |
| <span style="white-space: nowrap;">`checkpointer`</span> | Configuration for the checkpointer. Contains a `ttl` field which is an object with the following keys: <ul><li>`strategy`: How to handle expired checkpoints (e.g., `"delete"`).</li><li>`sweep_interval_minutes`: How often to check for expired checkpoints (integer).</li><li>`default_ttl`: Default time-to-live for checkpoints in **minutes** (integer). Defines how long checkpoints are kept before the specified strategy is applied.</li></ul> |
| <span style="white-space: nowrap;">`http`</span> | HTTP server configuration with the following fields: <ul><li>`app`: Path to custom Starlette/FastAPI app (e.g., `"./src/agent/webapp.py:app"`). See [custom routes guide](../../how-tos/http/custom_routes.md).</li><li>`disable_assistants`: Disable `/assistants` routes</li><li>`disable_threads`: Disable `/threads` routes</li><li>`disable_runs`: Disable `/runs` routes</li><li>`disable_store`: Disable `/store` routes</li><li>`disable_meta`: Disable `/ok`, `/info`, `/metrics`, and `/docs` routes</li><li>`disable_mcp`: Disable `/mcp` routes</li><li>`cors`: CORS configuration with fields for `allow_origins`, `allow_methods`, `allow_headers`, etc.</li><li>`configurable_headers`: Define which request headers to exclude or include as a run's configurable values.</li></ul> |
| <span style="white-space: nowrap;">`http`</span> | HTTP server configuration with the following fields: <ul><li>`app`: Path to custom Starlette/FastAPI app (e.g., `"./src/agent/webapp.py:app"`). See [custom routes guide](../../how-tos/http/custom_routes.md).</li><li>`disable_assistants`: Disable `/assistants` routes</li><li>`disable_threads`: Disable `/threads` routes</li><li>`disable_runs`: Disable `/runs` routes</li><li>`disable_store`: Disable `/store` routes</li><li>`disable_meta`: Disable `/ok`, `/info`, `/metrics`, and `/docs` routes</li><li>`cors`: CORS configuration with fields for `allow_origins`, `allow_methods`, `allow_headers`, etc.</li><li>`configurable_headers`: Define which request headers to exclude or include as a run's configurable values.</li></ul> |
=== "JS"
@@ -129,7 +128,7 @@ The LangGraph CLI requires a JSON configuration file that follows this [schema](
- `cohere:embed-english-v3.0`: 1024
- `cohere:embed-english-light-v3.0`: 384
- `cohere:embed-multilingual-v3.0`: 1024
- `cohere:embed-multilingual-light-v3.0`: 384
- `cohere:embed-multilingual-light-v3.0`: 384
#### Semantic search with a custom embedding function
@@ -362,8 +361,8 @@ The LangGraph CLI requires a JSON configuration file that follows this [schema](
**Options**
| Option | Default | Description |
| -------------------- | ---------------- | --------------------------------------------------------------------------------------------------------------- |
| Option | Default | Description |
| -------------------- | ---------------- | ---------------------------------------------------------------------------------------------------------------------------- |
| `--platform TEXT` | | Target platform(s) to build the Docker image for. Example: `langgraph build --platform linux/amd64,linux/arm64` |
| `-t, --tag TEXT` | | **Required**. Tag for the Docker image. Example: `langgraph build -t my-image` |
| `--pull / --no-pull` | `--pull` | Build with latest remote Docker image. Use `--no-pull` for running the LangGraph Platform API server with locally built images. |
@@ -382,8 +381,8 @@ The LangGraph CLI requires a JSON configuration file that follows this [schema](
**Options**
| Option | Default | Description |
| -------------------- | ---------------- | --------------------------------------------------------------------------------------------------------------- |
| Option | Default | Description |
| -------------------- | ---------------- | ---------------------------------------------------------------------------------------------------------------------------- |
| `--platform TEXT` | | Target platform(s) to build the Docker image for. Example: `langgraph build --platform linux/amd64,linux/arm64` |
| `-t, --tag TEXT` | | **Required**. Tag for the Docker image. Example: `langgraph build -t my-image` |
| `--no-pull` | | Use locally built images. Defaults to `false` to build with latest remote Docker image. |
+3 -2
View File
@@ -50,9 +50,10 @@ Set this environment variable to have a deployment send traces to a self-hosted
## `LANGSMITH_TRACING`
Set `LANGSMITH_TRACING` to `false` to disable tracing to LangSmith.
!!! info "Only for Self-Hosted Data Plane, Self-Hosted Control Plane, and Standalone Container"
Disabling LangSmith tracing is only available for [Self-Hosted Data Plane](../../concepts/langgraph_self_hosted_data_plane.md), [Self-Hosted Control Plane](../../concepts/langgraph_self_hosted_control_plane.md), and [Standalone Container](../../concepts/langgraph_standalone_container.md) deployments.
Defaults to `true`.
Set `LANGSMITH_TRACING` to `false` to disable tracing to LangSmith.
## `LOG_LEVEL`
+15 -13
View File
@@ -1,29 +1,31 @@
# Assistants
!!! info "Prerequisites"
**Assistants** allow you to manage configurations (like prompts, LLM selection, tools) separately from your graph's core logic, enabling rapid changes that don't alter the graph architecture. It is a way to create multiple specialized versions of the same graph architecture, each optimized for different use cases through configuration variations rather than structural changes.
- [LangGraph Server](./langgraph_server.md)
- [Configuration](./low_level.md#configuration)
When building agents, it is common to make rapid changes that _do not_ alter the graph logic. For example, simply changing prompts or the LLM selection can have significant impacts on the behavior of the agent but does not require updating your graph's architecture. Assistants offer a straightforward way to manage these configurations separately from your graph's core logic.
Imagine a general-purpose writing agent built on a common graph architecture. While the structure remains the same, different writing styles—such as blog posts and tweets—require tailored configurations to optimize performance. To support these variations, you can create multiple assistants (e.g., one for blogs and another for tweets) that share the underlying graph but differ in model selection and system prompt.
For example, imagine a general-purpose writing agent built on a common graph architecture. While the structure remains the same, different writing styles—such as blog posts and tweets—require tailored configurations to optimize performance. To support these variations, you can create multiple assistants (e.g., one for blogs and another for tweets) that share the underlying graph but differ in model selection and system prompt.
![assistant versions](img/assistants.png)
## Configuring assistants
The LangGraph Cloud API provides several endpoints for creating and managing assistants and their versions. See the [API reference](../cloud/reference/api/api_ref.html#tag/assistants) for more details.
!!! info
Assistants are a [LangGraph Platform](langgraph_platform.md) concept. They are not available in the open source LangGraph library.
## Configuration
Assistants build on the LangGraph open source concept of [configuration](low_level.md#configuration).
While configuration is available in the open source LangGraph library, assistants are only present in [LangGraph Platform](langgraph_platform.md).
This is due to the fact that assistants are tightly coupled to your deployed graph. Upon deployment, LangGraph Server will automatically create a default assistant for each graph using the graph's default configuration settings.
While configuration is available in the open source LangGraph library, assistants are only present in [LangGraph Platform](langgraph_platform.md). This is due to the fact that assistants are tightly coupled to your deployed graph. Upon deployment, LangGraph Server will automatically create a default assistant for each graph using the graph's default configuration settings.
In practice, an assistant is just an _instance_ of a graph with a specific configuration. Therefore, multiple assistants can reference the same graph but can contain different configurations (e.g. prompts, models, tools). The LangGraph Server API provides several endpoints for creating and managing assistants. See the [API reference](../cloud/reference/api/api_ref.html) and [this how-to](../cloud/how-tos/configuration_cloud.md) for more details on how to create assistants.
## Versioning assistants
## Versioning
Assistants support versioning to track changes over time.
Once you've created an assistant, subsequent edits to that assistant will create new versions. See [this how-to](../cloud/how-tos/configuration_cloud.md#create-a-new-version-for-your-assistant) for more details on how to manage assistant versions.
## Learn more
## Execution
* The LangGraph Cloud API provides several endpoints for creating and managing assistants and their versions. See the [API reference](../cloud/reference/api/api_ref.html#tag/assistants) for more details.
A **run** is an invocation of an assistant. Each run may have its own input, configuration, and metadata, which may affect execution and output of the underlying graph. A run can optionally be executed on a [thread](../../concepts/persistence.md#threads).
The LangGraph Platform API provides several endpoints for creating and managing runs. See the [API reference](../../cloud/reference/api/api_ref.html#tag/thread-runs/) for more details.
+11 -2
View File
@@ -5,7 +5,16 @@ search:
# Deployment Options
There are 4 main options for deploying with the LangGraph Platform:
## Free deployment
There are two free options for deploying LangGraph applications via the LangGraph Server:
1. [Local](../tutorials/langgraph-platform/local-server.md): Deploy for local testing and development.
1. [Standalone Container (Lite)](../concepts/langgraph_standalone_container.md): A limited version of Standalone Container for deployments unlikely to see more that 1 million node executions per year and that do not need crons and other enterprise features. Standalone Container (Lite) deployment option is free with a LangSmith API key.
## Production deployment
There are 4 main options for deploying with the [LangGraph Platform](langgraph_platform.md):
1. [Cloud SaaS](#cloud-saas)
@@ -22,7 +31,7 @@ A quick comparison:
|----------------------|----------------|----------------------------|-------------------------------|--------------------------|
| **[Control plane UI/API](../concepts/langgraph_control_plane.md)** | Yes | Yes | Yes | No |
| **CI/CD** | Managed internally by platform | Managed externally by you | Managed externally by you | Managed externally by you |
| **Data/compute residency** | LangChain’s cloud | Your cloud | Your cloud | Your cloud |
| **Data/compute residency** | LangChain's cloud | Your cloud | Your cloud | Your cloud |
| **LangSmith compatibility** | Trace to LangSmith SaaS | Trace to LangSmith SaaS | Trace to Self-Hosted LangSmith | Optional tracing |
| **[Server version compatibility](../concepts/langgraph_server.md#server-versions)** | Enterprise | Enterprise | Enterprise | Lite, Enterprise |
| **[Pricing](https://www.langchain.com/pricing-langgraph-platform)** | Plus | Enterprise | Enterprise | Developer |
+1 -6
View File
@@ -9,18 +9,13 @@ search:
## Installation
The LangGraph CLI can be installed via pip or [Homebrew](https://brew.sh/):
The LangGraph CLI can be installed via pip:
=== "pip"
```bash
pip install langgraph-cli
```
=== "Homebrew"
```bash
brew install langgraph-cli
```
## Commands
LangGraph CLI provides the following core functionality:
@@ -19,7 +19,7 @@ From the control plane UI, you can:
- Update a deployment.
- Update environment variables for a deployment.
- View build and server logs of a deployment.
- View deployment metrics such as CPU and memory usage.
- View deployment metrics like CPU and memory usage.
- Delete a deployment.
The Control Plane UI is embedded in [LangSmith](https://docs.smith.langchain.com/langgraph_cloud).
@@ -95,8 +95,6 @@ After a deployment is ready, the control plane monitors the deployment and recor
- CPU and memory usage of the deployment.
- Number of container restarts.
- Number of replicas (this will increase with [autoscaling](../concepts/langgraph_data_plane.md#autoscaling)).
- [Postgres](../concepts/langgraph_data_plane.md#postgres) CPU, memory usage, and disk usage.
These metrics are displayed as charts in the Control Plane UI.
+1 -1
View File
@@ -17,7 +17,7 @@ Develop, deploy, scale, and manage agents with **LangGraph Platform** — the pu
LangGraph Platform makes it easy to get your agent running in production — whether it’s built with LangGraph or another framework — so you can focus on your app logic, not infrastructure. Deploy with one click to get a live endpoint, and use our robust APIs and built-in task queues to handle production scale.
- **[Streaming Support](../cloud/concepts/streaming.md)**: As agents grow more sophisticated, they often benefit from streaming both token outputs and intermediate states back to the user. Without this, users are left waiting for potentially long operations with no feedback. LangGraph Server provides multiple streaming modes optimized for various application needs.
- **[Streaming Support](../cloud/how-tos/streaming.md)**: As agents grow more sophisticated, they often benefit from streaming both token outputs and intermediate states back to the user. Without this, users are left waiting for potentially long operations with no feedback. LangGraph Server provides multiple streaming modes optimized for various application needs.
- **[Background Runs](../cloud/how-tos/background_run.md)**: For agents that take longer to process (e.g., hours), maintaining an open connection can be impractical. The LangGraph Server supports launching agent runs in the background and provides both polling endpoints and webhooks to monitor run status effectively.
@@ -3,7 +3,7 @@
There are two versions of the self-hosted deployment: [Self-Hosted Data Plane](./deployment_options.md#self-hosted-data-plane) and [Self-Hosted Control Plane](./deployment_options.md#self-hosted-control-plane).
!!! info "Important"
The Self-Hosted Control Plane deployment option is currently in beta stage and requires an [Enterprise](../../concepts/plans.md) plan.
The Self-Hosted Control Plane deployment option is currently in beta stage and requires an [Enterprise](plans.md) plan.
## Requirements
@@ -8,7 +8,7 @@ search:
There are two versions of the self-hosted deployment: [Self-Hosted Data Plane](./deployment_options.md#self-hosted-data-plane) and [Self-Hosted Control Plane](./deployment_options.md#self-hosted-control-plane).
!!! info "Important"
The Self-Hosted Data Plane deployment option is currently in beta stage and requires an [Enterprise](../../concepts/plans.md) plan.
The Self-Hosted Data Plane deployment option is currently in beta stage and requires an [Enterprise](plans.md) plan.
## Requirements
+1 -1
View File
@@ -7,7 +7,7 @@ search:
**LangGraph Server** offers an API for creating and managing agent-based applications. It is built on the concept of [assistants](assistants.md), which are agents configured for specific tasks, and includes built-in [persistence](persistence.md#memory-store) and a **task queue**. This versatile API supports a wide range of agentic application use cases, from background processing to real-time interactions.
Use LangGraph Server to create and manage [assistants](assistants.md), [threads](../cloud/concepts/threads.md), [runs](../cloud/concepts/runs.md), [cron jobs](../cloud/concepts/cron_jobs.md), [webhooks](../cloud/concepts/webhooks.md), and more.
Use LangGraph Server to create and manage [assistants](assistants.md), [threads](./persistence.md#threads), [runs](../cloud/concepts/runs.md), [cron jobs](../cloud/concepts/cron_jobs.md), [webhooks](../cloud/concepts/webhooks.md), and more.
!!! tip "API reference"
+1 -2
View File
@@ -87,7 +87,6 @@ One of the most common agent types is a [tool-calling agent](../agents/overview.
```python
from langchain_core.tools import tool
@tool
def transfer_to_bob():
"""Transfer to bob."""
return Command(
@@ -415,4 +414,4 @@ There are two high-level approaches to achieve that:
An agent might need to have a different state schema from the rest of the agents. For example, a search agent might only need to keep track of queries and retrieved documents. There are two ways to achieve this in LangGraph:
- Define [subgraph](./subgraphs.md) agents with a separate state schema. If there are no shared state keys (channels) between the subgraph and the parent graph, it’s important to [add input / output transformations](../how-tos/subgraph.ipynb#different-state-schemas) so that the parent graph knows how to communicate with the subgraphs.
- Define agent node functions with a [private input state schema](../how-tos/graph-api.ipynb/#pass-private-state-between-nodes) that is distinct from the overall graph state schema. This allows passing information that is only needed for executing that particular agent.
- Define agent node functions with a [private input state schema](../how-tos/graph-api.ipynb/#pass-private-state-between-nodes) that is distinct from the overall graph state schema. This allows passing information that is only needed for executing that particular agent.
+8 -2
View File
@@ -15,15 +15,19 @@ LangGraph has a built-in persistence layer, implemented through checkpointers. W
## Threads
A thread is a unique ID or [thread identifier](#threads) assigned to each checkpoint saved by a checkpointer. When invoking graph with a checkpointer, you **must** specify a `thread_id` as part of the `configurable` portion of the config:
A thread is a unique ID or thread identifier assigned to each checkpoint saved by a checkpointer. It contains the accumulated state of a sequence of [runs](../cloud/concepts/runs.md). When a run is executed, the [state](../concepts/low_level.md#state) of the underlying graph of the assistant will be persisted to the thread.
When invoking graph with a checkpointer, you **must** specify a `thread_id` as part of the `configurable` portion of the config:
```python
{"configurable": {"thread_id": "1"}}
```
A thread's current and historical state can be retrieved. To persist state, a thread must be created prior to executing a run. The LangGraph Platform API provides several endpoints for creating and managing threads and thread state. See the [API reference](../cloud/reference/api/api_ref.html#tag/threads) for more details.
## Checkpoints
Checkpoint is a snapshot of the graph state saved at each super-step and is represented by `StateSnapshot` object with the following key properties:
The state of a thread at a particular point in time is called a checkpoint. Checkpoint is a snapshot of the graph state saved at each super-step and is represented by `StateSnapshot` object with the following key properties:
- `config`: Config associated with this checkpoint.
- `metadata`: Metadata associated with this checkpoint.
@@ -31,6 +35,8 @@ Checkpoint is a snapshot of the graph state saved at each super-step and is repr
- `next` A tuple of the node names to execute next in the graph.
- `tasks`: A tuple of `PregelTask` objects that contain information about next tasks to be executed. If the step was previously attempted, it will include error information. If a graph was interrupted [dynamically](../how-tos/human_in_the_loop/breakpoints.ipynb#dynamic-breakpoints) from within a node, tasks will contain additional data associated with interrupts.
Checkpoints are persisted and can be used to restore the state of a thread at a later time.
Let's see what checkpoints are saved when a simple graph is invoked as follows:
```python
+1 -1
View File
@@ -18,6 +18,6 @@ There are three main categories of data you can stream:
- [**Stream LLM tokens**](../how-tos/streaming.md#messages) — capture token streams from anywhere: inside nodes, subgraphs, or tools.
- [**Emit progress notifications from tools**](../how-tos/streaming.md#stream-custom-data) — send custom updates or progress signals directly from tool functions.
- [**Stream from subgraphs**](../how-tos/streaming.md#subgraphs) — include outputs from both the parent graph and any nested subgraphs.
- [**Stream from subgraphs**](../how-tos/streaming.md#stream-subgraph-outputs) — include outputs from both the parent graph and any nested subgraphs.
- [**Use any LLM**](../how-tos/streaming.md#use-with-any-llm) — stream tokens from any LLM, even if it's not a LangChain model using the `custom` streaming mode.
- [**Use multiple streaming modes**](../how-tos/streaming.md#stream-multiple-modes) — choose from `values` (full state), `updates` (state deltas), `messages` (LLM tokens + metadata), `custom` (arbitrary user data), or `debug` (detailed traces).
+40 -38
View File
@@ -1,62 +1,64 @@
# Tools
Many AI applications interact directly with humans. In these cases, it is appropriate for models to respond in natural language.
But what about cases where we want a model to also interact *directly* with systems, such as databases or an API?
These systems often have a particular input schema; for example, APIs frequently have a required payload structure. You can use [tool calling](https://platform.openai.com/docs/guides/function-calling/example-use-cases) to request model responses that match a particular schema.
Many AI applications interact with users via natural language. However, some use cases require models to interface directly with external systems—such as APIs, databases, or file systems—using structured input. In these scenarios, **tool calling** enables models to generate requests that conform to a specified input schema.
[Tools](https://python.langchain.com/docs/concepts/tools/) are a way to encapsulate a function and its input schema in a way that can be passed to a chat model that supports tool calling. This allows the model to request the execution of this function with specific inputs.
**Tools** can be passed to [chat models](https://python.langchain.com/docs/concepts/chat_models) that support [tool calling](https://python.langchain.com/docs/concepts/tool_calling) allowing the model to request the execution of a specific function with specific inputs.
You can [create custom tools](https://python.langchain.com/docs/how_to/custom_tools/) or use [prebuilt](#prebuilt-tools) tools.
[Tools](https://python.langchain.com/docs/concepts/tools/) encapsulate a callable function and its input schema. These can be passed to compatible [chat models](https://python.langchain.com/docs/concepts/chat_models), allowing the model to decide whether to invoke a tool and with what arguments.
## Tool calling
![Diagram of a tool call by a model](./img/tool_call.png)
A key principle of tool calling is that the model decides when to use a tool based on the input's relevance. The model doesn't always need to call a tool.
For example, given an input that is *irrelevant to the tool*, the model would not call the tool:
Tool calling is typically **conditional**. Based on the user input and available tools, the model may choose to issue a tool call request. This request is returned in an `AIMessage` object, which includes a `tool_calls` field that specifies the tool name and input arguments:
```python
result = llm_with_tools.invoke("Hello world!")
llm_with_tools.invoke("What is 2 multiplied by 3?")
# -> AIMessage(tool_calls=[{'name': 'multiply', 'args': {'a': 2, 'b': 3}, ...}])
```
The result would be an `AIMessage` containing the model's response in natural language (e.g., "Hello!").
However, if we pass an input *relevant to the tool*, the model should choose to call it:
If the input is unrelated to any tool, the model returns only a natural language message:
```python
result = llm_with_tools.invoke("What is 2 multiplied by 3?")
llm_with_tools.invoke("Hello world!") # -> AIMessage(content="Hello!")
```
As before, the output `result` will be an `AIMessage`.
But, if the tool was called, `result` will have a `tool_calls` attribute.
This attribute includes everything needed to execute the tool, including the tool name and input arguments:
Importantly, the model does not execute the tool—it only generates a request. A separate executor (such as a runtime or agent) is responsible for handling the tool call and returning the result.
```
result.tool_calls
{'name': 'multiply', 'args': {'a': 2, 'b': 3}, 'id': 'xxx', 'type': 'tool_call'}
```
For more details on usage, see the [how-to guide](../how-tos/tool-calling.ipynb).
## Execute tools
LangGraph offers pre-built components — [`ToolNode`][langgraph.prebuilt.tool_node.ToolNode] and [`create_react_agent`][langgraph.prebuilt.chat_agent_executor.create_react_agent] — that invoke the tools on behalf of the user.
See this [how-to guide](../how-tos/tool-calling.ipynb#use-prebuilt-toolnode) on tool calling.
See the [tool calling guide](../how-tos/tool-calling.md) for more details.
## Prebuilt tools
LangChain supports a wide range of prebuilt tool integrations for interacting with APIs, databases, file systems, web data, and more. These tools extend the functionality of agents and enable rapid development.
LangChain provides prebuilt tool integrations for common external systems including APIs, databases, file systems, and web data.
You can browse the full list of available integrations in the [LangChain integrations directory](https://python.langchain.com/docs/integrations/tools/).
Browse the [integrations directory](https://python.langchain.com/docs/integrations/tools/) for available tools.
Some commonly used tool categories include:
Common categories:
- **Search**: Bing, SerpAPI, Tavily
- **Code interpreters**: Python REPL, Node.js REPL
- **Databases**: SQL, MongoDB, Redis
- **Web data**: Web scraping and browsing
- **APIs**: OpenWeatherMap, NewsAPI, and others
* **Search**: Bing, SerpAPI, Tavily
* **Code execution**: Python REPL, Node.js REPL
* **Databases**: SQL, MongoDB, Redis
* **Web data**: Scraping and browsing
* **APIs**: OpenWeatherMap, NewsAPI, etc.
These integrations can be configured and added to your agents using the same `tools` parameter shown in the examples above.
## Custom tools
You can define custom tools using the `@tool` decorator or plain Python functions. For example:
```python
from langchain_core.tools import tool
@tool
def multiply(a: int, b: int) -> int:
"""Multiply two numbers."""
return a * b
```
See the [tool calling guide](../how-tos/tool-calling.md) for more details.
## Tool execution
While the model determines *when* to call a tool, **execution** of the tool call must be handled by a runtime component.
LangGraph provides prebuilt components for this:
* [`ToolNode`][oolNode]: Executes tools based on AI tool calls.
* [`create_react_agent`][create_react_agent]: Constructs a full agent that manages tool calling automatically.
+2 -2
View File
@@ -34,7 +34,7 @@ def read_root():
## Configure `langgraph.json`
Add the following to your `langgraph.json` configuration file. Make sure the path points to the FastAPI application instance `app` in the `webapp.py` file you created above.
Add the following to your `langgraph.json` configuration file. Make sure the path points to the `app.py` file you created above.
```json
{
@@ -71,4 +71,4 @@ You can deploy this app as-is to LangGraph Platform or to your self-hosted platf
## Next steps
Now that you've added a custom route to your deployment, you can use this same technique to further customize how your server behaves, such as defining custom [custom middleware](./custom_middleware.md) and [custom lifespan events](./custom_lifespan.md).
Now that you've added a custom route to your deployment, you can use this same technique to further customize how your server behaves, such as defining custom [custom middleware](./custom_middleware.md) and [custom lifespan events](./custom_lifespan.md).
+223 -25
View File
@@ -1,11 +1,220 @@
# Stream outputs
## Streaming API
You can [stream outputs](../concepts/streaming.md) from a LangGraph agent or workflow.
## Supported stream modes
Pass one or more of the following stream modes as a list to the [`stream()`][langgraph.graph.state.CompiledStateGraph.stream] or [`astream()`][langgraph.graph.state.CompiledStateGraph.astream] methods:
| Mode | Description |
|------|-------------|
| `values` | Streams the full value of the state after each step of the graph. |
| `updates` | Streams the updates to the state after each step of the graph. If multiple updates are made in the same step (e.g., multiple nodes are run), those updates are streamed separately. |
| `custom` | Streams custom data from inside your graph nodes. |
| `messages` | Streams 2-tuples (LLM token, metadata) from any graph nodes where an LLM is invoked. |
| `debug` | Streams as much information as possible throughout the execution of the graph.
## Stream from an agent
### Agent progress
To stream agent progress, use the [`stream()`][langgraph.graph.state.CompiledStateGraph.stream] or [`astream()`][langgraph.graph.state.CompiledStateGraph.astream] methods with `stream_mode="updates"`. This emits an event after every agent step.
For example, if you have an agent that calls a tool once, you should see the following updates:
* **LLM node**: AI message with tool call requests
* **Tool node**: Tool message with execution result
* **LLM node**: Final AI response
=== "Sync"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
for chunk in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="updates"
):
print(chunk)
print("\n")
```
=== "Async"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
async for chunk in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="updates"
):
print(chunk)
print("\n")
```
### LLM tokens
To stream tokens as they are produced by the LLM, use `stream_mode="messages"`:
=== "Sync"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
for token, metadata in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="messages"
):
print("Token", token)
print("Metadata", metadata)
print("\n")
```
=== "Async"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
# highlight-next-line
async for token, metadata in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="messages"
):
print("Token", token)
print("Metadata", metadata)
print("\n")
```
### Tool updates
To stream updates from tools as they are executed, you can use [get_stream_writer][langgraph.config.get_stream_writer].
=== "Sync"
```python
# highlight-next-line
from langgraph.config import get_stream_writer
def get_weather(city: str) -> str:
"""Get weather for a given city."""
# highlight-next-line
writer = get_stream_writer()
# stream any arbitrary data
# highlight-next-line
writer(f"Looking up data for city: {city}")
return f"It's always sunny in {city}!"
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
for chunk in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="custom"
):
print(chunk)
print("\n")
```
=== "Async"
```python
# highlight-next-line
from langgraph.config import get_stream_writer
def get_weather(city: str) -> str:
"""Get weather for a given city."""
# highlight-next-line
writer = get_stream_writer()
# stream any arbitrary data
# highlight-next-line
writer(f"Looking up data for city: {city}")
return f"It's always sunny in {city}!"
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
async for chunk in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode="custom"
):
print(chunk)
print("\n")
```
!!! Note
If you add `get_stream_writer` inside your tool, you won't be able to invoke the tool outside of a LangGraph execution context.
### Stream multiple modes
You can specify multiple streaming modes by passing stream mode as a list: `stream_mode=["updates", "messages", "custom"]`:
=== "Sync"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
for stream_mode, chunk in agent.stream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode=["updates", "messages", "custom"]
):
print(chunk)
print("\n")
```
=== "Async"
```python
agent = create_react_agent(
model="anthropic:claude-3-7-sonnet-latest",
tools=[get_weather],
)
async for stream_mode, chunk in agent.astream(
{"messages": [{"role": "user", "content": "what is the weather in sf"}]},
# highlight-next-line
stream_mode=["updates", "messages", "custom"]
):
print(chunk)
print("\n")
```
### Disable streaming
In some applications you might need to disable streaming of individual tokens for a given model. This is useful in [multi-agent](../agents/multi-agent.md) systems to control which agents stream their output.
See the [Models](../agents/models.md#disable-streaming) guide to learn how to disable streaming.
## Stream from a workflow
### Basic usage example
LangGraph graphs expose the [`.stream()`][langgraph.pregel.Pregel.stream] (sync) and [`.astream()`][langgraph.pregel.Pregel.astream] (async) methods to yield streamed outputs as iterators.
Basic usage example:
=== "Sync"
```python
@@ -61,18 +270,7 @@ Basic usage example:
```output
{'refine_topic': {'topic': 'ice cream and cats'}}
{'generate_joke': {'joke': 'This is a joke about ice cream and cats'}}
```
### Supported stream modes
| Mode | Description |
|----------------------------------|-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
| [`values`](#stream-graph-state) | Streams the full value of the state after each step of the graph. |
| [`updates`](#stream-graph-state) | Streams the updates to the state after each step of the graph. If multiple updates are made in the same step (e.g., multiple nodes are run), those updates are streamed separately. |
| [`custom`](#stream-custom-data) | Streams custom data from inside your graph nodes. |
| [`messages`](#messages) | Streams 2-tuples (LLM token, metadata) from any graph nodes where an LLM is invoked. |
| [`debug`](#debug) | Streams as much information as possible throughout the execution of the graph. |
``` |
### Stream multiple modes
@@ -94,7 +292,7 @@ The streamed outputs will be tuples of `(mode, chunk)` where `mode` is the name
print(chunk)
```
## Stream graph state
### Stream graph state
Use the stream modes `updates` and `values` to stream the state of the graph as it executes.
@@ -157,7 +355,7 @@ graph = (
```
## Subgraphs
### Stream subgraph outputs
To include outputs from [subgraphs](../concepts/subgraphs.md) in the streamed outputs, you can set `subgraphs=True` in the `.stream()` method of the parent graph. This will stream outputs from both the parent graph and any subgraphs.
@@ -233,7 +431,7 @@ for chunk in graph.stream(
**Note** that we are receiving not just the node updates, but we also the namespaces which tell us what graph (or subgraph) we are streaming from.
## Debugging {#debug}
### Debugging {#debug}
Use the `debug` streaming mode to stream as much information as possible throughout the execution of the graph. The streamed outputs include the name of the node as well as the full state.
@@ -247,7 +445,7 @@ for chunk in graph.stream(
```
## LLM tokens {#messages}
### LLM tokens {#messages}
Use the `messages` streaming mode to stream Large Language Model (LLM) outputs **token by token** from any part of your graph, including nodes, tools, subgraphs, or tasks.
@@ -307,7 +505,7 @@ for message_chunk, metadata in graph.stream( # (2)!
2. The "messages" stream mode returns an iterator of tuples `(message_chunk, metadata)` where `message_chunk` is the token streamed by the LLM and `metadata` is a dictionary with information about the graph node where the LLM was called and other information.
### Filter by LLM invocation
#### Filter by LLM invocation
You can associate `tags` with LLM invocations to filter the streamed tokens by LLM invocation.
@@ -391,7 +589,7 @@ async for msg, metadata in graph.astream( # (3)!
4. The `stream_mode` is set to "messages" to stream LLM tokens. The `metadata` contains information about the LLM invocation, including the tags.
### Filter by node
#### Filter by node
To stream tokens only from specific nodes, use `stream_mode="messages"` and filter the outputs by the `langgraph_node` field in the streamed metadata:
@@ -464,7 +662,7 @@ for msg, metadata in graph.stream( # (1)!
1. The "messages" stream mode returns a tuple of `(message_chunk, metadata)` where `message_chunk` is the token streamed by the LLM and `metadata` is a dictionary with information about the graph node where the LLM was called and other information.
2. Filter the streamed tokens by the `langgraph_node` field in the metadata to only include the tokens from the `write_poem` node.
## Stream custom data
### Stream custom data
To send **custom user-defined data** from inside a LangGraph node or tool, follow these steps:
@@ -541,7 +739,7 @@ To send **custom user-defined data** from inside a LangGraph node or tool, follo
3. Emit another custom key-value pair.
4. Set `stream_mode="custom"` to receive the custom data in the stream.
## Use with any LLM
### Use with any LLM
You can use `stream_mode="custom"` to stream data from **any LLM API** — even if that API does **not** implement the LangChain chat model interface.
@@ -701,7 +899,7 @@ for chunk in graph.stream(
```
## Disable streaming for specific chat models
### Disable streaming for specific chat models
If your application mixes models that support streaming with those that do not, you may need to explicitly disable streaming for
models that do not support it.
@@ -733,7 +931,7 @@ Set `disable_streaming=True` when initializing the model.
1. Set `disable_streaming=True` to disable streaming for the chat model.
## Async with Python < 3.11 { #async }
### Async with Python < 3.11 { #async }
In Python versions < 3.11, [asyncio tasks](https://docs.python.org/3/library/asyncio-task.html#asyncio.create_task) do not support the `context` parameter.
This limits LangGraph ability to automatically propagate context, and affects LangGraph’s streaming mechanisms in two key ways:
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
-22
View File
@@ -1,22 +0,0 @@
---
search:
boost: 2
---
# Deployment 🚀
There are two free options for deploying LangGraph applications via the LangGraph Server:
- [Local](./langgraph-platform/local-server.md): Deploy for local testing and development.
- [Standalone Container (Lite)](../concepts/langgraph_standalone_container.md): A limited version of Standalone Container for deployments unlikely to see more that 1 million node executions per year and that do not need crons and other enterprise features. Standalone Container (Lite) deployment option is free with a LangSmith API key.
## Other deployment options
Additionally, you can deploy to production with [LangGraph Platform](../concepts/langgraph_platform.md):
- [Cloud SaaS](../concepts/langgraph_cloud.md): Connect your GitHub repositories and deploy LangGraph Servers within LangChain's cloud. *We manage everything.*
- [Self-Hosted Data Plane<sup>(Beta)</sup>](../concepts/langgraph_self_hosted_data_plane.md): Create deployments from the [Control Plane UI](../concepts/langgraph_control_plane.md#control-plane-ui) and deploy LangGraph Servers to **your** cloud. *We manage the [control plane](../concepts/langgraph_control_plane.md). You manage the deployments.*
- [Self-Hosted Control Plane<sup>(Beta)</sup>](../concepts/langgraph_self_hosted_control_plane.md): Create deployments from a self-hosted [Control Plane UI](../concepts/langgraph_control_plane.md#control-plane-ui) and deploy LangGraph Servers to **your** cloud. *You manage everything.*
- [Standalone Container](../concepts/langgraph_standalone_container.md): Deploy LangGraph Server Docker images however you like.
For more information, see [Deployment options](../concepts/deployment_options.md).
+1 -1
View File
@@ -648,7 +648,7 @@ With orchestrator-worker, an orchestrator breaks down a task and delegates each
Because orchestrator-worker workflows are common, LangGraph **has the `Send` API to support this**. It lets you dynamically create worker nodes and send each one a specific input. Each worker has its own state, and all worker outputs are written to a *shared state key* that is accessible to the orchestrator graph. This gives the orchestrator access to all worker output and allows it to synthesize them into a final output. As you can see below, we iterate over a list of sections and `Send` each to a worker node. See further documentation [here](https://langchain-ai.github.io/langgraph/how-tos/map-reduce/) and [here](https://langchain-ai.github.io/langgraph/concepts/low_level/#send).
```python
from langgraph.types import Send
from langgraph.constants import Send
# Graph state
+140 -148
View File
@@ -89,163 +89,145 @@ plugins:
- "!^_"
nav:
- Guides:
- Get started:
- index.md
- Get started:
- Quickstart: agents/agents.md
- LangGraph basics:
- concepts/why-langgraph.md
- Build a basic chatbot: tutorials/get-started/1-build-basic-chatbot.md
- tutorials/get-started/2-add-tools.md
- tutorials/get-started/3-add-memory.md
- Add human-in-the-loop: tutorials/get-started/4-human-in-the-loop.md
- tutorials/get-started/5-customize-state.md
- tutorials/get-started/6-time-travel.md
- Deployment: tutorials/deployment.md
- Prebuilt agents:
- Overview: agents/overview.md
- agents/run_agents.md
- agents/streaming.md
- agents/models.md
- agents/tools.md
- agents/mcp.md
- agents/context.md
- agents/memory.md
- agents/human-in-the-loop.md
- agents/multi-agent.md
- agents/evals.md
- agents/deployment.md
- agents/ui.md
- LangGraph framework:
- Agent architectures:
- Overview: concepts/agentic_concepts.md
- Quickstarts:
- Agent: agents/agents.md
- Local server: tutorials/langgraph-platform/local-server.md
- Deployment: cloud/quick_start.md
- General concepts:
- Common patterns:
- Agent architectures: concepts/agentic_concepts.md
- Workflows & agents: tutorials/workflows.md
- Graphs:
- Overview: concepts/low_level.md
- Runtime overview: concepts/pregel.md
- Use the Graph API: how-tos/graph-api.ipynb
- Streaming:
- Overview: concepts/streaming.md
- "Stream outputs": how-tos/streaming.md
- Persistence:
- Overview: concepts/persistence.md
- concepts/durable_execution.md
- how-tos/persistence.ipynb
- Memory:
- Overview: concepts/memory.md
- Manage memory: how-tos/memory.ipynb
- Human-in-the-loop:
- Overview: concepts/human_in_the_loop.md
- how-tos/human_in_the_loop/add-human-in-the-loop.md
- Breakpoints:
- Overview: concepts/breakpoints.md
- how-tos/human_in_the_loop/breakpoints.ipynb
- Time travel:
- Overview: concepts/time-travel.md
- how-tos/human_in_the_loop/time-travel.ipynb
- Tools:
- Overview: concepts/tools.md
- how-tos/tool-calling.ipynb
- Subgraphs:
- Overview: concepts/subgraphs.md
- how-tos/subgraph.ipynb
- Multi-agent:
- Overview: concepts/multi_agent.md
- how-tos/multi_agent.ipynb
- Functional API:
- Overview: concepts/functional_api.md
- how-tos/use-functional-api.md
- LangGraph Platform:
- Overview: concepts/langgraph_platform.md
- Get started:
- Quickstart: tutorials/langgraph-platform/local-server.md
- Deployment quickstart: cloud/quick_start.md
- Components:
- Overview: concepts/langgraph_components.md
- LangGraph Server:
- Overview: concepts/langgraph_server.md
- Application structure:
- Overview: concepts/application_structure.md
- cloud/deployment/setup.md
- cloud/deployment/setup_pyproject.md
- cloud/deployment/setup_javascript.md
- cloud/deployment/custom_docker.md
- LangGraph CLI: concepts/langgraph_cli.md
- LangGraph Studio:
- Overview: concepts/langgraph_studio.md
- Quickstart: cloud/how-tos/studio/quick_start.md
- cloud/how-tos/invoke_studio.md
- cloud/how-tos/studio/manage_assistants.md
- cloud/how-tos/threads_studio.md
- cloud/how-tos/iterate_graph_studio.md
- cloud/how-tos/studio/run_evals.md
- cloud/how-tos/clone_traces_studio.md
- cloud/how-tos/datasets_studio.md
- LangGraph SDK: concepts/sdk.md
- Data management:
- Add semantic search: cloud/deployment/semantic_search.md
- Add TTLs: how-tos/ttl/configure_ttl.md
- Agent development: agents/overview.md
- Workflow orchestration:
- Graphs: concepts/low_level.md
- Subgraphs: concepts/subgraphs.md
- Runtime: concepts/pregel.md
- Functional API: concepts/functional_api.md
- Core capabilities:
- Streaming: concepts/streaming.md
- Persistence: concepts/persistence.md
- Durable execution: concepts/durable_execution.md
- Memory: concepts/memory.md
- Tools: concepts/tools.md
- Human-in-the-loop: concepts/human_in_the_loop.md
- Breakpoints: concepts/breakpoints.md
- Time travel: concepts/time-travel.md
- Multi-agent: concepts/multi_agent.md
- Platform capabilities:
- LangGraph Platform:
- Overview: concepts/langgraph_platform.md
- Components:
- Overview: concepts/langgraph_components.md
- LangGraph Server:
- Overview: concepts/langgraph_server.md
- Data plane: concepts/langgraph_data_plane.md
- Control plane: concepts/langgraph_control_plane.md
- LangGraph CLI: concepts/langgraph_cli.md
- LangGraph Studio: concepts/langgraph_studio.md
- LangGraph SDK: concepts/sdk.md
- Plans & pricing: concepts/plans.md
- Application structure: concepts/application_structure.md
- Scalability & resilience: concepts/scalability_and_resilience.md
- Authentication & access control: concepts/auth.md
- Assistants: concepts/assistants.md
- Double-texting: concepts/double_texting.md
- Webhooks: cloud/concepts/webhooks.md
- Cron jobs: cloud/concepts/cron_jobs.md
- Deployment:
- Overview: concepts/deployment_options.md
- Deployment options:
- Cloud SaaS: concepts/langgraph_cloud.md
- Self-Hosted Data Plane: concepts/langgraph_self_hosted_data_plane.md
- Self-Hosted Control Plane: concepts/langgraph_self_hosted_control_plane.md
- Standalone Container: concepts/langgraph_standalone_container.md
- Guides:
- LangGraph APIs:
- Use the Graph API: how-tos/graph-api.ipynb
- Use the Functional API: how-tos/use-functional-api.md
- Models:
- Configure model: agents/models.md
- Streaming:
- Stream outputs: how-tos/streaming.md
- Use Server API: cloud/how-tos/streaming.md
- Context:
- Use in agent: agents/context.md
- Memory:
- Basic implementation: agents/memory.md
- Persistence: how-tos/persistence.ipynb # MERGE
- Custom implementation: how-tos/memory.ipynb
- Human-in-the-loop:
- Add to agent: agents/human-in-the-loop.md
- Add to workflow: how-tos/human_in_the_loop/add-human-in-the-loop.md
- Use Server API: cloud/how-tos/add-human-in-the-loop.md
- Time travel:
- Use Server API: cloud/how-tos/human_in_the_loop_time_travel.md
- Breakpoints:
- Set breakpoints: how-tos/human_in_the_loop/breakpoints.ipynb
- Use Server API: cloud/how-tos/human_in_the_loop_breakpoint.md
- Tools:
- Call tools: how-tos/tool-calling.md
- Subgraphs:
- Use subgraphs: how-tos/subgraph.ipynb
- Multi-agent:
- Prebuilt implementation: agents/multi-agent.md
- Custom implementation: how-tos/multi_agent.ipynb
- MCP:
- Use MCP tools: agents/mcp.md
- Server deployment via MCP: concepts/server-mcp.md
- Deployment:
- Basic deployment: agents/deployment.md
- Set up your application:
- Use requirements.txt: cloud/deployment/setup.md
- Use pyproject.toml: cloud/deployment/setup_pyproject.md
- Use JavaScript: cloud/deployment/setup_javascript.md
- Use custom Docker: cloud/deployment/custom_docker.md
- Deploy to production:
- Cloud SaaS: cloud/deployment/cloud.md
- Self-Hosted Data Plane: cloud/deployment/self_hosted_data_plane.md
- Self-Hosted Control Plane: cloud/deployment/self_hosted_control_plane.md
- Standalone Container: cloud/deployment/standalone_container.md
- Evaluation:
- Basic implementation: agents/evals.md
- Platform capabilities:
- LangGraph Studio:
- Quickstart: cloud/how-tos/studio/quick_start.md
- cloud/how-tos/invoke_studio.md
- cloud/how-tos/studio/manage_assistants.md
- cloud/how-tos/threads_studio.md
- cloud/how-tos/iterate_graph_studio.md
- cloud/how-tos/studio/run_evals.md
- cloud/how-tos/clone_traces_studio.md
- cloud/how-tos/datasets_studio.md
- Authentication & access control:
- Overview: concepts/auth.md
- how-tos/auth/custom_auth.md
- how-tos/auth/openapi_security.md
- Assistants:
- Overview: concepts/assistants.md
- cloud/how-tos/configuration_cloud.md
- Threads:
- Overview: cloud/concepts/threads.md
- cloud/how-tos/use_threads.md
- Runs:
- Overview: cloud/concepts/runs.md
- cloud/how-tos/background_run.md
- cloud/how-tos/same-thread.md
- cloud/how-tos/cron_jobs.md
- cloud/how-tos/stateless_runs.md
- cloud/how-tos/configurable_headers.md
- Streaming:
- Overview: cloud/concepts/streaming.md
- cloud/how-tos/streaming.md
- Human-in-the-loop: cloud/how-tos/add-human-in-the-loop.md
- Breakpoints: cloud/how-tos/human_in_the_loop_breakpoint.md
- Time travel: cloud/how-tos/human_in_the_loop_time_travel.md
- MCP: concepts/server-mcp.md
- Threads: cloud/how-tos/use_threads.md
- Runs:
- cloud/how-tos/background_run.md
- cloud/how-tos/same-thread.md
- cloud/how-tos/cron_jobs.md
- cloud/how-tos/stateless_runs.md
- cloud/how-tos/configurable_headers.md
- Double-texting:
- Overview: concepts/double_texting.md
- cloud/how-tos/interrupt_concurrent.md
- cloud/how-tos/rollback_concurrent.md
- cloud/how-tos/reject_concurrent.md
- cloud/how-tos/enqueue_concurrent.md
- Webhooks:
- Overview: cloud/concepts/webhooks.md
- cloud/how-tos/webhooks.md
- Cron jobs:
- Overview: cloud/concepts/cron_jobs.md
- cloud/how-tos/cron_jobs.md
- Webhooks: cloud/how-tos/webhooks.md
- Cron jobs: cloud/how-tos/cron_jobs.md
- Server customization:
- how-tos/http/custom_lifespan.md
- how-tos/http/custom_middleware.md
- how-tos/http/custom_routes.md
- Deployment:
- Overview: concepts/deployment_options.md
- Data plane: concepts/langgraph_data_plane.md
- Control plane: concepts/langgraph_control_plane.md
- Deployment options:
- Cloud SaaS:
- Overview: concepts/langgraph_cloud.md
- Deploy Cloud SaaS: cloud/deployment/cloud.md
- Self-Hosted Data Plane:
- Overview: concepts/langgraph_self_hosted_data_plane.md
- Deploy Self-Hosted Data Plane: cloud/deployment/self_hosted_data_plane.md
- Self-Hosted Control Plane:
- Overview: concepts/langgraph_self_hosted_control_plane.md
- Deploy Self-Hosted Control Plane: cloud/deployment/self_hosted_control_plane.md
- Standalone Container:
- Overview: concepts/langgraph_standalone_container.md
- Deploy Standalone Container: cloud/deployment/standalone_container.md
- Scalability & resilience: concepts/scalability_and_resilience.md
- Plans & pricing: concepts/plans.md
- Data management:
- Add semantic search: cloud/deployment/semantic_search.md
- Add TTLs: how-tos/ttl/configure_ttl.md
- Reference:
- reference/index.md
- LangGraph:
@@ -274,9 +256,20 @@ nav:
- Environment variables: cloud/reference/env_var.md
- Examples:
- agents/run_agents.md
- LangGraph basics:
- concepts/why-langgraph.md
- Build a basic chatbot: tutorials/get-started/1-build-basic-chatbot.md
- tutorials/get-started/2-add-tools.md
- tutorials/get-started/3-add-memory.md
- Add human-in-the-loop: tutorials/get-started/4-human-in-the-loop.md
- tutorials/get-started/5-customize-state.md
- tutorials/get-started/6-time-travel.md
- Template applications: concepts/template_applications.md # TODO: make tutorial
- Agentic RAG: tutorials/rag/langgraph_agentic_rag.ipynb
- Agent Supervisor: tutorials/multi_agent/agent_supervisor.ipynb
- SQL agent: tutorials/sql-agent.ipynb
- Prebuilt chat UI: agents/ui.md
- Graph runs in LangSmith: how-tos/run-id-langsmith.ipynb
- LangGraph Platform:
- Authentication:
@@ -291,11 +284,12 @@ nav:
- Integrate LangGraph into a React app: cloud/how-tos/use_stream_react.md
- Implement generative UI with LangGraph: cloud/how-tos/generative_ui_react.md
- Resources:
- concepts/faq.md
- Template applications: concepts/template_applications.md # TODO: make tutorial
- llms.txt: llms-txt-overview.md
- Additional resources:
- agents/prebuilt.md # NOTE: prebuilt.md is auto-generated by `make build-prebuilt`
- LangGraph Academy course: https://academy.langchain.com/courses/intro-to-langgraph
- Case studies: adopters.md
- concepts/faq.md
- llms.txt: llms-txt-overview.md
- Troubleshooting:
- Errors:
- troubleshooting/errors/index.md
@@ -306,9 +300,7 @@ nav:
- troubleshooting/errors/INVALID_CHAT_HISTORY.md
- troubleshooting/errors/INVALID_LICENSE.md
- LangGraph Studio: troubleshooting/studio.md
- Learn:
- LangGraph Academy course: https://academy.langchain.com/courses/intro-to-langgraph
- Case studies: adopters.md
markdown_extensions:
- abbr
-41
View File
@@ -1,41 +0,0 @@
.lang-python,
.lang-javascript {
display: none;
}
.language-switcher-global {
display: flex;
align-items: center;
padding-left: 0.5rem;
margin-right: 0.5rem;
}
/* Style the select to match the header */
.language-switcher-global select {
appearance: none;
font: inherit;
border: none;
padding: 0.25rem 0.6rem;
cursor: pointer;
outline: none;
font-weight: bolder;
}
/* Hover/focus effect */
.language-switcher-global select:hover,
.language-switcher-global select:focus {
text-decoration: underline;
}
/* Theme-specific overrides */
html[data-md-color-scheme="default"] .language-switcher-global select,
html[data-md-color-scheme="default"] .language-switcher-global option {
color: #333;
background-color: transparent;
}
html[data-md-color-scheme="slate"] .language-switcher-global select,
html[data-md-color-scheme="slate"] .language-switcher-global option {
color: #eee;
background-color: transparent;
}
-38
View File
@@ -1,38 +0,0 @@
function applyLanguageSwitching() {
const selector = document.getElementById("global-language-selector");
const langBlocks = {
python: document.querySelectorAll(".lang-python"),
javascript: document.querySelectorAll(".lang-javascript"),
};
const setLanguage = (lang) => {
for (const [key, blocks] of Object.entries(langBlocks)) {
blocks.forEach((block) => {
block.style.display = key === lang ? "block" : "none";
});
}
localStorage.setItem("preferredLang", lang);
};
const saved = localStorage.getItem("preferredLang") || "python";
if (selector) {
selector.value = saved;
selector.addEventListener("change", (e) => setLanguage(e.target.value));
}
setLanguage(saved);
}
// Run on initial load
document.addEventListener("DOMContentLoaded", applyLanguageSwitching);
// Re-run after client-side navigation (MkDocs Material)
document.addEventListener("pjax:success", applyLanguageSwitching);
// Optional: observe DOM changes (e.g., for late-loaded content)
if (window.MutationObserver) {
const observer = new MutationObserver(() => applyLanguageSwitching());
observer.observe(document.body, { childList: true, subtree: true });
}
@@ -1,6 +0,0 @@
<div class="md-header__button language-switcher-global" title="Select Language">
<select id="global-language-selector" aria-label="Select Language">
<option value="python">🐍 Python</option>
<option value="javascript">⚡️ JavaScript</option>
</select>
</div>
@@ -1,22 +0,0 @@
from _scripts.notebook_hooks import _apply_conditional_rendering
CONDITIONAL_RENDERING = """
above
:::js
js-content
:::
between
:::python
python-content
:::
below
"""
def test_conditional_rendering() -> None:
"""Test logic for conditional rendering of content."""
output = _apply_conditional_rendering(CONDITIONAL_RENDERING, "js")
assert output.strip() == "above\njs-content\n\nbetween\n\nbelow"
output = _apply_conditional_rendering(CONDITIONAL_RENDERING, "python")
assert output.strip() == "above\n\nbetween\npython-content\n\nbelow"
Generated
+5 -3
View File
@@ -2590,7 +2590,7 @@ wheels = [
[[package]]
name = "langgraph"
version = "0.4.7"
version = "0.5.0rc1"
source = { editable = "../libs/langgraph" }
dependencies = [
{ name = "langchain-core" },
@@ -2641,7 +2641,7 @@ dev = [
[[package]]
name = "langgraph-checkpoint"
version = "2.0.26"
version = "2.1.0"
source = { editable = "../libs/checkpoint" }
dependencies = [
{ name = "langchain-core" },
@@ -2660,6 +2660,8 @@ dev = [
{ name = "dataclasses-json" },
{ name = "mypy" },
{ name = "numpy" },
{ name = "pandas" },
{ name = "pandas-stubs", specifier = ">=2.2.2.240807" },
{ name = "pytest" },
{ name = "pytest-asyncio" },
{ name = "pytest-mock" },
@@ -2892,7 +2894,7 @@ test = [
[[package]]
name = "langgraph-prebuilt"
version = "0.2.2"
version = "0.5.0rc0"
source = { editable = "../libs/prebuilt" }
dependencies = [
{ name = "langchain-core" },
@@ -228,6 +228,7 @@ async def test_combined_metadata(saver_name: str, test_data) -> None:
checkpoint = await saver.aget_tuple(config)
assert checkpoint.metadata == {
**metadata,
"thread_id": "thread-2",
"run_id": "my_run_id",
}
@@ -210,6 +210,7 @@ def test_combined_metadata(saver_name: str, test_data) -> None:
checkpoint = saver.get_tuple(config)
assert checkpoint.metadata == {
**metadata,
"thread_id": "thread-2",
"run_id": "my_run_id",
}
+1 -1
View File
@@ -308,7 +308,7 @@ wheels = [
[[package]]
name = "langgraph-checkpoint"
version = "2.1.0"
version = "2.0.26"
source = { editable = "../checkpoint" }
dependencies = [
{ name = "langchain-core" },
+10 -2
View File
@@ -71,6 +71,7 @@ class TestAsyncSqliteSaver:
checkpoint = await saver.aget_tuple(config)
assert checkpoint is not None and checkpoint.metadata == {
**self.metadata_2,
"thread_id": "thread-2",
"run_id": "my_run_id",
}
@@ -91,11 +92,18 @@ class TestAsyncSqliteSaver:
search_results_1 = [c async for c in saver.alist(None, filter=query_1)]
assert len(search_results_1) == 1
assert search_results_1[0].metadata == self.metadata_1
assert search_results_1[0].metadata == {
"thread_id": "thread-1",
"thread_ts": "1",
**self.metadata_1,
}
search_results_2 = [c async for c in saver.alist(None, filter=query_2)]
assert len(search_results_2) == 1
assert search_results_2[0].metadata == self.metadata_2
assert search_results_2[0].metadata == {
"thread_id": "thread-2",
**self.metadata_2,
}
search_results_3 = [c async for c in saver.alist(None, filter=query_3)]
assert len(search_results_3) == 3
+10 -2
View File
@@ -72,6 +72,7 @@ class TestSqliteSaver:
checkpoint = saver.get_tuple(config)
assert checkpoint is not None and checkpoint.metadata == {
**self.metadata_2,
"thread_id": "thread-2",
"run_id": "my_run_id",
}
@@ -94,11 +95,18 @@ class TestSqliteSaver:
search_results_1 = list(saver.list(None, filter=query_1))
assert len(search_results_1) == 1
assert search_results_1[0].metadata == self.metadata_1
assert search_results_1[0].metadata == {
"thread_id": "thread-1",
"thread_ts": "1",
**self.metadata_1,
}
search_results_2 = list(saver.list(None, filter=query_2))
assert len(search_results_2) == 1
assert search_results_2[0].metadata == self.metadata_2
assert search_results_2[0].metadata == {
"thread_id": "thread-2",
**self.metadata_2,
}
search_results_3 = list(saver.list(None, filter=query_3))
assert len(search_results_3) == 3
+1 -1
View File
@@ -320,7 +320,7 @@ wheels = [
[[package]]
name = "langgraph-checkpoint"
version = "2.1.0"
version = "2.0.26"
source = { editable = "../checkpoint" }
dependencies = [
{ name = "langchain-core" },
@@ -392,10 +392,11 @@ def get_checkpoint_metadata(
for obj in (config.get("metadata"), config.get("configurable")):
if not obj:
continue
for key, v in obj.items():
for key in obj:
if key in metadata or key in EXCLUDED_METADATA_KEYS or key.startswith("__"):
continue
elif isinstance(v, str):
v = obj[key]
if isinstance(v, str):
metadata[key] = v.replace("\u0000", "")
elif isinstance(v, (int, bool, float)):
metadata[key] = v
@@ -412,16 +413,9 @@ Each Checkpointer implementation should use this mapping in put_writes.
WRITES_IDX_MAP = {ERROR: -1, SCHEDULED: -2, INTERRUPT: -3, RESUME: -4}
EXCLUDED_METADATA_KEYS = {
"thread_id",
"thread_ts",
"checkpoint_id",
"checkpoint_ns",
"checkpoint_map",
"langgraph_step",
"langgraph_node",
"langgraph_triggers",
"langgraph_path",
"langgraph_checkpoint_ns",
}
# --- below are deprecated utilities used by past versions of LangGraph ---
+19 -4
View File
@@ -75,6 +75,7 @@ class TestMemorySaver:
assert checkpoint is not None
assert checkpoint.metadata == {
**self.metadata_2,
"thread_id": "thread-2",
"run_id": "my_run_id",
}
@@ -111,11 +112,18 @@ class TestMemorySaver:
search_results_1 = list(self.memory_saver.list(None, filter=query_1))
assert len(search_results_1) == 1
assert search_results_1[0].metadata == self.metadata_1
assert search_results_1[0].metadata == {
"thread_id": "thread-1",
"thread_ts": "1",
**self.metadata_1,
}
search_results_2 = list(self.memory_saver.list(None, filter=query_2))
assert len(search_results_2) == 1
assert search_results_2[0].metadata == self.metadata_2
assert search_results_2[0].metadata == {
"thread_id": "thread-2",
**self.metadata_2,
}
search_results_3 = list(self.memory_saver.list(None, filter=query_3))
assert len(search_results_3) == 3
@@ -170,13 +178,20 @@ class TestMemorySaver:
c async for c in self.memory_saver.alist(None, filter=query_1)
]
assert len(search_results_1) == 1
assert search_results_1[0].metadata == self.metadata_1
assert search_results_1[0].metadata == {
"thread_id": "thread-1",
"thread_ts": "1",
**self.metadata_1,
}
search_results_2 = [
c async for c in self.memory_saver.alist(None, filter=query_2)
]
assert len(search_results_2) == 1
assert search_results_2[0].metadata == self.metadata_2
assert search_results_2[0].metadata == {
"thread_id": "thread-2",
**self.metadata_2,
}
search_results_3 = [
c async for c in self.memory_saver.alist(None, filter=query_3)
-5
View File
@@ -330,11 +330,6 @@ class HttpConfig(TypedDict, total=False):
disable_store: bool
"""Optional. If True, /store routes are removed, disabling direct store interactions via HTTP.
Default is False.
"""
disable_mcp: bool
"""Optional. If True, /mcp routes are removed, disabling the MCP server.
Default is False.
"""
disable_meta: bool
-4
View File
@@ -499,10 +499,6 @@
"type": "boolean",
"description": "Optional. If True, /assistants routes are removed from the server.\n\nDefault is False (meaning /assistants is enabled).\n"
},
"disable_mcp": {
"type": "boolean",
"description": "Optional. If True, /mcp routes are removed, disabling the MCP server.\n\nDefault is False.\n"
},
"disable_meta": {
"type": "boolean",
"description": "Optional. If True, all meta endpoints (/ok, /info, /metrics, /docs) are disabled.\n\nDefault is False.\n"
-4
View File
@@ -499,10 +499,6 @@
"type": "boolean",
"description": "Optional. If True, /assistants routes are removed from the server.\n\nDefault is False (meaning /assistants is enabled).\n"
},
"disable_mcp": {
"type": "boolean",
"description": "Optional. If True, /mcp routes are removed, disabling the MCP server.\n\nDefault is False.\n"
},
"disable_meta": {
"type": "boolean",
"description": "Optional. If True, all meta endpoints (/ok, /info, /metrics, /docs) are disabled.\n\nDefault is False.\n"
+2 -11
View File
@@ -148,15 +148,6 @@ class _NodeWithConfigWriterStore(Protocol[StateT_contra]):
) -> Any: ...
class _Invokable(Protocol[StateT_contra]):
def invoke(
self,
input: StateT_contra,
config: RunnableConfig | None = None,
**kwargs: Any,
) -> Any: ...
# TODO: we probably don't want to explicitly support the config / store signatures once
# we move to adding a context arg. Maybe what we do is we add support for kwargs with param spec
# this is purely for typing purposes though, so can easily change in the coming weeks.
@@ -169,7 +160,7 @@ StateNode: TypeAlias = Union[
_NodeWithConfigWriter[StateT_contra],
_NodeWithConfigStore[StateT_contra],
_NodeWithConfigWriterStore[StateT_contra],
_Invokable[StateT_contra],
Runnable[StateT_contra, Any],
]
@@ -545,7 +536,7 @@ class StateGraph(Generic[StateT, InputT, OutputT]):
if input_schema is not None:
self._add_schema(input_schema)
self.nodes[node] = StateNodeSpec(
coerce_to_runnable(action, name=node, trace=False), # type: ignore[arg-type]
coerce_to_runnable(action, name=node, trace=False),
metadata,
input=input_schema or self.state_schema,
retry_policy=retry_policy,
+16 -8
View File
@@ -1415,8 +1415,10 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
)
},
)
checkpoint_metadata = config["metadata"]
if saved:
checkpoint_config = patch_configurable(config, saved.config[CONF])
checkpoint_metadata = {**saved.metadata, **checkpoint_metadata}
channels, managed = channels_from_checkpoint(
self.channels,
checkpoint,
@@ -1481,6 +1483,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
create_checkpoint(checkpoint, None, step),
{
**checkpoint_metadata,
"source": "update",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -1503,6 +1506,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
next_checkpoint,
{
**checkpoint_metadata,
"source": "update",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -1539,11 +1543,9 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
create_checkpoint(checkpoint, channels, next_step),
{
**checkpoint_metadata,
"source": "input",
"step": next_step,
"parents": saved.metadata.get("parents", {})
if saved
else {},
},
get_new_channel_versions(
checkpoint_previous_versions,
@@ -1579,6 +1581,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
saved.parent_config or saved.config if saved else checkpoint_config,
next_checkpoint,
{
**checkpoint_metadata,
"source": "fork",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -1740,6 +1743,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
checkpoint,
{
**checkpoint_metadata,
"source": "update",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -1833,8 +1837,10 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
)
},
)
checkpoint_metadata = config["metadata"]
if saved:
checkpoint_config = patch_configurable(config, saved.config[CONF])
checkpoint_metadata = {**saved.metadata, **checkpoint_metadata}
channels, managed = channels_from_checkpoint(
self.channels,
checkpoint,
@@ -1897,6 +1903,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
create_checkpoint(checkpoint, None, step),
{
**checkpoint_metadata,
"source": "update",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -1919,6 +1926,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
next_checkpoint,
{
**checkpoint_metadata,
"source": "update",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -1955,11 +1963,9 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
create_checkpoint(checkpoint, channels, next_step),
{
**checkpoint_metadata,
"source": "input",
"step": next_step,
"parents": saved.metadata.get("parents", {})
if saved
else {},
},
get_new_channel_versions(
checkpoint_previous_versions,
@@ -1995,6 +2001,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
saved.parent_config or saved.config if saved else checkpoint_config,
next_checkpoint,
{
**checkpoint_metadata,
"source": "fork",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -2154,6 +2161,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
checkpoint_config,
checkpoint,
{
**checkpoint_metadata,
"source": "update",
"step": step + 1,
"parents": saved.metadata.get("parents", {}) if saved else {},
@@ -2412,7 +2420,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
debug=debug,
checkpoint_during=checkpoint_during
if checkpoint_during is not None
else config[CONF].get(CONFIG_KEY_CHECKPOINT_DURING, True),
else config[CONF].get(CONFIG_KEY_CHECKPOINT_DURING, False),
trigger_to_nodes=self.trigger_to_nodes,
migrate_checkpoint=self._migrate_checkpoint,
retry_policy=self.retry_policy,
@@ -2656,7 +2664,7 @@ class Pregel(PregelProtocol[StateT, InputT, OutputT], Generic[StateT, InputT, Ou
debug=debug,
checkpoint_during=checkpoint_during
if checkpoint_during is not None
else config[CONF].get(CONFIG_KEY_CHECKPOINT_DURING, True),
else config[CONF].get(CONFIG_KEY_CHECKPOINT_DURING, False),
trigger_to_nodes=self.trigger_to_nodes,
migrate_checkpoint=self._migrate_checkpoint,
retry_policy=self.retry_policy,
+12 -18
View File
@@ -30,6 +30,7 @@ from typing_extensions import ParamSpec, Self
from langgraph.cache.base import BaseCache
from langgraph.channels.base import BaseChannel
from langgraph.checkpoint.base import (
EXCLUDED_METADATA_KEYS,
WRITES_IDX_MAP,
BaseCheckpointSaver,
ChannelVersions,
@@ -313,22 +314,10 @@ class PregelLoop:
# deduplicate writes to special channels, last write wins
if all(w[0] in WRITES_IDX_MAP for w in writes):
writes = list({w[0]: w for w in writes}.values())
if task_id == NULL_TASK_ID:
# writes for the null task are accumulated
self.checkpoint_pending_writes = [
w
for w in self.checkpoint_pending_writes
if w[0] != task_id or w[1] not in WRITES_IDX_MAP
]
writes_to_save: WritesT = [
w[1:] for w in self.checkpoint_pending_writes if w[0] == task_id
] + list(writes)
else:
# remove existing writes for this task
self.checkpoint_pending_writes = [
w for w in self.checkpoint_pending_writes if w[0] != task_id
]
writes_to_save = writes
# remove existing writes for this task
self.checkpoint_pending_writes = [
w for w in self.checkpoint_pending_writes if w[0] != task_id
]
# save writes
self.checkpoint_pending_writes.extend((task_id, c, v) for c, v in writes)
if self.checkpoint_during and self.checkpointer_put_writes is not None:
@@ -349,7 +338,7 @@ class PregelLoop:
self.submit(
self.checkpointer_put_writes,
config,
writes_to_save,
writes,
task_id,
task_path_str(task.path) if task else "",
)
@@ -357,7 +346,7 @@ class PregelLoop:
self.submit(
self.checkpointer_put_writes,
config,
writes_to_save,
writes,
task_id,
)
# output writes
@@ -732,6 +721,11 @@ class PregelLoop:
)
# bail if no checkpointer
if do_checkpoint and self._checkpointer_put_after_previous is not None:
for k, v in self.config["metadata"].items():
if k in EXCLUDED_METADATA_KEYS:
continue
metadata.setdefault(k, v) # type: ignore
self.prev_checkpoint_config = (
self.checkpoint_config
if CONFIG_KEY_CHECKPOINT_ID in self.checkpoint_config[CONF]
+1 -2
View File
@@ -104,8 +104,7 @@ class FuturesDict(Generic[F, E], dict[F, Optional[PregelExecutableTask]]):
fut: F,
) -> None:
try:
if cb := self.callback():
cb(task, _exception(fut))
self.callback()(task, _exception(fut)) # type: ignore[misc]
finally:
with self.lock:
self.done.add(fut)
+9 -3
View File
@@ -1,5 +1,7 @@
from __future__ import annotations
from typing import Union
from typing_extensions import TypeVar
from langgraph._typing import StateLike
@@ -17,8 +19,12 @@ InputT = TypeVar("InputT", bound=StateLike, default=StateT)
Defaults to `StateT`.
"""
OutputT = TypeVar("OutputT", bound=StateLike, default=StateT)
"""Type variable used to represent the output of a state graph.
ResolvedInputT = TypeVar("ResolvedInputT", bound=StateLike)
"""Type variable used to represent the resolved input to a state graph.
Defaults to `StateT`.
No default.
"""
OutputT = TypeVar("OutputT", bound=Union[StateLike, None], default=StateT)
"""Type variable used to represent the output of a state graph."""
@@ -43,6 +43,7 @@ def get_expected_history(*, exc_task_results: int = 0) -> list[StateSnapshot]:
"source": "loop",
"step": 4,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config={
@@ -72,6 +73,7 @@ def get_expected_history(*, exc_task_results: int = 0) -> list[StateSnapshot]:
"source": "loop",
"step": 3,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config={
@@ -129,6 +131,7 @@ def get_expected_history(*, exc_task_results: int = 0) -> list[StateSnapshot]:
"source": "loop",
"step": 2,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config={
@@ -165,6 +168,7 @@ def get_expected_history(*, exc_task_results: int = 0) -> list[StateSnapshot]:
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config={
@@ -214,6 +218,7 @@ def get_expected_history(*, exc_task_results: int = 0) -> list[StateSnapshot]:
"source": "loop",
"step": 0,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config={
@@ -252,6 +257,7 @@ def get_expected_history(*, exc_task_results: int = 0) -> list[StateSnapshot]:
"source": "input",
"step": -1,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -337,6 +343,7 @@ SAVED_CHECKPOINTS = {
"source": "loop",
"step": 4,
"parents": {},
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -397,6 +404,7 @@ SAVED_CHECKPOINTS = {
"source": "loop",
"step": 3,
"parents": {},
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -472,6 +480,7 @@ SAVED_CHECKPOINTS = {
"source": "loop",
"step": 2,
"parents": {},
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -523,6 +532,7 @@ SAVED_CHECKPOINTS = {
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -577,6 +587,7 @@ SAVED_CHECKPOINTS = {
"source": "loop",
"step": 0,
"parents": {},
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -625,6 +636,7 @@ SAVED_CHECKPOINTS = {
"source": "input",
"step": -1,
"parents": {},
"thread_id": "1",
},
parent_config=None,
pending_writes=[
@@ -708,6 +720,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 4,
"parents": {},
},
@@ -769,6 +782,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 3,
"parents": {},
},
@@ -847,6 +861,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 2,
"parents": {},
},
@@ -902,6 +917,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 1,
"parents": {},
},
@@ -961,6 +977,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 0,
"parents": {},
},
@@ -1009,6 +1026,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "input",
"thread_id": "1",
"step": -1,
"parents": {},
},
@@ -1094,6 +1112,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 4,
"parents": {},
},
@@ -1155,6 +1174,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 3,
"parents": {},
},
@@ -1233,6 +1253,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 2,
"parents": {},
},
@@ -1288,6 +1309,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 1,
"parents": {},
},
@@ -1347,6 +1369,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 0,
"parents": {},
},
@@ -1395,6 +1418,7 @@ SAVED_CHECKPOINTS = {
},
metadata={
"source": "input",
"thread_id": "1",
"step": -1,
"parents": {},
},
+179 -81
View File
@@ -126,6 +126,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 6,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[1].config,
@@ -146,6 +147,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 5,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[2].config,
@@ -166,6 +168,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "input",
"step": 4,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[3].config,
@@ -186,6 +189,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[4].config,
@@ -206,6 +210,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "input",
"step": 2,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[5].config,
@@ -226,6 +231,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[6].config,
@@ -246,6 +252,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[7].config,
@@ -266,6 +273,7 @@ def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "input",
"step": -1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -334,6 +342,7 @@ def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 5,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[1].config,
@@ -354,6 +363,7 @@ def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[2].config,
@@ -374,6 +384,7 @@ def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[3].config,
@@ -394,6 +405,7 @@ def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 2,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[4].config,
@@ -414,6 +426,7 @@ def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[5].config,
@@ -434,6 +447,7 @@ def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[6].config,
@@ -454,6 +468,7 @@ def test_fork_always_re_runs_nodes(
"parents": {},
"source": "input",
"step": -1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -664,10 +679,7 @@ def test_conditional_state_graph(
config = {"configurable": {"thread_id": "1"}}
assert [
c
for c in app_w_interrupt.stream(
{"input": "what is weather in sf"}, config, checkpoint_during=False
)
c for c in app_w_interrupt.stream({"input": "what is weather in sf"}, config)
] == [
{
"agent": {
@@ -702,6 +714,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
parent_config=None,
interrupts=(),
@@ -741,6 +754,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -816,6 +830,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -834,10 +849,7 @@ def test_conditional_state_graph(
llm.i = 0 # reset the llm
assert [
c
for c in app_w_interrupt.stream(
{"input": "what is weather in sf"}, config, checkpoint_during=False
)
c for c in app_w_interrupt.stream({"input": "what is weather in sf"}, config)
] == [
{
"agent": {
@@ -870,6 +882,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "2",
},
parent_config=None,
interrupts=(),
@@ -909,6 +922,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -984,6 +998,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -1001,10 +1016,7 @@ def test_conditional_state_graph(
llm.i = 0 # reset the llm
assert [
c
for c in app_w_interrupt.stream(
{"input": "what is weather in sf"}, config, checkpoint_during=False
)
c for c in app_w_interrupt.stream({"input": "what is weather in sf"}, config)
] == [
{"__interrupt__": ()},
]
@@ -1027,6 +1039,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "3",
},
parent_config=None,
interrupts=(),
@@ -1064,6 +1077,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "3",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -1119,6 +1133,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "loop",
"step": 2,
"thread_id": "3",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -1148,10 +1163,7 @@ def test_conditional_state_graph(
llm.i = 0 # reset the llm
assert [
c
for c in app_w_interrupt.stream(
{"input": "what is weather in sf"}, config, checkpoint_during=False
)
c for c in app_w_interrupt.stream({"input": "what is weather in sf"}, config)
] == [
{
"agent": {
@@ -1184,6 +1196,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "4",
},
parent_config=None,
interrupts=(),
@@ -1237,6 +1250,7 @@ def test_conditional_state_graph(
"parents": {},
"source": "loop",
"step": 2,
"thread_id": "4",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -1849,9 +1863,7 @@ def test_state_graph_packets(
assert [
c
for c in app_w_interrupt.stream(
{"messages": HumanMessage(content="what is weather in sf")},
config,
checkpoint_during=False,
{"messages": HumanMessage(content="what is weather in sf")}, config
)
] == [
{
@@ -1903,6 +1915,7 @@ def test_state_graph_packets(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
parent_config=None,
interrupts=(),
@@ -1947,6 +1960,7 @@ def test_state_graph_packets(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "1",
},
parent_config=(
[*app_w_interrupt.checkpointer.list(config, limit=2)][-1].config
@@ -2042,6 +2056,7 @@ def test_state_graph_packets(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
parent_config=(
[*app_w_interrupt.checkpointer.list(config, limit=2)][-1].config
@@ -2095,6 +2110,7 @@ def test_state_graph_packets(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "1",
},
parent_config=(
[*app_w_interrupt.checkpointer.list(config, limit=2)][-1].config
@@ -2114,9 +2130,7 @@ def test_state_graph_packets(
assert [
c
for c in app_w_interrupt.stream(
{"messages": HumanMessage(content="what is weather in sf")},
config,
checkpoint_during=False,
{"messages": HumanMessage(content="what is weather in sf")}, config
)
] == [
{
@@ -2168,6 +2182,7 @@ def test_state_graph_packets(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "2",
},
parent_config=None,
interrupts=(),
@@ -2206,6 +2221,7 @@ def test_state_graph_packets(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "2",
},
parent_config=(
[*app_w_interrupt.checkpointer.list(config, limit=2)][-1].config
@@ -2301,6 +2317,7 @@ def test_state_graph_packets(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "2",
},
parent_config=(
[*app_w_interrupt.checkpointer.list(config, limit=2)][-1].config
@@ -2354,6 +2371,7 @@ def test_state_graph_packets(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "2",
},
parent_config=(
[*app_w_interrupt.checkpointer.list(config, limit=2)][-1].config
@@ -2582,10 +2600,7 @@ def test_message_graph(
config = {"configurable": {"thread_id": "1"}}
assert [
c
for c in app_w_interrupt.stream(
("human", "what is weather in sf"), config, checkpoint_during=False
)
c for c in app_w_interrupt.stream(("human", "what is weather in sf"), config)
] == [
{
"agent": AIMessage(
@@ -2632,6 +2647,7 @@ def test_message_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
parent_config=None,
interrupts=(),
@@ -2666,6 +2682,7 @@ def test_message_graph(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -2744,6 +2761,7 @@ def test_message_graph(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -2792,6 +2810,7 @@ def test_message_graph(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -2806,12 +2825,7 @@ def test_message_graph(
config = {"configurable": {"thread_id": "2"}}
model.i = 0 # reset the llm
assert [
c
for c in app_w_interrupt.stream(
"what is weather in sf", config, checkpoint_during=False
)
] == [
assert [c for c in app_w_interrupt.stream("what is weather in sf", config)] == [
{
"agent": AIMessage(
content="",
@@ -2857,6 +2871,7 @@ def test_message_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "2",
},
parent_config=None,
interrupts=(),
@@ -2897,6 +2912,7 @@ def test_message_graph(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -2975,6 +2991,7 @@ def test_message_graph(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3024,6 +3041,7 @@ def test_message_graph(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3073,6 +3091,7 @@ def test_message_graph(
"parents": {},
"source": "update",
"step": 6,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3304,10 +3323,7 @@ def test_root_graph(
config = {"configurable": {"thread_id": "1"}}
assert [
c
for c in app_w_interrupt.stream(
("human", "what is weather in sf"), config, checkpoint_during=False
)
c for c in app_w_interrupt.stream(("human", "what is weather in sf"), config)
] == [
{
"agent": AIMessage(
@@ -3354,6 +3370,7 @@ def test_root_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
parent_config=None,
interrupts=(),
@@ -3388,6 +3405,7 @@ def test_root_graph(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3467,6 +3485,7 @@ def test_root_graph(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3516,6 +3535,7 @@ def test_root_graph(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "1",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3530,12 +3550,7 @@ def test_root_graph(
config = {"configurable": {"thread_id": "2"}}
model.i = 0 # reset the llm
assert [
c
for c in app_w_interrupt.stream(
"what is weather in sf", config, checkpoint_during=False
)
] == [
assert [c for c in app_w_interrupt.stream("what is weather in sf", config)] == [
{
"agent": AIMessage(
content="",
@@ -3581,6 +3596,7 @@ def test_root_graph(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "2",
},
parent_config=None,
interrupts=(),
@@ -3621,6 +3637,7 @@ def test_root_graph(
"parents": {},
"source": "update",
"step": 2,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3700,6 +3717,7 @@ def test_root_graph(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3748,6 +3766,7 @@ def test_root_graph(
"parents": {},
"source": "update",
"step": 5,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3797,6 +3816,7 @@ def test_root_graph(
"parents": {},
"source": "update",
"step": 6,
"thread_id": "2",
},
parent_config=(
list(app_w_interrupt.checkpointer.list(config, limit=2))[-1].config
@@ -3877,6 +3897,7 @@ def test_root_graph(
"parents": {},
"source": "update",
"step": 6,
"thread_id": "2",
},
parent_config=(list(new_app.checkpointer.list(config, limit=2))[-1].config),
interrupts=(),
@@ -4219,9 +4240,7 @@ def test_dynamic_interrupt(sync_checkpointer: BaseCheckpointSaver) -> None:
# flow: interrupt -> clear tasks
thread1 = {"configurable": {"thread_id": "1"}}
# stop when about to enter node
assert tool_two.invoke(
{"my_key": "value ⛰️", "market": "DE"}, thread1, checkpoint_during=False
) == {
assert tool_two.invoke({"my_key": "value ⛰️", "market": "DE"}, thread1) == {
"my_key": "value ⛰️",
"market": "DE",
"__interrupt__": [
@@ -4234,6 +4253,7 @@ def test_dynamic_interrupt(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
]
@@ -4266,6 +4286,7 @@ def test_dynamic_interrupt(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
parent_config=None,
interrupts=(
@@ -4295,6 +4316,7 @@ def test_dynamic_interrupt(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "update",
"step": 1,
"thread_id": "1",
},
parent_config=(list(tool_two.checkpointer.list(thread1, limit=2))[-1].config),
interrupts=(),
@@ -4386,9 +4408,7 @@ def test_copy_checkpoint(sync_checkpointer: BaseCheckpointSaver) -> None:
# flow: interrupt -> clear tasks
thread1 = {"configurable": {"thread_id": "1"}}
# stop when about to enter node
assert tool_two.invoke(
{"my_key": "value ⛰️", "market": "DE"}, thread1, checkpoint_during=False
) == {
assert tool_two.invoke({"my_key": "value ⛰️", "market": "DE"}, thread1) == {
"my_key": "value ⛰️ one",
"market": "DE",
"__interrupt__": [
@@ -4400,6 +4420,7 @@ def test_copy_checkpoint(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
]
@@ -4438,6 +4459,7 @@ def test_copy_checkpoint(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
parent_config=None,
interrupts=(
@@ -4483,6 +4505,7 @@ def test_copy_checkpoint(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "fork",
"step": 1,
"thread_id": "1",
},
parent_config=([*tool_two.checkpointer.list(thread1, limit=2)][-1].config),
interrupts=(),
@@ -4574,9 +4597,7 @@ def test_dynamic_interrupt_subgraph(sync_checkpointer: BaseCheckpointSaver) -> N
# flow: interrupt -> clear tasks
thread1 = {"configurable": {"thread_id": "1"}}
# stop when about to enter node
assert tool_two.invoke(
{"my_key": "value ⛰️", "market": "DE"}, thread1, checkpoint_during=False
) == {
assert tool_two.invoke({"my_key": "value ⛰️", "market": "DE"}, thread1) == {
"my_key": "value ⛰️",
"market": "DE",
"__interrupt__": [
@@ -4598,6 +4619,7 @@ def test_dynamic_interrupt_subgraph(sync_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
]
@@ -4636,6 +4658,7 @@ def test_dynamic_interrupt_subgraph(sync_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
parent_config=None,
interrupts=(
@@ -4665,6 +4688,7 @@ def test_dynamic_interrupt_subgraph(sync_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "update",
"step": 1,
"thread_id": "1",
},
parent_config=(
list(
@@ -4795,6 +4819,7 @@ def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 4,
"parents": {},
},
@@ -4830,6 +4855,7 @@ def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 3,
"parents": {},
},
@@ -4872,6 +4898,7 @@ def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 2,
"parents": {},
},
@@ -4926,6 +4953,7 @@ def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 1,
"parents": {},
},
@@ -4980,6 +5008,7 @@ def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 0,
"parents": {},
},
@@ -5016,6 +5045,7 @@ def test_send_dedupe_on_resume(
},
metadata={
"source": "input",
"thread_id": "1",
"step": -1,
"parents": {},
},
@@ -5093,7 +5123,7 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
app = graph.compile(checkpointer=sync_checkpointer)
config = {"configurable": {"thread_id": "1"}}
app.invoke({"my_key": "my value"}, config, checkpoint_during=False)
app.invoke({"my_key": "my value"}, config, debug=True)
# test state w/ nested subgraph state (right after interrupt)
# first get_state without subgraph state
expected = StateSnapshot(
@@ -5118,6 +5148,7 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -5162,6 +5193,12 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
},
"source": "loop",
"step": 1,
"thread_id": "1",
"langgraph_node": "inner",
"langgraph_path": [PULL, "inner"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:inner"],
"langgraph_checkpoint_ns": AnyStr("inner:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -5181,6 +5218,7 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -5207,6 +5245,12 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
"source": "loop",
"step": 1,
"parents": {"": AnyStr()},
"thread_id": "1",
"langgraph_node": "inner",
"langgraph_path": [PULL, "inner"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:inner"],
"langgraph_checkpoint_ns": AnyStr("inner:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -5217,7 +5261,7 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
assert child_history == expected_child_history
# resume
app.invoke(None, config, checkpoint_during=False)
app.invoke(None, config, debug=True)
# test state w/ nested subgraph state (after resuming from interrupt)
assert app.get_state(config) == StateSnapshot(
values={"my_key": "hi my value here and there and back again"},
@@ -5234,6 +5278,7 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=(
@@ -5265,6 +5310,7 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=(
@@ -5303,6 +5349,7 @@ def test_nested_graph_state(sync_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -5370,12 +5417,7 @@ def test_doubly_nested_graph_state(
# test invoke w/ nested interrupt
config = {"configurable": {"thread_id": "1"}}
assert [
c
for c in app.stream(
{"my_key": "my value"}, config, subgraphs=True, checkpoint_during=False
)
] == [
assert [c for c in app.stream({"my_key": "my value"}, config, subgraphs=True)] == [
((), {"parent_1": {"my_key": "hi my value"}}),
(
(AnyStr("child:"), AnyStr("child_1:")),
@@ -5412,6 +5454,7 @@ def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -5448,9 +5491,15 @@ def test_doubly_nested_graph_state(
}
},
metadata={
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child",
"langgraph_path": ["__pregel_pull", "child"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:child"],
"parents": {"": AnyStr()},
"source": "loop",
"step": 0,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -5490,6 +5539,12 @@ def test_doubly_nested_graph_state(
),
"source": "loop",
"step": 1,
"thread_id": "1",
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child_1",
"langgraph_path": [PULL, AnyStr("child_1")],
"langgraph_step": 1,
"langgraph_triggers": ["branch:to:child_1"],
},
created_at=AnyStr(),
parent_config=None,
@@ -5545,6 +5600,17 @@ def test_doubly_nested_graph_state(
),
"source": "loop",
"step": 1,
"thread_id": "1",
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child_1",
"langgraph_path": [
PULL,
AnyStr("child_1"),
],
"langgraph_step": 1,
"langgraph_triggers": [
"branch:to:child_1",
],
},
created_at=AnyStr(),
parent_config=None,
@@ -5567,6 +5633,12 @@ def test_doubly_nested_graph_state(
"parents": {"": AnyStr()},
"source": "loop",
"step": 0,
"thread_id": "1",
"langgraph_node": "child",
"langgraph_path": [PULL, AnyStr("child")],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:child"],
"langgraph_checkpoint_ns": AnyStr("child:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -5586,15 +5658,14 @@ def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
interrupts=(),
)
# # resume
assert [
c for c in app.stream(None, config, subgraphs=True, checkpoint_during=False)
] == [
assert [c for c in app.stream(None, config, subgraphs=True)] == [
(
(AnyStr("child:"), AnyStr("child_1:")),
{"grandchild_2": {"my_key": "hi my value here and there"}},
@@ -5622,6 +5693,7 @@ def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=(
@@ -5655,6 +5727,7 @@ def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config={
@@ -5694,6 +5767,7 @@ def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -5720,6 +5794,12 @@ def test_doubly_nested_graph_state(
"source": "loop",
"step": 0,
"parents": {"": AnyStr()},
"thread_id": "1",
"langgraph_node": "child",
"langgraph_path": [PULL, AnyStr("child")],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:child"],
"langgraph_checkpoint_ns": AnyStr("child:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -5769,6 +5849,15 @@ def test_doubly_nested_graph_state(
AnyStr("child:"): AnyStr(),
}
),
"thread_id": "1",
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child_1",
"langgraph_path": [
PULL,
AnyStr("child_1"),
],
"langgraph_step": 1,
"langgraph_triggers": ["branch:to:child_1"],
},
created_at=AnyStr(),
parent_config=None,
@@ -5951,9 +6040,7 @@ def test_send_react_interrupt(
foo_called = 0
graph = builder.compile(checkpointer=sync_checkpointer, interrupt_before=["foo"])
thread1 = {"configurable": {"thread_id": "2"}}
assert graph.invoke(
{"messages": [HumanMessage("hello")]}, thread1, checkpoint_during=False
) == {
assert graph.invoke({"messages": [HumanMessage("hello")]}, thread1) == {
"messages": [
_AnyIdHumanMessage(content="hello"),
_AnyIdAIMessage(
@@ -6002,6 +6089,7 @@ def test_send_react_interrupt(
"step": 1,
"source": "loop",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=None,
@@ -6047,6 +6135,7 @@ def test_send_react_interrupt(
"step": 2,
"source": "update",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=(
@@ -6075,9 +6164,7 @@ def test_send_react_interrupt(
foo_called = 0
graph = builder.compile(checkpointer=sync_checkpointer, interrupt_before=["foo"])
thread1 = {"configurable": {"thread_id": "3"}}
assert graph.invoke(
{"messages": [HumanMessage("hello")]}, thread1, checkpoint_during=False
) == {
assert graph.invoke({"messages": [HumanMessage("hello")]}, thread1) == {
"messages": [
_AnyIdHumanMessage(content="hello"),
_AnyIdAIMessage(
@@ -6126,6 +6213,7 @@ def test_send_react_interrupt(
"step": 1,
"source": "loop",
"parents": {},
"thread_id": "3",
},
created_at=AnyStr(),
parent_config=None,
@@ -6192,6 +6280,7 @@ def test_send_react_interrupt(
"step": 2,
"source": "update",
"parents": {},
"thread_id": "3",
},
created_at=AnyStr(),
parent_config=(
@@ -6341,9 +6430,7 @@ def test_send_react_interrupt_control(
foo_called = 0
graph = builder.compile(checkpointer=sync_checkpointer, interrupt_before=["foo"])
thread1 = {"configurable": {"thread_id": "2"}}
assert graph.invoke(
{"messages": [HumanMessage("hello")]}, thread1, checkpoint_during=False
) == {
assert graph.invoke({"messages": [HumanMessage("hello")]}, thread1) == {
"messages": [
_AnyIdHumanMessage(content="hello"),
_AnyIdAIMessage(
@@ -6392,6 +6479,7 @@ def test_send_react_interrupt_control(
"step": 1,
"source": "loop",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=None,
@@ -6437,6 +6525,7 @@ def test_send_react_interrupt_control(
"step": 2,
"source": "update",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=(
@@ -6597,11 +6686,7 @@ def test_weather_subgraph(
assert [
c
for c in graph.stream(
inputs,
config=config,
stream_mode="updates",
subgraphs=True,
checkpoint_during=False,
inputs, config=config, stream_mode="updates", subgraphs=True
)
] == [
((), {"router_node": {"route": "weather"}}),
@@ -6628,6 +6713,7 @@ def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -6684,11 +6770,7 @@ def test_weather_subgraph(
assert [
c
for c in graph.stream(
inputs,
config=config,
stream_mode="updates",
subgraphs=True,
checkpoint_during=False,
inputs, config=config, stream_mode="updates", subgraphs=True
)
] == [
((), {"router_node": {"route": "weather"}}),
@@ -6713,6 +6795,7 @@ def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "14",
},
created_at=AnyStr(),
parent_config=None,
@@ -6746,6 +6829,12 @@ def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {"": AnyStr()},
"thread_id": "14",
"langgraph_node": "weather_graph",
"langgraph_path": [PULL, "weather_graph"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:weather_graph"],
"langgraph_checkpoint_ns": AnyStr("weather_graph:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -6785,6 +6874,7 @@ def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "14",
},
created_at=AnyStr(),
parent_config=None,
@@ -6819,6 +6909,14 @@ def test_weather_subgraph(
"step": 2,
"source": "update",
"parents": {"": AnyStr()},
"thread_id": "14",
"checkpoint_id": AnyStr(),
"checkpoint_ns": AnyStr("weather_graph:"),
"langgraph_node": "weather_graph",
"langgraph_path": [PULL, "weather_graph"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:weather_graph"],
"langgraph_checkpoint_ns": AnyStr("weather_graph:"),
},
created_at=AnyStr(),
parent_config=(
+133 -33
View File
@@ -119,6 +119,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 6,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[1].config,
@@ -139,6 +140,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 5,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[2].config,
@@ -159,6 +161,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "input",
"step": 4,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[3].config,
@@ -179,6 +182,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[4].config,
@@ -199,6 +203,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "input",
"step": 2,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[5].config,
@@ -219,6 +224,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[6].config,
@@ -239,6 +245,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[7].config,
@@ -259,6 +266,7 @@ async def test_invoke_two_processes_in_out_interrupt(
"parents": {},
"source": "input",
"step": -1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -334,6 +342,7 @@ async def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 5,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[1].config,
@@ -354,6 +363,7 @@ async def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[2].config,
@@ -374,6 +384,7 @@ async def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[3].config,
@@ -394,6 +405,7 @@ async def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 2,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[4].config,
@@ -414,6 +426,7 @@ async def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[5].config,
@@ -434,6 +447,7 @@ async def test_fork_always_re_runs_nodes(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=history[6].config,
@@ -454,6 +468,7 @@ async def test_fork_always_re_runs_nodes(
"parents": {},
"source": "input",
"step": -1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -683,7 +698,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
assert [
c
async for c in app_w_interrupt.astream(
{"input": "what is weather in sf"}, config, checkpoint_during=False
{"input": "what is weather in sf"}, config
)
] == [
{
@@ -721,6 +736,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
parent_config=None,
interrupts=(),
@@ -760,6 +776,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
"parents": {},
"source": "update",
"step": 2,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -837,6 +854,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
"parents": {},
"source": "update",
"step": 5,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -858,7 +876,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
assert [
c
async for c in app_w_interrupt.astream(
{"input": "what is weather in sf"}, config, checkpoint_during=False
{"input": "what is weather in sf"}, config
)
] == [
{
@@ -894,6 +912,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "2",
},
parent_config=None,
interrupts=(),
@@ -933,6 +952,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
"parents": {},
"source": "update",
"step": 2,
"thread_id": "2",
},
parent_config=[
c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)
@@ -1008,6 +1028,7 @@ async def test_conditional_graph_state(async_checkpointer: BaseCheckpointSaver)
"parents": {},
"source": "update",
"step": 5,
"thread_id": "2",
},
parent_config=[
c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)
@@ -1574,9 +1595,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
assert [
c
async for c in app_w_interrupt.astream(
{"messages": HumanMessage(content="what is weather in sf")},
config,
checkpoint_during=False,
{"messages": HumanMessage(content="what is weather in sf")}, config
)
] == [
{
@@ -1628,6 +1647,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
parent_config=None,
interrupts=(),
@@ -1665,6 +1685,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "update",
"step": 2,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -1757,6 +1778,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -1804,6 +1826,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "update",
"step": 5,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -1825,9 +1848,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
assert [
c
async for c in app_w_interrupt.astream(
{"messages": HumanMessage(content="what is weather in sf")},
config,
checkpoint_during=False,
{"messages": HumanMessage(content="what is weather in sf")}, config
)
] == [
{
@@ -1873,6 +1894,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "2",
},
parent_config=None,
interrupts=(),
@@ -1910,6 +1932,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "update",
"step": 2,
"thread_id": "2",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -2002,6 +2025,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "2",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -2049,6 +2073,7 @@ async def test_state_graph_packets(async_checkpointer: BaseCheckpointSaver) -> N
"parents": {},
"source": "update",
"step": 5,
"thread_id": "2",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -2254,9 +2279,7 @@ async def test_message_graph(async_checkpointer: BaseCheckpointSaver) -> None:
assert [
c
async for c in app_w_interrupt.astream(
HumanMessage(content="what is weather in sf"),
config,
checkpoint_during=False,
HumanMessage(content="what is weather in sf"), config
)
] == [
{
@@ -2299,6 +2322,7 @@ async def test_message_graph(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
parent_config=None,
interrupts=(),
@@ -2334,6 +2358,7 @@ async def test_message_graph(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "update",
"step": 2,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -2409,6 +2434,7 @@ async def test_message_graph(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -2454,6 +2480,7 @@ async def test_message_graph(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "update",
"step": 5,
"thread_id": "1",
},
parent_config=(
[c async for c in app_w_interrupt.checkpointer.alist(config, limit=2)][
@@ -2739,7 +2766,7 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
app = graph.compile(checkpointer=async_checkpointer)
config = {"configurable": {"thread_id": "1"}}
await app.ainvoke({"my_key": "my value"}, config, checkpoint_during=False)
await app.ainvoke({"my_key": "my value"}, config, debug=True)
# test state w/ nested subgraph state (right after interrupt)
# first get_state without subgraph state
expected = StateSnapshot(
@@ -2764,6 +2791,7 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -2807,6 +2835,12 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
},
"source": "loop",
"step": 1,
"thread_id": "1",
"langgraph_node": "inner",
"langgraph_path": [PULL, "inner"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:inner"],
"langgraph_checkpoint_ns": AnyStr("inner:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -2826,6 +2860,7 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -2859,6 +2894,12 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
"source": "loop",
"step": 1,
"parents": {"": AnyStr()},
"thread_id": "1",
"langgraph_node": "inner",
"langgraph_path": [PULL, "inner"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:inner"],
"langgraph_checkpoint_ns": AnyStr("inner:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -2870,7 +2911,7 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
assert child_history == expected_child_history
# resume
await app.ainvoke(None, config, checkpoint_during=False)
await app.ainvoke(None, config, debug=True)
# test state w/ nested subgraph state (after resuming from interrupt)
assert await app.aget_state(config) == StateSnapshot(
values={"my_key": "hi my value here and there and back again"},
@@ -2887,6 +2928,7 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=(
@@ -2918,6 +2960,7 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=(
@@ -2959,6 +3002,7 @@ async def test_nested_graph_state(async_checkpointer: BaseCheckpointSaver) -> No
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -3027,10 +3071,7 @@ async def test_doubly_nested_graph_state(
# test invoke w/ nested interrupt
config = {"configurable": {"thread_id": "1"}}
assert [
c
async for c in app.astream(
{"my_key": "my value"}, config, subgraphs=True, checkpoint_during=False
)
c async for c in app.astream({"my_key": "my value"}, config, subgraphs=True)
] == [
((), {"parent_1": {"my_key": "hi my value"}}),
(
@@ -3068,6 +3109,7 @@ async def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -3104,9 +3146,15 @@ async def test_doubly_nested_graph_state(
}
},
metadata={
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child",
"langgraph_path": ["__pregel_pull", "child"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:child"],
"parents": {"": AnyStr()},
"source": "loop",
"step": 0,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -3146,6 +3194,14 @@ async def test_doubly_nested_graph_state(
),
"source": "loop",
"step": 1,
"thread_id": "1",
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child_1",
"langgraph_path": [PULL, AnyStr("child_1")],
"langgraph_step": 1,
"langgraph_triggers": [
"branch:to:child_1",
],
},
created_at=AnyStr(),
parent_config=None,
@@ -3201,6 +3257,17 @@ async def test_doubly_nested_graph_state(
),
"source": "loop",
"step": 1,
"thread_id": "1",
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child_1",
"langgraph_path": [
PULL,
AnyStr("child_1"),
],
"langgraph_step": 1,
"langgraph_triggers": [
"branch:to:child_1",
],
},
created_at=AnyStr(),
parent_config=None,
@@ -3223,6 +3290,14 @@ async def test_doubly_nested_graph_state(
"parents": {"": AnyStr()},
"source": "loop",
"step": 0,
"thread_id": "1",
"langgraph_node": "child",
"langgraph_path": [PULL, AnyStr("child")],
"langgraph_step": 2,
"langgraph_triggers": [
"branch:to:child",
],
"langgraph_checkpoint_ns": AnyStr("child:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -3242,18 +3317,14 @@ async def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
interrupts=(),
)
# resume
assert [
c
async for c in app.astream(
None, config, subgraphs=True, checkpoint_during=False
)
] == [
assert [c async for c in app.astream(None, config, subgraphs=True)] == [
(
(AnyStr("child:"), AnyStr("child_1:")),
{"grandchild_2": {"my_key": "hi my value here and there"}},
@@ -3284,6 +3355,7 @@ async def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=(
@@ -3317,6 +3389,7 @@ async def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 3,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config={
@@ -3355,6 +3428,7 @@ async def test_doubly_nested_graph_state(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -3383,6 +3457,12 @@ async def test_doubly_nested_graph_state(
"source": "loop",
"step": 0,
"parents": {"": AnyStr()},
"thread_id": "1",
"langgraph_node": "child",
"langgraph_path": [PULL, AnyStr("child")],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:child"],
"langgraph_checkpoint_ns": AnyStr("child:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -3434,6 +3514,17 @@ async def test_doubly_nested_graph_state(
AnyStr("child:"): AnyStr(),
}
),
"thread_id": "1",
"langgraph_checkpoint_ns": AnyStr("child:"),
"langgraph_node": "child_1",
"langgraph_path": [
PULL,
AnyStr("child_1"),
],
"langgraph_step": 1,
"langgraph_triggers": [
"branch:to:child_1",
],
},
created_at=AnyStr(),
parent_config=None,
@@ -3657,11 +3748,7 @@ async def test_weather_subgraph(
assert [
c
async for c in graph.astream(
inputs,
config=config,
stream_mode="updates",
subgraphs=True,
checkpoint_during=False,
inputs, config=config, stream_mode="updates", subgraphs=True
)
] == [
((), {"router_node": {"route": "weather"}}),
@@ -3688,6 +3775,7 @@ async def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=None,
@@ -3746,11 +3834,7 @@ async def test_weather_subgraph(
assert [
c
async for c in graph.astream(
inputs,
config=config,
stream_mode="updates",
subgraphs=True,
checkpoint_during=False,
inputs, config=config, stream_mode="updates", subgraphs=True
)
] == [
((), {"router_node": {"route": "weather"}}),
@@ -3775,6 +3859,7 @@ async def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "14",
},
created_at=AnyStr(),
parent_config=None,
@@ -3808,6 +3893,12 @@ async def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {"": AnyStr()},
"thread_id": "14",
"langgraph_node": "weather_graph",
"langgraph_path": [PULL, "weather_graph"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:weather_graph"],
"langgraph_checkpoint_ns": AnyStr("weather_graph:"),
},
created_at=AnyStr(),
parent_config=None,
@@ -3847,6 +3938,7 @@ async def test_weather_subgraph(
"source": "loop",
"step": 1,
"parents": {},
"thread_id": "14",
},
created_at=AnyStr(),
parent_config=None,
@@ -3882,6 +3974,14 @@ async def test_weather_subgraph(
"step": 2,
"source": "update",
"parents": {"": AnyStr()},
"thread_id": "14",
"checkpoint_id": AnyStr(),
"checkpoint_ns": AnyStr("weather_graph:"),
"langgraph_node": "weather_graph",
"langgraph_path": [PULL, "weather_graph"],
"langgraph_step": 2,
"langgraph_triggers": ["branch:to:weather_graph"],
"langgraph_checkpoint_ns": AnyStr("weather_graph:"),
},
created_at=AnyStr(),
parent_config=(
+19 -22
View File
@@ -904,6 +904,7 @@ def test_pending_writes_resume(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
}
# get_state with checkpoint_id should not apply any pending writes
state = graph.get_state(state.config)
@@ -993,6 +994,7 @@ def test_pending_writes_resume(
"parents": {},
"step": 1,
"source": "loop",
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -1040,6 +1042,7 @@ def test_pending_writes_resume(
"parents": {},
"step": 0,
"source": "loop",
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -1091,6 +1094,7 @@ def test_pending_writes_resume(
"parents": {},
"step": -1,
"source": "input",
"thread_id": "1",
},
parent_config=None,
pending_writes=UnsortedSequence(
@@ -2054,6 +2058,7 @@ def test_in_one_fan_out_state_graph_waiting_edge(
"parents": {},
"source": "update",
"step": 4,
"thread_id": "2",
},
parent_config=expected_parent_config,
interrupts=(),
@@ -2323,6 +2328,7 @@ def test_in_one_fan_out_state_graph_defer_node(
"parents": {},
"source": "update",
"step": 4,
"thread_id": "2",
},
parent_config=expected_parent_config,
interrupts=(),
@@ -3924,6 +3930,7 @@ def test_checkpoint_metadata(sync_checkpointer: BaseCheckpointSaver) -> None:
# assert that checkpoint metadata contains the run's configurable fields
chkpnt_metadata_1 = sync_checkpointer.get_tuple(config).metadata
assert chkpnt_metadata_1["thread_id"] == "1"
assert chkpnt_metadata_1["test_config_1"] == "foo"
assert chkpnt_metadata_1["test_config_2"] == "bar"
@@ -3932,6 +3939,7 @@ def test_checkpoint_metadata(sync_checkpointer: BaseCheckpointSaver) -> None:
# on how the graph is constructed.
chkpnt_tuples_1 = sync_checkpointer.list(config)
for chkpnt_tuple in chkpnt_tuples_1:
assert chkpnt_tuple.metadata["thread_id"] == "1"
assert chkpnt_tuple.metadata["test_config_1"] == "foo"
assert chkpnt_tuple.metadata["test_config_2"] == "bar"
@@ -3951,6 +3959,7 @@ def test_checkpoint_metadata(sync_checkpointer: BaseCheckpointSaver) -> None:
# assert that checkpoint metadata contains the run's configurable fields
chkpnt_metadata_2 = sync_checkpointer.get_tuple(config).metadata
assert chkpnt_metadata_2["thread_id"] == "2"
assert chkpnt_metadata_2["test_config_3"] == "foo"
assert chkpnt_metadata_2["test_config_4"] == "bar"
@@ -3968,6 +3977,7 @@ def test_checkpoint_metadata(sync_checkpointer: BaseCheckpointSaver) -> None:
# assert that checkpoint metadata contains the run's configurable fields
chkpnt_metadata_3 = sync_checkpointer.get_tuple(config).metadata
assert chkpnt_metadata_3["thread_id"] == "2"
assert chkpnt_metadata_3["test_config_3"] == "foo"
assert chkpnt_metadata_3["test_config_4"] == "bar"
@@ -3976,6 +3986,7 @@ def test_checkpoint_metadata(sync_checkpointer: BaseCheckpointSaver) -> None:
# on how the graph is constructed.
chkpnt_tuples_2 = sync_checkpointer.list(config)
for chkpnt_tuple in chkpnt_tuples_2:
assert chkpnt_tuple.metadata["thread_id"] == "2"
assert chkpnt_tuple.metadata["test_config_3"] == "foo"
assert chkpnt_tuple.metadata["test_config_4"] == "bar"
@@ -4811,9 +4822,7 @@ def test_parent_command(
config = {"configurable": {"thread_id": "1"}}
assert graph.invoke(
{"messages": [("user", "get user name")]}, config, checkpoint_during=False
) == {
assert graph.invoke({"messages": [("user", "get user name")]}, config) == {
"messages": [
_AnyIdHumanMessage(
content="get user name", additional_kwargs={}, response_metadata={}
@@ -4840,6 +4849,7 @@ def test_parent_command(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 1,
"parents": {},
},
@@ -4911,7 +4921,7 @@ def test_interrupt_multiple(sync_checkpointer: BaseCheckpointSaver):
assert [
event
for event in graph.stream(
Command(resume="answer 1", update={"my_key": " foofoo "}), thread1
Command(resume="answer 1", update={"my_key": "foofoo"}), thread1
)
] == [
{
@@ -4926,14 +4936,8 @@ def test_interrupt_multiple(sync_checkpointer: BaseCheckpointSaver):
}
]
assert [
event
for event in graph.stream(
Command(resume="answer 2"), thread1, stream_mode="values"
)
] == [
{"my_key": "DE foofoo "},
{"my_key": "DE foofoo answer 1 answer 2"},
assert [event for event in graph.stream(Command(resume="answer 2"), thread1)] == [
{"node": {"my_key": "answer 1 answer 2"}},
]
@@ -5541,10 +5545,7 @@ def test_falsy_return_from_task(sync_checkpointer: BaseCheckpointSaver):
configurable = {"configurable": {"thread_id": uuid.uuid4()}}
assert [
chunk
for chunk in graph.stream(
{"a": 5}, configurable, stream_mode="debug", checkpoint_during=False
)
chunk for chunk in graph.stream({"a": 5}, configurable, stream_mode="debug")
] == [
{
"payload": {
@@ -5646,12 +5647,7 @@ def test_falsy_return_from_task(sync_checkpointer: BaseCheckpointSaver):
]
assert [
c
for c in graph.stream(
Command(resume="123"),
configurable,
stream_mode="debug",
checkpoint_during=False,
)
for c in graph.stream(Command(resume="123"), configurable, stream_mode="debug")
] == [
{
"payload": {
@@ -5666,6 +5662,7 @@ def test_falsy_return_from_task(sync_checkpointer: BaseCheckpointSaver):
"parents": {},
"source": "input",
"step": -1,
"thread_id": AnyStr(),
},
"next": [
"graph",
+44 -30
View File
@@ -274,9 +274,7 @@ async def test_checkpoint_put_after_cancellation() -> None:
thread1 = {"configurable": {"thread_id": "1"}}
# start the task
t = asyncio.create_task(
graph.ainvoke({"hello": "world"}, thread1, checkpoint_during=False)
)
t = asyncio.create_task(graph.ainvoke({"hello": "world"}, thread1))
# cancel after 0.2 seconds
await asyncio.sleep(0.2)
t.cancel()
@@ -342,7 +340,7 @@ async def test_checkpoint_put_after_cancellation_stream_anext() -> None:
thread1 = {"configurable": {"thread_id": "1"}}
# start the task
s = graph.astream({"hello": "world"}, thread1, checkpoint_during=False)
s = graph.astream({"hello": "world"}, thread1)
t = asyncio.create_task(s.__anext__())
# cancel after 0.2 seconds
await asyncio.sleep(0.2)
@@ -410,11 +408,7 @@ async def test_checkpoint_put_after_cancellation_stream_events_anext() -> None:
# start the task
s = graph.astream_events(
{"hello": "world"},
thread1,
version="v2",
include_names=["LangGraph"],
checkpoint_during=False,
{"hello": "world"}, thread1, version="v2", include_names=["LangGraph"]
)
# skip first event (happens right away)
await s.__anext__()
@@ -601,9 +595,7 @@ async def test_dynamic_interrupt(async_checkpointer: BaseCheckpointSaver) -> Non
# stop when about to enter node
assert [
c
async for c in tool_two.astream(
{"my_key": "value ⛰️", "market": "DE"}, thread1, checkpoint_during=False
)
async for c in tool_two.astream({"my_key": "value ⛰️", "market": "DE"}, thread1)
] == [
{
"__interrupt__": (
@@ -620,6 +612,7 @@ async def test_dynamic_interrupt(async_checkpointer: BaseCheckpointSaver) -> Non
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
]
tup = await tool_two.checkpointer.aget_tuple(thread1)
@@ -646,6 +639,7 @@ async def test_dynamic_interrupt(async_checkpointer: BaseCheckpointSaver) -> Non
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
parent_config=None,
interrupts=(
@@ -671,6 +665,7 @@ async def test_dynamic_interrupt(async_checkpointer: BaseCheckpointSaver) -> Non
"parents": {},
"source": "update",
"step": 1,
"thread_id": "1",
},
parent_config=(
[c async for c in tool_two.checkpointer.alist(thread1, limit=2)][-1].config
@@ -773,9 +768,7 @@ async def test_dynamic_interrupt_subgraph(
# stop when about to enter node
assert [
c
async for c in tool_two.astream(
{"my_key": "value ⛰️", "market": "DE"}, thread1, checkpoint_during=False
)
async for c in tool_two.astream({"my_key": "value ⛰️", "market": "DE"}, thread1)
] == [
{
"__interrupt__": (
@@ -792,6 +785,7 @@ async def test_dynamic_interrupt_subgraph(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
]
tup = await tool_two.checkpointer.aget_tuple(thread1)
@@ -824,6 +818,7 @@ async def test_dynamic_interrupt_subgraph(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
parent_config=None,
interrupts=(
@@ -849,6 +844,7 @@ async def test_dynamic_interrupt_subgraph(
"parents": {},
"source": "update",
"step": 1,
"thread_id": "1",
},
parent_config=(
[c async for c in tool_two.checkpointer.alist(thread1root, limit=2)][
@@ -950,9 +946,7 @@ async def test_copy_checkpoint(async_checkpointer: BaseCheckpointSaver) -> None:
# flow: interrupt -> clear tasks
thread1 = {"configurable": {"thread_id": "1"}}
# stop when about to enter node
assert await tool_two.ainvoke(
{"my_key": "value ⛰️", "market": "DE"}, thread1, checkpoint_during=False
) == {
assert await tool_two.ainvoke({"my_key": "value ⛰️", "market": "DE"}, thread1) == {
"my_key": "value ⛰️ one",
"market": "DE",
"__interrupt__": [
@@ -969,6 +963,7 @@ async def test_copy_checkpoint(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
]
@@ -1005,6 +1000,7 @@ async def test_copy_checkpoint(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
},
parent_config=None,
interrupts=(
@@ -1043,6 +1039,7 @@ async def test_copy_checkpoint(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "fork",
"step": 1,
"thread_id": "1",
},
parent_config=(
[c async for c in tool_two.checkpointer.alist(thread1, limit=2)][-1].config
@@ -1220,6 +1217,7 @@ async def test_cancel_graph_astream(async_checkpointer: BaseCheckpointSaver) ->
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
}
@@ -1294,6 +1292,7 @@ async def test_cancel_graph_astream_events_v2(
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "2",
}
@@ -1816,6 +1815,7 @@ async def test_pending_writes_resume(
"parents": {},
"source": "loop",
"step": 0,
"thread_id": "1",
}
# get_state with checkpoint_id should not apply any pending writes
state = await graph.aget_state(state.config)
@@ -1905,6 +1905,7 @@ async def test_pending_writes_resume(
"parents": {},
"step": 1,
"source": "loop",
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -1952,6 +1953,7 @@ async def test_pending_writes_resume(
"parents": {},
"step": 0,
"source": "loop",
"thread_id": "1",
},
parent_config={
"configurable": {
@@ -1999,6 +2001,7 @@ async def test_pending_writes_resume(
"parents": {},
"step": -1,
"source": "input",
"thread_id": "1",
},
parent_config=None,
pending_writes=UnsortedSequence(
@@ -2598,6 +2601,7 @@ async def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 4,
"parents": {},
},
@@ -2633,6 +2637,7 @@ async def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 3,
"parents": {},
},
@@ -2675,6 +2680,7 @@ async def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 2,
"parents": {},
},
@@ -2729,6 +2735,7 @@ async def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 1,
"parents": {},
},
@@ -2783,6 +2790,7 @@ async def test_send_dedupe_on_resume(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 0,
"parents": {},
},
@@ -2819,6 +2827,7 @@ async def test_send_dedupe_on_resume(
},
metadata={
"source": "input",
"thread_id": "1",
"step": -1,
"parents": {},
},
@@ -2948,9 +2957,7 @@ async def test_send_react_interrupt(async_checkpointer: BaseCheckpointSaver) ->
foo_called = 0
graph = builder.compile(checkpointer=async_checkpointer, interrupt_before=["foo"])
thread1 = {"configurable": {"thread_id": "2"}}
assert await graph.ainvoke(
{"messages": [HumanMessage("hello")]}, thread1, checkpoint_during=False
) == {
assert await graph.ainvoke({"messages": [HumanMessage("hello")]}, thread1) == {
"messages": [
_AnyIdHumanMessage(content="hello"),
_AnyIdAIMessage(
@@ -2999,6 +3006,7 @@ async def test_send_react_interrupt(async_checkpointer: BaseCheckpointSaver) ->
"step": 1,
"source": "loop",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=None,
@@ -3044,6 +3052,7 @@ async def test_send_react_interrupt(async_checkpointer: BaseCheckpointSaver) ->
"step": 2,
"source": "update",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=(
@@ -3072,9 +3081,7 @@ async def test_send_react_interrupt(async_checkpointer: BaseCheckpointSaver) ->
foo_called = 0
graph = builder.compile(checkpointer=async_checkpointer, interrupt_before=["foo"])
thread1 = {"configurable": {"thread_id": "3"}}
assert await graph.ainvoke(
{"messages": [HumanMessage("hello")]}, thread1, checkpoint_during=False
) == {
assert await graph.ainvoke({"messages": [HumanMessage("hello")]}, thread1) == {
"messages": [
_AnyIdHumanMessage(content="hello"),
_AnyIdAIMessage(
@@ -3123,6 +3130,7 @@ async def test_send_react_interrupt(async_checkpointer: BaseCheckpointSaver) ->
"step": 1,
"source": "loop",
"parents": {},
"thread_id": "3",
},
created_at=AnyStr(),
parent_config=None,
@@ -3189,6 +3197,7 @@ async def test_send_react_interrupt(async_checkpointer: BaseCheckpointSaver) ->
"step": 2,
"source": "update",
"parents": {},
"thread_id": "3",
},
created_at=AnyStr(),
parent_config=(
@@ -3337,9 +3346,7 @@ async def test_send_react_interrupt_control(
foo_called = 0
graph = builder.compile(checkpointer=async_checkpointer, interrupt_before=["foo"])
thread1 = {"configurable": {"thread_id": "2"}}
assert await graph.ainvoke(
{"messages": [HumanMessage("hello")]}, thread1, checkpoint_during=False
) == {
assert await graph.ainvoke({"messages": [HumanMessage("hello")]}, thread1) == {
"messages": [
_AnyIdHumanMessage(content="hello"),
_AnyIdAIMessage(
@@ -3388,6 +3395,7 @@ async def test_send_react_interrupt_control(
"step": 1,
"source": "loop",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=None,
@@ -3433,6 +3441,7 @@ async def test_send_react_interrupt_control(
"step": 2,
"source": "update",
"parents": {},
"thread_id": "2",
},
created_at=AnyStr(),
parent_config=(
@@ -4378,6 +4387,7 @@ async def test_in_one_fan_out_state_graph_waiting_edge_custom_state_class(
"parents": {},
"source": "loop",
"step": 4,
"thread_id": "1",
},
created_at=AnyStr(),
parent_config=(
@@ -5636,6 +5646,7 @@ async def test_checkpoint_metadata(async_checkpointer: BaseCheckpointSaver) -> N
# assert that checkpoint metadata contains the run's configurable fields
chkpnt_metadata_1 = (await async_checkpointer.aget_tuple(config)).metadata
assert chkpnt_metadata_1["thread_id"] == "1"
assert chkpnt_metadata_1["test_config_1"] == "foo"
assert chkpnt_metadata_1["test_config_2"] == "bar"
@@ -5644,6 +5655,7 @@ async def test_checkpoint_metadata(async_checkpointer: BaseCheckpointSaver) -> N
# on how the graph is constructed.
chkpnt_tuples_1 = async_checkpointer.alist(config)
async for chkpnt_tuple in chkpnt_tuples_1:
assert chkpnt_tuple.metadata["thread_id"] == "1"
assert chkpnt_tuple.metadata["test_config_1"] == "foo"
assert chkpnt_tuple.metadata["test_config_2"] == "bar"
@@ -5663,6 +5675,7 @@ async def test_checkpoint_metadata(async_checkpointer: BaseCheckpointSaver) -> N
# assert that checkpoint metadata contains the run's configurable fields
chkpnt_metadata_2 = (await async_checkpointer.aget_tuple(config)).metadata
assert chkpnt_metadata_2["thread_id"] == "2"
assert chkpnt_metadata_2["test_config_3"] == "foo"
assert chkpnt_metadata_2["test_config_4"] == "bar"
@@ -5680,6 +5693,7 @@ async def test_checkpoint_metadata(async_checkpointer: BaseCheckpointSaver) -> N
# assert that checkpoint metadata contains the run's configurable fields
chkpnt_metadata_3 = (await async_checkpointer.aget_tuple(config)).metadata
assert chkpnt_metadata_3["thread_id"] == "2"
assert chkpnt_metadata_3["test_config_3"] == "foo"
assert chkpnt_metadata_3["test_config_4"] == "bar"
@@ -5688,6 +5702,7 @@ async def test_checkpoint_metadata(async_checkpointer: BaseCheckpointSaver) -> N
# on how the graph is constructed.
chkpnt_tuples_2 = async_checkpointer.alist(config)
async for chkpnt_tuple in chkpnt_tuples_2:
assert chkpnt_tuple.metadata["thread_id"] == "2"
assert chkpnt_tuple.metadata["test_config_3"] == "foo"
assert chkpnt_tuple.metadata["test_config_4"] == "bar"
@@ -6095,9 +6110,7 @@ async def test_parent_command(
config = {"configurable": {"thread_id": "1"}}
assert await graph.ainvoke(
{"messages": [("user", "get user name")]}, config, checkpoint_during=False
) == {
assert await graph.ainvoke({"messages": [("user", "get user name")]}, config) == {
"messages": [
_AnyIdHumanMessage(
content="get user name", additional_kwargs={}, response_metadata={}
@@ -6126,6 +6139,7 @@ async def test_parent_command(
},
metadata={
"source": "loop",
"thread_id": "1",
"step": 1,
"parents": {},
},
+1 -19
View File
@@ -1,6 +1,6 @@
from dataclasses import dataclass
from operator import add
from typing import Annotated, Any, Union
from typing import Annotated, Any
from langchain_core.runnables import RunnableConfig
from pydantic import BaseModel
@@ -103,21 +103,3 @@ def test_input_state_specified() -> None:
new_graph.invoke({"something": 1})
new_graph.invoke({"something": 2, "info": ["hello", "world"]}) # type: ignore[arg-type]
def test_invokeable_node_signature() -> None:
class State(TypedDict):
info: Annotated[list[str], add]
graph_builder = StateGraph(State)
class RunnableIsh:
def invoke(
self,
input: State,
config: Union[RunnableConfig, None] = None,
**kwargs: Any,
) -> dict[str, str]:
return {}
graph_builder.add_node("runnable", RunnableIsh())
+2
View File
@@ -91,6 +91,7 @@ def test_no_prompt(sync_checkpointer: BaseCheckpointSaver, version: str) -> None
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "123",
}
assert saved.pending_writes == []
@@ -117,6 +118,7 @@ async def test_no_prompt_async(async_checkpointer: BaseCheckpointSaver) -> None:
"parents": {},
"source": "loop",
"step": 1,
"thread_id": "123",
}
assert saved.pending_writes == []
+1 -1
View File
@@ -320,7 +320,7 @@ wheels = [
[[package]]
name = "langgraph"
version = "0.5.0rc1"
version = "0.5.0rc0"
source = { editable = "../langgraph" }
dependencies = [
{ name = "langchain-core" },
+1 -7
View File
@@ -894,13 +894,7 @@ export function useStream<
if (event === "events") options.onLangChainEvent?.(data);
if (event === "debug") options.onDebugEvent?.(data);
if (event === "values") {
if ("__interrupt__" in data) {
// don't update values on interrupt values event
continue;
}
setStreamValues(data);
}
if (event === "values") setStreamValues(data);
if (event === "messages") {
const [serialized] = data;