sequenceDiagram
autonumber
participant K as Broker: Apache Kafka
participant AudioAI as voice-speech-processor (AID)
participant S3 as MinIO Object Storage
K->>AudioAI: Step 4: Handler: ProcessVoiceStream()
activate AudioAI
critical Step 5: Download Multi-Media Audio Asset
AudioAI->>S3: HTTP GET /user-voice-streams/uploads/...
activate S3
S3-->>AudioAI: Binary audio stream byte buffer
deactivate S3
option Storage Failure (HTTP 404 Not Found)
Note over AudioAI: Terminate Cycle:<br/>Log error code IAD-VOICE-404
end
Note over AudioAI: Whisper Tiny Model Inference:<br/>Acoustic frequency analysis & text synthesis
critical Step 6: Route Transcribed Text to Censorship Queue
AudioAI->>K: Push to topic: bpds.inventory.in.receipt.upload
option Transcription Failure / Audio Corrupted
Note over AudioAI: Terminate Cycle:<br/>Log error code IAD-WHISPER-422
end
deactivate AudioAI
Метод ProcessVoiceStream
Домен: AID | Сервис: voice-speech-processor | Тип: Kafka Consumer
В открытом доступе представлена демонстрационная версия метода. В настоящей публичной документации отображены не все шаги, технические сценарии и приватные эндпоинты для системы цифровых симуляторов бизнес-процессов.
- Полная спецификация метода: Будет доступна только во внутреннем контуре разработки (Confluence / Swagger Enterprise).
Функциональное назначение
The asynchronous event handler ProcessVoiceStream is deployed within the voice-speech-processor microservice. It acts as a core Audio Context Provider, responsible for consuming voice message tasks from Kafka, streaming raw binary audio files (.mp3/.ogg) from MinIO object storage, and running neural network acoustic decoding (Whisper Tiny inference).
The method transforms sound waves into a raw text string. To maintain structural consistency and reuse core perimeter policies, the resulting text is seamlessly pushed into the unified text processing pipeline, triggering the existing censorship-control-worker downstream.
Core Tasks Handled by the Method
- Audio Asset Ingestion: Streams down the target voice note recording from the S3 bucket into local execution memory layers via an asynchronous data channel.
- Speech-to-Text Transcription (Whisper Inference): Decodes the multi-media audio matrix using hardware acceleration (GPU/CPU) to remove noise, analyze acoustic frequencies, and synthesize a raw text string.
- Perimeter Continuity Delegation: Packages the transcribed text dataset string and forwards it to the core receipt queue under the original trace context.
Interaction Protocol (Event Contract)
- Method Name:
ProcessVoiceStream - Message Broker:
Apache Kafka - Ingress Topic (Subscription):
bpds.inventory.in.voice.stream - Egress Topic (Publication):
bpds.inventory.in.receipt.upload - Data Format:
application/json
Ingress Message Payload Example
{
"audio_file_url": "https://minio.internal",
"app_lang": "ru-RU",
"x_request_id": "trace-stt-censor-uuid"
}Асинхронный метод-обработчик (Event Handler) ProcessVoiceStream развернут внутри микросервиса voice-speech-processor. Он выполняет архитектурную роль поставщика аудио-контекста (Audio Context Provider) и отвечает за вычитку задач из брокера Kafka, потоковую загрузку аудиосообщений (.mp3/.ogg) из объектного хранилища MinIO, акустическое декодирование звуковых волн и автоматический перевод речи в текст с помощью нейросети Whisper Tiny.
Метод транслирует звуковой поток в сырую текстовую строку. Для сохранения сквозной цепочки обработки и повторного использования единого контура защиты, извлеченный текст автоматически отправляется в общую текстовую очередь, активируя существующий воркер censorship-control-worker.
Критические задачи метода
- Скачивание медиа-ассета: Потоковая загрузка бинарной аудиозаписи из S3-бакета в локальные слои оперативной памяти с помощью асинхронного конвейера данных.
- Распознавание речи (Whisper Инференс): Декодирование аудиопотока с использованием аппаратного ускорения (GPU/CPU) для очистки от шумов, анализа частот звуковой волны и синтеза сырой строки текста.
- Делегирование в контур фильтрации: Упаковка транскрибированного текста в стандартизированное DTO и его публикация в общую текстовую очередь под оригинальным сквозным контекстом трассировки.
Протокол взаимодействия и триггеры (Event Contract)
- Имя метода:
ProcessVoiceStream - Брокер сообщений:
Apache Kafka - Входной топик (Subscription):
bpds.inventory.in.voice.stream - Выходной топик (Publication):
bpds.inventory.in.receipt.upload - Формат данных:
application/json
Пример входного события (Ingress Message Payload)
{
"audio_file_url": "https://minio.internal",
"app_lang": "ru-RU",
"x_request_id": "trace-stt-censor-uuid"
}Диаграмма последовательности (Mermaid)
sequenceDiagram
autonumber
participant K as Брокер: Apache Kafka
participant AudioAI as voice-speech-processor (AID)
participant S3 as MinIO Object Storage
K->>AudioAI: Шаг 4: Handler: ProcessVoiceStream()
activate AudioAI
critical Шаг 5: Скачивание мультимедийного аудиофайла
AudioAI->>S3: HTTP GET /user-voice-streams/uploads/...
activate S3
S3-->>AudioAI: Бинарный поток байт аудиофайла
deactivate S3
option Сбой хранилища (HTTP 404 Not Found)
Note over AudioAI: Экстренное прерывание:<br/>Фиксация ошибки с кодом IAD-VOICE-404
end
Note over AudioAI: Инференс модели Whisper Tiny:<br/>Анализ частот звука и синтез текста
critical Шаг 6: Перенаправление распознанного текста в топик цензуры
AudioAI->>K: Пуш в топик: bpds.inventory.in.receipt.upload
option Сбой транскрибации / Поврежден аудиопоток
Note over AudioAI: Экстренное прерывание:<br/>Фиксация ошибки с кодом IAD-WHISPER-422
end
deactivate AudioAI
Расшифровка шагов
| Step | Action | Parameters / Requests / DTO | Errors (Exceptions / Statuses) |
|---|---|---|---|
4 (K -> AudioAI) |
The worker polls the speech analysis event from the message broker queue, maps properties, and provisions resources to load the multi-media binary buffer. | Kafka Ingress Payload: The JSON request contract metadata properties specified above. |
No business errors |
5 (AudioAI -> S3) |
The microservice sends a request to the S3 bucket using secured internal tokens to stream the voice file directly into local thread memory layers. | HTTP GET Request: URL: https://minio.internal |
IAD-VOICE-404: Target source multi-media voice note file asset was not found or deleted from storage before processing. |
6 (AudioAI -> K) |
The hardware-accelerated Whisper Tiny core successfully decodes audio waves. The worker packages the raw text string into a standard input DTO and publishes it to the broker. | Kafka Egress Message: Topic: bpds.inventory.in.receipt.uploadPayload: { "raw_text_input": "Текст, который распознал Whisper из аудиозаписи", "app_lang": "ru-RU", "x_request_id": "trace-stt-censor-uuid" } |
IAD-WHISPER-422: Audio encoding is unreadable, corrupted due to extreme ambient noise, or inference core crashed. |
| Шаг | Действие | Параметры / Запросы / DTO | Ошибки (Исключения / Статусы) |
|---|---|---|---|
4 (K -> AudioAI) |
Воркер вычитывает событие анализа речи из очереди брокера сообщений, извлекает JSON-контекст и выделяет системные ресурсы под новую задачу декодирования. | Kafka Ingress Payload: JSON-контракт со свойствами и метаданными, описанными выше в примере. |
Бизнес-ошибки отсутствуют |
5 (AudioAI -> S3) |
Микросервис обращается к бакету S3, используя внутренние защищенные токены, для прямой загрузки голосовой заметки в память инференса. | HTTP GET Request: URL: https://minio.internal |
IAD-VOICE-404: Целевой мультимедийный файл голосовой заметки не найден или был удален из S3 до начала чтения воркером. |
6 (AudioAI -> K) |
Аппаратно-ускоренное ядро Whisper Tiny успешно декодирует звуковые волны. Скрипт упаковывает полученный текст в стандартный формат и публикует событие в брокер. | Kafka Egress Message: Топик: bpds.inventory.in.receipt.uploadPayload: { "raw_text_input": "Текст, который распознал Whisper из аудиозаписи", "app_lang": "ru-RU", "x_request_id": "trace-stt-censor-uuid" } |
IAD-WHISPER-422: Аудиопоток не поддается чтению, критически поврежден из-за экстремального уровня шума или произошел сбой ядра инференса. |