Imagine you're building a customer support AI that needs to: search your knowledge base, query customer databases, send emails, create support tickets, analyze sentiment, and route to human agents when needed.
This isn't just one tool—it's an entire ecosystem of capabilities working together seamlessly. This is what we're building today.
In this capstone guide, you'll create a complete end-to-end system combining LangGraph (for agent orchestration), MCP (for tool access), multiple specialized servers, authentication, error handling, monitoring, and structured output.
By the end, you'll have built something you can actually deploy and use in production! 🚀
🎯 What We're Building: Customer Support AI System
Let's define exactly what our complete system will do.
System Requirements
✅ Accept customer questions via API or chat
✅ Search knowledge base for answers
✅ Query customer database for account info
✅ Create support tickets when needed
✅ Send automated email responses
✅ Analyze sentiment and urgency
✅ Route complex issues to human agents
✅ Return structured responses with confidence scores
Non-Functional Requirements:
✅ Authentication and authorization
✅ Rate limiting and cost control
✅ Comprehensive logging
✅ Error handling and retries
✅ Performance monitoring
✅ Graceful degradation
System Architecture
• LangGraph orchestrates multi-step workflows
• Supervisor pattern routes to specialized agents
• Multiple MCP servers provide domain-specific tools
• Clean separation between orchestration and execution
• Scalable design allows adding new capabilities easily
🛠️ Step 1: Building the MCP Servers
First, let's create our specialized MCP servers. Each handles one domain.
Knowledge Base MCP Server
Create knowledge_server.py:
"""
Knowledge Base MCP Server
Provides semantic search over support documents
"""
from mcp.server.fastmcp import FastMCP, Context
import asyncio
from typing import List, Optional
import hashlib
mcp = FastMCP("Knowledge Base Server")
# Mock vector database
# In production: use Pinecone, Weaviate, or ChromaDB
KNOWLEDGE_BASE = [
{
"id": "kb001",
"title": "How to reset password",
"content": "To reset your password, click 'Forgot Password' on the login page. Enter your email and follow the instructions.",
"category": "account",
"tags": ["password", "login", "security"]
},
{
"id": "kb002",
"title": "Refund policy",
"content": "We offer full refunds within 30 days of purchase. Contact support@company.com to initiate a refund.",
"category": "billing",
"tags": ["refund", "billing", "money"]
},
{
"id": "kb003",
"title": "Shipping times",
"content": "Standard shipping: 5-7 business days. Express shipping: 2-3 business days. International: 10-14 business days.",
"category": "shipping",
"tags": ["delivery", "shipping", "time"]
}
]
def simple_similarity(query: str, document: dict) -> float:
"""Calculate simple keyword similarity"""
query_words = set(query.lower().split())
doc_words = set(
(document["title"] + " " + document["content"] + " " +
" ".join(document["tags"])).lower().split()
)
intersection = query_words & doc_words
union = query_words | doc_words
return len(intersection) / len(union) if union else 0.0
@mcp.tool()
async def search_knowledge_base(
query: str,
limit: int = 3,
ctx: Context = None
) -> dict:
"""
Search knowledge base for relevant articles
Args:
query: User's question or keywords
limit: Maximum number of results
Returns:
Ranked list of relevant articles
"""
if ctx:
await ctx.info(f"Searching knowledge base: {query}")
# Calculate similarity scores
results = []
for doc in KNOWLEDGE_BASE:
score = simple_similarity(query, doc)
if score > 0:
results.append({
**doc,
"relevance_score": round(score, 3)
})
# Sort by relevance
results.sort(key=lambda x: x["relevance_score"], reverse=True)
return {
"success": True,
"query": query,
"results": results[:limit],
"total_found": len(results)
}
@mcp.tool()
async def get_article_by_id(
article_id: str,
ctx: Context = None
) -> dict:
"""Get full article by ID"""
article = next((a for a in KNOWLEDGE_BASE if a["id"] == article_id), None)
if article:
return {
"success": True,
"article": article
}
else:
return {
"success": False,
"error": f"Article {article_id} not found"
}
if __name__ == "__main__":
import sys
print("🔍 Knowledge Base Server starting...", file=sys.stderr)
mcp.run()
Customer Data MCP Server
Create customer_server.py:
dict:
"""
Look up customer information by email
Args:
email: Customer's email address
Returns:
Customer profile and account details
"""
if ctx:
await ctx.info(f"Looking up customer: {email}")
customer = CUSTOMERS.get(email.lower())
if customer:
return {
"success": True,
"customer": customer
}
else:
return {
"success": False,
"error": "Customer not found",
"email": email
}
@mcp.tool()
async def get_customer_history(
customer_id: str,
ctx: Context = None
) -> dict:
"""Get customer's support ticket history"""
# Mock history
history = {
"cust_001": [
{"ticket": "TKT-123", "issue": "Billing question", "status": "resolved"},
{"ticket": "TKT-089", "issue": "Feature request", "status": "closed"}
],
"cust_002": []
}
return {
"success": True,
"customer_id": customer_id,
"history": history.get(customer_id, [])
}
if __name__ == "__main__":
import sys
print("👤 Customer Data Server starting...", file=sys.stderr)
mcp.run()
Ticketing MCP Server
Create ticket_server.py:
"""
Ticket Management MCP Server
Creates and manages support tickets
"""
from mcp.server.fastmcp import FastMCP, Context
from datetime import datetime
import random
mcp = FastMCP("Ticket Server")
# In-memory ticket storage
TICKETS = {}
@mcp.tool()
async def create_support_ticket(
customer_email: str,
subject: str,
description: str,
priority: str = "normal",
ctx: Context = None
) -> dict:
"""
Create a new support ticket
Args:
customer_email: Customer's email
subject: Brief description of issue
description: Detailed explanation
priority: low, normal, high, urgent
Returns:
Created ticket with ID and tracking info
"""
if ctx:
await ctx.info(f"Creating ticket for {customer_email}")
# Generate ticket ID
ticket_id = f"TKT-{random.randint(10000, 99999)}"
ticket = {
"ticket_id": ticket_id,
"customer_email": customer_email,
"subject": subject,
"description": description,
"priority": priority,
"status": "open",
"created_at": datetime.utcnow().isoformat(),
"assigned_to": "auto_assign",
"estimated_response_time": "2 hours" if priority == "high" else "24 hours"
}
TICKETS[ticket_id] = ticket
if ctx:
await ctx.info(f"Ticket created: {ticket_id}")
return {
"success": True,
"ticket": ticket
}
@mcp.tool()
async def get_ticket_status(
ticket_id: str,
ctx: Context = None
) -> dict:
"""Get current status of a ticket"""
ticket = TICKETS.get(ticket_id)
if ticket:
return {
"success": True,
"ticket": ticket
}
else:
return {
"success": False,
"error": f"Ticket {ticket_id} not found"
}
if __name__ == "__main__":
import sys
print("🎫 Ticket Server starting...", file=sys.stderr)
mcp.run()
🧠 Step 2: Building the LangGraph Agent
Now let's create the intelligent agent that orchestrates everything.
Install Dependencies
pip install langgraph langchain langchain-openai langchain-mcp-adapters python-dotenv
Create the Agent
Create support_agent.py:
"""
Complete Customer Support Agent with LangGraph
Integrates multiple MCP servers for full workflow
"""
import asyncio
import os
from typing import TypedDict, Annotated, Sequence
from datetime import datetime
from dotenv import load_dotenv
from langchain_openai import ChatOpenAI
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage, SystemMessage
from langchain_mcp_adapters.client import MultiServerMCPClient
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolNode
from langgraph.checkpoint.memory import MemorySaver
# Load environment variables
load_dotenv()
# ============================================
# STATE DEFINITION
# ============================================
class AgentState(TypedDict):
"""State that flows through the agent"""
messages: Annotated[Sequence[BaseMessage], "Conversation history"]
customer_email: str | None
customer_data: dict | None
knowledge_results: list | None
ticket_created: dict | None
confidence_score: float
requires_human: bool
next_action: str
# ============================================
# MCP SERVERS CONFIGURATION
# ============================================
MCP_SERVERS = {
"knowledge": {
"command": "python",
"args": ["knowledge_server.py"],
"transport": "stdio"
},
"customer": {
"command": "python",
"args": ["customer_server.py"],
"transport": "stdio"
},
"ticketing": {
"command": "python",
"args": ["ticket_server.py"],
"transport": "stdio"
}
}
# ============================================
# HELPER FUNCTIONS
# ============================================
def calculate_confidence(state: AgentState) -> float:
"""Calculate confidence score based on available data"""
score = 0.5 # Base score
# Boost if we found knowledge articles
if state.get("knowledge_results"):
score += 0.2
# Boost if we have customer context
if state.get("customer_data"):
score += 0.15
# Boost if answer is complete
last_message = state["messages"][-1]
if isinstance(last_message, AIMessage) and len(last_message.content) > 100:
score += 0.15
return min(score, 1.0)
def requires_human_intervention(state: AgentState) -> bool:
"""Determine if human agent is needed"""
last_message = state["messages"][-1].content.lower()
# Keywords that trigger human escalation
escalation_keywords = [
"angry", "furious", "lawsuit", "legal",
"unacceptable", "manager", "refund immediately"
]
for keyword in escalation_keywords:
if keyword in last_message:
return True
# Low confidence requires human
if state.get("confidence_score", 0) < 0.6:
return True
return False
# ============================================
# AGENT NODES
# ============================================
async def analyze_query(state: AgentState) -> AgentState:
"""
Initial analysis of customer query
Extracts email and determines query type
"""
# Extract email from query if present
query = state["messages"][-1].content
# Simple email extraction (in production, use regex)
words = query.split()
email = None
for word in words:
if "@" in word:
email = word.strip(".,;:!?")
break
state["customer_email"] = email
state["next_action"] = "search_knowledge"
return state
async def search_knowledge_node(state: AgentState) -> AgentState:
"""Search knowledge base for relevant articles"""
# This would be called by LangGraph's tool system
# For now, we'll mark that it should search
state["next_action"] = "check_customer"
return state
async def check_customer_node(state: AgentState) -> AgentState:
"""Look up customer data if email provided"""
# Customer lookup would happen via tools
state["next_action"] = "generate_response"
return state
async def generate_response_node(state: AgentState) -> AgentState:
"""Generate final response with all context"""
# Calculate confidence
state["confidence_score"] = calculate_confidence(state)
state["requires_human"] = requires_human_intervention(state)
if state["requires_human"]:
state["next_action"] = "escalate"
else:
state["next_action"] = "complete"
return state
async def escalate_node(state: AgentState) -> AgentState:
"""Escalate to human agent"""
# Create urgent ticket
state["next_action"] = "complete"
return state
# ============================================
# AGENT CONSTRUCTION
# ============================================
async def create_support_agent():
"""Build the complete LangGraph agent"""
# Initialize LLM
model = ChatOpenAI(model="gpt-4o", temperature=0)
# Connect to MCP servers and get tools
async with MultiServerMCPClient(MCP_SERVERS) as client:
tools = client.get_tools()
print(f"✅ Loaded {len(tools)} tools from MCP servers")
for tool in tools:
print(f" • {tool.name}")
# Create model with tools bound
model_with_tools = model.bind_tools(tools)
# Build the state graph
workflow = StateGraph(AgentState)
# Add nodes
workflow.add_node("analyze", analyze_query)
workflow.add_node("agent", lambda state: agent_node(state, model_with_tools))
workflow.add_node("tools", ToolNode(tools))
workflow.add_node("escalate", escalate_node)
# Add edges
workflow.set_entry_point("analyze")
workflow.add_edge("analyze", "agent")
# Conditional routing based on agent decision
workflow.add_conditional_edges(
"agent",
should_continue,
{
"tools": "tools",
"escalate": "escalate",
"end": END
}
)
workflow.add_edge("tools", "agent") # After tools, back to agent
workflow.add_edge("escalate", END)
# Compile with memory
memory = MemorySaver()
app = workflow.compile(checkpointer=memory)
return app, client
async def agent_node(state: AgentState, model):
"""Main agent reasoning node"""
# Add system message with instructions
system_msg = SystemMessage(content="""
You are a helpful customer support AI assistant.
Your capabilities:
- Search knowledge base for answers
- Look up customer information
- Create support tickets
- Escalate to human agents when needed
Guidelines:
1. Always search knowledge base first
2. Look up customer data if email is provided
3. Be empathetic and professional
4. Create tickets for complex issues
5. Escalate angry or legal matters to humans
Provide clear, concise answers with relevant article references.
""")
messages = [system_msg] + state["messages"]
response = await model.ainvoke(messages)
return {"messages": [response]}
def should_continue(state: AgentState) -> str:
"""Determine next step based on agent's response"""
last_message = state["messages"][-1]
# Check if agent wants to use tools
if hasattr(last_message, "tool_calls") and last_message.tool_calls:
return "tools"
# Check if escalation needed
if state.get("requires_human"):
return "escalate"
return "end"
# ============================================
# MAIN EXECUTION
# ============================================
async def main():
"""Run the complete support agent system"""
print("🚀 Starting Customer Support Agent System\n")
print("="*60)
# Create agent
app, mcp_client = await create_support_agent()
# Test queries
test_queries = [
"Hi, I forgot my password. My email is john@example.com",
"What's your refund policy?",
"This is unacceptable! I want a refund NOW or I'm calling my lawyer!"
]
for i, query in enumerate(test_queries, 1):
print(f"\n{'='*60}")
print(f"TEST QUERY {i}:")
print(f"{'='*60}")
print(f"User: {query}\n")
# Initialize state
initial_state = {
"messages": [HumanMessage(content=query)],
"customer_email": None,
"customer_data": None,
"knowledge_results": None,
"ticket_created": None,
"confidence_score": 0.0,
"requires_human": False,
"next_action": "start"
}
# Run agent
config = {"configurable": {"thread_id": f"test_{i}"}}
try:
final_state = await app.ainvoke(initial_state, config)
# Extract response
ai_response = final_state["messages"][-1].content
print(f"Agent: {ai_response}\n")
# Show metadata
print("📊 Metadata:")
print(f" Confidence: {final_state.get('confidence_score', 0):.2f}")
print(f" Human Needed: {final_state.get('requires_human', False)}")
print(f" Customer: {final_state.get('customer_email', 'Unknown')}")
except Exception as e:
print(f"❌ Error: {str(e)}")
print(f"\n{'='*60}")
print("🎉 Demo Complete!")
print(f"{'='*60}\n")
if __name__ == "__main__":
asyncio.run(main())
This is a simplified version for learning.
The next sections add authentication, monitoring, and production features!
🔐 Step 3: Adding Security & Authentication
Production systems must be secure. Let's add proper authentication.
API Key Authentication
Create auth.py:
"""
Authentication and Authorization for MCP Servers
"""
import os
import hashlib
import secrets
from typing import Optional
from functools import wraps
# In production: Use a database
API_KEYS = {
"sk_test_12345": {
"user_id": "user_001",
"name": "Development Key",
"permissions": ["read", "write"],
"rate_limit": 100
},
"sk_prod_67890": {
"user_id": "user_002",
"name": "Production Key",
"permissions": ["read"],
"rate_limit": 1000
}
}
class AuthenticationError(Exception):
"""Raised when authentication fails"""
pass
class AuthorizationError(Exception):
"""Raised when authorization fails"""
pass
def validate_api_key(api_key: str) -> dict:
"""
Validate API key and return user info
Args:
api_key: The API key to validate
Returns:
User information dict
Raises:
AuthenticationError: If key is invalid
"""
if not api_key:
raise AuthenticationError("API key required")
user_info = API_KEYS.get(api_key)
if not user_info:
raise AuthenticationError("Invalid API key")
return user_info
def check_permission(user_info: dict, required_permission: str) -> bool:
"""Check if user has required permission"""
return required_permission in user_info.get("permissions", [])
def require_auth(permission: str = "read"):
"""Decorator to require authentication for MCP tools"""
def decorator(func):
@wraps(func)
async def wrapper(*args, **kwargs):
# Extract API key from context
# In real implementation, this comes from HTTP headers
api_key = kwargs.get("api_key") or os.getenv("API_KEY")
try:
user_info = validate_api_key(api_key)
if not check_permission(user_info, permission):
raise AuthorizationError(
f"Permission '{permission}' required"
)
# Add user info to kwargs
kwargs["user_info"] = user_info
return await func(*args, **kwargs)
except (AuthenticationError, AuthorizationError) as e:
return {
"success": False,
"error": str(e),
"error_code": "AUTH_ERROR"
}
return wrapper
return decorator
def generate_api_key() -> str:
"""Generate a new API key"""
return f"sk_prod_{secrets.token_urlsafe(32)}"
# Usage in MCP tools:
# @mcp.tool()
# @require_auth(permission="write")
# async def sensitive_operation(...):
# ...
Rate Limiting
Create rate_limiter.py:
"""
Rate Limiting for MCP Tools
"""
import time
from collections import defaultdict
from typing import Optional
class RateLimiter:
"""Token bucket rate limiter"""
def __init__(self):
self.buckets = defaultdict(lambda: {
"tokens": 0,
"last_update": time.time()
})
def is_allowed(
self,
user_id: str,
limit: int = 100,
window: int = 60
) -> tuple[bool, Optional[str]]:
"""
Check if request is allowed
Args:
user_id: User identifier
limit: Max requests per window
window: Time window in seconds
Returns:
(allowed, error_message)
"""
now = time.time()
bucket = self.buckets[user_id]
# Refill tokens based on time passed
time_passed = now - bucket["last_update"]
tokens_to_add = (time_passed / window) * limit
bucket["tokens"] = min(
limit,
bucket["tokens"] + tokens_to_add
)
bucket["last_update"] = now
# Check if we have tokens available
if bucket["tokens"] >= 1:
bucket["tokens"] -= 1
return True, None
else:
wait_time = int((1 - bucket["tokens"]) * (window / limit))
return False, f"Rate limit exceeded. Try again in {wait_time}s"
# Global rate limiter instance
rate_limiter = RateLimiter()
def rate_limit(limit: int = 100, window: int = 60):
"""Decorator to enforce rate limits"""
def decorator(func):
async def wrapper(*args, **kwargs):
user_info = kwargs.get("user_info", {})
user_id = user_info.get("user_id", "anonymous")
allowed, error_msg = rate_limiter.is_allowed(
user_id,
limit,
window
)
if not allowed:
return {
"success": False,
"error": error_msg,
"error_code": "RATE_LIMIT_EXCEEDED"
}
return await func(*args, **kwargs)
return wrapper
return decorator
📊 Step 4: Monitoring & Observability
You can't improve what you can't measure. Let's add monitoring.
Create monitoring.py:
"""
Monitoring and Observability for MCP System
"""
import time
import json
from datetime import datetime
from typing import Any, Optional
from dataclasses import dataclass, asdict
from collections import defaultdict
@dataclass
class MetricEvent:
"""Single metric event"""
timestamp: str
event_type: str
tool_name: str
duration_ms: float
success: bool
user_id: Optional[str] = None
error: Optional[str] = None
metadata: Optional[dict] = None
class MetricsCollector:
"""Collect and aggregate metrics"""
def __init__(self):
self.events = []
self.counters = defaultdict(int)
self.timers = {}
def record_event(self, event: MetricEvent):
"""Record a metric event"""
self.events.append(event)
# Update counters
self.counters[f"{event.tool_name}_calls"] += 1
if event.success:
self.counters[f"{event.tool_name}_success"] += 1
else:
self.counters[f"{event.tool_name}_errors"] += 1
def start_timer(self, operation_id: str):
"""Start timing an operation"""
self.timers[operation_id] = time.time()
def end_timer(self, operation_id: str) -> float:
"""End timer and return duration in ms"""
if operation_id not in self.timers:
return 0.0
duration = (time.time() - self.timers[operation_id]) * 1000
del self.timers[operation_id]
return duration
def get_summary(self) -> dict:
"""Get metrics summary"""
return {
"total_events": len(self.events),
"counters": dict(self.counters),
"recent_events": [
asdict(e) for e in self.events[-10:]
]
}
def export_to_file(self, filepath: str):
"""Export metrics to JSON file"""
with open(filepath, 'w') as f:
json.dump({
"events": [asdict(e) for e in self.events],
"summary": self.get_summary()
}, f, indent=2)
# Global metrics collector
metrics = MetricsCollector()
def monitor_tool(func):
"""Decorator to monitor tool execution"""
async def wrapper(*args, **kwargs):
tool_name = func.__name__
operation_id = f"{tool_name}_{time.time()}"
metrics.start_timer(operation_id)
success = True
error = None
result = None
try:
result = await func(*args, **kwargs)
return result
except Exception as e:
success = False
error = str(e)
raise
finally:
duration = metrics.end_timer(operation_id)
# Extract user info if available
user_info = kwargs.get("user_info", {})
user_id = user_info.get("user_id")
# Record event
event = MetricEvent(
timestamp=datetime.utcnow().isoformat(),
event_type="tool_execution",
tool_name=tool_name,
duration_ms=duration,
success=success,
user_id=user_id,
error=error,
metadata={
"args_count": len(args),
"kwargs_count": len(kwargs)
}
)
metrics.record_event(event)
# Log to console
status = "✅" if success else "❌"
print(
f"{status} {tool_name} | {duration:.2f}ms | "
f"User: {user_id or 'anonymous'}",
flush=True
)
return wrapper
# Usage:
# @mcp.tool()
# @monitor_tool
# @rate_limit(limit=100)
# @require_auth(permission="read")
# async def my_tool(...):
# ...
🎯 Step 5: Complete Production Example
Let's put it all together with a fully production-ready example!
Create production_system.py:
"""
Complete Production-Ready Customer Support System
Combines: LangGraph + MCP + Auth + Monitoring + Error Handling
"""
import asyncio
import sys
from typing import TypedDict, Annotated, Sequence
from datetime import datetime
from langchain_openai import ChatOpenAI
from langchain_core.messages import BaseMessage, HumanMessage, SystemMessage
from langchain_mcp_adapters.client import MultiServerMCPClient
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolNode
from langgraph.checkpoint.memory import MemorySaver
# Our custom modules
from auth import validate_api_key, AuthenticationError
from rate_limiter import rate_limiter
from monitoring import metrics, MetricEvent
# ============================================
# STATE & CONFIG
# ============================================
class SupportState(TypedDict):
"""Complete agent state"""
messages: Sequence[BaseMessage]
customer_email: str | None
customer_tier: str | None
query_category: str | None
tools_used: list[str]
confidence: float
escalate: bool
response_time: float
user_id: str | None
MCP_CONFIG = {
"knowledge": {
"command": "python",
"args": ["knowledge_server.py"],
"transport": "stdio"
},
"customer": {
"command": "python",
"args": ["customer_server.py"],
"transport": "stdio"
},
"ticketing": {
"command": "python",
"args": ["ticket_server.py"],
"transport": "stdio"
}
}
# ============================================
# SYSTEM PROMPT
# ============================================
SYSTEM_PROMPT = """
You are an expert customer support AI assistant.
WORKFLOW:
1. Understand the customer's issue
2. Search knowledge base for relevant solutions
3. Look up customer account details if email provided
4. Provide clear, actionable answers
5. Create tickets for complex issues
6. Escalate when necessary
ESCALATION TRIGGERS:
- Legal threats or mentions of lawyers
- Extreme frustration or anger
- Request to speak with manager
- Issues beyond your knowledge
RESPONSE GUIDELINES:
- Be empathetic and professional
- Reference specific knowledge articles
- Personalize based on customer tier
- Offer proactive solutions
- End with clear next steps
Always prioritize customer satisfaction!
"""
# ============================================
# AGENT NODES
# ============================================
async def authenticate_user(state: SupportState) -> SupportState:
"""Validate user authentication"""
# In real system, extract from API headers
api_key = os.getenv("API_KEY", "sk_test_12345")
try:
user_info = validate_api_key(api_key)
state["user_id"] = user_info["user_id"]
# Check rate limit
allowed, error_msg = rate_limiter.is_allowed(
user_info["user_id"],
limit=user_info["rate_limit"]
)
if not allowed:
state["escalate"] = True
state["messages"].append(
SystemMessage(content=f"Rate limit: {error_msg}")
)
except AuthenticationError as e:
state["escalate"] = True
state["messages"].append(
SystemMessage(content=f"Auth error: {str(e)}")
)
return state
async def analyze_sentiment(state: SupportState) -> SupportState:
"""Analyze customer sentiment"""
query = state["messages"][-1].content.lower()
# Simple sentiment analysis
negative_words = ["angry", "furious", "terrible", "awful", "hate"]
urgent_words = ["urgent", "immediately", "asap", "now"]
negative_count = sum(1 for word in negative_words if word in query)
urgent_count = sum(1 for word in urgent_words if word in query)
if negative_count >= 2 or urgent_count >= 2:
state["escalate"] = True
return state
async def agent_node(state: SupportState, model) -> SupportState:
"""Main agent decision node"""
system_msg = SystemMessage(content=SYSTEM_PROMPT)
messages = [system_msg] + list(state["messages"])
response = await model.ainvoke(messages)
state["messages"] = list(state["messages"]) + [response]
return state
def should_continue(state: SupportState) -> str:
"""Routing logic"""
# Check for escalation
if state.get("escalate", False):
return "escalate"
# Check if agent wants to use tools
last_msg = state["messages"][-1]
if hasattr(last_msg, "tool_calls") and last_msg.tool_calls:
return "tools"
return "end"
async def escalation_node(state: SupportState) -> SupportState:
"""Handle escalation to human"""
escalation_msg = SystemMessage(content="""
This issue has been escalated to a human agent.
Our team will respond within 1 hour for priority customers,
or 24 hours for standard support.
""")
state["messages"] = list(state["messages"]) + [escalation_msg]
# Log escalation
metrics.record_event(MetricEvent(
timestamp=datetime.utcnow().isoformat(),
event_type="escalation",
tool_name="human_agent",
duration_ms=0,
success=True,
user_id=state.get("user_id"),
metadata={"reason": "automated_escalation"}
))
return state
# ============================================
# BUILD SYSTEM
# ============================================
async def build_production_system():
"""Construct complete production system"""
print("Building production system...\n")
# Initialize LLM
model = ChatOpenAI(
model="gpt-4o",
temperature=0,
request_timeout=30
)
# Connect to MCP servers
async with MultiServerMCPClient(MCP_CONFIG) as mcp_client:
tools = mcp_client.get_tools()
print(f"Connected to {len(MCP_CONFIG)} MCP servers")
print(f"Loaded {len(tools)} tools\n")
# Bind tools to model
model_with_tools = model.bind_tools(tools)
# Build graph
workflow = StateGraph(SupportState)
# Add nodes
workflow.add_node("auth", authenticate_user)
workflow.add_node("sentiment", analyze_sentiment)
workflow.add_node("agent", lambda s: agent_node(s, model_with_tools))
workflow.add_node("tools", ToolNode(tools))
workflow.add_node("escalate", escalation_node)
# Define flow
workflow.set_entry_point("auth")
workflow.add_edge("auth", "sentiment")
workflow.add_edge("sentiment", "agent")
workflow.add_conditional_edges(
"agent",
should_continue,
{
"tools": "tools",
"escalate": "escalate",
"end": END
}
)
workflow.add_edge("tools", "agent")
workflow.add_edge("escalate", END)
# Compile
memory = MemorySaver()
app = workflow.compile(checkpointer=memory)
return app, mcp_client
# ============================================
# RUN SYSTEM
# ============================================
async def run_support_query(app, query: str, thread_id: str):
"""Execute a support query"""
start_time = time.time()
initial_state = {
"messages": [HumanMessage(content=query)],
"customer_email": None,
"customer_tier": None,
"query_category": None,
"tools_used": [],
"confidence": 0.0,
"escalate": False,
"response_time": 0.0,
"user_id": None
}
config = {"configurable": {"thread_id": thread_id}}
try:
final_state = await app.ainvoke(initial_state, config)
response_time = (time.time() - start_time) * 1000
final_state["response_time"] = response_time
# Extract AI response
ai_response = final_state["messages"][-1].content
return {
"success": True,
"response": ai_response,
"metadata": {
"response_time_ms": response_time,
"confidence": final_state.get("confidence", 0),
"escalated": final_state.get("escalate", False),
"tools_used": final_state.get("tools_used", [])
}
}
except Exception as e:
return {
"success": False,
"error": str(e),
"metadata": {
"response_time_ms": (time.time() - start_time) * 1000
}
}
# ============================================
# MAIN
# ============================================
async def main():
"""Run production system demo"""
print("\n" + "="*70)
print(" PRODUCTION CUSTOMER SUPPORT SYSTEM")
print("="*70 + "\n")
# Build system
app, mcp_client = await build_production_system()
# Test queries
test_cases = [
{
"id": "case_1",
"query": "Hi, my email is john@example.com and I forgot my password",
"expected": "Should find password reset article and customer data"
},
{
"id": "case_2",
"query": "What's your refund policy?",
"expected": "Should return refund policy from knowledge base"
},
{
"id": "case_3",
"query": "This is unacceptable! I demand a refund NOW or I'm calling my lawyer!",
"expected": "Should escalate to human"
}
]
for test in test_cases:
print(f"\n{'─'*70}")
print(f"📝 TEST: {test['id']}")
print(f"{'─'*70}")
print(f"Query: {test['query']}")
print(f"Expected: {test['expected']}\n")
result = await run_support_query(app, test["query"], test["id"])
if result["success"]:
print(f"✅ Response: {result['response']}\n")
print("📊 Metadata:")
for key, value in result["metadata"].items():
print(f" {key}: {value}")
else:
print(f"❌ Error: {result['error']}")
# Print metrics summary
print(f"\n{'='*70}")
print(" SYSTEM METRICS")
print(f"{'='*70}")
summary = metrics.get_summary()
print(f"\nTotal Events: {summary['total_events']}")
print("\nCounters:")
for key, value in summary['counters'].items():
print(f" {key}: {value}")
# Export metrics
metrics.export_to_file("metrics.json")
print("\n✅ Metrics exported to metrics.json")
print(f"\n{'='*70}")
print("🎉 DEMO COMPLETE!")
print(f"{'='*70}\n")
if __name__ == "__main__":
import time
asyncio.run(main())
🎯 Exercise: Build Your Own System
Now it's your turn! Build a simplified version of this system.
Exercise Requirements
MCP Servers (3):
1. Menu Server - browse menu, get prices
2. Order Server - create orders, track status
3. Inventory Server - check item availability
LangGraph Agent:
• Take customer order via natural language
• Check menu and availability
• Calculate total with tax
• Create order
• Return order confirmation
Features Required:
✅ Error handling (out of stock items)
✅ Structured output (order summary JSON)
✅ Basic monitoring (print execution time)
Solution Template
Get started with this template:
# menu_server.py
from mcp.server.fastmcp import FastMCP
mcp = FastMCP("Menu Server")
MENU = {
"burger": {"price": 12.99, "category": "main"},
"fries": {"price": 4.99, "category": "side"},
"soda": {"price": 2.99, "category": "drink"}
}
@mcp.tool()
async def get_menu(category: str = None) -> dict:
"""Get menu items, optionally filtered by category"""
if category:
items = {k: v for k, v in MENU.items()
if v["category"] == category}
else:
items = MENU
return {"success": True, "items": items}
@mcp.tool()
async def get_item_price(item_name: str) -> dict:
"""Get price for specific item"""
item = MENU.get(item_name.lower())
if item:
return {"success": True, "price": item["price"]}
else:
return {"success": False, "error": "Item not found"}
if __name__ == "__main__":
mcp.run()
Complete the exercise by implementing:
- Order Server (create_order, get_order_status)
- Inventory Server (check_availability, update_stock)
- LangGraph agent that orchestrates the ordering workflow
📊 Production Deployment Checklist
Security:
□ API key authentication implemented
□ Rate limiting enabled
□ Input validation on all tools
□ SQL injection prevention
□ HTTPS for all connections
Reliability:
□ Error handling on every tool
□ Retry logic for transient failures
□ Timeout configuration
□ Circuit breakers for external services
□ Graceful degradation
Observability:
□ Comprehensive logging
□ Metrics collection
□ Performance monitoring
□ Error tracking (Sentry, etc.)
□ Distributed tracing
Scalability:
□ Load testing completed
□ Horizontal scaling configured
□ Database connection pooling
□ Caching strategy implemented
□ CDN for static assets
🎓 Key Takeaways
• Separation of concerns: LangGraph (orchestration) + MCP (execution)
• Domain-driven servers: Each server owns one responsibility
• Supervisor pattern: Central agent routes to specialists
• State management: LangGraph maintains context across steps
• Security first: Auth, rate limiting, input validation
• Observability: Logs, metrics, traces
• Error handling: Graceful failures, retries, fallbacks
• Testing: Unit, integration, load testing
• MultiServerMCPClient: Connect to multiple servers
• Tool binding: Expose tools to LLM
• State graphs: Define agent workflows
• Checkpointing: Persist conversation state
Comments
Post a Comment