Перейти к содержимому

AI Agent Service — реалтайм-чат с AI-агентом

Слой: libs/apis/providers/ai-agent-service · Приложение: apps/ai-agent-service

ai-agent-service — это шлюз между фронтендом заказчика (customer) и внешним AI-агентом Walrider. Сервис:

  • принимает сообщения пользователя по WebSocket (реалтайм-чат в контексте проекта);
  • сохраняет диалог в БД (модели WalriderThread / WalriderMessage / WalriderThreadState);
  • публикует запрос в RabbitMQ, откуда его забирает внешний AI-агент;
  • слушает очередь ответов RabbitMQ, сохраняет ответ агента и пушит его обратно клиенту по WebSocket;
  • отдаёт историю сообщений и текущее состояние диалога по HTTP REST.

Таким образом сервис совмещает три транспорта: HTTP (REST + Swagger), WebSocket (Socket.IO) и RabbitMQ (amqplib).

Приложение поднимается как HTTP-сервис на Express (NestFactory.create<NestExpressApplication>), а не как «чистый» микросервис:

  • глобальный префикс api, URI-версионирование (/api/v1/...), Swagger на /api;
  • APP_PORT (по умолчанию 3000, см. libs/apis/configs/shared/app);
  • CORS с credentials: true, origin: true.
  • WebSocket-шлюз поднимается автоматически поверх того же HTTP-сервера через стандартный IoAdapter Socket.IO (в main.ts кастомный адаптер не регистрируется — используется дефолтный) (предположительно, отдельный WS-порт не настраивается — сокет живёт на порту HTTP).
  • Подключение к RabbitMQ создаётся вручную в RabbitConsumerService.onModuleInit() через RabbitClientService (обёртка над amqplib), а не через NestFactory.createMicroservice.
flowchart LR
subgraph Client["Клиент (заказчик)"]
UI["WEB UI"]
end
subgraph AIS["ai-agent-service (NestJS, Express + Socket.IO + amqplib)"]
GW["AgentChatGateway<br/>ns: /agent-chat"]
SVC["AgentChatService"]
CONS["RabbitConsumerService"]
CTRL["AgentChatController<br/>REST /api/v1/agents"]
REDIS[("Redis<br/>ws:project:{id}")]
end
subgraph MQ["RabbitMQ"]
RX(["requestExchange"])
RQ["requestQueue<br/>(TTL 600s, DLX ai.dlx)"]
PX(["responsesExchange"])
PQ["responsesQueue"]
end
EXT["Внешний AI-агент<br/>WALRIDER"]
DB[("PostgreSQL<br/>walrider_threads /<br/>_messages / _thread_states")]
UI -- "WS sendMessage" --> GW
GW --> SVC
SVC -- "publish(requestExchange,'')" --> RX
RX --> RQ
RQ -- "consume" --> EXT
EXT -- "publish" --> PX
PX --> PQ
PQ -- "consume" --> CONS
CONS --> SVC
CONS -- "emit 'message'" --> GW
GW -- "WS message" --> UI
SVC <--> DB
SVC <--> REDIS
CTRL --> SVC
UI -- "HTTP history/state" --> CTRL

Роли в RabbitMQ:

  • Продюсер запросовai-agent-service (метод AgentChatService.sendMessagerabbitClient.publish(requestExchange, '', msg)).
  • Консьюмер запросов — внешний AI-агент Walrider (в этом репозитории его кода нет).
  • Продюсер ответов — внешний AI-агент Walrider.
  • Консьюмер ответовai-agent-service (RabbitConsumerService.consume(responsesQueue, ...)).

Шлюз: AgentChatGateway, namespace /agent-chat, транспорт Socket.IO, комнаты вида project:{projectId}.

Аутентификация при подключении (handleConnection): из handshake.query берутся token (JWT заказчика) и projectId. JWT проверяется секретом JwtCustomerConfigService.accessSecretKey; проверяется существование customer и что проект принадлежит этому заказчику. При провале — client.disconnect(). После успеха клиент кладётся в комнату project:{projectId}, а его socketId сохраняется в Redis по ключу ws:project:{projectId} (одно активное соединение на проект).

СобытиеНаправлениеPayloadОписание
connection (handshake)client → serverquery: { token, projectId }Подключение к ns /agent-chat. Проверка JWT + доступа к проекту, вход в комнату project:{projectId}, запись активного сокета в Redis.
disconnectclient → serverУдаление активного сокета проекта из Redis (removeActiveSocket).
sendMessageclient → server{ content: string }Пользователь отправляет сообщение. Сервис сохраняет WalriderMessage(role=USER) и публикует запрос в RabbitMQ.
messageSentserver → client (ack){ messageId: number }Синхронный ack на sendMessage — возвращает messageId созданного сообщения.
messageserver → client{ messageId, response, phase, isComplete, plan }Пуш ответа AI-агента в комнату project:{projectId} после прихода сообщения из responsesQueue.

Валидация payload события sendMessage описана DTO WsSendMessageDto (content: string) в data-access. AuthenticatedSocket.data хранит { customerId, projectId }.

Инициализация топологии — в RabbitConsumerService.onModuleInit(): объявляются оба exchange (durable), обе очереди (durable), биндинги с пустым routing key '', затем запускается consume только на очереди ответов.

ОбъектEnv-переменнаяНазначение
requestExchangeRABBITMQ_REQUEST_EXCHANGE (+ _TYPE)Exchange запросов к агенту. Продюсер — ai-agent-service.
requestQueueRABBITMQ_REQUEST_QUEUEОчередь запросов. Объявляется с x-message-ttl: 600000 (10 мин) и x-dead-letter-exchange: 'ai.dlx'. Читает её внешний агент.
responsesExchangeRABBITMQ_RESPONSES_EXCHANGE (+ _TYPE)Exchange ответов от агента. Продюсер — внешний агент.
responsesQueueRABBITMQ_RESPONSES_QUEUEОчередь ответов. Консьюмер — ai-agent-service (RabbitConsumerService).

Прочие env (libs/apis/configs/shared/rabbit, Joi-валидация): RABBITMQ_HOST, RABBITMQ_PORT, RABBITMQ_USERNAME, RABBITMQ_PASSWORD, RABBITMQ_PREFETCH (default 50).

DLX-обменник ai.dlx в топологии сервиса явно не создаётся отдельным assertExchange — предполагается его наличие на стороне инфраструктуры/внешнего агента (предположительно).

Структура сообщенийrabbit-messages.interface.ts:

// запрос: ai-agent-service → внешний агент
interface RabbitRequestMessage {
messageId: number; // WalriderMessage.messageId (autoincrement)
threadId: string; // WalriderThread.id (uuid)
content: string; // текст пользователя
}
// ответ: внешний агент → ai-agent-service
interface RabbitResponseMessage {
messageId: number; // id сообщения-запроса
threadId: string; // WalriderThread.id
response: string; // текст ответа агента
phase: string; // фаза/статус диалога → state.status
isComplete: boolean; // признак завершённости
plan: Record<string, unknown> | null; // структурированный «план» агента
}

Тело сообщений сериализуется в JSON (Buffer.from(JSON.stringify(...)), persistent: true, contentType: application/json). При обработке ответа консьюмер делает ack, при ошибке — nack(msg, false, false) (без requeue → в DLX).

  • @crewsforge-back/rabbit-client (RabbitClientService) — низкоуровневая обёртка над amqplib (connect / assertExchange / assertQueue / bindQueue / publish / consume / graceful shutdown).
  • @crewsforge-back/config-shared-rabbit (RabbitConfigService) — env-конфиг RabbitMQ.
  • @crewsforge-back/redis-client (RedisClientService) — хранит активный socketId на проект.
  • @crewsforge-back/apis/utils/prisma-client + @nestjs-cls/transactional (TransactionalAdapterPrisma) — транзакционный доступ к БД через TransactionHost.
  • @nestjs/jwt + @crewsforge-back/apis/configs/auth-api/jwt-customer — проверка JWT заказчика в WS-хэндшейке и в CustomerGuard для REST.
  • @crewsforge-back/apis/sharedCustomerGuard, MapperService, GetTokenPayload, PrismaExceptionFilter, пагинация.
  • Prisma-модели живут в apps/core-api/prisma/schema.prisma (общая схема).
  • apps/ai-agent-service/src/main.ts — bootstrap Express + Swagger + порт.
  • apps/ai-agent-service/src/app/app.module.ts — корневой модуль (ClsModule + Transactional + AgentChatModule).
  • libs/apis/providers/ai-agent-service/features/agent-chat/src/lib/agent-chat.module.ts — сборка фичи.
  • .../agent-chat.gateway.ts — WebSocket-шлюз, auth, комнаты, события.
  • .../agent-chat.service.ts — бизнес-логика: sendMessage, handleAgentResponse, getHistory, getState, работа с Redis.
  • .../rabbit-consumer.service.ts — инициализация топологии RabbitMQ и консьюмер ответов.
  • .../agent-chat.controller.ts — REST /api/v1/agents/:projectId/history и /state.
  • .../walrider-thread.repository.ts, .../walrider-message.repository.ts, .../walrider-state.repository.ts — репозитории (тонкие обёртки над Prisma через TransactionHost).
  • libs/apis/providers/ai-agent-service/data-access/src/lib/interfaces/rabbit-messages.interface.ts — контракты сообщений.
  • libs/apis/providers/ai-agent-service/data-access/src/lib/dtos/agent-chat.dto.ts — DTO WS/HTTP (WsSendMessageDto, AgentMessageDto, AgentStateDto, AgentHistoryQueryDto).
  • libs/apis/utils/rabbit-client/**, libs/apis/configs/shared/rabbit/** — клиент и конфиг RabbitMQ.
  • Новые события WS: добавить @SubscribeMessage(...) в AgentChatGateway + DTO в data-access.
  • Новые типы сообщений RabbitMQ: расширить RabbitRequestMessage / RabbitResponseMessage и логику в AgentChatService.handleAgentResponse.
  • Больше состояний диалога: phase/isComplete/plan кладутся в WalriderMessage.metadata и WalriderThreadState.data — можно расширять без миграций (поля Json).
  • Масштабирование WS: сейчас активный сокет — один на проект в Redis; для нескольких инстансов понадобится Socket.IO Redis-adapter (предположительно).
  • Внешний AI-агент: заменяем только консьюмера requestQueue / продюсера responsesQueue — контракт остаётся прежним.