Módulo 3: Pipelines secuenciales

El pipeline de Reservo, ejecutado de punta a punta

Descripción

Las dos lecciones anteriores construyeron cada pieza por separado: qué identifica una etapa (02) y cómo se encadena el dato entre dos de ellas (03). Esta lección las junta sobre el caso completo del módulo: las tres etapas del proceso real de reserva de Reservo —cotizar, validar la política de cancelación, confirmar— corriendo de punta a punta con run_pipeline, sin ningún cambio respecto a la lección 03. No hay ninguna pieza nueva acá — es la primera vez que ves el pipeline completo funcionando, con el historial de cada especialista impreso y el payload final citado.

Esta es la pieza central del módulo: la que sostiene la comparación de costo de la lección 05, la guarda de la lección 06, y los tres escenarios del mini-proyecto de la lección 08.

Conexión con el módulo

Esta lección reusa, sin cambios, PipelineStage (02), build_stage_task, extract_payload y run_pipeline (03). Lo único nuevo es la etapa de política insertada en el medio de la secuencia, y los tres guiones (concepto) que corresponden al caso completo. La lección 05 toma el resultado ejecutado de esta lección y le agrega el conteo de costo comparado contra un supervisor.


Analogía: la línea de ensamblaje, un turno completo

Las lecciones anteriores mostraron una estación de la línea a la vez: cortar, sola; cortar y soldar, en un pipeline pequeño de dos estaciones. Esta lección es la línea completa, en un turno real: la pieza entra por la primera estación, pasa por la segunda sin que nadie tenga que intervenir, y sale terminada por la tercera — cotizada, con su política validada, y reservada.


Ejemplo trabajado: el pipeline completo

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):
    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":
        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["cancellation_policy"] = result
        new_payload["cleared_to_book"] = True
    elif stage.kind == "confirm":
        new_payload["booking_id"] = result["booking_id"]
        new_payload["confirmed"] = result["confirmed"]
    return new_payload


def run_pipeline(stages, model_scripts, initial_payload):
    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, "history": history,
            "output": final["content"][0]["text"],
        })
    return payload, trace


# El pipeline completo de Reservo: cotizar -> validar -> confirmar.
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"),
]

INITIAL_PAYLOAD = {"room": "Focus", "tier": "pro", "hours": 3, "member": "Ana"}

# Los tres guiones (concepto, claude-sonnet-5) -- uno por etapa.
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)."}]},
]

model_scripts = [model_script_quote, model_script_policy, model_script_confirm]

payload, trace = run_pipeline(PIPELINE_STAGES, model_scripts, INITIAL_PAYLOAD)

for step in trace:
    print(f"--- etapa {step['stage']} ({step['label']}) -> {step['agent']} ---")
    print(f"tarea: {step['task']!r}")
    for i, m in enumerate(step["history"]):
        role, content = m["role"], m["content"]
        if isinstance(content, str):
            print(f"  [{i}] {role:<9} pregunta: {content!r}")
            continue
        for block in content:
            if block["type"] == "tool_use":
                print(f"  [{i}] {role:<9} tool_use({block['name']}): {block['input']}")
            elif block["type"] == "tool_result":
                print(f"  [{i}] {role:<9} tool_result: {block['content']}")
            elif block["type"] == "text":
                print(f"  [{i}] {role:<9} texto final: {block['text']!r}")
    print()

print("payload final:", payload)

Qué esperar (sobre un Reservo desechable, recién iniciado):

--- etapa 1 (cotizar) -> booking_agent ---
tarea: 'Cotiza Focus pro 3h.'
  [0] user      pregunta: 'Cotiza Focus pro 3h.'
  [1] assistant tool_use(get_quote): {'room': 'Focus', 'tier': 'pro', 'hours': 3}
  [2] user      tool_result: {'price_cents': 6000}
  [3] assistant texto final: 'Focus pro 3h cuesta 6000 centavos.'

--- etapa 2 (validar la política de cancelación) -> policy_agent ---
tarea: '¿Cuál es la política de cancelación para una reserva de Focus de 3h, antes de confirmarla?'
  [0] user      pregunta: '¿Cuál es la política de cancelación para una reserva de Focus de 3h, antes de confirmarla?'
  [1] assistant tool_use(search_docs): {'query': 'política de cancelación'}
  [2] user      tool_result: [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.
  [3] assistant texto final: 'Puedes cancelar sin cargo hasta 2 horas antes del horario reservado. No hay ningún impedimento para confirmar.'

--- etapa 3 (confirmar la reserva) -> booking_agent ---
tarea: 'Reserva Focus pro 3h para Ana -- la política de cancelación ya se validó.'
  [0] user      pregunta: 'Reserva Focus pro 3h para Ana -- la política de cancelación ya se validó.'
  [1] assistant tool_use(book_room): {'room': 'Focus', 'tier': 'pro', 'hours': 3, 'member': 'Ana'}
  [2] user      tool_result: {'booking_id': 1, 'confirmed': True}
  [3] assistant texto final: 'Reservé Focus pro 3h para Ana (confirmación #1).'

payload final: {'room': 'Focus', 'tier': 'pro', 'hours': 3, 'member': 'Ana', 'price_cents': 6000, 'cancellation_policy': '[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.', 'cleared_to_book': True, 'booking_id': 1, 'confirmed': True}

Tres etapas, tres especialistas invocados (dos de ellos el mismo, booking_agent, en momentos distintos), y ningún punto donde alguien tuvo que decidir "¿cuál sigue?". Fíjate en la tarea de la etapa 3: 'Reserva Focus pro 3h para Ana -- la política de cancelación ya se validó.' — esa frase se armó sola, porque payload['cleared_to_book'] ya era True cuando build_stage_task la construyó. Ningún socio ni ningún modelo tuvo que repetir esa información; viajó, real, desde la etapa 2.


Componiendo la respuesta final (concepto)

El pipeline terminó con un payload completo, pero el socio no necesita ver un diccionario de Python — necesita una respuesta en prosa. Igual que el supervisor del Módulo 2 tenía un Paso 4 (agregar), el pipeline necesita un paso final de síntesis:

# Paso final (concepto, claude-sonnet-5): compone la respuesta para el socio
# a partir del payload completo -- ningún dato nuevo, solo redacción.
final_response = (
    f"La cotización de {payload['room']} {payload['tier']} {payload['hours']}h "
    f"es {payload['price_cents']} centavos. Verifiqué la política de "
    f"cancelación: puedes cancelar sin cargo hasta 2 horas antes. Reservé la "
    f"sala para {payload['member']} (confirmación #{payload['booking_id']})."
)
print("--- respuesta compuesta para el socio (concepto) ---")
print(final_response)

Qué esperar:

--- respuesta compuesta para el socio (concepto) ---
La cotización de Focus pro 3h es 6000 centavos. Verifiqué la política de cancelación: puedes cancelar sin cargo hasta 2 horas antes. Reservé la sala para Ana (confirmación #1).

Esta respuesta combina información de las tres etapas —el precio de la 1, la política de la 2, la confirmación de la 3— en un solo mensaje. Igual que COMPOSE_CALLS en el Módulo 2, esta síntesis es concepto: un turno más del modelo, que la lección 05 cuenta como parte del costo total de coordinación.


Errores comunes

  1. Pensar que este pipeline maneja cualquier tarea de reserva. Solo maneja la secuencia exacta que PIPELINE_STAGES define — una tarea que, por ejemplo, solo pidiera cotizar sin reservar no encaja en este pipeline tal como está (el Módulo 2 ya resuelve ese caso más simple con un solo especialista).

  2. Correr esta lección después de otro ejemplo de la guía en el mismo intérprete. Si book_room no devuelve booking_id: 1, el proceso ya tenía reservas de una corrida anterior. Cada lección de esta guía asume su propio Reservo desechable, recién iniciado.

  3. Confundir "tres etapas" con "tres llamadas al modelo". Cada etapa por sí sola cuesta dos llamadas al modelo (un turno de tool_use, un turno de texto final) — la lección 05 hace la cuenta completa, incluyendo el costo de coordinación que rodea a las tres.

  4. Olvidar que cleared_to_book viene de la etapa 2, no de una suposición. El if payload.get("cleared_to_book") de build_stage_task depende de que la etapa 2 haya corrido ANTES —el orden fijo del pipeline lo garantiza—, pero si alguna vez reordenaras las etapas sin pensarlo, cleared_to_book no existiría todavía cuando la etapa 3 intentara leerlo.

  5. Pensar que la respuesta compuesta final "no hace nada nuevo". A diferencia del Módulo 2, lección 02 —donde la respuesta agregada resultaba ser, por coincidencia, el mismo texto que ya tenía el especialista—, acá la síntesis SÍ combina tres piezas de información distintas en un mensaje que ninguna etapa individual tenía completo.


Ejercicios

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

Ejecuta el pipeline completo de esta lección tú mismo y confirma, línea por línea, que tu salida coincide con el "Qué esperar" de arriba. Presta especial atención al booking_id y al cleared_to_book del payload final.

Ver solución

No hay una única "solución de código" para este ejercicio — es una verificación: si tu salida coincide exactamente con el bloque "Qué esperar" del ejemplo trabajado, tu Reservo desechable arrancó limpio y el pipeline corrió sin desvíos.

Ejercicio 2: Corre el mismo pipeline con Boardroom pro 4h (Medio)

Repite el pipeline completo, en la MISMA sesión donde ya corriste el ejemplo trabajado, con INITIAL_PAYLOAD = {"room": "Boardroom", "tier": "pro", "hours": 4, "member": "Sofía"} y sus tres guiones correspondientes. Confirma el precio (8000 * 4 * 80 // 100) y el booking_id.

Ver solución
INITIAL_PAYLOAD_BR = {"room": "Boardroom", "tier": "pro", "hours": 4, "member": "Sofía"}
script_quote_br = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
         "input": {"room": "Boardroom", "tier": "pro", "hours": 4}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Boardroom pro 4h cuesta 25600 centavos."}]},
]
script_confirm_br = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "book_room",
         "input": {"room": "Boardroom", "tier": "pro", "hours": 4, "member": "Sofía"}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Reservé Boardroom pro 4h para Sofía (confirmación #2)."}]},
]
payload_br, trace_br = run_pipeline(
    PIPELINE_STAGES, [script_quote_br, model_script_policy, script_confirm_br], INITIAL_PAYLOAD_BR,
)
print("precio calculado a mano:", 8000 * 4 * 80 // 100)
print("payload final:", payload_br)

Salida esperada:

precio calculado a mano: 25600
payload final: {'room': 'Boardroom', 'tier': 'pro', 'hours': 4, 'member': 'Sofía', 'price_cents': 25600, 'cancellation_policy': '[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.', 'cleared_to_book': True, 'booking_id': 2, 'confirmed': True}

Explicación: el mismo model_script_policy del ejemplo trabajado sirve sin cambios, porque la pregunta de política de esta corrida (armada por build_stage_task) sigue conteniendo "cancelación" — la keyword que el stub de search_docs reconoce, sin importar qué sala o cuántas horas se pidan. El booking_id sale 2, no 1: esta corrida comparte el mismo proceso —y el mismo reservo_tools.BOOKINGS— que el ejemplo trabajado, que ya había reservado Focus con booking_id: 1 antes de que corrieras este ejercicio. Si en cambio corres este bloque en un intérprete nuevo, sin haber ejecutado antes el ejemplo trabajado, el booking_id sale 1.

Ejercicio 3: ¿Qué pasa si se salta la etapa 2? (Difícil)

Construye un PIPELINE_STAGES alternativo con solo dos etapas —quote y confirm, sin validate_policy— y ejecútalo con INITIAL_PAYLOAD. Confirma que la tarea de la etapa confirm cambia de texto respecto al ejemplo trabajado, y explica por qué.

Ver solución
STAGES_NO_POLICY = [
    PipelineStage(kind="quote", name="booking_agent", label="cotizar"),
    PipelineStage(kind="confirm", name="booking_agent", label="confirmar la reserva"),
]
payload_np, trace_np = run_pipeline(
    STAGES_NO_POLICY, [model_script_quote, model_script_confirm], INITIAL_PAYLOAD,
)
print("tarea de la etapa confirm:", trace_np[1]["task"])

Salida esperada:

tarea de la etapa confirm: 'Reserva Focus pro 3h para Ana.'

Explicación: sin la etapa validate_policy, payload nunca tiene la clave cleared_to_book, así que payload.get("cleared_to_book") devuelve None (falsy) dentro de build_stage_task, y la rama if que menciona la política de cancelación nunca se ejecuta — la tarea queda más corta, sin mencionar ninguna validación. Esto confirma, ejecutando el código, que la mención a la política en la tarea de la etapa 3 depende genuinamente de que la etapa 2 haya corrido antes — no es un texto fijo, es un dato real que viaja o no viaja según el pipeline que se use.


Resumen y siguiente paso

  • Ejecutamos el pipeline completo de Reservo: cotizar → validar la política de cancelación → confirmar, con run_pipeline sin ningún cambio respecto a la lección 03.
  • El precio de la etapa 1 y la validación de la etapa 2 viajaron, reales, hasta la etapa 3 —la tarea de confirmación se armó sola, mencionando la política ya validada, sin que nadie repitiera esa información.
  • Compusimos una respuesta final (concepto) que combina las tres etapas en un solo mensaje para el socio — la primera vez en este módulo donde la síntesis final combina más de una fuente de información real.
  • En ningún momento de las tres etapas hizo falta decidir "¿cuál sigue?" — el orden estaba fijo en PIPELINE_STAGES desde antes de que existiera esta petición concreta.

Siguiente lección: 05 — Midiendo el costo de coordinación de un pipeline. Contamos, con números reales, cuánto ahorra este pipeline frente a un supervisor que tuviera que decidir en cada una de las tres etapas.


Recursos adicionales

  1. Anthropic — Building effective agents — El patrón "prompt chaining" con puertas de verificación entre pasos, ejecutado de punta a punta en esta lección con Reservo.
  2. Anthropic — Multi-agent research system — Un flujo con etapas obligatorias donde el resultado de una alimenta directamente a la siguiente, sin ningún punto de decisión intermedio.
  3. Anthropic — Messages API reference — La forma exacta de tool_use/tool_result/stop_reason que cada etapa de este pipeline respeta, sin cambios.
  4. Python — dataclasses — El módulo detrás de PipelineStage, reusado sin cambios desde la lección 02.