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
Метод POST /api/v2/stream/video/init
Документация API: Потоковая верификация продуктов через реальное время (Camera Stream) с использованием буферной очереди Apache Kafka
- 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 Функциональное назначение
Метод предназначен для обработки потока изображений (кадров) в реальном времени, поступающих непосредственно с камеры мобильного приложения во время наведения на продукт. В отличие от пакетной загрузки готовых файлов, данный контур минимизирует задержку ввода (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.
4 Расшифровка шагов
4.1 Расшифровка шагов сквозного процесса
| Шаг | Действие | Параметры / Запросы | Ошибки (Исключения / Статусы) |
|---|---|---|---|
Шаг 1 (User -> Nginx) |
Мобильный клиент (модуль камеры) запрашивает инициализацию медиа-сессии. | HTTP GET /api/v2/stream/video/initHeaders: Bearer JWT, Accept: text/event-stream |
DioException: connection timeoutHTTP 401 Unauthorized |
Шаг 2 (Nginx -> GW) |
Прокси-сервер пересылает входящий HTTP-запрос авторизации на шлюз backend-api. |
Пересылка исходных заголовков + инъекция X-Request-ID и app_lang |
Nginx: 502 Bad GatewayNginx: 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Бинарный поток фреймов высокого разрешения |
Обрыв связи: WebSocketDisconnectNetworkException |
Шаг 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 ResponsePayload: Boxes (x, y, w, h), class_names |
gRPC: Internal Error |
Шаг 12 (Worker -> User) |
Воркер возвращает метаданные детекции на фронтенд для отрисовки интерфейса. | WS StreamPayload: массив 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.uploadHandler: 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.