Módulo 7: Orchestrating The Full Reservo System

Las preguntas independientes corren en fan-out mientras el pipeline resuelve

Descripción

La lección 03 resolvió book_focus solo, de principio a fin. Esta lección agrega el segundo Track de PLANcompare_rooms, la comparación de Studio y Boardroom— y lo corre al mismo tiempo que book_focus, con el mismo mecanismo de ThreadPoolExecutor que ya usaste en el Módulo 4. La diferencia con run_fanout_parallel (M4) es mínima pero importante: ahí, cada sub-tarea del fan-out era siempre una llamada a run_specialist —un agente, una tarea—. Acá, una de las dos "sub-tareas" que corren en paralelo es un run_pipeline entero, de tres etapas y dos agentes. run_tracks_parallel, la única función nueva de esta lección, generaliza el mecanismo de ThreadPoolExecutor de "un agente por sub-tarea" a "un callable por sub-tarea" — sin cambiar en nada el principio de determinismo que ya conoces del Módulo 4.

Conexión con el módulo

Esta lección reusa, sin cambios, run_pipeline (03) y run_specialist (M2). Lo único nuevo es run_tracks_parallel, que extiende el mismo idioma de ThreadPoolExecutor + concurrent.futures.as_completed + sorted() que ya construiste en el Módulo 4, lección 06 — ahora aplicado a callables arbitrarios en vez de siempre a run_specialist. La lección 05 agrega un tercer job a este mismo mecanismo, uno que por dentro hace un handoff.


Analogía: dos mensajeros, uno de los cuales hace tres mandados en cadena

Vuelve a los dos mensajeros del Módulo 4. Hasta ahora, cada uno hacía un mandado —comprar pan, revisar el correo—, sencillo, de un solo paso. Esta lección manda al mismo tiempo a dos mensajeros otra vez, pero uno de ellos tiene un mandado de tres pasos encadenados: ir al banco, esperar la confirmación de un trámite, y recién después retirar un paquete que dependía de esa confirmación. El otro mensajero sigue con un mandado de un solo paso. A quien los mandó no le importa que uno tarde tres pasos y el otro uno solo — lo único que le importa es que ninguno depende del resultado del otro, así que puede mandarlos a la vez, sin que se estorben.


Ejemplo trabajado: run_tracks_parallel, con un pipeline y una llamada suelta a la vez

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}},
    "pricing_agent": {"tools": {"get_quote": rt.get_quote}},
}


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


def run_tracks_parallel(jobs):
    """Generaliza run_fanout_parallel (Módulo 4) de 'un agente por
    sub-tarea' a 'un callable por sub-tarea' -- mismo ThreadPoolExecutor,
    mismo as_completed, misma disciplina de leer con sorted(); lo único
    que cambia es que cada job puede ser, por dentro, cualquier cosa: un
    especialista suelto, o un pipeline entero de varias etapas."""
    results = {}
    with concurrent.futures.ThreadPoolExecutor(max_workers=len(jobs)) as pool:
        future_to_key = {pool.submit(fn): key for key, fn in jobs.items()}
        for future in concurrent.futures.as_completed(future_to_key):
            key = future_to_key[future]
            results[key] = future.result()
    return results


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


def job_book_focus():
    return run_pipeline(
        PIPELINE_STAGES,
        [model_script_quote, model_script_policy, model_script_confirm],
        INITIAL_PAYLOAD,
    )


model_script_pricing = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
         "input": {"room": "Studio", "tier": "pro", "hours": 3}},
        {"type": "tool_use", "id": "toolu_02", "name": "get_quote",
         "input": {"room": "Boardroom", "tier": "pro", "hours": 3}},
    ]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": (
            "Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. "
            "Studio es la opción más barata de las dos."
        )}]},
]


def job_compare_rooms():
    return run_specialist("pricing_agent", "Compara Studio y Boardroom pro 3h.", model_script_pricing)


print("--- dos tracks, sin ninguna dependencia entre sí ---")
print("book_focus (pipeline, 3 etapas) y compare_rooms (una llamada) no comparten ningún dato")
print("de entrada ni de salida -- compare_rooms no necesita que Focus ya esté reservado.")

print()
print("=== corriendo los dos jobs a la vez con run_tracks_parallel ===")
jobs = {"book_focus": job_book_focus, "compare_rooms": job_compare_rooms}
results = run_tracks_parallel(jobs)
print("claves según llegaron (informativo, NO se usa para imprimir):", list(results.keys()))

print()
print("=== salida SIEMPRE en el mismo orden (sorted por key del track) ===")
for key in sorted(results):
    print(f"--- {key} ---")
    if key == "book_focus":
        payload, trace = results[key]
        print("  payload final:", payload)
    else:
        final, history = results[key]
        print("  respuesta:", final["content"][0]["text"])

Qué esperar (corrida real; el orden de llegada de las claves puede variar entre ejecuciones):

--- dos tracks, sin ninguna dependencia entre sí ---
book_focus (pipeline, 3 etapas) y compare_rooms (una llamada) no comparten ningún dato
de entrada ni de salida -- compare_rooms no necesita que Focus ya esté reservado.

=== corriendo los dos jobs a la vez con run_tracks_parallel ===
claves según llegaron (informativo, NO se usa para imprimir): ['compare_rooms', 'book_focus']

=== salida SIEMPRE en el mismo orden (sorted por key del track) ===
--- book_focus ---
  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}
--- compare_rooms ---
  respuesta: Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. Studio es la opción más barata de las dos.

compare_rooms —una sola llamada, con un stop_reason: "tool_use" seguido de un "end_turn"— terminó de correr en un tiempo real muchísimo menor que book_focus —tres etapas, cada una con su propio run_specialist—. Y aun así, ninguno de los dos "esperó" al otro: ThreadPoolExecutor los lanzó a los dos a la vez, cada uno en su propio hilo, y run_tracks_parallel recogió el resultado de cada uno apenas terminó, sin bloquear al más rápido con el más lento.


Por qué run_tracks_parallel no necesita saber qué hay "adentro" de cada job

Fíjate en la firma de run_tracks_parallel(jobs): recibe un diccionario de {key: callable}, y punto. No recibe SubTask, no recibe una lista de agentes, no recibe ningún PipelineStage — cada valor del diccionario es, para esta función, una caja negra que no toma argumentos y devuelve algo. Eso es exactamente lo que permite que job_book_focus (un pipeline de tres etapas) y job_compare_rooms (una llamada suelta) convivan en el mismo jobs sin que run_tracks_parallel tenga que distinguir entre los dos casos.

Esta es la generalización completa del mecanismo de fan-out del Módulo 4: run_fanout_parallel sabía, de antemano, que cada sub-tarea era "un agente resolviendo una tarea" —por eso su firma recibía subtasks (con .agent y .task) y model_scripts por separado, y armaba la llamada a run_specialist por dentro—. run_tracks_parallel no asume nada sobre la forma de cada job — el código que arma la tarea completa (job_book_focus, job_compare_rooms) vive afuera, en un def de cero argumentos, y run_tracks_parallel solo se encarga de correrlos concurrentemente y recoger sus resultados.


Errores comunes

  1. Pensar que run_tracks_parallel reemplaza a run_fanout_parallel del Módulo 4. No lo reemplaza — lo generaliza para un caso que el Módulo 4 no necesitaba todavía: cuando una de las "sub-tareas" del fan-out es, ella misma, un mecanismo de varios pasos. Si todas tus sub-tareas son agentes sueltos, como en el Módulo 4, run_fanout_parallel sigue siendo la opción más directa y legible.

  2. Olvidar envolver cada job en una función de cero argumentos. run_tracks_parallel llama a pool.submit(fn), sin argumentos — si intentas pasarle directamente run_pipeline(...) (la llamada ya evaluada, en vez de la función sin evaluar), Python ejecutaría el pipeline de inmediato, en el hilo principal, antes de que ThreadPoolExecutor pudiera hacer nada con él. Por eso job_book_focus es un def que "recuerda" sus argumentos por cierre (closure), no una llamada ya resuelta.

  3. Medir el tiempo real de esta corrida y sacar conclusiones sobre el ahorro. Como ya advirtió el Módulo 4, con guiones concepto que no hacen ninguna llamada de red real, el tiempo de reloj de esta lección no refleja ningún ahorro real de producción — el punto de esta lección es la composición del mecanismo, no una medición de latencia.

  4. Usar las claves de results en el orden de inserción, sin sorted(). Exactamente el mismo error del Módulo 4, lección 06 — el orden de inserción de un diccionario, bajo ThreadPoolExecutor, refleja el orden de finalización de los hilos, no ningún orden con significado.


Ejercicios

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

Ejecuta el ejemplo trabajado de esta lección tres veces seguidas, en el mismo proceso. Confirma que la salida ordenada (la sección "salida SIEMPRE en el mismo orden") es idéntica en las tres corridas, aunque el orden de results.keys() pueda variar.

Ver solución

No hay una única "solución de código" para este ejercicio — es una verificación: si las tres corridas producen la misma salida ordenada (aunque la línea de "claves según llegaron" varíe entre ellas), confirmaste con tus propios ojos que la generalización de esta lección conserva la misma garantía de determinismo del Módulo 4.

Ejercicio 2: Agrega un tercer job suelto, sin handoff todavía (Medio)

Antes de que la lección 05 agregue el Track boardroom_no_show (que sí hace handoff), agrega un tercer job simple: job_focus_availability, que le pregunta a pricing_agent "¿Cuánto sale reservar dos veces el Focus pro, 3h cada vez?" (una sola llamada a get_quote, sin ninguna dependencia con los otros dos). Corre los tres jobs juntos con run_tracks_parallel.

Ver solución
script_double_focus = [
    {"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": "Cada reserva de Focus pro 3h cuesta 6000 centavos; dos veces son 12000 centavos."}]},
]


def job_focus_availability():
    return run_specialist("pricing_agent", "¿Cuánto sale reservar dos veces el Focus pro, 3h cada vez?", script_double_focus)


jobs_3 = {"book_focus": job_book_focus, "compare_rooms": job_compare_rooms,
          "focus_availability": job_focus_availability}
results_3 = run_tracks_parallel(jobs_3)
for key in sorted(results_3):
    print(f"--- {key} ---")
    if key == "book_focus":
        print("  payload final:", results_3[key][0])
    else:
        print("  respuesta:", results_3[key][0]["content"][0]["text"])

Salida esperada (continuando en el mismo proceso del ejemplo trabajado, así que booking_id sigue en 1 si es la primera reserva del proceso):

--- book_focus ---
  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}
--- compare_rooms ---
  respuesta: Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. Studio es la opción más barata de las dos.
--- focus_availability ---
  respuesta: Cada reserva de Focus pro 3h cuesta 6000 centavos; dos veces son 12000 centavos.

Explicación: los tres jobs conviven sin ningún cambio a run_tracks_parallel —la función no tiene ningún límite hardcodeado de "dos jobs"—, y la salida ordenada alfabéticamente ubica book_focus primero, compare_rooms segundo, focus_availability tercero, sin importar el orden real de finalización de los tres hilos.

Ejercicio 3: ¿Qué pasa si max_workers es menor que la cantidad de jobs? (Difícil)

Modifica run_tracks_parallel para que reciba un max_workers explícito, y corre los dos jobs del ejemplo trabajado con max_workers=1 (forzando a que compartan un único hilo). Confirma que la salida agregada es idéntica a la del ejemplo trabajado.

Ver solución
def run_tracks_parallel_workers(jobs, max_workers):
    results = {}
    with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as pool:
        future_to_key = {pool.submit(fn): key for key, fn in jobs.items()}
        for future in concurrent.futures.as_completed(future_to_key):
            key = future_to_key[future]
            results[key] = future.result()
    return results


results_seq = run_tracks_parallel_workers(jobs, max_workers=1)
same = all(
    (results_seq[k][0] == results[k][0]) if k == "book_focus"
    else (results_seq[k][0]["content"][0]["text"] == results[k][0]["content"][0]["text"])
    for k in results
)
print("salida con max_workers=1 igual a la del ejemplo trabajado:", same)

Salida esperada:

salida con max_workers=1 igual a la del ejemplo trabajado: True

Explicación: exactamente la misma conclusión que ya confirmó el Módulo 4, lección 06 —el determinismo de la salida agregada no depende de cuánto paralelismo real haya, solo de que cada key del diccionario results se llene por completo antes de leerla con sorted()—. Con max_workers=1, los dos jobs terminan corriendo, en los hechos, uno después del otro dentro del mismo hilo del pool; el resultado agregado no cambia en absoluto.


Resumen y siguiente paso

  • run_tracks_parallel generaliza run_fanout_parallel (M4) de "un agente por sub-tarea" a "un callable por sub-tarea", con el mismo ThreadPoolExecutor + as_completed + sorted().
  • book_focus (un run_pipeline de tres etapas) y compare_rooms (una sola llamada) corrieron a la vez, sin que ninguno esperara al otro, porque ninguno depende del resultado del otro.
  • run_tracks_parallel no necesita saber qué hay "adentro" de cada job — solo que sea un callable de cero argumentos.
  • La salida agregada, leída con sorted(results), es idéntica sin importar el orden real de finalización de los hilos — la misma garantía de determinismo del Módulo 4, confirmada de nuevo.

Siguiente lección: 05 — Un agente hace handoff a mitad de una rama del fan-out. Agregamos el tercer Track de PLANboardroom_no_show— que por dentro cede el turno de booking_agent a policy_agent, sin que run_tracks_parallel note ninguna diferencia.


Recursos adicionales

  1. Python — concurrent.futures.ThreadPoolExecutor — La clase detrás de run_tracks_parallel, la misma que ya usa dispatch_parallel y run_fanout_parallel sin ningún cambio.
  2. Python — Funciones como objetos de primera clase — La base de por qué job_book_focus y job_compare_rooms pueden pasarse a pool.submit como cualquier otro valor.
  3. Anthropic — Building effective agents — El principio de tratar cada unidad de trabajo como una caja negra componible, sin que el orquestador necesite conocer su implementación interna.
  4. Anthropic — Multi-agent research system — Un orquestador real donde algunas sub-tareas paralelas son, ellas mismas, flujos de varios pasos — la misma composición que construye esta lección.