Метод (Event Handler) ProcessUploadText

Спецификация асинхронного обработчика | Сервис: censorship-control-worker

WarningОграничение публичной документации

В открытом доступе представлена демонстрационная версия метода. В настоящей публичной документации отображены не все шаги, технические сценарии и приватные эндпоинты для системы цифровых симуляторов бизнес-процессов.

  • Полная спецификация метода: Будет доступна только во внутреннем контуре разработки (Confluence / Swagger Enterprise).

Функциональное назначение

Этот документ описывает внутреннюю логику, алгоритм ветвления и интеграции асинхронного обработчика (Event Handler) ProcessUploadText, развернутого внутри воркера censorship-control-worker.

Данный хэндлер является универсальной сквозной точкой очистки и цензурирования данных в системе. Он изолирован от синхронного жизненного цикла HTTP/gRPC запросов и повторно используется для обработки текста, поступающего из различных каналов ввода мобильного приложения (ручной текстовый ввод, распознавание фото-чеков, голосовой ввод и т.д.).

Метод решает следующие задачи:

  1. Оркестрация ИИ-цензуры (LLM Orchestration): Взаимодействие с локальным контейнером нейросети (Ollama) через gRPC для выявления обсценной лексики и спама.
  2. Анти-фрод и контроль лимитов (Abuse Control): Атомарный учет количества нарушений пользователя с помощью инкрементальных Lua-скриптов в Redis.
  3. Маршрутизация событий (Event Routing): Распределение результатов обработки по специализированным топикам Kafka в зависимости от флагов валидации и уровня уверенности (Confidence Score) нейросети.

Протокол взаимодействия и триггеры (Event Contract)

  • Тип метода: Event Handler / Kafka Consumer
  • Брокер сообщений: Apache Kafka
  • Входной топик (Subscription): bpds.inventory.in.receipt.upload
  • Формат данных: Avro / JSON (в зависимости от конфигурации сериализатора)

Спецификация входного события (Ingress Message Payload)

Хэндлер ожидает в топике стандартизированную структуру, сгенерированную шлюзом или другими сервисами-продюсерами:

Поле Тип Обязательный Описание Пример значения
trace_id String Да Сквозной ID запроса (из заголовка X-Request-ID) req-manual-text-99aa
user_id String Да Уникальный идентификатор пользователя usr_9876
source_type String Да Источник данных (MANUAL_TEXT, OCR_PHOTO, VOICE) MANUAL_TEXT
raw_text_input String Да Текст для анализа Молоко Юрта в ауле 3.2% 1.5 л
app_lang String Да Язык приложения для корректного выбора словаря ИИ ru

Пример JSON-события из топика:

{
  "trace_id": "req-manual-text-99aa",
  "user_id": "usr_9876",
  "source_type": "MANUAL_TEXT",
  "raw_text_input": "Молоко Юрта в ауле 3.2% 1.5 л 1029 тенге",
  "app_lang": "ru"
}

Диаграмма последовательности (Mermaid)

Ниже представлена изолированная логика работы хэндлера после вычитки сообщения из брокера.

sequenceDiagram
    autonumber
    participant K_In as Kafka: ...receipt.upload
    participant Censor as censorship-control-worker
    participant Ollama as Ollama Container (AID)
    participant Redis as In-Memory Redis (SECURITY)
    participant K_Out as Kafka Cluster (Out Topics)

    %% ВЫЧИТКА
    K_In->>Censor: Handler: ProcessUploadText()
    Note over Censor: Извлечение метаданных:<br/>trace_id, user_id, raw_text_input, app_lang

    %% БЛОК НЕЙРОСЕТИ
    critical Анализ текста через нейросеть (LLM Orchestration)
        Censor->>Ollama: gRPC: EvalTextConfidence(raw_text_input, app_lang)
        Ollama-->>Censor: Ответ: is_profane (bool), confidence_score (float)
    option Сбой Ollama (Таймаут / Падение контейнера)
        Note over Censor: Режим отказоустойчивости (Fallback):<br/>Аварийное перенаправление на ручной разбор
        Censor->>K_Out: Пуш в топик: bpds.mdm.in.product.process
    end

    %% БЛОК REDIS
    critical Проверка и инкремент лимитов нарушений в кэше
        Censor->>Redis: Вызов Lua-скрипта (Проверка словаря + INCR счетчик)
        Redis-->>Censor: Текущий count_violation для user_id
    option Ошибка Redis (Connection Refused / Сбой сети)
        Note over Censor: Режим отказоустойчивости:<br/>Игнорируем лимиты флуда,<br/>продолжаем по оценке ИИ
    end

    %% АЛГОРИТМИЧЕСКОЕ ВЕТВЛЕНИЕ
    alt Исход 5а: Текст чист + Высокая уверенность ИИ (Confidence >= 0.85)
        Note over Censor: Идеальный сценарий
        Censor->>K_Out: Пуш в топик: bpds.inventory.out.receipt.parsed
        
    else Исход 5б: Текст чист + ИИ сомневается (Confidence < 0.85)
        Note over Censor: Требуется верификация данных
        Censor->>K_Out: Пуш в топик: bpds.mdm.in.product.process
        Note over K_Out: Отправка draft_data<br/>на Вкладку модерации
        
    else Исход 5в: Обнаружен мат или спам (is_profane == true)
        Note over Censor: Зафиксировано нарушение правил платформы
        Censor->>K_Out: Пуш в топик: bpds.aid.out.profanity.violate<br/>(Payload: violation_count)
    end

sequenceDiagram
    autonumber
    actor User as Мобильный клиент (App)
    participant GW as API Gateway / Producer
    participant K_In as Kafka: bpds.inventory.in.receipt.upload
    participant Censor as censorship-control-worker
    participant Ollama as Ollama Container (AID)
    participant Redis as In-Memory Redis (SECURITY)
    participant K_Out as Kafka Cluster (Out Topics)

    %% ШАГ 1 & 2: ПУБЛИКАЦИЯ И ВЫЧИТКА
    User->>GW: Отправка текста (Запрос/Ввод)
    Note over GW: Шаг 1: Валидация контракта,<br/>выделение локали и обогащение метаданных
    GW->>K_In: Пуш события (Topic: ...receipt.upload)
    K_In->>Censor: Шаг 2: Handler: ProcessUploadText() (Вычитка из очереди)

    %% ШАГ 3: НЕЙРОСЕТЬ
    critical Шаг 3: Семантический анализ через ИИ
        Censor->>Ollama: HTTP POST /api/generate (Llama3 Model)
        Ollama-->>Censor: JSON: { is_clean: bool, confidence: float }
    option Ошибка Ollama (Timeout / Crash)
        Note over Censor: Fallback: Аварийное перенаправление в черновики<br/>Пуш в топик: bpds.mdm.in.product.process
    end

    %% ШАГ 4: REDIS СЛОВАРИ И ЛИМИТЫ
    critical Шаг 4: Локальные стоп-слова и инкремент счетчика в кэше
        Censor->>Redis: Redis Pipeline (SISMEMBER по dict:profanity + MULTI INCRBY)
        Redis-->>Censor: Текущий count_violation для user_id
    option Ошибка Redis (Connection Refused)
        Note over Censor: Режим отказоустойчивости:<br/>Игнорируем лимиты, доверяем только ИИ
    end

    %% ШАГ 5: АЛГОРИТМИЧЕСКОЕ ВЕТВЛЕНИЕ ИСХОДОВ
    alt Шаг 5а: Ветка «Успех» (Чисто + Уверенность ИИ >= 0.85)
        Censor->>K_Out: Пуш в топик: bpds.inventory.out.receipt.parsed
        Note over K_Out: Прямая запись в таблицу fridge_inventory (Direct Update)
        
    else Шаг 5б: Ветка «Подозрение» (Чисто + Уверенность ИИ < 0.85)
        Censor->>K_Out: Пуш в топик: bpds.mdm.in.product.process (draft_data)
        
    else Шаг 5в: Ветка «Нарушение» (Обнаружен мат или спам / is_clean == false)
        Note over Censor: Развилка Abuse_Check по метрике count_violation
        Censor->>K_Out: Пуш в топик: bpds.aid.out.profanity.violate<br/>(Payload: notification_type = FIRST_WARNING_PUSH)
    end

    %% ШАГ 6, 7 & 8: ЦИКЛ РАБОТЫ С ЧЕРНОВИКОМ И ПОВТОРНАЯ МОДЕРАЦИЯ
    opt Реализация сценария с Черновиком (Исходы из Шага 5б)
        K_Out-->>User: Шаг 6: gRPC/HTTP Stream трансляция черновика на вкладку модерации
        Note over User: Пользователь или модератор<br/>исправляет сомнительный текст руками
        User->>GW: Отправка исправленного текста
        GW->>K_In: Шаг 7: Повторный пуш в топик: ...receipt.upload (с тем же x_request_id)
        K_In->>Censor: Повторный цикл вычитки воркером
        
        alt Шаг 8а: Успех после модерации (Текст отредактирован чисто)
            Censor->>K_Out: Пуш в топик: bpds.inventory.out.receipt.parsed
        else Шаг 8б: Повторный мат при модерации (Счетчик >= 2)
            Censor->>K_Out: Пуш в топик: bpds.aid.out.profanity.violate<br/>(Payload: notification_type = BAN_OR_STRICT_ALERT)
        end
    end

Шаг Действие Параметры / Запросы / DTO Ошибки (Исключения / Статусы)
Шаг 1 (Gateway -> T_Results) API-шлюз проверяет длину входящей строки (len(text) <= 120), обогащает метаданными локали приложения, пробрасывает трассировочный ID и публикует событие в брокер. Kafka Message (Topic: bpds.inventory.in.receipt.upload):
Payload: { "raw_text_input": "Строка пользователя...", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
FastAPI.ValidationError (строка > 120 символов)
KafkaException: MessageTooLarge
KafkaException: QueueFullException
Шаг 2 (T_Results -> Censor) Фоновый воркер цензуры, подписанный на топик загрузки, вычитывает новую задачу из очереди для проведения лингвистического анализа. Kafka Consumer Poll Request:
Group_ID: "censorship-workers", Получен payload шага 1.
Offset: 452091
KafkaException: CommitFailedException (воркер завис, сработал max.poll.interval.ms)
SerializationException (битый json в пайлоаде)
Шаг 3 (Censor -> Ollama) Воркер отправляет сырой текст в локальный ИИ-контейнер для семантического анализа и определения вероятности наличия скрытого спама, мата или оскорблений. HTTP POST /api/generate (Ollama API):
{ "model": "llama3", "prompt": "Analyze the following text for profanity, hate speech, or spam. Return exact JSON: { 'is_clean': boolean, 'confidence': float }. Text: 'Строка пользователя...'", "stream": false }
ConnectTimeout: Ollama unreachable
HTTP 500 Internal Server Error (сбой аллокации VRAM внутри контейнера Ollama)
ReadTimeout: Inference duration exceeded
Шаг 4 (Censor -> Redis) Параллельно с ИИ воркер проверяет текст по жестким локальным словарям стоп-слов, загруженным под конкретный язык, и выполняет условный инкремент счетчика для пользователя. Redis Command Pipeline:
1. SISMEMBER "dict:profanity:ru-RU" "входящее_слово"
2. MULTI
3. INCRBY "abuse:counter:user_123" 1
4. EXPIRE "abuse:counter:user_123" 86400
5. EXEC
Redis.ConnectionError: Connection refused
Redis.TimeoutError: Command timed out
Redis.ClusterDownException (сбой репликации или сегментации нод Redis памяти)
Шаг 5а (Censor -> T_Clean_Text) Ветка «Успех»: Если стоп-слова не найдены, а ИИ выдал флаг чистоты с высокой степенью уверенности, текст отправляется в топик очищенных данных для дальнейшего парсинга. Kafka Message (Topic: bpds.inventory.out.receipt.parsed):
Payload: { "clean_text_input": "Строка пользователя...", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: DeliveryTimeout
KafkaException: NotCoordinatorForException (временный ребаланс брокеров во время пуша)
Шаг 5б (Censor -> T_Draft) Ветка «Подозрение»: Если локальные словари чисты, но ИИ сомневается в контексте (confidence < 0.85), воркер маркирует запись как сомнительный черновик для ручной модерации. Kafka Message (Topic: bpds.mdm.in.product.process):
Payload: { "draft_data": { "raw_id": "uuid", "text": "Строка пользователя..." }, "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: LeaderNotAvailableException
KafkaException: RecordTooLargeException
Шаг 5в (Censor -> Abuse_Check -> T_Notif) Первичное нарушение: Цензор обнаружил мат. Логический разветвитель проверяет счетчик в Redis. Если Счетчик == 1, генерируется событие первичного предупреждения с учетом языка. Kafka Message (Topic: bpds.aid.out.profanity.violate):
Payload: { "user_id": "user_123", "violation_count": 1, "app_lang": "ru-RU", "notification_type": "FIRST_WARNING_PUSH" }
KafkaException: BrokerNotAvailable
KafkaException: MessageTimedOut (сбой подтверждения доставки от реплик брокера)
Шаг 6 (T_Draft -> App) Мобильное приложение вычитывает из топика черновиков сомнительную запись и отображает её во вкладке модерации на родном языке пользователя для исправления. gRPC / HTTP Stream (вкладка модерации):
Payload: { "draft_id": "uuid-draft-555", "raw_text_input": "Сомнительный текст...", "app_lang": "ru-RU" }
DioException: connection error (клиент потерял сеть)
HTTP 401 Unauthorized (истекла сессия модератора в приложении)
Шаг 7 (App -> T_Results) Пользователь (или модератор) отправляет отредактированную, очищенную версию текста обратно в систему под тем же сквозным языковым контекстом. Kafka Message (Topic: bpds.inventory.in.receipt.upload):
Payload: { "raw_text_input": "Отредактированный чистый текст", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: QueueFullException
FastAPI.ValidationError (текст пустой или превысил лимит после редактирования)
Шаг 8а (Censor -> T_Clean_Text) Успешный исход модерации: Повторный цикл проверки цензором подтверждает, что отредактированный текст полностью чист. Запись уходит в финальный топик. Kafka Message (Topic: bpds.inventory.out.receipt.parsed):
Payload: { "clean_text_input": "Отредактированный чистый текст", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: DeliveryTimeout
KafkaException: ConcurrentModificationException
Шаг 8б (Abuse_Check -> T_Notif) Повторный мат при модерации: Если в отредактированном тексте снова найден мат, разветвитель фиксирует Счетчик >= 2. Пользователь помечается как ненадежный. Kafka Message (Topic: bpds.aid.out.profanity.violate):
Payload: { "user_id": "user_123", "violation_count": 2, "app_lang": "ru-RU", "notification_type": "BAN_OR_STRICT_ALERT" }
KafkaException: LeaderNotAvailableException
InvalidTopicException (топик заблокирован администратором кафки)
Шаг Действие Параметры / Запросы / DTO Ошибки (Исключения / Статусы)
Шаг 1 (Gateway -> T_Results) API-шлюз (или другой продюсер) валидирует первичный ввод, обогащает метаданными локали и публикует событие в брокер. Kafka Message (Topic: bpds.inventory.in.receipt.upload):
Payload: { "raw_text_input": "Строка текста...", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid", "user_id": "user_123" }
FastAPI.ValidationError
KafkaException: QueueFullException
Шаг 2 (T_Results -> Censor) Фоновый воркер цензуры вычитывает задачу из очереди для проведения лингвистического и семантического анализа. Kafka Consumer Poll Request:
Group_ID: "censorship-workers", Получен payload шага 1.
Offset: 452091
KafkaException: CommitFailedException
SerializationException
Шаг 3 (Censor -> Ollama) Воркер отправляет сырой текст в локальный ИИ-контейнер для определения вероятности наличия скрытого спама, мата или оскорблений. HTTP POST /api/generate (Ollama API):
{ "model": "llama3", "prompt": "Analyze for profanity. Return JSON: { 'is_clean': boolean, 'confidence': float }", "stream": false }
ConnectTimeout: Ollama unreachable
HTTP 500 Internal Server Error
ReadTimeout: Inference duration exceeded
Шаг 4 (Censor -> Redis) Параллельно воркер проверяет текст по жестким локальным словарям стоп-слов под конкретный язык и управляет счетчиком нарушений. Redis Command Pipeline:
1. SISMEMBER "dict:profanity:ru-RU"
2. MULTI
3. INCRBY "abuse:counter:user_123" 1
4. EXEC
Redis.ConnectionError
Redis.TimeoutError
Redis.ClusterDownException
Шаг 5а (Censor -> T_Clean_Text) Ветка «Успех»: Стоп-слова не найдены, ИИ уверен в чистоте (confidence >= 0.85). Текст уходит на прямой парсинг и добавление в холодильник. Kafka Message (Topic: bpds.inventory.out.receipt.parsed):
Payload: { "clean_text_input": "Чистая строка...", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: DeliveryTimeout
KafkaException: NotCoordinatorForException
Шаг 5б (Censor -> T_Draft) Ветка «Подозрение»: Локальные словари чисты, но ИИ сомневается (confidence < 0.85). Воркер маркирует запись как сомнительный черновик для модерации. Kafka Message (Topic: bpds.mdm.in.product.process):
Payload: { "draft_data": { "raw_id": "uuid", "text": "Сомнительная строка" }, "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: LeaderNotAvailableException
KafkaException: RecordTooLargeException
Шаг 5в (Censor -> Abuse_Check -> T_Notif) Ветка «Нарушение»: Обнаружен мат/спам. Проверяется счетчик нарушений из Redis. При Счетчик == 1 генерируется событие первичного предупреждения. Kafka Message (Topic: bpds.aid.out.profanity.violate):
Payload: { "user_id": "user_123", "violation_count": 1, "app_lang": "ru-RU", "notification_type": "FIRST_WARNING_PUSH" }
KafkaException: BrokerNotAvailable
KafkaException: MessageTimedOut
Шаг 6 (T_Draft -> App) Мобильное приложение вычитывает сомнительную запись из топика черновиков и отображает её во вкладке модерации для ручного исправления пользователем. gRPC / HTTP Stream (вкладка модерации):
Payload: { "draft_id": "uuid-draft-555", "raw_text_input": "Текст для редактирования", "app_lang": "ru-RU" }
DioException: connection error
HTTP 401 Unauthorized
Шаг 7 (App -> T_Results) Пользователь отправляет отредактированную версию текста обратно в систему под тем же сквозным x_request_id. Kafka Message (Topic: bpds.inventory.in.receipt.upload):
Payload: { "raw_text_input": "Отредактированный чистый текст", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: QueueFullException
FastAPI.ValidationError
Шаг 8а (Censor -> T_Clean_Text) Успех после модерации: Повторный цикл проверки подтверждает, что отредактированный текст полностью чист. Запись уходит в финальный топик. Kafka Message (Topic: bpds.inventory.out.receipt.parsed):
Payload: { "clean_text_input": "Отредактированный чистый текст", "app_lang": "ru-RU", "x_request_id": "req-censor-999-uuid" }
KafkaException: DeliveryTimeout
Шаг 8б (Abuse_Check -> T_Notif) Повторное нарушение: Если в отредактированном тексте снова найден мат, фиксируется Счетчик >= 2. Инициируется процедура блокировки/строгого алерта. Kafka Message (Topic: bpds.aid.out.profanity.violate):
Payload: { "user_id": "user_123", "violation_count": 2, "app_lang": "ru-RU", "notification_type": "BAN_OR_STRICT_ALERT" }
KafkaException: LeaderNotAvailableException
InvalidTopicException

Appendix

Общее описание компонента

Для обеспечения защиты платформы от флуда, обсценной лексики и спама на этапе ручного или автоматического (OCR/Voice) ввода, внутри воркера censorship-control-worker используется распределенный кэш Redis In-Memory DB.

Проверка локальных словарей стоп-слов и инкремент счетчика нарушений выполняются с помощью Lua-скрипта. Использование Lua-скрипта гарантирует атомарность (Atomicity) операции на стороне единого потока Redis (Single-Threaded). Это исключает состояние гонки (Race Conditions), когда один и тот же пользователь пытается одновременно отправить несколько спам-запросов через разные сетевые потоки шлюза.


Стратегия именования ключей и TTL (Key Space Design)

В рамках компонента используются два типа ключей:

  1. Словари стоп-слов (Read-Only Sets):
    • Шаблон ключа: dict:profanity:{app_lang}
    • Тип данных: Set (набор уникальных строк/корней слов).
    • Время жизни (TTL): Персистентно (без ограничения времени, обновляется при деплое словарей).
    • Пример: dict:profanity:ru-RU \(\rightarrow\) ["мат1", "мат2", "спам-фраза"]
  2. Счетчики нарушений пользователей (Dynamic Counters):
    • Шаблон ключа: abuse:counter:{user_id}
    • Тип данных: String (целочисленный счетчик).
    • Время жизни (TTL): 86400 секунд (24 часа с момента фиксации первого нарушения). Применяется скользящее окно блокировки.

Архитектурная схема выполнения скрипта

flowchart TD
    A([Вход: Массив слов + user_id]) --> B[Цикл по словам: Проверка через SISMEMBER]
    B -- Слово найдено в dict:profanity --> C[Флаг is_profane = true]
    B -- Текст чист --> D{is_profane == true?}
    C --> D
    D -- Да --> E[INCRBY abuse:counter:user_id 1]
    D -- Нет --> F([Выход: is_profane=false, counter=текущий])
    E --> G{Текущий count == 1?}
    G -- Да (Первый раз) --> H[EXPIRE abuse:counter 86400 сек]
    G -- Нет (Повторно) --> I[Оставить старый TTL]
    H --> J([Выход: is_profane=true, counter=обновленный])
    I --> J


Тело Lua-скрипта (Production Ready)

Этот скрипт передается в Redis через команду EVAL или вызывается по хэшу через EVALSHA.

-- KEYS[1]: Ключ словаря стоп-слов (например, 'dict:profanity:ru-RU')
-- KEYS[2]: Ключ счетчика пользователя (например, 'abuse:counter:usr_9876')
-- ARGV: Массив слов из очищенного текста пользователя (передается как аргументы)

local is_profane = false
local dict_key = KEYS[1]
local counter_key = KEYS[2]

-- 1. Лингвистический анализ по жесткому словарю
for i = 1, #ARGV do
    local word_exists = redis.call('SISMEMBER', dict_key, ARGV[i])
    if word_exists == 1 then
        is_profane = true
        break -- Завершаем цикл при первом же совпадении
    end
end

-- 2. Управление счетчиком нарушений (Abuse Control)
local current_violations = 0
if is_profane then
    -- Инкрементируем счетчик нарушений пользователя
    current_violations = redis.call('INCRBY', counter_key, 1)
    
    -- Если это самое первое нарушение, выставляем TTL на 24 часа
    if tonumber(current_violations) == 1 then
        redis.call('EXPIRE', counter_key, 86400)
    end
else
    -- Если текст чист, просто считываем текущее состояние (для логов)
    local val = redis.call('GET', counter_key)
    if val then
        current_violations = tonumber(val)
    end
end

-- 3. Возврат результата в воркер
-- Возвращает массив: [is_profane (0 или 1), текущее_кол_во_нарушений]
local status_flag = is_profane and 1 or 0
return {status_flag, current_violations}

Интеграция в общую таблицу шагов (Шаг 4.1)

При вызове из кода воркера censorship-control-worker (написанного на Python/Go/Node.js), логика обрабатывает исключения кэша по принципу Fail-Safe: если нода Redis падает, скрипт выплевывает инфраструктурную ошибку, но воркер не останавливает конвейер, а временно доверяет только оценке нейросети Ollama, чтобы не заблокировать пользовательский UI.