sequenceDiagram
autonumber
actor User as Мобильный клиент (Микрофон)
participant Nginx as Nginx Proxy
participant GW as FastAPI Gateway (backend-api)
participant Worker as purchase-stream-worker (INVENTORY)
participant AudioAI as grpc-analytics (AID)
participant K as Apache Kafka Cluster
participant Censor as censorship-control-worker (SECURITY)
%% ЭТАП 1: ИНИЦИАЦИЯ СЕССИИ И АПГРЕЙД ПРОТОКОЛА
User->>Nginx: GET /api/v2/stream/voice/init (Запрос WebSocket)
Nginx->>GW: Пересылка HTTP-запроса на авторизацию сессии
Note over GW: Шаг 2: Валидация сессии пользователя<br/>Инъекция X-Request-ID (trace_id) и app_lang
GW-->>Nginx: HTTP 200 OK (Разрешение на Upgrade)
Nginx-->>User: HTTP 101 Switching Protocols
%% ЭТАП 2: ПОТОКОВАЯ ПЕРЕДАЧА И СТРИМИНГ
Nginx->>Worker: Проброс и удержание WebSocket соединения
loop Передача бинарных чанков аудио
User-->>Worker: WS Stream: Бинарные аудио-данные (Chunks)
Worker->>AudioAI: Двунаправленный gRPC: TranscribeAudioChunkStream(Chunk)
Note over AudioAI: Шаг 7: Инференс локального Whisper Tiny<br/>Распознавание промежуточного текста фраз
AudioAI-->>Worker: gRPC Response: Промежуточный текст (Token text)
Worker-->>User: WS Stream: Отображение текста на лету в интерфейсе
end
%% ЭТАП 3: ФИНАЛИЗАЦИЯ И СБРОС В ЦЕНЗУРУ
User->>Worker: WS Control Event: CLOSE_STREAM (Юзер отпустил кнопку)
Note over Worker: Шаг 10: Сборка всех чанков в финальный текст сессии<br/>Закрытие WebSocket / gRPC каналов
Worker->>K: Пуш в топик: bpds.inventory.in.receipt.upload
Worker-->>User: WS Connection Closed (Успешно завершено)
%% ЭТАП 4: КОНТУР ЦЕНЗУРЫ
K->>Censor: Handler: ProcessUploadText()
Note over Censor: Исход 12: Текст чист (is_clean = True)
Censor->>K: Пуш в топик: bpds.inventory.out.receipt.parsed
Метод POST /api/v2/stream/voice/init
Документация API: Текстовое добавление номенклатур с автоматической gRPC ИИ-модерацией
- IAD-MIGRATION-116 Backlog — Краткое Описание задачи 1.
- IAD-MIGRATION-117 Refinement (Уточнение) — Краткое Описание задачи 2.
- GW-2 Ready for Development — Подключение эндпоинта в шлюзе.
- AID-012 Ready for Development — создать метож ProcessVoiceStream.
- AID-013 Ready for Development — создать метож ProcessUploadText.
- INV-FRONTEND-222 Refinement (Уточнение) — Экран управления B2B-интеграциями
- INV-FRONTEND-223 Refinement (Уточнение) — Экран управления B2B-интеграциями
1 Функциональное назначение
Эндпоинт реализует логику быстрого добавления продуктов питания в «цифровой холодильник» посредством ввода неструктурированного текста (например, через строку АРМ или интерфейс быстрого подбора мобильного приложения).
2 Протокол взаимодействия (HTTP Контракт)
- Метод:
POST - Маршрут:
/api/v2/stream/voice/init - Формат данных:
application/json
2.1 Спецификация тела запроса (Request Body)
| Поле | Тип | Обязательный | Описание |
|---|---|---|---|
raw_text_input |
String | Да | Неструктурированная текстовая строка ввода (например: "Груши 2.5") |
3 Схема обработки запроса пользователя (Mermaid)
В настоящей публичной документации отображены не все шаги и сценарии для приложения в частности и для системы цифровых симуляторов бизнес-процессов в общем.
4 Расшифровка шагов
- Шаг 1 (
User -> APP): Пользователь открывает экран «Умная камера» во Flutter-приложении и наводит объектив смартфона на продукт (например, лежащий на столе томат). - Шаг 2 (
APP -> API): Приложение инициируетHTTP GETзапрос на эндпоинт/stream-capture, передавая авторизационный токен и уникальныйX-Request-ID. Запрашивается обновление протокола до WebSockets. - Шаг 3 (
API -> APP):FastAPIпроизводит валидацию токена. Если сессия подтверждена, сервер возвращает статусHTTP 101 Switching Protocols. Устанавливается постоянное дуплексное TCP-соединение. Бэкенд переходит в режим ожидания фреймов (activate API). - Шаг 4 (
APP -> API): Внутри бесконечного цикла мобильное приложение захватывает кадры с видеобуфера камеры, сжимает их до разрешения640x480в форматJPEGи непрерывно отправляет в WebSocket-канал в виде бинарных фреймов (Binary Frame). - Шаг 5 (
API -> KAFKA): Бэкенд-шлюзFastAPIработает в асинхронном неблокирующем режиме. Он не сохраняет файлы на диск и не вызывает ИИ-модели. Получив бинарный фрейм, он моментально упаковывает его в массив байт и выполняет методproducer.send(), отправляя кадр в топикmedia.camera.stream.v2. На этом контекст шлюза освобождается (deactivate API). - Шаг 6 (
KAFKA -> GRPC): ИИ-сервис на базе PyTorch работает как изолированный потребитель (Consumer). Он вычитывает бинарные кадры из топикаmedia.camera.stream.v2. Если приложение шлет 30 кадров в секунду, а видеокарта (GPU) успевает обработать только 15, Kafka удерживает оставшиеся 15 кадров в очереди, не перегружая память ИИ-сервиса. - Шаг 7 (
GRPC -> GRPC): Внутренний шаг ИИ-сервиса: каждый извлеченный из очереди кадр пропускается через сверточную нейросетьMobileNetV2. Модель рассчитывает вектор вероятностей. - Шаг 8 (
GRPC -> GRPC): Модель фиксирует совпадение: на картинке обнаружен объект с кодом класса, соответствующим техническому SKU"TOMATO", с высоким коэффициентом уверенности (confidence = 0.94). - Шаг 9 (
GRPC -> API): ИИ-сервис открывает обратный gRPC-канал к шлюзу и отправляет сообщениеStreamResponse(predicted_sku='TOMATO', confidence=0.94). Воркер ИИ-сервиса уходит на чтение следующего кадра из Kafka. - Шаг 10 (
API -> APP):FastAPIпринимает gRPC-ответ, находит нужную WebSocket-сессию пользователя поsession_idи отправляет текстовый JSON-пакет с событиемSTREAM_PREDICTION_UPDATED. - Шаг 11 (
APP -> APP): Интерфейс мобильного приложения мгновенно реагирует на сигнал: поверх превью камеры отрисовывается зеленая рамка-таргет, а пользователю выводится подсказка: «Томат (94%). Удерживайте для фиксации». - Шаг 12 (
User -> APP): Пользователь видит, что ИИ правильно распознал продукт, и нажимает на экране физическую кнопку затвора «Зафиксировать продукт». - Шаг 13 (
APP -> API): Вместо отправки тяжелой фотографии, приложение шлет в открытый WebSocket легковесный текстовый JSON-фрейм команды:{"action": "FREEZE_FRAME", "sku": "TOMATO"}. - Шаг 14 (
API -> DB): ШлюзFastAPIперехватывает команду фиксации, останавливает трансляцию кадров из Kafka и обращается к PostgreSQL. Так как ИИ определил тип как одиночный продукт, срабатывает логика маршрутизации: бэкенд выполняетINSERT INTO pending_products_bufferс дефолтным статусомPROCESSING. - Шаг 15 (
API -> APP): Сервер отправляет в вебсокет финальный статусSESSION_FROZEN. Соединение закрывается. На стороне Flutter экран камеры блокируется, и приложение автоматически перенаправляет пользователя на форму ручного ввода, где ему остается только указать примерный вес томата (например,0.3 кг), как это описано в стандарте Версии 1.
5 Спецификация сообщений брокера и вилок исключений
5.1 Спецификация топика Apache Kafka: media.camera.stream.v2
В этот топик бэкенд-шлюз FastAPI непрерывно сбрасывает кадры, полученные из бинарных фреймов WebSocket. Для оптимизации пропускной способности данные передаются в бинарном виде (сериализация ByteArray), а метаданные упаковываются в заголовки сообщения (Kafka Headers).
- Стратегия партиционирования (Partition Key):
session_id(гарантирует, что все кадры в рамках одной сессии камеры попадают на одну партицию и обрабатываются ИИ строго в хронологическом порядке). - Формат Key:
String(значениеsession_id). - Формат Value:
Byte[](сырой бинарный поток сжатого JPEG-изображения). - Конфигурация заголовков сообщения (Kafka Headers):
| Ключ заголовка | Тип данных | Описание | Пример значения |
|---|---|---|---|
x-request-id |
String | Сквозной ID сессии для распределенного трейсинга | stream-7a2b-9c1d-00ef |
frame_timestamp |
String (ISO) | Время захвата кадра на мобильном устройстве | 2026-07-01T21:10:05.123Z |
user_id |
String | Идентификатор владельца сессии | usr-8822-fa41 |
5.2 Спецификация управляющих текстовых фреймов (JSON)
5.2.1 Команда фиксации кадра от клиента (Шаг 13)
Отправляется в WebSocket-канал в текстовом формате при нажатии пользователем кнопки затвора.
{
"action": "FREEZE_FRAME",
"timestamp": "2026-07-01T21:10:12.450Z",
"payload": {
"verified_sku": "TOMATO",
"confidence_at_freeze": 0.942
}
}5.2.2 Ответ сервера о заморозке сессии (Шаг 15)
Финальное системное сообщение от FastAPI, подтверждающее успешную запись в буфер и закрытие стрима.
{
"event_type": "SESSION_FROZEN",
"timestamp": "2026-07-01T21:10:13.110Z",
"data": {
"draft_id": "8a71d11e-9500-4b11-9a2c-d2b0d7b3dcba",
"target_buffer_table": "pending_products_buffer",
"next_ui_step": "NAVIGATE_TO_MANUAL_WEIGHT_INPUT"
}
}5.3 Вилки исключений и стратегии обработки ошибок (Failures & Fail-Safe)
Стриминговый контур чувствителен к сетевым задержкам и стабильности соединений. Система обрабатывает три основных аварийных сценария:
5.3.1 Сценарий А: Отказ ИИ-сервиса или падение gRPC (Зависание инференса)
Если ИИ-сервис уходит в offline или перестает отвечать на gRPC-запросы шлюза FastAPI, срабатывает стратегия Client-Side Timeout / Fallback.
- Поведение системы: Если в течение 3 секунд от сервера не приходит ни одного обновления
STREAM_PREDICTION_UPDATED, Flutter-приложение мягко переключает интерфейс в режим ручной съемки. На экран выводится тост-уведомление, а кнопка затвора меняет логику: при нажатии отправляется стандартный пакетныйHTTP POST /product-photoпо сценарию Версии 1. Система не блокирует пользователя.
5.3.2 Сценарий Б: Ошибка переполнения буфера или отказ Kafka (HTTP 503)
Выбрасывается на Шаге 5, если брокер Kafka недоступен или внутренний буфер бэкенда переполнен из-за аномально высокого FPS от клиента.
{
"error_code": "ERR-STREAM-BUFFER-OVERFLOW",
"message": "Поток камеры временно перегружен. Пожалуйста, удерживайте камеру неподвижно.",
"details": {
"reason": "Kafka broker partition leader is unavailable or local producer memory buffer is full."
}
}5.3.3 Сценарий В: Разрыв сетевого соединения в процессе стриминга (Close Code)
Если у пользователя пропадает 4G/Wi-Fi соединение на Шаге 4, WebSocket-сессия аварийно закрывается операционной системой.
- Поведение системы:
FastAPIловит системное событиеWebSocketDisconnect. Бэкенд немедленно отправляет команду в Kafka на удаление промежуточных необработанных кадров этой сессии, очищая партицию. В базе данныхPostgreSQLтранзакция черновика аннулируется. При восстановлении сети Flutter-приложение генерирует новыйsession_idи начинает процесс с Шага 2.