Módulo 3: Pipelines secuenciales

Cuándo una etapa falla

Descripción

Las lecciones 03, 04 y 05 asumieron, sin decirlo, que cada etapa deja siempre lo que la siguiente necesita: quote siempre entrega un price_cents, validate_policy siempre entrega un cleared_to_book. Mira de nuevo extract_payload de la lección 04: la rama de validate_policy pone new_payload["cleared_to_book"] = True sin ninguna condición — como si el resultado de search_docs siempre alcanzara para confirmar. No siempre alcanza. El stub de policy_agent tiene solo dos entradas (cancellation-policy, no-show-policy); cualquier pregunta que no matchee ninguna de las dos devuelve "No se encontró una política relevante para esa pregunta." — un resultado real, sin ninguna excepción, que no es una base válida para confirmar una reserva.

Esta lección construye la pieza que faltaba: una guarda que revisa, después de cada etapa, si lo que dejó alcanza para seguir — y que detiene el pipeline ahí mismo, antes de que la etapa confirm llegue a correr sobre una política sin validar. Es la propiedad que un supervisor tiene casi gratis —puede devolver None y no delegar nada, como ya viste en el Módulo 2— y que un pipeline, por no decidir nunca nada, tiene que ganarse a propósito.

Conexión con el módulo

Esta lección reusa, sin cambios, PipelineStage, build_stage_task y PIPELINE_STAGES de las lecciones 02 a 04. Extiende extract_payload con una condición real en vez de un True incondicional, y agrega dos piezas nuevas: stage_ok (la guarda) y run_pipeline_guarded (el runner que la aplica). La lección 07 retoma esta misma distinción —qué depende genuinamente de qué— desde otro ángulo: si quote y validate_policy no dependen entre sí, la mitad del riesgo que esta lección cubre podría, en realidad, correr en paralelo (Módulo 4). El mini-proyecto de la lección 08 usa run_pipeline_guarded completo, sin ningún cambio, sobre tres escenarios nuevos.


Analogía: el control de calidad entre estaciones, no solo al final

Vuelve a la línea de ensamblaje. Hasta ahora, cada estación entregaba su pieza a la siguiente sin que nadie la revisara en el camino — se asumía que, si la estación de corte terminó, la pieza está lista para soldar. Una línea de ensamblaje seria no asume eso: entre cada par de estaciones hay un control de calidad puntual, no solo al final de toda la línea. Si el control detecta que la pieza que sale de soldadura tiene una fisura, la línea se detiene ahí — no sigue mandando esa pieza a pintura, esperando que alguien la note en la inspección final. stage_ok es ese control de calidad puntual: revisa, etapa por etapa, no solo al final.


Ejemplo trabajado: la guarda, primero en el camino feliz, después deteniendo el pipeline

import ast
import concurrent.futures
from dataclasses import dataclass
import reservo_tools as rt


def dispatch_parallel(tool_use_blocks, tools):
    with concurrent.futures.ThreadPoolExecutor(max_workers=len(tool_use_blocks)) as pool:
        futures = [pool.submit(tools[b["name"]], **b["input"]) for b in tool_use_blocks]
        results = [f.result() for f in futures]
    return [
        {"type": "tool_result", "tool_use_id": b["id"], "content": str(r)}
        for b, r in zip(tool_use_blocks, results)
    ]


def run_agent_parallel(question, model_script, tools, max_iterations=10):
    messages = [{"role": "user", "content": question}]
    for step in range(max_iterations):
        turn = model_script[step]
        messages.append({"role": "assistant", "content": turn["content"]})
        if turn["stop_reason"] != "tool_use":
            return turn, messages
        tool_result_blocks = dispatch_parallel(turn["content"], tools)
        messages.append({"role": "user", "content": tool_result_blocks})
    raise RuntimeError(f"max_iterations alcanzado ({max_iterations})")


POLICY_DOCS = {
    "no-show-policy": (
        "Si un miembro no se presenta a una reserva confirmada y no cancela "
        "con al menos 2 horas de anticipación, Reservo cobra el 50% del "
        "precio cotizado como cargo por no-presentación."
    ),
    "cancellation-policy": (
        "Las reservas se pueden cancelar sin cargo hasta 2 horas antes del "
        "horario reservado. Cancelaciones dentro de esas 2 horas aplican "
        "el cargo de no-presentación."
    ),
}


def search_docs(query):
    """STUB de policy_agent: solo dos entradas. Cualquier query que no
    matchee ninguna devuelve el mensaje de 'no encontrado' -- un resultado
    REAL, no una excepción."""
    q = query.lower()
    if "cancela" in q:
        return f"[cancellation-policy] {POLICY_DOCS['cancellation-policy']}"
    if "no" in q and ("present" in q or "show" in q):
        return f"[no-show-policy] {POLICY_DOCS['no-show-policy']}"
    return "No se encontró una política relevante para esa pregunta."


SPECIALISTS = {
    "booking_agent": {
        "tools": {
            "list_rooms": rt.list_rooms, "get_quote": rt.get_quote,
            "book_room": rt.book_room, "cancel_booking": rt.cancel_booking,
        },
        "expertise": "cotizar, reservar y cancelar salas",
    },
    "policy_agent": {
        "tools": {"search_docs": search_docs},
        "expertise": "responder preguntas de política (cancelación, no-presentación)",
    },
}


def run_specialist(name, task, model_script):
    tools = SPECIALISTS[name]["tools"]
    return run_agent_parallel(task, model_script, tools)


@dataclass
class PipelineStage:
    kind: str
    name: str
    label: str


def last_tool_result(history):
    for m in reversed(history):
        if isinstance(m["content"], list):
            for b in m["content"]:
                if b["type"] == "tool_result":
                    try:
                        return ast.literal_eval(b["content"])
                    except (ValueError, SyntaxError):
                        return b["content"]
    return None


def build_stage_task(stage, payload):
    if stage.kind == "quote":
        return f"Cotiza {payload['room']} {payload['tier']} {payload['hours']}h."
    if stage.kind == "validate_policy":
        if payload.get("client_type") == "corporate":
            return (f"¿Qué política aplica para una reserva corporativa grande de "
                     f"{payload['room']}, antes de confirmarla para {payload['member']}?")
        return (f"¿Cuál es la política de cancelación para una reserva de "
                 f"{payload['room']} de {payload['hours']}h, antes de confirmarla?")
    if stage.kind == "confirm":
        if payload.get("cleared_to_book"):
            return (f"Reserva {payload['room']} {payload['tier']} {payload['hours']}h "
                     f"para {payload['member']} -- la política de cancelación ya se validó.")
        return (f"Reserva {payload['room']} {payload['tier']} {payload['hours']}h "
                 f"para {payload['member']}.")
    raise ValueError(f"no sé armar la tarea de la etapa {stage.kind!r}")


def extract_payload(stage, history, payload):
    new_payload = dict(payload)
    result = last_tool_result(history)
    if stage.kind == "quote":
        new_payload["price_cents"] = result["price_cents"]
    elif stage.kind == "validate_policy":
        new_payload["policy_result"] = result
        # LA GUARDA de esta lección: cleared_to_book SOLO es True si
        # search_docs encontró una política real -- se reconoce porque el
        # stub la prefija con "[algo]". En las lecciones 03/04, esta rama
        # ponía cleared_to_book = True sin ninguna condición.
        new_payload["cleared_to_book"] = isinstance(result, str) and result.startswith("[")
    elif stage.kind == "confirm":
        new_payload["booking_id"] = result["booking_id"]
        new_payload["confirmed"] = result["confirmed"]
    return new_payload


class PipelineHalted(Exception):
    """Se lanza cuando una etapa deja un payload que no alcanza para que
    la etapa siguiente continúe con seguridad. Lleva el trace PARCIAL --
    lo que sí alcanzó a correr -- para que quien la capture pueda
    inspeccionar qué pasó antes de detenerse."""
    def __init__(self, reason, trace):
        super().__init__(reason)
        self.trace = trace


def stage_ok(stage, payload):
    """La guarda: ejecutada, sin ningún modelo de por medio. Decide si el
    payload que dejó ESTA etapa alcanza para seguir con la próxima."""
    if stage.kind == "validate_policy" and not payload.get("cleared_to_book"):
        return False, "no se encontró una política relevante -- no hay base para confirmar"
    return True, None


def run_pipeline_guarded(stages, model_scripts, initial_payload):
    """Igual que run_pipeline (lecciones 03/04), MÁS una guarda: después
    de extraer el payload de cada etapa, verifica con stage_ok si alcanza
    para seguir. Si no alcanza, detiene el pipeline ahí mismo -- nunca
    llega a construir la tarea de la etapa siguiente."""
    payload = dict(initial_payload)
    trace = []
    for i, stage in enumerate(stages):
        task = build_stage_task(stage, payload)
        final, history = run_specialist(stage.name, task, model_scripts[i])
        payload = extract_payload(stage, history, payload)
        trace.append({
            "stage": i + 1, "kind": stage.kind, "agent": stage.name,
            "label": stage.label, "task": task,
            "output": final["content"][0]["text"],
        })
        ok, reason = stage_ok(stage, payload)
        if not ok:
            raise PipelineHalted(
                f"pipeline detenido después de la etapa {i + 1} ({stage.label}): {reason}",
                trace,
            )
    return payload, trace


# El PIPELINE_STAGES de siempre, sin ningún cambio. La etiqueta
# "validar la política de cancelación" describe el ROL de la etapa 2 --
# no cambia aunque la pregunta real (build_stage_task) sea otra para un
# cliente corporativo, como ves más abajo.
PIPELINE_STAGES = [
    PipelineStage(kind="quote", name="booking_agent", label="cotizar"),
    PipelineStage(kind="validate_policy", name="policy_agent",
                  label="validar la política de cancelación"),
    PipelineStage(kind="confirm", name="booking_agent", label="confirmar la reserva"),
]

print("--- primero, confirmamos que la guarda NO estorba el camino feliz ---")
INITIAL_PAYLOAD = {"room": "Focus", "tier": "pro", "hours": 3, "member": "Ana"}
model_script_quote = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
         "input": {"room": "Focus", "tier": "pro", "hours": 3}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Focus pro 3h cuesta 6000 centavos."}]},
]
model_script_policy = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "search_docs",
         "input": {"query": "política de cancelación"}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": (
            "Puedes cancelar sin cargo hasta 2 horas antes del horario "
            "reservado. No hay ningún impedimento para confirmar."
        )}]},
]
model_script_confirm = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "book_room",
         "input": {"room": "Focus", "tier": "pro", "hours": 3, "member": "Ana"}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Reservé Focus pro 3h para Ana (confirmación #1)."}]},
]
payload, trace = run_pipeline_guarded(
    PIPELINE_STAGES, [model_script_quote, model_script_policy, model_script_confirm], INITIAL_PAYLOAD,
)
print("cleared_to_book:", payload["cleared_to_book"], "| booking_id:", payload["booking_id"])

print()
print("--- ahora, una reserva corporativa que activa la guarda ---")
INITIAL_PAYLOAD_CORP = {
    "room": "Boardroom", "tier": "pro", "hours": 6, "member": "TechCorp",
    "client_type": "corporate",
}
model_script_quote_corp = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
         "input": {"room": "Boardroom", "tier": "pro", "hours": 6}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Boardroom pro 6h cuesta 38400 centavos."}]},
]
model_script_policy_miss = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "search_docs",
         "input": {"query": "política para reservas corporativas grandes"}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": (
            "No encontré una política específica para reservas corporativas "
            "grandes en nuestra base -- te recomiendo escribir al equipo de "
            "cuentas corporativas antes de confirmar."
        )}]},
]
# Solo 2 guiones: la etapa 3 (confirmar) nunca llega a correr.
model_scripts_corp = [model_script_quote_corp, model_script_policy_miss]

try:
    payload_corp, trace_corp = run_pipeline_guarded(
        PIPELINE_STAGES, model_scripts_corp, INITIAL_PAYLOAD_CORP,
    )
except PipelineHalted as e:
    print(f"PipelineHalted capturado: {e}")
    trace_corp = e.trace

print()
print("--- el trace parcial: la etapa 1 SÍ corrió, la etapa 3 nunca se intentó ---")
for step in trace_corp:
    print(f"etapa {step['stage']} ({step['label']}) -> {step['agent']}: {step['output']!r}")
print("BOOKINGS después del intento corporativo:", rt.BOOKINGS)

Qué esperar:

--- primero, confirmamos que la guarda NO estorba el camino feliz ---
cleared_to_book: True | booking_id: 1

--- ahora, una reserva corporativa que activa la guarda ---
PipelineHalted capturado: pipeline detenido después de la etapa 2 (validar la política de cancelación): no se encontró una política relevante -- no hay base para confirmar

--- el trace parcial: la etapa 1 SÍ corrió, la etapa 3 nunca se intentó ---
etapa 1 (cotizar) -> booking_agent: 'Boardroom pro 6h cuesta 38400 centavos.'
etapa 2 (validar la política de cancelación) -> policy_agent: 'No encontré una política específica para reservas corporativas grandes en nuestra base -- te recomiendo escribir al equipo de cuentas corporativas antes de confirmar.'
BOOKINGS después del intento corporativo: {1: {'booking_id': 1, 'room': 'Focus', 'tier': 'pro', 'hours': 3, 'member': 'Ana', 'price_cents': 6000}}

Dos cosas para notar. Primero: el camino feliz (Focus pro 3h para Ana) produce exactamente el mismo resultado que ya viste en la lección 04 —cleared_to_book: True, booking_id: 1— la guarda no le agrega ningún costo ni ningún paso extra cuando todo va bien; solo actúa cuando hace falta. Segundo: BOOKINGS al final solo tiene una reserva —la de Ana—. El intento corporativo nunca llegó a book_room: la excepción se lanzó apenas stage_ok detectó que cleared_to_book seguía siendo False después de la etapa 2, antes de que run_pipeline_guarded construyera siquiera la tarea de la etapa 3.


Por qué la excepción lleva el trace parcial

Fíjate en PipelineHalted.__init__: recibe reason y trace, no solo reason. Si la excepción solo llevara el mensaje, quien la capture sabría que el pipeline se detuvo, pero no qué sí alcanzó a pasar antes — información valiosa para decidir qué hacer después (¿reintentar solo la etapa 2? ¿escalar a un humano con el contexto ya armado?). Guardar el trace parcial en la excepción misma evita el problema que ya viste con el NameError al capturar mal el resultado de una función que lanza antes de retornar: sin el trace adentro de la excepción, la variable trace_corp nunca llegaría a existir, porque run_pipeline_guarded nunca alcanza su return.


Errores comunes

  1. Pensar que un pipeline sin excepciones significa un pipeline sin fallas. El ejemplo trabajado de esta lección corrió sin lanzar ningún KeyError ni TypeError en el camino feliz — pero el escenario corporativo también corrió "sin errores de Python" hasta que stage_ok decidió, a propósito, lanzar PipelineHalted. La ausencia de una excepción de Python no prueba que el resultado de una etapa haya alcanzado para seguir.

  2. Confundir cleared_to_book = False con un bug. No lo es — es exactamente el resultado correcto cuando search_docs no encuentra nada relevante. El bug estaba en las lecciones 03/04, que ponían True sin condición, no en que el resultado real de la tool sea, a veces, un "no encontrado".

  3. Poner la guarda dentro de build_stage_task en vez de después de extract_payload. La guarda necesita el payload YA actualizado con lo que dejó la etapa que acaba de correr — revisarlo antes, dentro de build_stage_task, llegaría demasiado tarde: esa función arma la tarea de la etapa actual, no decide si la etapa anterior alcanzó.

  4. Olvidar que stage_ok solo conoce el kind de la etapa que acaba de terminar, no el pipeline completo. Igual que build_stage_task (lección 03, error común 4), cada guarda nueva que agregues solo necesita saber de su propio kind — no hace falta tocar la guarda de validate_policy para agregar, por ejemplo, una guarda para quote.

  5. Pensar que esta guarda cubre TODOS los modos de falla posibles. No — el Ejercicio 3 de esta lección construye, a propósito, un guion mal escrito que la guarda no detecta: una etapa que nunca pide ninguna tool. Conocer los límites de una guarda es tan importante como construirla.


Ejercicios

Ejercicio 1: Confirma tu propia ejecución (Fácil)

Ejecuta el ejemplo trabajado completo tú mismo y confirma, línea por línea, que tu salida coincide con el "Qué esperar" de arriba. Presta atención a BOOKINGS al final: debería tener una sola reserva, no dos.

Ver solución

No hay una única "solución de código" para este ejercicio — es una verificación: si BOOKINGS al final tiene exactamente una entrada (booking_id: 1, la de Ana), confirmaste que el intento corporativo nunca llegó a book_room — la prueba más directa de que la guarda detuvo el pipeline donde dijo que lo detendría.

Ejercicio 2: Activa la guarda SIN el camino corporativo (Medio)

stage_ok no sabe nada de client_type — solo mira cleared_to_book. Confirma esto: usa el INITIAL_PAYLOAD estándar (sin client_type="corporate") para un socio que pregunta, por error, algo que el stub de search_docs tampoco cubre —por ejemplo, "¿cuáles son los horarios de acceso al edificio?"— y confirma que run_pipeline_guarded se detiene igual.

Ver solución
payload_std = {"room": "Focus", "tier": "pro", "hours": 2, "member": "Marcos"}
script_quote_std = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
         "input": {"room": "Focus", "tier": "pro", "hours": 2}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Focus pro 2h cuesta 4000 centavos."}]},
]
script_policy_miss_std = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "search_docs",
         "input": {"query": "horarios de acceso al edificio"}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "No encontré una política de horarios de acceso en nuestra base."}]},
]
try:
    run_pipeline_guarded(PIPELINE_STAGES, [script_quote_std, script_policy_miss_std], payload_std)
except PipelineHalted as e:
    print(f"PipelineHalted capturado: {e}")

Salida real:

PipelineHalted capturado: pipeline detenido después de la etapa 2 (validar la política de cancelación): no se encontró una política relevante -- no hay base para confirmar

Explicación: el pipeline se detiene exactamente igual, sin ningún client_type="corporate" en el payload. Esto confirma el punto central de la lección: la guarda mira el resultado real de search_docs —si empieza con "[", si trae una etiqueta de política reconocida—, no una bandera que alguien puso de antemano. Cualquier pregunta que el stub no cubra activa la misma guarda, sin importar por qué el socio la hizo.

Ejercicio 3: Un guion mal escrito que la guarda NO detecta (Difícil)

stage_ok solo revisa cleared_to_book después de validate_policy. Construye un model_script_quote roto a propósito —uno que termina en end_turn sin pedir get_quote primero— y ejecuta la etapa quote con él a través de run_pipeline_guarded. ¿Qué excepción se produce, y por qué NO es un PipelineHalted limpio como el del ejemplo trabajado?

Ver solución
payload_bad = {"room": "Focus", "tier": "pro", "hours": 2, "member": "Nora"}
script_quote_bad = [
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Cotizando Focus pro 2h..."}]},
]
try:
    run_pipeline_guarded(PIPELINE_STAGES, [script_quote_bad], payload_bad)
except TypeError as e:
    print(f"TypeError (NO PipelineHalted) capturado: {e!r}")

Salida real:

TypeError (NO PipelineHalted) capturado: TypeError("'NoneType' object is not subscriptable")

Explicación: este guion nunca pide get_quote —termina en end_turn de inmediato—, así que last_tool_result no encuentra ningún tool_result en el historial y devuelve None (exactamente el caso que la lección 03, error común 5, ya había señalado como posible). extract_payload intenta entonces result["price_cents"] sobre None, y Python lanza un TypeError antes de que stage_ok tenga oportunidad de correr — la guarda de esta lección solo revisa el payload después de que extract_payload termina con éxito; si extract_payload mismo falla, la guarda nunca se ejecuta. Esta es una distinción importante y honesta: stage_ok protege contra resultados válidos que no alcanzan (el modo de falla que diseñó esta lección), no contra guiones mal escritos que rompen el mecanismo de extracción antes de llegar ahí — ese segundo modo de falla necesitaría una guarda distinta, en un punto distinto del código, fuera del alcance de esta lección.


Resumen y siguiente paso

  • Las lecciones 03/04 asumían, sin decirlo, que validate_policy siempre alcanza para confirmar — extract_payload ponía cleared_to_book = True sin ninguna condición. Esta lección corrige eso: cleared_to_book ahora depende del resultado REAL de search_docs.
  • stage_ok es la guarda: revisa, después de cada etapa, si lo que dejó alcanza para seguir. run_pipeline_guarded la aplica, y lanza PipelineHalted —con el trace parcial adentro— apenas una etapa no alcanza, sin llegar a construir la tarea de la etapa siguiente.
  • Confirmamos, con salida real, que la guarda no le agrega costo al camino feliz (el mismo booking_id: 1 de la lección 04) y que sí detiene, a tiempo, una reserva corporativa sin política validada — sin que BOOKINGS termine con una reserva que nunca debió confirmarse.
  • El Ejercicio 3 dejó un límite honesto: la guarda protege contra resultados válidos que no alcanzan, no contra un guion tan mal escrito que rompe la extracción del dato antes de que la guarda tenga oportunidad de correr.

Siguiente lección: 07 — Eligiendo pipeline o supervisor. Con el riesgo de un pipeline ya medido (lección 05) y su fragilidad ya cubierta (esta lección), cerramos el criterio completo — y confirmamos, ejecutando el código, que dos de las tres etapas de este módulo no dependían entre sí desde el principio.


Recursos adicionales

  1. Anthropic — Building effective agents — El principio de agregar puertas de verificación programáticas entre pasos de una cadena, exactamente lo que stage_ok construye para este pipeline.
  2. Python — Excepciones definidas por el usuario — La base de PipelineHalted, que además de heredar de Exception lleva su propio atributo (trace) para no perder el estado parcial.
  3. Anthropic — Multi-agent research system — Un sistema real donde detectar, a tiempo, que un sub-resultado no alcanza evita que el error se propague silenciosamente hasta el resultado final.
  4. Python — try/except y el flujo de control — El mecanismo que separa, en esta lección, el "camino feliz" de la detención a tiempo.