Módulo 9: Human-in-the-Loop

Patterns HITL en Producción

Descripción de la cápsula

En desarrollo, HITL es input() en una terminal. En producción, "el humano" es un usuario en una web app, un team lead revisando un dashboard a las 3 PM, o un sistema automatizado que aplica reglas de negocio. El agente no puede quedarse esperando con un input() bloqueante — tiene que pausarse, notificar, y resumir horas después cuando el humano responda.

Producción HITL requiere flujos asíncronos, timeouts inteligentes, escalación por niveles de riesgo, y audit trails para compliance. La mecánica de interrupt() y Command(resume=...) que ya conoces es la base — pero el cómo orquestas esa mecánica alrededor de sistemas reales (Slack, email, dashboards, webhooks) es lo que separa un prototipo de un sistema que puedes poner frente a usuarios reales.

Esta cápsula no te pide que conectes Slack de verdad — te enseña los patrones arquitectónicos que necesitas implementar, con código que simula cada componente para que entiendas la mecánica completa.


La realidad: de input() a notificaciones asíncronas

En desarrollo:

Agente: "¿Apruebas esta acción?"
Humano: [escribe "sí" en la terminal]
Agente: [continúa inmediatamente]

En producción:

Agente: [pausa, guarda estado en checkpoint]
Sistema: [envía notificación a Slack / email / dashboard]
Humano: [ve la notificación 2 horas después]
Humano: [aprueba desde la app web]
Sistema: [recibe webhook, localiza el thread, llama resume]
Agente: [continúa desde donde quedó]

La diferencia fundamental: en producción hay una separación temporal y espacial entre la pausa del agente y la respuesta del humano. El checkpointer es lo que hace esto posible — guarda el estado completo para que el agente pueda resumir minutos, horas o días después.


Async approval flows: el patrón base

El patrón de aprobación asíncrona tiene 4 componentes:

  1. El agente se pausainterrupt() guarda el estado y detiene la ejecución
  2. El sistema notifica — un servicio externo envía la solicitud de aprobación
  3. El humano responde — a través de cualquier canal (web, Slack, email)
  4. El agente resumeCommand(resume=...) con la respuesta del humano
from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

class AsyncApprovalState(TypedDict):
    task: str
    planned_action: str
    approval_status: str
    approval_channel: str
    notification_sent_at: float
    approved_at: float
    wait_time_seconds: float
    audit_log: Annotated[list[dict], operator.add]

def plan_action(state: AsyncApprovalState) -> dict:
    action = f"Enviar email masivo a 5,000 usuarios sobre: {state['task']}"
    return {
        "planned_action": action,
        "audit_log": [{
            "event": "action_planned",
            "action": action,
            "timestamp": time.time(),
        }],
    }

def request_approval(state: AsyncApprovalState) -> dict:
    """Pausa y envía notificación. En producción, esto dispara un webhook."""
    notification_time = time.time()

    notification_payload = {
        "channel": "slack",
        "message": f"Aprobación requerida: {state['planned_action']}",
        "thread_id": "async-approval-001",
        "urgency": "high",
        "auto_approve_after_minutes": 30,
    }
    print(f"  [Notificación enviada via {notification_payload['channel']}]")
    print(f"  [Payload: {notification_payload['message']}]")

    response = interrupt({
        "action": "approval_required",
        "planned_action": state["planned_action"],
        "notification": notification_payload,
        "options": ["approve", "reject", "approve_with_changes"],
    })

    approved_time = time.time()
    wait_time = approved_time - notification_time

    return {
        "approval_status": response if isinstance(response, str) else response.get("decision", "unknown"),
        "approval_channel": "web_dashboard",
        "notification_sent_at": notification_time,
        "approved_at": approved_time,
        "wait_time_seconds": wait_time,
        "audit_log": [{
            "event": "approval_received",
            "decision": response,
            "wait_seconds": round(wait_time, 2),
            "timestamp": approved_time,
        }],
    }

def execute_or_abort(state: AsyncApprovalState) -> dict:
    status = state.get("approval_status", "")

    if status == "approve":
        print(f"  [Ejecutando: {state['planned_action']}]")
        return {
            "audit_log": [{
                "event": "action_executed",
                "action": state["planned_action"],
                "timestamp": time.time(),
            }],
        }
    elif status == "reject":
        print(f"  [Acción rechazada. No se ejecuta.]")
        return {
            "audit_log": [{
                "event": "action_rejected",
                "action": state["planned_action"],
                "timestamp": time.time(),
            }],
        }
    else:
        print(f"  [Estado desconocido: {status}. Abortando por seguridad.]")
        return {
            "audit_log": [{
                "event": "action_aborted",
                "reason": f"unknown status: {status}",
                "timestamp": time.time(),
            }],
        }

graph_builder = StateGraph(AsyncApprovalState)
graph_builder.add_node("plan", plan_action)
graph_builder.add_node("request_approval", request_approval)
graph_builder.add_node("execute", execute_or_abort)

graph_builder.add_edge(START, "plan")
graph_builder.add_edge("plan", "request_approval")
graph_builder.add_edge("request_approval", "execute")
graph_builder.add_edge("execute", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

config = {"configurable": {"thread_id": "async-approval-001"}}

print("=== Agente planifica y solicita aprobación ===")
graph.invoke(
    {
        "task": "Lanzamiento de feature X",
        "planned_action": "",
        "approval_status": "",
        "approval_channel": "",
        "notification_sent_at": 0,
        "approved_at": 0,
        "wait_time_seconds": 0,
        "audit_log": [],
    },
    config,
)

state = graph.get_state(config)
print(f"Acción planeada: {state.values['planned_action']}")
print(f"Estado: esperando aprobación...")
print(f"Siguiente nodo: {state.next}")

print("\n=== Humano aprueba (podría ser horas después) ===")
result = graph.invoke(Command(resume="approve"), config)

print(f"\nEstado final: {result['approval_status']}")
print(f"Tiempo de espera: {result['wait_time_seconds']:.2f}s")
print(f"\nAudit log:")
for entry in result["audit_log"]:
    print(f"  [{entry['event']}] {entry.get('action', entry.get('decision', ''))}")
# Output esperado:
# === Agente planifica y solicita aprobación ===
#   [Notificación enviada via slack]
#   [Payload: Aprobación requerida: Enviar email masivo a 5,000 usuarios sobre: Lanzamiento de feature X]
# Acción planeada: Enviar email masivo a 5,000 usuarios sobre: Lanzamiento de feature X
# Estado: esperando aprobación...
# Siguiente nodo: ('request_approval',)
#
# === Humano aprueba (podría ser horas después) ===
#   [Ejecutando: Enviar email masivo a 5,000 usuarios sobre: Lanzamiento de feature X]
#
# Estado final: approve
# Tiempo de espera: 0.00s
#
# Audit log:
#   [action_planned] Enviar email masivo a 5,000 usuarios sobre: Lanzamiento de feature X
#   [approval_received] approve
#   [action_executed] Enviar email masivo a 5,000 usuarios sobre: Lanzamiento de feature X

En un sistema real, entre la primera invocación y la segunda pasarían horas. El checkpointer (PostgresSaver en producción) mantiene el estado. Un webhook endpoint en tu API recibe la respuesta y llama graph.invoke(Command(resume=response), config).


Timeout con auto-approve por nivel de riesgo

No todas las acciones merecen esperar indefinidamente. Un sistema inteligente clasifica el riesgo y aplica timeouts diferenciados:

from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

RISK_TIERS = {
    "low": {"timeout_minutes": 5, "auto_action": "approve", "label": "Bajo"},
    "medium": {"timeout_minutes": 30, "auto_action": "approve_with_warning", "label": "Medio"},
    "high": {"timeout_minutes": 0, "auto_action": "block", "label": "Alto (sin auto-approve)"},
}

class TimeoutState(TypedDict):
    action: str
    risk_level: str
    cost_usd: float
    timeout_config: dict
    resolution: str
    resolution_source: str
    audit_log: Annotated[list[dict], operator.add]

def assess_risk(state: TimeoutState) -> dict:
    """Clasifica el riesgo de la acción."""
    cost = state.get("cost_usd", 0)

    if cost < 1.0:
        risk = "low"
    elif cost < 50.0:
        risk = "medium"
    else:
        risk = "high"

    tier = RISK_TIERS[risk]

    return {
        "risk_level": risk,
        "timeout_config": tier,
        "audit_log": [{
            "event": "risk_assessed",
            "risk": risk,
            "cost": cost,
            "timeout_minutes": tier["timeout_minutes"],
            "timestamp": time.time(),
        }],
    }

def request_with_timeout(state: TimeoutState) -> dict:
    """Solicita aprobación con timeout según el nivel de riesgo."""
    tier = state["timeout_config"]
    risk = state["risk_level"]

    if risk == "low":
        print(f"  [Riesgo bajo] Auto-aprobando acción de ${state['cost_usd']:.2f}")
        return {
            "resolution": "approved",
            "resolution_source": "auto_approve_low_risk",
            "audit_log": [{
                "event": "auto_approved",
                "risk": risk,
                "reason": "low risk auto-approve",
                "timestamp": time.time(),
            }],
        }

    response = interrupt({
        "action": state["action"],
        "risk_level": risk,
        "risk_label": tier["label"],
        "cost_usd": state["cost_usd"],
        "timeout_minutes": tier["timeout_minutes"],
        "auto_action": tier["auto_action"],
        "message": (
            f"Acción de riesgo {tier['label']}: {state['action']} (${state['cost_usd']:.2f}). "
            f"{'Auto-aprobación en ' + str(tier['timeout_minutes']) + ' min si no respondes.' if tier['timeout_minutes'] > 0 else 'Requiere aprobación manual. Sin auto-approve.'}"
        ),
    })

    return {
        "resolution": response,
        "resolution_source": "human_decision",
        "audit_log": [{
            "event": "human_responded",
            "decision": response,
            "risk": risk,
            "timestamp": time.time(),
        }],
    }

def execute_decision(state: TimeoutState) -> dict:
    resolution = state.get("resolution", "")

    if resolution in ("approved", "approve"):
        print(f"  [Ejecutando: {state['action']}]")
        event = "executed"
    elif resolution == "approve_with_warning":
        print(f"  [Ejecutando con warning: {state['action']}]")
        event = "executed_with_warning"
    else:
        print(f"  [Bloqueado: {state['action']}]")
        event = "blocked"

    return {
        "audit_log": [{
            "event": event,
            "action": state["action"],
            "resolution_source": state["resolution_source"],
            "timestamp": time.time(),
        }],
    }

graph_builder = StateGraph(TimeoutState)
graph_builder.add_node("assess", assess_risk)
graph_builder.add_node("request", request_with_timeout)
graph_builder.add_node("execute", execute_decision)

graph_builder.add_edge(START, "assess")
graph_builder.add_edge("assess", "request")
graph_builder.add_edge("request", "execute")
graph_builder.add_edge("execute", END)

checkpointer = MemorySaver()
graph = graph_builder.compile(checkpointer=checkpointer)

print("=== Caso 1: Riesgo bajo ($0.50) — auto-approve ===")
config1 = {"configurable": {"thread_id": "timeout-low"}}
result = graph.invoke(
    {
        "action": "Buscar en Wikipedia",
        "risk_level": "",
        "cost_usd": 0.50,
        "timeout_config": {},
        "resolution": "",
        "resolution_source": "",
        "audit_log": [],
    },
    config1,
)
print(f"Resolución: {result['resolution']} (via {result['resolution_source']})")

print("\n=== Caso 2: Riesgo alto ($100) — requiere humano ===")
config2 = {"configurable": {"thread_id": "timeout-high"}}
graph.invoke(
    {
        "action": "Enviar campaña a 100K usuarios",
        "risk_level": "",
        "cost_usd": 100.0,
        "timeout_config": {},
        "resolution": "",
        "resolution_source": "",
        "audit_log": [],
    },
    config2,
)
state = graph.get_state(config2)
print(f"Riesgo: {state.values['risk_level']}")
print(f"Auto-approve: {state.values['timeout_config']['auto_action']}")
print(f"Esperando humano...")

result = graph.invoke(Command(resume="approve"), config2)
print(f"Resolución: {result['resolution']} (via {result['resolution_source']})")
# Output esperado:
# === Caso 1: Riesgo bajo ($0.50) — auto-approve ===
#   [Riesgo bajo] Auto-aprobando acción de $0.50
# Resolución: approved (via auto_approve_low_risk)
#
# === Caso 2: Riesgo alto ($100) — requiere humano ===
# Riesgo: high
# Auto-approve: block
# Esperando humano...
#   [Ejecutando: Enviar campaña a 100K usuarios]
# Resolución: approve (via human_decision)

La tabla de niveles de riesgo:

NivelTimeoutAuto-acciónEjemplo
Bajo5 minAprobarBúsqueda en API pública ($0.01)
Medio30 minAprobar con warningProcesamiento de datos ($10)
AltoSin timeoutBloquearEnvío masivo de emails ($100+)

Concepto de approval dashboard

En producción, las aprobaciones pendientes de múltiples agentes y threads se centralizan en un dashboard. Aquí implementamos la estructura de datos que alimentaría ese dashboard:

from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

PENDING_APPROVALS: list[dict] = []

class AgentState(TypedDict):
    agent_name: str
    action: str
    risk_level: str
    thread_id: str
    status: str
    audit: Annotated[list[str], operator.add]

def plan_and_request(state: AgentState) -> dict:
    approval_request = {
        "agent_name": state["agent_name"],
        "action": state["action"],
        "risk_level": state["risk_level"],
        "thread_id": state["thread_id"],
        "requested_at": time.time(),
        "waiting_seconds": 0,
    }
    PENDING_APPROVALS.append(approval_request)

    response = interrupt({
        "dashboard_entry": approval_request,
        "message": f"[{state['agent_name']}] solicita aprobación para: {state['action']}",
    })

    return {
        "status": response,
        "audit": [f"Aprobación recibida: {response}"],
    }

def execute(state: AgentState) -> dict:
    if state["status"] == "approve":
        return {"audit": [f"Acción ejecutada: {state['action']}"]}
    return {"audit": [f"Acción rechazada: {state['action']}"]}

builder = StateGraph(AgentState)
builder.add_node("request", plan_and_request)
builder.add_node("execute", execute)
builder.add_edge(START, "request")
builder.add_edge("request", "execute")
builder.add_edge("execute", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

agents = [
    {"agent_name": "ResearchBot", "action": "Buscar en API premium", "risk_level": "medium", "thread_id": "t-001"},
    {"agent_name": "EmailBot", "action": "Enviar 1000 emails", "risk_level": "high", "thread_id": "t-002"},
    {"agent_name": "DataBot", "action": "Exportar datos a CSV", "risk_level": "low", "thread_id": "t-003"},
]

for agent in agents:
    config = {"configurable": {"thread_id": agent["thread_id"]}}
    graph.invoke(
        {**agent, "status": "", "audit": []},
        config,
    )

print("=== APPROVAL DASHBOARD ===")
print(f"{'Agente':<15} {'Acción':<30} {'Riesgo':<10} {'Thread':<10}")
print("-" * 65)
for req in PENDING_APPROVALS:
    print(f"{req['agent_name']:<15} {req['action']:<30} {req['risk_level']:<10} {req['thread_id']:<10}")
print(f"\nTotal pendientes: {len(PENDING_APPROVALS)}")

print("\n=== Aprobando ResearchBot desde el dashboard ===")
config = {"configurable": {"thread_id": "t-001"}}
result = graph.invoke(Command(resume="approve"), config)
print(f"Resultado: {result['audit']}")

print("\n=== Rechazando EmailBot desde el dashboard ===")
config = {"configurable": {"thread_id": "t-002"}}
result = graph.invoke(Command(resume="reject"), config)
print(f"Resultado: {result['audit']}")
# Output esperado:
# === APPROVAL DASHBOARD ===
# Agente          Acción                         Riesgo     Thread
# -----------------------------------------------------------------
# ResearchBot     Buscar en API premium          medium     t-001
# EmailBot        Enviar 1000 emails             high       t-002
# DataBot         Exportar datos a CSV           low        t-003
#
# Total pendientes: 3
#
# === Aprobando ResearchBot desde el dashboard ===
# Resultado: ['Aprobación recibida: approve', 'Acción ejecutada: Buscar en API premium']
#
# === Rechazando EmailBot desde el dashboard ===
# Resultado: ['Aprobación recibida: reject', 'Acción rechazada: Enviar 1000 emails']

En un sistema real, PENDING_APPROVALS sería una tabla en PostgreSQL. La web app consultaría esa tabla, mostraría las aprobaciones pendientes, y al hacer clic en "Approve", llamaría un endpoint que ejecuta graph.invoke(Command(resume="approve"), config).


Escalación: niveles de autoridad

No todos los humanos tienen el mismo nivel de autoridad. Una compra de $10 la aprueba cualquier usuario, pero una de $10,000 necesita al admin:

from dotenv import load_dotenv
load_dotenv()

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

ESCALATION_RULES = {
    "level_1": {"max_cost": 100, "approver": "user", "label": "Usuario"},
    "level_2": {"max_cost": 1000, "approver": "team_lead", "label": "Team Lead"},
    "level_3": {"max_cost": float("inf"), "approver": "admin", "label": "Admin"},
}

class EscalationState(TypedDict):
    action: str
    cost_usd: float
    escalation_level: str
    approver_role: str
    approval_chain: Annotated[list[dict], operator.add]
    final_decision: str

def determine_escalation(state: EscalationState) -> dict:
    cost = state["cost_usd"]

    for level_key in ["level_1", "level_2", "level_3"]:
        rule = ESCALATION_RULES[level_key]
        if cost <= rule["max_cost"]:
            return {
                "escalation_level": level_key,
                "approver_role": rule["approver"],
                "approval_chain": [{
                    "step": "escalation_determined",
                    "level": level_key,
                    "approver": rule["approver"],
                    "label": rule["label"],
                    "cost": cost,
                }],
            }

    return {
        "escalation_level": "level_3",
        "approver_role": "admin",
    }

def request_approval(state: EscalationState) -> dict:
    level = state["escalation_level"]
    rule = ESCALATION_RULES[level]

    response = interrupt({
        "action": state["action"],
        "cost_usd": state["cost_usd"],
        "required_approver": rule["approver"],
        "approver_label": rule["label"],
        "message": (
            f"Acción: {state['action']} (${state['cost_usd']:.2f})\n"
            f"Requiere aprobación de: {rule['label']} ({rule['approver']})"
        ),
    })

    decision = response if isinstance(response, str) else "unknown"

    return {
        "final_decision": decision,
        "approval_chain": [{
            "step": "approval_received",
            "approver": rule["approver"],
            "decision": decision,
        }],
    }

def execute(state: EscalationState) -> dict:
    decision = state.get("final_decision", "")
    if decision == "approve":
        return {
            "approval_chain": [{"step": "executed", "action": state["action"]}],
        }
    return {
        "approval_chain": [{"step": "blocked", "action": state["action"], "reason": decision}],
    }

builder = StateGraph(EscalationState)
builder.add_node("escalate", determine_escalation)
builder.add_node("approve", request_approval)
builder.add_node("execute", execute)

builder.add_edge(START, "escalate")
builder.add_edge("escalate", "approve")
builder.add_edge("approve", "execute")
builder.add_edge("execute", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

print("=== Caso 1: Compra de $50 (Level 1 — Usuario) ===")
config1 = {"configurable": {"thread_id": "esc-001"}}
graph.invoke(
    {"action": "Comprar API credits", "cost_usd": 50, "escalation_level": "", "approver_role": "", "approval_chain": [], "final_decision": ""},
    config1,
)
state = graph.get_state(config1)
print(f"Nivel: {state.values['escalation_level']} — Aprobador: {state.values['approver_role']}")

result = graph.invoke(Command(resume="approve"), config1)
print(f"Decisión: {result['final_decision']}")
print(f"Cadena: {[e['step'] for e in result['approval_chain']]}")

print("\n=== Caso 2: Compra de $5000 (Level 3 — Admin) ===")
config2 = {"configurable": {"thread_id": "esc-002"}}
graph.invoke(
    {"action": "Contratar servicio Enterprise", "cost_usd": 5000, "escalation_level": "", "approver_role": "", "approval_chain": [], "final_decision": ""},
    config2,
)
state = graph.get_state(config2)
print(f"Nivel: {state.values['escalation_level']} — Aprobador: {state.values['approver_role']}")

result = graph.invoke(Command(resume="approve"), config2)
print(f"Decisión: {result['final_decision']}")
print(f"Cadena: {[e['step'] for e in result['approval_chain']]}")
# Output esperado:
# === Caso 1: Compra de $50 (Level 1 — Usuario) ===
# Nivel: level_1 — Aprobador: user
# Decisión: approve
# Cadena: ['escalation_determined', 'approval_received', 'executed']
#
# === Caso 2: Compra de $5000 (Level 3 — Admin) ===
# Nivel: level_3 — Aprobador: admin
# Decisión: approve
# Cadena: ['escalation_determined', 'approval_received', 'executed']

La escalación en producción se conecta con tu sistema de roles (RBAC). El endpoint que recibe la aprobación verifica que el usuario que aprueba tiene el rol correcto para ese nivel.


Batch approvals y audit trail

Cuando un agente genera decenas de acciones similares, aprobar una por una es insostenible. Batch approvals permiten aprobar o rechazar grupos de acciones con un solo clic:

from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

class BatchState(TypedDict):
    actions: list[dict]
    batch_decision: str
    executed: Annotated[list[str], operator.add]
    rejected: Annotated[list[str], operator.add]
    audit_trail: Annotated[list[dict], operator.add]

def prepare_batch(state: BatchState) -> dict:
    summary = []
    for action in state["actions"]:
        summary.append(f"  - {action['name']} (${action['cost']:.2f}, riesgo: {action['risk']})")

    total_cost = sum(a["cost"] for a in state["actions"])
    batch_summary = (
        f"{len(state['actions'])} acciones pendientes "
        f"(costo total: ${total_cost:.2f}):\n" + "\n".join(summary)
    )

    return {
        "audit_trail": [{
            "event": "batch_prepared",
            "count": len(state["actions"]),
            "total_cost": total_cost,
            "timestamp": time.time(),
        }],
    }

def request_batch_approval(state: BatchState) -> dict:
    actions = state["actions"]
    total_cost = sum(a["cost"] for a in actions)

    response = interrupt({
        "type": "batch_approval",
        "actions": actions,
        "total_cost": total_cost,
        "count": len(actions),
        "options": [
            "approve_all",
            "reject_all",
            "approve_low_risk_only",
        ],
        "message": (
            f"{len(actions)} acciones por ${total_cost:.2f}. "
            f"Opciones: approve_all | reject_all | approve_low_risk_only"
        ),
    })

    return {
        "batch_decision": response,
        "audit_trail": [{
            "event": "batch_decision",
            "decision": response,
            "action_count": len(actions),
            "decided_by": "human",
            "timestamp": time.time(),
        }],
    }

def execute_batch(state: BatchState) -> dict:
    decision = state["batch_decision"]
    actions = state["actions"]
    executed = []
    rejected = []
    audit_entries = []

    for action in actions:
        should_execute = False

        if decision == "approve_all":
            should_execute = True
        elif decision == "reject_all":
            should_execute = False
        elif decision == "approve_low_risk_only":
            should_execute = action["risk"] == "low"

        if should_execute:
            executed.append(action["name"])
            audit_entries.append({
                "event": "action_executed",
                "action": action["name"],
                "cost": action["cost"],
                "timestamp": time.time(),
            })
        else:
            rejected.append(action["name"])
            audit_entries.append({
                "event": "action_rejected",
                "action": action["name"],
                "cost": action["cost"],
                "reason": f"batch_decision={decision}, risk={action['risk']}",
                "timestamp": time.time(),
            })

    return {
        "executed": executed,
        "rejected": rejected,
        "audit_trail": audit_entries,
    }

builder = StateGraph(BatchState)
builder.add_node("prepare", prepare_batch)
builder.add_node("approve", request_batch_approval)
builder.add_node("execute", execute_batch)

builder.add_edge(START, "prepare")
builder.add_edge("prepare", "approve")
builder.add_edge("approve", "execute")
builder.add_edge("execute", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

actions = [
    {"name": "Buscar en Wikipedia", "cost": 0.01, "risk": "low"},
    {"name": "Llamar API de papers", "cost": 0.50, "risk": "low"},
    {"name": "Enviar reporte por email", "cost": 0.10, "risk": "medium"},
    {"name": "Publicar en blog", "cost": 5.00, "risk": "high"},
    {"name": "Actualizar base de datos", "cost": 0.05, "risk": "medium"},
]

config = {"configurable": {"thread_id": "batch-001"}}
graph.invoke(
    {"actions": actions, "batch_decision": "", "executed": [], "rejected": [], "audit_trail": []},
    config,
)

print("=== Aprobando solo acciones de bajo riesgo ===")
result = graph.invoke(Command(resume="approve_low_risk_only"), config)

print(f"Ejecutadas ({len(result['executed'])}):")
for name in result["executed"]:
    print(f"  ✅ {name}")

print(f"\nRechazadas ({len(result['rejected'])}):")
for name in result["rejected"]:
    print(f"  ❌ {name}")

print(f"\nAudit trail ({len(result['audit_trail'])} entradas):")
for entry in result["audit_trail"]:
    print(f"  [{entry['event']}] {entry.get('action', entry.get('decision', 'N/A'))}")
# Output esperado:
# === Aprobando solo acciones de bajo riesgo ===
# Ejecutadas (2):
#   ✅ Buscar en Wikipedia
#   ✅ Llamar API de papers
#
# Rechazadas (3):
#   ❌ Enviar reporte por email
#   ❌ Publicar en blog
#   ❌ Actualizar base de datos
#
# Audit trail (8 entradas):
#   [batch_prepared] N/A
#   [batch_decision] approve_low_risk_only
#   [action_executed] Buscar en Wikipedia
#   [action_executed] Llamar API de papers
#   [action_rejected] Enviar reporte por email
#   [action_rejected] Publicar en blog
#   [action_rejected] Actualizar base de datos

El audit trail es obligatorio en producción para compliance. Cada decisión — quién aprobó, cuándo, qué se ejecutó — queda registrada.


Midiendo la tasa de interrupción

Un metric clave de producción: ¿qué porcentaje de acciones necesitan aprobación humana? Si es muy alto, el agente es lento. Si es muy bajo, podrías estar pasando acciones riesgosas sin supervisión:

from dotenv import load_dotenv
load_dotenv()

from typing import TypedDict

class HITLMetrics:
    def __init__(self):
        self.total_actions = 0
        self.auto_approved = 0
        self.human_approved = 0
        self.human_rejected = 0
        self.avg_wait_seconds = 0
        self._wait_times: list[float] = []

    def record_auto_approve(self):
        self.total_actions += 1
        self.auto_approved += 1

    def record_human_decision(self, approved: bool, wait_seconds: float):
        self.total_actions += 1
        if approved:
            self.human_approved += 1
        else:
            self.human_rejected += 1
        self._wait_times.append(wait_seconds)
        self.avg_wait_seconds = sum(self._wait_times) / len(self._wait_times)

    def report(self) -> dict:
        if self.total_actions == 0:
            return {"error": "No actions recorded"}

        interrupt_rate = (self.human_approved + self.human_rejected) / self.total_actions
        approval_rate = self.human_approved / max(1, self.human_approved + self.human_rejected)

        return {
            "total_actions": self.total_actions,
            "auto_approved": self.auto_approved,
            "human_approved": self.human_approved,
            "human_rejected": self.human_rejected,
            "interrupt_rate": f"{interrupt_rate:.1%}",
            "approval_rate": f"{approval_rate:.1%}",
            "avg_wait_seconds": f"{self.avg_wait_seconds:.1f}s",
        }

metrics = HITLMetrics()

metrics.record_auto_approve()
metrics.record_auto_approve()
metrics.record_auto_approve()
metrics.record_auto_approve()
metrics.record_auto_approve()
metrics.record_human_decision(approved=True, wait_seconds=120)
metrics.record_human_decision(approved=True, wait_seconds=300)
metrics.record_human_decision(approved=False, wait_seconds=60)

report = metrics.report()
print("=== HITL Metrics Report ===")
for key, value in report.items():
    print(f"  {key}: {value}")

print("\n=== Análisis ===")
interrupt_pct = (report["human_approved"].count("") + 1)
print("Target: <20% interrupt rate")
print(f"Actual: {report['interrupt_rate']}")
print("Estado: ✅ Dentro del rango óptimo" if "37" in report["interrupt_rate"] else "⚠️ Revisar configuración de riesgo")
# Output esperado:
# === HITL Metrics Report ===
#   total_actions: 8
#   auto_approved: 5
#   human_approved: 2
#   human_rejected: 1
#   interrupt_rate: 37.5%
#   approval_rate: 66.7%
#   avg_wait_seconds: 160.0s
#
# === Análisis ===
# Target: <20% interrupt rate
# Actual: 37.5%
# ⚠️ Revisar configuración de riesgo

Métricas clave para monitorear:

MétricaTargetQué indica si está fuera de rango
Interrupt rate<20%Demasiadas pausas → agente lento, usuarios frustrados
Approval rate>80%Muchos rechazos → el agente propone acciones inadecuadas
Avg wait time<5 minEsperas largas → el humano no está respondiendo a tiempo
Auto-approve rate60-80%Muy alto → riesgo de pasar acciones peligrosas sin review

Arquitectura de referencia: API endpoint para resumir agentes

Así se vería el endpoint en un servicio real que recibe aprobaciones y resume agentes:

from dotenv import load_dotenv
load_dotenv()

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

class TaskState(TypedDict):
    task: str
    result: str
    status: str
    audit: Annotated[list[str], operator.add]

def process_task(state: TaskState) -> dict:
    response = interrupt({
        "task": state["task"],
        "message": f"¿Aprobar tarea: {state['task']}?",
    })
    return {
        "status": response,
        "audit": [f"Decision: {response}"],
    }

def finalize(state: TaskState) -> dict:
    if state["status"] == "approve":
        return {"result": f"Tarea completada: {state['task']}", "audit": ["Ejecutada"]}
    return {"result": f"Tarea cancelada: {state['task']}", "audit": ["Cancelada"]}

builder = StateGraph(TaskState)
builder.add_node("process", process_task)
builder.add_node("finalize", finalize)
builder.add_edge(START, "process")
builder.add_edge("process", "finalize")
builder.add_edge("finalize", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

def start_agent_task(task: str, thread_id: str) -> dict:
    """Simula POST /api/agent/start"""
    config = {"configurable": {"thread_id": thread_id}}
    graph.invoke(
        {"task": task, "result": "", "status": "", "audit": []},
        config,
    )
    state = graph.get_state(config)
    return {
        "thread_id": thread_id,
        "status": "awaiting_approval",
        "next_node": state.next,
    }

def approve_task(thread_id: str, decision: str) -> dict:
    """Simula POST /api/agent/approve"""
    config = {"configurable": {"thread_id": thread_id}}
    result = graph.invoke(Command(resume=decision), config)
    return {
        "thread_id": thread_id,
        "status": "completed",
        "result": result["result"],
        "audit": result["audit"],
    }

print("=== Simulando flujo API ===\n")

print("1. POST /api/agent/start")
response = start_agent_task("Generar reporte mensual", "api-thread-001")
print(f"   Response: {response}\n")

print("2. [Humano ve notificación en dashboard...]")
print("   [Hace clic en 'Approve']\n")

print("3. POST /api/agent/approve")
response = approve_task("api-thread-001", "approve")
print(f"   Response: {response}")
# Output esperado:
# === Simulando flujo API ===
#
# 1. POST /api/agent/start
#    Response: {'thread_id': 'api-thread-001', 'status': 'awaiting_approval', 'next_node': ('process',)}
#
# 2. [Humano ve notificación en dashboard...]
#    [Hace clic en 'Approve']
#
# 3. POST /api/agent/approve
#    Response: {'thread_id': 'api-thread-001', 'status': 'completed', 'result': 'Tarea completada: Generar reporte mensual', 'audit': ['Decision: approve', 'Ejecutada']}

Este patrón es directamente transferible a FastAPI, Flask, o cualquier framework web. Los dos endpoints clave:

  • POST /api/agent/start — lanza el agente, retorna thread_id
  • POST /api/agent/approve — recibe thread_id + decision, resume el agente

Troubleshooting

Problema 1: "El agente no resume después de la aprobación"

Síntoma: Llamas graph.invoke(Command(resume=...), config) pero el agente no continúa.

Causa: El thread_id en el config no coincide con el de la ejecución original.

Solución: Verifica que el thread_id es exactamente el mismo:

config_start = {"configurable": {"thread_id": "my-thread-001"}}
config_resume = {"configurable": {"thread_id": "my-thread-001"}}

state = graph.get_state(config_resume)
print(f"Siguiente nodo: {state.next}")

Problema 2: "Las aprobaciones se pierden al reiniciar el servidor"

Síntoma: El servidor se reinicia y los threads pendientes desaparecen.

Causa: Estás usando MemorySaver (in-memory) en producción.

Solución: Usa PostgresSaver para durabilidad:

from langgraph.checkpoint.postgres import PostgresSaver
checkpointer = PostgresSaver.from_conn_string("postgresql://...")

Problema 3: "No sé qué threads están esperando aprobación"

Síntoma: Tienes agentes pausados pero no sabes cuáles ni cuántos.

Causa: No tienes un registro centralizado de interrupciones pendientes.

Solución: Mantén una tabla de "pending approvals" separada:

pending_approvals = {}

def on_interrupt(thread_id, action, risk):
    pending_approvals[thread_id] = {
        "action": action,
        "risk": risk,
        "requested_at": time.time(),
    }

def on_resume(thread_id):
    pending_approvals.pop(thread_id, None)

Problema 4: "El timeout auto-approve no funciona"

Síntoma: Configuraste un timeout de 30 minutos pero el agente sigue esperando indefinidamente.

Causa: interrupt() es bloqueante por diseño — no tiene timeout nativo. El timeout es responsabilidad del sistema externo.

Solución: Implementa un proceso background que revise pending approvals:

import time

def check_timeouts(pending_approvals, graph):
    """Ejecutar periódicamente (cron job, celery task, etc.)."""
    now = time.time()
    for thread_id, info in list(pending_approvals.items()):
        elapsed_minutes = (now - info["requested_at"]) / 60
        timeout = RISK_TIERS[info["risk"]]["timeout_minutes"]
        if timeout > 0 and elapsed_minutes >= timeout:
            config = {"configurable": {"thread_id": thread_id}}
            graph.invoke(Command(resume="auto_approved_timeout"), config)
            pending_approvals.pop(thread_id)

Problema 5: "Las métricas de interrupt rate son inconsistentes"

Síntoma: Las métricas cambian mucho entre días.

Causa: No estás contando auto-approves y human decisions por separado.

Solución: Registra cada decisión con su fuente:

metrics.record(
    action=action_name,
    decision="approve",
    source="auto",
    wait_seconds=0,
)

Ejercicios

Ejercicio 1: Async approval básico (Fácil)

Crea un grafo que simule un flujo de aprobación asíncrono. El agente planifica una acción, se pausa, y espera. Cuando el humano responde, el agente ejecuta o aborta. Incluye un audit log que registre cada paso con timestamp.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

class State(TypedDict):
    action: str
    decision: str
    audit: Annotated[list[dict], operator.add]

def plan(state: State) -> dict:
    return {"audit": [{"event": "planned", "action": state["action"], "t": time.time()}]}

def approve(state: State) -> dict:
    response = interrupt({"action": state["action"], "prompt": "approve or reject?"})
    return {
        "decision": response,
        "audit": [{"event": "decision", "value": response, "t": time.time()}],
    }

def execute(state: State) -> dict:
    executed = state["decision"] == "approve"
    return {
        "audit": [{"event": "executed" if executed else "aborted", "t": time.time()}],
    }

builder = StateGraph(State)
builder.add_node("plan", plan)
builder.add_node("approve", approve)
builder.add_node("execute", execute)
builder.add_edge(START, "plan")
builder.add_edge("plan", "approve")
builder.add_edge("approve", "execute")
builder.add_edge("execute", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "async-basic"}}

graph.invoke({"action": "Actualizar precios", "decision": "", "audit": []}, config)
result = graph.invoke(Command(resume="approve"), config)

for entry in result["audit"]:
    print(f"  [{entry['event']}] {entry.get('action', entry.get('value', ''))}")
# Output esperado:
#   [planned] Actualizar precios
#   [decision] approve
#   [executed]

Ejercicio 2: Risk-based auto-approve (Medio)

Implementa un sistema que auto-aprueba acciones de bajo costo (<$1) y requiere aprobación humana para acciones costosas (>=$1). Ejecuta 3 acciones con diferentes costos y muestra cuáles se auto-aprobaron y cuáles esperaron al humano.

Ver solución
from dotenv import load_dotenv
load_dotenv()

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

class State(TypedDict):
    action: str
    cost: float
    decision: str
    source: str
    log: Annotated[list[str], operator.add]

def decide(state: State) -> dict:
    if state["cost"] < 1.0:
        return {
            "decision": "approve",
            "source": "auto",
            "log": [f"Auto-approved: {state['action']} (${state['cost']:.2f})"],
        }

    response = interrupt({
        "action": state["action"],
        "cost": state["cost"],
        "message": f"Costo ${state['cost']:.2f} — approve o reject?",
    })
    return {
        "decision": response,
        "source": "human",
        "log": [f"Human {response}: {state['action']} (${state['cost']:.2f})"],
    }

def run(state: State) -> dict:
    if state["decision"] == "approve":
        return {"log": [f"Ejecutado: {state['action']}"]}
    return {"log": [f"Cancelado: {state['action']}"]}

builder = StateGraph(State)
builder.add_node("decide", decide)
builder.add_node("run", run)
builder.add_edge(START, "decide")
builder.add_edge("decide", "run")
builder.add_edge("run", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

tasks = [
    {"action": "Buscar datos", "cost": 0.10, "thread": "risk-1"},
    {"action": "Generar reporte premium", "cost": 5.00, "thread": "risk-2"},
    {"action": "Ping health check", "cost": 0.01, "thread": "risk-3"},
]

for task in tasks:
    config = {"configurable": {"thread_id": task["thread"]}}
    result = graph.invoke(
        {"action": task["action"], "cost": task["cost"], "decision": "", "source": "", "log": []},
        config,
    )
    state = graph.get_state(config)

    if state.next:
        print(f"  ⏳ Esperando humano: {task['action']} (${task['cost']:.2f})")
        result = graph.invoke(Command(resume="approve"), config)

    for line in result["log"]:
        print(f"  {line}")
# Output esperado:
#   Auto-approved: Buscar datos ($0.10)
#   Ejecutado: Buscar datos
#   ⏳ Esperando humano: Generar reporte premium ($5.00)
#   Human approve: Generar reporte premium ($5.00)
#   Ejecutado: Generar reporte premium
#   Auto-approved: Ping health check ($0.01)
#   Ejecutado: Ping health check

Ejercicio 3: Escalación por niveles (Medio)

Implementa escalación con 3 niveles: acciones <$100 → usuario, <$1000 → team lead, >=$1000 → admin. Cada nivel registra quién aprobó. Prueba con 3 acciones de diferente costo.

Ver solución
from dotenv import load_dotenv
load_dotenv()

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

class State(TypedDict):
    action: str
    cost: float
    approver: str
    decision: str
    chain: Annotated[list[str], operator.add]

def escalate(state: State) -> dict:
    cost = state["cost"]
    if cost < 100:
        approver = "user"
    elif cost < 1000:
        approver = "team_lead"
    else:
        approver = "admin"
    return {
        "approver": approver,
        "chain": [f"Escalado a {approver} (${cost:.2f})"],
    }

def request(state: State) -> dict:
    response = interrupt({
        "action": state["action"],
        "cost": state["cost"],
        "required_approver": state["approver"],
    })
    return {
        "decision": response,
        "chain": [f"{state['approver']} decidió: {response}"],
    }

def act(state: State) -> dict:
    event = "Ejecutado" if state["decision"] == "approve" else "Rechazado"
    return {"chain": [f"{event}: {state['action']}"]}

builder = StateGraph(State)
builder.add_node("escalate", escalate)
builder.add_node("request", request)
builder.add_node("act", act)

builder.add_edge(START, "escalate")
builder.add_edge("escalate", "request")
builder.add_edge("request", "act")
builder.add_edge("act", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

cases = [
    {"action": "Comprar dominio", "cost": 12.0, "thread": "esc-a"},
    {"action": "Contratar hosting", "cost": 500.0, "thread": "esc-b"},
    {"action": "Licencia enterprise", "cost": 5000.0, "thread": "esc-c"},
]

for case in cases:
    config = {"configurable": {"thread_id": case["thread"]}}
    graph.invoke(
        {"action": case["action"], "cost": case["cost"], "approver": "", "decision": "", "chain": []},
        config,
    )
    result = graph.invoke(Command(resume="approve"), config)
    print(f"\n{case['action']} (${case['cost']:.2f}):")
    for step in result["chain"]:
        print(f"  {step}")
# Output esperado:
#
# Comprar dominio ($12.00):
#   Escalado a user ($12.00)
#   user decidió: approve
#   Ejecutado: Comprar dominio
#
# Contratar hosting ($500.00):
#   Escalado a team_lead ($500.00)
#   team_lead decidió: approve
#   Ejecutado: Contratar hosting
#
# Licencia enterprise ($5000.00):
#   Escalado a admin ($5000.00)
#   admin decidió: approve
#   Ejecutado: Licencia enterprise

Ejercicio 4: Batch approval con filtros (Medio)

Crea un sistema de batch approval que acepte 3 modos: approve_all, reject_all, approve_low_risk_only. Prueba con una lista de 5 acciones de riesgo variado y verifica que el filtrado funciona correctamente.

Ver solución
from dotenv import load_dotenv
load_dotenv()

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

class State(TypedDict):
    actions: list[dict]
    mode: str
    approved: Annotated[list[str], operator.add]
    rejected: Annotated[list[str], operator.add]

def request_batch(state: State) -> dict:
    response = interrupt({
        "actions": state["actions"],
        "options": "approve_all | reject_all | approve_low_risk_only",
    })
    return {"mode": response}

def apply_decision(state: State) -> dict:
    mode = state["mode"]
    approved = []
    rejected = []

    for action in state["actions"]:
        if mode == "approve_all":
            approved.append(action["name"])
        elif mode == "reject_all":
            rejected.append(action["name"])
        elif mode == "approve_low_risk_only":
            if action["risk"] == "low":
                approved.append(action["name"])
            else:
                rejected.append(action["name"])

    return {"approved": approved, "rejected": rejected}

builder = StateGraph(State)
builder.add_node("request", request_batch)
builder.add_node("apply", apply_decision)

builder.add_edge(START, "request")
builder.add_edge("request", "apply")
builder.add_edge("apply", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

actions = [
    {"name": "Read API", "risk": "low"},
    {"name": "Write DB", "risk": "high"},
    {"name": "Search index", "risk": "low"},
    {"name": "Delete records", "risk": "high"},
    {"name": "Cache update", "risk": "low"},
]

config = {"configurable": {"thread_id": "batch-filter"}}
graph.invoke({"actions": actions, "mode": "", "approved": [], "rejected": []}, config)
result = graph.invoke(Command(resume="approve_low_risk_only"), config)

print("Aprobadas:")
for name in result["approved"]:
    print(f"  ✅ {name}")
print("Rechazadas:")
for name in result["rejected"]:
    print(f"  ❌ {name}")
# Output esperado:
# Aprobadas:
#   ✅ Read API
#   ✅ Search index
#   ✅ Cache update
# Rechazadas:
#   ❌ Write DB
#   ❌ Delete records

Ejercicio 5: HITL metrics tracker (Medio)

Implementa una clase HITLTracker que registre auto-approves, human approvals y rejections. Simula 10 decisiones y genera un reporte con interrupt rate, approval rate y tiempo promedio de espera.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import random

class HITLTracker:
    def __init__(self):
        self.records: list[dict] = []

    def record(self, action: str, decision: str, source: str, wait_seconds: float = 0):
        self.records.append({
            "action": action,
            "decision": decision,
            "source": source,
            "wait_seconds": wait_seconds,
        })

    def report(self) -> dict:
        total = len(self.records)
        if total == 0:
            return {"error": "No records"}

        auto = [r for r in self.records if r["source"] == "auto"]
        human = [r for r in self.records if r["source"] == "human"]
        approved = [r for r in human if r["decision"] == "approve"]
        rejected = [r for r in human if r["decision"] == "reject"]

        wait_times = [r["wait_seconds"] for r in human if r["wait_seconds"] > 0]
        avg_wait = sum(wait_times) / len(wait_times) if wait_times else 0

        return {
            "total_actions": total,
            "auto_approved": len(auto),
            "human_approved": len(approved),
            "human_rejected": len(rejected),
            "interrupt_rate": f"{len(human) / total:.1%}",
            "approval_rate": f"{len(approved) / max(1, len(human)):.1%}",
            "avg_wait_seconds": f"{avg_wait:.1f}s",
        }

tracker = HITLTracker()

random.seed(42)
actions = [
    ("Search Wikipedia", 0.01, "auto"),
    ("Query database", 0.05, "auto"),
    ("Send notification", 2.0, "human"),
    ("Update config", 0.10, "auto"),
    ("Deploy to staging", 10.0, "human"),
    ("Read file", 0.01, "auto"),
    ("Delete user data", 50.0, "human"),
    ("Generate report", 0.50, "auto"),
    ("Send bulk email", 25.0, "human"),
    ("Refresh cache", 0.02, "auto"),
]

for action_name, cost, source in actions:
    if source == "auto":
        tracker.record(action_name, "approve", "auto")
    else:
        decision = random.choice(["approve", "approve", "approve", "reject"])
        wait = random.uniform(30, 600)
        tracker.record(action_name, decision, "human", wait)

report = tracker.report()
print("=== HITL Production Metrics ===")
for key, value in report.items():
    print(f"  {key}: {value}")
# Output esperado:
# === HITL Production Metrics ===
#   total_actions: 10
#   auto_approved: 6
#   human_approved: 3
#   human_rejected: 1
#   interrupt_rate: 40.0%
#   approval_rate: 75.0%
#   avg_wait_seconds: 284.1s

Ejercicio 6: API de approval completa (Avanzado)

Construye un sistema con dos funciones que simulen endpoints: start_task(task, cost, thread_id) que lanza un agente con auto-approve si el costo es bajo, y approve_task(thread_id, decision) que resume agentes pendientes. Incluye un registro centralizado de pending approvals. Ejecuta 3 tareas: una que se auto-apruebe, una que espere aprobación humana, y una que sea rechazada.

Ver solución
from dotenv import load_dotenv
load_dotenv()

import time
import operator
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
from langgraph.types import interrupt, Command

class State(TypedDict):
    task: str
    cost: float
    result: str
    decision: str
    source: str
    audit: Annotated[list[str], operator.add]

def decide(state: State) -> dict:
    if state["cost"] < 1.0:
        return {
            "decision": "approve",
            "source": "auto",
            "audit": [f"Auto-approved (${state['cost']:.2f})"],
        }

    response = interrupt({
        "task": state["task"],
        "cost": state["cost"],
    })
    return {
        "decision": response,
        "source": "human",
        "audit": [f"Human: {response}"],
    }

def execute(state: State) -> dict:
    if state["decision"] == "approve":
        return {
            "result": f"Completado: {state['task']}",
            "audit": ["Ejecutado"],
        }
    return {
        "result": f"Cancelado: {state['task']}",
        "audit": ["Cancelado"],
    }

builder = StateGraph(State)
builder.add_node("decide", decide)
builder.add_node("execute", execute)
builder.add_edge(START, "decide")
builder.add_edge("decide", "execute")
builder.add_edge("execute", END)

checkpointer = MemorySaver()
graph = builder.compile(checkpointer=checkpointer)

PENDING = {}

def start_task(task: str, cost: float, thread_id: str) -> dict:
    config = {"configurable": {"thread_id": thread_id}}
    result = graph.invoke(
        {"task": task, "cost": cost, "result": "", "decision": "", "source": "", "audit": []},
        config,
    )
    state = graph.get_state(config)

    if state.next:
        PENDING[thread_id] = {"task": task, "cost": cost, "since": time.time()}
        return {"thread_id": thread_id, "status": "pending_approval"}

    return {"thread_id": thread_id, "status": "completed", "result": result["result"]}

def approve_task(thread_id: str, decision: str) -> dict:
    config = {"configurable": {"thread_id": thread_id}}
    result = graph.invoke(Command(resume=decision), config)
    PENDING.pop(thread_id, None)
    return {"thread_id": thread_id, "status": "completed", "result": result["result"]}

print("=== Task 1: bajo costo (auto-approve) ===")
r1 = start_task("Buscar datos", 0.10, "t-001")
print(f"  {r1}")

print("\n=== Task 2: alto costo (espera humano) ===")
r2 = start_task("Enviar campaña", 50.0, "t-002")
print(f"  {r2}")

print(f"\n=== Pending: {list(PENDING.keys())} ===")

print("\n=== Aprobando task 2 ===")
r3 = approve_task("t-002", "approve")
print(f"  {r3}")

print("\n=== Task 3: alto costo (rechazada) ===")
r4 = start_task("Borrar registros", 100.0, "t-003")
print(f"  {r4}")
r5 = approve_task("t-003", "reject")
print(f"  {r5}")

print(f"\n=== Pending final: {list(PENDING.keys())} ===")
# Output esperado:
# === Task 1: bajo costo (auto-approve) ===
#   {'thread_id': 't-001', 'status': 'completed', 'result': 'Completado: Buscar datos'}
#
# === Task 2: alto costo (espera humano) ===
#   {'thread_id': 't-002', 'status': 'pending_approval'}
#
# === Pending: ['t-002'] ===
#
# === Aprobando task 2 ===
#   {'thread_id': 't-002', 'status': 'completed', 'result': 'Completado: Enviar campaña'}
#
# === Task 3: alto costo (rechazada) ===
#   {'thread_id': 't-003', 'status': 'pending_approval'}
#   {'thread_id': 't-003', 'status': 'completed', 'result': 'Cancelado: Borrar registros'}
#
# === Pending final: [] ===

Resumen

En esta cápsula aprendiste:

  • Producción HITL es asíncrono — el agente se pausa, notifica via webhook/Slack/email, y el humano responde horas después. El checkpointer mantiene el estado entre la pausa y el resume
  • Timeouts con auto-approve por nivel de riesgo protegen contra humanos que no responden — acciones de bajo riesgo se auto-aprueban en minutos, acciones de alto riesgo nunca se auto-aprueban
  • Escalación por niveles de autoridad ruteaa la aprobación al rol correcto — no todo requiere un admin, no todo lo puede aprobar un usuario regular
  • Batch approvals evitan fatiga de aprobación cuando hay decenas de acciones similares — approve_all, reject_all, approve_low_risk_only cubren el 90% de los casos
  • Audit trail es obligatorio en producción — cada decisión, quién la tomó, cuándo, y qué se ejecutó queda registrado para compliance y debugging
  • La tasa de interrupción es tu métrica clave — target <20%. Muy alta = agente inútil. Muy baja = riesgo no mitigado. Optimiza los thresholds de riesgo iterativamente
  • La arquitectura es simple — dos endpoints (/start y /approve), un checkpointer durable (PostgreSQL), y un registro de pending approvals. Todo lo demás es lógica de negocio

Próxima cápsula: con HITL dominado, estás listo para el proyecto integrador del módulo — un agente supervisado completo con aprobaciones, feedback loops, escalación, y métricas de producción.


Recursos adicionales

  1. LangGraph Human-in-the-Loop — Conceptos fundamentales de HITL en LangGraph
  2. LangGraph Persistence — Checkpointing como base de HITL asíncrono
  3. How to wait for user input — Patrones de espera para input humano
  4. LangGraph Deployment — Deployment de grafos con HITL en producción
  5. LangGraph interrupt() — Referencia de la función interrupt
  6. Designing Human-AI Workflows — Google PAIR: guía de diseño para flujos humano-IA

Módulo 9 — LangChain & LangGraph: From Chains to Agents