Метод POST /api/v2/stream/voice/init

Документация API: Текстовое добавление номенклатур с автоматической gRPC ИИ-модерацией

Published

June 11, 2026

1 Функциональное назначение

Эндпоинт реализует логику быстрого добавления продуктов питания в «цифровой холодильник» посредством ввода неструктурированного текста (например, через строку АРМ или интерфейс быстрого подбора мобильного приложения).

2 Протокол взаимодействия (HTTP Контракт)

  • Метод: POST
  • Маршрут: /api/v2/stream/voice/init
  • Формат данных: application/json

2.1 Спецификация тела запроса (Request Body)

Поле Тип Обязательный Описание
raw_text_input String Да Неструктурированная текстовая строка ввода (например: "Груши 2.5")

3 Схема обработки запроса пользователя (Mermaid)

WarningВажное примечание

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

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

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.