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

Документация API: Потоковая верификация продуктов через реальное время (Camera Stream) с использованием буферной очереди Apache Kafka

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

Метод предназначен для обработки потока изображений (кадров) в реальном времени, поступающих непосредственно с камеры мобильного приложения во время наведения на продукт. В отличие от пакетной загрузки готовых файлов, данный контур минимизирует задержку ввода (Latency) и защищает вычислительные ресурсы ИИ-сервера (PyTorch/GPU) от перегрузки при высокой частоте кадров (FPS).

Бэкенд-шлюз FastAPI в этой схеме выступает в роли ультрабыстрого легковесного ретранслятора, который сгружает бинарные кадры в брокер сообщений Apache Kafka, выполняющий роль амортизатора нагрузки (Load Smoother). ИИ-сервис асинхронно вычитывает кадры из очереди, производит инференс и возвращает результат распознавания обратно в WebSocket-сессию пользователя.

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

  • Протокол: WebSockets (WSS)
  • Маршрут (Эндпоинт): /api/v2/fridge/stream-capture/{session_id}
  • Тип обмена: Дуплексный (Двунаправленный асинхронный поток)
  • Базовый формат кадров (Frame Payload): Бинарный поток сжатых изображений (image/jpeg или image/webp) + JSON-управляющие сообщения.

2.1 Спецификация заголовков при установлении соединения (Handshake Headers)

Перед переключением протокола (Upgrade) с HTTP на WebSockets клиент Dio / WebSocketChannel обязан передать валидные метаданные.

Заголовок Обязательный Описание Пример значения
Upgrade Да Системный заголовок для переключения протокола websocket
Connection Да Системный заголовок для переключения протокола Upgrade
Authorization Да Короткоживущий Access Token для валидации сессии пользователя Bearer eyJhbGciOiJIUzI1Ni...
X-Request-ID Да Сквозной ID WebSocket-сессии для агрегации логов в Jaeger/Loki stream-7a2b-9c1d-00ef

2.2 Спецификация структуры потоковых сообщений (Stream Payload)

В рамках одной открытой сессии клиент отправляет бинарные данные, а сервер возвращает структурированный JSON с результатами ИИ-классификации.

2.2.1 Входящий поток от клиента (Client-to-Server, Binary/Text)

Мобильное приложение отправляет кадры в бинарном формате для экономии трафика. Кадры сжимаются на стороне Flutter до разрешения 640x480 (дефолт для MobileNetV2).

Тип фрейма Обязательный Описание Формат данных
Binary Frame Да Сырые байты текущего кадра с камеры смартфона ByteBuf (image/jpeg)
Text Frame (Ping) Нет Системный фрейм поддержания соединения (каждые 5 сек) {"type": "PING"}

2.2.2 Исходящий поток от сервера (Server-to-Client, JSON)

Бэкенд возвращает результаты распознавания только тогда, когда ИИ-сервис преодолевает порог уверенности (Confidence Threshold > 0.85).

Поле Тип Обязательный Описание
event_type String Да Тип события потока
confidence Float Да Степень уверенности ИИ в распознавании объекта
predicted_sku String Да Распознанный технический код продукта питания
frame_status String Да Инструкция для UI (HOLD — удерживать, DETECTED — зафиксирован)

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

На диаграмме представлена сквозная архитектура стрим-контура (Версия 2). Бэкенд-шлюз FastAPI изолирован от ИИ-нагрузки: он работает как транслятор фреймов в Apache Kafka. Kafka выступает амортизатором (буфером), из которого ИИ-сервис асинхронно забирает кадры на инференс, предотвращая падение системы при высоком FPS.

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 Vision as product-vision-processor (AID)
    participant Redis as In-Memory Redis (STATE STORAGE)
    participant K as Apache Kafka Cluster
    participant Censor as censorship-control-worker (SECURITY)

    %% ЭТАП 1: ИНИЦИАЦИЯ МЕДИА-СЕССИИ
    User->>Nginx: GET /api/v2/stream/video/init (Установление WebRTC/WS)
    Nginx->>GW: Пересылка HTTP-запроса авторизации
    Note over GW: Шаг 2: Валидация прав и сессии пользователя<br/>Инъекция X-Request-ID (trace_id) и app_lang
    GW-->>Nginx: HTTP 200 OK (Разрешение на удержание медиа-потока)
    Nginx-->>User: HTTP 101 Switching Protocols

    %% ЭТАП 2: НЕПРЕРЫВНЫЙ ИНФЕРЕНС И ТРЕКИНГ ОБЪЕКТОВ
    Nginx->>Worker: Проброс и фиксация стабильного потока кадров
    loop Потоковая трансляция видеофреймов
        User-->>Worker: WS/WebRTC Stream: Сырые видеокадры (Raw Frames)
        Worker->>Vision: gRPC: TrackVideoFrames(FrameRequest)
        Vision->>Redis: EX / INCR: Атомарный скоринг ID детекций и треков
        Redis-->>Vision: Подтвержденное состояние уникальных объектов в кадре
        Vision-->>Worker: gRPC Response: Координаты рамок (Boxes) + Имена классов
        Worker-->>User: WS Stream: Подсветка Bounding Boxes на экране смартфона
    end

    %% ЭТАП 3: СТОП-СИГНАЛ И ФИНАЛИЗАЦИЯ В КАФКУ
    User->>Worker: WS Control Event: STOP_STREAM (Фиксировать содержимое)
    Worker->>Redis: Метод ReadAndClearSessionState() (Чтение накопленного списка)
    Redis-->>Worker: Агрегированный текстовый массив найденных классов продуктов
    Note over Worker: Шаг 12: Формирование плоской строки состава инвентаря<br/>Закрытие WebSocket / gRPC соединения сессии
    Worker->>K: Пуш в топик: bpds.inventory.in.receipt.upload
    Worker-->>User: WS Connection Closed (Сканирование успешно завершено)

    %% ЭТАП 4: ЦЕНЗУРА СВЕДЕННОГО РЕЗУЛЬТАТА
    K->>Censor: Handler: ProcessUploadText()
    Note over Censor: Исход 14: Проверка чиста (is_clean = True)
    Censor->>K: Пуш в топик: bpds.inventory.out.receipt.parsed

4 Расшифровка шагов

4.1 Расшифровка шагов сквозного процесса

Шаг Действие Параметры / Запросы Ошибки (Исключения / Статусы)
Шаг 1 (User -> Nginx) Мобильный клиент (модуль камеры) запрашивает инициализацию медиа-сессии. HTTP GET /api/v2/stream/video/init
Headers: Bearer JWT, Accept: text/event-stream
DioException: connection timeout
HTTP 401 Unauthorized
Шаг 2 (Nginx -> GW) Прокси-сервер пересылает входящий HTTP-запрос авторизации на шлюз backend-api. Пересылка исходных заголовков + инъекция X-Request-ID и app_lang Nginx: 502 Bad Gateway
Nginx: 504 Gateway Timeout
Шаг 3 (GW -> GW) Внутреннее действие: Шлюз валидирует права/сессию и генерирует контекст трассировки. Проверка токена в Redis Token Blacklist.
Инъекция trace_id
Отказ Redis: периметр заблокирован (Fail-Close) → HTTP 503 Service Unavailable
Шаг 4 (GW -> Nginx) Шлюз подтверждает успешную валидацию периметра и авторизацию сессии. HTTP 200 OK
Разрешение на удержание медиа-потока
HTTP 500 Internal Server Error
Шаг 5 (Nginx -> User) Nginx выполняет апгрейд протокола, устанавливая постоянное двустороннее соединение. HTTP 101 Switching Protocols
Переключение на WebRTC / WebSocket транспорт
Nginx: 500 Protocol Switch Failed
Шаг 6 (Nginx -> Worker) Прокси пробрасывает и жестко закрепляет открытый сокет за свободным инстансом воркера. Изоляция сетевого сокета на уровне инфраструктуры Nginx: 503 Service Unavailable (нет свободных воркеров)
Шаг 7 (User -> Worker) Камера клиента начинает непрерывную трансляцию сырых видеокадров. WS/WebRTC Stream: Raw Frames
Бинарный поток фреймов высокого разрешения
Обрыв связи: WebSocketDisconnect
NetworkException
Шаг 8 (Worker -> Vision) Воркер транслирует видеокадры ИИ-процессору для непрерывного поиска объектов. Метод: gRPC TrackVideoFrames(FrameRequest)
Контекст: trace_id, gRPC Deadline = 4.0s
Превышение лимита: gRPC: DeadlineExceeded
Включение ИИ-Fallback (уход на модерацию)
Шаг 9 (Vision -> Redis) ИИ-процессор производит атомарный скоринг ID детекций и траекторий уникальных продуктов. Вызов Lua-скрипта:
Команды EX (TTL) / INCR (Инкремент счетчиков)
Отказ кэша: RedisConnectionError
Блокировка инференса трекера
Шаг 10 (Redis -> Vision) Оперативная память возвращает подтвержденное состояние уникальных объектов, убирая дребезг. Результат выполнения Lua-скрипта: актуальный стейт RedisDataCorruptedException
Шаг 11 (Vision -> Worker) ИИ-процессор возвращает воркеру координаты рамок и распознанные классы. gRPC Response
Payload: Boxes (x, y, w, h), class_names
gRPC: Internal Error
Шаг 12 (Worker -> User) Воркер возвращает метаданные детекции на фронтенд для отрисовки интерфейса. WS Stream
Payload: массив Bounding Boxes для рендеринга
WS Connection Broken
Шаг 13 (User -> Worker) Пользователь завершает съемку и отправляет управляющий стоп-сигнал фиксации. WS Control Event: STOP_STREAM
Команда на финализацию сессии
WS Event Parsing Error
Шаг 14 (Worker -> Redis) Воркер инициирует вычитку итогового состава и очистку временных ключей сессии. Метод: ReadAndClearSessionState() Redis: KeyNotFoundError
Шаг 15 (Redis -> Worker) Redis отдает накопленный за время стрима текстовый массив найденных классов продуктов. Payload: Array of unique product classes RedisConnectionError
Шаг 16 (Worker -> Worker) Внутреннее действие: Воркер сводит массив в плоскую строку состава инвентаря и высвобождает ресурсы. Конкатенация строки инвентаря,
закрытие gRPC и WebSocket соединений
StringFormattingException
Шаг 17 (Worker -> KAFKA) Воркер публикует плоскую строку состава инвентаря во входящий топик чеков брокера. Топик: bpds.inventory.in.receipt.upload
Параметры: max.block.ms = 1000
Сбой шины: Kafka: TimeoutException
Сброс логов в Promtail для ручного разбора
Шаг 18 (Worker -> User) Воркер закрывает WebSocket соединение с клиентом, подтверждая успех операции. WS Connection Closed
Статус: Сканирование успешно завершено
WS Close Frame Error
Шаг 19 (KAFKA -> Censor) Воркер безопасности вычитывает задачу из топика для проведения контент-контроля. Топик: bpds.inventory.in.receipt.upload
Handler: ProcessUploadText()
Kafka: CommitFailedException
Шаг 20 (Censor -> Censor) Внутреннее действие: Цензор проверяет текст на спам и маты, определяя Confidence. Валидация текста через Llama (gRPC: EvalTextConfidence). Исход: Чисто / Уверен (is_clean = True, Сonfidence >= 0.85) Таймаут Ollama (4 сек): MDM Fallback (сброс задачи в топик модерации и WebSocket на фронт)
Шаг 21 (Censor -> KAFKA) Очищенный результат отправляется дальше по контуру веерной дистрибуции. Топик: bpds.inventory.out.receipt.parsed
Триггер для веера: Billing, Fridge, Audit
Сбой шины: KafkaProducerError → уход сообщения в DLQ нарушений

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.