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 PLAN —compare_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
-
Pensar que
run_tracks_parallelreemplaza arun_fanout_paralleldel 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_parallelsigue siendo la opción más directa y legible. -
Olvidar envolver cada job en una función de cero argumentos.
run_tracks_parallelllama apool.submit(fn), sin argumentos — si intentas pasarle directamenterun_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 queThreadPoolExecutorpudiera hacer nada con él. Por esojob_book_focuses undefque "recuerda" sus argumentos por cierre (closure), no una llamada ya resuelta. -
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.
-
Usar las claves de
resultsen el orden de inserción, sinsorted(). Exactamente el mismo error del Módulo 4, lección 06 — el orden de inserción de un diccionario, bajoThreadPoolExecutor, 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_parallelgeneralizarun_fanout_parallel(M4) de "un agente por sub-tarea" a "un callable por sub-tarea", con el mismoThreadPoolExecutor+as_completed+sorted().book_focus(unrun_pipelinede tres etapas) ycompare_rooms(una sola llamada) corrieron a la vez, sin que ninguno esperara al otro, porque ninguno depende del resultado del otro.run_tracks_parallelno 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 PLAN —boardroom_no_show— que por dentro cede el turno de booking_agent a
policy_agent, sin que run_tracks_parallel note ninguna diferencia.
Recursos adicionales
- Python —
concurrent.futures.ThreadPoolExecutor— La clase detrás derun_tracks_parallel, la misma que ya usadispatch_parallelyrun_fanout_parallelsin ningún cambio. - Python — Funciones como objetos de primera clase — La base de por qué
job_book_focusyjob_compare_roomspueden pasarse apool.submitcomo cualquier otro valor. - 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.
- 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.