Módulo 3: Pipelines secuenciales

Encadenando el dato entre etapas

Descripción

La lección 02 cerró con un problema expuesto a propósito: armar el payload que la etapa 2 necesita, copiando price_cents a mano desde el tool_result de la etapa 1. Funciona para una etapa. No escala a tres, y mucho menos a un pipeline que alguien más tenga que mantener. Esta lección resuelve ese problema con dos funciones —build_stage_task, que arma la tarea de la etapa actual a partir de lo que dejaron las anteriores, y extract_payload, que lee el resultado real de esa etapa y agrega solo el dato nuevo que la siguiente necesita— generalizadas en run_pipeline, el runner que encadena N etapas sin volver a tocar nada a mano.

El punto más importante de esta lección no es el código en sí: es dónde se ancla el dato que viaja entre etapas. No se ancla en el texto libre que el modelo redactó —"Focus pro 3h cuesta 6000 centavos" es una frase, no un valor confiable para que la etapa 3 reserve—. Se ancla en el tool_result real que dispatch_parallel ya ejecutó, parseado con ast.literal_eval — nunca con eval() — para recuperar el dato estructurado exacto que la tool devolvió.

Conexión con el módulo

build_stage_task, extract_payload y run_pipeline son la pieza que sostiene el resto del módulo: la lección 04 los usa, sin ningún cambio, para correr el pipeline completo de tres etapas de Reservo; la lección 06 los extiende con una guarda que detiene el pipeline cuando una etapa no alcanza para continuar.


Analogía: la orden de trabajo que viaja con la pieza, no el recuerdo de quien la hizo

Vuelve a la línea de ensamblaje de la lección 02. Cuando una pieza pasa de la estación de corte a la de soldadura, no viaja con el recuerdo verbal de lo que hizo el operario de corte —"corté algo más o menos así"—; viaja con una orden de trabajo física, sujeta a la pieza, con las medidas exactas que el corte produjo. El operario de soldadura no necesita confiar en la memoria de nadie: lee la orden de trabajo, que es un dato verificable, no un relato. extract_payload construye exactamente esa orden de trabajo — un dato anclado en lo que la tool devolvió, no en lo que alguien dijo que hizo.


Ejemplo trabajado: el mecanismo completo, sobre un pipeline de dos etapas

Para aislar el mecanismo de encadenar, este ejemplo usa solo dos etapas —cotizar y confirmar—, sin la etapa de política todavía. La lección 04 agrega la tercera.

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})")


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",
    },
}


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):
    """Devuelve el dict o texto REAL del último tool_result del historial,
    parseado con ast.literal_eval -- NUNCA con eval() (nunca se ejecuta
    texto arbitrario). dispatch_parallel guarda cada resultado como un
    string (str(r)); si ese string es un literal de Python (un dict), se
    reconstruye; si ya era texto plano, se devuelve tal cual."""
    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):
    """Arma el texto que recibe la etapa actual, usando SOLO lo que dejaron
    las etapas anteriores en el payload -- ejecutado, sin ningún modelo de
    por medio. Nunca vuelve a preguntarle nada al socio."""
    if stage.kind == "quote":
        return f"Cotiza {payload['room']} {payload['tier']} {payload['hours']}h."
    if stage.kind == "confirm":
        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):
    """Lee el resultado REAL (no el texto libre del modelo) de la etapa
    que acaba de correr y agrega SOLO los datos nuevos que la etapa
    siguiente necesita. El payload nunca pierde lo que ya tenía -- solo
    crece."""
    new_payload = dict(payload)
    result = last_tool_result(history)
    if stage.kind == "quote":
        new_payload["price_cents"] = result["price_cents"]
    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):
    """Corre las etapas en el ORDEN FIJO de `stages` -- nunca decide cuál
    sigue, nunca vuelve a consultar cuál es la próxima. El payload que
    arma cada etapa es la ENTRADA de la etapa siguiente -- la propiedad
    central del patrón pipeline."""
    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,
            "output": final["content"][0]["text"],
        })
    return payload, trace


STAGES_SMALL = [
    PipelineStage(kind="quote", name="booking_agent", label="cotizar"),
    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_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)."}]},
]

payload, trace = run_pipeline(STAGES_SMALL, [model_script_quote, model_script_confirm], INITIAL_PAYLOAD)

for step in trace:
    print(f"--- etapa {step['stage']} ({step['label']}) -> {step['agent']} ---")
    print("  tarea armada:", repr(step["task"]))
    print("  salida:", repr(step["output"]))

print()
print("payload final:", payload)

Qué esperar:

--- etapa 1 (cotizar) -> booking_agent ---
  tarea armada: 'Cotiza Focus pro 3h.'
  salida: 'Focus pro 3h cuesta 6000 centavos.'
--- etapa 2 (confirmar la reserva) -> booking_agent ---
  tarea armada: 'Reserva Focus pro 3h para Ana.'
  salida: 'Reservé Focus pro 3h para Ana (confirmación #1).'

payload final: {'room': 'Focus', 'tier': 'pro', 'hours': 3, 'member': 'Ana', 'price_cents': 6000, 'booking_id': 1, 'confirmed': True}

Fíjate en la tarea armada de la etapa 2: 'Reserva Focus pro 3h para Ana.'build_stage_task la construyó leyendo payload['room'], payload['tier'], payload['hours'] y payload['member'], sin que nadie volviera a preguntarle nada al socio ni al modelo. El price_cents: 6000 del payload final tampoco vino del texto "Focus pro 3h cuesta 6000 centavos" —vino de parsear, con ast.literal_eval, el tool_result real que get_quote devolvió en la etapa 1—. Ese es el mecanismo completo: cada etapa lee el payload que dejaron las anteriores, y cada etapa deja un payload más grande para la siguiente.


Por qué ast.literal_eval y no eval()

dispatch_parallel guarda cada resultado de tool como str(r) — el mismo formato que exige el protocolo real de tool use, donde el content de un tool_result siempre es texto. Para recuperar el dict real detrás de ese texto ("{'price_cents': 6000}"{'price_cents': 6000}), hace falta parsear ese string de vuelta a un objeto de Python.

eval() haría eso —y mucho más—: ejecutaría cualquier expresión de Python contenida en el string, incluyendo código malicioso si ese string viniera de una fuente no confiable. ast.literal_eval, de la librería estándar, hace exactamente lo que este mecanismo necesita y nada más: reconstruye únicamente literales de Python —dicts, listas, números, strings, booleanos— y rechaza con una excepción cualquier cosa que no sea un literal seguro. En este módulo, el string siempre viene de reservo_tools, una fuente propia y confiable — pero la disciplina de usar ast.literal_eval en vez de eval() es la misma que aplicarías si ese dato viniera de una tool externa, y vale la pena practicarla desde acá.


run_pipeline: el for que reemplaza al if del supervisor

Compara run_pipeline con el dispatcher del Módulo 2, lección 06: ahí, la pieza central era route_deterministic(text) —una función que decide, leyendo la petición, a cuál especialista delegar—. Acá, la pieza central es un for i, stage in enumerate(stages) — no hay ninguna decisión adentro del loop, solo una iteración sobre una lista ya ordenada. run_pipeline nunca le pregunta a nadie "¿cuál sigue?" — la respuesta ya estaba en stages[i + 1] desde antes de que la primera petición llegara.


Errores comunes

  1. Anclar el payload en el texto libre del modelo, no en el tool_result. Si extract_payload leyera final["content"][0]["text"] para sacar el precio (por ejemplo, con una expresión regular buscando números), el pipeline quedaría atado a que el modelo redacte el texto de una forma exacta y predecible — frágil, y contrario al hábito de grounding que ya viste en agent-fundamentals M5 L07.

  2. Usar eval() en vez de ast.literal_eval. Aunque en este módulo el string siempre viene de una fuente propia, eval() ejecuta cualquier expresión de Python — un hábito peligroso que no vale la pena adoptar ni en un ejemplo controlado.

  3. Olvidar que el payload SOLO crece, nunca se reemplaza. extract_payload empieza con new_payload = dict(payload) — una copia completa del payload anterior— y le agrega claves nuevas. Si en cambio construyera un diccionario nuevo desde cero, la etapa 3 perdería acceso a room/tier/hours/member, que la etapa 1 ya había dejado.

  4. Pensar que build_stage_task necesita saber el orden completo del pipeline. No — cada rama del if solo sabe construir la tarea de su propio kind, sin ninguna referencia a qué etapa viene después. Esto es intencional: agregar una etapa nueva al final del pipeline no obliga a tocar las ramas de las etapas anteriores.

  5. Ejecutar last_tool_result sobre un historial sin ningún tool_result. Si una etapa terminara en end_turn sin haber pedido ninguna tool antes, last_tool_result devolvería None, y la línea result["price_cents"] de extract_payload fallaría con un TypeError. Esto solo pasaría si el guion de esa etapa estuviera mal escrito —la lección 06 profundiza en fallas de etapa de un tipo distinto, donde la tool SÍ corre pero su resultado no alcanza para continuar.


Ejercicios

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

Repite el pipeline de dos etapas de esta lección, pero con INITIAL_PAYLOAD = {"room": "Studio", "tier": "basic", "hours": 2, "member": "Carlos"} (Studio basic 2h = 4000 * 2 = 8000 centavos, sin descuento). Confirma el payload final.

Ver solución
INITIAL_PAYLOAD_STUDIO = {"room": "Studio", "tier": "basic", "hours": 2, "member": "Carlos"}
script_quote_studio = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "get_quote",
         "input": {"room": "Studio", "tier": "basic", "hours": 2}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Studio basic 2h cuesta 8000 centavos."}]},
]
script_confirm_studio = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "book_room",
         "input": {"room": "Studio", "tier": "basic", "hours": 2, "member": "Carlos"}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Reservé Studio basic 2h para Carlos (confirmación #1)."}]},
]
payload_studio, trace_studio = run_pipeline(
    STAGES_SMALL, [script_quote_studio, script_confirm_studio], INITIAL_PAYLOAD_STUDIO,
)
print("payload final:", payload_studio)

Salida esperada:

payload final: {'room': 'Studio', 'tier': 'basic', 'hours': 2, 'member': 'Carlos', 'price_cents': 8000, 'booking_id': 1, 'confirmed': True}

Explicación: el mecanismo es idéntico —build_stage_task y extract_payload no cambian—, solo cambia el INITIAL_PAYLOAD de entrada. 8000 = 4000 * 2 porque basic no aplica el descuento del 20% que sí aplica pro.

Ejercicio 2: Agrega extract_payload para un kind de cancelación (Medio)

Extiende extract_payload (en una copia nueva de la función) para que reconozca un tercer kind, "cancel", que llama a cancel_booking y agrega payload["cancelled"] con el resultado. Pruébala sobre el booking_id: 1 que dejó el ejemplo trabajado de esta lección.

Ver solución
def extract_payload_v2(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 == "confirm":
        new_payload["booking_id"] = result["booking_id"]
        new_payload["confirmed"] = result["confirmed"]
    elif stage.kind == "cancel":
        new_payload["cancelled"] = result["cancelled"]
    return new_payload


cancel_stage = PipelineStage(kind="cancel", name="booking_agent", label="cancelar la reserva")
script_cancel = [
    {"stop_reason": "tool_use", "content": [
        {"type": "tool_use", "id": "toolu_01", "name": "cancel_booking", "input": {"id": 1}}]},
    {"stop_reason": "end_turn", "content": [
        {"type": "text", "text": "Cancelé la reserva número 1."}]},
]
final_cancel, history_cancel = run_specialist(cancel_stage.name, "Cancela la reserva 1.", script_cancel)
payload_cancelled = extract_payload_v2(cancel_stage, history_cancel, payload)
print("payload con cancelación:", payload_cancelled)

Salida esperada:

payload con cancelación: {'room': 'Focus', 'tier': 'pro', 'hours': 3, 'member': 'Ana', 'price_cents': 6000, 'booking_id': 1, 'confirmed': True, 'cancelled': True}

Explicación: el patrón se repite sin sorpresas: una rama nueva de if, un kind nuevo, un dato nuevo agregado al payload que ya tenía todo lo anterior. Esto es lo que hace que extract_payload escale a pipelines más largos sin que las ramas existentes tengan que tocarse.

Ejercicio 3: ¿Qué pasa si build_stage_task recibe un kind que no reconoce? (Difícil)

Construye un PipelineStage(kind="notify", name="booking_agent", label="notificar") e intenta armar su tarea con build_stage_task. ¿Qué excepción se produce, y por qué esa falla —en vez de devolver un string vacío o genérico— es el comportamiento correcto?

Ver solución
notify_stage = PipelineStage(kind="notify", name="booking_agent", label="notificar")
try:
    build_stage_task(notify_stage, payload)
except ValueError as e:
    print(f"ValueError capturado: {e}")

Salida real:

ValueError capturado: no sé armar la tarea de la etapa 'notify'

Explicación: build_stage_task termina con un raise ValueError explícito para cualquier kind que ninguna rama del if reconozca — el mismo principio de "fallar ruidoso, no silencioso" que ya viste con el KeyError de SPECIALISTS en el Módulo 2. Si en cambio devolviera un string vacío o un texto genérico como "continúa", el pipeline seguiría corriendo con una tarea sin sentido para esa etapa, y el error real —un kind nuevo sin implementar— quedaría escondido hasta que alguien notara, mucho más tarde, que esa etapa nunca hizo lo que debía.


Resumen y siguiente paso

  • build_stage_task arma la tarea de cada etapa leyendo SOLO el payload acumulado; extract_payload lee el tool_result real de esa etapa —con ast.literal_eval, nunca eval()— y agrega el dato nuevo, sin perder lo que ya tenía.
  • run_pipeline reemplaza el if de decisión del supervisor por un simple for sobre una lista ya ordenada — no hay ninguna pregunta "¿cuál sigue?" adentro del loop.
  • Ejecutamos un pipeline de dos etapas —cotizar, confirmar— con el precio viajando de la primera a la segunda sin que nadie lo copiara a mano, resolviendo el problema que la lección 02 dejó pendiente.
  • El payload se ancla en el resultado real de cada tool, nunca en el texto libre del modelo — el mismo hábito de grounding que ya conocías, aplicado ahora al traspaso de datos entre agentes completos.

Siguiente lección: 04 — El pipeline de Reservo, ejecutado de punta a punta. Agregamos la etapa de validación de política en el medio y corremos las tres etapas completas, con el historial de cada especialista impreso.


Recursos adicionales

  1. Python — ast.literal_eval — La función central de esta lección: reconstruye literales de Python de forma segura, sin ejecutar código arbitrario.
  2. Anthropic — Tool use (function calling) overview — La forma exacta de tool_result (siempre texto) que hace necesario parsear el dato de vuelta a un objeto de Python.
  3. Python — Diccionarios — La estructura detrás del payload, que crece por copia (dict(payload)) en cada etapa, nunca se reemplaza.
  4. Anthropic — Building effective agents — El patrón "prompt chaining" con verificación programática entre pasos — la misma disciplina que extract_payload aplica al anclar cada dato en un resultado real.