Agentic AI #langgraph#human-in-the-loop#fastapi#websocket#approval-queue#agents#production

Implementing Dynamic Human-in-the-Loop Approval Queues in LangGraph with FastAPI

S

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.
✉ Newsletter

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):

  1. Compile the graph with interrupt_before=["risky_node"]
  2. Pause — when the graph hits the interrupt, it commits state to Postgres and halts
  3. Resume — after human approval, call graph.ainvoke again with the same thread_id
graph = 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

PackageVersion
langgraph0.2.14
langgraph-checkpoint-postgres2.0.1
fastapi0.115+
asyncpg0.29.0
websockets12.0
Python3.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.

✉ Newsletter

Want to build production-ready AI?

Subscribe to StackMindset to receive actionable systems engineering checklists and code walkthroughs. No spam, only technical insights.

S

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

Agentic AI
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).

Agentic AI
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.

Agentic AI
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.