Módulo 3: Pipelines secuenciales
Cuándo una etapa falla
Descripción
Las lecciones 03, 04 y 05 asumieron, sin decirlo, que cada etapa deja siempre lo que la siguiente
necesita: quote siempre entrega un price_cents, validate_policy siempre entrega un
cleared_to_book. Mira de nuevo extract_payload de la lección 04: la rama de validate_policy
pone new_payload["cleared_to_book"] = True sin ninguna condición — como si el resultado de
search_docs siempre alcanzara para confirmar. No siempre alcanza. El stub de policy_agent tiene
solo dos entradas (cancellation-policy, no-show-policy); cualquier pregunta que no matchee
ninguna de las dos devuelve "No se encontró una política relevante para esa pregunta." — un
resultado real, sin ninguna excepción, que no es una base válida para confirmar una reserva.
Esta lección construye la pieza que faltaba: una guarda que revisa, después de cada etapa, si lo
que dejó alcanza para seguir — y que detiene el pipeline ahí mismo, antes de que la etapa confirm
llegue a correr sobre una política sin validar. Es la propiedad que un supervisor tiene casi gratis
—puede devolver None y no delegar nada, como ya viste en el Módulo 2— y que un pipeline, por no
decidir nunca nada, tiene que ganarse a propósito.
Conexión con el módulo
Esta lección reusa, sin cambios, PipelineStage, build_stage_task y PIPELINE_STAGES de las
lecciones 02 a 04. Extiende extract_payload con una condición real en vez de un True incondicional,
y agrega dos piezas nuevas: stage_ok (la guarda) y run_pipeline_guarded (el runner que la aplica).
La lección 07 retoma esta misma distinción —qué depende genuinamente de qué— desde otro ángulo: si
quote y validate_policy no dependen entre sí, la mitad del riesgo que esta lección cubre podría,
en realidad, correr en paralelo (Módulo 4). El mini-proyecto de la lección 08 usa run_pipeline_guarded
completo, sin ningún cambio, sobre tres escenarios nuevos.
Analogía: el control de calidad entre estaciones, no solo al final
Vuelve a la línea de ensamblaje. Hasta ahora, cada estación entregaba su pieza a la siguiente sin que
nadie la revisara en el camino — se asumía que, si la estación de corte terminó, la pieza está lista
para soldar. Una línea de ensamblaje seria no asume eso: entre cada par de estaciones hay un control
de calidad puntual, no solo al final de toda la línea. Si el control detecta que la pieza que sale de
soldadura tiene una fisura, la línea se detiene ahí — no sigue mandando esa pieza a pintura,
esperando que alguien la note en la inspección final. stage_ok es ese control de calidad puntual:
revisa, etapa por etapa, no solo al final.
Ejemplo trabajado: la guarda, primero en el camino feliz, después deteniendo el pipeline
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 = {
"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):
"""STUB de policy_agent: solo dos entradas. Cualquier query que no
matchee ninguna devuelve el mensaje de 'no encontrado' -- un resultado
REAL, no una excepción."""
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,
},
"expertise": "cotizar, reservar y cancelar salas",
},
"policy_agent": {
"tools": {"search_docs": search_docs},
"expertise": "responder preguntas de política (cancelación, no-presentación)",
},
}
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":
if payload.get("client_type") == "corporate":
return (f"¿Qué política aplica para una reserva corporativa grande de "
f"{payload['room']}, antes de confirmarla para {payload['member']}?")
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["policy_result"] = result
# LA GUARDA de esta lección: cleared_to_book SOLO es True si
# search_docs encontró una política real -- se reconoce porque el
# stub la prefija con "[algo]". En las lecciones 03/04, esta rama
# ponía cleared_to_book = True sin ninguna condición.
new_payload["cleared_to_book"] = isinstance(result, str) and result.startswith("[")
elif stage.kind == "confirm":
new_payload["booking_id"] = result["booking_id"]
new_payload["confirmed"] = result["confirmed"]
return new_payload
class PipelineHalted(Exception):
"""Se lanza cuando una etapa deja un payload que no alcanza para que
la etapa siguiente continúe con seguridad. Lleva el trace PARCIAL --
lo que sí alcanzó a correr -- para que quien la capture pueda
inspeccionar qué pasó antes de detenerse."""
def __init__(self, reason, trace):
super().__init__(reason)
self.trace = trace
def stage_ok(stage, payload):
"""La guarda: ejecutada, sin ningún modelo de por medio. Decide si el
payload que dejó ESTA etapa alcanza para seguir con la próxima."""
if stage.kind == "validate_policy" and not payload.get("cleared_to_book"):
return False, "no se encontró una política relevante -- no hay base para confirmar"
return True, None
def run_pipeline_guarded(stages, model_scripts, initial_payload):
"""Igual que run_pipeline (lecciones 03/04), MÁS una guarda: después
de extraer el payload de cada etapa, verifica con stage_ok si alcanza
para seguir. Si no alcanza, detiene el pipeline ahí mismo -- nunca
llega a construir la tarea de la etapa siguiente."""
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"],
})
ok, reason = stage_ok(stage, payload)
if not ok:
raise PipelineHalted(
f"pipeline detenido después de la etapa {i + 1} ({stage.label}): {reason}",
trace,
)
return payload, trace
# El PIPELINE_STAGES de siempre, sin ningún cambio. La etiqueta
# "validar la política de cancelación" describe el ROL de la etapa 2 --
# no cambia aunque la pregunta real (build_stage_task) sea otra para un
# cliente corporativo, como ves más abajo.
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"),
]
print("--- primero, confirmamos que la guarda NO estorba el camino feliz ---")
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)."}]},
]
payload, trace = run_pipeline_guarded(
PIPELINE_STAGES, [model_script_quote, model_script_policy, model_script_confirm], INITIAL_PAYLOAD,
)
print("cleared_to_book:", payload["cleared_to_book"], "| booking_id:", payload["booking_id"])
print()
print("--- ahora, una reserva corporativa que activa la guarda ---")
INITIAL_PAYLOAD_CORP = {
"room": "Boardroom", "tier": "pro", "hours": 6, "member": "TechCorp",
"client_type": "corporate",
}
model_script_quote_corp = [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "get_quote",
"input": {"room": "Boardroom", "tier": "pro", "hours": 6}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "Boardroom pro 6h cuesta 38400 centavos."}]},
]
model_script_policy_miss = [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "search_docs",
"input": {"query": "política para reservas corporativas grandes"}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": (
"No encontré una política específica para reservas corporativas "
"grandes en nuestra base -- te recomiendo escribir al equipo de "
"cuentas corporativas antes de confirmar."
)}]},
]
# Solo 2 guiones: la etapa 3 (confirmar) nunca llega a correr.
model_scripts_corp = [model_script_quote_corp, model_script_policy_miss]
try:
payload_corp, trace_corp = run_pipeline_guarded(
PIPELINE_STAGES, model_scripts_corp, INITIAL_PAYLOAD_CORP,
)
except PipelineHalted as e:
print(f"PipelineHalted capturado: {e}")
trace_corp = e.trace
print()
print("--- el trace parcial: la etapa 1 SÍ corrió, la etapa 3 nunca se intentó ---")
for step in trace_corp:
print(f"etapa {step['stage']} ({step['label']}) -> {step['agent']}: {step['output']!r}")
print("BOOKINGS después del intento corporativo:", rt.BOOKINGS)
Qué esperar:
--- primero, confirmamos que la guarda NO estorba el camino feliz ---
cleared_to_book: True | booking_id: 1
--- ahora, una reserva corporativa que activa la guarda ---
PipelineHalted capturado: pipeline detenido después de la etapa 2 (validar la política de cancelación): no se encontró una política relevante -- no hay base para confirmar
--- el trace parcial: la etapa 1 SÍ corrió, la etapa 3 nunca se intentó ---
etapa 1 (cotizar) -> booking_agent: 'Boardroom pro 6h cuesta 38400 centavos.'
etapa 2 (validar la política de cancelación) -> policy_agent: 'No encontré una política específica para reservas corporativas grandes en nuestra base -- te recomiendo escribir al equipo de cuentas corporativas antes de confirmar.'
BOOKINGS después del intento corporativo: {1: {'booking_id': 1, 'room': 'Focus', 'tier': 'pro', 'hours': 3, 'member': 'Ana', 'price_cents': 6000}}
Dos cosas para notar. Primero: el camino feliz (Focus pro 3h para Ana) produce exactamente el mismo
resultado que ya viste en la lección 04 —cleared_to_book: True, booking_id: 1— la guarda no le
agrega ningún costo ni ningún paso extra cuando todo va bien; solo actúa cuando hace falta. Segundo:
BOOKINGS al final solo tiene una reserva —la de Ana—. El intento corporativo nunca llegó a
book_room: la excepción se lanzó apenas stage_ok detectó que cleared_to_book seguía siendo
False después de la etapa 2, antes de que run_pipeline_guarded construyera siquiera la tarea de
la etapa 3.
Por qué la excepción lleva el trace parcial
Fíjate en PipelineHalted.__init__: recibe reason y trace, no solo reason. Si la excepción
solo llevara el mensaje, quien la capture sabría que el pipeline se detuvo, pero no qué sí
alcanzó a pasar antes — información valiosa para decidir qué hacer después (¿reintentar solo la
etapa 2? ¿escalar a un humano con el contexto ya armado?). Guardar el trace parcial en la excepción
misma evita el problema que ya viste con el NameError al capturar mal el resultado de una función
que lanza antes de retornar: sin el trace adentro de la excepción, la variable trace_corp nunca
llegaría a existir, porque run_pipeline_guarded nunca alcanza su return.
Errores comunes
-
Pensar que un pipeline sin excepciones significa un pipeline sin fallas. El ejemplo trabajado de esta lección corrió sin lanzar ningún
KeyErrorniTypeErroren el camino feliz — pero el escenario corporativo también corrió "sin errores de Python" hasta questage_okdecidió, a propósito, lanzarPipelineHalted. La ausencia de una excepción de Python no prueba que el resultado de una etapa haya alcanzado para seguir. -
Confundir
cleared_to_book = Falsecon un bug. No lo es — es exactamente el resultado correcto cuandosearch_docsno encuentra nada relevante. El bug estaba en las lecciones 03/04, que poníanTruesin condición, no en que el resultado real de la tool sea, a veces, un "no encontrado". -
Poner la guarda dentro de
build_stage_tasken vez de después deextract_payload. La guarda necesita el payload YA actualizado con lo que dejó la etapa que acaba de correr — revisarlo antes, dentro debuild_stage_task, llegaría demasiado tarde: esa función arma la tarea de la etapa actual, no decide si la etapa anterior alcanzó. -
Olvidar que
stage_oksolo conoce elkindde la etapa que acaba de terminar, no el pipeline completo. Igual quebuild_stage_task(lección 03, error común 4), cada guarda nueva que agregues solo necesita saber de su propiokind— no hace falta tocar la guarda devalidate_policypara agregar, por ejemplo, una guarda paraquote. -
Pensar que esta guarda cubre TODOS los modos de falla posibles. No — el Ejercicio 3 de esta lección construye, a propósito, un guion mal escrito que la guarda no detecta: una etapa que nunca pide ninguna tool. Conocer los límites de una guarda es tan importante como construirla.
Ejercicios
Ejercicio 1: Confirma tu propia ejecución (Fácil)
Ejecuta el ejemplo trabajado completo tú mismo y confirma, línea por línea, que tu salida coincide
con el "Qué esperar" de arriba. Presta atención a BOOKINGS al final: debería tener una sola
reserva, no dos.
Ver solución
No hay una única "solución de código" para este ejercicio — es una verificación: si BOOKINGS al
final tiene exactamente una entrada (booking_id: 1, la de Ana), confirmaste que el intento
corporativo nunca llegó a book_room — la prueba más directa de que la guarda detuvo el pipeline
donde dijo que lo detendría.
Ejercicio 2: Activa la guarda SIN el camino corporativo (Medio)
stage_ok no sabe nada de client_type — solo mira cleared_to_book. Confirma esto: usa el
INITIAL_PAYLOAD estándar (sin client_type="corporate") para un socio que pregunta, por
error, algo que el stub de search_docs tampoco cubre —por ejemplo, "¿cuáles son los horarios de
acceso al edificio?"— y confirma que run_pipeline_guarded se detiene igual.
Ver solución
payload_std = {"room": "Focus", "tier": "pro", "hours": 2, "member": "Marcos"}
script_quote_std = [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "get_quote",
"input": {"room": "Focus", "tier": "pro", "hours": 2}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "Focus pro 2h cuesta 4000 centavos."}]},
]
script_policy_miss_std = [
{"stop_reason": "tool_use", "content": [
{"type": "tool_use", "id": "toolu_01", "name": "search_docs",
"input": {"query": "horarios de acceso al edificio"}}]},
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "No encontré una política de horarios de acceso en nuestra base."}]},
]
try:
run_pipeline_guarded(PIPELINE_STAGES, [script_quote_std, script_policy_miss_std], payload_std)
except PipelineHalted as e:
print(f"PipelineHalted capturado: {e}")
Salida real:
PipelineHalted capturado: pipeline detenido después de la etapa 2 (validar la política de cancelación): no se encontró una política relevante -- no hay base para confirmar
Explicación: el pipeline se detiene exactamente igual, sin ningún client_type="corporate" en
el payload. Esto confirma el punto central de la lección: la guarda mira el resultado real de
search_docs —si empieza con "[", si trae una etiqueta de política reconocida—, no una bandera
que alguien puso de antemano. Cualquier pregunta que el stub no cubra activa la misma guarda, sin
importar por qué el socio la hizo.
Ejercicio 3: Un guion mal escrito que la guarda NO detecta (Difícil)
stage_ok solo revisa cleared_to_book después de validate_policy. Construye un
model_script_quote roto a propósito —uno que termina en end_turn sin pedir get_quote
primero— y ejecuta la etapa quote con él a través de run_pipeline_guarded. ¿Qué excepción se
produce, y por qué NO es un PipelineHalted limpio como el del ejemplo trabajado?
Ver solución
payload_bad = {"room": "Focus", "tier": "pro", "hours": 2, "member": "Nora"}
script_quote_bad = [
{"stop_reason": "end_turn", "content": [
{"type": "text", "text": "Cotizando Focus pro 2h..."}]},
]
try:
run_pipeline_guarded(PIPELINE_STAGES, [script_quote_bad], payload_bad)
except TypeError as e:
print(f"TypeError (NO PipelineHalted) capturado: {e!r}")
Salida real:
TypeError (NO PipelineHalted) capturado: TypeError("'NoneType' object is not subscriptable")
Explicación: este guion nunca pide get_quote —termina en end_turn de inmediato—, así que
last_tool_result no encuentra ningún tool_result en el historial y devuelve None (exactamente
el caso que la lección 03, error común 5, ya había señalado como posible). extract_payload intenta
entonces result["price_cents"] sobre None, y Python lanza un TypeError antes de que
stage_ok tenga oportunidad de correr — la guarda de esta lección solo revisa el payload después
de que extract_payload termina con éxito; si extract_payload mismo falla, la guarda nunca se
ejecuta. Esta es una distinción importante y honesta: stage_ok protege contra resultados válidos
que no alcanzan (el modo de falla que diseñó esta lección), no contra guiones mal escritos que
rompen el mecanismo de extracción antes de llegar ahí — ese segundo modo de falla necesitaría una
guarda distinta, en un punto distinto del código, fuera del alcance de esta lección.
Resumen y siguiente paso
- Las lecciones 03/04 asumían, sin decirlo, que
validate_policysiempre alcanza para confirmar —extract_payloadponíacleared_to_book = Truesin ninguna condición. Esta lección corrige eso:cleared_to_bookahora depende del resultado REAL desearch_docs. stage_okes la guarda: revisa, después de cada etapa, si lo que dejó alcanza para seguir.run_pipeline_guardedla aplica, y lanzaPipelineHalted—con el trace parcial adentro— apenas una etapa no alcanza, sin llegar a construir la tarea de la etapa siguiente.- Confirmamos, con salida real, que la guarda no le agrega costo al camino feliz (el mismo
booking_id: 1de la lección 04) y que sí detiene, a tiempo, una reserva corporativa sin política validada — sin queBOOKINGStermine con una reserva que nunca debió confirmarse. - El Ejercicio 3 dejó un límite honesto: la guarda protege contra resultados válidos que no alcanzan, no contra un guion tan mal escrito que rompe la extracción del dato antes de que la guarda tenga oportunidad de correr.
Siguiente lección: 07 — Eligiendo pipeline o supervisor. Con el riesgo de un pipeline ya medido (lección 05) y su fragilidad ya cubierta (esta lección), cerramos el criterio completo — y confirmamos, ejecutando el código, que dos de las tres etapas de este módulo no dependían entre sí desde el principio.
Recursos adicionales
- Anthropic — Building effective agents — El principio de agregar puertas de verificación programáticas entre pasos de una cadena, exactamente lo que
stage_okconstruye para este pipeline. - Python — Excepciones definidas por el usuario — La base de
PipelineHalted, que además de heredar deExceptionlleva su propio atributo (trace) para no perder el estado parcial. - Anthropic — Multi-agent research system — Un sistema real donde detectar, a tiempo, que un sub-resultado no alcanza evita que el error se propague silenciosamente hasta el resultado final.
- Python —
try/excepty el flujo de control — El mecanismo que separa, en esta lección, el "camino feliz" de la detención a tiempo.