
How to Ship a Real Autonomous Agent: MCP Architecture Patterns That Actually Work in Production
Production-ready MCP architecture patterns for autonomous agents: tool orchestration, context management, error resilience, security sandboxing, and observability.
Most autonomous agent prototypes die in production. They hallucinate tool calls, blow through context windows, silently fail on partial tool errors, and become impossible to debug at scale. The Model Context Protocol (MCP) provides the standardization layer that makes production agents tractable — but the protocol alone isn't enough. You need architecture patterns that handle concurrency, resilience, security boundaries, and observability at the system level.
This deep-dive covers the architectural patterns that separate agents that survive contact with real users from those that live forever in a Jupyter notebook.
Table of Contents
- 1. MCP Fundamentals: What the Protocol Actually Gives You
- 2. The Core Architecture: Client-Server Separation
- 3. Tool Orchestration Patterns
- 4. Context Management and Window Budgeting
- 5. Error Resilience and Retry Semantics
- 6. Security and Sandboxing Boundaries
- 7. Observability and Debugging Production Agents
- 8. Scaling: Multi-Agent Topologies
- 9. Frequently Asked Questions
1. MCP Fundamentals: What the Protocol Actually Gives You
MCP (Model Context Protocol) defines a JSON-RPC 2.0-based transport layer between an LLM-powered client and stateful tool servers. The critical architectural insight is that MCP separates tool definition from tool execution — each MCP server owns its tools, schemas, and state, while the client (your agent runtime) handles planning, context assembly, and LLM inference.
sequenceDiagram
participant User
participant Agent as Agent Runtime
participant LLM
participant MCP_S1 as MCP Server (Database)
participant MCP_S2 as MCP Server (API)
User->>Agent: Natural language request
Agent->>MCP_S1: tools/list
Agent->>MCP_S2: tools/list
MCP_S1-->>Agent: Tool schemas
MCP_S2-->>Agent: Tool schemas
Agent->>LLM: Context + tools + user input
LLM-->>Agent: Tool call decision
Agent->>MCP_S1: tools/call {name, arguments}
MCP_S1-->>Agent: Result
Agent->>LLM: Result + updated context
LLM-->>Agent: Final response
Agent-->>User: Answer
The three core MCP capabilities are:
- Tools — Stateful operations the LLM can invoke (CRUD, API calls, file operations)
- Resources — Read-only context (files, databases, knowledge bases) exposed to the client
- Prompts — Pre-defined interaction templates that structure multi-turn workflows
A production agent rarely connects to a single MCP server. You'll typically have 3-15 MCP servers, each owning a domain of tools. The architectural challenge is orchestration: deciding which server to query, how to parallelize calls, and how to handle the combinatorial explosion of possible tool sequences.
2. The Core Architecture: Client-Server Separation
The Agent Runtime
The agent runtime is the brain — it owns the LLM connection, conversation state, tool routing, and execution loop. In production, this is a long-running service (not a script), typically built as a stateless compute unit backed by a persistent session store.
// Agent runtime skeleton — the core execution loop
import { MCPClient } from '@modelcontextprotocol/sdk/client';
import { StreamableHTTPClientTransport } from '@modelcontextprotocol/sdk/client/streamableHttp';
import { randomUUID } from 'crypto';
interface AgentConfig {
model: string;
maxIterations: number;
contextBudget: number; // tokens
toolTimeoutMs: number;
circuitBreakerThreshold: number;
}
interface AgentSession {
sessionId: string;
messages: Message[];
toolRegistry: ToolDefinition[];
activeTools: string[];
iteration: number;
startedAt: Date;
}
class AgentRuntime {
private mcpClients: Map<string, MCPClient> = new Map();
private sessionStore: SessionStore;
private circuitBreakers: Map<string, CircuitBreaker> = new Map();
async initialize(servers: MCPServerConfig[]): Promise<void> {
for (const server of servers) {
const client = new MCPClient({
capabilities: { tools: {}, resources: {}, prompts: {} }
});
const transport = new StreamableHTTPClientTransport(
new URL(server.endpoint)
);
await client.connect(transport);
this.mcpClients.set(server.id, client);
}
}
async execute(sessionId: string, userMessage: string): Promise<string> {
const session = await this.sessionStore.load(sessionId);
const config = this.getConfig();
session.messages.push({ role: 'user', content: userMessage });
for (let i = 0; i < config.maxIterations; i++) {
session.iteration = i;
const context = this.assembleContext(session, config.contextBudget);
const toolSchemas = this.getActiveToolSchemas(session);
const response = await this.llm.generate({
model: config.model,
messages: context,
tools: toolSchemas,
temperature: 0.2,
});
if (response.toolCalls.length === 0) {
session.messages.push({ role: 'assistant', content: response.content });
await this.sessionStore.save(session);
return response.content;
}
// Execute tool calls
const results = await this.executeToolCalls(
response.toolCalls,
config.toolTimeoutMs
);
session.messages.push({
role: 'assistant',
content: response.content,
toolCalls: response.toolCalls,
});
for (const result of results) {
session.messages.push({
role: 'tool',
toolCallId: result.toolCallId,
content: result.output,
});
}
}
throw new AgentMaxIterationsError(sessionId, config.maxIterations);
}
}
MCP Server Topology
Production deployments use a hub-and-spoke topology where each MCP server is a dedicated microservice:
┌─────────────────────────────────────────────────────────┐
│ Agent Runtime Pool │
│ (stateless, auto-scaled, load-balanced) │
└──────────┬──────────────┬──────────────┬─────────────────┘
│ │ │
┌─────▼─────┐ ┌────▼─────┐ ┌───▼──────┐
│ MCP: DB │ │ MCP: API │ │MCP: Files│
│ Server │ │ Server │ │ Server │
└───────────┘ └──────────┘ └──────────┘
Key design decisions for the server topology:
| Decision | Recommendation | Rationale |
|---|---|---|
| Transport | Streamable HTTP (not stdio) | Survives restarts, supports horizontal scaling |
| Auth | Per-server API keys + mTLS | Blast radius containment |
| State | Server-side with TTL | Client stays stateless |
| Deployment | K8s with pod-per-server | Independent scaling and rollback |
3. Tool Orchestration Patterns
Pattern 1: Sequential Tool Chaining
The simplest pattern — the LLM decides one tool call at a time, results feed back into context, and the next decision incorporates the prior result. This is the MCP default and works for 80% of cases.
// The execution loop naturally supports sequential chaining
async executeToolCallsSequential(
toolCalls: ToolCall[],
timeoutMs: number
): Promise<ToolResult[]> {
const results: ToolResult[] = [];
for (const call of toolCalls) {
const result = await this.executeWithCircuitBreaker(call, timeoutMs);
results.push(result);
}
return results;
}
When to use: When tool calls have data dependencies (output of tool A is input to tool B). Most agentic workflows follow this pattern.
Pattern 2: Parallel Fan-Out
When multiple independent tool calls can execute concurrently, parallelization dramatically reduces latency:
async executeToolCallsParallel(
toolCalls: ToolCall[],
timeoutMs: number
): Promise<ToolResult[]> {
// Partition into independent groups
const groups = this.partitionIndependentCalls(toolCalls);
const allResults: ToolResult[] = [];
// Execute groups sequentially, calls within group in parallel
for (const group of groups) {
const results = await Promise.allSettled(
group.map(call => this.executeWithCircuitBreaker(call, timeoutMs))
);
for (const [i, settled] of results.entries()) {
if (settled.status === 'fulfilled') {
allResults.push(settled.value);
} else {
allResults.push(this.buildErrorResult(group[i], settled.reason));
}
}
}
return allResults;
}
private partitionIndependentCalls(toolCalls: ToolCall[]): ToolCall[][] {
// Topological sort by declared dependencies
// Calls without dependencies form the first group
const dependencyMap = new Map<string, string[]>();
for (const call of toolCalls) {
if (call.dependencies) {
dependencyMap.set(call.id, call.dependencies);
}
}
// Kahn's algorithm for topological grouping
const groups: ToolCall[][] = [];
let remaining = [...toolCalls];
while (remaining.length > 0) {
const ready = remaining.filter(call => {
const deps = dependencyMap.get(call.id) ?? [];
return deps.every(dep =>
groups.flat().some(completed => completed.id === dep)
);
});
if (ready.length === 0) {
// Circular dependency — fail fast
throw new CircularDependencyError(toolCalls);
}
groups.push(ready);
remaining = remaining.filter(c => !ready.includes(c));
}
return groups;
}
When to use: When the LLM emits multiple independent tool calls (e.g., fetch user data AND check order status simultaneously). MCP servers that handle the same domain but different operations are prime candidates.
Pattern 3: Router-Dispatcher
A pre-LLM router classifies the user intent and activates only the relevant MCP server subset, reducing tool schema bloat:
interface IntentRouter {
route(userMessage: string): Promise<RouteDecision>;
}
interface RouteDecision {
activeServerIds: string[];
intent: string;
confidence: number;
fallbackServers: string[];
}
class LightweightRouter implements IntentRouter {
async route(userMessage: string): Promise<RouteDecision> {
// Use a small, fast model for classification (not your main agent model)
const classification = await this.smallModel.classify({
model: 'claude-3-haiku',
messages: [{
role: 'user',
content: `Classify this request into one of: database, api, files, search, compute\n\nRequest: ${userMessage}`
}],
maxTokens: 50,
});
const intent = this.parseIntent(classification);
const serverMap: Record<string, string[]> = {
database: ['mcp-postgres', 'mcp-redis'],
api: ['mcp-weather', 'mcp-stripe'],
files: ['mcp-gcs', 'mcp-s3'],
search: ['mcp-vector', 'mcp-elastic'],
compute: ['mcp-python', 'mcp-shell'],
};
return {
activeServerIds: serverMap[intent] ?? [],
intent,
confidence: this.computeConfidence(classification),
fallbackServers: ['mcp-general'],
};
}
}
// In the agent runtime, before the main loop:
async execute(sessionId: string, userMessage: string): Promise<string> {
const session = await this.sessionStore.load(sessionId);
const route = await this.intentRouter.route(userMessage);
// Activate only relevant tools — reduces schema tokens by 60-80%
session.activeTools = route.activeServerIds;
if (route.confidence < 0.6) {
// Low confidence — include fallback tools
session.activeTools.push(...route.fallbackServers);
}
// ... proceed with main loop using only active tools
}
When to use: When you have 10+ MCP servers and the full tool schema exceeds 40% of your context window. The router adds ~200ms latency but saves thousands of tokens per turn.
4. Context Management and Window Budgeting
Context window exhaustion is the #1 production failure mode for agents. The strategy is to treat context as a budgeted resource with explicit accounting.
The Token Budget Framework
interface ContextBudget {
totalWindow: number; // e.g., 200,000 for Claude 3.5 Sonnet
systemPrompt: number; // Fixed overhead
toolSchemas: number; // Dynamic — varies by active tools
conversationHistory: number; // Grows with each turn
toolResults: number; // Grows with each tool call
outputReserve: number; // Always reserve for the response
}
function calculateBudget(
window: number,
activeTools: ToolDefinition[],
historyMessages: Message[]
): ContextBudget {
const toolSchemaTokens = activeTools.reduce(
(sum, tool) => sum + estimateTokens(tool.jsonSchema), 0
);
const historyTokens = historyMessages.reduce(
(sum, msg) => sum + estimateTokens(msg.content), 0
);
return {
totalWindow: window,
systemPrompt: 500,
toolSchemas: toolSchemaTokens,
conversationHistory: historyTokens,
toolResults: 0, // Computed per iteration
outputReserve: Math.max(4096, window * 0.05), // 5% minimum
};
}
Compaction Strategies
When the budget approaches exhaustion, apply progressive compaction:
class ContextManager {
private compactionThreshold = 0.85; // Compact at 85% usage
async prepareContext(
session: AgentSession,
budget: ContextBudget
): Promise<PreparedContext> {
const currentUsage = this.calculateUsage(session, budget);
if (currentUsage / budget.totalWindow > this.compactionThreshold) {
await this.compact(session);
}
return {
systemPrompt: this.buildSystemPrompt(session),
messages: session.messages,
tools: this.getActiveToolSchemas(session),
};
}
private async compact(session: AgentSession): Promise<void> {
// Strategy 1: Summarize old messages
const cutoff = session.messages.length - this.keepRecentCount;
if (cutoff > 0) {
const oldMessages = session.messages.slice(0, cutoff);
const summary = await this.llm.summarize({
model: 'claude-3-haiku',
messages: [{
role: 'user',
content: `Summarize this conversation preserving all key facts, decisions, and tool results:\n\n${oldMessages.map(m => `[${m.role}] ${m.content}`).join('\n')}`
}],
maxTokens: 1000,
});
session.messages = [
{ role: 'system', content: `[COMPACTED HISTORY]\n${summary}` },
...session.messages.slice(cutoff),
];
}
// Strategy 2: Truncate large tool results
for (const msg of session.messages) {
if (msg.role === 'tool' && this.estimateTokens(msg.content) > 2000) {
msg.content = this.truncateWithSummary(msg.content, 2000);
msg.truncated = true;
}
}
// Strategy 3: Prune inactive tool schemas
session.activeTools = session.activeTools.filter(toolId =>
this.recentlyUsed(toolId, session.messages, 3)
);
}
}
The Sliding Window Anti-Pattern
Do NOT use a simple sliding window that drops old messages. Agents rely on earlier context to avoid repeating tool calls, maintain consistency, and reference prior results. Instead:
- Summarize old messages (preserving facts, dropping prose)
- Tier messages by importance (user messages > tool results > intermediate thoughts)
- Externalize large data to a retrieval store, keeping only references in context
5. Error Resilience and Retry Semantics
Production agents fail constantly. The difference between a production system and a demo is how it handles those failures.
Circuit Breaker Pattern
class CircuitBreaker {
private state: 'CLOSED' | 'OPEN' | 'HALF_OPEN' = 'CLOSED';
private failureCount = 0;
private lastFailureTime = 0;
private successCount = 0;
constructor(
private serverId: string,
private threshold: number = 5,
private resetTimeoutMs: number = 30_000,
private halfOpenMax: number = 3
) {}
async execute<T>(fn: () => Promise<T>): Promise<T> {
if (this.state === 'OPEN') {
if (Date.now() - this.lastFailureTime > this.resetTimeoutMs) {
this.state = 'HALF_OPEN';
this.successCount = 0;
} else {
throw new CircuitOpenError(this.serverId);
}
}
try {
const result = await fn();
this.onSuccess();
return result;
} catch (error) {
this.onFailure(error);
throw error;
}
}
private onSuccess(): void {
this.failureCount = 0;
if (this.state === 'HALF_OPEN') {
this.successCount++;
if (this.successCount >= this.halfOpenMax) {
this.state = 'CLOSED';
}
}
}
private onFailure(error: Error): void {
this.failureCount++;
this.lastFailureTime = Date.now();
if (this.failureCount >= this.threshold) {
this.state = 'OPEN';
}
// Classify for retry decisions
const isRetryable = this.isRetryableError(error);
if (!isRetryable && this.state !== 'OPEN') {
// Non-retryable — open immediately
this.failureCount = this.threshold;
this.state = 'OPEN';
}
}
private isRetryableError(error: Error): boolean {
// Rate limits, timeouts, 5xx — retryable
// 4xx (auth, validation), schema errors — NOT retryable
if (error instanceof RateLimitError) return true;
if (error instanceof TimeoutError) return true;
if (error instanceof ServerError && error.status >= 500) return true;
if (error instanceof AuthError) return false;
if (error instanceof ValidationError) return false;
return false;
}
}
Retry with Exponential Backoff and Jitter
class RetryPolicy {
constructor(
private maxRetries: number = 3,
private baseDelayMs: number = 1000,
private maxDelayMs: number = 30_000
) {}
async executeWithRetry<T>(
fn: () => Promise<T>,
context: RetryContext
): Promise<T> {
let lastError: Error;
for (let attempt = 0; attempt <= this.maxRetries; attempt++) {
try {
return await fn();
} catch (error) {
lastError = error as Error;
if (!this.isRetryable(error) || attempt === this.maxRetries) {
break;
}
const delay = this.calculateDelay(attempt);
context.onRetry({
attempt: attempt + 1,
delay,
error: lastError,
toolCallId: context.toolCallId,
});
await sleep(delay);
}
}
throw lastError;
}
private calculateDelay(attempt: number): number {
// Exponential backoff with full jitter
const exponential = Math.min(
this.maxDelayMs,
this.baseDelayMs * Math.pow(2, attempt)
);
return Math.random() * exponential; // Full jitter
}
}
Agent-Level Error Recovery
When a tool call fails irrecoverably, the agent should not crash. Instead, feed the failure back to the LLM as a tool result and let it adapt:
private buildErrorResult(
toolCall: ToolCall,
error: unknown
): ToolResult {
const errorMessage = error instanceof Error ? error.message : String(error);
return {
toolCallId: toolCall.id,
output: JSON.stringify({
error: true,
tool: toolCall.name,
message: errorMessage,
recoverable: this.isRecoverable(error),
suggestion: this.getRecoverySuggestion(toolCall, error),
}),
isError: true,
};
}
private getRecoverySuggestion(toolCall: ToolCall, error: unknown): string {
if (error instanceof ValidationError) {
return `The arguments for '${toolCall.name}' were invalid. ` +
`Check the schema and try again with corrected arguments.`;
}
if (error instanceof RateLimitError) {
return `The '${toolCall.name}' service is rate-limited. ` +
`Try a different approach or wait before retrying.`;
}
if (error instanceof TimeoutError) {
return `The '${toolCall.name}' call timed out. ` +
`Consider breaking the operation into smaller steps.`;
}
return `The '${toolCall.name}' call failed. Try an alternative approach.`;
}
This pattern is critical: by returning structured error information to the LLM, you enable self-healing — the agent can retry with different parameters, switch to an alternative tool, or inform the user of the limitation.
6. Security and Sandboxing Boundaries
Autonomous agents that can execute code, access databases, and call external APIs are powerful — and dangerous. Production architectures enforce multiple security layers.
Layer 1: Tool-Level Authorization
Each MCP server enforces its own authorization. The agent runtime passes a scoped token that limits what the agent can do:
interface ToolAuthorization {
agentId: string;
userId: string;
scope: ToolScope[];
rateLimit: RateLimitPolicy;
dataClassification: DataLevel;
expiresAt: Date;
}
interface ToolScope {
toolName: string;
permissions: ('read' | 'write' | 'delete')[];
resourceFilter?: Record<string, unknown>; // e.g., { orgId: 'acme' }
}
// Example: Agent can only read/write its own user's data
const agentAuth: ToolAuthorization = {
agentId: 'agent-7x2k',
userId: 'user-123',
scope: [
{ toolName: 'db.query', permissions: ['read'], resourceFilter: { userId: 'user-123' } },
{ toolName: 'db.insert', permissions: ['write'], resourceFilter: { userId: 'user-123' } },
{ toolName: 'api.weather', permissions: ['read'] },
],
rateLimit: { requestsPerMinute: 60, tokensPerHour: 100_000 },
dataClassification: 'internal',
expiresAt: new Date(Date.now() + 30 * 60 * 1000), // 30 min TTL
};
Layer 2: Execution Sandboxing
For tools that execute code (Python, shell, etc.), never run in the same process as the agent runtime:
class SandboxedExecutor {
private sandboxPool: SandboxInstance[];
async execute(toolCall: ToolCall, auth: ToolAuthorization): Promise<ToolResult> {
const sandbox = await this.sandboxPool.acquire();
try {
// Apply resource limits
const result = await sandbox.execute({
code: toolCall.arguments.code,
timeoutMs: 30_000,
memoryLimitMb: 256,
cpuLimit: 0.5, // 50% of one core
networkAccess: auth.scope.some(s => s.toolName.startsWith('api.')),
filesystemAccess: auth.scope.some(s => s.toolName.startsWith('file.')),
environment: {
// Inject only whitelisted env vars
OPENAI_API_KEY: this.secretStore.get('OPENAI_API_KEY'),
},
});
return { toolCallId: toolCall.id, output: result.stdout };
} finally {
await this.sandboxPool.release(sandbox);
}
}
}
Layer 3: Input/Output Filtering
Filter both directions to prevent prompt injection from tool results and data exfiltration from agent output:
class SecurityFilter {
async sanitizeToolOutput(output: string, source: string): Promise<string> {
// Remove potential prompt injection patterns from tool results
const sanitized = output
.replace(/<\|.*?\|>/g, '[REDACTED]') // LLM system tokens
.replace(/\b(system|developer)\b\s*:/gi, '[FILTERED]')
.replace(/ignore\s+(all\s+)?(previous|above)\s+instructions/gi, '[FILTERED]');
// Truncate to prevent context flooding
if (this.estimateTokens(sanitized) > 4000) {
return sanitized.slice(0, this.findSentenceBoundary(sanitized, 4000)) +
'\n... [output truncated]';
}
return sanitized;
}
async validateAgentOutput(output: string, policy: OutputPolicy): Promise<string> {
// Check for data exfiltration patterns
const sensitivePatterns = [
/\b\d{4}[-\s]?\d{4}[-\s]?\d{4}[-\s]?\d{4}\b/g, // Credit cards
/\b\d{3}[-\s]?\d{2}[-\s]?\d{3}[-\s]?\d{4}\b/g, // SSNs
];
for (const pattern of sensitivePatterns) {
const matches = output.match(pattern);
if (matches && !policy.allowSensitiveData) {
throw new OutputPolicyViolation(
`Output contains ${matches.length} potential sensitive data matches`
);
}
}
return output;
}
}
7. Observability and Debugging Production Agents
Debugging agents is fundamentally different from debugging deterministic code. The same input can produce different tool sequences across runs. You need specialized observability.
Tracing Every Decision Point
import { trace, span, SpanStatus } from '@opentelemetry/api';
const tracer = trace.getTracer('agent-runtime', '1.0.0');
async execute(sessionId: string, userMessage: string): Promise<string> {
const rootSpan = tracer.startSpan('agent.execute', {
attributes: {
'agent.session_id': sessionId,
'agent.model': this.config.model,
'agent.max_iterations': this.config.maxIterations,
},
});
try {
const session = await this.sessionStore.load(sessionId);
for (let i = 0; i < this.config.maxIterations; i++) {
const iterationSpan = tracer.startSpan(`agent.iteration.${i}`, {
attributes: {
'agent.iteration': i,
'agent.context_tokens': this.estimateContextTokens(session),
'agent.active_tools': session.activeTools.length,
},
});
// Trace LLM call
const llmSpan = tracer.startSpan('llm.generate', {
attributes: {
'llm.model': this.config.model,
'llm.input_tokens': this.estimateInputTokens(session),
},
});
const response = await this.llm.generate(/* ... */);
llmSpan.setAttributes({
'llm.output_tokens': response.usage?.outputTokens ?? 0,
'llm.tool_calls': response.toolCalls.length,
'llm.finish_reason': response.finishReason,
});
llmSpan.end();
// Trace each tool call
for (const toolCall of response.toolCalls) {
const toolSpan = tracer.startSpan(`tool.${toolCall.name}`, {
attributes: {
'tool.name': toolCall.name,
'tool.server': this.getServerForTool(toolCall.name),
'tool.arguments_size': JSON.stringify(toolCall.arguments).length,
},
});
try {
const result = await this.executeWithCircuitBreaker(toolCall);
toolSpan.setAttributes({
'tool.status': 'success',
'tool.result_size': result.output.length,
});
} catch (error) {
toolSpan.setStatus(SpanStatus.ERROR, String(error));
toolSpan.setAttributes({
'tool.status': 'error',
'tool.error_type': error.constructor.name,
});
} finally {
toolSpan.end();
}
}
iterationSpan.end();
}
rootSpan.setStatus(SpanStatus.OK);
return response.content;
} catch (error) {
rootSpan.setStatus(SpanStatus.ERROR, String(error));
throw error;
} finally {
rootSpan.end();
}
}
Structured Logging for Agent Events
interface AgentLogEvent {
timestamp: string;
sessionId: string;
iteration: number;
event: 'tool_call' | 'tool_result' | 'llm_response' | 'error' | 'compaction';
data: Record<string, unknown>;
}
// Every meaningful event gets logged — this is your debugging lifeline
private logEvent(sessionId: string, event: AgentLogEvent): void {
this.logger.info(JSON.stringify({
level: 'info',
event: event.event,
session: sessionId,
iteration: event.iteration,
...event.data,
timestamp: new Date().toISOString(),
}));
}
The Debug Replay Pattern
The most powerful debugging tool for agents is replay — storing every LLM request/response pair so you can reproduce any execution:
class ExecutionRecorder {
async record(
sessionId: string,
iteration: number,
request: LLMRequest,
response: LLMResponse
): Promise<void> {
await this.store.put(`replay/${sessionId}/${iteration}`, {
request: {
model: request.model,
messages: request.messages,
tools: request.tools,
temperature: request.temperature,
},
response: {
content: response.content,
toolCalls: response.toolCalls,
usage: response.usage,
finishReason: response.finishReason,
},
metadata: {
timestamp: new Date().toISOString(),
latencyMs: response.latency,
contextTokens: request.messages.reduce((s, m) => s + m.tokens, 0),
},
});
}
}
// In test/debug mode, swap the real LLM for a replay provider
class ReplayLLMProvider implements LLMProvider {
constructor(private recorder: ExecutionRecorder) {}
async generate(request: LLMRequest): Promise<LLMResponse> {
// Find the matching recorded response
const recorded = await this.recorder.findMatch(request);
if (!recorded) {
throw new Error('No replay data for this request — run in live mode first');
}
return recorded;
}
}
This enables deterministic testing of agent behavior without calling the LLM — critical for CI/CD and regression testing.
8. Scaling: Multi-Agent Topologies
Single agents hit limits. Production systems use multi-agent topologies where specialized agents collaborate.
Pattern: Orchestrator-Worker
interface AgentTopology {
orchestrator: OrchestratorConfig;
workers: WorkerConfig[];
}
interface WorkerConfig {
id: string;
name: string;
mcpServers: string[];
maxIterations: number;
model: string;
contextBudget: number;
}
class OrchestratorAgent {
async execute(userMessage: string): Promise<string> {
// Step 1: Plan — decompose into subtasks
const plan = await this.plan(userMessage);
// Step 2: Dispatch — send subtasks to specialized workers
const results = await Promise.all(
plan.subtasks.map(async (subtask) => {
const worker = this.workers.get(subtask.workerId);
return worker.execute(subtask.instructions, subtask.context);
})
);
// Step 3: Synthesize — combine results
return this.synthesize(userMessage, plan, results);
}
private async plan(userMessage: string): Promise<ExecutionPlan> {
return this.llm.generateStructured<ExecutionPlan>({
model: 'claude-3-opus',
systemPrompt: `You are a task planner. Decompose the request into parallel subtasks.
Each subtask must specify which worker agent should handle it.
Available workers: ${this.getWorkerDescriptions()}`,
userMessage,
schema: ExecutionPlanSchema,
});
}
}
Pattern: Agent-to-Agent via MCP
Advanced topologies connect agents directly via MCP — one agent exposes itself as an MCP server to another:
// Worker agent exposes itself as an MCP server
import { MCPServer } from '@modelcontextprotocol/sdk/server';
class AgentAsMCPServer {
private server: MCPServer;
constructor(private agent: WorkerAgent) {
this.server = new MCPServer({
name: agent.config.name,
version: '1.0.0',
});
// Expose the agent's capabilities as MCP tools
this.server.registerTool({
name: `${agent.config.name}.execute`,
description: `Execute a task using the ${agent.config.name} agent`,
inputSchema: {
type: 'object',
properties: {
task: { type: 'string', description: 'The task to execute' },
context: { type: 'string', description: 'Additional context' },
maxSteps: { type: 'number', description: 'Maximum execution steps' },
},
required: ['task'],
},
handler: async (args) => {
const result = await this.agent.execute({
userMessage: args.task,
additionalContext: args.context,
maxIterations: args.maxSteps ?? 5,
});
return { content: [{ type: 'text', text: result }] };
},
});
}
async start(port: number): Promise<void> {
await this.server.listen({ port, host: '0.0.0.0' });
}
}
Scaling Considerations
| Dimension | Strategy | Threshold |
|---|---|---|
| Tool count | Router pattern + lazy loading | >20 tools active |
| Context size | Progressive compaction + external memory | >100K tokens used |
| Latency | Parallel tool execution + streaming | >5s per iteration |
| Concurrency | Session-per-pod + queue | >50 concurrent sessions |
| Cost | Model routing (haiku for classification, sonnet for execution) | >$0.50/session |
Cost Optimization: Model Tiering
interface ModelRouter {
route(task: RoutingCriteria): string;
}
class CostAwareModelRouter implements ModelRouter {
route(criteria: RoutingCriteria): string {
if (criteria.isClassification) return 'claude-3-haiku'; // $0.25/MTok
if (criteria.isSummarization) return 'claude-3-haiku'; // $0.25/MTok
if (criteria.complexity === 'high') return 'claude-3-opus'; // $15/MTok
if (criteria.complexity === 'medium') return 'claude-3-sonnet'; // $3/MTok
return 'claude-3-haiku'; // Default cheap
}
}
// In practice: a router agent uses haiku, the main agent uses sonnet,
// and only complex reasoning steps escalate to opus.
// This typically reduces costs by 60-80% vs. using one model for everything.
9. Frequently Asked Questions
Q: How do I handle MCP server failures without breaking the agent loop?
Never let an MCP server failure crash the agent. Wrap every tool call in a circuit breaker. When a server fails, return a structured error as a tool result and let the LLM decide whether to retry, use an alternative tool, or inform the user. The agent's execution loop should treat tool failures as information, not as exceptions.
Q: What's the right context window strategy for long-running agents?
Use a three-tier approach: (1) keep the last 5-10 messages verbatim, (2) summarize older messages into a running summary, (3) externalize large data blobs to a retrieval store (vector DB or file system) and keep only references in context. Never use a simple sliding window — agents need memory of prior decisions to avoid circular tool calls.
Q: How do I test agents when LLM outputs are non-deterministic?
Implement the replay pattern: record every LLM request/response pair during live execution, then use those recordings in test mode. Combine this with property-based testing that validates invariants ("agent never calls a tool without valid auth", "agent always terminates within N iterations") rather than asserting exact outputs. For CI, test the agent's behavioral properties — not its specific words.
For more on production AI system architecture, see Tamiz's Insights for additional deep-dives on agent engineering patterns and infrastructure design.
MCP Server Lifecycle Management in Production
One of the most overlooked aspects of shipping MCP-based agents is managing the lifecycle of MCP servers themselves. Unlike traditional microservices, MCP servers maintain stateful sessions with language models, and improper lifecycle management leads to resource leaks, stale connections, and degraded response quality.
The Session-Scoped Server Pattern
The most reliable production pattern treats each MCP server instance as session-scoped rather than globally shared. This eliminates the race conditions and state contamination that plague shared server pools.
# mcp_lifecycle.py
import asyncio
from contextlib import asynccontextmanager
from typing import AsyncGenerator
from mcp import ClientSession, StdioServerParameters
class MCPServerPool:
"""
Manages MCP server instances with automatic lifecycle handling.
Each session gets an isolated server instance to prevent state leakage.
"""
def __init__(self, max_instances: int = 10, idle_timeout: float = 300.0):
self.max_instances = max_instances
self.idle_timeout = idle_timeout
self._semaphore = asyncio.Semaphore(max_instances)
self._instances: dict[str, ClientSession] = {}
self._lock = asyncio.Lock()
@asynccontextmanager
async def acquire(self, server_params: StdioServerParameters) -> AsyncGenerator[ClientSession, None]:
"""Acquire a session-scoped MCP server instance."""
async with self._semaphore:
session_id = str(uuid4())
session = None
try:
session = await self._create_session(server_params)
async with session:
yield session
except Exception as e:
logger.error(f"MCP server error for session {session_id}: {e}")
raise
finally:
await self._cleanup(session_id)
async def _create_session(self, params: StdioServerParameters) -> ClientSession:
"""Spawn a fresh MCP server process."""
from mcp import stdio_client
async with stdio_client(params) as (read_stream, write_stream):
session = ClientSession(read_stream, write_stream)
await session.initialize()
return session
async def _cleanup(self, session_id: str):
"""Ensure resources are released."""
async with self._lock:
self._instances.pop(session_id, None)
The key insight here is that asynccontextmanager guarantees cleanup even when exceptions propagate. In production, we've seen cases where a hung tool execution would leak server processes indefinitely without this pattern.
Health Checking and Automatic Recovery
MCP servers can crash silently—especially when they wrap external services that time out. Production systems need proactive health checking:
# mcp_health.py
import time
from dataclasses import dataclass, field
from enum import Enum
class ServerHealth(Enum):
HEALTHY = "healthy"
DEGRADED = "degraded"
UNRESPONSIVE = "unresponsive"
DEAD = "dead"
@dataclass
class ServerHealthState:
last_heartbeat: float = field(default_factory=time.time)
consecutive_failures: int = 0
total_requests: int = 0
health: ServerHealth = ServerHealth.HEALTHY
def record_success(self):
self.consecutive_failures = 0
self.health = ServerHealth.HEALTHY
self.last_heartbeat = time.time()
self.total_requests += 1
def record_failure(self):
self.consecutive_failures += 1
self.total_requests += 1
if self.consecutive_failures >= 3:
self.health = ServerHealth.UNRESPONSIVE
elif self.consecutive_failures >= 1:
self.health = ServerHealth.DEGRADED
class ResilientMCPServer:
"""Wraps an MCP server with automatic health monitoring and recovery."""
def __init__(self, params: StdioServerParameters, circuit_breaker_threshold: int = 5):
self.params = params
self.health_state = ServerHealthState()
self.circuit_breaker_threshold = circuit_breaker_threshold
self.circuit_open_until: float = 0
self._session: ClientSession | None = None
async def execute_tool(self, tool_name: str, arguments: dict) -> str:
"""Execute a tool with circuit breaker protection."""
now = time.time()
# Check circuit breaker
if now < self.circuit_open_until:
raise MCPCircuitOpenError(
f"Circuit breaker open, retry after {self.circuit_open_until - now:.1f}s"
)
# Check for stale connection
if self._is_stale():
await self._reconnect()
try:
result = await asyncio.wait_for(
self._session.call_tool(tool_name, arguments),
timeout=30.0
)
self.health_state.record_success()
return result.content[0].text
except (asyncio.TimeoutError, ConnectionError) as e:
self.health_state.record_failure()
if self.health_state.consecutive_failures >= self.circuit_breaker_threshold:
self.circuit_open_until = time.time() + 60 # Open for 60 seconds
raise MCPCircuitOpenError(str(e))
raise
def _is_stale(self) -> bool:
return (time.time() - self.health_state.last_heartbeat) > 120
async def _reconnect(self):
"""Tear down and recreate the session."""
if self._session:
try:
await self._session.close()
except Exception:
pass
# Reconnection logic follows the same pattern as initial setup
Observability for MCP Agents
Production MCP agents fail in ways that are invisible without proper instrumentation. The MCP protocol itself is relatively opaque—unlike REST APIs, there's no standardized logging format, no built-in tracing, and no metrics endpoint.
Structured Event Logging
Every interaction between your agent and an MCP server should produce a structured log entry:
# mcp_observability.py
import structlog
from dataclasses import dataclass
from datetime import datetime
from uuid import UUID
logger = structlog.get_logger()
@dataclass
class MCPCallEvent:
session_id: str
tool_name: str
arguments: dict
result: str | None
error: str | None
duration_ms: float
timestamp: datetime
model_token_usage: int = 0
class MCPTracer:
"""Provides distributed tracing for MCP interactions."""
def __init__(self, service_name: str = "mcp-agent"):
self.service_name = service_name
self._traces: list[MCPCallEvent] = []
async def traced_call(self, session: ClientSession, tool_name: str,
arguments: dict, session_id: str):
"""Wrap an MCP tool call with full observability."""
start_time = time.perf_counter()
trace_id = str(uuid4())
logger.info("mcp.call.start",
tool=tool_name,
session_id=session_id,
trace_id=trace_id,
arguments_keys=list(arguments.keys()))
try:
result = await session.call_tool(tool_name, arguments)
duration_ms = (time.perf_counter() - start_time) * 1000
event = MCPCallEvent(
session_id=session_id,
tool_name=tool_name,
arguments=arguments,
result=result.content[0].text,
error=None,
duration_ms=duration_ms,
timestamp=datetime.utcnow()
)
self._traces.append(event)
logger.info("mcp.call.success",
tool=tool_name,
trace_id=trace_id,
duration_ms=round(duration_ms, 2),
result_length=len(result.content[0].text))
return result
except Exception as e:
duration_ms = (time.perf_counter() - start_time) * 1000
event = MCPCallEvent(
session_id=session_id,
tool_name=tool_name,
arguments=arguments,
result=None,
error=str(e),
duration_ms=duration_ms,
timestamp=datetime.utcnow()
)
self._traces.append(event)
logger.error("mcp.call.failed",
tool=tool_name,
trace_id=trace_id,
duration_ms=round(duration_ms, 2),
error=str(e))
raise
def get_p99_latency(self, tool_name: str | None = None) -> float:
"""Calculate p99 latency for a specific tool or all tools."""
events = [e for e in self._traces
if e.error is None and (tool_name is None or e.tool_name == tool_name)]
if not events:
return 0.0
sorted_latencies = sorted(e.duration_ms for e in events)
p99_index = int(len(sorted_latencies) * 0.99)
return sorted_latencies[p99_index]
Prometheus Metrics Integration
For production monitoring, expose MCP-specific metrics that map to SLOs:
# mcp_metrics.py
from prometheus_client import Counter, Histogram, Gauge
# Tool call counters
mcp_tool_calls_total = Counter(
'mcp_tool_calls_total',
'Total MCP tool calls',
['tool_name', 'status'] # status: success, failure, timeout
)
# Latency histograms
mcp_tool_latency_seconds = Histogram(
'mcp_tool_latency_seconds',
'MCP tool call latency',
['tool_name'],
buckets=[0.1, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0]
)
# Connection health
mcp_server_connections_active = Gauge(
'mcp_server_connections_active',
'Number of active MCP server connections'
)
mcp_circuit_breaker_open = Gauge(
'mcp_circuit_breaker_open',
'Circuit breaker status (1=open, 0=closed)',
['server_name']
)
# Token efficiency
mcp_tokens_per_successful_call = Histogram(
'mcp_tokens_per_successful_call',
'LLM tokens consumed per successful MCP tool call',
['tool_name'],
buckets=[50, 100, 200, 500, 1000, 2000, 5000]
)
These metrics enable meaningful alerting. For example, an alert on mcp_tool_latency_seconds p99 exceeding 10 seconds for any tool indicates a systemic problem. A sudden spike in mcp_tokens_per_successful_call suggests the model is struggling with tool output formatting.
Security Patterns for Production MCP
MCP servers execute code on behalf of agents. In production, this creates a trust boundary that must be explicitly managed. The following patterns address the most common security concerns.
Tool Allowlisting and Argument Validation
Never trust tool arguments blindly. MCP servers can accept arbitrary JSON, and malicious or malformed input can cause unexpected behavior:
# mcp_security.py
from pydantic import BaseModel, field_validator
from typing import Any
class ToolPermissionPolicy:
"""Defines what tools an agent can access and with what constraints."""
def __init__(self):
self._allowed_tools: dict[str, ToolPolicy] = {}
def register_tool(self, tool_name: str, policy: ToolPolicy):
self._allowed_tools[tool_name] = policy
def validate_call(self, tool_name: str, arguments: dict) -> dict:
"""Validate and sanitize tool arguments before execution."""
if tool_name not in self._allowed_tools:
raise ToolPermissionDenied(f"Tool '{tool_name}' is not in the allowlist")
policy = self._allowed_tools[tool_name]
validated = policy.schema.model_validate(arguments)
return validated.model_dump()
def get_allowed_tools(self) -> list[str]:
return list(self._allowed_tools.keys())
class ToolPolicy(BaseModel):
schema: type[BaseModel]
max_arguments_size: int = 10_000
requires_approval: bool = False
approval_timeout: float = 300.0
# Example: A file system tool with strict validation
class ReadFileArgs(BaseModel):
path: str
@field_validator("path")
def validate_path(cls, v: str) -> str:
if ".." in v:
raise ValueError("Path traversal not allowed")
if not v.startswith("/allowed/base/path"):
raise ValueError("Path must be within allowed directory")
return v
class WriteFileArgs(BaseModel):
path: str
content: str
@field_validator("path")
def validate_path(cls, v: str) -> str:
if ".." in v:
raise ValueError("Path traversal not allowed")
if not v.startswith("/allowed/base/path"):
raise ValueError("Path must be within allowed directory")
return v
@field_validator("content")
def validate_content_size(cls, v: str) -> str:
if len(v) > 1_000_000:
raise ValueError("Content exceeds 1MB limit")
return v
Human-in-the-Loop Approval Gates
For destructive or high-impact operations, implement approval gates that pause agent execution:
# mcp_approval.py
import asyncio
from dataclasses import dataclass
from enum import Enum
from typing import Callable, Awaitable
class ApprovalDecision(Enum):
APPROVED = "approved"
DENIED = "denied"
TIMEOUT = "timeout"
@dataclass
class ApprovalRequest:
tool_name: str
arguments: dict
risk_level: str # "low", "medium", "high", "critical"
description: str
agent_id: str
session_id: str
class ApprovalGate:
"""Blocks agent execution until a human approves sensitive operations."""
def __init__(self, approval_callback: Callable[[ApprovalRequest], Awaitable[ApprovalDecision]]):
self._callback = approval_callback
self._pending: dict[str, asyncio.Future[ApprovalDecision]] = {}
async def check(self, tool_name: str, arguments: dict,
risk_level: str, description: str,
agent_id: str, session_id: str) -> bool:
"""Check if a tool call requires approval and obtain it if needed."""
if risk_level not in ("high", "critical"):
return True
request = ApprovalRequest(
tool_name=tool_name,
arguments=arguments,
risk_level=risk_level,
description=description,
agent_id=agent_id,
session_id=session_id
)
loop = asyncio.get_event_loop()
future = loop.create_future()
request_id = str(uuid4())
self._pending[request_id] = future
try:
decision = await asyncio.wait_for(future, timeout=300.0)
return decision == ApprovalDecision.APPROVED
except asyncio.TimeoutError:
logger.warning(f"Approval timeout for request {request_id}")
return False
finally:
self._pending.pop(request_id, None)
async def resolve(self, request_id: str, decision: ApprovalDecision):
"""Called by the approval UI/webhook when a human makes a decision."""
if request_id in self._pending:
self._pending[request_id].set_result(decision)
Sandboxed MCP Server Execution
For untrusted MCP servers, run them in isolated containers:
# mcp_sandbox.py
import subprocess
import os
from pathlib import Path
class SandboxedMCPServer:
"""Runs MCP servers in Docker containers with strict resource limits."""
def __init__(self, server_command: str, container_name: str,
memory_limit: str = "256m", cpu_limit: float = 0.5,
network_access: bool = False):
self.server_command = server_command
self.container_name = container_name
self.memory_limit = memory_limit
self.cpu_limit = cpu_limit
self.network_access = network_access
self._process: subprocess.Popen | None = None
def start(self) -> tuple[int, int]:
"""Start the sandboxed server and return (stdin_fd, stdout_fd)."""
docker_args = [
"docker", "run", "--rm",
f"--name={self.container_name}",
f"--memory={self.memory_limit}",
f"--cpus={self.cpu_limit}",
"--read-only",
"--cap-drop=ALL",
"--security-opt=no-new-privileges",
]
if not self.network_access:
docker_args.append("--network=none")
# Mount necessary files read-only
docker_args.extend([
"-v", f"{self._get_server_path()}:/server:ro",
"--workdir", "/server",
"python", "-m", "mcp_server"
])
self._process = subprocess.Popen(
docker_args,
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
bufsize=0
)
return self._process.stdin.fileno(), self._process.stdout.fileno()
def stop(self):
"""Gracefully stop the sandboxed server."""
if self._process:
self._process.terminate()
try:
self._process.wait(timeout=10)
except subprocess.TimeoutExpired:
self._process.kill()
Scaling Patterns for High-Throughput MCP Agents
When you move from single-agent development to multi-tenant production systems, the architecture must evolve to handle concurrent sessions, shared tool access, and resource contention.
Connection Pooling with Session Isolation
# mcp_pool.py
import asyncio
from collections import deque
from contextlib import asynccontextmanager
class MCPSessionPool:
"""
Connection pool for MCP servers that maintains session isolation.
Critical for multi-tenant deployments where each tenant needs their own context.
"""
def __init__(self, server_params: StdioServerParameters,
pool_size: int = 20, max_queue_size: int = 100):
self.server_params = server_params
self.pool_size = pool_size
self.max_queue_size = max_queue_size
self._pool: asyncio.Queue[ClientSession] = asyncio.Queue(maxsize=pool_size)
self._initialized = False
self._lock = asyncio.Lock()
async def initialize(self):
"""Pre-warm the pool with initial connections."""
async with self._lock:
if self._initialized:
return
tasks = [self._create_session() for _ in range(self.pool_size)]
results = await asyncio.gather(*tasks, return_exceptions=True)
for result in results:
if not isinstance(result, Exception):
await self._pool.put(result)
self._initialized = True
async def _create_session(self) -> ClientSession:
"""Create a new MCP client session."""
from mcp import stdio_client
async with stdio_client(self.server_params) as (read, write):
session = ClientSession(read, write)
await session.initialize()
return session
@asynccontextmanager
async def get_session(self) -> AsyncGenerator[ClientSession, None]:
"""
Borrow a session from the pool.
Sessions are returned to the pool after use for reuse.
"""
try:
session = await asyncio.wait_for(self._pool.get(), timeout=10.0)
except asyncio.TimeoutError:
# Pool exhausted - create an overflow session
session = await self._create_session()
overflow = True
else:
overflow = False
try:
yield session
except Exception:
# Session may be corrupted - replace it
if not overflow:
await self._replace_session(session)
raise
else:
if not overflow:
await self._pool.put(session)
async def _replace_session(self, old_session: ClientSession):
"""Replace a corrupted session with a fresh one."""
try:
await old_session.close()
except Exception:
pass
new_session = await self._create_session()
await self._pool.put(new_session)
async def shutdown(self):
"""Gracefully shut down all pooled sessions."""
while not self._pool.empty():
session = self._pool.get_nowait()
try:
await session.close()
except Exception:
pass
Multi-Server Agent Architecture
Production agents often need access to multiple MCP servers simultaneously. The routing layer determines which server handles which tool call:
# mcp_router.py
from typing import Protocol
class MCPToolProvider(Protocol):
"""Interface for any MCP server that provides tools."""
async def list_tools(self) -> list[dict]:
...
async def call_tool(self, tool_name: str, arguments: dict) -> str:
...
class MCPRouter:
"""Routes tool calls to the appropriate MCP server."""
def __init__(self):
self._providers: dict[str, MCPToolProvider] = {}
self._tool_index: dict[str, str] = {} # tool_name -> provider_name
def register_provider(self, name: str, provider: MCPToolProvider):
self._providers[name] = provider
async def build_tool_index(self):
"""Build an index of all available tools across all providers."""
for name, provider in self._providers.items():
tools = await provider.list_tools()
for tool in tools:
tool_name = tool["name"]
# Handle naming conflicts with prefixing
if tool_name in self._tool_index:
prefixed_name = f"{name}__{tool_name}"
self._tool_index[prefixed_name] = name
tool["name"] = prefixed_name
else:
self._tool_index[tool_name] = name
async def route_call(self, tool_name: str, arguments: dict) -> str:
"""Route a tool call to the correct provider."""
provider_name = self._tool_index.get(tool_name)
if not provider_name:
raise UnknownToolError(f"Tool '{tool_name}' not found")
provider = self._providers[provider_name]
return await provider.call_tool(tool_name, arguments)
def get_all_tools(self) -> list[dict]:
"""Return all tools with their provider information."""
all_tools = []
for tool_name, provider_name in self._tool_index.items():
provider = self._providers[provider_name]
# Fetch full tool definition
tools = provider.list_tools()
for tool in tools:
if tool["name"] == tool_name:
tool["provider"] = provider_name
all_tools.append(tool)
return all_tools
State Management Across Sessions
Autonomous agents often need to maintain state across multiple interactions. MCP's stateless protocol design means this responsibility falls entirely on the agent side.
Checkpointing Pattern
# mcp_state.py
from dataclasses import dataclass, field
from typing import Any
@dataclass
class AgentCheckpoint:
session_id: str
step: int
completed_actions: list[dict] = field(default_factory=list)
pending_actions: list[dict] = field(default_factory=list)
context: dict[str, Any] = field(default_factory=dict)
tool_results_cache: dict[str, str] = field(default_factory=dict)
def to_dict(self) -> dict:
return {
"session_id": self.session_id,
"step": self.step,
"completed_actions": self.completed_actions,
"pending_actions": self.pending_actions,
"context": self.context,
"tool_results_cache": self.tool_results_cache
}
@classmethod
def from_dict(cls, data: dict) -> "AgentCheckpoint":
return cls(**data)
class CheckpointStore:
"""Persists agent state for crash recovery and session resumption."""
def __init__(self, backend: str = "redis"):
self._backend = backend
self._prefix = "mcp:agent:checkpoint:"
async def save(self, checkpoint: AgentCheckpoint):
"""Save a checkpoint to persistent storage."""
data = checkpoint.to_dict()
if self._backend == "redis":
await self._redis.set(
self._prefix + checkpoint.session_id,
json.dumps(data),
ex=86400 # 24 hour TTL
)
async def load(self, session_id: str) -> AgentCheckpoint | None:
"""Load the most recent checkpoint for a session."""
if self._backend == "redis":
data = await self._redis.get(self._prefix + session_id)
if data:
return AgentCheckpoint.from_dict(json.loads(data))
return None
async def resume_execution(self, session_id: str,
agent: "MCPAgent") -> None:
"""Resume agent execution from the last checkpoint."""
checkpoint = await self.load(session_id)
if not checkpoint:
raise NoCheckpointError(f"No checkpoint found for {session_id}")
# Restore tool results cache to avoid redundant API calls
agent.tool_results_cache = checkpoint.tool_results_cache
agent.completed_actions = checkpoint.completed_actions
# Continue from the next pending action
for action in checkpoint.pending_actions:
result = await agent.execute_action(action)
await self.save(agent.current_checkpoint())
Putting It All Together: A Production MCP Agent
The following example demonstrates how all these patterns integrate into a cohesive production agent:
# production_agent.py
import asyncio
import json
from dataclasses import dataclass
from typing import Any, Callable
@dataclass
class AgentConfig:
model: str = "gpt-4-turbo"
max_steps: int = 20
temperature: float = 0.0
tool_timeout: float = 30.0
enable_checkpointing: bool = True
require_approval_above: str = "medium" # risk level threshold
class ProductionMCPAgent:
"""
A production-ready MCP agent with lifecycle management,
observability, security, and state persistence.
"""
def __init__(self, config: AgentConfig,
mcp_params: StdioServerParameters,
approval_gate: ApprovalGate | None = None):
self.config = config
self.mcp_params = mcp_params
self.approval_gate = approval_gate
# Infrastructure components
self.pool = MCPSessionPool(mcp_params, pool_size=5)
self.tracer = MCPTracer(service_name="production-agent")
self.checkpoint_store = CheckpointStore()
self.policy = ToolPermissionPolicy()
# Agent state
self.session_id: str | None = None
self.step: int = 0
self.tool_results_cache: dict[str, str] = {}
self.completed_actions: list[dict] = []
async def run(self, user_input: str, session_id: str | None = None) -> str:
"""Execute the agent loop with full production safeguards."""
self.session_id = session_id or str(uuid4())
# Attempt to resume from checkpoint
if self.config.enable_checkpointing and session_id:
checkpoint = await self.checkpoint_store.load(session_id)
if checkpoint:
self.step = checkpoint.step
self.tool_results_cache = checkpoint.tool_results_cache
self.completed_actions = checkpoint.completed_actions
await self.pool.initialize()
messages = [{"role": "user", "content": user_input}]
final_answer = None
while self.step < self.config.max_steps:
self.step += 1
# Get model response
response = await self._call_model(messages)
if response.tool_calls:
for tool_call in response.tool_calls:
result = await self._execute_tool_call(
tool_call,
self.session_id
)
messages.append(result)
else:
final_answer = response.content
break
# Checkpoint after each step
if self.config.enable_checkpointing:
await self._save_checkpoint()
return final_answer or "Max steps reached without final answer"
async def _execute_tool_call(self, tool_call, session_id: str) -> dict:
"""Execute a tool call with security checks and observability."""
tool_name = tool_call.function.name
arguments = json.loads(tool_call.function.arguments)
# Security: validate against policy
validated_args = self.policy.validate_call(tool_name, arguments)
# Security: check approval gate
risk_level = self.policy.get_risk_level(tool_name)
if self.approval_gate and risk_level in ("high", "critical"):
approved = await self.approval_gate.check(
tool_name=tool_name,
arguments=validated_args,
risk_level=risk_level,
description=f"Agent requests {tool_name}",
agent_id=self.session_id,
session_id=session_id
)
if not approved:
return {
"role": "tool",
"tool_call_id": tool_call.id,
"content": "Action denied by approval gate"
}
# Cache check
cache_key = f"{tool_name}:{json.dumps(validated_args, sort_keys=True)}"
if cache_key in self.tool_results_cache:
return {
"role": "tool",
"tool_call_id": tool_call.id,
"content": self.tool_results_cache[cache_key]
}
# Execute with tracing
async with self.pool.get_session() as session:
result = await self.tracer.traced_call(
session, tool_name, validated_args, session_id
)
result_text = result.content[0].text
self.tool_results_cache[cache_key] = result_text
self.completed_actions.append({
"tool": tool_name,
"arguments": validated_args,
"result_length": len(result_text)
})
return {
"role": "tool",
"tool_call_id": tool_call.id,
"content": result_text
}
async def _call_model(self, messages: list[dict]) -> Any:
"""Call the LLM with current message history."""
# Implementation depends on your LLM provider
...
async def _save_checkpoint(self):
"""Persist current state for crash recovery."""
checkpoint = AgentCheckpoint(
session_id=self.session_id,
step=self.step,
completed_actions=self.completed_actions,
tool_results_cache=self.tool_results_cache
)
await self.checkpoint_store.save(checkpoint)
Deployment Considerations
Kubernetes Configuration
For containerized deployments, the MCP server typically runs as a sidecar or separate pod:
# k8s-mcp-agent.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: mcp-agent
spec:
replicas: 3
selector:
matchLabels:
app: mcp-agent
template:
metadata:
labels:
app: mcp-agent
spec:
containers:
- name: agent
image: your-registry/mcp-agent:latest
resources:
requests:
cpu: "500m"
memory: "512Mi"
limits:
cpu: "2"
memory: "1Gi"
env:
- name: MCP_SERVER_COMMAND
value: "python -m mcp_server"
- name: POOL_SIZE
value: "10"
- name: REDIS_URL
valueFrom:
secretKeyRef:
name: redis-credentials
key: url
livenessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 30
periodSeconds: 10
readinessProbe:
httpGet:
path: /ready
port: 8080
initialDelaySeconds: 15
periodSeconds: 5
---
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: mcp-agent-hpa
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: mcp-agent
minReplicas: 3
maxReplicas: 20
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
- type: Pods
pods:
metric:
name: mcp_active_sessions
target:
type: AverageValue
averageValue: "50"
Cost Optimization
MCP agents can be expensive to run. The following patterns reduce costs without sacrificing reliability:
-
Tool result caching: Identical tool calls with identical arguments should return cached results. This alone can reduce LLM token usage by 30-50% in typical workloads.
-
Tiered model routing: Use a smaller, cheaper model for routine tool calls and escalate to a larger model only when complex reasoning is needed.
-
Batch processing: When multiple independent tool calls are needed, batch them into a single model response rather than making sequential calls.
-
Context compression: Periodically summarize long conversation histories to reduce token consumption on subsequent calls.
Common Pitfalls and How to Avoid Them
Pitfall 1: Infinite Agent Loops
Agents that don't converge waste resources and frustrate users. Implement hard limits:
# Loop detection
class LoopDetector:
def __init__(self, max_repeats: int = 3):
self.max_repeats = max_repeats
self._action_history: list[str] = []
def check(self, action: dict) -> bool:
"""Returns True if this action appears to be looping."""
action_key = json.dumps(action, sort_keys=True)
self._action_history.append(action_key)
# Check for repeated sequences
recent = self._action_history[-self.max_repeats * 2:]
for i in range(len(recent) - self.max_repeats):
if (recent[i:i+self.max_repeats] ==
recent[i+self.max_repeats:i+self.max_repeats*2]):
return True
return False
Pitfall 2: Tool Argument Drift
Models occasionally hallucinate argument names or types. Always validate against the tool schema:
def validate_tool_arguments(tool_name: str, arguments: dict,
tool_schema: dict) -> dict:
"""Validate arguments against the tool's JSON schema."""
errors = []
# Check for unknown arguments
known_args = tool_schema.get("properties", {})
for key in arguments:
if key not in known_args:
errors.append(f"Unknown argument: '{key}'")
# Check for missing