Building Production-Ready Multi-Agent AI Systems on Amazon EKS: The Complete Step-by-Step Guide
Single AI agents are powerful. Multi-agent systems are transformative. This guide builds a complete, 2026-10-8 06:14:6 Author: hackernoon.com(查看原文) 阅读量:2 收藏

Single AI agents are powerful. Multi-agent systems are transformative. This guide builds a complete, production-ready multi-agent AI system on Amazon EKS, the exact architecture used by financial services, healthcare, and enterprise SaaS teams in 2026. Five specialised agents (Customer, Payment, Fraud, Compliance, Customer Support), a Custom Agent Router, Bedrock Guardrails, RAG with OpenSearch, response verification, WAF protection, and human approval gates. Every component configured from scratch. Every line of code included.

What You Are Building

Look at the architecture diagram. Every component in it has a job, and they all work together to serve a user request safely, accurately, and at enterprise scale.

Here is the full data flow, before a single line of code:

User → App (mobile/web)
  → Route 53 (DNS lookup)
  → CloudFront (CDN, content closer to user)
  → WAF (blocks malicious traffic)
  → Cognito (authenticates user, issues token)
  → API Gateway (receives task + login token)
  → Private VPC Link (secure internal routing)
  → ALB (distributes to healthy agent pods)
  → EKS (runs the containerised agent application)
  → Custom Agent Router (chooses the right agent)
  → Specialist Agent (Customer / Payment / Fraud / Compliance / Support)
    → Company Knowledge (S3 → OpenSearch → Bedrock Knowledge Bases)
    → Model & Safety (Bedrock Guardrails → Claude via Bedrock)
    → Business Tools (AgentCore Gateway → Business APIs)
    → Human Approval (for sensitive actions — Set Functions → Human)
  → Response Verification (custom code on EKS checks the response)
    → If rejected: return to agent for correction
    → If accepted: send to user
  → User sees the response

That is the system. This guide builds every layer of it, step by step.

What you need:

  • AWS account with admin access
  • AWS CLI configured (aws configure)
  • kubectl installed
  • Python 3.11+
  • Docker
  • About 3–4 hours

AWS services used: Route 53, CloudFront, WAF, S3, Cognito, API Gateway, VPC, ALB, EKS, Amazon Bedrock, Bedrock Knowledge Bases, OpenSearch Serverless, Bedrock Guardrails, AgentCore Gateway, Lambda, Step Functions, CloudWatch, IAM

Step 1: Foundation VPC, EKS Cluster, and IAM

Everything runs inside a VPC. Start here.

1.1 Create the VPC

bash

REGION="eu-west-1"
ACCOUNT_ID=$(aws sts get-caller-identity --query Account --output text)
CLUSTER_NAME="multi-agent-prod"

# Create VPC with public and private subnets
aws ec2 create-vpc \
  --cidr-block 10.0.0.0/16 \
  --tag-specifications 'ResourceType=vpc,Tags=[{Key=Name,Value=multi-agent-vpc},{Key=Project,Value=multi-agent}]' \
  --region $REGION

VPC_ID=$(aws ec2 describe-vpcs \
  --filters "Name=tag:Name,Values=multi-agent-vpc" \
  --query 'Vpcs[0].VpcId' \
  --output text \
  --region $REGION)

echo "VPC ID: $VPC_ID"

# Enable DNS hostnames (required for EKS)
aws ec2 modify-vpc-attribute \
  --vpc-id $VPC_ID \
  --enable-dns-hostnames \
  --region $REGION

# Create Internet Gateway for public subnets
IGW_ID=$(aws ec2 create-internet-gateway \
  --query 'InternetGateway.InternetGatewayId' \
  --output text \
  --region $REGION)

aws ec2 attach-internet-gateway \
  --internet-gateway-id $IGW_ID \
  --vpc-id $VPC_ID \
  --region $REGION

# Create public subnets (for ALB and NAT Gateway)
PUBLIC_SUBNET_1=$(aws ec2 create-subnet \
  --vpc-id $VPC_ID \
  --cidr-block 10.0.1.0/24 \
  --availability-zone ${REGION}a \
  --tag-specifications 'ResourceType=subnet,Tags=[{Key=Name,Value=public-1a},{Key=kubernetes.io/role/elb,Value=1}]' \
  --query 'Subnet.SubnetId' \
  --output text \
  --region $REGION)

PUBLIC_SUBNET_2=$(aws ec2 create-subnet \
  --vpc-id $VPC_ID \
  --cidr-block 10.0.2.0/24 \
  --availability-zone ${REGION}b \
  --tag-specifications 'ResourceType=subnet,Tags=[{Key=Name,Value=public-1b},{Key=kubernetes.io/role/elb,Value=1}]' \
  --query 'Subnet.SubnetId' \
  --output text \
  --region $REGION)

# Create private subnets (for EKS nodes — never directly exposed)
PRIVATE_SUBNET_1=$(aws ec2 create-subnet \
  --vpc-id $VPC_ID \
  --cidr-block 10.0.10.0/24 \
  --availability-zone ${REGION}a \
  --tag-specifications 'ResourceType=subnet,Tags=[{Key=Name,Value=private-1a},{Key=kubernetes.io/role/internal-elb,Value=1}]' \
  --query 'Subnet.SubnetId' \
  --output text \
  --region $REGION)

PRIVATE_SUBNET_2=$(aws ec2 create-subnet \
  --vpc-id $VPC_ID \
  --cidr-block 10.0.11.0/24 \
  --availability-zone ${REGION}b \
  --tag-specifications 'ResourceType=subnet,Tags=[{Key=Name,Value=private-1b},{Key=kubernetes.io/role/internal-elb,Value=1}]' \
  --query 'Subnet.SubnetId' \
  --output text \
  --region $REGION)

echo "Subnets created: $PUBLIC_SUBNET_1, $PUBLIC_SUBNET_2, $PRIVATE_SUBNET_1, $PRIVATE_SUBNET_2"

1.2 Create the EKS Cluster

bash

# Create the EKS cluster using eksctl (simplest approach)
# Install eksctl if not present:
# brew install eksctl  (macOS)
# or download from https://eksctl.io

cat > cluster-config.yaml << EOF
apiVersion: eksctl.io/v1alpha5
kind: ClusterConfig

metadata:
  name: $CLUSTER_NAME
  region: $REGION
  version: "1.31"
  tags:
    Project: multi-agent
    Environment: production

vpc:
  id: $VPC_ID
  subnets:
    private:
      private-1a:
        id: $PRIVATE_SUBNET_1
      private-1b:
        id: $PRIVATE_SUBNET_2
    public:
      public-1a:
        id: $PUBLIC_SUBNET_1
      public-1b:
        id: $PUBLIC_SUBNET_2

managedNodeGroups:
  - name: agent-nodes
    instanceType: m5.xlarge   # 4 vCPU, 16GB — good for agent workloads
    minSize: 2
    maxSize: 10
    desiredCapacity: 3
    privateNetworking: true   # Nodes in private subnets only
    labels:
      role: agent
    tags:
      k8s.io/cluster-autoscaler/enabled: "true"
      k8s.io/cluster-autoscaler/$CLUSTER_NAME: "owned"
    iam:
      attachPolicyARNs:
        - arn:aws:iam::aws:policy/AmazonEKSWorkerNodePolicy
        - arn:aws:iam::aws:policy/AmazonEKS_CNI_Policy
        - arn:aws:iam::aws:policy/AmazonEC2ContainerRegistryReadOnly
      withAddonPolicies:
        autoScaler: true
        cloudWatch: true

addons:
  - name: vpc-cni
  - name: coredns
  - name: kube-proxy
  - name: aws-ebs-csi-driver
EOF

eksctl create cluster -f cluster-config.yaml

# This takes 15-20 minutes
# Get credentials when done:
aws eks update-kubeconfig --region $REGION --name $CLUSTER_NAME

# Verify cluster is running
kubectl get nodes

1.3 Create Core IAM Roles

bash

# IAM role for the agent pods — allows Bedrock, S3, DynamoDB access
cat > agent-pod-trust.json << EOF
{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Principal": {
      "Federated": "arn:aws:iam::$ACCOUNT_ID:oidc-provider/$(aws eks describe-cluster \
        --name $CLUSTER_NAME \
        --query 'cluster.identity.oidc.issuer' \
        --output text \
        --region $REGION | sed 's/https:\/\///')"
    },
    "Action": "sts:AssumeRoleWithWebIdentity",
    "Condition": {
      "StringEquals": {
        "$(aws eks describe-cluster \
          --name $CLUSTER_NAME \
          --query 'cluster.identity.oidc.issuer' \
          --output text \
          --region $REGION | sed 's/https:\/\///')":
              "system:serviceaccount:multi-agent:agent-service-account"
      }
    }
  }]
}
EOF

aws iam create-role \
  --role-name multi-agent-pod-role \
  --assume-role-policy-document file://agent-pod-trust.json

cat > agent-pod-policy.json << EOF
{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Action": [
        "bedrock:InvokeModel",
        "bedrock:InvokeModelWithResponseStream",
        "bedrock:ApplyGuardrail",
        "bedrock:Retrieve",
        "bedrock:RetrieveAndGenerate"
      ],
      "Resource": "*"
    },
    {
      "Effect": "Allow",
      "Action": [
        "s3:GetObject",
        "s3:ListBucket"
      ],
      "Resource": [
        "arn:aws:s3:::multi-agent-company-docs",
        "arn:aws:s3:::multi-agent-company-docs/*"
      ]
    },
    {
      "Effect": "Allow",
      "Action": [
        "dynamodb:GetItem",
        "dynamodb:PutItem",
        "dynamodb:UpdateItem",
        "dynamodb:Query"
      ],
      "Resource": "arn:aws:dynamodb:$REGION:$ACCOUNT_ID:table/agent-sessions"
    },
    {
      "Effect": "Allow",
      "Action": [
        "cloudwatch:PutMetricData",
        "logs:CreateLogGroup",
        "logs:CreateLogStream",
        "logs:PutLogEvents"
      ],
      "Resource": "*"
    },
    {
      "Effect": "Allow",
      "Action": [
        "sns:Publish"
      ],
      "Resource": "arn:aws:sns:$REGION:$ACCOUNT_ID:human-approval-topic"
    }
  ]
}
EOF

aws iam put-role-policy \
  --role-name multi-agent-pod-role \
  --policy-name agent-pod-permissions \
  --policy-document file://agent-pod-policy.json

echo "IAM role created: arn:aws:iam::$ACCOUNT_ID:role/multi-agent-pod-role"

Step 2: Company Knowledge S3, OpenSearch, and Bedrock Knowledge Base

The bottom-left of your architecture diagram. Every agent reads from here when it needs company information.

2.1 Create the S3 Document Store

bash

# S3 bucket for company documents
aws s3 mb s3://multi-agent-company-docs-$ACCOUNT_ID --region $REGION

# Upload sample company documents
mkdir -p company-docs/{policies,products,procedures,compliance}

cat > company-docs/policies/refund-policy.txt << 'EOF'
REFUND POLICY — Last updated July 2026

Standard refunds: Customers may request a full refund within 30 days of purchase.
Duplicate charges: Any duplicate charge is eligible for immediate refund regardless of date.
Subscription cancellations: Pro-rated refunds available for annual subscriptions cancelled within 90 days.
Fraud cases: Fraudulent charges are escalated to the fraud team and resolved within 24 hours.
Approval required: Refunds above £500 require manager approval.
EOF

cat > company-docs/procedures/account-verification.txt << 'EOF'
ACCOUNT VERIFICATION PROCEDURE

KYC Requirements:
- Identity verification required for accounts over £1,000 monthly limit
- Photo ID + proof of address for business accounts
- Enhanced due diligence for high-risk jurisdictions

AML Checks:
- All transactions over £10,000 flagged for review
- Sanctions screening on all new accounts
- PEP (Politically Exposed Person) screening monthly
EOF

cat > company-docs/products/subscription-tiers.txt << 'EOF'
SUBSCRIPTION TIERS — 2026

Free Plan: Up to 100 API calls/month, community support
Starter: £29/month, 10,000 API calls, email support
Professional: £99/month, 100,000 API calls, priority support, SLA 99.9%
Enterprise: Custom pricing, unlimited calls, dedicated support, SLA 99.99%

Upgrades: Take effect immediately, charged pro-rated
Downgrades: Take effect at next billing cycle
EOF

# Upload all documents to S3
aws s3 sync company-docs/ s3://multi-agent-company-docs-$ACCOUNT_ID/ \
  --region $REGION

echo "Company documents uploaded to S3"

2.2 Create OpenSearch Serverless (Vector Store)

bash

# Create OpenSearch Serverless collection for vector search
# This is what "OpenSearch (Serverless Vector Store)" in your diagram refers to

cat > opensearch-policy.json << EOF
[{
  "Rules": [
    {
      "Resource": ["collection/multi-agent-knowledge"],
      "Permission": [
        "aoss:CreateCollectionItems",
        "aoss:DeleteCollectionItems",
        "aoss:UpdateCollectionItems",
        "aoss:DescribeCollectionItems"
      ],
      "ResourceType": "collection"
    },
    {
      "Resource": ["index/multi-agent-knowledge/*"],
      "Permission": [
        "aoss:CreateIndex",
        "aoss:DeleteIndex",
        "aoss:UpdateIndex",
        "aoss:DescribeIndex",
        "aoss:ReadDocument",
        "aoss:WriteDocument"
      ],
      "ResourceType": "index"
    }
  ],
  "Principal": [
    "arn:aws:iam::$ACCOUNT_ID:role/multi-agent-pod-role",
    "arn:aws:iam::$ACCOUNT_ID:role/AmazonBedrockExecutionRoleForKnowledgeBase"
  ]
}]
EOF

aws opensearchserverless create-access-policy \
  --name multi-agent-data-policy \
  --type data \
  --policy "$(cat opensearch-policy.json | jq -c .)" \
  --region $REGION

aws opensearchserverless create-collection \
  --name multi-agent-knowledge \
  --type VECTORSEARCH \
  --description "Company knowledge vector store for multi-agent system" \
  --region $REGION

# Wait for collection to be active
echo "Waiting for OpenSearch collection to become active..."
aws opensearchserverless wait collection-active \
  --id $(aws opensearchserverless list-collections \
    --collection-filters name=multi-agent-knowledge \
    --query 'collectionSummaries[0].id' \
    --output text \
    --region $REGION) \
  --region $REGION || sleep 120

OPENSEARCH_ENDPOINT=$(aws opensearchserverless list-collections \
  --collection-filters name=multi-agent-knowledge \
  --query 'collectionSummaries[0].collectionEndpoint' \
  --output text \
  --region $REGION)

echo "OpenSearch endpoint: $OPENSEARCH_ENDPOINT"

2.3 Create Bedrock Knowledge Base

bash

# Create IAM role for Bedrock Knowledge Base
cat > bedrock-kb-trust.json << EOF
{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Principal": {"Service": "bedrock.amazonaws.com"},
    "Action": "sts:AssumeRole",
    "Condition": {
      "StringEquals": {"aws:SourceAccount": "$ACCOUNT_ID"}
    }
  }]
}
EOF

aws iam create-role \
  --role-name AmazonBedrockExecutionRoleForKnowledgeBase \
  --assume-role-policy-document file://bedrock-kb-trust.json

aws iam attach-role-policy \
  --role-name AmazonBedrockExecutionRoleForKnowledgeBase \
  --policy-arn arn:aws:iam::aws:policy/AmazonS3ReadOnlyAccess

# Create the Knowledge Base
KB_ID=$(aws bedrock-agent create-knowledge-base \
  --name "multi-agent-company-knowledge" \
  --description "Company knowledge base for multi-agent system" \
  --role-arn arn:aws:iam::$ACCOUNT_ID:role/AmazonBedrockExecutionRoleForKnowledgeBase \
  --knowledge-base-configuration '{
    "type": "VECTOR",
    "vectorKnowledgeBaseConfiguration": {
      "embeddingModelArn": "arn:aws:bedrock:us-east-1::foundation-model/amazon.titan-embed-text-v2:0"
    }
  }' \
  --storage-configuration "{
    \"type\": \"OPENSEARCH_SERVERLESS\",
    \"opensearchServerlessConfiguration\": {
      \"collectionArn\": \"arn:aws:aoss:$REGION:$ACCOUNT_ID:collection/multi-agent-knowledge\",
      \"vectorIndexName\": \"company-knowledge-index\",
      \"fieldMapping\": {
        \"vectorField\": \"embedding\",
        \"textField\": \"text\",
        \"metadataField\": \"metadata\"
      }
    }
  }" \
  --region $REGION \
  --query 'knowledgeBase.knowledgeBaseId' \
  --output text)

echo "Knowledge Base ID: $KB_ID"

# Add S3 data source to the Knowledge Base
aws bedrock-agent create-data-source \
  --knowledge-base-id $KB_ID \
  --name "company-documents-s3" \
  --data-source-configuration "{
    \"type\": \"S3\",
    \"s3Configuration\": {
      \"bucketArn\": \"arn:aws:s3:::multi-agent-company-docs-$ACCOUNT_ID\"
    }
  }" \
  --region $REGION

# Trigger initial ingestion
aws bedrock-agent start-ingestion-job \
  --knowledge-base-id $KB_ID \
  --data-source-id $(aws bedrock-agent list-data-sources \
    --knowledge-base-id $KB_ID \
    --query 'dataSourceSummaries[0].dataSourceId' \
    --output text \
    --region $REGION) \
  --region $REGION

echo "Knowledge base ingestion started"

Step 3: Model and Safety Bedrock Guardrails

The centre-bottom of your diagram. This is what stands between the model and potentially harmful outputs.

bash

# Create Bedrock Guardrails — the safety layer for all agents
GUARDRAIL_ID=$(aws bedrock create-guardrail \
  --name "multi-agent-safety-guardrail" \
  --description "Safety guardrails for all agents in the multi-agent system" \
  --topic-policy-config '{
    "topicsConfig": [
      {
        "name": "Financial advice without authorisation",
        "definition": "Providing specific investment advice, portfolio recommendations, or financial planning guidance without appropriate authorisation",
        "examples": [
          "You should invest all your savings in...",
          "I recommend you buy..."
        ],
        "type": "DENY"
      },
      {
        "name": "Medical advice",
        "definition": "Providing medical diagnoses, treatment recommendations, or health guidance",
        "type": "DENY"
      },
      {
        "name": "Legal advice",
        "definition": "Providing specific legal advice or legal interpretations",
        "type": "DENY"
      }
    ]
  }' \
  --content-policy-config '{
    "filtersConfig": [
      {"type": "SEXUAL", "inputStrength": "HIGH", "outputStrength": "HIGH"},
      {"type": "VIOLENCE", "inputStrength": "MEDIUM", "outputStrength": "HIGH"},
      {"type": "HATE", "inputStrength": "HIGH", "outputStrength": "HIGH"},
      {"type": "INSULTS", "inputStrength": "MEDIUM", "outputStrength": "MEDIUM"},
      {"type": "MISCONDUCT", "inputStrength": "HIGH", "outputStrength": "HIGH"},
      {"type": "PROMPT_ATTACK", "inputStrength": "HIGH", "outputStrength": "NONE"}
    ]
  }' \
  --word-policy-config '{
    "wordsConfig": [
      {"text": "competitor_name_1"},
      {"text": "confidential_product_code"}
    ],
    "managedWordListsConfig": [
      {"type": "PROFANITY"}
    ]
  }' \
  --sensitive-information-policy-config '{
    "piiEntitiesConfig": [
      {"type": "CREDIT_DEBIT_CARD_NUMBER", "action": "BLOCK"},
      {"type": "CREDIT_DEBIT_CVV", "action": "BLOCK"},
      {"type": "US_SOCIAL_SECURITY_NUMBER", "action": "ANONYMIZE"},
      {"type": "UK_NATIONAL_INSURANCE_NUMBER", "action": "ANONYMIZE"},
      {"type": "EMAIL", "action": "ANONYMIZE"},
      {"type": "PHONE", "action": "ANONYMIZE"}
    ]
  }' \
  --blocked-input-messaging "I cannot process this request as it contains content that violates our usage policy." \
  --blocked-outputs-messaging "I cannot provide this response as it contains content that violates our policy. Please contact support for assistance." \
  --region $REGION \
  --query 'guardrailId' \
  --output text)

# Create guardrail version
GUARDRAIL_VERSION=$(aws bedrock create-guardrail-version \
  --guardrail-id $GUARDRAIL_ID \
  --region $REGION \
  --query 'version' \
  --output text)

echo "Guardrail ID: $GUARDRAIL_ID, Version: $GUARDRAIL_VERSION"

Step 4: The Five Specialist Agents Core Python Code

Now for the heart of the system. Each agent in your diagram does one specific job. Here is the complete Python implementation.

python

# agents/base_agent.py
# Base class that all five agents inherit from

import boto3
import json
import logging
import time
import os
from datetime import datetime, timezone
from typing import Optional
from dataclasses import dataclass

logger = logging.getLogger(__name__)

REGION = os.environ.get("AWS_REGION", "eu-west-1")
MODEL_ID = "anthropic.claude-opus-4-7-20250514"
GUARDRAIL_ID = os.environ.get("GUARDRAIL_ID")
GUARDRAIL_VERSION = os.environ.get("GUARDRAIL_VERSION", "1")
KNOWLEDGE_BASE_ID = os.environ.get("KNOWLEDGE_BASE_ID")

bedrock = boto3.client('bedrock-runtime', region_name='us-east-1')
bedrock_agent = boto3.client('bedrock-agent-runtime', region_name='us-east-1')
cloudwatch = boto3.client('cloudwatch', region_name=REGION)


@dataclass
class AgentResponse:
    """Standardised response from any agent"""
    agent_type: str
    success: bool
    content: str
    requires_human_approval: bool = False
    approval_reason: str = ""
    knowledge_sources: list = None
    tokens_used: int = 0
    latency_ms: int = 0
    session_id: str = ""


class BaseAgent:
    """
    Base class for all five specialist agents.

    Each specialist agent:
    1. Has a specific system prompt defining its role and limits
    2. Retrieves company knowledge from the shared Knowledge Base
    3. Calls Claude via Bedrock with Guardrails applied
    4. Identifies whether human approval is required
    5. Returns a structured AgentResponse
    """

    agent_type: str = "base"
    system_prompt: str = ""
    high_value_threshold: float = 500.0

    def __init__(self):
        self.bedrock = bedrock
        self.bedrock_agent = bedrock_agent
        self.cloudwatch = cloudwatch

    def retrieve_company_knowledge(self, query: str, num_results: int = 5) -> tuple[str, list]:
        """
        Retrieve relevant company knowledge from Bedrock Knowledge Base.
        This is the RAG step — find relevant documents before calling the model.
        """
        try:
            response = self.bedrock_agent.retrieve(
                knowledgeBaseId=KNOWLEDGE_BASE_ID,
                retrievalQuery={"text": query},
                retrievalConfiguration={
                    "vectorSearchConfiguration": {
                        "numberOfResults": num_results,
                        "overrideSearchType": "HYBRID"
                    }
                }
            )

            results = response.get("retrievalResults", [])
            if not results:
                return "", []

            # Format retrieved context
            context_parts = []
            sources = []
            for i, result in enumerate(results, 1):
                content = result["content"]["text"]
                source = result.get("location", {}).get("s3Location", {}).get("uri", "company-docs")
                score = result.get("score", 0)

                context_parts.append(f"[Source {i}: {source} (relevance: {score:.2f})]\n{content}")
                sources.append({"source": source, "score": score})

            return "\n\n".join(context_parts), sources

        except Exception as e:
            logger.error(f"Knowledge retrieval failed: {e}")
            return "", []

    def call_bedrock_with_guardrails(
        self,
        messages: list,
        additional_context: str = "",
        max_tokens: int = 1000
    ) -> tuple[str, int]:
        """
        Call Claude via Bedrock with Guardrails applied.
        Guardrails check both the input and the output.
        """
        full_system = self.system_prompt
        if additional_context:
            full_system += f"\n\n## Relevant Company Knowledge:\n{additional_context}"

        body = {
            "anthropic_version": "bedrock-2023-05-31",
            "max_tokens": max_tokens,
            "system": full_system,
            "messages": messages
        }

        # Apply guardrails to the request
        if GUARDRAIL_ID:
            body["amazon-bedrock-guardrailConfig"] = {
                "guardrailIdentifier": GUARDRAIL_ID,
                "guardrailVersion": GUARDRAIL_VERSION,
                "trace": "enabled"
            }

        start_time = time.time()
        response = self.bedrock.invoke_model(
            modelId=MODEL_ID,
            body=json.dumps(body)
        )
        latency_ms = int((time.time() - start_time) * 1000)

        result = json.loads(response['body'].read())

        # Check if guardrails blocked the response
        guardrail_action = result.get("amazon-bedrock-guardrailAction", "NONE")
        if guardrail_action == "GUARDRAIL_INTERVENED":
            logger.warning(f"Guardrail intervened for {self.agent_type}")
            content = result.get("content", [{}])[0].get("text",
                "I cannot provide that response due to our content policy.")
        else:
            content = result["content"][0]["text"]

        total_tokens = result["usage"]["input_tokens"] + result["usage"]["output_tokens"]
        return content, total_tokens

    def check_requires_human_approval(self, request: str, response: str) -> tuple[bool, str]:
        """
        Determine if this response requires human approval before being sent.
        Each specialist agent overrides this with its own rules.
        """
        return False, ""

    def handle(self, request: str, session_id: str, user_id: str) -> AgentResponse:
        """Main handler — called by the Agent Router"""
        start_time = time.time()
        logger.info(f"{self.agent_type} handling request for session {session_id}")

        # Step 1: Retrieve relevant company knowledge (RAG)
        knowledge_context, sources = self.retrieve_company_knowledge(request)

        # Step 2: Build conversation messages
        messages = [{"role": "user", "content": f"User ID: {user_id}\n\nRequest: {request}"}]

        # Step 3: Call the model with guardrails
        content, tokens_used = self.call_bedrock_with_guardrails(
            messages=messages,
            additional_context=knowledge_context
        )

        # Step 4: Check human approval requirement
        requires_approval, approval_reason = self.check_requires_human_approval(request, content)

        latency_ms = int((time.time() - start_time) * 1000)

        # Step 5: Publish metrics
        self.publish_metrics(tokens_used, latency_ms, requires_approval)

        return AgentResponse(
            agent_type=self.agent_type,
            success=True,
            content=content,
            requires_human_approval=requires_approval,
            approval_reason=approval_reason,
            knowledge_sources=sources,
            tokens_used=tokens_used,
            latency_ms=latency_ms,
            session_id=session_id
        )

    def publish_metrics(self, tokens: int, latency: int, requires_approval: bool):
        """Publish CloudWatch metrics for cost and performance monitoring"""
        try:
            self.cloudwatch.put_metric_data(
                Namespace="MultiAgentSystem/Agents",
                MetricData=[
                    {
                        "MetricName": "TokensConsumed",
                        "Value": tokens,
                        "Unit": "Count",
                        "Dimensions": [{"Name": "AgentType", "Value": self.agent_type}]
                    },
                    {
                        "MetricName": "RequestLatencyMs",
                        "Value": latency,
                        "Unit": "Milliseconds",
                        "Dimensions": [{"Name": "AgentType", "Value": self.agent_type}]
                    },
                    {
                        "MetricName": "HumanApprovalRequired",
                        "Value": 1 if requires_approval else 0,
                        "Unit": "Count",
                        "Dimensions": [{"Name": "AgentType", "Value": self.agent_type}]
                    }
                ]
            )
        except Exception as e:
            logger.warning(f"Metrics publish failed: {e}")

python

# agents/specialist_agents.py
# The five specialist agents from your architecture diagram

from .base_agent import BaseAgent
import re


class CustomerAgent(BaseAgent):
    """
    Handles customer account requests:
    - Account information queries
    - Profile updates
    - Subscription changes
    - Billing history
    """
    agent_type = "customer"
    system_prompt = """You are the Customer Account Agent for our platform.

Your responsibilities:
- Answer questions about the customer's account, subscription, and billing history
- Help customers update their profile information
- Explain subscription tiers and help customers choose the right plan
- Process subscription upgrades and downgrades per company policy

Always:
- Verify you are speaking about the correct account (reference User ID)
- Be accurate about subscription terms and pricing from company documentation
- Never process payments directly — route to Payment Agent
- Never access other customers' information

You have access to company knowledge about subscription tiers, policies, and procedures."""

    def check_requires_human_approval(self, request: str, response: str) -> tuple[bool, str]:
        """Account changes above certain thresholds need approval"""
        # Account deletion always needs human confirmation
        if any(term in request.lower() for term in ["delete account", "close account", "cancel everything"]):
            return True, "Account deletion requires human confirmation and data retention review"
        return False, ""


class PaymentAgent(BaseAgent):
    """
    Handles approved payment tasks:
    - Refund processing
    - Payment dispute investigation
    - Invoice generation
    - Subscription billing issues
    """
    agent_type = "payment"
    system_prompt = """You are the Payment Agent for our platform.

Your responsibilities:
- Process refund requests according to company refund policy
- Investigate payment disputes
- Clarify billing questions and invoice details
- Process approved payment tasks

Rules you MUST follow:
- Refunds up to £500: you can approve and process directly
- Refunds above £500: flag for human manager approval with full justification
- Duplicate charges: always refund immediately, no approval needed
- Suspected fraud: immediately escalate to Fraud Agent — do NOT process payment
- Always reference the specific transaction ID in your response
- Log all payment actions with timestamp and reason

Company refund policy is available in your knowledge base."""

    def check_requires_human_approval(self, request: str, response: str) -> tuple[bool, str]:
        """Large refunds and unusual payment patterns need human review"""
        # Extract amounts from request and response
        amounts = re.findall(r'£(\d+(?:,\d{3})*(?:\.\d{2})?)', request + response)
        for amount_str in amounts:
            amount = float(amount_str.replace(',', ''))
            if amount > self.high_value_threshold:
                return True, f"Refund of £{amount:,.2f} exceeds £{self.high_value_threshold:.0f} threshold — manager approval required"
        return False, ""


class FraudAgent(BaseAgent):
    """
    Investigates suspicious activity:
    - Unusual transaction patterns
    - Suspicious login attempts
    - KYC/AML checks
    - Sanctions screening
    """
    agent_type = "fraud"
    system_prompt = """You are the Fraud Investigation Agent for our platform.

Your responsibilities:
- Investigate reports of suspicious account activity
- Identify patterns consistent with fraud, money laundering, or identity theft
- Check accounts against KYC (Know Your Customer) requirements
- Flag potential AML (Anti-Money Laundering) concerns
- Perform sanctions and PEP (Politically Exposed Person) screening

Investigation framework:
1. Review the transaction pattern or suspicious activity reported
2. Check against company AML and KYC policies (in knowledge base)
3. Assign a risk level: LOW / MEDIUM / HIGH / CRITICAL
4. For MEDIUM and above: recommend specific actions
5. For HIGH and CRITICAL: ALWAYS escalate for human review

You do NOT have authority to freeze accounts or block transactions.
You gather evidence and make recommendations. Human team executes.

Your knowledge base contains AML procedures, KYC requirements, and compliance documentation."""

    def check_requires_human_approval(self, request: str, response: str) -> tuple[bool, str]:
        """Fraud findings always need human review before action"""
        risk_indicators = ["HIGH", "CRITICAL", "suspicious", "fraud", "money laundering",
                          "sanctions", "freeze", "block", "escalate"]
        response_lower = response.lower()
        for indicator in risk_indicators:
            if indicator.lower() in response_lower:
                return True, f"Fraud investigation flagged risk indicator: '{indicator}' — compliance team review required"
        return False, ""


class ComplianceAgent(BaseAgent):
    """
    Handles regulatory compliance checks:
    - KYC and AML verification
    - Sanctions screening
    - Regulatory reporting requirements
    - Data protection queries (GDPR, etc.)
    """
    agent_type = "compliance"
    system_prompt = """You are the Compliance Agent for our platform.

Your responsibilities:
- Verify customer compliance with KYC and AML requirements
- Assess regulatory obligations for specific transactions or business activities
- Advise on data protection requirements (GDPR, UK GDPR)
- Check against sanctions lists and PEP databases
- Advise on regulatory reporting thresholds

Important boundaries:
- You provide compliance INFORMATION and ASSESSMENT, not legal advice
- For definitive legal interpretations, recommend customers consult a qualified lawyer
- High-risk compliance findings MUST be reviewed by the compliance team (human)
- Always cite the specific regulation or company policy you are referencing

Your knowledge base contains regulatory guidelines, compliance procedures, and policy documents."""

    def check_requires_human_approval(self, request: str, response: str) -> tuple[bool, str]:
        """Compliance actions always require human sign-off"""
        compliance_actions = ["report to", "file with", "notify regulator", "submit to fca",
                              "suspicious activity report", "sar", "freeze", "restrict"]
        response_lower = response.lower()
        for action in compliance_actions:
            if action in response_lower:
                return True, f"Compliance action '{action}' requires regulatory compliance team approval"
        return False, ""


class CustomerSupportAgent(BaseAgent):
    """
    Handles general customer support requests:
    - Product questions
    - Technical troubleshooting
    - Feature guidance
    - Escalation routing
    """
    agent_type = "customer_support"
    system_prompt = """You are the Customer Support Agent for our platform.

Your responsibilities:
- Answer product and feature questions
- Guide customers through technical issues step by step
- Explain how to use platform features effectively
- Identify when issues need escalation to specialist agents

Routing rules:
- Billing and payment questions → route to Payment Agent
- Account access problems → route to Customer Agent
- Suspicious activity → route to Fraud Agent
- Compliance questions → route to Compliance Agent
- General product help → handle yourself

Response guidelines:
- Be friendly, clear, and patient
- Use numbered steps for technical instructions
- Always confirm the customer understood the solution
- If you cannot resolve in 3 exchanges, offer to escalate

Your knowledge base contains product documentation, FAQs, and troubleshooting guides."""

    def check_requires_human_approval(self, request: str, response: str) -> tuple[bool, str]:
        """Escalation requests need human agent involvement"""
        escalation_terms = ["speak to a human", "talk to someone", "escalate",
                           "manager", "supervisor", "not helpful", "frustrated"]
        request_lower = request.lower()
        for term in escalation_terms:
            if term in request_lower:
                return True, "Customer has requested human agent involvement"
        return False, ""

Step 5: The Custom Agent Router

The blue brain icon in the centre of your diagram. This is what receives every request and decides which specialist agent handles it.

python

# agents/router.py
# The Custom Agent Router — decides which agent handles each request

import boto3
import json
import logging
import os
from .specialist_agents import (
    CustomerAgent, PaymentAgent, FraudAgent,
    ComplianceAgent, CustomerSupportAgent
)
from .base_agent import AgentResponse, bedrock

logger = logging.getLogger(__name__)

MODEL_ID = "anthropic.claude-3-5-haiku-20241022"  # Fast model for routing decisions


class AgentRouter:
    """
    The Custom Agent Router coordinates the specialist agents.

    Responsibilities:
    1. Receive user requests from the API
    2. Classify the request and select the right agent
    3. Dispatch to the selected agent
    4. Receive the agent's response
    5. Pass the response to Response Verification
    6. Handle human approval workflows
    """

    def __init__(self):
        self.agents = {
            "customer": CustomerAgent(),
            "payment": PaymentAgent(),
            "fraud": FraudAgent(),
            "compliance": ComplianceAgent(),
            "customer_support": CustomerSupportAgent()
        }

        # Use a fast, cheap model for routing — no need for Opus here
        self.bedrock = bedrock

    def classify_request(self, request: str, user_id: str) -> str:
        """
        Classify which specialist agent should handle this request.
        Uses Claude Haiku for speed — routing is a simple classification task.
        """
        response = self.bedrock.invoke_model(
            modelId=MODEL_ID,
            body=json.dumps({
                "anthropic_version": "bedrock-2023-05-31",
                "max_tokens": 50,
                "temperature": 0.0,  # No randomness — routing must be deterministic
                "system": """Classify this customer request into exactly ONE category.
Return ONLY the category name, nothing else.

Categories:
- customer: Account info, profile updates, subscription changes, login help
- payment: Refunds, billing disputes, invoice questions, payment processing
- fraud: Suspicious activity, unusual transactions, security concerns, identity theft
- compliance: KYC, AML, regulatory questions, data protection, GDPR
- customer_support: General product help, technical issues, feature questions, anything else

If multiple categories apply, pick the PRIMARY one.""",
                "messages": [{"role": "user", "content": f"Request: {request}"}]
            })
        )

        result = json.loads(response['body'].read())
        agent_type = result["content"][0]["text"].strip().lower()

        # Validate the classification
        if agent_type not in self.agents:
            logger.warning(f"Unknown agent type '{agent_type}' — defaulting to customer_support")
            return "customer_support"

        logger.info(f"Request classified as: {agent_type}")
        return agent_type

    def route(self, request: str, session_id: str, user_id: str) -> AgentResponse:
        """
        Main routing method.
        Classifies the request and dispatches to the correct specialist agent.
        """
        logger.info(f"Router received request for session {session_id}")

        # Step 1: Classify the request
        agent_type = self.classify_request(request, user_id)

        # Step 2: Dispatch to the selected agent
        selected_agent = self.agents[agent_type]
        response = selected_agent.handle(request, session_id, user_id)

        return response

Step 6: Response Verification The Safety Net

The centre-right box in your diagram: "Response Verification Custom Code on Amazon EKS". This checks every response before it reaches the user.

python

# verification/response_verifier.py
# Custom verification code running on EKS

import re
import json
import logging
from dataclasses import dataclass
from typing import Optional

logger = logging.getLogger(__name__)


@dataclass
class VerificationResult:
    accepted: bool
    rejection_reason: Optional[str] = None
    modified_content: Optional[str] = None


class ResponseVerifier:
    """
    Checks agent responses before they reach the user.

    This is a second safety layer — after Bedrock Guardrails,
    this code performs custom business-logic checks that are
    specific to the application.

    From the architecture: 
    "Checks the proposed response before releasing it to the user"
    "Send rejection reason" back to agent if rejected
    "Agent-generated response accepted" if it passes
    """

    # Patterns that indicate a response should be blocked
    BLOCKED_PATTERNS = [
        r'\bpassword\b.{0,50}\bis\b',       # Never reveal passwords
        r'\b4[0-9]{12}(?:[0-9]{3})?\b',     # Credit card numbers
        r'\b[0-9]{3}-[0-9]{2}-[0-9]{4}\b',  # SSN format
        r'<script[^>]*>',                    # XSS attempts in output
        r'SELECT .+ FROM .+ WHERE',          # SQL in output (injection indicator)
        r'internal server error',            # Exposing system errors to users
        r'stack trace',                      # Exposing technical details
        r'(api_key|secret_key)\s*[:=]',     # Exposing credentials
    ]

    # Financial accuracy patterns — dollar amounts must be precise
    FINANCIAL_PATTERNS = [
        r'approximately\s+£',               # "Approximately £500" is not okay for financial responses
        r'around\s+£',
        r'roughly\s+£',
    ]

    # Patterns that require human review despite agent passing them
    ESCALATION_PATTERNS = [
        r'legal action',
        r'regulatory authority',
        r'media',
        r'going public',
        r'formal complaint',
    ]

    def verify(self, response_content: str, agent_type: str, original_request: str) -> VerificationResult:
        """
        Verify an agent response before sending it to the user.

        Returns:
        - VerificationResult with accepted=True if the response passes
        - VerificationResult with accepted=False and rejection_reason if blocked
        """

        # Check 1: Blocked patterns — hard stop
        for pattern in self.BLOCKED_PATTERNS:
            if re.search(pattern, response_content, re.IGNORECASE):
                logger.warning(f"Response blocked: pattern '{pattern}' detected in {agent_type} response")
                return VerificationResult(
                    accepted=False,
                    rejection_reason=f"Response contains blocked content pattern. Rewrite without including {pattern}"
                )

        # Check 2: Financial responses must be precise
        if agent_type == "payment":
            for pattern in self.FINANCIAL_PATTERNS:
                if re.search(pattern, response_content, re.IGNORECASE):
                    return VerificationResult(
                        accepted=False,
                        rejection_reason="Financial responses must use exact amounts, not approximations. Rewrite with precise figures."
                    )

        # Check 3: Minimum quality — response must be substantive
        if len(response_content.strip()) < 50:
            return VerificationResult(
                accepted=False,
                rejection_reason="Response is too short to be helpful. Provide a complete, substantive answer."
            )

        # Check 4: Response must be in correct language (English)
        # Simple heuristic — production should use a proper language detector
        english_word_ratio = sum(
            1 for word in response_content.split()
            if word.lower() in {"the", "a", "is", "are", "was", "were", "have", "has", "be", "to", "of", "and", "in"}
        ) / max(len(response_content.split()), 1)

        # Check 5: No response should tell the user the agent "can't" do something
        # without offering an alternative path
        if "cannot" in response_content.lower() or "can't" in response_content.lower():
            if not any(word in response_content.lower()
                      for word in ["instead", "however", "alternatively", "suggest", "recommend", "please"]):
                return VerificationResult(
                    accepted=False,
                    rejection_reason="Response declines without offering an alternative. Rewrite to include a constructive path forward."
                )

        # All checks passed
        logger.info(f"Response from {agent_type} accepted by verifier")
        return VerificationResult(accepted=True)

Step 7: Human Approval Workflow

The bottom-right of your diagram. When an agent flags that human review is needed, this system handles it.

python

# approval/human_approval.py
# Human approval workflow using SNS + Step Functions

import boto3
import json
import uuid
import logging
import os
from datetime import datetime, timezone

logger = logging.getLogger(__name__)

REGION = os.environ.get("AWS_REGION", "eu-west-1")
ACCOUNT_ID = os.environ.get("AWS_ACCOUNT_ID")

sns = boto3.client('sns', region_name=REGION)
dynamodb = boto3.resource('dynamodb', region_name=REGION)
approval_table = dynamodb.Table('human-approvals')


class HumanApprovalWorkflow:
    """
    Manages the human approval gate in the multi-agent system.

    From the architecture diagram:
    "Use human approval when required"
    "Set Functions → Human approval"
    "Return allow or block decision"
    """

    APPROVAL_TOPIC_ARN = f"arn:aws:sns:{REGION}:{ACCOUNT_ID}:human-approval-topic"

    def request_approval(
        self,
        session_id: str,
        user_id: str,
        agent_type: str,
        original_request: str,
        proposed_response: str,
        approval_reason: str
    ) -> str:
        """
        Submit a response for human approval.
        Returns an approval_id that can be used to check status.
        """
        approval_id = str(uuid.uuid4())

        # Store the approval request in DynamoDB
        approval_table.put_item(Item={
            "approvalId": approval_id,
            "sessionId": session_id,
            "userId": user_id,
            "agentType": agent_type,
            "originalRequest": original_request,
            "proposedResponse": proposed_response,
            "approvalReason": approval_reason,
            "status": "PENDING",
            "createdAt": datetime.now(timezone.utc).isoformat(),
            "ttl": int(datetime.now(timezone.utc).timestamp()) + 3600  # 1 hour TTL
        })

        # Notify the human approval team via SNS
        sns.publish(
            TopicArn=self.APPROVAL_TOPIC_ARN,
            Subject=f"APPROVAL REQUIRED: {agent_type.upper()} Agent — {approval_reason[:50]}",
            Message=json.dumps({
                "approval_id": approval_id,
                "approval_url": f"https://admin.yourcompany.com/approvals/{approval_id}",
                "session_id": session_id,
                "user_id": user_id,
                "agent_type": agent_type,
                "reason": approval_reason,
                "original_request": original_request[:500],
                "proposed_response": proposed_response[:500],
                "created_at": datetime.now(timezone.utc).isoformat()
            }, indent=2),
            MessageAttributes={
                "agent_type": {
                    "DataType": "String",
                    "StringValue": agent_type
                },
                "priority": {
                    "DataType": "String",
                    "StringValue": "HIGH" if agent_type in ["fraud", "compliance"] else "MEDIUM"
                }
            }
        )

        logger.info(f"Human approval requested: {approval_id} for session {session_id}")
        return approval_id

    def check_approval_status(self, approval_id: str) -> dict:
        """Check the status of a pending approval"""
        response = approval_table.get_item(Key={"approvalId": approval_id})
        return response.get("Item", {"status": "NOT_FOUND"})

    def process_approval_decision(self, approval_id: str, decision: str, approver_id: str, notes: str = "") -> bool:
        """
        Record a human approval decision.
        Called by the approval UI when a human makes a decision.

        decision: "APPROVED" or "REJECTED"
        """
        if decision not in ["APPROVED", "REJECTED"]:
            raise ValueError(f"Invalid decision: {decision}. Must be APPROVED or REJECTED")

        approval_table.update_item(
            Key={"approvalId": approval_id},
            UpdateExpression="SET #s = :status, approverId = :approver, notes = :notes, decidedAt = :decided",
            ExpressionAttributeNames={"#s": "status"},
            ExpressionAttributeValues={
                ":status": decision,
                ":approver": approver_id,
                ":notes": notes,
                ":decided": datetime.now(timezone.utc).isoformat()
            }
        )

        logger.info(f"Approval {approval_id} {decision} by {approver_id}")
        return decision == "APPROVED"

Step 8: The Main API Handler Bringing It All Together

python

# main.py
# FastAPI application running on EKS — the entry point from API Gateway

from fastapi import FastAPI, HTTPException, Header, Depends
from pydantic import BaseModel
from typing import Optional
import logging
import uuid
import time
import os

from agents.router import AgentRouter
from verification.response_verifier import ResponseVerifier
from approval.human_approval import HumanApprovalWorkflow

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

app = FastAPI(
    title="Multi-Agent AI System",
    description="Production multi-agent AI on Amazon EKS",
    version="1.0.0"
)

# Initialise core components
router = AgentRouter()
verifier = ResponseVerifier()
approval_workflow = HumanApprovalWorkflow()

MAX_VERIFICATION_RETRIES = 3


class AgentRequest(BaseModel):
    message: str
    user_id: str
    session_id: Optional[str] = None


class AgentApiResponse(BaseModel):
    session_id: str
    agent_type: str
    response: str
    requires_human_approval: bool
    approval_id: Optional[str] = None
    knowledge_sources: list = []
    processing_time_ms: int


@app.get("/health")
async def health_check():
    """Health check endpoint for ALB target group"""
    return {"status": "healthy", "service": "multi-agent-system"}


@app.post("/api/v1/agent", response_model=AgentApiResponse)
async def handle_agent_request(
    request: AgentRequest,
    authorization: str = Header(...)
):
    """
    Main endpoint — receives requests from API Gateway.

    Flow:
    1. Validate authentication token
    2. Route to correct specialist agent
    3. Verify the response
    4. Handle human approval if required
    5. Return response to user
    """
    start_time = time.time()

    # Generate session ID if not provided
    session_id = request.session_id or str(uuid.uuid4())

    logger.info(f"Request received: session={session_id}, user={request.user_id}")

    # Step 1: Route to specialist agent
    agent_response = router.route(
        request=request.message,
        session_id=session_id,
        user_id=request.user_id
    )

    # Step 2: Verify the response (up to 3 attempts)
    final_content = agent_response.content
    verification_passed = False

    for attempt in range(MAX_VERIFICATION_RETRIES):
        verification = verifier.verify(
            response_content=final_content,
            agent_type=agent_response.agent_type,
            original_request=request.message
        )

        if verification.accepted:
            verification_passed = True
            break
        else:
            logger.warning(
                f"Verification failed (attempt {attempt + 1}): {verification.rejection_reason}"
            )

            if attempt < MAX_VERIFICATION_RETRIES - 1:
                # Ask the agent to try again with the rejection reason
                correction_response = router.agents[agent_response.agent_type].call_bedrock_with_guardrails(
                    messages=[
                        {"role": "user", "content": request.message},
                        {"role": "assistant", "content": final_content},
                        {"role": "user", "content": f"Your previous response was rejected: {verification.rejection_reason}. Please provide a corrected response."}
                    ]
                )
                final_content = correction_response[0]

    if not verification_passed:
        raise HTTPException(
            status_code=500,
            detail="Unable to generate an acceptable response. Please try again or contact support."
        )

    # Step 3: Handle human approval if required
    approval_id = None
    if agent_response.requires_human_approval:
        approval_id = approval_workflow.request_approval(
            session_id=session_id,
            user_id=request.user_id,
            agent_type=agent_response.agent_type,
            original_request=request.message,
            proposed_response=final_content,
            approval_reason=agent_response.approval_reason
        )

        # Return pending message to user while approval is in progress
        final_content = (
            f"Your request has been received and is pending review by our team. "
            f"Reference number: {approval_id}. "
            f"We will respond within 1 business hour. "
            f"\n\nPreliminary assessment: {final_content}"
        )

    processing_time_ms = int((time.time() - start_time) * 1000)

    logger.info(
        f"Request completed: session={session_id}, agent={agent_response.agent_type}, "
        f"time={processing_time_ms}ms, approval_required={agent_response.requires_human_approval}"
    )

    return AgentApiResponse(
        session_id=session_id,
        agent_type=agent_response.agent_type,
        response=final_content,
        requires_human_approval=agent_response.requires_human_approval,
        approval_id=approval_id,
        knowledge_sources=agent_response.knowledge_sources or [],
        processing_time_ms=processing_time_ms
    )

Step 9: Deploy to EKS — Kubernetes Manifests

yaml

# k8s/namespace.yaml
apiVersion: v1
kind: Namespace
metadata:
  name: multi-agent
  labels:
    app: multi-agent-system

yaml

# k8s/service-account.yaml
# IRSA — IAM Roles for Service Accounts
# Links the Kubernetes service account to the IAM role we created in Step 1
apiVersion: v1
kind: ServiceAccount
metadata:
  name: agent-service-account
  namespace: multi-agent
  annotations:
    # Replace with your actual role ARN
    eks.amazonaws.com/role-arn: arn:aws:iam::YOUR_ACCOUNT_ID:role/multi-agent-pod-role

yaml

# k8s/configmap.yaml
apiVersion: v1
kind: ConfigMap
metadata:
  name: agent-config
  namespace: multi-agent
data:
  AWS_REGION: "eu-west-1"
  BEDROCK_REGION: "us-east-1"
  LOG_LEVEL: "INFO"
  MAX_VERIFICATION_RETRIES: "3"

yaml

# k8s/secrets.yaml
# In production: use AWS Secrets Manager with External Secrets Operator
# For initial setup:
apiVersion: v1
kind: Secret
metadata:
  name: agent-secrets
  namespace: multi-agent
type: Opaque
stringData:
  GUARDRAIL_ID: "YOUR_GUARDRAIL_ID"
  GUARDRAIL_VERSION: "1"
  KNOWLEDGE_BASE_ID: "YOUR_KB_ID"

yaml

# k8s/deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: multi-agent-system
  namespace: multi-agent
  labels:
    app: multi-agent
    version: v1
spec:
  replicas: 3                    # 3 replicas across 2 AZs
  selector:
    matchLabels:
      app: multi-agent
  template:
    metadata:
      labels:
        app: multi-agent
        version: v1
    spec:
      serviceAccountName: agent-service-account  # IRSA for AWS API access
      containers:
        - name: multi-agent
          image: YOUR_ECR_REGISTRY/multi-agent:latest
          ports:
            - containerPort: 8000
          envFrom:
            - configMapRef:
                name: agent-config
            - secretRef:
                name: agent-secrets
          resources:
            requests:
              cpu: "500m"
              memory: "1Gi"
            limits:
              cpu: "2000m"
              memory: "4Gi"
          readinessProbe:
            httpGet:
              path: /health
              port: 8000
            initialDelaySeconds: 10
            periodSeconds: 5
          livenessProbe:
            httpGet:
              path: /health
              port: 8000
            initialDelaySeconds: 30
            periodSeconds: 10
          securityContext:
            runAsNonRoot: true      # Never run as root
            runAsUser: 1000
            readOnlyRootFilesystem: true
      affinity:
        podAntiAffinity:
          preferredDuringSchedulingIgnoredDuringExecution:
            - weight: 100
              podAffinityTerm:
                labelSelector:
                  matchLabels:
                    app: multi-agent
                topologyKey: kubernetes.io/hostname  # Spread across nodes

yaml

# k8s/service.yaml
apiVersion: v1
kind: Service
metadata:
  name: multi-agent-service
  namespace: multi-agent
  annotations:
    # Internal ALB — NOT internet-facing
    service.beta.kubernetes.io/aws-load-balancer-type: "external"
    service.beta.kubernetes.io/aws-load-balancer-scheme: "internal"
    service.beta.kubernetes.io/aws-load-balancer-name: "multi-agent-internal-alb"
spec:
  selector:
    app: multi-agent
  ports:
    - port: 80
      targetPort: 8000
  type: LoadBalancer

---
# k8s/hpa.yaml — Horizontal Pod Autoscaler
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
  name: multi-agent-hpa
  namespace: multi-agent
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: multi-agent-system
  minReplicas: 3
  maxReplicas: 20
  metrics:
    - type: Resource
      resource:
        name: cpu
        target:
          type: Utilization
          averageUtilization: 70
    - type: Resource
      resource:
        name: memory
        target:
          type: Utilization
          averageUtilization: 80

Deploy everything:

bash

# Apply all manifests
kubectl apply -f k8s/namespace.yaml
kubectl apply -f k8s/service-account.yaml
kubectl apply -f k8s/configmap.yaml
kubectl apply -f k8s/secrets.yaml
kubectl apply -f k8s/deployment.yaml
kubectl apply -f k8s/service.yaml
kubectl apply -f k8s/hpa.yaml

# Verify deployment
kubectl get pods -n multi-agent
kubectl get svc -n multi-agent

# Check logs
kubectl logs -n multi-agent -l app=multi-agent --tail=50

Step 10: Front-End Security Layer WAF, CloudFront, API Gateway

This is the top layer of your diagram, protecting the system from harmful traffic before it reaches EKS.

bash

# Create WAF WebACL
WAF_ACL_ID=$(aws wafv2 create-web-acl \
  --name "multi-agent-waf" \
  --scope CLOUDFRONT \
  --default-action Allow={} \
  --rules '[
    {
      "Name": "AWSManagedRulesCommonRuleSet",
      "Priority": 1,
      "OverrideAction": {"None": {}},
      "Statement": {
        "ManagedRuleGroupStatement": {
          "VendorName": "AWS",
          "Name": "AWSManagedRulesCommonRuleSet"
        }
      },
      "VisibilityConfig": {
        "SampledRequestsEnabled": true,
        "CloudWatchMetricsEnabled": true,
        "MetricName": "CommonRuleSet"
      }
    },
    {
      "Name": "AWSManagedRulesKnownBadInputsRuleSet",
      "Priority": 2,
      "OverrideAction": {"None": {}},
      "Statement": {
        "ManagedRuleGroupStatement": {
          "VendorName": "AWS",
          "Name": "AWSManagedRulesKnownBadInputsRuleSet"
        }
      },
      "VisibilityConfig": {
        "SampledRequestsEnabled": true,
        "CloudWatchMetricsEnabled": true,
        "MetricName": "KnownBadInputs"
      }
    },
    {
      "Name": "RateLimitRule",
      "Priority": 3,
      "Action": {"Block": {}},
      "Statement": {
        "RateBasedStatement": {
          "Limit": 100,
          "AggregateKeyType": "IP"
        }
      },
      "VisibilityConfig": {
        "SampledRequestsEnabled": true,
        "CloudWatchMetricsEnabled": true,
        "MetricName": "RateLimit"
      }
    }
  ]' \
  --visibility-config SampledRequestsEnabled=true,CloudWatchMetricsEnabled=true,MetricName=multi-agent-waf \
  --region us-east-1 \
  --query 'Summary.Id' \
  --output text)

echo "WAF ACL ID: $WAF_ACL_ID"

Step 11: Testing the Complete System

bash

# Get the API Gateway endpoint
API_ENDPOINT="https://your-api-gateway-id.execute-api.eu-west-1.amazonaws.com/prod"

# Test 1: Customer Support Query (should route to customer_support agent)
echo "=== Test 1: Customer Support ==="
curl -s -X POST "$API_ENDPOINT/api/v1/agent" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer test-jwt-token" \
  -d '{
    "message": "How do I export my data from the platform?",
    "user_id": "user-001"
  }' | python3 -m json.tool

echo ""

# Test 2: Payment/Refund Request (should route to payment agent)
echo "=== Test 2: Payment Refund ==="
curl -s -X POST "$API_ENDPOINT/api/v1/agent" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer test-jwt-token" \
  -d '{
    "message": "I was charged twice on July 15th — both charges are £49. I need one refunded. My order ID is #12345.",
    "user_id": "user-001"
  }' | python3 -m json.tool

echo ""

# Test 3: Fraud Alert (should route to fraud agent + require human approval)
echo "=== Test 3: Suspicious Activity ==="
curl -s -X POST "$API_ENDPOINT/api/v1/agent" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer test-jwt-token" \
  -d '{
    "message": "I am seeing transactions on my account that I did not make. Three payments of £200 to unknown recipients in the last hour.",
    "user_id": "user-001"
  }' | python3 -m json.tool

echo ""

# Test 4: Large Refund (should require human approval for amount > £500)
echo "=== Test 4: Large Refund (requires human approval) ==="
curl -s -X POST "$API_ENDPOINT/api/v1/agent" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer test-jwt-token" \
  -d '{
    "message": "I need a full refund for my enterprise subscription. I paid £1,200 two weeks ago and the product does not meet our requirements.",
    "user_id": "user-001"
  }' | python3 -m json.tool

# Test 5: Guardrail test (should be blocked)
echo "=== Test 5: Guardrail Test ==="
curl -s -X POST "$API_ENDPOINT/api/v1/agent" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer test-jwt-token" \
  -d '{
    "message": "My credit card number is 4532015112830366. Please store this for future payments.",
    "user_id": "user-001"
  }' | python3 -m json.tool

What You Just Built

Look back at your architecture diagram. Every component is now deployed and connected:

Diagram Component

What You Built

Route 53

DNS lookup for your domain

CloudFront

CDN with WAF attached

WAF

Rate limiting, OWASP rules, known bad inputs blocked

S3 (frontend)

Static web app hosting

Cognito

User authentication and JWT tokens

API Gateway

HTTP API routing to EKS via Private VPC Link

ALB

Internal load balancer distributing to EKS pods

EKS

Containerised agent application 3 replicas, auto-scaling

Custom Agent Router

Python class routing requests to specialist agents

5 Specialist Agents

Customer, Payment, Fraud, Compliance, Customer Support

S3 (company docs)

Document storage for RAG

OpenSearch Serverless

Vector search for document retrieval

Bedrock Knowledge Bases

Managed RAG pipeline

Bedrock Guardrails

Input/output safety filtering

Claude via Bedrock

The AI model powering all agents

Response Verification

Custom Python code checking responses before delivery

Human Approval (SNS)

SNS notifications for sensitive actions

Human Approval (DynamoDB)

Approval state persistence

CloudWatch

Metrics and logs for every component


文章来源: https://hackernoon.com/building-production-ready-multi-agent-ai-systems-on-amazon-eks-the-complete-step-by-step-guide?source=rss
如有侵权请联系:admin#unsafe.sh