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
Метод POST /api/v2/stream/voice/init
Документация API: Текстовое добавление номенклатур с автоматической gRPC ИИ-модерацией
- 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).
1 Функциональное назначение
Эндпоинт реализует логику быстрого добавления продуктов питания в «цифровой холодильник» посредством ввода неструктурированного текста (например, через строку АРМ или интерфейс быстрого подбора мобильного приложения).
2 Протокол взаимодействия (HTTP Контракт)
- Метод:
POST - Маршрут:
/api/v2/stream/voice/init - Формат данных:
application/json
2.1 Спецификация тела запроса (Request Body)
| Поле | Тип | Обязательный | Описание |
|---|---|---|---|
raw_text_input |
String | Да | Неструктурированная текстовая строка ввода (например: "Груши 2.5") |
3 Схема обработки запроса пользователя (Mermaid)
4 Расшифровка шагов
| Шаг | Действие | Параметры / Запросы / 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 (агрегированный текстовый список превысил лимиты брокера) |
Шаг 13 (T_Results -> Censor) |
Воркер цензуры считывает из топика загрузок сформированный текстовый список классов объектов видеострима для проведения лингвистической проверки. | Kafka Consumer Poll Request:Group_ID: "censorship-workers", Topic: bpds.inventory.in.receipt.uploadPartition: 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: DeliveryTimeoutKafkaException: 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.