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

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

Published

June 11, 2026

WarningОграничение публичной документации

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

  • Полная спецификация метода: Будет доступна только во внутреннем контуре разработки (Confluence / Swagger Enterprise).

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

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

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

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

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

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

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

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 Расшифровка шагов

Шаг Действие Параметры / Запросы / DTO Ошибки (Исключения / Статусы)
Шаг 1 (App -> Nginx) Мобильное приложение отправляет запрос на согласование параметров медиа-сессии и инициацию сигнального сокета для трансляции видео. HTTP GET /api/v2/video-stream/init
Headers: Upgrade: websocket, Connection: Upgrade, Authorization: Bearer <JWT>
DioException: connection timeout
HTTP 400 Bad Request (некорректные заголовки или параметры рукопожатия)
HTTP 401 Unauthorized
Шаг 2 (Nginx -> FastAPI_Gateway) Прокси-сервер перенаправляет запрос на шлюз, выполняя инъекцию сквозного X-Request-ID и пробрасывая метаданные языкового контекста. HTTP GET /internal/v2/video-stream/authorize
Added Headers:
X-Request-ID: "trace-video-stream-uuid", X-App-Lang: "ru-RU"
HTTP 502 Bad Gateway (контейнер backend-api API Gateway недоступен/упал)
HTTP 504 Gateway Timeout
Шаг 3 (FastAPI_Gateway -> Nginx) Шлюз проводит валидацию прав и сессии пользователя, возвращая апрув на переключение протокола и трансляцию потока. HTTP Статус: 101 Switching Protocols (в схеме транслирован как HTTP 200 разрешения на трансляцию). FastAPI.ValidationError (сбой схемы валидации сессии шлюзом)
HTTP 403 Forbidden (исчерпан лимит видео-трафика сессии)
Шаг 4 (Nginx -> StreamWorker) Nginx пробрасывает и удерживает стабильное WebSocket/медиа-соединение с воркером стримов для непрерывной передачи данных. Проброс потока на purchase-stream-worker: ws://stream-worker.internal/v2/video-live. NginxError: 502 Bad Gateway (воркер стримов недоступен или перезагружается)
WebSocketException: Connection refused
Шаг 5 (App -> StreamWorker) Клиент начинает непрерывную потоковую передачу сырых бинарных видеофреймов (Raw Frames/YUV/JPEG) высокого разрешения. WebSocket Binary Frame / WebRTC Media Chunk:
Payload: raw_video_frame_bytes, Частота: ~25-30 fps.
WebSocketDisconnect (аварийный обрыв связи на стороне смартфона)
FrameTooLargeException (размер кадра превысил буфер воркера)
Шаг 6 (StreamWorker -> ProductAI) Воркер стримов выполняет пофреймовую пересылку бинарных матриц изображений в сервис компьютерного зрения по RPC-каналу. gRPC Streaming Request:
ProcessVideoFrames(stream VideoFrameRequest)
Payload: frame_content: bytes, width: 1280, height: 720
gRPC Status: UNAVAILABLE (сервис product-vision-processor недоступен)
gRPC Status: DEADLINE_EXCEEDED (падение FPS из-за троттлинга GPU)
Шаг 7 (ProductAI -> Redis) Модель YOLO детектирует объекты, выполняет скоринг детекций и обновляет трекинг-состояние сессии в кэше, связывая кадры в динамике. Redis Command Pipeline:
1. HSET "stream:track:trace-video-stream-uuid" "pear_01" "{"score":0.94,"count":12}"
2. HSET "stream:track:trace-video-stream-uuid" "apple_02" "{"score":0.89,"count":5}"
3. EXPIRE "stream:track:trace-video-stream-uuid" 1800
Redis.ConnectionError: Connection refused
Redis.TimeoutError: Command timed out
Redis.ClusterDownException
Шаг 8 (ProductAI -> StreamWorker) Сервис компьютерного зрения возвращает воркеру текущий список успешно идентифицированных продуктов и координаты их рамок. gRPC Streaming Response:
VideoFrameResponse(detected_objects=[{"class": "Груша", "box":, "confidence": 0.94}])
YOLOInferenceError (сбой аллокации памяти CUDA)
gRPC Status: INTERNAL (ошибка маршаллинга структуры ответа)
Шаг 9 (StreamWorker -> App) Воркер мгновенно транслирует координаты рамок обратно клиенту для динамической подсветки Bounding Boxes поверх видеопотока. WebSocket Text Frame:
Payload: { "event": "BOUNDING_BOXES", "objects": [{"name": "Груша", "coords": [100, 150, 200, 300]}], "x_request_id": "trace-video-stream-uuid" }
WebSocketException: Broken pipe (приложение свернуто, или ОС заблокировала фоновый сокетный поток)
Шаг 10 (App -> StreamWorker) Пользователь нажимает кнопку «Фиксировать содержимое», и приложение отправляет команду завершения сканирования и остановки трансляции. WebSocket Text Frame (Control Message):
Payload: { "event": "STOP_STREAM", "session_id": "ws-video-session-uuid" }
WebSocketException: ConnectionReset (клиент принудительно отключился или потерял сеть без отправки финальной команды)
Шаг 11 (StreamWorker -> Redis) Воркер стримов запрашивает из оперативной памяти финальный агрегированный список всех затреканных классов объектов для данной сессии. Redis Command:
HGETALL "stream:track:trace-video-stream-uuid"
Redis.ConnectionError: Connection refused
Redis.TimeoutError: Command timed out (превышено время ожидания ответа от кэша Redis)
Шаг 12 (StreamWorker -> T_Results) Модуль агрегирует вычитанные данные, формирует текстовый список найденных классов продуктов и публикует итоговый payload в брокер. Kafka Message (Topic: bpds.inventory.in.receipt.upload):
Payload: { "raw_text_input": "Груша, Яблоко", "app_lang": "ru-RU", "x_request_id": "trace-video-stream-uuid" }
KafkaException: DeliveryTimeout
KafkaException: MessageTooLargeException (агрегированный текстовый список превысил лимиты брокера)
Шаг 13 (T_Results -> Censor) Воркер цензуры считывает из топика загрузок сформированный текстовый список классов объектов видеострима для проведения лингвистической проверки. Kafka Consumer Poll Request:
Group_ID: "censorship-workers", Topic: bpds.inventory.in.receipt.upload
Partition: 5, Offset: 104210
KafkaException: CommitFailedException (зависание обработчика)
SerializationException (ошибка десериализации агрегированной структуры)
Шаг 14 (Censor -> T_Clean_Text) Успешный исход: Цензор успешно верифицирует список классов, подтверждает флаг чистоты is_clean = True и пересылает проверенный текст в топик парсинга. Kafka Message (Topic: bpds.inventory.out.receipt.parsed):
Payload: { "clean_text_input": "Груша, Яблоко", "app_lang": "ru-RU", "x_request_id": "trace-video-stream-uuid" }
KafkaException: DeliveryTimeout
KafkaException: LeaderNotAvailableException

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.