Продвинутый уровень

Фоновый AI-агент с удалённым MCP на Gemini API

Время чтения и практики: 60 минут Результат: фоновая задача, статус и удалённые инструменты

Долгая агентная задача не должна жить внутри одного HTTP-запроса. Клиент может закрыть вкладку, прокси — оборвать соединение, а модель — запросить несколько инструментов подряд. Надёжная схема принимает работу за миллисекунды, сохраняет её в очереди, возвращает идентификатор, выполняет цикл модели в отдельном процессе и подключает внутренние системы через удалённый MCP-сервер.

Что мы построим

В этой лабораторной работе мы создадим AI-агента, который использует Gemini API для планирования и генерации ответа, а внутренний инструмент получает через удалённый Model Context Protocol, или MCP. Пользователь запускает работу через короткий HTTP-запрос и затем проверяет состояние отдельным запросом.

Готовая система будет состоять из четырёх частей:

  1. API принимает описание задачи, сохраняет её и немедленно отвечает кодом 202 Accepted.
  2. Фоновый worker забирает задачу, обращается к Gemini и выполняет цикл вызовов инструментов.
  3. Удалённый MCP-сервер публикует разрешённые операции над внутренними данными.
  4. Клиент получает состояние queued, running, succeeded или failed через отдельный endpoint.
Клиент
  │
  │ POST /v1/jobs
  ▼
API ───────► SQLite: queued
  │              │
  │ 202 + job_id │ claim
  ▼              ▼
Клиент        Worker ─────► Gemini API
  │              ▲             │
  │ GET status   │ tool result │ function call
  ▼              │             ▼
API ◄──────── SQLite       удалённый MCP
                                  │
                                  ▼
                         внутренний источник

Для демонстрации очередь хранится в SQLite. Это позволяет повторить работу без Redis, брокера сообщений и облачной инфраструктуры. В разделе об ограничениях разберём, когда SQLite следует заменить.

Конкретный кейс: анализ инцидента

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

Внутренний каталог инструкций нельзя напрямую открыть из публичного API. Поэтому мы публикуем узкий MCP-инструмент lookup_runbook. Он принимает имя сервиса и возвращает только заранее разрешённые поля. В учебном примере данные статичны, но граница интеграции остаётся такой же для базы знаний, CMDB, Git-репозитория или внутреннего REST API.

Такой кейс показывает две независимые проблемы:

  • выполнение может занять дольше тайм-аута браузера, балансировщика или serverless-платформы;
  • доступ к внутренним данным требует отдельной доверенной точки, а не прямого сетевого доступа со стороны модели.

Почему обычный синхронный endpoint ненадёжен

Наивная реализация вызывает Gemini прямо внутри обработчика POST /analyze, затем там же обращается к инструментам и удерживает соединение до готового ответа. Она работает на коротких демонстрациях, но плохо переносит реальную эксплуатацию.

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

Дополнительные источники сбоев:

  • тайм-аут reverse proxy короче времени выполнения агента;
  • перезапуск API-процесса уничтожает задачу, хранившуюся только в памяти;
  • повторный вызов модели после сетевой ошибки расходует токены и может дать другой результат;
  • инструмент отвечает медленно или возвращает слишком большой документ;
  • worker завершается после получения задания, но до записи результата;
  • модель циклически вызывает один инструмент с одинаковыми аргументами.

Правильная граница HTTP-запроса заканчивается сразу после надёжной постановки работы в очередь. Фактическое выполнение получает собственный жизненный цикл.

Контракт фоновой задачи

Создание задачи возвращает идентификатор и URL состояния:

POST /v1/jobs
Idempotency-Key: incident-2026-07-26-billing-503
Content-Type: application/json

{
  "prompt": "После релиза у billing выросло число 503. Подготовь план диагностики."
}
HTTP/1.1 202 Accepted
Location: /v1/jobs/01K123...

{
  "id": "01K123...",
  "status": "queued",
  "status_url": "/v1/jobs/01K123..."
}

Endpoint состояния возвращает одну и ту же структуру на всех стадиях:

{
  "id": "01K123...",
  "status": "running",
  "created_at": "2026-07-26T10:00:00+00:00",
  "started_at": "2026-07-26T10:00:01+00:00",
  "finished_at": null,
  "result": null,
  "error": null,
  "attempt": 1
}

Клиент использует polling с увеличивающимся интервалом. Для интерфейса с большим количеством обновлений можно добавить Server-Sent Events или webhook, но состояние в базе всё равно остаётся источником истины.

Подготовка проекта

Нужны Python 3.11 или новее, ключ Gemini API и возможность запустить два процесса: приложение с API и worker. Удалённый MCP-сервер запускается третьим процессом или разворачивается на отдельном внутреннем хосте.

mkdir managed-agent
cd managed-agent
python -m venv .venv
source .venv/bin/activate
pip install fastapi "uvicorn[standard]" google-genai mcp httpx pydantic-settings

Структура проекта:

managed-agent/
├── app/
│   ├── __init__.py
│   ├── api.py
│   ├── agent.py
│   ├── config.py
│   ├── db.py
│   └── worker.py
├── mcp_server/
│   └── server.py
└── .env

Создайте конфигурацию окружения. Настоящие значения не включайте в репозиторий:

GEMINI_API_KEY=replace-me
GEMINI_MODEL=gemini-2.5-flash
MCP_URL=http://127.0.0.1:8001/mcp
MCP_BEARER_TOKEN=
DATABASE_PATH=./agent_jobs.sqlite3
MAX_AGENT_STEPS=8
TOOL_TIMEOUT_SECONDS=20
JOB_LEASE_SECONDS=120

Название модели вынесено в переменную, потому что доступность моделей зависит от аккаунта и может меняться. Перед запуском укажите модель Gemini, доступную вашему проекту. Секрет передаётся только через окружение.

Шаг 1. Конфигурация

Файл app/config.py:

from functools import lru_cache

from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    gemini_api_key: str
    gemini_model: str = "gemini-2.5-flash"
    mcp_url: str = "http://127.0.0.1:8001/mcp"
    mcp_bearer_token: str = ""
    database_path: str = "./agent_jobs.sqlite3"
    max_agent_steps: int = 8
    tool_timeout_seconds: float = 20.0
    job_lease_seconds: int = 120

    model_config = SettingsConfigDict(
        env_file=".env",
        env_file_encoding="utf-8",
        extra="ignore",
    )


@lru_cache
def get_settings() -> Settings:
    return Settings()

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

Шаг 2. Персистентная очередь в SQLite

Нам нужны атомарное получение задания, сохранение результата и lease — срок владения работой. Если worker погибнет, задача не должна навсегда остаться в состоянии running.

Файл app/db.py:

import json
import sqlite3
import uuid
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone
from typing import Any

from app.config import get_settings


def utc_now() -> datetime:
    return datetime.now(timezone.utc)


def iso(value: datetime | None = None) -> str:
    return (value or utc_now()).isoformat()


@contextmanager
def connect():
    settings = get_settings()
    connection = sqlite3.connect(
        settings.database_path,
        timeout=10,
        isolation_level=None,
    )
    connection.row_factory = sqlite3.Row
    connection.execute("PRAGMA journal_mode=WAL")
    connection.execute("PRAGMA foreign_keys=ON")
    try:
        yield connection
    finally:
        connection.close()


def init_db() -> None:
    with connect() as db:
        db.executescript("""
        CREATE TABLE IF NOT EXISTS jobs (
            id TEXT PRIMARY KEY,
            idempotency_key TEXT UNIQUE,
            prompt TEXT NOT NULL,
            status TEXT NOT NULL
                CHECK (status IN ('queued', 'running', 'succeeded', 'failed')),
            result_json TEXT,
            error_json TEXT,
            attempt INTEGER NOT NULL DEFAULT 0,
            created_at TEXT NOT NULL,
            started_at TEXT,
            heartbeat_at TEXT,
            lease_expires_at TEXT,
            finished_at TEXT
        );

        CREATE INDEX IF NOT EXISTS jobs_claim_idx
        ON jobs(status, created_at);
        """)


def row_to_job(row: sqlite3.Row | None) -> dict[str, Any] | None:
    if row is None:
        return None

    job = dict(row)
    job["result"] = (
        json.loads(job.pop("result_json"))
        if job["result_json"] is not None
        else None
    )
    job["error"] = (
        json.loads(job.pop("error_json"))
        if job["error_json"] is not None
        else None
    )
    job.pop("prompt", None)
    job.pop("heartbeat_at", None)
    job.pop("lease_expires_at", None)
    job.pop("idempotency_key", None)
    return job


def create_job(prompt: str, idempotency_key: str | None) -> dict[str, Any]:
    job_id = uuid.uuid4().hex
    created_at = iso()

    with connect() as db:
        db.execute("BEGIN IMMEDIATE")
        try:
            if idempotency_key:
                existing = db.execute(
                    "SELECT * FROM jobs WHERE idempotency_key = ?",
                    (idempotency_key,),
                ).fetchone()
                if existing:
                    db.execute("COMMIT")
                    return row_to_job(existing)

            db.execute(
                """
                INSERT INTO jobs (
                    id, idempotency_key, prompt, status, created_at
                ) VALUES (?, ?, ?, 'queued', ?)
                """,
                (job_id, idempotency_key, prompt, created_at),
            )
            row = db.execute(
                "SELECT * FROM jobs WHERE id = ?",
                (job_id,),
            ).fetchone()
            db.execute("COMMIT")
            return row_to_job(row)
        except Exception:
            db.execute("ROLLBACK")
            raise


def get_job(job_id: str) -> dict[str, Any] | None:
    with connect() as db:
        row = db.execute(
            "SELECT * FROM jobs WHERE id = ?",
            (job_id,),
        ).fetchone()
        return row_to_job(row)


def claim_job() -> dict[str, Any] | None:
    settings = get_settings()
    now = utc_now()
    lease_until = now + timedelta(seconds=settings.job_lease_seconds)

    with connect() as db:
        db.execute("BEGIN IMMEDIATE")
        try:
            row = db.execute(
                """
                SELECT *
                FROM jobs
                WHERE status = 'queued'
                   OR (
                       status = 'running'
                       AND lease_expires_at IS NOT NULL
                       AND lease_expires_at < ?
                   )
                ORDER BY created_at
                LIMIT 1
                """,
                (iso(now),),
            ).fetchone()

            if row is None:
                db.execute("COMMIT")
                return None

            started_at = row["started_at"] or iso(now)
            db.execute(
                """
                UPDATE jobs
                SET status = 'running',
                    started_at = ?,
                    heartbeat_at = ?,
                    lease_expires_at = ?,
                    attempt = attempt + 1,
                    error_json = NULL
                WHERE id = ?
                """,
                (started_at, iso(now), iso(lease_until), row["id"]),
            )
            claimed = db.execute(
                "SELECT * FROM jobs WHERE id = ?",
                (row["id"],),
            ).fetchone()
            db.execute("COMMIT")
            return dict(claimed)
        except Exception:
            db.execute("ROLLBACK")
            raise


def heartbeat(job_id: str) -> None:
    settings = get_settings()
    now = utc_now()
    lease_until = now + timedelta(seconds=settings.job_lease_seconds)

    with connect() as db:
        db.execute(
            """
            UPDATE jobs
            SET heartbeat_at = ?, lease_expires_at = ?
            WHERE id = ? AND status = 'running'
            """,
            (iso(now), iso(lease_until), job_id),
        )


def complete_job(job_id: str, result: dict[str, Any]) -> None:
    with connect() as db:
        db.execute(
            """
            UPDATE jobs
            SET status = 'succeeded',
                result_json = ?,
                finished_at = ?,
                heartbeat_at = NULL,
                lease_expires_at = NULL
            WHERE id = ? AND status = 'running'
            """,
            (
                json.dumps(result, ensure_ascii=False),
                iso(),
                job_id,
            ),
        )


def fail_job(job_id: str, error: dict[str, Any]) -> None:
    with connect() as db:
        db.execute(
            """
            UPDATE jobs
            SET status = 'failed',
                error_json = ?,
                finished_at = ?,
                heartbeat_at = NULL,
                lease_expires_at = NULL
            WHERE id = ? AND status = 'running'
            """,
            (
                json.dumps(error, ensure_ascii=False),
                iso(),
                job_id,
            ),
        )

BEGIN IMMEDIATE сериализует операцию выбора и захвата задания. Два worker-процесса не смогут одновременно забрать одну строку. Lease восстанавливает потерянную работу, но создаёт важное следствие: выполнение в общем случае имеет семантику «как минимум один раз». Поэтому изменяющие инструменты обязаны принимать собственный ключ идемпотентности.

Шаг 3. API постановки задачи и получения статуса

Файл app/api.py:

from contextlib import asynccontextmanager
from typing import Literal

from fastapi import FastAPI, Header, HTTPException, Response, status
from pydantic import BaseModel, Field

from app.db import create_job, get_job, init_db


class CreateJobRequest(BaseModel):
    prompt: str = Field(min_length=10, max_length=20_000)


class JobCreated(BaseModel):
    id: str
    status: Literal["queued", "running", "succeeded", "failed"]
    status_url: str


class JobStatus(BaseModel):
    id: str
    status: Literal["queued", "running", "succeeded", "failed"]
    result: dict | None = None
    error: dict | None = None
    attempt: int
    created_at: str
    started_at: str | None = None
    finished_at: str | None = None


@asynccontextmanager
async def lifespan(_: FastAPI):
    init_db()
    yield


app = FastAPI(
    title="Managed Background Agent",
    version="1.0.0",
    lifespan=lifespan,
)


@app.post(
    "/v1/jobs",
    response_model=JobCreated,
    status_code=status.HTTP_202_ACCEPTED,
)
def submit_job(
    body: CreateJobRequest,
    response: Response,
    idempotency_key: str | None = Header(
        default=None,
        alias="Idempotency-Key",
        max_length=200,
    ),
):
    job = create_job(body.prompt, idempotency_key)
    status_url = f"/v1/jobs/{job['id']}"
    response.headers["Location"] = status_url
    response.headers["Retry-After"] = "2"
    return {
        "id": job["id"],
        "status": job["status"],
        "status_url": status_url,
    }


@app.get("/v1/jobs/{job_id}", response_model=JobStatus)
def read_job(job_id: str):
    job = get_job(job_id)
    if job is None:
        raise HTTPException(status_code=404, detail="Job not found")
    return job


@app.get("/health/live")
def live():
    return {"status": "ok"}

Endpoint не запускает asyncio.create_task(). Такая «фоновая» задача принадлежала бы процессу API и исчезла бы при перезапуске. Здесь HTTP-ответ отправляется только после фиксации строки в базе.

Шаг 4. Удалённый MCP-сервер

Вызов инструментов отделён от модели. MCP-сервер объявляет имя операции, описание и входную JSON Schema. Gemini решает, когда вызвать функцию, а наше приложение проверяет запрос и пересылает его MCP-серверу.

Файл mcp_server/server.py:

from mcp.server.fastmcp import FastMCP


mcp = FastMCP(
    "internal-runbooks",
    host="0.0.0.0",
    port=8001,
)


RUNBOOKS = {
    "billing": {
        "service": "billing",
        "owner": "payments-platform",
        "safe_checks": [
            "Сравнить долю 503 до и после релиза",
            "Проверить доступность upstream payment-gateway",
            "Сопоставить ошибки по версии приложения",
            "Проверить насыщение пула соединений",
        ],
        "rollback_condition": (
            "Рост 503 статистически и по времени связан с новой версией, "
            "а upstream остаётся здоровым"
        ),
    },
    "catalog": {
        "service": "catalog",
        "owner": "commerce-platform",
        "safe_checks": [
            "Проверить задержку базы данных",
            "Проверить исчерпание пула соединений",
            "Сравнить ошибки по зонам доступности",
        ],
        "rollback_condition": (
            "Ошибки появились после релиза и локализованы в новой версии"
        ),
    },
}


@mcp.tool()
def lookup_runbook(service: str) -> dict:
    """Возвращает безопасные диагностические шаги для внутреннего сервиса."""
    normalized = service.strip().lower()
    if normalized not in RUNBOOKS:
        return {
            "found": False,
            "service": normalized,
            "message": "Инструкция для сервиса не найдена",
        }

    return {
        "found": True,
        "runbook": RUNBOOKS[normalized],
    }


if __name__ == "__main__":
    mcp.run(transport="streamable-http")

Запустите сервер:

source .venv/bin/activate
python -m mcp_server.server

В зависимости от версии Python SDK точный путь Streamable HTTP может задаваться настройкой сервера. Проверьте фактический endpoint при старте и приведите MCP_URL к нему. Для приведённой конфигурации ожидается http://127.0.0.1:8001/mcp.

Что изменить перед реальным развёртыванием

Учебный сервер слушает локальный адрес и не реализует аутентификацию. Не публикуйте его в интернет в таком виде. В рабочей сети поставьте перед MCP endpoint API gateway или reverse proxy, который:

  • проверяет короткоживущий сервисный токен;
  • завершает TLS;
  • ограничивает размер тела запроса и ответа;
  • разрешает доступ только с подсети worker;
  • пишет аудит имени инструмента, субъекта и результата;
  • не пересылает секреты в прикладные журналы.

Для production предпочтителен OAuth 2.0 с отдельной машинной учётной записью или взаимный TLS. Статический bearer-токен в примере клиента предусмотрен как минимальный механизм для закрытой тестовой среды.

Шаг 5. Агентный цикл Gemini ↔ MCP

Главная часть системы выполняет пять действий:

  1. подключается к MCP-серверу;
  2. получает список доступных инструментов;
  3. преобразует их схемы в объявления функций Gemini;
  4. передаёт вызовы модели MCP-серверу;
  5. возвращает результаты инструментов модели до получения финального текста.

Файл app/agent.py:

import asyncio
import json
from typing import Any

from google import genai
from google.genai import types
from mcp import ClientSession
from mcp.client.streamable_http import streamablehttp_client

from app.config import get_settings


SYSTEM_INSTRUCTION = """
Ты — помощник дежурного инженера.

Правила:
1. Для внутренних инструкций используй доступный инструмент.
2. Не утверждай, что проверка выполнена, если инструмент вернул только инструкцию.
3. Не придумывай метрики, события, версии, владельцев или результаты команд.
4. Сначала предложи безопасные проверки, затем условия эскалации или отката.
5. Если данных недостаточно, явно перечисли, что требуется уточнить.
6. Ответ должен быть на русском языке.
""".strip()


def mcp_tools_to_gemini(mcp_tools: list[Any]) -> list[types.FunctionDeclaration]:
    declarations: list[types.FunctionDeclaration] = []

    for tool in mcp_tools:
        schema = tool.inputSchema or {
            "type": "object",
            "properties": {},
        }
        declarations.append(
            types.FunctionDeclaration(
                name=tool.name,
                description=tool.description or "",
                parameters_json_schema=schema,
            )
        )

    return declarations


def serialize_tool_result(result: Any) -> dict[str, Any]:
    if hasattr(result, "model_dump"):
        return result.model_dump(mode="json")
    return {"content": str(result)}


def final_text(content: types.Content) -> str:
    chunks: list[str] = []
    for part in content.parts or []:
        if getattr(part, "text", None):
            chunks.append(part.text)
    return "\n".join(chunks).strip()


async def execute_agent(prompt: str) -> dict[str, Any]:
    settings = get_settings()
    client = genai.Client(api_key=settings.gemini_api_key)

    headers = {}
    if settings.mcp_bearer_token:
        headers["Authorization"] = (
            f"Bearer {settings.mcp_bearer_token}"
        )

    async with streamablehttp_client(
        settings.mcp_url,
        headers=headers,
    ) as (read_stream, write_stream, _):
        async with ClientSession(
            read_stream,
            write_stream,
        ) as session:
            await session.initialize()

            listed = await asyncio.wait_for(
                session.list_tools(),
                timeout=settings.tool_timeout_seconds,
            )
            declarations = mcp_tools_to_gemini(listed.tools)
            allowed_names = {tool.name for tool in listed.tools}

            contents: list[types.Content] = [
                types.Content(
                    role="user",
                    parts=[types.Part.from_text(text=prompt)],
                )
            ]

            for step in range(1, settings.max_agent_steps + 1):
                response = await client.aio.models.generate_content(
                    model=settings.gemini_model,
                    contents=contents,
                    config=types.GenerateContentConfig(
                        system_instruction=SYSTEM_INSTRUCTION,
                        tools=[
                            types.Tool(
                                function_declarations=declarations
                            )
                        ],
                        temperature=0.2,
                    ),
                )

                if not response.candidates:
                    raise RuntimeError(
                        "Gemini не вернул ни одного кандидата"
                    )

                candidate = response.candidates[0]
                model_content = candidate.content
                if model_content is None:
                    raise RuntimeError(
                        "Gemini вернул кандидата без содержимого"
                    )

                contents.append(model_content)

                calls = [
                    part.function_call
                    for part in (model_content.parts or [])
                    if getattr(part, "function_call", None) is not None
                ]

                if not calls:
                    text = final_text(model_content)
                    if not text:
                        raise RuntimeError(
                            "Агент завершился без текстового ответа"
                        )
                    return {
                        "answer": text,
                        "steps": step,
                    }

                tool_response_parts: list[types.Part] = []

                for call in calls:
                    if call.name not in allowed_names:
                        payload = {
                            "is_error": True,
                            "error": "Инструмент не разрешён",
                        }
                    else:
                        arguments = dict(call.args or {})
                        try:
                            result = await asyncio.wait_for(
                                session.call_tool(
                                    call.name,
                                    arguments=arguments,
                                ),
                                timeout=settings.tool_timeout_seconds,
                            )
                            payload = serialize_tool_result(result)
                        except asyncio.TimeoutError:
                            payload = {
                                "is_error": True,
                                "error": "Тайм-аут MCP-инструмента",
                            }
                        except Exception as exc:
                            payload = {
                                "is_error": True,
                                "error": (
                                    f"Ошибка MCP-инструмента: "
                                    f"{type(exc).__name__}"
                                ),
                            }

                    tool_response_parts.append(
                        types.Part.from_function_response(
                            name=call.name,
                            response={"result": payload},
                        )
                    )

                contents.append(
                    types.Content(
                        role="user",
                        parts=tool_response_parts,
                    )
                )

    raise RuntimeError(
        "MCP-сессия завершилась до ответа агента"
    )

После цикла добавьте явную ошибку превышения лимита. В данном варианте она должна стоять сразу после блока for, но до выхода из MCP-сессии:

            raise RuntimeError(
                f"Агент превысил лимит {settings.max_agent_steps} шагов"
            )

Отступ этого фрагмента должен совпадать с отступом строки for step. Он выполняется, если модель продолжала вызывать инструменты на каждом шаге.

Почему список инструментов нельзя брать из ответа модели

Разрешённые имена формируются только из результата session.list_tools(). Модель не получает возможность задать произвольный URL или выбрать внутренний метод, которого нет в allowlist. Аргументы всё равно должны валидироваться на MCP-сервере: модель считается недоверенным источником ввода.

Ограничение размера результата

Учебный инструмент возвращает небольшой объект. В реальном приложении перед отправкой результата модели ограничьте размер, удалите бинарные поля и персональные данные. Большой результат быстро заполняет контекстное окно и увеличивает стоимость каждого следующего шага.

Шаг 6. Worker с heartbeat

Worker работает отдельно от веб-приложения. Он периодически пытается захватить задачу и продлевает lease, пока агент выполняется.

Файл app/worker.py:

import asyncio
import logging
import signal

from app.agent import execute_agent
from app.config import get_settings
from app.db import (
    claim_job,
    complete_job,
    fail_job,
    heartbeat,
    init_db,
)


logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(message)s",
)
logger = logging.getLogger("agent-worker")


async def heartbeat_loop(
    job_id: str,
    stop: asyncio.Event,
) -> None:
    settings = get_settings()
    interval = max(5, settings.job_lease_seconds // 3)

    while not stop.is_set():
        try:
            await asyncio.wait_for(stop.wait(), timeout=interval)
        except asyncio.TimeoutError:
            await asyncio.to_thread(heartbeat, job_id)


async def process_job(job: dict) -> None:
    stop_heartbeat = asyncio.Event()
    heartbeat_task = asyncio.create_task(
        heartbeat_loop(job["id"], stop_heartbeat)
    )

    try:
        result = await execute_agent(job["prompt"])
        await asyncio.to_thread(complete_job, job["id"], result)
        logger.info("job_succeeded id=%s", job["id"])
    except Exception as exc:
        logger.exception("job_failed id=%s", job["id"])
        error = {
            "code": type(exc).__name__,
            "message": str(exc)[:1000],
            "retryable": False,
        }
        await asyncio.to_thread(fail_job, job["id"], error)
    finally:
        stop_heartbeat.set()
        await heartbeat_task


async def run_worker() -> None:
    init_db()
    stopping = asyncio.Event()
    loop = asyncio.get_running_loop()

    for signal_name in (signal.SIGINT, signal.SIGTERM):
        loop.add_signal_handler(signal_name, stopping.set)

    logger.info("worker_started")

    while not stopping.is_set():
        job = await asyncio.to_thread(claim_job)

        if job is None:
            try:
                await asyncio.wait_for(
                    stopping.wait(),
                    timeout=1.0,
                )
            except asyncio.TimeoutError:
                pass
            continue

        logger.info(
            "job_started id=%s attempt=%s",
            job["id"],
            job["attempt"],
        )
        await process_job(job)

    logger.info("worker_stopped")


if __name__ == "__main__":
    asyncio.run(run_worker())

Worker обрабатывает одну задачу за раз. Это осознанно упрощает контроль нагрузки. Горизонтальное масштабирование достигается запуском нескольких процессов, а не неограниченным созданием coroutine внутри одного процесса.

Heartbeat продлевает lease через каждую треть его длительности. Значение lease должно быть заметно больше интервала heartbeat и ожидаемых коротких пауз планировщика.

Шаг 7. Запуск всей системы

Откройте три терминала в каталоге проекта.

Первый терминал — MCP-сервер:

source .venv/bin/activate
python -m mcp_server.server

Второй терминал — API:

source .venv/bin/activate
uvicorn app.api:app --host 127.0.0.1 --port 8000

Третий терминал — worker:

source .venv/bin/activate
python -m app.worker

Проверка liveness API:

curl --fail http://127.0.0.1:8000/health/live

Ожидается структурный ответ:

{"status":"ok"}

Он подтверждает только то, что API-процесс отвечает. Он не доказывает доступность Gemini, MCP или worker.

Шаг 8. Запуск фоновой задачи

Создайте задание:

curl -i \
  -X POST http://127.0.0.1:8000/v1/jobs \
  -H 'Content-Type: application/json' \
  -H 'Idempotency-Key: billing-503-demo-1' \
  --data '{
    "prompt": "После релиза у сервиса billing выросло число ответов 503. Получи внутреннюю инструкцию и подготовь безопасный план диагностики."
  }'

Скопируйте поле id и запросите состояние:

curl http://127.0.0.1:8000/v1/jobs/JOB_ID

Повторяйте запрос до терминального статуса. Для ручной проверки достаточно следующего shell-цикла:

JOB_ID="вставьте-id"

while true; do
  BODY="$(curl --silent --fail \
    "http://127.0.0.1:8000/v1/jobs/${JOB_ID}")" || exit 1

  printf '%s\n' "$BODY"

  STATUS="$(printf '%s' "$BODY" | python -c \
    'import json,sys; print(json.load(sys.stdin)["status"])')"

  case "$STATUS" in
    succeeded|failed) break ;;
  esac

  sleep 2
done

Для production-клиента используйте экспоненциальную задержку с jitter, например 1, 2, 4, 8 секунд с верхней границей 15–30 секунд. Учитывайте заголовок Retry-After, если сервер его возвращает.

Проверка результата без выдуманных ожиданий

Нельзя заранее утверждать, какой именно текст сгенерирует модель. Проверять следует контракт и наблюдаемые свойства системы, а не дословный ответ.

1. API действительно отвечает асинхронно

  • POST /v1/jobs возвращает 202;
  • в ответе есть непустые id и status_url;
  • задача сохранена до отправки ответа;
  • закрытие клиента не меняет жизненный цикл задания.

2. Состояния имеют допустимый порядок

Нормальная последовательность:

queued → running → succeeded
queued → running → failed

После терминального состояния задача не должна снова становиться running. Поле finished_at должно быть заполнено только для succeeded или failed.

3. Идемпотентность работает

Дважды отправьте одинаковый запрос с одним Idempotency-Key:

for n in 1 2; do
  curl --silent \
    -X POST http://127.0.0.1:8000/v1/jobs \
    -H 'Content-Type: application/json' \
    -H 'Idempotency-Key: same-operation-42' \
    --data '{"prompt":"Получить runbook для billing и составить план диагностики."}'
  printf '\n'
done

Оба ответа должны ссылаться на один идентификатор. Если содержимое запроса может различаться при одинаковом ключе, сохраните также хеш нормализованного тела и возвращайте 409 Conflict при несовпадении.

4. MCP действительно вызывается

Временно добавьте в MCP-сервер структурный журнал имени инструмента и сервиса без текста пользовательского запроса и секретов. Затем убедитесь, что задание с упоминанием billing приводит к вызову lookup_runbook. Не считайте наличие текста о runbook доказательством вызова: модель могла сформулировать общий совет самостоятельно.

5. Ответ не приписывает системе несуществующие действия

Результат может предложить команды и проверки, но не должен утверждать, что метрики уже исследованы или откат выполнен. Наш инструмент возвращает инструкцию, а не телеметрию.

Проверка отказов

Надёжность агента определяется не успешной демонстрацией, а поведением на границах.

MCP-сервер недоступен

Остановите MCP-процесс и создайте задачу. Ожидаемое свойство: HTTP-запрос создания по-прежнему получает 202, а worker переводит задачу в failed с диагностическим кодом. Пользовательский ответ не должен содержать traceback, URL с секретом или токен.

Неверный ключ Gemini

Укажите заведомо тестовое нерабочее значение и перезапустите worker. Задача должна завершиться контролируемой ошибкой. Не записывайте сам ключ в поле error или журнал.

Инструмент отвечает слишком долго

Добавьте в учебный инструмент задержку, превышающую TOOL_TIMEOUT_SECONDS. Агент получит структурированную ошибку инструмента и сможет либо сформировать ограниченный ответ, либо закончить работу ошибкой. Тайм-аут должен действовать на каждый вызов, а не только на задачу целиком.

Worker завершён посреди задачи

Завершите worker во время статуса running. После истечения lease другой worker должен суметь захватить строку. Это проверяет восстановление, но не гарантирует однократное выполнение внешнего действия. Изменяющие MCP-инструменты всё равно требуют идемпотентного ключа.

Модель вызывает инструмент бесконечно

Установите MAX_AGENT_STEPS=1 и дайте задачу, для которой инструмент обязателен. Если после первого вызова модель не успевает сформировать ответ, задача должна завершиться ошибкой лимита, а не продолжаться бесконечно.

Повреждённый или огромный ответ инструмента

MCP-клиент должен считать результат недоверенным. Перед отправкой в Gemini добавьте:

  • лимит сериализованного размера;
  • удаление лишних полей;
  • нормализацию кодировки;
  • маркер усечения;
  • запрет на автоматическую интерпретацию результата как системной инструкции.

Безопасность удалённого MCP

Удалённый MCP решает проблему интеграции, но одновременно создаёт новую сетевую границу. Её следует проектировать как привилегированный внутренний API.

Не позволяйте модели выбирать адрес

MCP_URL задаётся оператором, а не пользовательским prompt и не ответом модели. Иначе появляется риск SSRF: агент может обратиться к metadata endpoint, панели администрирования или другому внутреннему адресу.

Разделяйте инструменты чтения и изменения

Для чтения достаточно строгой схемы, лимитов и аудита. Инструменты, меняющие состояние, требуют дополнительных мер:

  • явного подтверждения человеком для опасных операций;
  • идемпотентного operation ID;
  • проверки прав на стороне инструмента;
  • короткого срока действия авторизации;
  • журнала до и после выполнения;
  • отдельного лимита частоты;
  • возможности запуска в режиме dry-run.

Считайте содержимое инструмента недоверенным

Документ из базы знаний может содержать текст, похожий на инструкцию модели. Это форма prompt injection. Системная инструкция должна явно говорить, что данные инструмента — источник фактов, а не новые правила. Для чувствительных операций одной инструкции недостаточно: политика разрешений должна исполняться обычным кодом.

Минимизируйте полномочия

Worker для анализа инцидента не должен получать универсальный токен администратора. Лучше опубликовать отдельный read-only инструмент с узким ответом, чем давать агенту общий SQL endpoint или shell.

Повторы и классификация ошибок

В учебной реализации любая ошибка переводит задачу в failed. Автоматические повторы следует добавлять только после классификации причин.

Ситуация Повтор Почему
Сетевой сбой до начала ответа Обычно да Кратковременная инфраструктурная ошибка
Ограничение частоты Да, с backoff Нужно учитывать рекомендованную задержку
Невалидные аргументы инструмента Нет автоматически Повтор тех же данных ничего не изменит
Отказ в доступе Нет Требуется исправление конфигурации или прав
Неизвестен результат изменяющей операции Только с operation ID Первый вызов мог выполниться до разрыва
Превышен лимит шагов Нет Нужно менять prompt, инструменты или лимит

Для повторов добавьте состояния retry_wait и поле next_attempt_at. Не удерживайте worker через длительный sleep: верните задачу в очередь с будущим временем запуска.

Наблюдаемость

Каждая задача должна иметь единый correlation ID, совпадающий с job_id. Передавайте его в журналы worker и, если протокол инструмента допускает метаданные, в MCP-слой.

Минимальные метрики:

  • число задач по статусам;
  • возраст самой старой задачи queued;
  • время от создания до начала;
  • время выполнения;
  • число шагов модели;
  • число и длительность вызовов каждого инструмента;
  • ошибки Gemini и MCP по классу;
  • количество задач, повторно захваченных после истечения lease.

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

Отмена задачи

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

Для поддержки отмены добавьте состояния cancel_requested и cancelled. Worker должен проверять флаг:

  • перед вызовом Gemini;
  • перед каждым инструментом;
  • после возврата внешнего вызова;
  • перед сохранением финального результата.

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

Когда заменить SQLite

SQLite подходит для одного хоста, умеренного потока заданий и учебного развёртывания. WAL улучшает параллельное чтение, но не превращает файл в распределённую очередь.

Переходите на PostgreSQL или специализированный брокер, если:

  • API и worker работают на разных хостах;
  • несколько worker часто конкурируют за запись;
  • нужны приоритеты и планирование на будущее;
  • требуется высокая пропускная способность;
  • нужны dead-letter queue и развитая политика повторов;
  • локальный диск контейнера не является постоянным.

В PostgreSQL захват задания обычно строится на транзакции с SELECT … FOR UPDATE SKIP LOCKED. При использовании брокера статус и результат всё равно удобно хранить в отдельной базе, потому что очередь не обязана быть пользовательским API состояния.

Ограничения решения

  • Нет строгого exactly-once. Lease обеспечивает восстановление, но одна работа может начаться повторно.
  • Нет встроенного планировщика повторов. Ошибки сейчас терминальны.
  • Нет потоковой выдачи частичного ответа. Клиент получает только статус и итог.
  • Один worker выполняет одну задачу. Масштабирование требует дополнительных процессов и контроля квот.
  • История диалога не хранится отдельно. Система реализует задачу, а не многоходовую пользовательскую сессию.
  • MCP-сервер из примера не готов для публичной сети. Нужны TLS, авторизация, аудит и сетевые ограничения.
  • Схемы SDK могут меняться. Зафиксируйте проверенные версии зависимостей после успешного воспроизводимого запуска.
  • Стоимость не ограничена токенами напрямую. Лимит шагов помогает, но следует также ограничить длину prompt, объём результатов и бюджет на задачу.

Production-чеклист

  • API возвращает 202 только после надёжной записи задания.
  • Есть уникальный ключ идемпотентности и проверка хеша тела.
  • Worker отделён от API-процесса.
  • Захват задания атомарен.
  • Lease и heartbeat проверены принудительным завершением worker.
  • Есть лимиты времени, шагов, размера prompt и ответа инструмента.
  • MCP URL нельзя задать пользовательским вводом.
  • Имена инструментов сверяются с allowlist.
  • Аргументы валидируются на MCP-сервере.
  • Изменяющие инструменты имеют operation ID и отдельную авторизацию.
  • Секреты не попадают в результат задачи и журналы.
  • Клиент использует polling с backoff и прекращает его на терминальном статусе.
  • Есть метрики очереди, длительности, ошибок и повторного захвата.
  • Версии Python-зависимостей зафиксированы после проверки.
  • Проведены тесты недоступности Gemini, MCP и базы.

Итог

Надёжный управляемый агент — это не один длинный запрос к модели. Это конечный автомат задачи, персистентная очередь, отдельный worker, ограниченный агентный цикл и защищённая граница инструментов.

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

Продолжить практику можно в разделе гайдов Agent Lab Journal. Определения использованных понятий и протоколов собраны в глоссарии.