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
-
Anclar el payload en el texto libre del modelo, no en el
tool_result. Siextract_payloadleyerafinal["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 enagent-fundamentalsM5 L07. -
Usar
eval()en vez deast.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. -
Olvidar que el payload SOLO crece, nunca se reemplaza.
extract_payloadempieza connew_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 aroom/tier/hours/member, que la etapa 1 ya había dejado. -
Pensar que
build_stage_tasknecesita saber el orden completo del pipeline. No — cada rama delifsolo sabe construir la tarea de su propiokind, 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. -
Ejecutar
last_tool_resultsobre un historial sin ningúntool_result. Si una etapa terminara enend_turnsin haber pedido ninguna tool antes,last_tool_resultdevolveríaNone, y la línearesult["price_cents"]deextract_payloadfallaría con unTypeError. 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_taskarma la tarea de cada etapa leyendo SOLO el payload acumulado;extract_payloadlee eltool_resultreal de esa etapa —conast.literal_eval, nuncaeval()— y agrega el dato nuevo, sin perder lo que ya tenía.run_pipelinereemplaza elifde decisión del supervisor por un simpleforsobre 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
- Python —
ast.literal_eval— La función central de esta lección: reconstruye literales de Python de forma segura, sin ejecutar código arbitrario. - 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. - Python — Diccionarios — La estructura detrás del payload, que crece por copia (
dict(payload)) en cada etapa, nunca se reemplaza. - Anthropic — Building effective agents — El patrón "prompt chaining" con verificación programática entre pasos — la misma disciplina que
extract_payloadaplica al anclar cada dato en un resultado real.