Stream graph responses
In shortStream token-by-token output using astream(), handle StreamChunk events, and pause with interrupts
- 9 min read
- 14 sections
- Updated
- v0.10.0
- Markdown
Overview
Streaming lets you emit graph responses incrementally as they execute, token by token. Instead of waiting for the full result, your frontend receives chunks in real time, enabling responsive chat interfaces, visible progress during tool use, and the ability to cancel long operations. CompiledGraph.astream() yields StreamChunk objects that you process based on the event type (message, state, tool update, error).
When to stream
- Chat interfaces: Display assistant responses as they arrive, improving perceived latency.
- Long-running workflows: Show progress during extended operations (search, file processing, multiple tool calls).
- Cancellation: Stop the graph mid-execution if the user requests it or closes the connection.
When to prefer ainvoke() instead: if you only need the final result, do not need to show progress, and are willing to wait for the entire run to finish.
Set up a graph for streaming
Streaming requires no special configuration. Any compiled graph can be streamed (compile() uses an InMemoryCheckpointer by default). Install a provider extra for the model, for example pip install "10xgraph[openai]". If you want the graph to pause at specific nodes (for human-in-the-loop workflows), configure interrupt_before or interrupt_after:
from tenxgraph.core.graph import StateGraph, Agent
from tenxgraph.storage.checkpointer import InMemoryCheckpointer
from tenxgraph.utils import END
# Build a simple graph
agent = Agent(model="gpt-4o")
graph = StateGraph()
graph.add_node("MAIN", agent)
graph.set_entry_point("MAIN")
graph.add_edge("MAIN", END)
# Compile with checkpointer; streaming works by default
app = graph.compile(
checkpointer=InMemoryCheckpointer(),
interrupt_before=["MAIN"], # optional: pause before MAIN runs
)The checkpointer is required to save the paused state when using interrupts. For basic streaming without pauses, InMemoryCheckpointer is sufficient for development; use PgCheckpointer in production for durable, multi-replica persistence.
Stream and handle chunks
Call astream() with your input and iterate over chunks. Every chunk is a StreamChunk with an event field that tells you what kind of data it carries.
import asyncio
from tenxgraph.core.state import Message, StreamEvent
from tenxgraph.utils import ResponseGranularity
async def stream_chat(app, user_input: str, thread_id: str):
"""Stream a single turn of conversation."""
print(f"User: {user_input}")
print("Assistant: ", end="", flush=True)
async for chunk in app.astream(
input_data={"messages": [Message.text_message(user_input)]},
config={"thread_id": thread_id},
response_granularity=ResponseGranularity.LOW,
):
# Always check chunk.event first
if chunk.event == StreamEvent.MESSAGE and chunk.message:
# The stream first echoes your input message; skip it
if chunk.message.role == "user":
continue
# Print the text (whether partial or complete)
print(chunk.message.text(), end="", flush=True)
print() # newline after stream ends
# Run a conversation
asyncio.run(stream_chat(app, "What is 2 + 2?", thread_id="session-1"))Key patterns:
- Always branch on
chunk.eventbefore reading data. chunk.message,chunk.state, andchunk.dataare mutually exclusive; only one is set per event.chunk.message.text()extracts text from a message’s content blocks (do not access.contentdirectly).- Use
async forto iterate;astream()is an async generator. - The first
MESSAGEchunk is an echo of your input messages (roleuser). Skip it if you do not want to display it. - With the default
LOWgranularity onlyMESSAGEandERRORchunks are yielded.
Understand StreamChunk structure
Every chunk yielded by astream() is a StreamChunk with these fields:
| Field | Type | When set | What it contains |
|---|---|---|---|
event |
StreamEvent |
Always | One of MESSAGE, STATE, UPDATES, ERROR (string values are lowercase: "message", "state", "updates", "error") |
message |
Message | None |
event == MESSAGE |
A Message object; use .text() to extract text |
state |
AgentState | None |
event == STATE |
Current graph state (depends on response_granularity) |
data |
dict | None |
UPDATES, ERROR, and some tool MESSAGE chunks |
Event-specific metadata (tool status, error reason) |
thread_id |
str | None |
Sometimes | The thread ID from config |
run_id |
str | None |
Sometimes | The run ID from config |
metadata |
dict | None |
Sometimes | Extra metadata (e.g., "node" and "step" keys on STATE chunks) |
Event types
StreamEvent.MESSAGE: A Message from the graph (assistant response, tool result, or tool progress). Readchunk.message.StreamEvent.STATE: A snapshot of the graph state (only whenresponse_granularityisPARTIALorFULL). Readchunk.state.StreamEvent.UPDATES: Lifecycle updates such asinvoking_node,node_invoked,invoking_toolandgraph_invoked(only withFULL). Readchunk.datafor{"status": "...", "node": "..."}; tool updates also carrytool_name.StreamEvent.ERROR: A tool failure. Readchunk.datafor{"status": "tool_failed", "reason": "...", "error": "..."}.
Differentiate partial from complete messages
The chunk.message.delta field is a boolean flag:
delta=True: The message is partial; you are receiving tokens as they arrive. Append to your display buffer.delta=False: The message is complete; replace any streaming placeholder with the final version.
Both cases use chunk.message.text() to get the text:
async for chunk in app.astream(
{"messages": [Message.text_message("Explain gravity.")]},
config={"thread_id": "physics-session"}
):
if chunk.event != StreamEvent.MESSAGE or not chunk.message:
continue
if chunk.message.delta:
# Partial message: update the live display buffer
display_buffer.append(chunk.message.text())
update_ui_stream(display_buffer.join(""))
else:
# Complete message: finalize
final_text = chunk.message.text()
show_final_message(final_text)
display_buffer.clear()Control output detail with ResponseGranularity
The response_granularity parameter controls how much state information is included in chunks:
| Value | What you get | Use case |
|---|---|---|
ResponseGranularity.LOW (default) |
MESSAGE and ERROR chunks only |
Chat UI; minimal overhead |
ResponseGranularity.PARTIAL |
Adds STATE chunks |
Monitor conversation history without lifecycle noise |
ResponseGranularity.FULL |
Adds STATE and UPDATES chunks (node and tool lifecycle, final graph_invoked) |
Debugging, observability, rebuilding state on the client |
async for chunk in app.astream(
{"messages": [Message.text_message("Summarize our conversation.")]},
config={"thread_id": "session-2"},
response_granularity=ResponseGranularity.FULL,
):
if chunk.event == StreamEvent.STATE and chunk.state:
# Full state; useful for complex client-side state management
print(f"Context: {len(chunk.state.context)} messages")
print(f"Summary: {chunk.state.context_summary}")
elif chunk.event == StreamEvent.MESSAGE and chunk.message:
# Still get message tokens
print(chunk.message.text(), end="", flush=True)FULL mode increases network traffic and processing; use it only when you need full observability.
Observe tool calls and progress
UPDATES chunks are only yielded with ResponseGranularity.FULL. When the graph invokes a tool, you receive:
- An
UPDATESchunk when the tool starts (status = “invoking_tool”). - A
MESSAGEchunk with the tool result (role = “tool”). - An
ERRORchunk with statustool_failedif the tool returned an error.
The producing node name is in chunk.metadata.get("node") for state chunks and chunk.data.get("node") for update and error chunks:
from tenxgraph.core.state import StreamEvent
async for chunk in app.astream(
{"messages": [Message.text_message("What is 15 * 23?")]},
config={"thread_id": "calc-session"},
response_granularity=ResponseGranularity.FULL,
):
# Extract node name (present in most chunks)
node = (chunk.metadata or {}).get("node") or (chunk.data or {}).get("node") or "unknown"
if chunk.event == StreamEvent.UPDATES and chunk.data:
# Tool lifecycle event
status = chunk.data.get("status")
tool_name = chunk.data.get("tool_name", "")
if status == "invoking_tool":
print(f"[{node}] Invoking {tool_name}...")
elif chunk.event == StreamEvent.MESSAGE and chunk.message:
# Text or tool result
if chunk.message.role == "tool":
print(f"[{node}] Tool result: {chunk.message.text()}")
elif chunk.message.role == "user":
continue
else:
print(chunk.message.text(), end="", flush=True)
elif chunk.event == StreamEvent.ERROR and chunk.data:
reason = chunk.data.get("reason", "unknown error")
print(f"[{node}] Error: {reason}")If your tools emit progress updates using StreamEmitter, those arrive as UPDATES chunks (or MESSAGE/ERROR chunks for progress and failure helpers).
Collect the full response
If you want to display tokens during the stream but also keep the final messages for later use (e.g., saving to a database), collect them as chunks arrive:
from tenxgraph.core.state import StreamEvent, Message
async def stream_to_messages(
app,
input_messages: list[Message],
thread_id: str
) -> list[Message]:
"""Run streaming and return the final messages."""
final_messages: list[Message] = []
seen: set[str] = set()
async for chunk in app.astream(
{"messages": input_messages},
config={"thread_id": thread_id},
):
if chunk.event != StreamEvent.MESSAGE or not chunk.message:
continue
# Skip partial deltas and the echoed user input
if chunk.message.delta or chunk.message.role == "user":
continue
# Avoid duplicates by message ID
msg_id = str(chunk.message.message_id)
if msg_id in seen:
continue
seen.add(msg_id)
final_messages.append(chunk.message)
return final_messages
# Use it
messages = await stream_to_messages(
app,
[Message.text_message("Generate a report.")],
thread_id="report-session"
)
# messages now contains all final messages from the runAlternatively, run with ResponseGranularity.FULL and extract the context list from the last STATE chunk.
Stop a stream early
To cancel a running stream (e.g., if the user closes the chat connection), call astop() from another task. The graph checks the stop flag between nodes and halts:
import asyncio
from tenxgraph.core.state import Message, StreamEvent
thread_id = "long-running-session"
async def run_stream():
"""Stream a long-running query."""
async for chunk in app.astream(
{"messages": [Message.text_message("Write a 5000-word essay on climate change.")]},
config={"thread_id": thread_id},
):
if chunk.event == StreamEvent.MESSAGE and chunk.message:
print(chunk.message.text(), end="", flush=True)
async def stop_after_delay():
"""Request stop after 10 seconds."""
await asyncio.sleep(10.0)
result = await app.astop({"thread_id": thread_id})
print(f"\nStop requested: {result}")
async def main():
# Run both concurrently
await asyncio.gather(
run_stream(),
stop_after_delay()
)
asyncio.run(main())The sync equivalent is app.stop(config) for non-async contexts.
Stop behavior
- Stop is checked between node executions; a running node or tool call finishes before the graph halts.
astop()needs a thread with a checkpointer. It returns{"ok": True, "running": True}when a stop was recorded, or{"ok": True, "running": False, "reason": "not-running"}when the thread is idle or paused.
Pause and resume with interrupts
For human-in-the-loop workflows, pause the graph at specific nodes to let a human review or approve the state before continuing.
Compile with interrupts
Declare which nodes should pause:
app = graph.compile(
checkpointer=checkpointer,
interrupt_before=["review_node"], # pause before this node
interrupt_after=["approval_node"], # pause after this node
)Run until the pause
async for chunk in app.astream(
{"messages": [Message.text_message("Start the review workflow.")]},
config={"thread_id": "review-workflow-1"},
):
if chunk.event == StreamEvent.MESSAGE and chunk.message:
print(chunk.message.text(), end="", flush=True)
# Graph is now paused before "review_node"
print("\nPaused. Human review happening...")Resume after approval
On the same thread_id, call astream() or ainvoke() again with new input. The graph resumes from the pause point. (If the pause came from an interrupt() call inside a node, resume with {"resume": value} instead; see Add human approval.)
# Simulate human approval
await asyncio.sleep(2.0)
# Resume on the same thread
async for chunk in app.astream(
{"messages": [Message.text_message("Approved. Continue with execution.")]},
config={"thread_id": "review-workflow-1"},
):
if chunk.event == StreamEvent.MESSAGE and chunk.message:
print(chunk.message.text(), end="", flush=True)
print("\nWorkflow completed.")The new message is appended to the conversation history and the graph continues from where it paused. The checkpoint preserves everything between pauses, so your state is consistent.
Handle errors in streams
Errors from nodes or tools emit StreamEvent.ERROR chunks. The chunk.data dict contains:
"reason": Human-readable error description."error": The error message from the failed tool."status","tool_name","tool_call_id","node": which tool failed.
async for chunk in app.astream(
{"messages": input_messages},
config={"thread_id": "error-demo"},
):
if chunk.event == StreamEvent.ERROR and chunk.data:
reason = chunk.data.get("reason", "unknown error")
print(f"Error occurred: {reason}")
# Decide whether to retry, fallback, or fail
if "rate_limit" in reason.lower():
print("Rate limited. Retrying in 5 seconds...")
await asyncio.sleep(5.0)
else:
print("Fatal error. Stopping.")
break
elif chunk.event == StreamEvent.MESSAGE and chunk.message:
print(chunk.message.text(), end="", flush=True)ERROR chunks report tool failures. Node exceptions are raised from the async for loop, so wrap the loop in try/except if you need to handle them.
Complete example: streaming chat app
import asyncio
from tenxgraph.core.graph import StateGraph, Agent, ToolNode
from tenxgraph.core.state import AgentState, Message, StreamEvent
from tenxgraph.storage.checkpointer import InMemoryCheckpointer
from tenxgraph.utils import ResponseGranularity, END
# Define a tool
def get_current_time() -> str:
from datetime import datetime
return datetime.now().isoformat()
# Build the graph
tool_node = ToolNode([get_current_time])
agent = Agent(
model="gpt-4o",
system_prompt=[
{
"role": "system",
"content": "You are a helpful assistant. Use tools when the user asks for the time."
}
],
tool_node=tool_node,
)
graph = StateGraph()
graph.add_node("MAIN", agent)
graph.add_node("TOOL", tool_node)
def should_use_tools(state: AgentState) -> str:
if state.context and state.context[-1].role == "assistant":
if state.context[-1].tools_calls:
return "TOOL"
return END
graph.add_conditional_edges("MAIN", should_use_tools, {"TOOL": "TOOL", END: END})
graph.add_edge("TOOL", "MAIN")
graph.set_entry_point("MAIN")
# Compile and stream
app = graph.compile(checkpointer=InMemoryCheckpointer())
async def chat_session():
"""Multi-turn streaming conversation."""
thread_id = "user-session-1"
messages = []
async def turn(user_input: str):
"""Run a single conversation turn."""
print(f"\nYou: {user_input}")
print("Assistant: ", end="", flush=True)
async for chunk in app.astream(
{"messages": [Message.text_message(user_input)]},
config={"thread_id": thread_id},
response_granularity=ResponseGranularity.LOW,
):
if chunk.event == StreamEvent.MESSAGE and chunk.message:
if chunk.message.role == "user":
continue # skip the echoed input
print(chunk.message.text(), end="", flush=True)
print() # newline
# Run multiple turns
await turn("Hello! What's your name?")
await turn("What time is it right now?")
await turn("Thanks for your help!")
asyncio.run(chat_session())Output varies by model. Each turn prints the assistant text as it streams, and the second turn also runs the get_current_time tool before the final answer.
Related guides
- Build a graph: Construct and run graphs without streaming.
- Add human approval: Use the
interrupt()function for human-in-the-loop patterns. - Set up checkpointing: Choose a checkpointer for production persistence.
- Concepts: Streaming: Why streaming, transports (REST, WebSocket), and client-side handling.
Frequently asked questions
- When should I use astream() instead of ainvoke()?
- Use astream() to display tokens as they arrive (lower latency in chat UIs), observe intermediate progress during long runs, or cancel execution early. Use ainvoke() when you need the final result and can wait for the entire graph to complete.
- What's the difference between chunk.message.delta and a complete message?
- delta=True means the message is partial (token-by-token text arriving). delta=False means the complete message is ready. Both carry text in chunk.message.text().
- Do I need to configure streaming, or does it work with any graph?
- Streaming works with any compiled graph without configuration. Just call astream() instead of ainvoke(). To pause at specific nodes, pass interrupt_before or interrupt_after to compile().