Módulo 4: Fan-out paralelo y agregación

Mini-proyecto: fan-out en Reservo

Descripción

Siete lecciones dejaron el patrón fan-out completo: cómo confirmar independencia (02), el caso base con orden fijo (03), la agregación determinista (04), la variante dentro de un solo agente (05), la extensión con concurrencia real (06), y el ahorro medido en rondas (07). Este mini-proyecto no agrega ningún concepto nuevo — te da tres escenarios de Reservo que nunca viste y te pide aplicar el patrón completo: reconocer qué tipo de fan-out corresponde (entre agentes, dentro de un agente, o ninguno), despacharlo con concurrencia real cuando aplica, y citar el costo de coordinación de cada uno.

El escenario más importante del lote es el B — a propósito, NO es fan-out entre agentes. Distinguir correctamente cuándo el mecanismo de las lecciones 03-06 NO hace falta es tan parte del criterio de este módulo como saber construirlo.

Conexión con el módulo

Este mini-proyecto es la síntesis de las siete lecciones anteriores, no una lección nueva. De la 02 usas el criterio de independencia. De las 03-04, SubTask y la agregación determinista. De la 05, el criterio para reconocer cuándo el fan-out es interno a un agente, no entre agentes. De la 06, run_fanout_parallel. De la 07, fanout_rounds para citar el ahorro de cada escenario. Cuando termines, el Módulo 5 toma este mismo Reservo y construye el patrón siguiente: handoff, donde un agente EN CURSO decide, a mitad de tarea, transferirle el control a otro.


El encargo

Reservo te pasa tres peticiones que llegaron la misma semana:

Escenario A: "Cotiza Studio pro 2h y dime si hay cargo por no presentarme."
Escenario B: "Compara Focus, Studio y Boardroom, todos pro, 3h."
Escenario C: "Reserva Boardroom pro 4h para Marta, dime la política de cancelación,
              y compara Focus, Studio y Boardroom, todos pro, 3h."

Tu encargo tiene tres partes:

a) Para cada escenario, decide qué patrón corresponde: fan-out entre agentes (lecciones 03-06), fan-out dentro de un solo agente (lección 05), o ninguno de los dos.

b) Despacha cada escenario con el mecanismo correcto, citando el historial y la respuesta de cada sub-tarea (o del agente único, si no hace falta repartir entre varios).

c) Cita el costo de coordinación de cada escenario — llamadas al modelo, hops, y rondas secuencial vs. paralelo, usando fanout_rounds de la lección 07 donde aplique.


La solución completa (el entregable)

Ver la solución completa
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,
        },
    },
    "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)


def count_model_calls(history):
    return sum(1 for m in history if m["role"] == "assistant")


def count_tool_calls(history):
    total = 0
    for m in history:
        if isinstance(m["content"], list):
            total += sum(1 for b in m["content"] if b["type"] == "tool_use")
    return total


@dataclass
class SubTask:
    agent: str
    task: str


def run_fanout_parallel(subtasks, model_scripts):
    results = {}
    with concurrent.futures.ThreadPoolExecutor(max_workers=len(subtasks)) as pool:
        future_to_subtask = {
            pool.submit(run_specialist, sub.agent, sub.task, model_scripts[sub.agent]): sub
            for sub in subtasks
        }
        for future in concurrent.futures.as_completed(future_to_subtask):
            sub = future_to_subtask[future]
            final, history = future.result()
            results[sub.agent] = {"task": sub.task, "output": final["content"][0]["text"], "history": history}
    return results


def fanout_rounds(call_counts_list, surrounding_calls):
    sequential = sum(call_counts_list) + surrounding_calls
    parallel = max(call_counts_list) + surrounding_calls
    return sequential, parallel


def report(label, results, split_calls=1, compose_calls=1):
    print(f"=== {label} ===")
    for agent in sorted(results):
        print(f"[{agent}] tarea: {results[agent]['task']!r}")
        print(f"[{agent}] salida: {results[agent]['output']!r}")
    specialist_calls = sum(count_model_calls(results[a]["history"]) for a in results)
    specialist_tools = sum(count_tool_calls(results[a]["history"]) for a in results)
    call_counts = [count_model_calls(results[a]["history"]) for a in results]
    total_calls = split_calls + specialist_calls + compose_calls
    seq_rounds = split_calls + sum(call_counts) + compose_calls
    par_rounds = split_calls + max(call_counts) + compose_calls
    print(f"llamadas al modelo TOTAL: {total_calls}  |  llamadas a tools: {specialist_tools}")
    print(f"rondas secuencial: {seq_rounds}  |  rondas paralelo: {par_rounds}  |  ahorradas: {seq_rounds - par_rounds}")
    print()


# --- Escenario A: cotizar Y política de no-presentación -- fan-out ENTRE dos agentes ---
task_a = "Cotiza Studio pro 2h y dime si hay cargo por no presentarme."
subtasks_a = [
    SubTask(agent="booking_agent", task="Cotiza Studio pro 2h."),
    SubTask(agent="policy_agent", task="¿Hay cargo por no presentarme a mi reserva?"),
]
model_scripts_a = {
    "booking_agent": [
        {"stop_reason": "tool_use", "content": [
            {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
             "input": {"room": "Studio", "tier": "pro", "hours": 2}}]},
        {"stop_reason": "end_turn", "content": [
            {"type": "text", "text": "Studio pro 2h cuesta 6400 centavos."}]},
    ],
    "policy_agent": [
        {"stop_reason": "tool_use", "content": [
            {"type": "tool_use", "id": "toolu_01", "name": "search_docs",
             "input": {"query": "no me presento a mi reserva"}}]},
        {"stop_reason": "end_turn", "content": [
            {"type": "text", "text": (
                "Sí -- si no te presentas y no cancelas con al menos 2 horas "
                "de anticipación, se cobra el 50% del precio cotizado."
            )}]},
    ],
}
print("--- petición compuesta A ---")
print(repr(task_a))
results_a = run_fanout_parallel(subtasks_a, model_scripts_a)
report("Escenario A (fan-out real: booking_agent + policy_agent)", results_a)


# --- Escenario B: comparar TRES salas -- fan-out DENTRO de un solo agente ---
task_b = "Compara Focus, Studio y Boardroom, todos pro, 3h."
model_script_pricing_b = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
         "input": {"room": "Focus", "tier": "pro", "hours": 3}},
        {"type": "tool_use", "id": "toolu_02", "name": "get_quote",
         "input": {"room": "Studio", "tier": "pro", "hours": 3}},
        {"type": "tool_use", "id": "toolu_03", "name": "get_quote",
         "input": {"room": "Boardroom", "tier": "pro", "hours": 3}},
    ]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": (
            "Focus pro 3h: 6000 centavos. Studio pro 3h: 9600 centavos. "
            "Boardroom pro 3h: 19200 centavos. Focus es la opción más "
            "barata de las tres."
        )}]},
]
print("--- petición B (NO es fan-out entre agentes) ---")
print(repr(task_b))
final_b, hist_b = run_specialist("pricing_agent", task_b, model_script_pricing_b)
print("¿hace falta run_fanout_sequential/run_fanout_parallel? No -- un solo agente "
      "(pricing_agent) resuelve todo, con fan-out DENTRO de su propio turno (lección 05).")
print("respuesta:", final_b["content"][0]["text"])
print(f"llamadas al modelo: {count_model_calls(hist_b)}  |  llamadas a tools: {count_tool_calls(hist_b)}")
print()


# --- Escenario C: reservar + política + comparar -- fan-out de TRES agentes ---
task_c = ("Reserva Boardroom pro 4h para Marta, dime la política de cancelación, "
          "y compara Focus, Studio y Boardroom, todos pro, 3h.")
subtasks_c = [
    SubTask(agent="booking_agent", task="Reserva Boardroom pro 4h para Marta."),
    SubTask(agent="policy_agent", task="¿Cuál es la política de cancelación?"),
    SubTask(agent="pricing_agent", task="Compara Focus, Studio y Boardroom, todos pro, 3h."),
]
model_scripts_c = {
    "booking_agent": [
        {"stop_reason": "tool_use", "content": [
            {"type": "tool_use", "id": "toolu_01", "name": "book_room",
             "input": {"room": "Boardroom", "tier": "pro", "hours": 4, "member": "Marta"}}]},
        {"stop_reason": "end_turn", "content": [
            {"type": "text", "text": "Reservé Boardroom pro 4h para Marta (confirmación #1), 25600 centavos."}]},
    ],
    "policy_agent": [
        {"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. Después de eso aplica el cargo de no-presentación."
            )}]},
    ],
    "pricing_agent": model_script_pricing_b,
}
print("--- petición compuesta C (fan-out de TRES agentes) ---")
print(repr(task_c))
results_c = run_fanout_parallel(subtasks_c, model_scripts_c)
report("Escenario C (fan-out real: booking_agent + policy_agent + pricing_agent)", results_c)

print("--- confirmación de la reserva creada en el Escenario C ---")
print(rt.BOOKINGS)

Qué esperar:

--- petición compuesta A ---
'Cotiza Studio pro 2h y dime si hay cargo por no presentarme.'
=== Escenario A (fan-out real: booking_agent + policy_agent) ===
[booking_agent] tarea: 'Cotiza Studio pro 2h.'
[booking_agent] salida: 'Studio pro 2h cuesta 6400 centavos.'
[policy_agent] tarea: '¿Hay cargo por no presentarme a mi reserva?'
[policy_agent] salida: 'Sí -- si no te presentas y no cancelas con al menos 2 horas de anticipación, se cobra el 50% del precio cotizado.'
llamadas al modelo TOTAL: 6  |  llamadas a tools: 2
rondas secuencial: 6  |  rondas paralelo: 4  |  ahorradas: 2

--- petición B (NO es fan-out entre agentes) ---
'Compara Focus, Studio y Boardroom, todos pro, 3h.'
¿hace falta run_fanout_sequential/run_fanout_parallel? No -- un solo agente (pricing_agent) resuelve todo, con fan-out DENTRO de su propio turno (lección 05).
respuesta: Focus pro 3h: 6000 centavos. Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. Focus es la opción más barata de las tres.
llamadas al modelo: 2  |  llamadas a tools: 3

--- petición compuesta C (fan-out de TRES agentes) ---
'Reserva Boardroom pro 4h para Marta, dime la política de cancelación, y compara Focus, Studio y Boardroom, todos pro, 3h.'
=== Escenario C (fan-out real: booking_agent + policy_agent + pricing_agent) ===
[booking_agent] tarea: 'Reserva Boardroom pro 4h para Marta.'
[booking_agent] salida: 'Reservé Boardroom pro 4h para Marta (confirmación #1), 25600 centavos.'
[policy_agent] tarea: '¿Cuál es la política de cancelación?'
[policy_agent] salida: 'Puedes cancelar sin cargo hasta 2 horas antes del horario reservado. Después de eso aplica el cargo de no-presentación.'
[pricing_agent] tarea: 'Compara Focus, Studio y Boardroom, todos pro, 3h.'
[pricing_agent] salida: 'Focus pro 3h: 6000 centavos. Studio pro 3h: 9600 centavos. Boardroom pro 3h: 19200 centavos. Focus es la opción más barata de las tres.'
llamadas al modelo TOTAL: 8  |  llamadas a tools: 5
rondas secuencial: 8  |  rondas paralelo: 4  |  ahorradas: 4

--- confirmación de la reserva creada en el Escenario C ---
{1: {'booking_id': 1, 'room': 'Boardroom', 'tier': 'pro', 'hours': 4, 'member': 'Marta', 'price_cents': 25600}}

El razonamiento por escenario:

Escenario A — fan-out entre dos agentes, sin sorpresas. "Cotiza Studio pro 2h" y "¿hay cargo por no presentarme?" son dos preguntas completamente separadas, a especialistas distintos — el mismo patrón exacto de las lecciones 03-06. 6400 = 4000 * 2 * 80 // 100. Costo: 6 llamadas, 4 rondas en paralelo (ahorrando 2 frente a las 6 que costaría en secuencial).

Escenario B — la trampa del lote, y el caso que NO es fan-out entre agentes. Aunque "compara Focus, Studio y Boardroom" reparte trabajo independiente (ninguna cotización depende de otra), las tres viven dentro del mismo turno del mismo agente (pricing_agent) — exactamente la variante de la lección 05, no el mecanismo de las lecciones 03-06. Envolver esto en run_fanout_sequential o run_fanout_parallel no estaría mal técnicamente (el Ejercicio 3 de la lección 05 ya confirmó que con una sola sub-tarea el resultado es idéntico) — pero sería agregar código innecesario para un caso que dispatch_parallel, ya construido en agent-fundamentals, resuelve solo.

Escenario C — fan-out real de tres agentes, con una reserva de verdad adentro. A diferencia de los escenarios A y B, acá booking_agent no solo cotiza — reserva de verdad, con book_room. Eso no cambia nada del mecanismo de fan-out: sigue siendo una sub-tarea independiente de las otras dos (ni la política de cancelación ni la comparación de precios necesitan que la reserva de Marta ya exista). El costo escala como predijo la fórmula de la lección 07: con tres sub-tareas de 2 llamadas cada una, 8 rondas en secuencial contra 4 en paralelo — el doble de ahorro que el Escenario A, con una sub-tarea más.


Errores comunes

  1. Forzar el Escenario B a pasar por run_fanout_sequential/run_fanout_parallel "para ser consistente" con A y C. El punto de este mini-proyecto es justamente reconocer cuándo el mecanismo NO hace falta — envolver el Escenario B de todos modos no produce un resultado incorrecto, pero agrega código sin ningún beneficio.

  2. Confundir el Escenario B (fan-out dentro de un agente) con "no hay ningún paralelismo". Sí lo hay — las tres cotizaciones de pricing_agent corren en el mismo turno, repartidas por dispatch_parallel. Lo que NO hay es fan-out entre agentes, que es lo que este módulo construyó de cero.

  3. Pensar que el Escenario C "falló" porque tuvo más llamadas que A. No falló — tiene una sub-tarea más (tres especialistas en vez de dos), así que es esperable que cueste más llamadas en total. Lo relevante no es el total absoluto, es la comparación secuencial-vs-paralelo dentro del MISMO escenario.

  4. Ejecutar los tres escenarios en el mismo proceso sin reiniciar el estado entre corridas. Si corres este mini-proyecto después de otro ejemplo de la guía en el mismo intérprete, el booking_id del Escenario C puede no salir 1 — cada escenario de esta guía asume su propio Reservo desechable, recién iniciado.

  5. Olvidar sorted(results) al reportar el Escenario C. Con tres agentes en fan-out real, la tentación de imprimir en el orden en que run_fanout_parallel los devolvió es mayor que con dos — la disciplina de la lección 06 aplica exactamente igual, sin importar cuántas sub-tareas haya.


Ejercicios

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

Ejecuta los tres escenarios del encargo tú mismo y confirma, línea por línea, que tu salida coincide con la solución completa. Para el Escenario B, escribe en una frase adicional por qué NO usaste run_fanout_parallel, a pesar de que la petición también reparte trabajo independiente.

Ver solución

No hay una única "solución de código" para la primera parte — es una verificación: si tu salida coincide con la de la solución completa, tu Reservo desechable arrancó limpio.

Sobre el Escenario B: la independencia de las tres cotizaciones es real, pero vive dentro del turno de un único agente (pricing_agent), no entre agentes distintos con historiales separados — dispatch_parallel, ya construido en agent-fundamentals, resuelve ese paralelismo sin necesitar ninguna pieza de este módulo.

Ejercicio 2: Un cuarto escenario, D (Medio)

Diseña un cuarto escenario de Reservo: "Cotiza Boardroom pro 3h y dime si puedo cancelar sin cargo con 3 horas de anticipación." Decide qué patrón corresponde y ejecútalo con el mecanismo correcto.

Ver solución
task_ex2 = "Cotiza Boardroom pro 3h y dime si puedo cancelar sin cargo con 3 horas de anticipación."
subtasks_ex2 = [
    SubTask(agent="booking_agent", task="Cotiza Boardroom pro 3h."),
    SubTask(agent="policy_agent", task="¿Puedo cancelar sin cargo con 3 horas de anticipación?"),
]
model_scripts_ex2 = {
    "booking_agent": [
        {"stop_reason": "tool_use", "content": [
            {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
             "input": {"room": "Boardroom", "tier": "pro", "hours": 3}}]},
        {"stop_reason": "end_turn", "content": [{"type": "text", "text": "Boardroom pro 3h cuesta 19200 centavos."}]},
    ],
    "policy_agent": [
        {"stop_reason": "tool_use", "content": [
            {"type": "tool_use", "id": "toolu_01", "name": "search_docs",
             "input": {"query": "cancelación con anticipación"}}]},
        {"stop_reason": "end_turn", "content": [
            {"type": "text", "text": "Sí -- 3 horas de anticipación superan las 2 horas mínimas, así que no aplica cargo."}]},
    ],
}
print("petición:", repr(task_ex2))
results_ex2 = run_fanout_parallel(subtasks_ex2, model_scripts_ex2)
for agent in sorted(results_ex2):
    print(f"[{agent}] {results_ex2[agent]['output']}")
print("Boardroom pro 3h a mano:", 8000 * 3 * 80 // 100)

Salida esperada:

petición: 'Cotiza Boardroom pro 3h y dime si puedo cancelar sin cargo con 3 horas de anticipación.'
[booking_agent] Boardroom pro 3h cuesta 19200 centavos.
[policy_agent] Sí -- 3 horas de anticipación superan las 2 horas mínimas, así que no aplica cargo.
Boardroom pro 3h a mano: 19200

Explicación: dos preguntas independientes, a dos especialistas distintos — el mismo patrón del Escenario A, con datos nuevos. 19200 = 8000 * 3 * 80 // 100.

Ejercicio 3: Verifica el ahorro de rondas del Escenario C con la fórmula (Difícil)

Usando fanout_rounds de la lección 07, predice las rondas secuencial y paralelo del Escenario C (tres sub-tareas, cada una con 2 llamadas internas, surrounding_calls=2) y confirma que coincide con lo medido en la solución completa.

Ver solución
call_counts_c = [2, 2, 2]  # booking_agent, policy_agent, pricing_agent -- 2 llamadas cada uno
seq_c, par_c = fanout_rounds(call_counts_c, surrounding_calls=2)
print(f"predicho con la fórmula: secuencial={seq_c}, paralelo={par_c}, ahorradas={seq_c - par_c}")
print("medido en el Escenario C del encargo: secuencial=8, paralelo=4, ahorradas=4")
print("¿coinciden?", (seq_c, par_c) == (8, 4))

Salida esperada:

predicho con la fórmula: secuencial=8, paralelo=4, ahorradas=4
medido en el Escenario C del encargo: secuencial=8, paralelo=4, ahorradas=4
¿coinciden? True

Explicación: la fórmula de la lección 07 predice exactamente lo que report() midió sobre el Escenario C real — confirma que el modelo de costo (sum() para secuencial, max() para paralelo, más los surrounding_calls fijos) generaliza sin ajustes al caso de tres agentes, no solo al de dos que usaron las lecciones 03-07.


Resumen y siguiente paso

  • El mini-proyecto no agregó ningún concepto nuevo: aplicó las siete lecciones anteriores —criterio de independencia, caso base, agregación determinista, la variante dentro de un agente, concurrencia real, y el ahorro medido en rondas— sobre tres escenarios de Reservo nuevos.
  • Dos de los tres escenarios (A, C) fueron fan-out real entre agentes, con costo de coordinación medido y citado; el Escenario B confirmó, sobre un caso nuevo, que "trabajo independiente" no siempre significa "fan-out entre agentes" — a veces vive dentro del turno de uno solo.
  • El Escenario C mostró que el patrón escala sin cambios a tres agentes, con una reserva real (book_room) como una de las sub-tareas — el fan-out no distingue entre sub-tareas de solo lectura y sub-tareas con efectos reales, siempre que sean genuinamente independientes.
  • Con este módulo completo, tienes tres de los cinco patrones de la guía construidos: supervisor (M2, decide quién), pipeline (M3, encadena sin decidir), fan-out (este módulo, reparte y agrega).

Con esto termina el Módulo 4. Construiste el patrón fan-out paralelo completo: el criterio de independencia, el caso base con orden fijo, la agregación determinista, la distinción con el fan-out interno de un agente, la extensión con concurrencia real, y la medición del ahorro en rondas. En el Módulo 5 construimos el cuarto patrón: handoff y delegación — cuando un agente que YA está trabajando se da cuenta, a mitad de tarea, de que necesita a otro especialista, y le transfiere el control directamente, sin volver a un supervisor externo.


Recursos adicionales

  1. Anthropic — Building effective agents — El patrón "parallelization" completo: sectioning (repartir sub-tareas) y aggregator (combinarlas), el eje de todo este mini-proyecto.
  2. Anthropic — Multi-agent research system — Un caso real donde reconocer qué trabajo es genuinamente paralelizable —y qué trabajo no lo es— determinó el diseño final del sistema, igual que el Escenario B de este encargo.
  3. Python — concurrent.futures — El módulo detrás de run_fanout_parallel, corriendo sin cambios sobre dos y sobre tres sub-tareas en este mini-proyecto.
  4. Python — Diccionarios y funciones — La estructura detrás de results y SPECIALISTS, la base de todo el fan-out de este módulo.