Skip to main content

Building a Complete MCP System — The Ultimate Integration

Calculating read time…
You've learned MCP piece by piece: resources, tools, servers, multi-model architectures. Now it's time to put it all together into something real—a production-grade system that actually works.

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

Functional 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

┌─────────────────────────────────────────────────────────┐ │ CLIENT LAYER │ │ (Web UI, Mobile App, API Clients) │ └────────────────────┬────────────────────────────────────┘ │ ⬍ HTTP/WebSocket │ ┌────────────────────▼────────────────────────────────────┐ │ LANGGRAPH AGENT │ │ ┌──────────────────────────────────────────────┐ │ │ │ Supervisor Node (Routes to specialists) │ │ │ │ │ │ │ │ │ ├─→ Knowledge Base Agent │ │ │ │ ├─→ Customer Data Agent │ │ │ │ ├─→ Ticketing Agent │ │ │ │ └─→ Email Agent │ │ │ └──────────────────────────────────────────────┘ │ └──────┬──────────┬──────────┬──────────┬────────────────┘ │ │ │ │ ⬍ ⬍ ⬍ ⬍ ┌──────▼────┬─────▼────┬─────▼────┬─────▼──────┐ │ Knowledge │ Customer │ Ticket │ Email │ │ MCP │ MCP │ MCP │ MCP │ │ Server │ Server │ Server │ Server │ └───────┬───┴──────┬───┴──────┬───┴──────┬─────┘ │ │ │ │ ⬍ ⬍ ⬍ ⬍ ┌───────▼──────────▼──────────▼──────────▼─────┐ │ DATA & SERVICES LAYER │ │ Vector DB | PostgreSQL | SMTP | Ticket API │ └───────────────────────────────────────────────┘
Architecture Highlights:

• 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())
⚠️ Important Note:

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

Build a Restaurant Ordering System:

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

Before Deploying to Production:

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

Architecture Patterns:

• 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
Production Essentials:

• Security first: Auth, rate limiting, input validation
• Observability: Logs, metrics, traces
• Error handling: Graceful failures, retries, fallbacks
• Testing: Unit, integration, load testing
Integration Patterns:

• MultiServerMCPClient: Connect to multiple servers
• Tool binding: Expose tools to LLM
• State graphs: Define agent workflows
• Checkpointing: Persist conversation state

Comments