sequenceDiagram
autonumber
actor User as Mobile Client (Microphone)
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)
%% PHASE 1: SESSION INITIATION AND PROTOCOL UPGRADE
User->>Nginx: GET /api/v2/stream/voice/init (WebSocket Request)
Nginx->>GW: Forward HTTP Request for Session Authorization
Note over GW: Step 2: User Session Validation<br/>Injected X-Request-ID (trace_id) and app_lang
GW-->>Nginx: HTTP 200 OK (Upgrade Permission Granted)
Nginx-->>User: HTTP 101 Switching Protocols
%% PHASE 2: STREAMING AND CHUNK PROCESSING
Nginx->>Worker: Establish and Maintain Active WebSocket Connection
loop Binary Audio Chunk Streaming
User-->>Worker: WS Stream: Binary Audio Data (Chunks)
Worker->>AudioAI: Bi-directional gRPC: TranscribeAudioChunkStream(Chunk)
Note over AudioAI: Step 7: Local Whisper Tiny Inference<br/>Intermediate Phrase Text Recognition
AudioAI-->>Worker: gRPC Response: Intermediate Text Tokens
Worker-->>User: WS Stream: On-the-fly UI Text Rendering
end
%% PHASE 3: FINALIZATION AND CENSORSHIP HANDOFF
User->>Worker: WS Control Event: CLOSE_STREAM (User Released Button)
Note over Worker: Step 10: Aggregate Chunks into Final Session Text<br/>Terminate WebSocket / gRPC Channels
Worker->>K: Push to Topic: bpds.inventory.in.receipt.upload
Worker-->>User: WS Connection Closed (Successfully Completed)
Метод 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).
Функциональное назначение
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.
The endpoint implements the logic for rapid structural injection of food assets into the “digital refrigerator” via unstructured textual inputs (e.g., via workspace fields or native mobile quick-select components).
Communication Protocol (HTTP Contract)
- Method:
POST - Route:
/api/v2/stream/voice/init - Data Format:
application/json
Request Body Specification
| Field | Type | Mandatory | Description |
|---|---|---|---|
raw_text_input |
String | Yes | Unstructured text input string (e.g., "Груши 2.5") |
Эндпоинт реализует логику быстрого добавления продуктов питания в «цифровой холодильник» посредством ввода неструктурированного текста (например, через строку АРМ или интерфейс быстрого подбора мобильного приложения).
Протокол взаимодействия (HTTP Контракт)
- Метод:
POST - Маршрут:
/api/v2/stream/voice/init - Формат данных:
application/json
Спецификация тела запроса (Request Body)
| Поле | Тип | Обязательный | Описание |
|---|---|---|---|
raw_text_input |
String | Да | Неструктурированная текстовая строка ввода (например: "Груши 2.5") |
Схема обработки запроса пользователя (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.
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 (Успешно завершено)
Расшифровка шагов
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.
| Step | Component (From) | Direction | Component (To) | Operation / Event Payload | Business Exceptions |
|---|---|---|---|---|---|
| 1 | Mobile Client | --> |
Nginx Proxy | WebSocket Request: GET /api/v2/stream/voice/init |
|
| 2 | Nginx Proxy | --> |
FastAPI Gateway | Forward HTTP Request for Session Authorization | |
| 3 | FastAPI Gateway | --> |
FastAPI Gateway | User Session Validation, Inject trace_id and app_lang |
ERR-INVALID-SESSION |
| 4 | FastAPI Gateway | <-- |
Nginx Proxy | Return HTTP 200 OK (Upgrade Permission Granted) |
|
| 5 | Nginx Proxy | <-- |
Mobile Client | Return HTTP 101 Switching Protocols |
|
| 6 | Nginx Proxy | --> |
purchase-stream-worker | Establish and Maintain Active WebSocket Connection | |
| 7 | Mobile Client | --> |
purchase-stream-worker | WS Stream: Ingest Binary Audio Data Chunks | |
| 8 | purchase-stream-worker | --> |
grpc-analytics | gRPC Stream: TranscribeAudioChunkStream(Chunk) |
|
| 9 | grpc-analytics | --> |
grpc-analytics | Whisper Tiny Inference, Intermediate Text Token Recognition | ERR-AUDIO-DECODE-FAILED |
| 10 | grpc-analytics | <-- |
purchase-stream-worker | gRPC Response: Return Intermediate Text Tokens | |
| 11 | purchase-stream-worker | <-- |
Mobile Client | WS Stream: On-the-fly UI Text Rendering | |
| 12 | Mobile Client | --> |
purchase-stream-worker | WS Control Event: Send CLOSE_STREAM Signals |
|
| 13 | purchase-stream-worker | --> |
purchase-stream-worker | Aggregate Chunks into Final Text, Terminate WS/gRPC Channels | No business errors |
| 14 | purchase-stream-worker | --> |
Apache Kafka Cluster | Push Final Session Text to Topic bpds.inventory.in.receipt.upload |
|
| 15 | purchase-stream-worker | <-- |
Mobile Client | Terminate WebSocket Connection (Successfully Completed) |
| Шаг | Компонент (Из) | Направление | Компонент (В) | Операция / Тело события | Бизнес-исключения |
|---|---|---|---|---|---|
| 1 | Мобильный клиент | --> |
Nginx Proxy | Запрос WebSocket: GET /api/v2/stream/voice/init |
|
| 2 | Nginx Proxy | --> |
FastAPI Gateway | Пересылка HTTP-запроса на авторизацию сессии | |
| 3 | FastAPI Gateway | --> |
FastAPI Gateway | Валидация сессии пользователя, инъекция trace_id и app_lang |
ERR-INVALID-SESSION |
| 4 | FastAPI Gateway | <-- |
Nginx Proxy | Возврат HTTP 200 OK (Разрешение на Upgrade) |
|
| 5 | Nginx Proxy | <-- |
Мобильный клиент | Возврат HTTP 101 Switching Protocols |
|
| 6 | Nginx Proxy | --> |
purchase-stream-worker | Проброс и удержание WebSocket соединения | |
| 7 | Мобильный клиент | --> |
purchase-stream-worker | WS Stream: Передача бинарных чанков аудио-данных | |
| 8 | purchase-stream-worker | --> |
grpc-analytics | gRPC-стрим: TranscribeAudioChunkStream(Chunk) |
|
| 9 | grpc-analytics | --> |
grpc-analytics | Инференс Whisper Tiny, распознавание промежуточного текста | ERR-AUDIO-DECODE-FAILED |
| 10 | grpc-analytics | <-- |
purchase-stream-worker | gRPC Response: Возврат промежуточных токенов текста | |
| 11 | purchase-stream-worker | <-- |
Мобильный клиент | WS Stream: Отображение распознанного текста на лету в интерфейсе | |
| 12 | Мобильный клиент | --> |
purchase-stream-worker | WS Control Event: Управляющее событие CLOSE_STREAM |
|
| 13 | purchase-stream-worker | --> |
purchase-stream-worker | Сборка всех чанков в финальный текст, закрытие каналов WS/gRPC | No business errors |
| 14 | purchase-stream-worker | --> |
Apache Kafka Cluster | Пуш финального текста в топик bpds.inventory.in.receipt.upload |
|
| 15 | purchase-stream-worker | <-- |
Мобильный клиент | Успешное закрытие WebSocket соединения |
Спецификация сообщений брокера и вилок исключений
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 Topic Specification: media.camera.stream.v2
The FastAPI backend gateway continuously flushes frames received from binary WebSocket frames into this topic. To optimize throughput, data is transmitted in binary format (ByteArray serialization), and metadata is packed into message headers (Kafka Headers).
- Partitioning Strategy (Partition Key):
session_id(guarantees that all frames within a single camera session route to the same partition and are processed by the AI in strict chronological order). - Key Format:
String(the value ofsession_id). - Value Format:
Byte[](raw binary stream of the compressed JPEG image). - Message Headers Configuration (Kafka Headers):
| Header Key | Data Type | Description | Example Value |
|---|---|---|---|
x-request-id |
String | End-to-end session ID for distributed tracing | stream-7a2b-9c1d-00ef |
frame_timestamp |
String (ISO) | Frame capture timestamp from the mobile device | 2026-07-01T21:10:05.123Z |
user_id |
String | Identifier of the session owner | usr-8822-fa41 |
Control Text Frames Specification (JSON)
Client Frame Freeze Command (Step 13)
Sent to the WebSocket channel in text format when the user presses the shutter button.
{
"action": "FREEZE_FRAME",
"timestamp": "2026-07-01T21:10:12.450Z",
"payload": {
"verified_sku": "TOMATO",
"confidence_at_freeze": 0.942
}
}Server Response on Session Freeze (Step 15)
The final system message from FastAPI, confirming a successful write to the buffer and stream closure.
{
"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"
}
}Спецификация топика 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 |
Спецификация управляющих текстовых фреймов (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"
}
}Failures & Fail-Safe Strategies
The streaming loop is sensitive to network latency and connection stability. The system handles three primary emergency scenarios:
Scenario A: AI Service Outage or gRPC Failure (Inference Hang)
If the AI service goes offline or stops responding to FastAPI gateway gRPC requests, a Client-Side Timeout / Fallback strategy is triggered.
- System Behavior: If no
STREAM_PREDICTION_UPDATEDevent is received from the server within 3 seconds, the Flutter application smoothly switches the interface to manual capture mode. A toast notification is displayed on the screen, and the shutter button changes its logic: upon click, it sends a standard batchHTTP POST /product-photomatching the Version 1 scenario. The system does not block the user.
Scenario B: Buffer Overflow or Kafka Outage (HTTP 503)
Thrown at Step 5 if the Kafka broker is unavailable or the backend’s internal buffer overflows due to an abnormally high client FPS.
{
"error_code": "ERR-STREAM-BUFFER-OVERFLOW",
"message": "Camera stream is temporarily overloaded. Please hold the camera still.",
"details": {
"reason": "Kafka broker partition leader is unavailable or local producer memory buffer is full."
}
}Scenario C: Network Disconnection During Streaming (Close Code)
If the user loses 4G/Wi-Fi connection at Step 4, the WebSocket session is abnormally terminated by the operating system.
- System Behavior:
FastAPIcatches the systemWebSocketDisconnectevent. The backend immediately pushes a command to Kafka to purge the intermediate unprocessed frames of this session, cleaning up the partition. The draft transaction in thePostgreSQLdatabase is rolled back. Once the network recovers, the Flutter application generates a newsession_idand restarts the process from Step 2.
Вилки исключений и стратегии обработки ошибок (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.