Главная / Статьи / Практические заметки: Агентные архитектуры — Статья 6: Оркестрация множества агентов

Практические заметки: Агентные архитектуры — Статья 6: Оркестрация множества агентов

Пошаговое руководство по практическим заметкам: агентные архитектуры — Статья 6: Оркестрация множества агентов: контракты, проверки и готовые блоки кода для команд, использующих эту паттерн-архитектуру.

5399 слов

В следующих примечаниях описывается практический подход к изучению темы «Агентные архитектуры — Статья 6: Шаблоны оркестрации множества агентов». Основное внимание уделяется контрактам, проверкам и местам для вставки кода, а не мотивирующим аспектам. На этапе обзора сначала запишите контракт: необходимые входные данные, сигнал успешного выполнения и действия при частичной неудаче. Такой список поможет сохранять честность при последующих изменениях кода. Рядом с функциональными результатами записывайте время выполнения и стоимость токенов или запросов. Отображение затрат с самого начала предотвращает неожиданные расходы при переходе от демо-версии к общедоступным средам.

Что вы здесь найдете

Этап «Что вы найдете» работает наилучшим образом, если рассматривать его как измеримую поверхность. Соберите один идеальный пример работы, один случай сбоя и запись о возврате к предыдущему состоянию перед расширением объема работ. Храните конфигурацию вне кода приложения. Файлы среды, хранилища секретов и флаги функций должны находиться в одном месте, чтобы операторы могли их проверять, не читая весь граф. Сохраняйте состояние графа простым и типизированным. Вложенные структуры данных скрывают информацию о том, какой узел заполнил тот или иной поле, и приводят к нарушению возобновления работы после перерывов.

Почему один агент достигает своего предела

Этап «Почему один агент» работает наилучшим образом, когда его рассматривают как измеримую поверхность. Соберите один идеальный пример работы, один случай сбоя и записку о возврате к предыдущему состоянию перед расширением объёма работ. Документируйте одновременно успешный и восстановительный пути выполнения. Повторные попытки, проверки человеком и обработка неработающих сообщений являются частью продукта, а не этапом последующей доработки. Сохраняйте структуру графа простой и типизированной; вложенные структуры скрывают информацию о том, какой узел заполнил тот или иной поле, и мешают возобновлению работы после прерываний.

Таксономия шаблонов многократных агентов

Стадия «Таксономия многоклиентских систем» работает наилучшим образом, когда рассматривается как измеримая структура. Соберите один идеальный пример работы, один случай сбоя и запись о возврате к предыдущему состоянию перед расширением объёма работ. Предпочитайте небольшие, тестируемые единицы вместо обширных скриптов. При сбое какого-либо шага причина должна указывать на конкретную ответственность, а не на запутанную цепочку операций. Сохраняйте простую и типизированную структуру графа. Вложенные элементы скрывают информацию о том, какой узел заполнил тот или иной поле, и мешают возобновлению работы после прерываний. Стадия «Таксономия многоклиентских систем» работает наилучшим образом, когда рассматривается как измеримая структура. Соберите один идеальный пример работы, один случай сбоя и запись о возврате к предыдущему состоянию перед расширением объёма работ. Записывайте временные показатели, а также стоимость токенов или запросов рядом с функциональными результатами. Отображение стоимости на раннем этапе предотвращает неожиданные расходы при переходе от демо-версии к общим средам.

+--------------------+--------------------------------------------------+------------------+
| Pattern            | Structure                                        | Best For         |
+--------------------+--------------------------------------------------+------------------+
| Supervisor-Worker  | One orchestrator decomposes + delegates          | Complex tasks    |
|                    | N workers execute specialized subtasks           | needing expert   |
|                    |                                                  | decomposition    |
+--------------------+--------------------------------------------------+------------------+
| Pipeline           | Agent A -> Agent B -> Agent C (sequential)       | Transformation   |
|                    | Each processes the previous output               | chains, ETL-like |
|                    |                                                  | workflows        |
+--------------------+--------------------------------------------------+------------------+
| Parallel Fan-out   | Orchestrator sends same/related task to N agents | Research,        |
|                    | Results aggregated into single output            | analysis tasks   |
|                    |                                                  | that parallelize |
+--------------------+--------------------------------------------------+------------------+
| Debate / Critique  | Agent A produces solution, Agent B critiques     | High-stakes      |
|                    | Agent C synthesizes or adjudicates               | outputs needing  |
|                    |                                                  | adversarial QA   |
+--------------------+--------------------------------------------------+------------------+

Шаблон 1: Надзиратель-рабочий

Для этапа «Надзорщик-рабочий» в шаблоне 1 необходимо заранее определить входные данные, ответственного за выполнение шага и критерии завершения перед изменением кода. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Храните конфигурацию вне кода приложения. Файлы среды, хранилища секретов и флаги функций должны находиться в одном месте, чтобы операторы могли их проверять, не читая весь код. Вносите ручное утверждение для операций, связанных с тратой денег или изменением производственных данных. Подключение компонентов во время компиляции не гарантирует полноты бизнес-логики.

          ┌─────────────────────────────┐
          │       SUPERVISOR AGENT      │
          │  (task decomposition +      │
          │   result synthesis)         │
          └──────┬───────────┬──────────┘
                 │           │
        ┌────────▼──┐   ┌────▼──────────┐  ┌──────────────┐
        │  Worker A  │  │   Worker B    │  │   Worker C   │
        │  (code     │  │  (security    │  │  (test       │
        │   analysis)│  │   review)     │  │   coverage)  │
        └────────────┘  └───────────────┘  └──────────────┘
# harness/multi_agent/supervisor.py
from typing import Literal, TypedDict, Annotated, List
from langgraph.graph import StateGraph, END, START
from langgraph.graph.message import add_messages
from langchain_core.messages import BaseMessage, HumanMessage, SystemMessage
from langchain_aws import ChatBedrock
from pydantic import BaseModel
import boto3


class SupervisorState(TypedDict):
    messages: Annotated[List[BaseMessage], add_messages]
    task_spec: str
    subtasks: List[dict]           # decomposed work items
    worker_results: dict           # keyed by subtask ID
    current_worker: str            # which worker is active
    synthesis_complete: bool
    agent_run_id: str


class SubtaskAssignment(BaseModel):
    """Structured output from supervisor decomposition."""
    subtask_id: str
    worker_type: Literal["code_analyst", "security_reviewer", "test_evaluator"]
    description: str
    depends_on: List[str]          # subtask IDs this depends on
    priority: int


SUPERVISOR_SYSTEM_PROMPT = """
You are an orchestrator agent. You do not write code or perform analysis yourself.
Your job is to:

1. Break the task into discrete subtasks
2. Assign each subtask to the correct specialist worker
3. Track dependencies between subtasks
4. Synthesize worker outputs into a coherent final result

Available workers:
- code_analyst: Reads and analyzes code structure, dependencies, patterns
- security_reviewer: Evaluates security implications, checks against CVEs
- test_evaluator: Assesses test coverage, identifies gaps

When decomposing, be specific. A subtask description like "analyze the auth module"
is useful. "analyze the code" is not.

Respond in JSON when asked to decompose. Respond in prose when asked to synthesize.
"""


def build_supervisor(region: str = "us-east-1") -> callable:
    bedrock = boto3.client("bedrock-runtime", region_name=region)

    # Supervisor uses the heavier model — it's doing strategic reasoning
    supervisor_model = ChatBedrock(
        client=bedrock,
        model_id="anthropic.claude-3-7-sonnet-20250219-v1:0",
        model_kwargs={
            "temperature": 0.2,
            "max_tokens": 8000,
            "thinking": {"type": "enabled", "budget_tokens": 5000}
        }
    )

    def supervisor_node(state: SupervisorState) -> SupervisorState:
        if not state.get("subtasks"):
            # First pass: decompose the task
            response = supervisor_model.invoke([
                SystemMessage(content=SUPERVISOR_SYSTEM_PROMPT),
                HumanMessage(content=f"Decompose this task into subtasks:\n{state['task_spec']}")
            ])
            import json
            try:
                subtasks = json.loads(response.content)
                if isinstance(subtasks, dict) and "subtasks" in subtasks:
                    subtasks = subtasks["subtasks"]
            except json.JSONDecodeError:
                subtasks = []

            return {**state, "subtasks": subtasks}

        # All workers done: synthesize
        results_summary = "\n\n".join([
            f"=== {worker} ===\n{result}"
            for worker, result in state["worker_results"].items()
        ])
        synthesis_prompt = f"""
Original task: {state['task_spec']}

Worker results:
{results_summary}

Synthesize these into a coherent final report. Highlight conflicts between
worker findings and make clear recommendations.
"""
        response = supervisor_model.invoke([
            SystemMessage(content=SUPERVISOR_SYSTEM_PROMPT),
            HumanMessage(content=synthesis_prompt)
        ])
        return {
            **state,
            "messages": state["messages"] + [response],
            "synthesis_complete": True,
        }

    return supervisor_node


def route_to_worker(state: SupervisorState) -> str:
    """
    Routes to the next worker with unfinished subtasks.
    Returns END when all subtasks are complete and synthesis is done.
    """
    if state.get("synthesis_complete"):
        return END

    # Find next unfinished subtask whose dependencies are met
    completed = set(state.get("worker_results", {}).keys())
    for subtask in state.get("subtasks", []):
        sid = subtask["subtask_id"]
        if sid in completed:
            continue
        deps = set(subtask.get("depends_on", []))
        if deps.issubset(completed):
            return subtask["worker_type"]   # route to this worker

    # All subtasks done, back to supervisor for synthesis
    return "supervisor"
# harness/multi_agent/workers.py
from langchain_aws import ChatBedrock
from langchain_core.messages import SystemMessage, HumanMessage
import boto3


CODE_ANALYST_PROMPT = """
You are a code analysis specialist. You receive specific, bounded analysis tasks.
Focus only on what you were asked to analyze. Do not expand scope.
Return structured findings: what you found, confidence level, and specific evidence.
"""

SECURITY_REVIEWER_PROMPT = """
You are a security review specialist. You look for vulnerabilities, insecure patterns,
and CVE-relevant code. Reference specific CWE numbers when applicable.
Return structured findings with severity levels (CRITICAL, HIGH, MEDIUM, LOW).
"""

TEST_EVALUATOR_PROMPT = """
You are a test coverage specialist. You assess test quality and identify gaps.
Focus on: coverage percentage where available, missing edge cases, and untested paths.
Return structured findings with specific test gaps and suggested test cases.
"""

WORKER_PROMPTS = {
    "code_analyst": CODE_ANALYST_PROMPT,
    "security_reviewer": SECURITY_REVIEWER_PROMPT,
    "test_evaluator": TEST_EVALUATOR_PROMPT,
}


def build_worker(worker_type: str, region: str = "us-east-1") -> callable:
    bedrock = boto3.client("bedrock-runtime", region_name=region)

    # Workers use the faster model — they execute a specific, bounded task
    worker_model = ChatBedrock(
        client=bedrock,
        model_id="anthropic.claude-3-5-sonnet-20241022-v2:0",
        model_kwargs={"temperature": 0, "max_tokens": 4000}
    )

    system_prompt = WORKER_PROMPTS[worker_type]

    def worker_node(state: SupervisorState) -> SupervisorState:
        # Find the subtask assigned to this worker type
        completed = set(state.get("worker_results", {}).keys())
        current_subtask = None
        for subtask in state["subtasks"]:
            if subtask["worker_type"] == worker_type and subtask["subtask_id"] not in completed:
                deps = set(subtask.get("depends_on", []))
                if deps.issubset(completed):
                    current_subtask = subtask
                    break

        if not current_subtask:
            return state

        # Pass relevant prior results as context if this task has dependencies
        context = ""
        if current_subtask.get("depends_on"):
            for dep_id in current_subtask["depends_on"]:
                if dep_id in state.get("worker_results", {}):
                    context += f"\nPrevious analysis ({dep_id}):\n{state['worker_results'][dep_id]}\n"

        prompt = f"Task: {current_subtask['description']}"
        if context:
            prompt = f"Prior context:{context}\n\n{prompt}"

        response = worker_model.invoke([
            SystemMessage(content=system_prompt),
            HumanMessage(content=prompt)
        ])

        updated_results = {**state.get("worker_results", {})}
        updated_results[current_subtask["subtask_id"]] = response.content

        return {**state, "worker_results": updated_results}

    return worker_node
# harness/multi_agent/supervisor_graph.py
from langgraph.graph import StateGraph, END, START
from langgraph.checkpoint.memory import MemorySaver
from harness.multi_agent.supervisor import SupervisorState, build_supervisor, route_to_worker
from harness.multi_agent.workers import build_worker


def build_supervisor_graph(region: str = "us-east-1"):
    graph = StateGraph(SupervisorState)

    graph.add_node("supervisor", build_supervisor(region))
    graph.add_node("code_analyst", build_worker("code_analyst", region))
    graph.add_node("security_reviewer", build_worker("security_reviewer", region))
    graph.add_node("test_evaluator", build_worker("test_evaluator", region))

    graph.add_edge(START, "supervisor")

    graph.add_conditional_edges(
        "supervisor",
        route_to_worker,
        {
            "code_analyst": "code_analyst",
            "security_reviewer": "security_reviewer",
            "test_evaluator": "test_evaluator",
            END: END,
        }
    )

    # All workers route back to supervisor after completing their subtask
    for worker in ["code_analyst", "security_reviewer", "test_evaluator"]:
        graph.add_edge(worker, "supervisor")

    return graph.compile(checkpointer=MemorySaver())

Шаблон 2: Оркестрация потоков обработки

Для этапа оркестрации Pipeline Orchestration Pattern 2 необходимо заранее определить входные данные, ответственного за выполнение шага и критерии завершения перед изменением кода. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Необходимо одновременно задокументировать успешный сценарий выполнения и сценарий восстановления. Повторные попытки, проверки человеком и обработка неработоспособных сообщений являются частью продукта, а не элементами последующей доработки. Внедрять утверждение человеком для тех операций, которые влекут за собой расходы или изменение производственных данных. Подключения на этапе компиляции не заменяют полноты бизнес-логики.

  ┌───────────────┐     ┌───────────────┐     ┌───────────────┐
  │   Agent A     │     │   Agent B     │     │   Agent C     │
  │  (extraction) │────>│ (enrichment)  │────>│ (validation)  │
  └───────────────┘     └───────────────┘     └───────────────┘
       output A               output B               output C
    becomes input B        becomes input C         final result
# harness/multi_agent/pipeline.py
from typing import TypedDict, Annotated, List, Optional, Any
from langgraph.graph import StateGraph, END, START
from langchain_core.messages import BaseMessage, HumanMessage, SystemMessage
from langchain_aws import ChatBedrock
import boto3


class PipelineState(TypedDict):
    original_input: str
    stage_outputs: List[dict]       # accumulates each stage's compressed result
    current_stage: int
    final_output: Optional[str]
    agent_run_id: str


def compress_for_handoff(full_output: str, model: ChatBedrock) -> str:
    """
    Compresses a stage's full output into a structured handoff summary.
    This is the core of pipeline context management — the next agent gets
    the substance, not the reasoning trace.
    """
    response = model.invoke([
        SystemMessage(content="""
Compress the following agent output into a structured handoff summary.
Include: key findings, decisions made, artifacts produced, and what the
next stage needs to know. Discard reasoning traces and intermediate steps.
Target length: 20% of original. Use bullet points for clarity.
"""),
        HumanMessage(content=f"Compress this:\n\n{full_output}")
    ])
    return response.content


def build_pipeline_stage(
    stage_name: str,
    system_prompt: str,
    region: str = "us-east-1"
) -> callable:
    bedrock = boto3.client("bedrock-runtime", region_name=region)
    model = ChatBedrock(
        client=bedrock,
        model_id="anthropic.claude-3-5-sonnet-20241022-v2:0",
        model_kwargs={"temperature": 0, "max_tokens": 6000}
    )
    # Cheaper model for compression — this is mechanical, not creative
    compressor = ChatBedrock(
        client=bedrock,
        model_id="anthropic.claude-haiku-3-5",
        model_kwargs={"temperature": 0, "max_tokens": 2000}
    )

    def stage_node(state: PipelineState) -> PipelineState:
        # Build context from compressed prior stage outputs only
        prior_context = ""
        for past_stage in state.get("stage_outputs", []):
            prior_context += f"\n### {past_stage['stage']} output:\n{past_stage['compressed_output']}\n"

        prompt = f"Original task: {state['original_input']}\n"
        if prior_context:
            prompt += f"\nPrior stage results:\n{prior_context}\n"
        prompt += f"\nNow perform your stage: {stage_name}"

        response = model.invoke([
            SystemMessage(content=system_prompt),
            HumanMessage(content=prompt)
        ])
        full_output = response.content

        # Compress before storing — next stage won't see raw output
        compressed = compress_for_handoff(full_output, compressor)

        updated_outputs = list(state.get("stage_outputs", []))
        updated_outputs.append({
            "stage": stage_name,
            "full_output": full_output,
            "compressed_output": compressed,
        })

        return {
            **state,
            "stage_outputs": updated_outputs,
            "current_stage": state.get("current_stage", 0) + 1,
        }

    return stage_node

Pattern 3: Параллельное расширение/сужение

Для этапа параллельного распространения шаблона 3 необходимо заранее определить входные данные, ответственного за выполнение шага и критерии завершения перед внесением изменений в код. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Лучше использовать небольшие, тестируемые единицы кода вместо обширных скриптов. При сбое шага причина должна быть связана с конкретной областью ответственности, а не с запутанной структурой обработки данных. Внедрять человеческое утверждение для операций, связанных с тратой денег или изменением производственных данных. Компиляционная настройка не гарантирует полноты бизнес-процессов. Для этапа параллельного распространения шаблона 3 необходимо заранее определить входные данные, ответственного за выполнение шага и критерии завершения перед внесением изменений в код. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Регистрировать время выполнения, а также стоимость токенов или запросов вместе с функциональными результатами. Отображение стоимости на ранних этапах предотвращает неожиданные счета при переходе от демо-режима к общей среде.

Среды.

                    ┌──────────────────┐
                    │   ORCHESTRATOR   │
                    │  (task splitter) │
                    └──┬───┬───┬───┬──┘
                       │   │   │   │
            ┌──────────▼┐ ┌▼─┐ ┌▼──┐ ┌▼──────────┐
            │  Agent 1  │ │A2│ │A3 │ │  Agent 4  │
            │ (region A)│ │  │ │   │ │ (region D)│
            └──────────┬┘ └┬─┘ └┬──┘ └┬──────────┘
                       │   │    │     │
                    ┌──▼───▼────▼─────▼──┐
                    │     AGGREGATOR     │
                    │  (result merger)   │
                    └────────────────────┘
# harness/multi_agent/fanout.py
import asyncio
from typing import List, TypedDict, Annotated, Optional
from langchain_aws import ChatBedrock
from langchain_core.messages import SystemMessage, HumanMessage
import boto3


class FanoutResult(TypedDict):
    agent_id: str
    input_slice: str
    output: str
    success: bool
    error: Optional[str]


async def run_agent_async(
    agent_id: str,
    input_slice: str,
    system_prompt: str,
    model: ChatBedrock,
) -> FanoutResult:
    """Run a single agent asynchronously."""
    try:
        response = await model.ainvoke([
            SystemMessage(content=system_prompt),
            HumanMessage(content=input_slice)
        ])
        return FanoutResult(
            agent_id=agent_id,
            input_slice=input_slice,
            output=response.content,
            success=True,
            error=None,
        )
    except Exception as e:
        return FanoutResult(
            agent_id=agent_id,
            input_slice=input_slice,
            output="",
            success=False,
            error=str(e),
        )


async def fan_out(
    task_slices: List[str],
    system_prompt: str,
    region: str = "us-east-1",
    max_concurrent: int = 5,   # don't hammer Bedrock rate limits
) -> List[FanoutResult]:
    """
    Runs agents concurrently with a semaphore to cap parallelism.
    max_concurrent protects against Bedrock throttling — you
    will hit rate limits if you fire 20 concurrent requests.
    """
    bedrock = boto3.client("bedrock-runtime", region_name=region)
    model = ChatBedrock(
        client=bedrock,
        model_id="anthropic.claude-3-5-sonnet-20241022-v2:0",
        model_kwargs={"temperature": 0, "max_tokens": 4000}
    )

    semaphore = asyncio.Semaphore(max_concurrent)

    async def bounded_run(agent_id, slice_content):
        async with semaphore:
            return await run_agent_async(agent_id, slice_content, system_prompt, model)

    tasks = [
        bounded_run(f"agent_{i}", slice_content)
        for i, slice_content in enumerate(task_slices)
    ]
    return await asyncio.gather(*tasks)


def aggregate_results(
    results: List[FanoutResult],
    aggregator_model: ChatBedrock,
    aggregation_strategy: str = "synthesize",
) -> str:
    """
    Merges parallel agent outputs.

    aggregation_strategy options:
    - "synthesize": ask a model to merge findings coherently
    - "concat": simple concatenation (fast, no model call needed)
    - "vote": majority-vote for classification tasks
    """
    if aggregation_strategy == "concat":
        successful = [r for r in results if r["success"]]
        return "\n\n---\n\n".join(r["output"] for r in successful)

    failed = [r for r in results if not r["success"]]
    successful = [r for r in results if r["success"]]

    if failed:
        # Log partial failures but don't crash — partial results are usually useful
        for f in failed:
            print(f"Agent {f['agent_id']} failed: {f['error']}")

    outputs_for_synthesis = "\n\n".join([
        f"[Agent {r['agent_id']}]:\n{r['output']}"
        for r in successful
    ])

    response = aggregator_model.invoke([
        SystemMessage(content="""
You are a results aggregator. You receive outputs from multiple parallel agents
that each analyzed a different slice of the same problem. Your job is to:
1. Identify common findings across agents
2. Surface unique findings from individual agents
3. Flag any contradictions between agents and explain which to trust
4. Produce a single coherent output as if one expert had analyzed everything

Do not simply concatenate. Actively synthesize.
"""),
        HumanMessage(content=f"Synthesize these {len(successful)} agent outputs:\n\n{outputs_for_synthesis}")
    ])
    return response.content


# Convenience wrapper for synchronous callers
def run_parallel_analysis(
    task_slices: List[str],
    system_prompt: str,
    region: str = "us-east-1",
) -> str:
    results = asyncio.run(fan_out(task_slices, system_prompt, region))
    bedrock = boto3.client("bedrock-runtime", region_name=region)
    aggregator = ChatBedrock(
        client=bedrock,
        model_id="anthropic.claude-3-7-sonnet-20250219-v1:0",
        model_kwargs={"temperature": 0, "max_tokens": 8000}
    )
    return aggregate_results(results, aggregator)

Шаблон 4: Дебаты/критика

При работе с этапом дебатов и критики по шаблону 4 сначала запишите условия работы: необходимые входные данные, сигнал успешного выполнения и действия при частичной неудаче. Такой список помогает сохранять честность при последующих изменениях кода. Храните конфигурацию вне кода приложения. Файлы с настройками, хранилища секретов и флаги функций должны находиться в одном месте, чтобы операторы могли их проверять, не читая весь код. Создавайте контрольные точки после дорогостоящих операций. Система возобновления работы не должна снова взимать плату за один и тот же вызов большой языковой модели, когда оператор пытается выполнить следующий шаг.

  ┌─────────────┐          ┌─────────────┐
  │  Proposer   │          │  Challenger │
  │  (solution  │          │  (solution  │
  │   attempt 1)│          │   attempt 2)│
  └──────┬──────┘          └──────┬──────┘
         │                        │
         └──────────┬─────────────┘
                    ▼
          ┌─────────────────┐
          │    ADJUDICATOR  │
          │  (critique +    │
          │   synthesis)    │
          └─────────────────┘
# harness/multi_agent/debate.py
from dataclasses import dataclass
from typing import Optional
from langchain_aws import ChatBedrock
from langchain_core.messages import SystemMessage, HumanMessage
import boto3


@dataclass
class DebateResult:
    proposer_solution: str
    challenger_solution: str
    adjudication: str
    final_recommendation: str
    agreement_level: str        # HIGH | MEDIUM | LOW | CONTRADICTION


PROPOSER_PROMPT = """
You are a solution proposer. Approach the problem carefully and produce your best
solution. Explain your reasoning. Do not hedge excessively — commit to a specific answer.
"""

CHALLENGER_PROMPT = """
You are a solution challenger. You will receive a problem that another agent has
already attempted. Produce your own independent solution WITHOUT seeing their work.
Approach this fresh. Your goal is not to contradict — it's to find the best solution.
"""

ADJUDICATOR_PROMPT = """
You are an adjudicator reviewing two independent solutions to the same problem.
Your job:
1. Identify where the two solutions agree (these are likely correct)
2. Identify where they diverge (these need careful evaluation)
3. For each divergence, evaluate which solution is better and why
4. Produce a final synthesis that takes the best of both

Be direct about contradictions. Do not smooth over genuine disagreements —
surface them clearly so the human reviewer can make a judgment call.

Rate the agreement level: HIGH (minor differences), MEDIUM (some significant
divergences), LOW (fundamentally different approaches), or CONTRADICTION
(mutually exclusive conclusions).
"""


def run_debate(
    problem: str,
    region: str = "us-east-1",
) -> DebateResult:
    bedrock = boto3.client("bedrock-runtime", region_name=region)

    heavy_model = ChatBedrock(
        client=bedrock,
        model_id="anthropic.claude-3-7-sonnet-20250219-v1:0",
        model_kwargs={
            "temperature": 0.3,        # slight temperature for independent solutions
            "max_tokens": 6000,
            "thinking": {"type": "enabled", "budget_tokens": 4000}
        }
    )

    # Proposer works the problem
    proposer_response = heavy_model.invoke([
        SystemMessage(content=PROPOSER_PROMPT),
        HumanMessage(content=problem)
    ])
    proposer_solution = proposer_response.content

    # Challenger works the same problem independently
    # Note: Challenger does NOT see Proposer's solution
    challenger_response = heavy_model.invoke([
        SystemMessage(content=CHALLENGER_PROMPT),
        HumanMessage(content=problem)
    ])
    challenger_solution = challenger_response.content

    # Adjudicator sees both and synthesizes
    adjudicator_response = heavy_model.invoke([
        SystemMessage(content=ADJUDICATOR_PROMPT),
        HumanMessage(content=f"""
Problem: {problem}

Solution A (Proposer):
{proposer_solution}

Solution B (Challenger):
{challenger_solution}

Adjudicate and synthesize.
""")
    ])

    adjudication = adjudicator_response.content

    # Extract agreement level from adjudication
    agreement_level = "MEDIUM"
    for level in ["CONTRADICTION", "LOW", "HIGH", "MEDIUM"]:
        if level in adjudication.upper():
            agreement_level = level
            break

    return DebateResult(
        proposer_solution=proposer_solution,
        challenger_solution=challenger_solution,
        adjudication=adjudication,
        final_recommendation=adjudication,
        agreement_level=agreement_level,
    )

Передача данных между агентами и передача контекста

При работе над этапами передачи данных между агентами и управления контекстом сначала запишите условия работы: необходимые входные данные, сигнал о успешном выполнении и действия при частичной неудаче. Такой чек-лист поможет сохранять честность при последующих изменениях кода. Документируйте одновременно успешный и восстановительный сценарии работы. Повторные попытки, проверки человеком и обработка неработоспособных сообщений являются частью продукта, а не элементами последующей доработки. Выполняйте контрольные точки после дорогостоящих операций. Механизм возобновления работы не должен повторно взимать плату за один и тот же вызов большой языковой модели, когда оператор пытается выполнить последующий шаг.

# harness/multi_agent/handoff.py
from dataclasses import dataclass, field
from typing import Any, Optional
from langchain_core.messages import HumanMessage


@dataclass
class AgentHandoff:
    """
    Structured context passed between agents.

    The split between result and trace is deliberate: downstream agents
    need the result, not the full reasoning history. Keeping them separate
    lets each agent decide how much context it wants to consume.
    """
    source_agent: str
    task_completed: str
    result_summary: str              # compressed, structured result
    artifacts: dict = field(default_factory=dict)  # files, code, structured data
    reasoning_trace: Optional[str] = None          # full trace, passed only if downstream needs it
    confidence: str = "MEDIUM"       # HIGH | MEDIUM | LOW
    flags: list = field(default_factory=list)      # NEEDS_REVIEW, PARTIAL_RESULT, etc.

    def to_context_message(self, include_trace: bool = False) -> HumanMessage:
        """
        Converts handoff to a HumanMessage for injection into next agent's context.
        include_trace=True only when the downstream agent genuinely needs the reasoning.
        """
        content = f"""
[HANDOFF FROM: {self.source_agent}]
Task completed: {self.task_completed}
Confidence: {self.confidence}
Flags: {', '.join(self.flags) if self.flags else 'none'}

Result summary:
{self.result_summary}
"""
        if self.artifacts:
            content += f"\nArtifacts available:\n"
            for key, value in self.artifacts.items():
                if isinstance(value, str) and len(value) < 500:
                    content += f"  {key}: {value}\n"
                else:
                    content += f"  {key}: [available, {type(value).__name__}]\n"

        if include_trace and self.reasoning_trace:
            content += f"\nFull reasoning trace:\n{self.reasoning_trace}"

        return HumanMessage(content=content)

Усиление надежности в производственной среде

При работе над этапом укрепления системы для производственного использования сначала запишите условия работы: необходимые входные данные, сигнал о успешном выполнении и действия при частичной неудаче. Такой список помогает сохранять честность при последующих изменениях кода. Предпочитайте небольшие, тестируемые модули вместо обширных скриптов. Если какой-то шаг не сработает, причина неудачи должна указывать на конкретную ответственность, а не на запутанную структуру обработки данных. Выполняйте проверки после дорогостоящих шагов. Система возобновления работы не должна снова взимать плату за один и тот же вызов большой языковой модели, когда оператор пытается выполнить следующий этап. При работе над этапом укрепления системы для производственного использования сначала запишите условия работы: необходимые входные данные, сигнал о успешном выполнении и действия при частичной неудаче. Такой список помогает сохранять честность при последующих изменениях кода. Записывайте время выполнения и стоимость токенов или запросов рядом с функциональными результатами. Отслеживание затрат на раннем этапе предотвращает неожиданные счета при переходе от демо-среды к общедоступным средам.

Предотвращение резкого увеличения затрат

Этап предотвращения взрыва затрат наиболее эффективен, когда рассматривается как измеримая структура. Соберите один идеальный пример работы, один случай сбоя и запись о возврате к предыдущему состоянию перед расширением объёма работы. Храните конфигурацию вне кода приложения. Файлы среды, хранилища секретов и флаги функций должны находиться в одном месте, чтобы операторы могли их проверять, не читая всю структуру. Сохраняйте состояние структуры простым и типизированным. Вложенные объекты скрывают информацию о том, какой узел заполнил тот или иной поле, и мешают возобновлению работы после прерываний.

# harness/multi_agent/budget.py
import boto3
import time
from decimal import Decimal


class AgentBudgetGuard:
    """
    Tracks cumulative cost and agent count per run.
    Hard-stops execution when limits are exceeded.
    """

    def __init__(
        self,
        max_agents_per_run: int = 10,
        max_total_tokens: int = 500_000,
        table_name: str = "agent-budget-state",
        region: str = "us-east-1",
    ):
        self.max_agents = max_agents_per_run
        self.max_tokens = max_total_tokens
        self.table = boto3.resource("dynamodb", region_name=region).Table(table_name)

    def register_agent_spawn(self, run_id: str, agent_id: str) -> bool:
        """
        Returns True if spawn is allowed, False if budget exceeded.
        Call this before spawning any sub-agent.
        """
        response = self.table.update_item(
            Key={"run_id": run_id},
            UpdateExpression="SET agent_count = if_not_exists(agent_count, :z) + :inc",
            ExpressionAttributeValues={":z": 0, ":inc": 1},
            ReturnValues="UPDATED_NEW",
        )
        new_count = int(response["Attributes"]["agent_count"])

        if new_count > self.max_agents:
            raise AgentBudgetExceededError(
                f"Run {run_id} attempted to spawn agent #{new_count}, "
                f"but max_agents_per_run is {self.max_agents}. "
                f"Either the supervisor is over-decomposing, or there is a spawn loop."
            )
        return True

    def record_token_usage(self, run_id: str, tokens_used: int):
        response = self.table.update_item(
            Key={"run_id": run_id},
            UpdateExpression="SET total_tokens = if_not_exists(total_tokens, :z) + :inc",
            ExpressionAttributeValues={":z": 0, ":inc": tokens_used},
            ReturnValues="UPDATED_NEW",
        )
        total = int(response["Attributes"]["total_tokens"])
        if total > self.max_tokens:
            raise AgentBudgetExceededError(
                f"Run {run_id} consumed {total:,} tokens, exceeding limit of {self.max_tokens:,}."
            )


class AgentBudgetExceededError(Exception):
    pass

Обнаружение заедания

Этап обнаружения заедания работает наилучшим образом, когда его рассматривают как измеримую структуру. Соберите один идеальный пример работы, один случай сбоя и записку о возврате к предыдущему состоянию перед расширением объёма исследования. Документируйте одновременно успешный путь выполнения и путь восстановления. Повторные попытки, проверки человеком и обработка неработоспособных сообщений являются частью продукта, а не этапом последующей доработки. Сохраняйте структуру графа простой и типизированной; вложенные структуры скрывают информацию о том, какой узел заполнил тот или иной поле, и мешают возобновлению работы после прерываний.

# harness/multi_agent/deadlock.py
import boto3
import time
from typing import List


class DeadlockDetector:
    """
    Tracks the delegation chain per run and detects cycles.
    Stored in DynamoDB so it works across parallel agent branches.
    """

    def __init__(self, table_name: str = "agent-delegation-chain", region: str = "us-east-1"):
        self.table = boto3.resource("dynamodb", region_name=region).Table(table_name)

    def record_delegation(self, run_id: str, from_agent: str, to_agent: str):
        """
        Records a delegation event and checks for cycles.
        Raises DeadlockDetectedError if a cycle is found.
        """
        self.table.put_item(Item={
            "run_id": run_id,
            "delegation_id": f"{from_agent}->{to_agent}-{int(time.time())}",
            "from_agent": from_agent,
            "to_agent": to_agent,
            "timestamp": int(time.time()),
        })

        chain = self._get_delegation_chain(run_id)
        if self._has_cycle(chain):
            cycle_description = self._describe_cycle(chain)
            raise DeadlockDetectedError(
                f"Delegation cycle detected in run {run_id}: {cycle_description}. "
                f"Check supervisor decomposition logic for circular dependencies."
            )

    def _get_delegation_chain(self, run_id: str) -> List[tuple]:
        response = self.table.query(
            KeyConditionExpression="run_id = :rid",
            ExpressionAttributeValues={":rid": run_id}
        )
        return [(item["from_agent"], item["to_agent"]) for item in response.get("Items", [])]

    def _has_cycle(self, chain: List[tuple]) -> bool:
        graph = {}
        for from_a, to_a in chain:
            graph.setdefault(from_a, set()).add(to_a)

        visited, rec_stack = set(), set()

        def dfs(node):
            visited.add(node)
            rec_stack.add(node)
            for neighbor in graph.get(node, []):
                if neighbor not in visited:
                    if dfs(neighbor): return True
                elif neighbor in rec_stack:
                    return True
            rec_stack.discard(node)
            return False

        return any(dfs(node) for node in graph if node not in visited)

    def _describe_cycle(self, chain: List[tuple]) -> str:
        return " -> ".join(f"{f}->{t}" for f, t in chain[-5:])


class DeadlockDetectedError(Exception):
    pass

Наблюдаемость между агентами

Этап наблюдаемости межагентных систем работает наилучшим образом, когда рассматривается как измеримая поверхность. Соберите один эталонный протокол, один пример сбоя и записку о возврате к предыдущему состоянию перед расширением объема работ. Предпочитайте небольшие, тестируемые единицы кода вместо обширных скриптов. При сбое какого-либо шага причина должна быть связана с конкретной функцией, а не с запутанной цепочкой операций. Сохраняйте структуру графа простой и типизированной. Вложенные структуры данных маскируют информацию о том, какой узел заполнил тот или иной поле, что приводит к нарушению возобновления работы после перерывов. Этап наблюдаемости межагентных систем работает наилучшим образом, когда рассматривается как измеримая поверхность. Соберите один эталонный протокол, один пример сбоя и записку о возврате к предыдущему состоянию перед расширением объема работ. Записывайте временные показатели, а также стоимость токенов или запросов рядом с функциональными результатами. Отслеживание затрат на раннем этапе предотвращает неожиданные расходы при переходе с демо-среды в общедоступные среды.

# harness/multi_agent/observability.py
import os
from contextlib import contextmanager
from langsmith import Client
from langsmith.run_trees import RunTree


class MultiAgentTracer:
    """
    Maintains a run tree across all agents in a multi-agent system.
    Pass the parent_run_id to each agent so their traces nest correctly.
    """

    def __init__(self):
        self.client = Client()

    @contextmanager
    def agent_span(self, parent_run_id: str, agent_name: str, inputs: dict):
        """
        Context manager for an individual agent's trace span.

        Usage:
            with tracer.agent_span(parent_run_id, "security_reviewer", {...}) as span:
                result = run_security_review(...)
                span.end(outputs={"result": result})
        """
        run = self.client.create_run(
            name=agent_name,
            run_type="chain",
            inputs=inputs,
            parent_run_id=parent_run_id,
        )
        try:
            yield run
        except Exception as e:
            self.client.update_run(run.id, error=str(e))
            raise
        finally:
            self.client.update_run(run.id, end_time=None)  # auto-sets end time

Проверка реальности в производственных условиях

На этапе проверки реальности в производственной среде необходимо заранее определить входные данные, ответственного за выполнение шага и критерии завершения перед внесением изменений в код. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии системы. Конфигурацию следует хранить отдельно от кода приложения. Файлы с настройками окружения, хранилища конфиденциальных данных и флаги функций должны находиться в одном месте, чтобы операторы могли их проверять, не читая весь кодовый базис. Для операций, связанных с тратой денег или изменением данных в производстве, необходимо предусмотреть утверждение человеком. Подключения, созданные во время компиляции, не гарантируют полноты обработки бизнес-логики.

Архитектура для справки

На этапе архитектуры для справки необходимо определить входные данные, ответственного за выполнение шага и критерии завершения перед изменением кода. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Необходимо задокументировать как успешный, так и восстановительный сценарии работы. Повторные попытки, проверки человеком и обработка неработоспособных сообщений являются частью продукта, а не элементами последующей доработки. Внедрять утверждение человеком для операций, связанных с тратой денег или изменением производственных данных. Подключение компонентов во время компиляции не гарантирует полноты функционала продукта.

                  ┌───────────────────────────────────────────┐
                  │            User Request                   │
                  └──────────────────┬────────────────────────┘
                                     │
                  ┌──────────────────▼────────────────────────┐
                  │           AgentHarness Runtime            │
                  │  (budget guard, deadlock detector, tracer)│
                  └──────────────────┬────────────────────────┘
                                     │
                  ┌──────────────────▼────────────────────────┐
                  │         SUPERVISOR / ORCHESTRATOR         │
                  │   Claude 3.7 + extended thinking          │
                  │   Task decomposition + result synthesis   │
                  └────┬──────────────┬──────────────┬────────┘
                       │              │              │
           ┌───────────▼──┐  ┌────────▼──┐  ┌───────▼───────┐
           │   Worker A   │  │  Worker B  │  │   Worker C   │
           │  Claude 3.5  │  │ Claude 3.5 │  │  Claude 3.5  │
           │  specialist  │  │ specialist │  │  specialist  │
           └───────┬──────┘  └─────┬─────┘  └──────┬────────┘
                   │               │                │
           ┌───────▼───────────────▼────────────────▼────────┐
           │              Tool Execution Layer               │
           │     (auth, retry, circuit breaker from Art.5)   │
           └───────────────────────┬─────────────────────────┘
                                   │
           ┌───────────────────────▼──────────────────────────┐
           │                  AWS Services                    │
           │  Bedrock │ DynamoDB (budget+deadlock+loop+circ)  │
           │  Secrets Manager │ Knowledge Bases               │
           └──────────────────────────────────────────────────┘

  Observability: LangSmith run trees — full hierarchy per user request

Стек справочной инфраструктуры

На этапе инфраструктуры для справочных целей необходимо заранее определить входные данные, ответственного за выполнение шага и критерии завершения перед внесением изменений в код. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не пытаясь угадать скрытое состояние системы. Лучше использовать небольшие, тестируемые единицы кода вместо обширных скриптов. При сбое шага он должен указывать на конкретную причину, а не на сложную структуру обработки данных. Внедрять человеческое утверждение для операций, связанных с расходованием средств или изменением производственных данных. Простая настройка во время компиляции не гарантирует полноты реализации бизнес-логики.

+-----------------------------+---------------------+------------------------------+
| Component                   | Technology          | Role                         |
+-----------------------------+---------------------+------------------------------+
| Orchestration               | LangGraph 0.2+      | Multi-agent graph, routing,  |
|                             |                     | supervisor-worker topology   |
+-----------------------------+---------------------+------------------------------+
| Supervisor Model            | Claude 3.7 Sonnet   | Task decomposition,          |
|                             | (extended thinking) | result synthesis             |
+-----------------------------+---------------------+------------------------------+
| Worker Models               | Claude 3.5 Sonnet   | Specialized execution,       |
|                             |                     | bounded tasks                |
+-----------------------------+---------------------+------------------------------+
| Parallel Execution          | asyncio + Bedrock   | Concurrent agent runs with   |
|                             |                     | semaphore-gated concurrency  |
+-----------------------------+---------------------+------------------------------+
| Context Compression         | Claude Haiku 3.5    | Pipeline stage handoffs,     |
|                             |                     | summary generation           |
+-----------------------------+---------------------+------------------------------+
| Budget Guard                | DynamoDB            | Agent count + token limits   |
|                             |                     | per run                      |
+-----------------------------+---------------------+------------------------------+
| Deadlock Detection          | DynamoDB            | Delegation cycle detection   |
+-----------------------------+---------------------+------------------------------+
| Loop Detection              | DynamoDB (Art. 5)   | Per-resource edit tracking   |
+-----------------------------+---------------------+------------------------------+
| Circuit Breaker State       | DynamoDB (Art. 5)   | Shared across all agents     |
|                             |                     | in a run                     |
+-----------------------------+---------------------+------------------------------+
| Cross-Agent Observability   | LangSmith run trees | Full hierarchy per request   |
+-----------------------------+---------------------+------------------------------+
| Auth Propagation            | CredentialManager   | JWT passed to all workers    |
|                             | (Art. 5)            | via execution context        |
+-----------------------------+---------------------+------------------------------+
| Local Dev Alternative       | Ollama + Docker     | All patterns testable        |
|                             | Compose             | without Bedrock costs        |
+-----------------------------+---------------------+------------------------------+
| Infrastructure as Code      | Terraform           | DynamoDB tables, IAM roles   |
+-----------------------------+---------------------+------------------------------+

Краткое описание будущих шагов

На этапе «Примечание к A» необходимо определить входные данные, ответственного за выполнение шага и критерии завершения перед изменением кода. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии. Рассматривайте этот этап как контракт между входными данными и проверенными выходными результатами. Укажите названия результатов работы, определите критерии успеха и не допускайте молчаливого частичного завершения задачи. Внедряйте утверждение человека для операций, связанных с тратой денег или изменением производственных данных. Подключение компонентов во время компиляции не гарантирует полноты выполнения бизнес-задач.

Чек-лист операций

На этапе «Чек-лист операций» необходимо определить входные данные, ответственного за выполнение шага и критерии завершения перед изменением кода. Операторы должны иметь возможность перезапустить шаг с известной точки контроля, не догадываясь о скрытом состоянии.

Рассматривайте этот этап как контракт между входными данными и проверенными выходными результатами. Дайте названия элементам, определите критерии успешности и не соглашайтесь на молчаливое частичное выполнение задачи.

Обеспечьте человеческое утверждение для операций, связанных с тратой денег или изменением производственных данных. Подключение на этапе компиляции не гарантирует полноты решения в бизнес-среде.

Напишите краткий руководство: как обновлять ключи, как опустошать очередь, как откатить последнюю загрузку данных.

Записывайте время выполнения операций, а также стоимость токенов или запросов рядом с функциональными результатами. Ясность стоимости заранее предотвращает неожиданные счета при переходе от демо-среды к общедоступным средам.

Обеспечьте человеческое утверждение для операций, связанных с тратой денег или изменением производственных данных. Подключение на этапе компиляции не гарантирует полноты решения в бизнес-среде.

Перед внедрением стека заморозьте версии, сделайте копию «золотого» отчета для критической цепочки задач и уточните шаги возврата к предыдущему состоянию. В совместных средах необходимы ограничения на частоту запросов, проверки принадлежности пользователя и четко определенный ответственный за обновление секретов. Лучше выбирать надежность, чем креативные одноразовые демонстрации.

Примечание для a0dc7ff1211b: не храните ключи поставщика в репозитории, установите лимит токенов на сессию и сохраняйте отчеты рядом с фикстурами для оценки, чтобы последующие замены моделей оставались сопоставимыми.