State Management & Recovery: Durability & Consistency¶
Overview¶
Agents maintain state across multiple steps. If an agent crashes mid-execution, all progress is lost unless state is persisted.
State management is the difference between a recoverable system and data loss.
Agent State & Checkpointing¶
What is Agent State?¶
class AgentState:
"""All data needed to resume execution"""
def __init__(self):
self.components = {
'task': {}, # What agent is solving
'memory': {}, # What agent remembers
'plan': [], # Steps planned
'progress': [], # Steps completed
'context': {}, # Current context
'tool_results': {}, # Previous tool outputs
'decisions': [], # Decisions made
'timestamp': None # When captured
}
Checkpointing Strategies¶
Strategy 1: Full Snapshot¶
class FullSnapshotCheckpoint:
"""Save complete state at each step"""
def checkpoint(self, agent):
"""Save entire agent state"""
checkpoint = {
'id': generate_id(),
'timestamp': time.time(),
'agent_state': deep_copy(agent.state),
'memory': deep_copy(agent.memory.dump()),
'context': deep_copy(agent.context),
'step': agent.step_count
}
# Persist to durable storage
self.storage.save(checkpoint)
return checkpoint['id']
def recover(self, checkpoint_id):
"""Restore from checkpoint"""
checkpoint = self.storage.load(checkpoint_id)
agent = Agent()
agent.restore_from_checkpoint(checkpoint)
return agent
Pros: Complete recovery, no data loss
Cons: Large storage, slow checkpoints
Strategy 2: Incremental Checkpoint¶
class IncrementalCheckpoint:
"""Save only what changed"""
def __init__(self):
self.last_checkpoint = None
self.changes = []
def checkpoint(self, agent):
"""Save only changed data"""
changes = self.detect_changes(agent)
checkpoint = {
'id': generate_id(),
'parent': self.last_checkpoint, # Delta from previous
'changes': changes,
'timestamp': time.time()
}
self.storage.save(checkpoint)
self.last_checkpoint = checkpoint['id']
return checkpoint['id']
def recover(self, checkpoint_id):
"""Reconstruct from checkpoint chain"""
# Walk back chain of checkpoints
chain = []
current = checkpoint_id
while current:
checkpoint = self.storage.load(current)
chain.append(checkpoint)
current = checkpoint.get('parent')
# Apply changes in order
agent = Agent()
for checkpoint in reversed(chain):
agent.apply_changes(checkpoint['changes'])
return agent
Pros: Smaller storage, faster
Cons: More complex recovery
Strategy 3: Event Sourcing¶
class EventSourcingCheckpoint:
"""Save events, replay to reconstruct state"""
def __init__(self):
self.events = []
def record_event(self, event):
"""Record action as immutable event"""
event_record = {
'type': event.type,
'data': event.data,
'timestamp': time.time(),
'sequence': len(self.events)
}
self.storage.append_event(event_record)
self.events.append(event_record)
def recover(self, up_to_sequence=None):
"""Replay events to reconstruct state"""
events = self.storage.get_events(up_to_sequence)
agent = Agent()
for event in events:
agent.process_event(event)
return agent
Pros: Complete audit trail, easy to debug
Cons: Can be slow on large event streams
Recovery from Failures¶
Recovery Pattern: Checkpoint + Retry¶
class CheckpointedExecution:
def execute_with_recovery(self, task):
"""Execute with checkpoints for recovery"""
agent = Agent()
checkpoint_id = None
try:
# Step 1: Setup
agent.initialize(task)
checkpoint_id = self.checkpoint_manager.checkpoint(agent)
# Step 2: Execute
while not agent.is_done():
agent.step()
# Checkpoint after each step
checkpoint_id = self.checkpoint_manager.checkpoint(agent)
except Exception as e:
# Failed! Recover and retry
if checkpoint_id:
agent = self.checkpoint_manager.recover(checkpoint_id)
# Wait before retry (exponential backoff)
wait_time = self.backoff_strategy(e)
time.sleep(wait_time)
# Retry the step
try:
agent.step()
except Exception as retry_error:
# Still failing? Give up
self.log_failure(task, retry_error)
raise
else:
raise
return agent.result
Distributed State Coordination¶
Multi-Agent State¶
class DistributedStateCoordination:
"""Coordinate state across multiple agents"""
def __init__(self):
self.state_store = RedisStateStore() # Distributed
def get_shared_state(self, agent_id):
"""Agents read from shared store"""
state = self.state_store.get(f"agent:{agent_id}")
return state
def update_shared_state(self, agent_id, updates):
"""Atomic updates to shared state"""
# Use transactions to prevent race conditions
with self.state_store.transaction():
current = self.state_store.get(f"agent:{agent_id}")
current.update(updates)
self.state_store.set(f"agent:{agent_id}", current)
def acquire_lock(self, resource_id, timeout=10):
"""Prevent concurrent modifications"""
lock = self.state_store.lock(f"lock:{resource_id}", timeout)
return lock
Saga Pattern for Multi-Step Workflows¶
class SagaOrchestration:
"""Coordinate multi-step workflows with rollback"""
def execute_saga(self, saga_definition):
"""Execute steps with compensations"""
completed_steps = []
try:
for step in saga_definition.steps:
# Execute step
result = self.execute_step(step)
completed_steps.append((step, result))
# Checkpoint after each step
self.save_checkpoint(completed_steps)
except Exception as e:
# Rollback in reverse order
for step, result in reversed(completed_steps):
if step.has_compensation:
self.execute_compensation(step, result)
raise SagaFailure(f"Saga failed at step: {step}")
return result
def execute_step(self, step):
"""Execute single saga step"""
for attempt in range(step.max_retries):
try:
return step.action()
except TransientError:
time.sleep(2 ** attempt)
except PermanentError:
raise
def execute_compensation(self, step, previous_result):
"""Undo side effects of failed step"""
try:
step.compensation(previous_result)
except Exception as e:
# Compensation failed! Manual intervention needed
self.alert_ops(f"Compensation failed for {step.name}: {e}")
3 Warnings ⚠️¶
Warning 1: Forgetting to Checkpoint¶
# ❌ WRONG
result = agent.execute(task)
# No checkpoints = full retry on failure
# Agent progresses to step 8 of 10
# Crashes! Start from scratch
# ✅ RIGHT
result = agent.execute_with_checkpoints(task)
# Crash at step 8 → recover from step 8 → finish
# Recovery time: 20 seconds instead of 2 minutes
Warning 2: Stale Checkpoints¶
# ❌ WRONG
old_checkpoint = checkpoint_store.load(old_id)
# But memory/memory store has changed!
agent = recover_from_checkpoint(old_checkpoint)
# Agent has outdated information
# ✅ RIGHT
latest_checkpoint = checkpoint_store.get_latest(agent_id)
agent = recover_from_checkpoint(latest_checkpoint)
# Always use most recent
Warning 3: Orphaned State¶
# ❌ WRONG
agent = Agent()
agent.execute(task)
# State saved in agent's memory
# Agent process dies
# State never persisted to durable storage
# ✅ RIGHT
agent = Agent()
state_manager = StateManager()
agent.execute(task)
state_manager.persist(agent.state)
# State always in durable storage
Last Updated: August 9, 2026