Implementing Dynamic Human-in-the-Loop Approval Queues in LangGraph with FastAPI
S L Manikanta
Sep 13, 2026 • 7 min read
bolt Key Takeaways
- LangGraph's interrupt_before=['node_name'] pauses execution before a node and persists state to the checkpointer — the graph resumes when you call graph.ainvoke with the same thread_id.
- Combine interrupt_before with AsyncPostgresSaver so interrupted state survives process restarts — approval workflows can wait hours or days without losing context.
- Use a FastAPI WebSocket endpoint to push real-time approval requests to a reviewer UI, and a POST endpoint to submit decisions and resume the graph.
- Store pending approvals in a Postgres queue table — not in LangGraph state — so you can query 'all pending approvals' without scanning checkpoints.
Want to build production-ready AI?
Subscribe to StackMindset to receive actionable systems engineering checklists and code walkthroughs. No spam, only technical insights.
[!NOTE] Core Pattern (3 steps):
- Compile the graph with
interrupt_before=["risky_node"]- Pause — when the graph hits the interrupt, it commits state to Postgres and halts
- Resume — after human approval, call
graph.ainvokeagain with the samethread_idgraph = builder.compile( checkpointer=postgres_checkpointer, interrupt_before=["send_email"], # pause before this node ) config = {"configurable": {"thread_id": "run-99"}} # First invocation — pauses at "send_email" await graph.ainvoke({"messages": [...]}, config=config) # Human reviews and approves — update state, then resume await graph.aupdate_state(config, {"approved": True}) await graph.ainvoke(None, config=config) # resumes from "send_email"
Building reliable human-in-the-loop is harder than tutorials suggest: the naive demo shows a synchronous approve/reject in the same Python process. Production means: the graph pauses, the reviewer might not respond for 6 hours, the server restarts in between, and the graph must resume cleanly anyway. Here’s the full implementation.
Environment
| Package | Version |
|---|---|
langgraph | 0.2.14 |
langgraph-checkpoint-postgres | 2.0.1 |
fastapi | 0.115+ |
asyncpg | 0.29.0 |
websockets | 12.0 |
| Python | 3.11+ |
1. Architecture: Pause, Queue, Resume
sequenceDiagram
participant Agent as LangGraph Agent
participant PG as Postgres\n(checkpoints + approval queue)
participant API as FastAPI
participant UI as Reviewer UI
Agent->>Agent: executes up to interrupt point
Agent->>PG: commit state checkpoint
Agent->>PG: insert approval request (thread_id, payload)
Agent-->>API: raise Interrupt (HTTP response returns)
API->>UI: WebSocket push: "new approval pending"
UI->>API: GET /approvals/pending
API->>PG: SELECT pending approvals
PG-->>API: [{thread_id, payload, created_at}]
API-->>UI: approval details
UI->>API: POST /approvals/{thread_id}/decide {approved: true}
API->>PG: UPDATE approval status = approved
API->>Agent: graph.aupdate_state + graph.ainvoke
Agent->>Agent: resumes from interrupted node
Agent-->>API: completed run
2. State Design
from typing import TypedDict, Annotated
from langgraph.graph import add_messages
from langchain_core.messages import BaseMessage
class AgentState(TypedDict):
messages: Annotated[list[BaseMessage], add_messages]
# Human review fields — populated before resuming
approved: bool | None # None = not yet reviewed
rejection_reason: str | None
reviewer_id: str | None
# Action payload that needs approval
pending_action: dict | None # e.g., {"type": "send_email", "to": "...", "body": "..."}
3. Graph with Interrupt
from langgraph.graph import StateGraph, END
from langchain_anthropic import ChatAnthropic
from langchain_core.messages import SystemMessage
llm = ChatAnthropic(model="claude-3-5-sonnet-20241022")
def plan_action(state: AgentState) -> AgentState:
"""Decide what action to take. Sets pending_action."""
messages = [SystemMessage(content="You are an email drafting assistant.")] + state["messages"]
response = llm.invoke(messages)
# In real code, parse the LLM response to extract the action
pending_action = {
"type": "send_email",
"to": "[email protected]",
"subject": "Weekly AI Summary",
"body": response.content,
}
return {"messages": [response], "pending_action": pending_action, "approved": None}
def send_email(state: AgentState) -> AgentState:
"""Send the email. Only runs after human approval."""
if not state.get("approved"):
return {"messages": [{"role": "assistant", "content": "Email sending was rejected or not approved."}]}
action = state["pending_action"]
# Actually send the email
_send_email_impl(action["to"], action["subject"], action["body"])
return {"messages": [{"role": "assistant", "content": f"Email sent to {action['to']}."}]}
def route_after_approval(state: AgentState) -> str:
if state.get("approved") is False:
return "handle_rejection"
return "send_email"
def handle_rejection(state: AgentState) -> AgentState:
reason = state.get("rejection_reason", "No reason provided")
return {"messages": [{"role": "assistant", "content": f"Email rejected: {reason}. Stopping."}]}
builder = StateGraph(AgentState)
builder.add_node("plan_action", plan_action)
builder.add_node("send_email", send_email)
builder.add_node("handle_rejection", handle_rejection)
builder.set_entry_point("plan_action")
# Interrupt BEFORE send_email — this is where we pause for human review
builder.add_conditional_edges("plan_action", route_after_approval)
builder.add_edge("send_email", END)
builder.add_edge("handle_rejection", END)
graph = builder.compile(
checkpointer=postgres_checkpointer,
interrupt_before=["send_email"], # pause before executing send_email
)
4. Approval Queue in Postgres
Create a separate table to track pending approvals (do not rely on scanning checkpoint tables):
CREATE TABLE approval_requests (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
thread_id TEXT NOT NULL,
action_type TEXT NOT NULL,
payload JSONB NOT NULL,
status TEXT NOT NULL DEFAULT 'pending', -- 'pending', 'approved', 'rejected'
reviewer_id TEXT,
reviewer_note TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
decided_at TIMESTAMPTZ
);
CREATE INDEX idx_approval_requests_status ON approval_requests(status);
CREATE INDEX idx_approval_requests_thread_id ON approval_requests(thread_id);
5. FastAPI Endpoints
import uuid
from contextlib import asynccontextmanager
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
import asyncpg
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
pool: asyncpg.Pool = None
checkpointer: AsyncPostgresSaver = None
# Track connected reviewer WebSockets for push notifications
connected_reviewers: list[WebSocket] = []
@asynccontextmanager
async def lifespan(app: FastAPI):
global pool, checkpointer
pool = await asyncpg.create_pool(dsn="postgresql://...", min_size=2, max_size=10)
checkpointer = AsyncPostgresSaver(pool)
await checkpointer.setup()
yield
await pool.close()
app = FastAPI(lifespan=lifespan)
@app.post("/agent/run")
async def start_agent_run(user_message: str):
"""Start a new agent run. Returns immediately after the interrupt fires."""
thread_id = str(uuid.uuid4())
config = {"configurable": {"thread_id": thread_id}}
# Run graph — will pause at interrupt_before=["send_email"]
result = await graph.ainvoke(
{"messages": [{"role": "user", "content": user_message}]},
config=config,
)
# After interrupt, graph state is saved. Insert approval request.
state = await graph.aget_state(config)
pending_action = state.values.get("pending_action")
if pending_action:
async with pool.acquire() as conn:
await conn.execute(
"""
INSERT INTO approval_requests (thread_id, action_type, payload)
VALUES ($1, $2, $3)
""",
thread_id,
pending_action["type"],
pending_action, # asyncpg serializes dict to JSONB
)
# Push notification to all connected reviewer WebSockets
for ws in connected_reviewers:
try:
await ws.send_json({"event": "new_approval", "thread_id": thread_id})
except Exception:
pass
return {"thread_id": thread_id, "status": "awaiting_approval"}
@app.get("/approvals/pending")
async def list_pending_approvals():
"""List all pending approval requests."""
async with pool.acquire() as conn:
rows = await conn.fetch(
"SELECT id, thread_id, action_type, payload, created_at FROM approval_requests WHERE status = 'pending' ORDER BY created_at ASC"
)
return [dict(row) for row in rows]
@app.post("/approvals/{thread_id}/decide")
async def decide_approval(
thread_id: str,
approved: bool,
reviewer_id: str = "system",
reviewer_note: str = "",
):
"""Accept or reject a pending approval and resume the graph."""
config = {"configurable": {"thread_id": thread_id}}
# Update the approval queue record
async with pool.acquire() as conn:
await conn.execute(
"""
UPDATE approval_requests
SET status = $1, reviewer_id = $2, reviewer_note = $3, decided_at = NOW()
WHERE thread_id = $4 AND status = 'pending'
""",
"approved" if approved else "rejected",
reviewer_id,
reviewer_note,
thread_id,
)
# Update LangGraph state with the human decision
await graph.aupdate_state(
config,
{
"approved": approved,
"rejection_reason": reviewer_note if not approved else None,
"reviewer_id": reviewer_id,
},
)
# Resume the graph — runs from the interrupted node with updated state
result = await graph.ainvoke(None, config=config)
return {"status": "resumed", "approved": approved}
@app.websocket("/ws/approvals")
async def approval_websocket(websocket: WebSocket):
"""WebSocket endpoint for real-time approval notifications."""
await websocket.accept()
connected_reviewers.append(websocket)
try:
while True:
await websocket.receive_text() # keep connection alive
except WebSocketDisconnect:
connected_reviewers.remove(websocket)
6. Timeout Handling for Stale Approvals
Add a background job to auto-reject approvals that haven’t been reviewed within your SLA:
import asyncio
async def timeout_stale_approvals(timeout_hours: int = 24):
"""Auto-reject approvals older than timeout_hours."""
async with pool.acquire() as conn:
stale_approvals = await conn.fetch(
"""
SELECT thread_id FROM approval_requests
WHERE status = 'pending'
AND created_at < NOW() - INTERVAL '1 hour' * $1
""",
timeout_hours,
)
for row in stale_approvals:
thread_id = row["thread_id"]
config = {"configurable": {"thread_id": thread_id}}
await graph.aupdate_state(
config,
{"approved": False, "rejection_reason": "Auto-rejected: review timeout exceeded"},
)
await graph.ainvoke(None, config=config)
async with pool.acquire() as conn:
await conn.execute(
"UPDATE approval_requests SET status = 'timed_out' WHERE thread_id = $1",
thread_id,
)
# Run every hour with APScheduler or asyncio.create_task in lifespan
Next Steps
For the durable Postgres checkpointer that makes interrupted state survive process restarts, see Building Resilient LangGraph Workflows with Async Postgres Checkpointing.
For comparing LangGraph’s human-in-the-loop capabilities against PydanticAI, see LangGraph vs PydanticAI: State Machine Architecture, Latency, and Memory Footprint.
The LangGraph Complete Guide covers the full StateGraph API including interrupt_after, astream_events, and subgraph composition.
Want to build production-ready AI?
Subscribe to StackMindset to receive actionable systems engineering checklists and code walkthroughs. No spam, only technical insights.
Written by S L Manikanta
AI Engineer specializing in agentic workflows, multi-step LLM validation pipelines, and secure cloud environments. Sharing practical lessons from building software.
Related Articles
AI Agent Architecture Patterns: A Guide for Platform Engineers (2026)
A technical comparison of AI agent architectures. Learn when to use Prompt Chaining, Routing, Orchestrator-Workers, and Cyclic State Graphs (LangGraph).
What is an AI Agent Harness? Complete Technical Reference (2026)
A comprehensive technical reference on AI Agent Harnesses. Learn architecture, security, cost optimization, and how to deploy LangGraph agents into production with custom harnesses.
Building StoxFlow: Hybrid Local/Cloud AI Agents Architecture Complete Guide (2026)
How to build a decoupled three-tier AI stock research agent using LangGraph, FastAPI, and hybrid LLM routing (Ollama/Gemini) for optimal cost and performance.