I was on call at 02:13 am when the onboarding flow for a new customer stalled at step 4—our recommendation LLM kept hitting a *vector‑DB search timeout* while the billing microservice had already charged the user. The dashboard showed “partial success”: the UI displayed a recommendation, the invoice was sent, but the personal‑finance summary was blank. Debug logs were a mess of `ToolCallError: timeout` and `ContextWindowOverflow`. It turned into a full‑blown incident because the fallback we’d sketched on a whiteboard never made it past our test environment.
That night taught me a hard truth: **multi‑step AI agents are fragile by design, and partial failures will surface in production the moment you stitch together a LLM, a tool, and a database.** If you don’t plan for them, you’ll end up with angry users, exploding costs, and a night‑shift that lasts far too long.
- Treat every step in an agent pipeline as a potentially failing micro‑service.
- Use the Saga pattern + compensating transactions to roll back side effects.
- Persist checkpoints with LangGraph or Temporal; hydrate state on retry.
- Apply circuit breakers and idempotent retries with jitter to external tools.
- Bucket errors semantically, monitor SLIs, and set alerts on cost per successful run.
Before you start: Python ≥ 3.12, LangGraph 0.5+, Temporal Python SDK 1.3.0, OpenAI v1.45.0+ Assistant API, Pydantic 2.6, Redis 7.0 (for state store), and a basic understanding of async/await.
Partial failures in multi‑step AI agent pipelines: what you need to know
Partially failing AI pipelines require deliberate architectural patterns like the Saga pattern for rollback, event‑driven checkpointing for state recovery, and circuit breakers for fault isolation. The 2026 approach implements these with frameworks like LangGraph, focusing on semantic error classification, cost‑aware retries, and graceful degradation to maintain service reliability despite individual component failures.
The unavoidable reality: why AI agent pipelines fail partially
The 2026 state of multi‑step agent architecture
Two years ago most teams still glued together LangChain chains with a single `run()` call. In 2026 the consensus has shifted to **orchestrated, stateful graphs**—each node is a distinct tool or LLM call that can be retried, compensated, or replaced without tearing down the whole workflow. The most popular stacks now mix:
| Layer | Typical tech (2026) | Why it matters |
|---|---|---|
| Orchestration | **LangGraph**, Temporal, Microsoft Semantic Kernel | Guarantees exactly‑once execution, checkpointing, and replay. |
| Vector store | LlamaIndex v0.12+, Haystack 2.0 | Provides versioned embeddings and fallback search paths. |
| LLM access | OpenAI v1.45.0+ Assistant API, Anthropic, Groq | New streaming+tool‑calling semantics require explicit error handling. |
| Observability | Honeycomb, Helicone, Argilla | Surface semantic error buckets, not just HTTP codes. |
These components are **independent services**. A timeout in any one of them propagates downstream, turning a “nice‑to‑have” feature into a hard failure. The more steps you chain, the larger the surface area for partial breakdowns.
Beyond LLM hallucinations: new failure modes
The classic “hallucination” still haunts us, but we now see **four extra categories** that surface only in production pipelines:
| Failure mode | Symptom | Typical source |
|---|---|---|
| Structured output parsing failure | JSONDecodeError, missing fields | LLM returns malformed schema despite `response_format`. |
| Vector DB search degradation | `SearchTimeout`, low recall | High QPS, cold index, or hardware throttling. |
| Context window overflow after step N | `ContextWindowOverflow` error from OpenAI | Cumulative token budget exceeded due to untrimmed intermediate results. |
| Tool‑call side‑effect leakage | Duplicate payment, stale cache entry | Non‑idempotent external APIs without compensating transactions. |
Ignoring these leads to exactly the incident I described earlier: the LLM succeeded, the vector search failed, and the downstream billing system kept rolling.
Architectural patterns for resilience: 2026 best practices
Saga pattern with compensating transactions for AI agents
A saga breaks a long‑running workflow into **atomic steps**. If step k fails, the system runs **compensating actions** for steps 1…k‑1. In AI pipelines that means:
- **Persist intent** (user request) in a durable store.
- **Execute step 1** (e.g., fetch user profile). On success, write a *completion flag*.
- **Execute step 2** (LLM generation). If the LLM times out, trigger a *fallback generation* or **abort**.
- **Execute step 3** (tool call). On failure, run a *compensating transaction* (e.g., refund a provisional charge).
With **Temporal**, this looks like a workflow definition where each activity returns a result struct. If an activity raises `temporalio.exceptions.ActivityError`, Temporal automatically schedules the defined compensation activities.
# temporal_workflow.py - Temporal v1.3.0
from temporalio import workflow, activity
from pydantic import BaseModel
class Profile(BaseModel):
user_id: str
tier: str
# Compensation activities -------------------------------------------------
@activity.defn
async def refund_payment(payment_id: str) -> None:
# idempotent refund; should be safe to call multiple times
...
@activity.defn
async def rollback_vector_upsert(doc_id: str) -> None:
# delete the just‑added embedding
...
# Main saga ---------------------------------------------------------------
@workflow.defn
class RecommendationSaga:
@workflow.run
async def run(self, request: dict) -> dict:
profile: Profile = await workflow.execute_activity(
fetch_user_profile,
request["user_id"],
schedule_to_close_timeout=5,
)
# If profile fetch fails, Temporal will retry automatically (see retry config)
llm_out = await workflow.execute_activity(
generate_recommendation,
profile,
schedule_to_close_timeout=15,
)
try:
payment_id = await workflow.execute_activity(
charge_user,
request["order_id"],
llm_out["price"],
schedule_to_close_timeout=8,
)
except Exception as e:
# Compensate previous successful steps
await workflow.execute_activity(refund_payment, payment_id)
await workflow.execute_activity(rollback_vector_upsert, llm_out["doc_id"])
raise workflow.FailureError("Charge failed") from e
return {"status": "ok", "recommendation": llm_out["text"]}
The code uses **dependency‑injected activities** (you can replace `charge_user` with a mock during tests). Each activity has its own **retry policy** (exponential backoff + jitter) defined in `temporalio.common.RetryPolicy`. The saga guarantees that even if the LLM works but the payment service is flaky, we never leave a dangling charge.
Event‑driven checkpointing & state hydration
When you’re not on Temporal, **LangGraph** gives you first‑class checkpointing. After each node you call `graph.save_checkpoint(state)`; the state includes:
- LLM’s last message list
- Tool call history (`tool_name`, `args`, `output`)
- Intermediate embeddings (if any)
- A monotonically increasing `step_id`
# checkpointing_example.py - LangGraph v0.5.1
from langgraph.checkpoint import RedisCheckpoint
from langgraph import Graph
ckpt = RedisCheckpoint(url="redis://localhost:6379/0")
graph = Graph(name="recommendation")
graph.add_node("fetch_profile", fetch_profile)
graph.add_node("generate", generate_with_llm)
graph.add_node("charge", charge_user)
# Wire edges ---------------------------------------------------------------
graph.add_edge("fetch_profile", "generate")
graph.add_edge("generate", "charge")
@graph.entrypoint
async def run(request: dict):
ctx = {"request": request}
# Load last known checkpoint if this run is a retry
state = await ckpt.load(request["run_id"])
if state:
ctx.update(state) # hydrate previous context
# Execute graph, persisting after each node
for node in graph.traverse(ctx):
await node.run(ctx)
await ckpt.save(request["run_id"], ctx) # checkpoint
return ctx["final_result"]
If a downstream tool crashes, you **rehydrate** the graph from the last checkpoint and re‑run only the failing node. The rest of the pipeline stays untouched, saving both latency and money.
Circuit breakers and graceful feature‑degradation paths
External tools—search APIs, payment gateways, even LLM endpoints—can become unavailable. A **circuit breaker** prevents a flood of retries from exhausting your connection pool and lets you fall back to a cached or simplified path.
I wrote an extensive guide on circuit breakers in Go; the same ideas apply in Python. Below is a tiny but production‑ready implementation using the **circuitbreaker** PyPI package (v2.2.0) with exponential backoff and jitter:
# circuit_breaker.py - circuitbreaker 2.2.0
import random
from circuitbreaker import circuit, CircuitBreakerError
# 5‑second cooldown, max 3 failures before opening
@circuit(failure_threshold=3, recovery_timeout=5, expected_exception=Exception)
def call_vector_search(query: str) -> list:
# jittered back‑off is handled by the decorator automatically
if random.random() < 0.2: # simulate flakiness
raise TimeoutError("Vector DB timed out")
return ["doc1", "doc2"]
def safe_search(query: str) -> list:
try:
return call_vector_search(query)
except CircuitBreakerError:
# Graceful degradation: return cached results or a static fallback
return get_cached_results(query)
You can read more about the theory behind circuit breakers in the classic *Hystrix* paper; the Python library mirrors the same semantics. For an in‑depth look at wiring this into an orchestrator, see my **[Circuit Breaker Go: 5 Ways to Safeguard AI APIs (2026)](https://nileshblog.tech/?p=6938)**.
Step‑by‑step implementation guide with production code (2026)
Scenario 1: handling tool/API dependency failure
Imagine a recommendation pipeline that calls an external **pricing service** after the LLM suggests a product. The pricing API sometimes returns 429 due to rate limiting.
# pricing_service.py - httpx 0.27.0
import httpx
from tenacity import retry, stop_after_attempt, wait_exponential_jitter
client = httpx.AsyncClient(timeout=5)
@retry(stop=stop_after_attempt(4), wait=wait_exponential_jitter(multiplier=0.5))
async def fetch_price(product_id: str) -> float:
resp = await client.get(f"https://prices.example.com/v1/{product_id}")
resp.raise_for_status()
return resp.json()["price"]
The `tenacity` decorator gives us **idempotent retries** with jitter, which cuts down on thundering‑herd effects. If the retry policy exhausts, we fall back to a **cached price** stored in Redis. The fallback is wired into the saga as a compensating transaction that **does not charge** the user until a real price is confirmed.
Scenario 2: managing LLM context window overflows mid‑pipeline
Our chain builds a “conversation history” that grows with each step. By step 4 we hit the 128 k‑token limit of the OpenAI `gpt‑4o‑2024‑08‑06` model.
# context_manager.py - OpenAI v1.45.0+
from openai import OpenAI
from pydantic import BaseModel
client = OpenAI()
def trim_messages(messages: list[dict], max_tokens: int = 115_000) -> list[dict]:
"""Simple sliding‑window trim; keep system prompt + last N turns."""
token_count = sum(len(m["content"]) for m in messages) # rough estimate
while token_count > max_tokens and len(messages) > 2:
# drop the oldest user/assistant pair
messages.pop(1) # remove first user message
messages.pop(1) # remove its assistant reply
token_count = sum(len(m["content"]) for m in messages)
return messages
async def safe_chat_completion(messages):
safe = trim_messages(messages)
resp = await client.chat.completions.create(
model="gpt-4o-2024-08-06",
messages=safe,
temperature=0.7,
max_completion_tokens=2048,
)
return resp.choices[0].message
You can plug `safe_chat_completion` into any LangGraph node. The trimming logic is deterministic, enabling **replayability**—the same input always yields the same trimmed set, which is vital for checkpoint‑based retries.
Scenario 3: recovering from parallel step timeout cascades
A common pattern is to **fetch multiple knowledge sources in parallel** (e.g., a vector search, a SQL lookup, and an external API). If one branch hangs, the whole workflow stalls.
# parallel_fetch.py - asyncio 3.12, httpx 0.27.0
import asyncio
import httpx
async def fetch_vector(query):
# Could raise TimeoutError
...
async def fetch_sql(query):
...
async def fetch_external_api(query):
...
async def run_parallel(query):
tasks = [
asyncio.wait_for(fetch_vector(query), timeout=3),
asyncio.wait_for(fetch_sql(query), timeout=3),
asyncio.wait_for(fetch_external_api(query), timeout=3),
]
# gather with return_exceptions to avoid cascade cancel
results = await asyncio.gather(*tasks, return_exceptions=True)
# Separate successes from failures
successes = [r for r in results if not isinstance(r, Exception)]
failures = [r for r in results if isinstance(r, Exception)]
# Log failures with semantic bucket (see later)
for err in failures:
log_semantic_error(err)
# If critical path missing, raise a saga‑compatible exception
if not any(isinstance(r, dict) and r.get("type") == "essential" for r in successes):
raise RuntimeError("Critical data missing after parallel fetch")
return successes
By wrapping each branch in `asyncio.wait_for` **and** using `return_exceptions=True`, we prevent one timeout from cancelling the others—a common source of cascading