оркестрация, протокол A2A, self-hosted развёртывание на bare-metal VM + Kubernetes, роли участников
Multi-Agent System (MAS) — это архитектура, в которой несколько LLM-агентов работают совместно, передавая задачи и контекст друг другу. Каждый агент специализирован: один планирует, другой пишет код, третий проверяет, четвёртый работает с базой. Это аналог команды разработки.
Контекстное окно ограничено. Длинные цепочки рассуждений деградируют. Нельзя параллелить. Одна точка отказа.
Planner → Researcher → Coder → Critic → Executor. Каждый агент делает одно, но хорошо. Параллельное выполнение независимых задач.
Горизонтальный скейлинг воркеров. Независимые деплои. Разные модели для разных задач. Fault isolation.
Orchestrator принимает задачу, декомпозирует её на подзадачи, делегирует Worker-агентам, собирает результаты. Классический boss–worker.
Задачи и агенты описываются как граф (DAG). LangGraph / Temporal. Каждый узел — агент или функция, рёбра — условия перехода.
Агент генерирует ответ → Critic-агент критикует → агент улучшает. Цикл повторяется N раз или до достижения качества.
Оркестратор раздаёт независимые подзадачи параллельно нескольким агентам. Ждёт всех (join) или первого готового.
from langgraph.graph import StateGraph, END from langgraph.prebuilt import ToolNode from typing import TypedDict, List class AgentState(TypedDict): messages: List task: str plan: str result: str attempts: int # Узлы графа — каждый вызывает своего агента def planner_node(state: AgentState) -> AgentState: plan = planner_agent.invoke(state["task"]) return {**state, "plan": plan} def researcher_node(state: AgentState) -> AgentState: findings = researcher_agent.invoke(state["plan"]) state["messages"].append({"role": "researcher", "content": findings}) return state def coder_node(state: AgentState) -> AgentState: code = coder_agent.invoke(state["messages"]) return {**state, "result": code} def critic_node(state: AgentState) -> AgentState: verdict = critic_agent.invoke(state["result"]) state["attempts"] += 1 return {**state, "messages": [*state["messages"], verdict]} # Условный переход: повторить или завершить def should_retry(state: AgentState) -> str: if "APPROVED" in state["messages"][-1]["content"]: return "end" if state["attempts"] >= 3: return "end" # предохранитель return "retry" # Сборка графа graph = StateGraph(AgentState) graph.add_node("planner", planner_node) graph.add_node("researcher", researcher_node) graph.add_node("coder", coder_node) graph.add_node("critic", critic_node) graph.set_entry_point("planner") graph.add_edge("planner", "researcher") graph.add_edge("researcher", "coder") graph.add_edge("coder", "critic") graph.add_conditional_edges("critic", should_retry, { "retry": "coder", "end": END }) app = graph.compile(checkpointer=MemorySaver())
A2A (Agent-to-Agent) — открытый протокол от Google DeepMind (апрель 2025). Решает проблему
несовместимости агентов из разных фреймворков. Агент публикует Agent Card
— JSON-манифест своих возможностей. Другие агенты находят его и вызывают стандартным способом.
Каждый агент отдаёт /.well-known/agent.json с описанием своих capabilities, эндпоинтов и схем входа/выхода.
Статусы задачи: submitted → working → input-required → completed / failed. Streaming через SSE или WebSocket.
{
"name": "CodeReviewAgent",
"description": "Проверяет Python-код на ошибки и стиль",
"version": "1.2.0",
"url": "https://agents.internal/code-review",
"capabilities": {
"streaming": true,
"pushNotifications": false,
"stateTransitionHistory": true
},
"skills": [
{
"id": "review_python",
"name": "Review Python code",
"inputModes": ["text"],
"outputModes": ["text", "data"],
"examples": ["review this function for bugs"]
}
],
"authentication": {
"schemes": ["Bearer"]
}
}
import httpx import json class A2AClient: def __init__(self, agent_url: str): self.base = agent_url self.client = httpx.AsyncClient() async def discover(self) -> dict: # Получить Agent Card r = await self.client.get(f"{self.base}/.well-known/agent.json") return r.json() async def send_task(self, message: str, session_id: str = None) -> dict: payload = { "jsonrpc": "2.0", "method": "tasks/send", "id": 1, "params": { "sessionId": session_id, "message": { "role": "user", "parts": [{"text": message}] } } } r = await self.client.post( f"{self.base}/a2a", json=payload, headers={"Authorization": "Bearer <token>"} ) return r.json()["result"] # Использование: оркестратор вызывает воркер-агента async def orchestrate(): client = A2AClient("http://code-review-agent:8080") card = await client.discover() print(f"Agent: {card['name']}, skills: {[s['id'] for s in card['skills']]}") result = await client.send_task("Review this code: def f(x): return x*x") print(result)
google-adk или a2a-sdk — прямая реализация сложнее.
Параллельно существует протокол MCP (Anthropic) для инструментов; A2A — для агент–агент вызовов.
| Фреймворк | Модель оркестрации | Self-hosted | A2A / MCP | Когда выбирать |
|---|---|---|---|---|
| LangGraph рек. | Граф (DAG + циклы), StateGraph | полностью | A2A MCP | Сложные workflow, нужен полный контроль, production |
| CrewAI | Ролевые агенты, иерархия, процессы | полностью | A2A | Быстрый старт, команды агентов, HR-метафора |
| AutoGen v0.4 | Conversation-driven, GroupChat | полностью | частично | Multi-turn диалоги, research tasks, Microsoft stack |
| Google ADK | Иерархия + A2A нативно | частично | нативно | A2A-экосистема, Vertex AI + свои агенты |
| Temporal + LLM | Workflow engine, durable execution | полностью | нет | Критичные к надёжности задачи, retry, saga-паттерн |
| Celery + агенты | Очереди задач, воркеры | полностью | нет | Простые асинхронные задачи, уже есть Celery в стеке |
apiVersion: apps/v1 kind: Deployment metadata: name: mas-orchestrator namespace: agents spec: replicas: 2 selector: matchLabels: { app: orchestrator } template: metadata: labels: { app: orchestrator } annotations: prometheus.io/scrape: "true" prometheus.io/port: "9090" spec: containers: - name: orchestrator image: registry.internal/mas-orchestrator:v1.4.2 env: - { name: LLM_BASE_URL, value: "http://vllm-service:8000/v1" } - { name: REDIS_URL, value: "redis://redis-cluster:6379" } - { name: VECTOR_DB_URL, value: "http://qdrant:6333" } - { name: LOG_LEVEL, value: "INFO" } resources: requests: { cpu: "500m", memory: "1Gi" } limits: { cpu: "2000m", memory: "4Gi" } livenessProbe: httpGet: { path: /health, port: 8080 } initialDelaySeconds: 15 --- apiVersion: autoscaling/v2 kind: HorizontalPodAutoscaler metadata: name: worker-hpa namespace: agents spec: scaleTargetRef: apiVersion: apps/v1 kind: Deployment name: mas-worker minReplicas: 2 maxReplicas: 20 metrics: - type: External external: metric: name: redis_queue_length # кастомная метрика из Prometheus target: { type: Value, value: "5" } # 1 воркер на 5 задач
#!/bin/bash # GPU VM: 2x A100 80GB, запуск vLLM как сервиса # 1. Установка pip install vllm --break-system-packages # 2. Запуск модели с tensor parallelism python -m vllm.entrypoints.openai.api_server \ --model mistralai/Mixtral-8x7B-Instruct-v0.1 \ --tensor-parallel-size 2 \ --max-model-len 32768 \ --quantization awq \ --host 0.0.0.0 \ --port 8000 \ --served-model-name mixtral \ --enable-prefix-caching \ # кешировать системный промпт --max-num-seqs 64 # concurrent requests # 3. Systemd unit для автозапуска cat > /etc/systemd/system/vllm.service <<EOF [Unit] Description=vLLM Inference Server After=network.target [Service] User=ubuntu WorkingDirectory=/home/ubuntu ExecStart=/home/ubuntu/.local/bin/python -m vllm.entrypoints.openai.api_server \ --model /models/mixtral --tensor-parallel-size 2 --port 8000 Restart=always RestartSec=10 [Install] WantedBy=multi-user.target EOF systemctl enable --now vllm
Прежде всего нужен LLM endpoint. Устанавливаем vLLM, запускаем модель, убеждаемся что GET /v1/models отвечает из K8s кластера.
# С любого K8s node или пода: curl http://<GPU-VM-IP>:8000/v1/models # Должен ответить: {"data": [{"id": "mixtral", ...}]} # Если закрыто — открыть в firewall: ufw allow from <K8s-CIDR> to any port 8000
Создаём namespace agents, деплоим Redis (для очередей и state) и Qdrant (vector search для RAG).
kubectl create namespace agents # Redis с persistence helm repo add bitnami https://charts.bitnami.com/bitnami helm install redis bitnami/redis \ --namespace agents \ --set auth.enabled=true \ --set auth.password=changeme \ --set master.persistence.size=10Gi # Qdrant vector DB helm repo add qdrant https://qdrant.github.io/qdrant-helm helm install qdrant qdrant/qdrant \ --namespace agents \ --set persistence.size=50Gi # Проверка kubectl get pods -n agents # redis-master-0 1/1 Running # qdrant-0 1/1 Running
Каждый агент — отдельный Docker-образ. Минимальный агент: FastAPI сервер + LangGraph логика + подключение к LLM. Публикуем в внутренний registry.
from fastapi import FastAPI from langchain_openai import ChatOpenAI from langchain_core.messages import HumanMessage import os, json app = FastAPI() # Подключаемся к vLLM через OpenAI-совместимый API llm = ChatOpenAI( base_url=os.getenv("LLM_BASE_URL", "http://vllm-service:8000/v1"), model="mixtral", api_key="not-needed", temperature=0.3, max_tokens=2048, ) # A2A Agent Card AGENT_CARD = { "name": os.getenv("AGENT_NAME", "GenericWorker"), "version": "1.0.0", "url": os.getenv("AGENT_URL"), "capabilities": {"streaming": False}, "skills": [{"id": "process_task", "name": "Process task"}] } @app.get("/.well-known/agent.json") def agent_card(): return AGENT_CARD @app.post("/a2a") async def handle_task(req: dict): msg = req["params"]["message"]["parts"][0]["text"] resp = await llm.ainvoke([HumanMessage(content=msg)]) return { "jsonrpc": "2.0", "id": req["id"], "result": { "status": {"state": "completed"}, "artifacts": [{"parts": [{"text": resp.content}]}] } } @app.get("/health") def health(): return {"status": "ok"}
Применяем Deployment + Service для каждого агента. Оркестратор делает Service Discovery через DNS (http://coder-agent.agents.svc.cluster.local).
apiVersion: apps/v1 kind: Deployment metadata: { name: coder-agent, namespace: agents } spec: replicas: 3 selector: matchLabels: { app: coder-agent } template: metadata: labels: { app: coder-agent } spec: containers: - name: agent image: registry.internal/coder-agent:latest env: - { name: AGENT_NAME, value: CoderAgent } - { name: AGENT_URL, value: "http://coder-agent.agents.svc.cluster.local" } - { name: LLM_BASE_URL, value: "http://<GPU-VM-IP>:8000/v1" } ports: - { containerPort: 8080 } --- apiVersion: v1 kind: Service metadata: { name: coder-agent, namespace: agents } spec: selector: { app: coder-agent } ports: - { port: 80, targetPort: 8080 }
Критически важно с первого дня. Деплоим kube-prometheus-stack, добавляем кастомные метрики агентов.
# Prometheus + Grafana через Helm helm install monitoring prometheus-community/kube-prometheus-stack \ --namespace monitoring --create-namespace \ --set grafana.adminPassword=changeme # В коде агента — экспозиция метрик # pip install prometheus-fastapi-instrumentator
from prometheus_client import Counter, Histogram, Gauge from prometheus_fastapi_instrumentator import Instrumentator # Кастомные метрики task_counter = Counter( "agent_tasks_total", "Total tasks processed", ["agent_name", "status"] ) task_duration = Histogram( "agent_task_duration_seconds", "Task processing time", buckets=[1, 5, 15, 30, 60, 120] ) llm_tokens = Counter( "agent_llm_tokens_total", "LLM tokens consumed", ["agent_name", "type"] # prompt / completion ) active_tasks = Gauge("agent_active_tasks", "Currently running tasks") # Автоинструментация HTTP Instrumentator().instrument(app).expose(app)
Настроить Terraform для авто-создания VM в облаке при перегрузке. Облачные воркеры подключаются к тому же Redis через WireGuard туннель.
#!/bin/bash # Запускается автоматически при создании облачной VM # Установка зависимостей apt-get update -qq pip install fastapi uvicorn langchain-openai redis # Настройка WireGuard туннеля до K8s кластера wg-quick up /etc/wireguard/wg0.conf # Запуск воркера, подключённого к локальному Redis REDIS_URL="redis://:changeme@10.0.0.1:6379" \ LLM_BASE_URL="http://10.0.0.5:8000/v1" \ AGENT_NAME="CloudWorker" \ uvicorn main:app --host 0.0.0.0 --port 8080 & # Регистрация в Redis как доступный воркер redis-cli -u $REDIS_URL SADD workers:available "$(curl -s ifconfig.me):8080"
ORCHESTRATOR_SYSTEM = """ Ты — Orchestrator в multi-agent системе. Твоя роль: планирование и делегирование. Никогда не решай задачи напрямую. ДОСТУПНЫЕ АГЕНТЫ: - researcher: поиск информации, веб-поиск, анализ документов - coder: написание и исправление кода на Python/JS/Go - critic: проверка кода, поиск ошибок, code review - executor: запуск команд, работа с файловой системой - writer: написание текстов, документации ПРОЦЕСС: 1. Проанализируй задачу. Ответь себе: что нужно на входе? что на выходе? 2. Разбей на шаги. Определи зависимости. 3. Создай план в JSON формате с полями: - task_id, agent, input, depends_on[], priority 4. Передай план через tool: execute_plan(plan) 5. После получения результатов — агрегируй их. ПРАВИЛА: - Не превышай 7 подзадач в одном плане - Параллели всё, что можно (depends_on: []) - При ошибке воркера — попробуй другого или уточни задачу - Всегда запрашивай critic после coder - Устанавливай timeout на каждую задачу """ DECOMPOSE_PROMPT = """ Задача: {task} Создай план выполнения. Формат JSON: {{ "plan": [ {{ "task_id": "t1", "agent": "researcher", "input": "найди информацию о ...", "depends_on": [], "priority": 1, "timeout_seconds": 60 }}, ... ], "expected_output": "описание финального результата" }} """
CODER_SYSTEM = """ Ты — Coder Agent. Специализируешься на написании Python-кода. ВХОД: Описание задачи + контекст от researcher ВЫХОД (строго JSON): { "code": "...", // готовый Python код "explanation": "...", // что делает код (1-3 предложения) "tests": "...", // pytest тесты "dependencies": [], // pip пакеты "confidence": 0.0-1.0 // уверенность в решении } ИНСТРУМЕНТЫ: - execute_python(code: str) → stdout, stderr — запустить и проверить - search_docs(query: str) → str — найти в документации ПРАВИЛА: 1. ВСЕГДА запускай код через execute_python перед отправкой 2. Если confidence < 0.6 — верни {"error": "недостаточно контекста", "needs": [...]} 3. Максимум 5 попыток выполнения 4. Никогда не используй subprocess.run с shell=True ПРИМЕР хорошего ответа: { "code": "def parse_csv(path: str) -> list[dict]:\\n ...", "explanation": "Читает CSV файл и возвращает список словарей", "tests": "def test_parse_csv():\\n ...", "dependencies": ["pandas"], "confidence": 0.95 } """ from langchain.agents import create_react_agent from langchain_core.tools import tool @tool def execute_python(code: str) -> str: """Выполнить Python код в sandbox и вернуть вывод""" import subprocess result = subprocess.run( ["python", "-c", code], capture_output=True, text=True, timeout=30 ) return result.stdout + (f"\nERROR: {result.stderr}" if result.stderr else "") coder_agent = create_react_agent( model=llm, tools=[execute_python, search_docs], prompt=ChatPromptTemplate.from_messages([ ("system", CODER_SYSTEM), ("human", "{input}"), ("placeholder", "{agent_scratchpad}"), ]) )
| Что мониторить | Метрика | Порог тревоги | Что делать |
|---|---|---|---|
| LLM latency | vllm_request_duration_seconds | > 30s P95 | Добавить GPU, снизить max_tokens, квантизация |
| Очередь задач | redis_queue_length | > 50 задач | Autoscale workers (HPA), burst на облако |
| Ошибки агентов | agent_tasks_total{status="error"} | > 5% | Смотреть трейсы, улучшить промпты, добавить retry |
| Зависшие задачи | agent_task_duration_seconds | > timeout | Принудительный kill, алерт оператору |
| Память воркеров | container_memory_usage_bytes | > 80% limit | Очистка контекста, увеличить limits |
| GPU VRAM | DCGM_FI_DEV_FB_USED | > 90% | Снизить batch size, добавить GPU |
| Токены / стоимость | agent_llm_tokens_total | бюджет | Оптимизировать промпты, сократить контекст |
from opentelemetry import trace from opentelemetry.sdk.trace import TracerProvider from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter import time, uuid tracer = trace.get_tracer("mas.agent") class TracedAgent: def __init__(self, name: str, agent): self.name = name self.agent = agent async def invoke(self, task: str, parent_trace_id: str = None): task_id = str(uuid.uuid4())[:8] with tracer.start_as_current_span( f"agent.{self.name}", attributes={ "agent.name": self.name, "task.id": task_id, "task.input_length": len(task), "parent.trace_id": parent_trace_id or "root", } ) as span: start = time.time() try: result = await self.agent.ainvoke({"input": task}) span.set_attribute("task.status", "success") span.set_attribute("task.output_length", len(str(result))) return result except Exception as e: span.set_attribute("task.status", "error") span.set_attribute("error.message", str(e)) span.record_exception(e) raise finally: duration = time.time() - start span.set_attribute("task.duration_s", duration) task_duration.observe(duration) # prometheus
helm install jaeger jaegertracing/jaeger --namespace monitoring
/health отвечает, Prometheus метрики доступны, systemd автозапуск
agents изолирован NetworkPolicy
max_iterations=10, timeout=120s, watchdog-горутина/тред.
import asyncio, time from enum import Enum class State(Enum): CLOSED = "closed" # норма OPEN = "open" # блокируем вызовы HALF = "half" # пробный вызов class LLMCircuitBreaker: def __init__(self, failure_threshold=5, recovery_timeout=60): self.state = State.CLOSED self.failures = 0 self.threshold = failure_threshold self.timeout = recovery_timeout self.opened_at = None self.fallback_url = "https://api.openai.com/v1" # облачный fallback async def call(self, llm_fn, *args, **kwargs): if self.state == State.OPEN: if time.time() - self.opened_at > self.timeout: self.state = State.HALF else: # Переключиться на fallback LLM return await self._fallback_call(*args, **kwargs) try: result = await asyncio.wait_for(llm_fn(*args, **kwargs), timeout=90) self._on_success() return result except (Exception, asyncio.TimeoutError) as e: self._on_failure() raise def _on_success(self): self.failures = 0 self.state = State.CLOSED def _on_failure(self): self.failures += 1 if self.failures >= self.threshold: self.state = State.OPEN self.opened_at = time.time()