Módulo 7: Orchestrating The Full Reservo System

El supervisor decide, el pipeline valida antes de reservar

Descripción

La lección 02 dejó PLAN con tres Track, sin ejecutar nada todavía. Esta lección toma el primero —book_focus, etiquetado pipeline— y lo resuelve de verdad, con run_pipeline del Módulo 3, sin tocarle una sola línea. No hay ninguna pieza nueva de bajo nivel en esta lección: PipelineStage, build_stage_task, extract_payload y run_pipeline son exactamente los mismos que dejaste funcionando al final del Módulo 3.

Lo que sí es nuevo es el rol que juega la decisión del supervisor aquí, comparado con el Módulo 2 aislado: ahí, el supervisor decidía a qué especialista delegar la petición completa. Acá, ya decidió (en la lección 02, concepto) que esta parte de la petición necesita un pipeline entero —tres etapas, dos agentes distintos, en un orden fijo— y le entrega esa sub-tarea al mecanismo completo, no a un solo especialista.

Conexión con el módulo

Esta lección reusa, sin cambios, PipelineStage, build_stage_task, extract_payload y run_pipeline del Módulo 3 — el mismo pipeline de tres etapas (cotizar → validar política → confirmar) que ya ejecutaste ahí, ahora aplicado como una pieza de una petición más grande. La lección 04 toma este mismo resultado y lo corre junto con una segunda sub-tarea independiente, usando un mecanismo de fan-out que trata a este pipeline entero como una sola unidad de trabajo.


Analogía: la estación de montaje, adentro de un pedido más grande

Retoma la línea de ensamblaje del Módulo 3. Esta lección no cambia nada de cómo funciona esa línea —la pieza sigue entrando por la primera estación, pasando por la segunda sin que nadie intervenga, saliendo terminada por la tercera—. Lo único distinto es el contexto: esa línea de montaje ya no es el pedido completo de la noche — es una estación más dentro de un pedido de mesa grande, que también incluye cosas que ni siquiera pasan por esa línea. El maître (supervisor) no necesita saber cómo funciona la línea de montaje por dentro — solo necesita saber que ESTA parte del pedido tiene que pasar por ahí, completa, antes de considerarse resuelta.


Ejemplo trabajado: el Track book_focus, resuelto con run_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 = {
    "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']}"
    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,
        },
    },
    "policy_agent": {"tools": {"search_docs": search_docs}},
}


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


print("--- sub-tarea del track 'book_focus' (la parte de la petición que necesita pipeline) ---")
print("'Cotiza y reserva Focus pro 3h para el lanzamiento, validando la política de "
      "cancelación antes de 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"}

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)."}]},
]

print()
print("=== Track 'book_focus' resuelto con run_pipeline (Módulo 3, sin cambios) ===")
payload, trace = run_pipeline(
    PIPELINE_STAGES,
    [model_script_quote, model_script_policy, model_script_confirm],
    INITIAL_PAYLOAD,
)
for step in trace:
    print(f"etapa {step['stage']} ({step['label']}) -> {step['agent']}")
    print(f"  tarea: {step['task']!r}")
    print(f"  salida: {step['output']!r}")
print()
print("payload final:", payload)

Qué esperar:

--- sub-tarea del track 'book_focus' (la parte de la petición que necesita pipeline) ---
'Cotiza y reserva Focus pro 3h para el lanzamiento, validando la política de cancelación antes de confirmar.'

=== Track 'book_focus' resuelto con run_pipeline (Módulo 3, sin cambios) ===
etapa 1 (cotizar) -> booking_agent
  tarea: 'Cotiza Focus pro 3h.'
  salida: '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?'
  salida: '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ó.'
  salida: '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}

6000 confirma el ancla del ecosistema (get_quote("Focus", "pro", 3)6000), y booking_id: 1 confirma que esta corrida arrancó sobre un Reservo desechable recién iniciado — sin ninguna reserva previa en reservo_tools.BOOKINGS. Nota algo que ya viste en el Módulo 3, pero que vale la pena repetir con este payload: nadie decidió, en ningún punto de esta ejecución, "¿cuál etapa sigue?" — el orden estaba fijo en PIPELINE_STAGES desde antes de que existiera esta petición concreta de Ana.


Lo que decidió el supervisor, y lo que decidió run_pipeline

Vale la pena separar con precisión dos decisiones que se ven parecidas pero son de naturaleza distinta:

El supervisor (concepto, lección 02) decidió:
  "esta parte de la petición necesita un PIPELINE, con estas 3 etapas"
  -- UNA decisión, tomada UNA vez, antes de que corriera nada.

run_pipeline (ejecutado, esta lección) decidió:
  NADA -- no tomó ninguna decisión de ruteo. Simplemente recorrió
  PIPELINE_STAGES en el orden en que ya estaba escrito, y para cada
  etapa, construyó la tarea con build_stage_task y la despachó al
  agente que esa etapa ya traía asignado (stage.name).

Esta distinción es la misma que ya viste, en otra forma, en el Módulo 3, lección 02: un pipeline no decide nada en tiempo de ejecución — toda la decisión (qué etapas, en qué orden, con qué agente) ya estaba tomada quando PIPELINE_STAGES se escribió. Lo que este módulo agrega es de dónde viene esa decisión: en el Módulo 3 aislado, PIPELINE_STAGES era simplemente "el pipeline de la lección"; acá, es la traducción directa de un Track que el supervisor etiquetó como pipeline en la lección 02.


Errores comunes

  1. Pensar que hace falta una versión nueva de run_pipeline para usarla "dentro de" una petición compuesta. No hace falta ninguna — la función es exactamente la misma, con la misma firma, recibiendo los mismos tres argumentos (stages, model_scripts, initial_payload). "Ser parte de una petición más grande" no cambia nada de cómo se ejecuta un pipeline por dentro.

  2. Olvidar que este INITIAL_PAYLOAD viene de leer la petición completa de Ana, no de un ejemplo aislado. {"room": "Focus", "tier": "pro", "hours": 3, "member": "Ana"} no apareció de la nada — es la lectura (concepto) de la frase "cotiza y reserva Focus pro 3h para el lanzamiento" dentro de la petición completa de la lección 01.

  3. Ejecutar 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 — esta lección, como todas las de la guía, asume un Reservo desechable, recién iniciado.

  4. Pensar que, porque esta sub-tarea usa un pipeline, el resto de la petición de Ana también necesita uno. No — la lección 02 ya etiquetó las otras dos partes (compare_rooms, boardroom_no_show) con patrones distintos. Cada Track se resuelve con el mecanismo que le corresponde a ÉL, no al que resolvió al primero.


Ejercicios

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

Ejecuta el ejemplo trabajado de esta lección tú mismo y confirma, línea por línea, que tu salida coincide con el "Qué esperar". Presta especial atención a booking_id: 1 — si tu Reservo no arrancó limpio, va a salir un número distinto.

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 track book_focus corrió sin desvíos.

Ejercicio 2: Resuelve el Track book_focus del Ejercicio 2 de la lección 02 (Medio)

En la lección 02, el Ejercicio 2 te pedía aplicar el criterio a la petición de Diego: "reserva el Focus pro 2h para la entrevista de mañana, validando la política de cancelación antes de confirmar". Construye el INITIAL_PAYLOAD y los guiones correspondientes, y ejecuta ese track con run_pipeline, sin cambiar la función.

Ver solución
INITIAL_PAYLOAD_D = {"room": "Focus", "tier": "pro", "hours": 2, "member": "Diego"}
script_quote_d = [
    {"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_confirm_d = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "book_room",
         "input": {"room": "Focus", "tier": "pro", "hours": 2, "member": "Diego"}}]},
    {"stop_reason": "end_turn", "content": [{"type": "text", "text": "Reservé Focus pro 2h para Diego (confirmación #2)."}]},
]
payload_d, trace_d = run_pipeline(
    PIPELINE_STAGES, [script_quote_d, model_script_policy, script_confirm_d], INITIAL_PAYLOAD_D,
)
print("precio calculado a mano:", 2500 * 2 * 80 // 100)
print("payload final:", payload_d)

Salida esperada (en el mismo proceso donde ya corriste el ejemplo trabajado, así que el booking_id sigue desde 1):

precio calculado a mano: 4000
payload final: {'room': 'Focus', 'tier': 'pro', 'hours': 2, 'member': 'Diego', 'price_cents': 4000, '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 que arma build_stage_task sigue conteniendo "cancelación" — la keyword que reconoce el stub de search_docs. 4000 = 2500 * 2 * 80 // 100 confirma el precio; booking_id: 2 confirma que esta corrida comparte el proceso con el ejemplo trabajado, que ya había reservado Focus con booking_id: 1.

Ejercicio 3: ¿Qué pasaría si book_focus se etiquetara, por error, como fanout? (Difícil)

Sin ejecutar código todavía, explica por qué correr las tres etapas de book_focusquote, validate_policy, confirm— como si fueran independientes entre sí (por ejemplo, con run_specialist suelto para cada una, todas al mismo tiempo) produciría un resultado incorrecto, y en qué etapa específica se notaría el error primero.

Ver solución

La etapa confirm depende de dos cosas que solo existen después de que las etapas anteriores corrieron: payload['price_cents'] (de quote) no se usa directamente en confirm, pero payload.get("cleared_to_book") —que sí determina qué texto arma build_stage_task para la etapa confirm— depende de que validate_policy ya haya corrido y haya escrito cleared_to_book: True. Si las tres etapas corrieran "a la vez", como si fueran tracks independientes de un fan-out, la etapa confirm tendría que construirse ANTES de que validate_policy hubiera terminado —porque, en un fan-out real, ningún job espera a otro—. En ese caso, payload.get("cleared_to_book") devolvería None (la clave todavía no existiría), y build_stage_task armaría la tarea SIN la frase "la política de cancelación ya se validó" — exactamente el mismo comportamiento que ya viste, ejecutado, en el Ejercicio 3 de la lección 04 del Módulo 3, cuando se saltaba la etapa de política a propósito. El error se notaría primero en el TEXTO de la tarea de confirmación, no en un traceback — el pipeline seguiría corriendo, pero sin la garantía de que la política ya se validó antes de reservar, que es justamente el punto por el que esta sub-tarea necesitaba un pipeline y no un fan-out.


Resumen y siguiente paso

  • El Track book_focus, etiquetado pipeline en la lección 02, se resolvió con run_pipeline del Módulo 3, sin ningún cambio de código.
  • El supervisor decidió QUÉ patrón le correspondía a esta sub-tarea (una vez, antes de ejecutar nada) — run_pipeline no decidió nada en tiempo de ejecución, solo recorrió etapas ya fijas.
  • Confirmado ejecutando: Focus pro 3h = 6000 centavos, booking_id: 1, y la etapa de confirmación mencionando la política ya validada — la misma garantía de orden que el Módulo 3 construyó, ahora aplicada como una pieza de una petición más grande.

Siguiente lección: 04 — Las preguntas independientes corren en fan-out mientras el pipeline resuelve. Agregamos el segundo Track de PLANcompare_rooms— y lo corremos AL MISMO TIEMPO que book_focus, con un mecanismo de fan-out generalizado a "cualquier callable", no solo a agentes sueltos.


Recursos adicionales

  1. Anthropic — Building effective agents — El patrón "prompt chaining" reusado, sin cambios, como una pieza de una orquestación más grande.
  2. Anthropic — Messages API reference — La forma exacta de tool_use/tool_result/stop_reason que cada etapa de este pipeline sigue respetando, sin cambios, dentro de una petición compuesta.
  3. Python — dataclasses — El módulo detrás de PipelineStage y Track, reusados juntos en esta lección.
  4. Anthropic — Multi-agent research system — Un flujo con etapas obligatorias que forma parte de una tarea más grande, exactamente el rol que juega book_focus dentro de la petición completa de Ana.