When a shopper says *“Add the silk blouse in size M, and keep it in my cart while I look at the shoes”* the voice‑AI has to remember that choice, keep the cart alive while the user asks unrelated questions, and survive a possible network hiccup or a server restart. In my last deployment, a silent container crash erased the in‑memory checkout state, leaving the customer dangling in the middle of a purchase. The fix? Separate the *reasoning* (Claude) from the *source of truth* (persistent state) and make the two speak through an async, fault‑tolerant pipeline.
That’s what this guide builds: a production‑ready voice‑AI retail agent that never loses a cart, can pick up a conversation hours later, and stays under the latency budget required for natural‑speech turn‑taking. We’ll stitch together Python 3.12, Anthropic’s Claude 3.5/4.0, a Redis 7.2 cache, PostgreSQL 16 with JSONB, and a Streamlit 1.35 desktop UI for monitoring.
—
- Use a Redis + PostgreSQL hybrid layer to balance speed and durability.
- Model agent state with Pydantic v2 and store it as JSONB rows.
- Leverage Claude’s tool‑calling to read/write state without leaving the LLM’s prompt.
- Wrap every I/O call in async retries and exponential back‑off.
- Garbage‑collect idle sessions and provide GDPR‑compliant forget‑functions.
Before you start: Python 3.12+, Anthropic Python SDK v1.4+, Redis 7.2+, PostgreSQL 16, SQLAlchemy 2.0, Pydantic v2, Streamlit 1.35, Docker 20.10+. You’ll need an Anthropic API key, a Redis endpoint (or Docker container) and a PostgreSQL instance. Optional: LangChain 0.2 for orchestration.
Voice‑AI Agent State Management for Long Retail Conversations
Voice‑AI agent state management for long retail conversations involves persisting shopping context, cart data, and user intent across multiple sessions. A 2026 architecture uses a hybrid model: Redis for fast, in‑memory state during active dialogue and PostgreSQL for durable storage and recovery, synchronized via Python’s asyncio and integrated with LLMs like Claude via tool‑calling APIs.
Core Concepts: Defining Agent State for Retail
Conversation Thread & Memory
A *thread* is the ordered list of user utterances and system replies. We keep the raw transcript (for audit) and a condensed summary that fits Claude’s 100 k token window.
User Context & Purchase Intent
Context bundles demographics, loyalty tier, and inferred intent (browsing vs checkout). This influences product recommendations without hard‑coding business logic.
Session, Cart, and Transaction State
- **Session ID** – UUID that survives client disconnects.
- **Cart** – List of `{sku, quantity, price}` objects stored as JSONB.
- **Transaction** – Payment token, shipping address, order status.
External System Handles (ERP, CRM)
When the user finalizes checkout, the agent must fire a **create‑order** call to the ERP. The handle (API token, endpoint) lives in a separate encrypted table; the agent never writes it back to the cache.
Architectural Blueprint: Hybrid State Management Model (2026)
| Layer | Responsibility | Tech |
|---|---|---|
| **In‑Memory Cache** | Turn‑by‑turn read/write, sub‑millisecond latency | Redis 7.2 (TTL, Pub/Sub) |
| **Durable Store** | Long‑term persistence, recovery after crash | PostgreSQL 16 (JSONB, row‑level security) |
| **Sync Service** | Async batch flush, conflict resolution | asyncio, aioredis, asyncpg |
| **LLM Engine** | Stateless reasoning, tool‑calling | Anthropic Claude 3.5/4.0 via SDK |
| **API Gateway** | Expose `/webhook` for voice platform | FastAPI 0.110 (WebSocket) |
| **Observability** | Metrics, logs, trace IDs | OpenTelemetry, Grafana |
flowchart LR
User[Voice Client] -->|WebSocket| API[FastAPI Gateway]
API -->|Tool Call| Claude[Claude LLM]
Claude -->|Read/Write| Redis[Redis Cache]
Claude -->|Persist| PG[PostgreSQL DB]
Redis -->|Pub/Sub| Sync[Sync Service]
Sync -->|Batch Write| PG
PG -->|Recovery| API
API -->|State Stream| Dashboard[Streamlit UI]
**My take:** The *source of truth* must always be the relational store; Redis is merely a speed‑boost. Treat the cache as a write‑through buffer, not a source of permanent data.
Step‑by‑Step Build: Python Agent with Claude & Persistent State
1. Environment Setup
# python 3.12+ required
python -m venv .venv
source .venv/bin/activate
pip install "anthropic==1.4.*" \
"redis==7.2.*" \
"asyncpg==0.29.*" \
"sqlalchemy==2.0.*" \
"pydantic==2.*" \
"fastapi==0.110.*" \
"uvicorn[standard]" \
"streamlit==1.35.*"
Make sure the Anthropic key is exported:
export ANTHROPIC_API_KEY=sk-ant-...
2. Designing the State Schema (Pydantic Models)
# version: 2026-10-03
from pydantic import BaseModel, Field
from typing import List, Optional
from uuid import UUID, uuid4
from datetime import datetime
class CartItem(BaseModel):
sku: str
quantity: int = Field(gt=0)
price_cents: int
class SessionState(BaseModel):
session_id: UUID = Field(default_factory=uuid4)
user_id: str
created_at: datetime = Field(default_factory=datetime.utcnow)
last_active: datetime = Field(default_factory=datetime.utcnow)
conversation: List[dict] = Field(default_factory=list) # raw turns
summary: Optional[str] = None # Claude‑generated short context
cart: List[CartItem] = Field(default_factory=list)
intent: Optional[str] = None
3. Building the `StateManager` Class with Async/Await
# version: 2026-10-03
import json
import asyncio
import logging
from uuid import UUID
import aioredis
import asyncpg
from pydantic import ValidationError
logger = logging.getLogger(__name__)
class StateManager:
def __init__(self, redis_url: str, pg_dsn: str):
self.redis = aioredis.from_url(redis_url, decode_responses=True)
self.pg_pool = None
self.pg_dsn = pg_dsn
async def init(self):
self.pg_pool = await asyncpg.create_pool(dsn=self.pg_dsn,
min_size=1, max_size=10)
async def _pg_fetch(self, session_id: UUID) -> dict | None:
async with self.pg_pool.acquire() as conn:
row = await conn.fetchrow(
"SELECT state FROM sessions WHERE session_id=$1",
str(session_id)
)
return json.loads(row["state"]) if row else None
async def load_state(self, session_id: UUID) -> SessionState:
# Try Redis first
cached = await self.redis.get(f"session:{session_id}")
if cached:
try:
return SessionState.model_validate_json(cached)
except ValidationError:
logger.warning("Corrupt cache for %s, falling back to DB", session_id)
# Fallback to PostgreSQL
raw = await self._pg_fetch(session_id)
if raw:
state = SessionState(**raw)
await self.redis.set(f"session:{session_id}", state.model_dump_json(),
ex=300) # 5‑minute TTL for active sessions
return state
# New session
raise KeyError(f"Session {session_id} not found")
async def save_state(self, state: SessionState, flush: bool = False):
key = f"session:{state.session_id}"
payload = state.model_dump_json()
await self.redis.set(key, payload, ex=300)
if flush:
# Immediate durable write (used on checkout)
async with self.pg_pool.acquire() as conn:
await conn.execute(
"""
INSERT INTO sessions (session_id, state, updated_at)
VALUES ($1, $2::jsonb, now())
ON CONFLICT (session_id) DO UPDATE
SET state = EXCLUDED.state,
updated_at = now()
""",
str(state.session_id), payload
)
else:
# Schedule asynchronous batch write
asyncio.create_task(self._delayed_flush(state.session_id))
async def _delayed_flush(self, session_id: UUID, delay: float = 2.0):
await asyncio.sleep(delay)
try:
raw = await self.redis.get(f"session:{session_id}")
if raw:
async with self.pg_pool.acquire() as conn:
await conn.execute(
"""
INSERT INTO sessions (session_id, state, updated_at)
VALUES ($1, $2::jsonb, now())
ON CONFLICT (session_id) DO UPDATE
SET state = EXCLUDED.state,
updated_at = now()
""",
str(session_id), raw
)
except Exception as e:
logger.error("Failed delayed flush for %s: %s", session_id, e)
async def delete_state(self, session_id: UUID):
await self.redis.delete(f"session:{session_id}")
async with self.pg_pool.acquire() as conn:
await conn.execute("DELETE FROM sessions WHERE session_id=$1",
str(session_id))
async def close(self):
await self.redis.close()
await self.pg_pool.close()
**Why async?** Voice turns arrive every 800 ms. Blocking the event loop to write to PostgreSQL would add 150‑200 ms per turn—unacceptable for a natural conversation.
4. Integrating Claude API with Tool Calling for State Updates
Claude’s tool‑calling works like a function schema sent in the request. We expose three tools: `read_state`, `write_state`, `summarize_state`.
# version: 2026-10-03
from anthropic import Anthropic, AsyncAnthropic
from anthropic.types import MessageParam, Tool
anthropic = AsyncAnthropic(api_key=... ) # already set in env
def tool_definitions():
return [
Tool(
name="read_state",
description="Return the JSON blob for a given session_id",
input_schema={"type": "object", "properties": {"session_id": {"type": "string"}}}
),
Tool(
name="write_state",
description="Persist a new JSON blob for a session_id",
input_schema={"type":"object","properties":{
"session_id":{"type":"string"},
"state":{"type":"object"}
}}
),
Tool(
name="summarize_state",
description="Create a short summary (≤200 tokens) from conversation history",
input_schema={"type":"object","properties":{
"session_id":{"type":"string"},
"history":{"type":"array","items":{"type":"string"}}
}}
),
]
async def call_claude(messages: list[MessageParam],
session_id: str,
state_manager: StateManager):
response = await anthropic.messages.create(
model="claude-3-5-sonnet-20240620",
max_tokens=1024,
messages=messages,
tools=tool_definitions(),
temperature=0.0,
)
# Detect tool calls
if response.content[0].type == "tool_use":
tool_call = response.content[0]
if tool_call.name == "read_state":
state = await state_manager.load_state(UUID(tool_call.input["session_id"]))
return {"role":"assistant","content":state.model_dump_json()}
elif tool_call.name == "write_state":
new_state = SessionState(**tool_call.input["state"])
await state_manager.save_state(new_state, flush=False)
return {"role":"assistant","content":"State written"}
elif tool_call.name == "summarize_state":
# naive summarizer using Claude itself
summary_resp = await anthropic.messages.create(
model="claude-3-5-sonnet-20240620",
max_tokens=200,
messages=[{
"role":"user",
"content":f"Summarize this conversation: {tool_call.input['history']}"
}],
temperature=0.0
)
summary = summary_resp.content[0].text
# patch state
state = await state_manager.load_state(UUID(tool_call.input["session_id"]))
state.summary = summary
await state_manager.save_state(state, flush=False)
return {"role":"assistant","content":"Summary updated"}
else:
return response
The key point: **Claude never holds the cart**—it only asks the cache to read/write. This eliminates token‑budget blow‑up and guarantees durability.
5. Implementing Redis + PostgreSQL Hybrid Storage Layer
Create the PostgreSQL table:
-- version: 2026-10-03
CREATE TABLE IF NOT EXISTS sessions (
session_id UUID PRIMARY KEY,
state JSONB NOT NULL,
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Row‑level security for GDPR compliance
ALTER TABLE sessions ENABLE ROW LEVEL SECURITY;
CREATE POLICY user_can_read ON sessions
USING (session_id = current_setting('app.session_id')::uuid);
Redis TTL strategy:
# In Docker compose (redis service)
redis:
image: redis:7.2-alpine
command: ["redis-server", "--save", "", "--appendonly", "yes"]
ports:
- "6379:6379"
**Tip:** Set `appendonly yes` in Redis to survive a process crash while still offering sub‑millisecond reads.
6. FastAPI Webhook Endpoint (Voice Platform Bridge)
# version: 2026-10-03
from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Depends
from fastapi.responses import JSONResponse
import uuid
app = FastAPI()
state_manager = StateManager(redis_url="redis://localhost:6379",
pg_dsn="postgresql://user:pass@localhost:5432/retail")
await state_manager.init()
@app.websocket("/ws/{session_id}")
async def voice_ws(ws: WebSocket, session_id: str):
await ws.accept()
try:
while True:
data = await ws.receive_json()
user_msg = data["text"]
# Load or create state
try:
state = await state_manager.load_state(uuid.UUID(session_id))
except KeyError:
state = SessionState(user_id=data["user_id"], session_id=uuid.UUID(session_id))
await state_manager.save_state(state)
# Append user turn
state.conversation.append({"role":"user","content":user_msg})
state.last_active = datetime.utcnow()
await state_manager.save_state(state)
# Build Claude payload
msgs = [{"role":"user","content":user_msg}]
claude_resp = await call_claude(msgs, session_id, state_manager)
# Append assistant turn and push to WS
state.conversation.append({"role":"assistant","content":claude_resp["content"]})
await state_manager.save_state(state)
await ws.send_json({"reply": claude_resp["content"]})
except WebSocketDisconnect:
logger.info("User %s disconnected", session_id)
finally:
# Optional idle cleanup handled elsewhere
pass
7. Advanced GUI & Event Loop Integration (Streamlit Desktop)
# version: 2026-10-03
import streamlit as st
import