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-интеграциями
В открытом доступе представлена демонстрационная версия метода. В настоящей публичной документации отображены не все шаги, технические сценарии и приватные эндпоинты для системы цифровых симуляторов бизнес-процессов.
- Полная спецификация метода: Будет доступна только во внутреннем контуре разработки (Confluence / Swagger Enterprise).
Функциональное назначение
This documentation section is currently under development and may contain incomplete data. Some technical descriptions, parameters, and system operation scenarios for the digital business process simulator are subject to change.
- Stable documentation version: Will be published upon final completion and validation of the method code.
Метод предназначен для обработки потока изображений (кадров) в реальном времени, поступающих непосредственно с камеры мобильного приложения во время наведения на продукт. В отличие от пакетной загрузки готовых файлов, данный контур минимизирует задержку ввода (Latency) и защищает вычислительные ресурсы ИИ-сервера (PyTorch/GPU) от перегрузки при высокой частоте кадров (FPS).
Бэкенд-шлюз FastAPI в этой схеме выступает в роли ультрабыстрого легковесного ретранслятора, который сгружает бинарные кадры в брокер сообщений Apache Kafka, выполняющий роль амортизатора нагрузки (Load Smoother). ИИ-сервис асинхронно вычитывает кадры из очереди, производит инференс и возвращает результат распознавания обратно в WebSocket-сессию пользователя.
Протокол взаимодействия (Потоковый Контракт)
- Протокол:
WebSockets (WSS) - Маршрут (Эндпоинт):
/api/v2/fridge/stream-capture/{session_id} - Тип обмена: Дуплексный (Двунаправленный асинхронный поток)
- Базовый формат кадров (Frame Payload): Бинарный поток сжатых изображений (
image/jpegилиimage/webp) + JSON-управляющие сообщения.
Спецификация заголовков при установлении соединения (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 |
Спецификация структуры потоковых сообщений (Stream Payload)
В рамках одной открытой сессии клиент отправляет бинарные данные, а сервер возвращает структурированный JSON с результатами ИИ-классификации.
Входящий поток от клиента (Client-to-Server, Binary/Text)
Мобильное приложение отправляет кадры в бинарном формате для экономии трафика. Кадры сжимаются на стороне Flutter до разрешения 640x480 (дефолт для MobileNetV2).
| Тип фрейма | Обязательный | Описание | Формат данных |
|---|---|---|---|
Binary Frame |
Да | Сырые байты текущего кадра с камеры смартфона | ByteBuf (image/jpeg) |
Text Frame (Ping) |
Нет | Системный фрейм поддержания соединения (каждые 5 сек) | {"type": "PING"} |
Исходящий поток от сервера (Server-to-Client, JSON)
Бэкенд возвращает результаты распознавания только тогда, когда ИИ-сервис преодолевает порог уверенности (Confidence Threshold > 0.85).
| Поле | Тип | Обязательный | Описание |
|---|---|---|---|
event_type |
String | Да | Тип события потока |
confidence |
Float | Да | Степень уверенности ИИ в распознавании объекта |
predicted_sku |
String | Да | Распознанный технический код продукта питания |
frame_status |
String | Да | Инструкция для UI (HOLD — удерживать, DETECTED — зафиксирован) |
Диаграмма последовательности (Mermaid)
This documentation section is currently under development and may contain incomplete data. Some technical descriptions, parameters, and system operation scenarios for the digital business process simulator are subject to change.
- Stable documentation version: Will be published upon final completion and validation of the method code.
На диаграмме представлена сквозная архитектура стрим-контура (Версия 2). Бэкенд-шлюз FastAPI изолирован от ИИ-нагрузки: он работает как транслятор фреймов в Apache Kafka. Kafka выступает амортизатором (буфером), из которого ИИ-сервис асинхронно забирает кадры на инференс, предотвращая падение системы при высоком FPS.
Расшифровка шагов
This documentation section is currently under development and may contain incomplete data. Some technical descriptions, parameters, and system operation scenarios for the digital business process simulator are subject to change.
- Stable documentation version: Will be published upon final completion and validation of the method code.
| Шаг | Действие | Параметры / Запросы / DTO | Ошибки (Исключения / Статусы) |
|---|---|---|---|
Шаг 1 (App -> Nginx) |
Мобильное приложение отправляет запрос на согласование параметров медиа-сессии и инициацию сигнального сокета для трансляции видео. | HTTP GET /api/v2/video-stream/initHeaders: Upgrade: websocket, Connection: Upgrade, Authorization: Bearer <JWT> |
DioException: connection timeoutHTTP 400 Bad Request (некорректные заголовки или параметры рукопожатия)HTTP 401 Unauthorized |
Шаг 2 (Nginx -> FastAPI_Gateway) |
Прокси-сервер перенаправляет запрос на шлюз, выполняя инъекцию сквозного X-Request-ID и пробрасывая метаданные языкового контекста. |
HTTP GET /internal/v2/video-stream/authorizeAdded 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 refusedRedis.TimeoutError: Command timed outRedis.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 refusedRedis.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: DeliveryTimeoutKafkaException: MessageTooLargeException (агрегированный текстовый список превысил лимиты брокера) |
Спецификация сообщений брокера
This documentation section is currently under development and may contain incomplete data. Some technical descriptions, parameters, and system operation scenarios for the digital business process simulator are subject to change.
- Stable documentation version: Will be published upon final completion and validation of the method code.
Спецификация топика 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 |
This documentation section is currently under development and may contain incomplete data. Some technical descriptions, parameters, and system operation scenarios for the digital business process simulator are subject to change.
- Stable documentation version: Will be published upon final completion and validation of the method code.
Спецификация управляющих текстовых фреймов (JSON)
Команда фиксации кадра от клиента (Шаг 13)
Отправляется в WebSocket-канал в текстовом формате при нажатии пользователем кнопки затвора.
{
"action": "FREEZE_FRAME",
"timestamp": "2026-07-01T21:10:12.450Z",
"payload": {
"verified_sku": "TOMATO",
"confidence_at_freeze": 0.942
}
}Ответ сервера о заморозке сессии (Шаг 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"
}
}This documentation section is currently under development and may contain incomplete data. Some technical descriptions, parameters, and system operation scenarios for the digital business process simulator are subject to change.
- Stable documentation version: Will be published upon final completion and validation of the method code.
Вилки исключений и стратегии обработки ошибок (Failures & Fail-Safe)
Стриминговый контур чувствителен к сетевым задержкам и стабильности соединений. Система обрабатывает три основных аварийных сценария:
Сценарий А: Отказ ИИ-сервиса или падение gRPC (Зависание инференса)
Если ИИ-сервис уходит в offline или перестает отвечать на gRPC-запросы шлюза FastAPI, срабатывает стратегия Client-Side Timeout / Fallback.
- Поведение системы: Если в течение 3 секунд от сервера не приходит ни одного обновления
STREAM_PREDICTION_UPDATED, Flutter-приложение мягко переключает интерфейс в режим ручной съемки. На экран выводится тост-уведомление, а кнопка затвора меняет логику: при нажатии отправляется стандартный пакетныйHTTP POST /product-photoпо сценарию Версии 1. Система не блокирует пользователя.
Сценарий Б: Ошибка переполнения буфера или отказ 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."
}
}Сценарий В: Разрыв сетевого соединения в процессе стриминга (Close Code)
Если у пользователя пропадает 4G/Wi-Fi соединение на Шаге 4, WebSocket-сессия аварийно закрывается операционной системой.
- Поведение системы:
FastAPIловит системное событиеWebSocketDisconnect. Бэкенд немедленно отправляет команду в Kafka на удаление промежуточных необработанных кадров этой сессии, очищая партицию. В базе данныхPostgreSQLтранзакция черновика аннулируется. При восстановлении сети Flutter-приложение генерирует новыйsession_idи начинает процесс с Шага 2.