
Scaling 'Build for a Friend' Tools into Production Systems: A DevOps & AI Guide
Learn how to transform Hacktoberfest-style personal projects into resilient production systems using modern DevOps pipelines and AI-driven automation agents.
Intro
The culture of open-source contribution is often defined by events like Hacktoberfest, where thousands of developers contribute to public repositories in a single month. While these efforts are invaluable for community building, the resulting codebases frequently share a common DNA: they are designed for a single user (the developer), lack robust error handling, ignore security concerns, and rely on manual deployment. These "Build for a Friend" tools are excellent prototypes but disastrous production systems.
Bridging the gap between a personal script and a scalable production system requires a fundamental shift in engineering mindset. In this deep dive, we explore the architectural and operational patterns necessary to refine these projects. We focus on two modern pillars: Resilient DevOps Pipelines and AI Agents that act as automated operational co-pilots.
1. Architectural Refactoring: From Monolith to Service
Before optimizing the deployment pipeline, the code itself must be ready for production. Personal projects often conflate data storage, business logic, and presentation layers into a single entry point.
Separation of Concerns
Production systems require clear boundaries. A common pattern for evolving a simple API is to decouple the entry point from the logic. Instead of a single main.py handling requests, database connections, and business rules, we adopt a modular structure:
# Old: Personal Project Style
# def handle_request(req):
# # Fetch DB
# # Business Logic
# # Return Response
# New: Production Style
from .services.user_service import User
from .db.database import get_connection
class UserController:
def get_user(self, user_id):
# Validation
if not user_id: raise ValueError("Missing ID")
# Business Logic (Delegated)
service = User()
return service.fetch(user_id)
Why this matters for production:
- Testability: By separating
UserControllerfromUserservice, you can write unit tests for the controller without mocking the database. - Scalability: If the
Userservice becomes a bottleneck, you can offload it to a queue or a separate microservice without touching the controller. - Statelessness: The controller holds no state, making it easy to scale horizontally behind a load balancer.
Configuration Management
Personal projects hardcode database URLs and API keys. Production systems abstract these into environment-specific configurations. We recommend using a library like python-dotenv or pydantic-settings to manage secrets and config.
from pydantic_settings import BaseSettings
class Settings(BaseSettings):
database_url: str
redis_url: str
is_production: bool = False
class Config:
env_file = ".env"
2. The DevOps Foundation: CI/CD Pipelines
A personal project is typically deployed by running git push and manually running docker build. A production system requires a Continuous Integration and Continuous Deployment (CI/CD) pipeline that guarantees code quality before any artifact is built.
Defining the Pipeline Stages
A robust pipeline includes four critical stages:
- Static Analysis & Linting: Runs before unit tests to catch style and potential security issues early. (e.g.,
eslint,flake8,semgrep). - Unit & Integration Tests: Ensures that the new code does not break existing functionality. Coverage thresholds are enforced.
- Security Scanning: Scans dependencies for known vulnerabilities and container images for misconfigurations. (e.g.,
trivy,snyk). - Build & Push: Compiles the code into a container image and pushes it to a registry (E.g., Docker Hub, ECR).
Example: GitHub Actions Pipeline
Below is a workflow that transforms a simple project into a pipeline-guarded system. Notice the integration of security scanning, which is rarely present in personal projects.
name: CI/CD Pipeline
on:
push:
branches: [ main ]
pull_request:
branches: [ main ]
jobs:
build-and-test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Set up Node
uses: actions/setup-node@v4
with:
node-version: '18'
- name: Install dependencies
run: npm ci
- name: Run Linters
run: npm run lint
- name: Run Tests
run: npm run test
- name: Security Scan
run: npm audit --audit-level=high
- name: Build Docker Image
run: |
docker build -t my-service:${{ github.sha }} .
docker tag my-service:${{ github.sha }} my-service:latest
docker push my-service:${{ github.sha }}
Production Best Practices:
- Immutability: Deploy a specific image tag (e.g.,
abc123), notlatest. This allows for instant rollback. - Parallelization: Run linters, tests, and security scans in parallel jobs to reduce pipeline duration.
- Artifact Storage: Store the Docker image in a secure registry, not just on the runner.
3. Infrastructure as Code: Beyond docker run
Manual infrastructure changes are a source of "snowflake servers." For production, your infrastructure must be defined as code. Terraform or Pulumi are the gold standards, but for smaller projects, container orchestrators (Kubernetes or Docker Swarm) provide the necessary abstraction.
Containerization Best Practices
Personal projects often use base images like python:3 or node:latest. Production images must be optimized for size and security.
# Production-optimized Dockerfile
FROM node:18-alpine AS base
# Use non-root user
USER node
WORKDIR /app
COPY package*.json ./
# Cache dependencies
RUN npm ci --omit=dev
COPY . .
# Health Check
HEALTHCHECK --interval=30s --timeout=3s --start-period=5s --retries=3 \
CMD node scripts/healthcheck.js
EXPOSE 3000
CMD ["node", "server.js"]
Key Differences:
--omit=dev: Removes development dependencies (linter, typescript, etc.), reducing image size by up to 40%.USER node: Running as non-root prevents container escape vulnerabilities.HEALTHCHECK: Instructs the orchestrator on how to check service liveness, enabling automatic restarts on failure.
4. Observability: Moving from "It Worked on My Machine" to "It is Working Now"
Production systems are distributed and fail in ways developers don't anticipate. You must implement the Three Pillars of Observability: Logs, Metrics, and Traces.
Structured Logging
Personal projects use console.log("User logged in"). Production systems use structured JSON logs. This allows log aggregators (Elasticsearch, CloudWatch, Datadog) to parse and query logs efficiently.
// Bad
console.log(`User ${user.id} failed to login`);
// Good
logger.info('login.failed', {
userId: user.id,
reason: 'invalid_password',
timestamp: new Date().toISOString()
});
Distributed Tracing
In a microservice architecture, a request might pass through 5 services. A tracer (like OpenTelemetry) attaches a unique traceId to each request, allowing you to visualize the entire journey and identify which service is slow.
5. The AI Agent Revolution: Automated Operations
Traditional DevOps requires a human to monitor dashboards and intervene when metrics spike. AI Agents are changing this paradigm. An AI Agent is an LLM-based system that has access to tools (APIs, CLI, database) and an objective (e.g., "keep the error rate below 0.1%").
Defining the Agent Role
Instead of just receiving an alert, the Agent receives the context, investigates the root cause, and proposes or executes a remediation.
The Agent Loop:
- Perception: The Agent ingests an alert (e.g., "500 errors on /api/users").
- Reasoning: The LLM analyzes the alert, accesses recent logs and code diffs.
- Action: The Agent calls a tool (e.g.,
get_logs,rollback_version,scale_replicas). - Verification: The Agent monitors the metrics to ensure the error rate drops.
Architecture of an SRE Agent
Here is a conceptual architecture for an SRE Agent built using LangChain or similar frameworks.
// Pseudo-code for an SRE Agent using ReAct pattern
import { Agent } from 'agent-framework';
import { LLM } from 'llm-provider';
import { Tool } from 'tool-definitions';
const sreAgent = new Agent({
llm: new LLM({ model: 'gpt-4o', temperature: 0.0 }),
tools: [
Tool.logs, // Query Elasticsearch/Datadog
Tool.k8s, // Kubectl operations
Tool.metrics, // Query Prometheus
Tool.deploy, // Trigger CI/CD rollback
],
objective: "Resolve any alerts for the `checkout-service`.",
maxSteps: 5
});
Security & Permission Boundaries
Giving an AI agent access to production systems is risky. The agent must be bound by Principle of Least Privilege.
- Read-Only by Default: The agent should only have read access to logs and metrics initially.
- Dry-Run Mode: When the agent proposes an action like
k8s.scale, the system should require a human-in-the-loop approval for write operations, or enforce a
strict policy engine that validates every mutation against a predefined schema before execution.
class PolicyEngine:
"""
Validates agent-proposed actions against a policy schema.
All write operations must pass here before reaching the cluster.
"""
def __init__(self, policy_file: str):
with open(policy_file) as f:
self.policies = json.load(f)
def validate(self, action: dict) -> PolicyResult:
"""
Returns PolicyResult(allowed=True/False, reason=str)
"""
action_type = action.get("type", "")
resource = action.get("resource", "")
# Check if action type is in the allow-list
if action_type not in self.policies.get("allowed_actions", []):
return PolicyResult(
allowed=False,
reason=f"Action type '{action_type}' is not in the allow-list"
)
# Check resource-level restrictions
resource_policy = self.policies.get("resource_policies", {}).get(resource)
if resource_policy:
if action_type == "scale":
max_replicas = resource_policy.get("max_replicas", 50)
requested = action.get("params", {}).get("replicas", 0)
if requested > max_replicas:
return PolicyResult(
allowed=False,
reason=f"Requested replicas ({requested}) exceeds max ({max_replicas})"
)
return PolicyResult(allowed=True, reason="All checks passed")
3. The Observability Layer: Making AI Decisions Auditable
One of the most common failures when scaling AI-driven DevOps tools is the black-box problem: the system makes a decision, something goes wrong, and nobody can explain why the decision was made. The observability layer must address three dimensions:
- Decision Tracing – Every prompt sent to the LLM, every tool call proposed, and every policy evaluation must be logged with full context.
- Metric Correlation – The system must link the agent's actions to the resulting changes in infrastructure metrics (latency, error rates, resource utilization).
- Anomaly Detection – The system should detect when its own actions produce unexpected outcomes and escalate.
Implementing Decision Tracing
import uuid
from datetime import datetime
from dataclasses import dataclass, field
from typing import Any
import json
import logging
logger = logging.getLogger("agent.trace")
@dataclass
class DecisionTrace:
"""Immutable record of a single agent decision cycle."""
trace_id: str = field(default_factory=lambda: str(uuid.uuid4()))
timestamp: str = field(default_factory=lambda: datetime.utcnow().isoformat())
user_query: str = ""
llm_prompt: str = ""
llm_response: str = ""
proposed_actions: list = field(default_factory=list)
policy_results: list = field(default_factory=list)
executed_actions: list = field(default_factory=list)
metadata: dict = field(default_factory=dict)
def to_json(self) -> str:
return json.dumps(self.__dict__, indent=2, default=str)
class TracingMiddleware:
"""
Wraps the agent's execution loop to capture every decision step.
Integrates with OpenTelemetry for distributed tracing.
"""
def __init__(self, trace_store: TraceStore):
self.trace_store = trace_store
async def trace_cycle(
self,
user_query: str,
llm_client: LLMClient,
tool_registry: ToolRegistry,
policy_engine: PolicyEngine
) -> DecisionTrace:
trace = DecisionTrace(user_query=user_query)
# Step 1: Build and send prompt
prompt = self._build_prompt(user_query)
trace.llm_prompt = prompt
# Step 2: Get LLM response
response = await llm_client.complete(prompt)
trace.llm_response = response.content
# Step 3: Parse proposed actions
proposed = self._parse_actions(response.content)
trace.proposed_actions = proposed
# Step 4: Validate against policy
for action in proposed:
result = policy_engine.validate(action)
trace.policy_results.append({
"action": action,
"result": {"allowed": result.allowed, "reason": result.reason}
})
# Step 5: Execute allowed actions
for action in proposed:
policy_result = next(
r for r in trace.policy_results if r["action"] == action
)
if policy_result["result"]["allowed"]:
execution_result = await tool_registry.execute(action)
trace.executed_actions.append({
"action": action,
"result": execution_result
})
# Step 6: Persist trace
await self.trace_store.save(trace)
logger.info(f"Trace {trace.trace_id} completed: "
f"{len(trace.executed_actions)} actions executed")
return trace
Integrating with OpenTelemetry
For production deployments, the tracing system should integrate with OpenTelemetry to enable distributed tracing across the entire request lifecycle:
from opentelemetry import trace
from opentelemetry.trace import SpanKind
tracer = trace.get_tracer("build_for_a_friend.agent")
class OpenTelemetryTracer:
"""
Bridges the agent's decision trace with OpenTelemetry spans.
"""
def __init__(self, endpoint: str = "http://localhost:4318/v1/traces"):
self.endpoint = endpoint
self._setup_exporter()
def _setup_exporter(self):
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.sdk.trace import TracerProvider
provider = TracerProvider()
exporter = OTLPSpanExporter(endpoint=self.endpoint)
provider.add_span_processor(BatchSpanProcessor(exporter))
trace.set_tracer_provider(provider)
def create_agent_span(self, trace: DecisionTrace):
with tracer.start_as_current_span(
"agent.decision_cycle",
kind=SpanKind.INTERNAL,
attributes={
"agent.trace_id": trace.trace_id,
"agent.num_actions_proposed": len(trace.proposed_actions),
"agent.num_actions_executed": len(trace.executed_actions),
"agent.query": trace.user_query[:200], # Truncate for storage
}
) as span:
return span
4. Scaling the Tool Registry: From Monolith to Plugin Architecture
When you first build the tool, you probably wrote everything in a single file. Scaling requires a plugin architecture where each tool is independently deployable, testable, and versionable.
Tool Registry with Dynamic Loading
import importlib
import pkgutil
from abc import ABC, abstractmethod
from dataclasses import dataclass
from typing import Type, Dict, Optional
import asyncio
@dataclass
class ToolSchema:
"""Schema describing a tool's capabilities for the LLM."""
name: str
description: str
parameters: Dict[str, dict] # JSON Schema format
permissions: str # "read" or "write"
version: str
class BaseTool(ABC):
"""Abstract base class for all agent tools."""
@abstractmethod
def schema(self) -> ToolSchema:
"""Return the tool's schema for LLM consumption."""
...
@abstractmethod
async def execute(self, params: dict) -> ToolResult:
"""Execute the tool with given parameters."""
...
@abstractmethod
def health_check(self) -> bool:
"""Verify the tool's dependencies are available."""
...
class ToolRegistry:
"""
Dynamic tool registry with hot-loading, versioning, and health monitoring.
"""
def __init__(self, config_path: str):
self.config = self._load_config(config_path)
self._tools: Dict[str, BaseTool] = {}
self._tool_versions: Dict[str, str] = {}
self._health_status: Dict[str, bool] = {}
def _load_config(self, path: str) -> dict:
import yaml
with open(path) as f:
return yaml.safe_load(f)
async def discover_tools(self, plugin_path: str):
"""
Scan a directory for tool plugins and register them.
Each plugin must expose a `create_tool()` function.
"""
tools_found = 0
for module_info in pkgutil.iter_modules([plugin_path]):
if module_info.name.startswith("_"):
continue
try:
module = importlib.import_module(f"{plugin_path}.{module_info.name}")
if hasattr(module, "create_tool"):
tool = module.create_tool()
if isinstance(tool, BaseTool):
self.register(tool)
tools_found += 1
logger.info(f"Registered tool: {tool.schema().name}")
except Exception as e:
logger.error(f"Failed to load tool plugin "
f"'{module_info.name}': {e}")
logger.info(f"Tool discovery complete: {tools_found} tools registered")
def register(self, tool: BaseTool):
"""Register a tool and run initial health check."""
schema = tool.schema()
self._tools[schema.name] = tool
self._tool_versions[schema.name] = schema.version
self._health_status[schema.name] = tool.health_check()
if not self._health_status[schema.name]:
logger.warning(f"Tool '{schema.name}' registered but health check failed")
def get_tools_for_llm(self) -> list:
"""
Return tool schemas formatted for the LLM's function-calling interface.
Only includes healthy tools.
"""
return [
tool.schema().__dict__
for name, tool in self._tools.items()
if self._health_status.get(name, False)
]
async def execute(self, action: dict) -> ToolResult:
"""
Execute a tool action with timeout and error handling.
"""
tool_name = action["tool"]
params = action.get("params", {})
if tool_name not in self._tools:
return ToolResult(success=False, error=f"Tool '{tool_name}' not found")
if not self._health_status.get(tool_name, False):
return ToolResult(
success=False,
error=f"Tool '{tool_name}' is unhealthy"
)
tool = self._tools[tool_name]
timeout = self.config.get("execution_timeout_seconds", 30)
try:
result = await asyncio.wait_for(
tool.execute(params),
timeout=timeout
)
return result
except asyncio.TimeoutError:
return ToolResult(
success=False,
error=f"Tool '{tool_name}' timed out after {timeout}s"
)
except Exception as e:
return ToolResult(
success=False,
error=f"Tool '{tool_name}' execution error: {str(e)}"
)
async def run_health_checks(self):
"""Periodic health check for all registered tools."""
for name, tool in self._tools.items():
try:
self._health_status[name] = tool.health_check()
except Exception:
self._health_status[name] = False
unhealthy = [n for n, h in self._health_status.items() if not h]
if unhealthy:
logger.warning(f"Unhealthy tools: {unhealthy}")
Example Tool Implementation: Kubernetes Scale
from kubernetes import client, config
from .base import BaseTool, ToolResult, ToolSchema
class KubernetesScaleTool(BaseTool):
"""Scales a Kubernetes deployment."""
def __init__(self, kube_config_path: Optional[str] = None):
self._client = None
self._config_path = kube_config_path
self._initialized = False
def _ensure_client(self):
if not self._initialized:
if self._config_path:
config.load_kube_config(config_file=self._config_path)
else:
config.load_incluster_config()
self._client = client.AppsV1Api()
self._initialized = True
def schema(self) -> ToolSchema:
return ToolSchema(
name="k8s.scale",
description="Scale a Kubernetes deployment to a specified number of replicas",
parameters={
"namespace": {"type": "string", "description": "Kubernetes namespace"},
"deployment": {"type": "string", "description": "Deployment name"},
"replicas": {"type": "integer", "description": "Target replica count",
"minimum": 0, "maximum": 50},
"reason": {"type": "string", "description": "Reason for scaling (for audit trail)"}
},
permissions="write",
version="1.2.0"
)
async def execute(self, params: dict) -> ToolResult:
try:
self._ensure_client()
namespace = params["namespace"]
deployment = params["deployment"]
replicas = params["replicas"]
reason = params.get("reason", "Agent-initiated scale")
# Get current state
current = self._client.read_namespaced_deployment(
name=deployment, namespace=namespace
)
current_replicas = current.spec.replicas
# Scale
self._client.patch_namespaced_deployment_scale(
name=deployment,
namespace=namespace,
body={"spec": {"replicas": replicas}}
)
return ToolResult(
success=True,
data={
"previous_replicas": current_replicas,
"new_replicas": replicas,
"namespace": namespace,
"deployment": deployment,
"reason": reason
}
)
except client.ApiException as e:
return ToolResult(success=False, error=f"Kubernetes API error: {e.reason}")
except Exception as e:
return ToolResult(success=False, error=str(e))
def health_check(self) -> bool:
try:
self._ensure_client()
self._client.list_namespaced_deployment(namespace="default")
return True
except Exception:
return False
def create_tool() -> BaseTool:
"""Factory function for plugin discovery."""
return KubernetesScaleTool()
5. Human-in-the-Loop Approval Workflows
Not every action should be automatically approved. The system needs a graduated trust model:
| Risk Level | Examples | Approval Required |
|---|---|---|
| Low | Read logs, fetch metrics, list resources | No |
| Medium | Scale within bounds, restart pod | Auto-approve with notification |
| High | Scale beyond bounds, modify config | Human approval required |
| Critical | Delete resources, modify IAM, change network rules | Dual approval required |
Approval Workflow Implementation
from enum import Enum
from typing import Callable, Optional
import asyncio
import webhooks # e.g., slack_sdk or custom webhook client
class RiskLevel(Enum):
LOW = "low"
MEDIUM = "medium"
HIGH = "high"
CRITICAL = "critical"
class ApprovalWorkflow:
"""
Manages human-in-the-loop approval for agent actions.
Supports Slack, PagerDuty, and webhook-based notifications.
"""
def __init__(self, config: dict):
self.config = config
self._pending_approvals: dict = {}
self._approval_timeout = config.get("timeout_seconds", 300)
def classify_risk(self, action: dict) -> RiskLevel:
"""Classify an action's risk level based on its type and parameters."""
action_type = action.get("type", "")
params = action.get("params", {})
if action_type in ("delete", "destroy", "purge"):
return RiskLevel.CRITICAL
if action_type == "scale":
replicas = params.get("replicas", 0)
if replicas > self.config.get("auto_scale_max", 10):
return RiskLevel.HIGH
return RiskLevel.MEDIUM
if action_type in ("modify_config", "update_secret"):
return RiskLevel.HIGH
if action_type in ("restart", "rollout_restart"):
return RiskLevel.MEDIUM
if action_type.startswith("get_") or action_type.startswith("list_"):
return RiskLevel.LOW
return RiskLevel.HIGH # Default to high for unknown actions
async def request_approval(self, action: dict, trace_id: str) -> bool:
"""
Request human approval for a high/critical risk action.
Returns True if approved, False if denied or timed out.
"""
risk = self.classify_risk(action)
if risk in (RiskLevel.LOW, RiskLevel.MEDIUM):
# Low/medium: auto-approve but notify
await self._notify(action, trace_id, risk)
return True
# High/Critical: require human approval
approval_id = f"approval_{trace_id}_{len(self._pending_approvals)}"
self._pending_approvals[approval_id] = {
"action": action,
"risk": risk,
"trace_id": trace_id,
"status": "pending",
"created_at": asyncio.get_event_loop().time()
}
# Send notification with approve/deny buttons
await self._send_approval_request(approval_id, action, risk, trace_id)
# Wait for approval with timeout
try:
async with asyncio.timeout(self._approval_timeout):
return await self._wait_for_approval(approval_id)
except asyncio.TimeoutError:
self._pending_approvals[approval_id]["status"] = "expired"
await self._notify_timeout(action, trace_id, risk)
return False
async def _send_approval_request(self, approval_id: str, action: dict,
risk: RiskLevel, trace_id: str):
"""Send approval request via configured channels."""
message = {
"approval_id": approval_id,
"trace_id": trace_id,
"risk_level": risk.value,
"action": action,
"message": (
f"⚠️ Agent Action Approval Required\n\n"
f"**Risk Level:** {risk.value.upper()}\n"
f"**Action:** {action.get('type', 'unknown')}\n"
f"**Parameters:** {json.dumps(action.get('params', {}), indent=2)}\n"
f"**Trace ID:** {trace_id}\n\n"
f"Approve or deny within {self._approval_timeout} seconds."
)
}
# Send via Slack
if self.config.get("slack_webhook"):
await webhooks.post(self.config["slack_webhook"], message)
# Send via PagerDuty for critical actions
if risk == RiskLevel.CRITICAL and self.config.get("pagerduty_key"):
await self._create_pagerduty_incident(message)
async def _wait_for_approval(self, approval_id: str) -> bool:
"""Poll for approval status until resolved or timeout."""
while True:
await asyncio.sleep(2)
approval = self._pending_approvals.get(approval_id)
if approval and approval["status"] in ("approved", "denied"):
return approval["status"] == "approved"
def resolve_approval(self, approval_id: str, approved: bool):
"""Called by the webhook handler when a human responds."""
if approval_id in self._pending_approvals:
self._pending_approvals[approval_id]["status"] = (
"approved" if approved else "denied"
)
async def _notify(self, action: dict, trace_id: str, risk: RiskLevel):
"""Send notification for auto-approved actions."""
if self.config.get("slack_webhook"):
await webhooks.post(self.config["slack_webhook"], {
"message": (
f"✅ Auto-approved: {action.get('type')} "
f"(risk: {risk.value}, trace: {trace_id})"
)
})
async def _notify_timeout(self, action: dict, trace_id: str, risk: RiskLevel):
"""Notify when approval times out."""
if self.config.get("slack_webhook"):
await webhooks.post(self.config["slack_webhook"], {
"message": (
f"⏰ Approval timed out for: {action.get('type')} "
f"(risk: {risk.value}, trace: {trace_id}). "
f"Action was NOT executed."
)
})
6. Production Deployment Architecture
Scaling from a single developer's laptop to a production system requires careful infrastructure planning.
System Architecture
┌─────────────────────────────────────────────────────────────────┐
│ Client Layer │
│ ┌──────────┐ ┌──────────┐ ┌──────────────┐ ┌───────────┐ │
│ │ Slack │ │ Web UI │ │ CLI Client │ │ API │ │
│ └────┬─────┘ └────┬─────┘ └──────┬───────┘ └─────┬─────┘ │
│ │ │ │ │ │
└───────┼──────────────┼───────────────┼────────────────┼─────────┘
│ │ │ │
▼ ▼ ▼ ▼
┌─────────────────────────────────────────────────────────────────┐
│ API Gateway / Load Balancer │
│ (Rate limiting, Auth, Request routing) │
└──────────────────────────┬──────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ Agent Orchestrator │
│ ┌────────────┐ ┌──────────────┐ ┌───────────────────────┐ │
│ │ Prompt │ │ Tool │ │ Policy Engine + │ │
│ │ Builder │ │ Registry │ │ Approval Workflow │ │
│ └────────────┘ └──────────────┘ └───────────────────────┘ │
│ ┌────────────┐ ┌──────────────┐ ┌───────────────────────┐ │
│ │ LLM │ │ Decision │ │ Observability / │ │
│ │ Client │ │ Tracer │ │ Tracing Middleware │ │
│ └────────────┘ └──────────────┘ └───────────────────────┘ │
└──────────────────────────┬──────────────────────────────────────┘
│
┌────────────┼────────────┐
▼ ▼ ▼
┌─────────────────┐ ┌──────────┐ ┌─────────────────┐
│ Tool Plugins │ │ Vector │ │ Trace Store │
│ (K8s, AWS, │ │ Store │ │ (TimescaleDB / │
│ Datadog, etc) │ │ (Pinecone│ │ ClickHouse) │
│ │ │ / pgvec)│ │ │
└─────────────────┘ └──────────┘ └─────────────────┘
Docker Compose for Development
version: "3.8"
services:
agent-orchestrator:
build:
context: .
dockerfile: Dockerfile.agent
ports:
- "8000:8000"
environment:
- OPENAI_API_KEY=${OPENAI_API_KEY}
- KUBE_CONFIG_PATH=/mnt/secrets/kubeconfig
- TOOL_CONFIG_PATH=/app/config/tools.yaml
- POLICY_CONFIG_PATH=/app/config/policies.yaml
- SLACK_WEBHOOK_URL=${SLACK_WEBHOOK_URL}
- DATABASE_URL=postgresql://agent:secret@postgres:5432/agent_db
- VECTOR_STORE_URL=http://qdrant:6333
- OTEL_EXPORTER_OTLP_ENDPOINT=http://otel-collector:4318
volumes:
- ./config:/app/config:readonly
- ./plugins:/app/plugins:readonly
- kubeconfig:/mnt/secrets/kubeconfig:readonly
depends_on:
- postgres
- qdrant
- redis
restart: unless-stopped
postgres:
image: timescale/timescaledb:latest-pg15
environment:
- POSTGRES_DB=agent_db
- POSTGRES_USER=agent
- POSTGRES_PASSWORD=secret
volumes:
- pgdata:/var/lib/postgresql/data
ports:
- "5432:5432"
qdrant:
image: qdrant/qdrant:latest
volumes:
- qdrant_data:/qdrant/storage
ports:
- "6333:6333"
redis:
image: redis:7-alpine
ports:
- "6379:6379"
otel-collector:
image: otel/opentelemetry-collector-contrib:latest
volumes:
- ./otel-config.yaml:/etc/otelcol/config.yaml
command: ["--config=/etc/otelcol/config.yaml"]
ports:
- "4318:4318"
grafana:
image: grafana/grafana:latest
ports:
- "3000:3000"
volumes:
- ./grafana-dashboards:/etc/grafana/provisioning/dashboards
depends_on:
- otel-collector
volumes:
pgdata:
qdrant_data:
Kubernetes Deployment Manifest
apiVersion: apps/v1
kind: Deployment
metadata:
name: agent-orchestrator
namespace: devops-agent
spec:
replicas: 3
selector:
matchLabels:
app: agent-orchestrator
template:
metadata:
labels:
app: agent-orchestrator
spec:
containers:
- name: agent
image: registry.example.com/agent-orchestrator:latest
ports:
- containerPort: 8000
env:
- name: OPENAI_API_KEY
valueFrom:
secretKeyRef:
name: agent-secrets
key: openai-api-key
- name: DATABASE_URL
valueFrom:
secretKeyRef:
name: agent-secrets
key: database-url
resources:
requests:
cpu: 500m
memory: 512Mi
limits:
cpu: "2"
memory: 2Gi
livenessProbe:
httpGet:
path: /healthz
port: 8000
initialDelaySeconds: 30
periodSeconds: 10
readinessProbe:
httpGet:
path: /readyz
port: 8000
initialDelaySeconds: 15
periodSeconds: 5
volumeMounts:
- name: config
mountPath: /app/config
readOnly: true
- name: plugins
mountPath: /app/plugins
readOnly: true
volumes:
- name: config
configMap:
name: agent-config
- name: plugins
persistentVolumeClaim:
claimName: agent-plugins-pvc
---
apiVersion: v1
kind: Service
metadata:
name: agent-orchestrator
namespace: devops-agent
spec:
selector:
app: agent-orchestrator
ports:
- port: 80
targetPort: 8000
type: ClusterIP
---
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: agent-orchestrator-hpa
namespace: devops-agent
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: agent-orchestrator
minReplicas: 3
maxReplicas: 10
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 70
- type: Resource
resource:
name: memory
target:
type: Utilization
averageUtilization: 80
7. Monitoring and Alerting for the Agent Itself
The agent system needs its own observability stack — you can't fly blind on the system that's supposed to monitor your infrastructure.
Key Metrics to Track
from prometheus_client import Counter, Histogram, Gauge, Info
# Request-level metrics
agent_requests_total = Counter(
"agent_requests_total",
"Total number of agent requests processed",
["endpoint", "status"]
)
agent_request_duration = Histogram(
"agent_request_duration_seconds",
"Time spent processing an agent request",
buckets=[0.5, 1, 2, 5, 10, 30, 60]
)
# LLM-specific metrics
llm_tokens_used = Counter(
"llm_tokens_used_total",
"Total tokens consumed by LLM calls",
["model", "direction"] # direction: "input" or "output"
)
llm_response_latency = Histogram(
"llm_response_latency_seconds",
"Time to get a response from the LLM",
buckets=[0.5, 1, 2, 5, 10, 30]
)
# Tool execution metrics
tool_executions_total = Counter(
"tool_executions_total",
"Total tool executions",
["tool_name", "status"] # status: "success" or "error"
)
tool_execution_latency = Histogram(
"tool_execution_latency_seconds",
"Time to execute a tool",
["tool_name"],
buckets=[0.1, 0.5, 1, 5, 10, 30]
)
# Approval workflow metrics
approvals_pending = Gauge(
"approvals_pending",
"Number of pending approval requests"
)
approval_timeout_total = Counter(
"approval_timeouts_total",
"Number of approval requests that timed out"
)
# System health
active_traces = Gauge(
"active_traces",
"Number of currently active decision traces"
)
# Cost tracking
daily_cost_estimate = Gauge(
"daily_cost_estimate_usd",
"Estimated daily cost of LLM API calls"
)
Alert Rules (Prometheus)
# prometheus-rules.yaml
groups:
- name: agent-alerts
rules:
- alert: AgentHighErrorRate
expr: |
rate(agent_requests_total{status="error"}[5m])
/ rate(agent_requests_total[5m]) > 0.1
for: 5m
labels:
severity: warning
annotations:
summary: "Agent error rate exceeds 10%"
description: "Error rate is {{ $value | humanizePercentage }}"
- alert: AgentHighLatency
expr: |
histogram_quantile(0.95, rate(agent_request_duration_seconds_bucket[5m])) > 30
for: 10m
labels:
severity: warning
annotations:
summary: "Agent P95 latency exceeds 30s"
- alert: ApprovalBacklog
expr: approvals_pending > 5
for: 15m
labels:
severity: critical
annotations:
summary: "Approval backlog growing"
description: "{{ $value }} approvals pending for over 15 minutes"
- alert: ToolHealthDegraded
expr: |
rate(tool_executions_total{status="error"}[5m])
/ rate(tool_executions_total[5m]) > 0.2
for: 5m
labels:
severity: warning
annotations:
summary: "Tool '{{ $labels.tool_name }}' has high error rate"
- alert: HighTokenConsumption
expr: |
rate(llm_tokens_used_total{direction="input"}[1h]) > 100000
for: 1h
labels:
severity: info
annotations:
summary: "High LLM token consumption"
description: "Over 100K input tokens/hour — check for prompt loops"
8. Cost Optimization Strategies
LLM API costs can spiral quickly in production. Here are proven strategies:
Prompt Caching and Reuse
from functools import lru_cache
import hashlib
import json
class PromptCache:
"""
Caches LLM responses for identical or near-identical queries.
Uses semantic similarity to detect cacheable requests.
"""
def __init__(self, vector_store, cache_ttl: int = 3600):
self.vector_store = vector_store
self.cache_ttl = cache_ttl
def get_cached_response(self, prompt: str) -> Optional[str]:
"""Check if a similar prompt has been answered recently."""
results = self.vector_store.search(
query=prompt,
collection="prompt_cache",
top_k=1,
filter={"min_similarity": 0.95}
)
if results:
cached = results[0]
if time.time() - cached["timestamp"] < self.cache_ttl:
return cached["response"]
else:
# Expired — remove from cache
self.vector_store.delete(cached["id"])
return None
def cache_response(self, prompt: str, response: str):
"""Store a prompt-response pair for future reuse."""
self.vector_store.upsert(
collection="prompt_cache",
id=hashlib.sha256(prompt.encode()).hexdigest(),
vector=self._embed(prompt),
payload={
"prompt": prompt,
"response": response,
"timestamp": time.time()
}
)
def _embed(self, text: str) -> list:
"""Generate embedding for the prompt."""
return self.embedding_model.embed(text)
Model Tiering
class ModelRouter:
"""
Routes requests to different LLM models based on complexity.
Simple queries go to cheap/fast models; complex ones to capable models.
"""
TIER_CONFIG = {
"simple": {
"model": "gpt-4o-mini",
"max_tokens": 500,
"temperature": 0.3
},
"moderate": {
"model": "gpt-4o",
"max_tokens": 2000,
"temperature": 0.5
},
"