🔥 0
0
Lesson 8 of 10 20 min +200 XP

Human-in-the-Loop Workflows

Not every decision should be automated. A $5,000 refund request, a first-time customer ordering 50 laptops, or an AI-generated response to a legal complaint - these need human oversight. LangGraph's interrupt mechanism makes it easy to pause workflows, get human input, and resume seamlessly.

When to Involve Humans

HIGH-RISK ACTIONS

  • Large refunds (> $500)
  • Account deletion
  • Credit extensions
  • Fraud overrides

QUALITY CONTROL

  • AI response review
  • Content moderation
  • Legal/compliance
  • VIP customer handling

EDGE CASES

  • Low confidence decisions
  • Unusual patterns
  • Complex multi-issue tickets
  • Escalation requests

How LangGraph Interrupt Works

Customer: "I want a refund for order #123"
                │
                ▼
        ┌───────────────┐
        │ Check Order   │
        │ & Calculate   │
        └───────┬───────┘
                │ Refund = $750 (> $500 threshold)
                ▼
        ┌───────────────┐
        │   INTERRUPT   │ ◄── Graph pauses here
        │  "Approve     │     State saved to checkpointer
        │   $750?"      │
        └───────┬───────┘
                │
                ▼
        [WAITING FOR HUMAN]
                │
        Manager clicks "Approve"
                │
                ▼
        ┌───────────────┐
        │   RESUME      │ ◄── Graph continues
        │  with input   │
        └───────┬───────┘
                │
                ▼
        ┌───────────────┐
        │  Process      │
        │  Refund       │
        └───────────────┘

Basic Interrupt Example

from langgraph.graph import StateGraph, START, END
from langgraph.types import interrupt, Command
from langgraph.checkpoint.memory import MemorySaver
from typing import TypedDict, Literal, Optional

class RefundState(TypedDict):
    order_id: str
    customer_id: str
    refund_amount: float
    refund_reason: str
    approval_status: Optional[Literal["pending", "approved", "rejected"]]
    result: Optional[str]

def calculate_refund(state: RefundState) -> dict:
    """Calculate the refund amount"""
    # In production: look up order, calculate based on items returned
    return {
        "refund_amount": 750.00,  # Example amount
        "refund_reason": "Item defective"
    }

def check_approval_needed(state: RefundState) -> dict:
    """Determine if human approval is required"""
    APPROVAL_THRESHOLD = 500

    if state["refund_amount"] > APPROVAL_THRESHOLD:
        # Request human approval
        approval = interrupt({
            "type": "refund_approval",
            "order_id": state["order_id"],
            "amount": state["refund_amount"],
            "reason": state["refund_reason"],
            "message": f"Approve refund of ${state['refund_amount']:.2f} for order {state['order_id']}?",
            "options": ["approve", "reject", "modify"]
        })

        # This code runs AFTER human responds
        if approval == "approve":
            return {"approval_status": "approved"}
        elif approval == "reject":
            return {"approval_status": "rejected"}
        else:
            # Human modified the amount
            return {
                "approval_status": "approved",
                "refund_amount": float(approval.get("new_amount", state["refund_amount"]))
            }
    else:
        # Auto-approve small refunds
        return {"approval_status": "approved"}

def process_refund(state: RefundState) -> dict:
    """Process the approved refund"""
    if state["approval_status"] == "approved":
        # In production: call payment gateway
        return {"result": f"Refund of ${state['refund_amount']:.2f} processed successfully"}
    else:
        return {"result": "Refund rejected by manager"}

# Build graph
workflow = StateGraph(RefundState)
workflow.add_node("calculate", calculate_refund)
workflow.add_node("check_approval", check_approval_needed)
workflow.add_node("process", process_refund)

workflow.add_edge(START, "calculate")
workflow.add_edge("calculate", "check_approval")
workflow.add_edge("check_approval", "process")
workflow.add_edge("process", END)

# Compile with checkpointer (REQUIRED for interrupt)
memory = MemorySaver()
refund_app = workflow.compile(checkpointer=memory)

Running an Interruptible Workflow

# Start the refund request
config = {"configurable": {"thread_id": "refund_001"}}

result = refund_app.invoke({
    "order_id": "ORD-12345",
    "customer_id": "CUST-789",
    "refund_amount": 0,
    "refund_reason": "",
    "approval_status": None,
    "result": None
}, config)

# Check if interrupted
if "__interrupt__" in result:
    interrupt_info = result["__interrupt__"][0]
    print(f"Approval needed: {interrupt_info.value}")
    # Output: Approval needed: {'type': 'refund_approval', 'amount': 750.0, ...}

Resuming After Human Decision

# Manager approves the refund
final_result = refund_app.invoke(
    Command(resume="approve"),  # Human's decision
    config  # Same thread_id!
)

print(final_result["result"])
# Output: "Refund of $750.00 processed successfully"

Real-World Example: Order Review Workflow

Let's build a comprehensive order review system with multiple approval points:

from langgraph.graph import StateGraph, START, END
from langgraph.types import interrupt, Command
from typing import TypedDict, Literal, Optional, Annotated
import operator

class OrderReviewState(TypedDict):
    # Order details
    order_id: str
    customer_id: str
    customer_tier: str
    order_total: float
    items: list[dict]

    # Risk factors
    is_first_order: bool
    fraud_score: float
    shipping_billing_mismatch: bool

    # Review status
    risk_level: Literal["low", "medium", "high"]
    review_notes: Annotated[list[str], operator.add]
    approval_status: Optional[Literal["approved", "rejected", "modified"]]

    # Final result
    processing_decision: str

def assess_risk(state: OrderReviewState) -> dict:
    """Assess order risk level"""
    risk_factors = 0
    notes = []

    # High-value order
    if state["order_total"] > 1000:
        risk_factors += 2
        notes.append(f"High value order: ${state['order_total']:.2f}")

    # First-time customer
    if state["is_first_order"]:
        risk_factors += 1
        notes.append("First-time customer")

    # Fraud score
    if state["fraud_score"] > 0.5:
        risk_factors += 2
        notes.append(f"Elevated fraud score: {state['fraud_score']}")
    elif state["fraud_score"] > 0.3:
        risk_factors += 1

    # Address mismatch
    if state["shipping_billing_mismatch"]:
        risk_factors += 1
        notes.append("Shipping/billing address mismatch")

    # Determine risk level
    if risk_factors >= 4:
        risk_level = "high"
    elif risk_factors >= 2:
        risk_level = "medium"
    else:
        risk_level = "low"

    return {
        "risk_level": risk_level,
        "review_notes": notes
    }

def route_by_risk(state: OrderReviewState) -> str:
    """Route based on risk level"""
    if state["risk_level"] == "high":
        return "senior_review"
    elif state["risk_level"] == "medium":
        return "standard_review"
    else:
        return "auto_approve"

def request_senior_review(state: OrderReviewState) -> dict:
    """Request senior reviewer approval for high-risk orders"""
    decision = interrupt({
        "review_type": "senior_approval",
        "priority": "high",
        "order_id": state["order_id"],
        "total": state["order_total"],
        "risk_factors": state["review_notes"],
        "customer_tier": state["customer_tier"],
        "fraud_score": state["fraud_score"],
        "message": "HIGH RISK ORDER - Senior approval required",
        "actions": {
            "approve": "Approve order for processing",
            "reject": "Reject and flag customer",
            "verify": "Request additional verification from customer"
        }
    })

    if decision["action"] == "approve":
        return {
            "approval_status": "approved",
            "review_notes": [f"Senior approved by {decision.get('reviewer', 'unknown')}"]
        }
    elif decision["action"] == "reject":
        return {
            "approval_status": "rejected",
            "review_notes": [f"Rejected: {decision.get('reason', 'No reason provided')}"]
        }
    else:
        return {
            "approval_status": "modified",
            "review_notes": [f"Verification requested: {decision.get('verification_type', 'ID check')}"]
        }

def request_standard_review(state: OrderReviewState) -> dict:
    """Request standard reviewer approval for medium-risk orders"""
    decision = interrupt({
        "review_type": "standard_approval",
        "priority": "medium",
        "order_id": state["order_id"],
        "total": state["order_total"],
        "risk_factors": state["review_notes"],
        "message": "Order flagged for review",
        "actions": ["approve", "reject"]
    })

    return {
        "approval_status": "approved" if decision == "approve" else "rejected",
        "review_notes": [f"Standard review: {decision}"]
    }

def auto_approve(state: OrderReviewState) -> dict:
    """Automatically approve low-risk orders"""
    return {
        "approval_status": "approved",
        "review_notes": ["Auto-approved: Low risk order"]
    }

def finalize_order(state: OrderReviewState) -> dict:
    """Finalize order based on approval status"""
    if state["approval_status"] == "approved":
        return {"processing_decision": "Order approved for fulfillment"}
    elif state["approval_status"] == "rejected":
        return {"processing_decision": "Order cancelled - customer notified"}
    else:
        return {"processing_decision": "Order held pending verification"}

# Build the graph
workflow = StateGraph(OrderReviewState)

workflow.add_node("assess_risk", assess_risk)
workflow.add_node("senior_review", request_senior_review)
workflow.add_node("standard_review", request_standard_review)
workflow.add_node("auto_approve", auto_approve)
workflow.add_node("finalize", finalize_order)

workflow.add_edge(START, "assess_risk")
workflow.add_conditional_edges(
    "assess_risk",
    route_by_risk,
    {
        "senior_review": "senior_review",
        "standard_review": "standard_review",
        "auto_approve": "auto_approve"
    }
)
workflow.add_edge("senior_review", "finalize")
workflow.add_edge("standard_review", "finalize")
workflow.add_edge("auto_approve", "finalize")
workflow.add_edge("finalize", END)

memory = MemorySaver()
order_review_app = workflow.compile(checkpointer=memory)

Building a Review Dashboard

# Simulating a review queue system
class ReviewQueue:
    def __init__(self, app):
        self.app = app
        self.pending_reviews = {}

    def submit_order(self, order_data: dict) -> str:
        """Submit an order for review"""
        thread_id = f"order_{order_data['order_id']}"
        config = {"configurable": {"thread_id": thread_id}}

        result = self.app.invoke(order_data, config)

        if "__interrupt__" in result:
            # Store for human review
            interrupt_data = result["__interrupt__"][0].value
            self.pending_reviews[thread_id] = {
                "interrupt": interrupt_data,
                "state": result,
                "config": config
            }
            return f"Order {order_data['order_id']} pending review"
        else:
            return f"Order {order_data['order_id']} auto-approved"

    def get_pending_reviews(self) -> list:
        """Get all orders pending review"""
        return [
            {
                "thread_id": tid,
                "order_id": data["interrupt"].get("order_id"),
                "priority": data["interrupt"].get("priority"),
                "risk_factors": data["interrupt"].get("risk_factors"),
                "total": data["interrupt"].get("total")
            }
            for tid, data in self.pending_reviews.items()
        ]

    def approve_order(self, thread_id: str, reviewer: str) -> str:
        """Approve a pending order"""
        if thread_id not in self.pending_reviews:
            return "Order not found in review queue"

        review_data = self.pending_reviews[thread_id]

        result = self.app.invoke(
            Command(resume={"action": "approve", "reviewer": reviewer}),
            review_data["config"]
        )

        del self.pending_reviews[thread_id]
        return result["processing_decision"]

    def reject_order(self, thread_id: str, reviewer: str, reason: str) -> str:
        """Reject a pending order"""
        if thread_id not in self.pending_reviews:
            return "Order not found in review queue"

        review_data = self.pending_reviews[thread_id]

        result = self.app.invoke(
            Command(resume={"action": "reject", "reviewer": reviewer, "reason": reason}),
            review_data["config"]
        )

        del self.pending_reviews[thread_id]
        return result["processing_decision"]

# Usage
queue = ReviewQueue(order_review_app)

# Submit a high-risk order
status = queue.submit_order({
    "order_id": "ORD-99999",
    "customer_id": "CUST-NEW",
    "customer_tier": "standard",
    "order_total": 2500.00,
    "items": [{"sku": "LAPTOP-001", "qty": 2}],
    "is_first_order": True,
    "fraud_score": 0.6,
    "shipping_billing_mismatch": True,
    "risk_level": "low",  # Will be recalculated
    "review_notes": [],
    "approval_status": None,
    "processing_decision": ""
})

print(status)  # "Order ORD-99999 pending review"

# View pending reviews
print(queue.get_pending_reviews())

# Manager approves
result = queue.approve_order("order_ORD-99999", "manager_jane")
print(result)  # "Order approved for fulfillment"

Best Practices

DO: Set Clear Thresholds

Define explicit rules for what needs human review: amounts, risk scores, customer types.

DO: Provide Context

Include all relevant information in the interrupt payload so reviewers can decide quickly.

DON'T: Interrupt Everything

Too many interrupts cause alert fatigue. Reserve for genuinely important decisions.

DON'T: Forget Timeouts

Implement TTLs for interrupted workflows. Auto-escalate or auto-reject if not reviewed.

Key Takeaways

  • interrupt() pauses execution and saves state until human responds
  • Command(resume=...) continues the workflow with human input
  • Same thread_id is required to resume the correct workflow
  • Checkpointer required - interrupt needs persistence to save state

Next up: Production Patterns - Error handling, retries, observability, and deployment best practices.

🧠 Quick Quiz

Test your understanding of this lesson.

1

What does LangGraph's interrupt() function do?

2

Which e-commerce actions should typically require human approval?

3

How do you resume a paused LangGraph workflow after human approval?

Inventory Management Agent