[DRAFT] Microservices: Сервис уведомлений

Спецификация асинхронного конвейера i18n-локализации и подсистемы доставки PUSH / Web-Socket

Published

June 11, 2026

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

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

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

1 Назначение и целевая бизнес-логика компонента

Создадим универсальный сервис отправки уведомлений по событию. Это будет асинхронный сервис читающий сообщения из топика и отправляющий пользователю уведомление по справочнику используя один из доступных каналов. Справочник будет содержать сообщения пользователю определяемые по айди и языку используемому пользователем.

Микросервис Notification Service является асинхронным конвейером, отвечающим за гарантированную доставку системных, транзакционных и маркетинговых уведомлений пользователям. Сервис считывает события из брокера очередей Kafka, определяет целевой язык интерфейса пользователя, извлекает шаблоны текстов и осуществляет отправку через два независимых канала: Web-Socket соединение (основной канал быстрого информирования внутри приложения) или Firebase Cloud Messaging (FCM) (резервный фоновый канал доставки пушей).

1.1 Архитектурные требования для миграции на Go:

  1. Слой хранения (База данных и Кэш): В целевой архитектуре локализованные шаблоны уведомлений выносятся из кода в СУБД (таблица notification_templates). Соответствие user_id -> fcm_token и признак его текущего активного Web-Socket соединения кэшируются в Redis для мгновенной проверки рантайм-статуса «клиент в сети».
  2. Многопоточность и Highload-пайплайн: На языке Go обработка сообщений из топика организуется через пул воркеров (Worker Pool). Каждый воркер обрабатывает событие в отдельной горутине (Goroutine), выполняя неблокирующие вызовы к сетевой инфраструктуре Google.
  3. Оптимизация массовых рассылок (Маркетинговые кампании): Для исключения перегрузки брокера и CPU при отправке уведомлений на миллион пользователей, сервис поддерживает пакетные сообщения. На вход подается вектор идентификаторов (user_ids) и мультиязычный словарь текстов. На уровне транспорта Go разбивает вектор на чанки по 500 адресатов и отправляет их в Google через метод messaging.SendMulticast(), перекладывая задачу массовой дистрибуции на инфраструктуру Firebase.

2 Структура и конфигурация слоя баз данных (Реестр СУБД)

Для перевода сервиса уведомлений на Go доменный слой изолирует хардкод мультиязычного словаря (I18N_ERROR_CATALOG) в реляционную структуру таблиц СУБД PostgreSQL, а оперативные девайс-токены и статусы сессий — в кэш-слой Redis.

2.1 Физическая спецификация таблиц (PostgreSQL)

Таблица notification_templates (Справочник локализованных шаблонов)• id (SERIAL, PRIMARY KEY) — Внутренний суррогатный ключ записи.

  • system_code (VARCHAR(100), NOT NULL) — Системный инвариантный код инцидента (Уникальный индекс: например, ‘SECURITY_BLOCKED’, ‘ERR_OCR_FAILED’).
  • lang (VARCHAR(2), NOT NULL) — Двухсимвольный маркер языка интерфейса (‘ru’, ‘en’).
  • title (VARCHAR(255), NOT NULL) — Локализованный заголовок пуш-уведомления.
  • body (TEXT, NOT NULL) — Сформированное тело текстового сообщения.
  • Составной индекс: UNIQUE INDEX (system_code, lang) для мгновенной выборки шаблона за один шаг.

2.2 Спецификация оперативного кэша (Redis / DB 1)

Для исключения паразитной нагрузки на основную БД при поштучной отправке пушей, воркеры вычитывают девайс-токены напрямую из оперативной памяти:

  • Ключ структуры: user:user_id:fcm_token (String) — Содержит актуальный регистрационный токен девайса, полученный от мобильного приложения.
  • Ключ структуры: user:user_id:ws_session (String/Enum) — Флаг или идентификатор активного Web-Socket соединения. Если ключ присутствует в кэше — пользователь в сети.

3 Спецификация gRPC-интерфейсов (Интеграционный контракт)

Хотя базовый конвейер доставки сообщений работает асинхронно через Kafka, микросервис уведомлений на Go предоставляет внутренний синхронный gRPC-интерфейс для административной панели, CRM-систем и ручных сервисных вызовов.

3.1 Метод SendDirectNotification (Протокол pb.NotificationService)

Используется для мгновенной точечной отправки сервисного пуша в обход общей очереди.

  • Входящий контракт (pb.DirectNotificationRequest):

  • user_id (int64, Required) — Идентификатор целевого пользователя.

  • type (string, Required) — Системный инвариантный код шаблона (например, ‘SECURITY_BLOCKED’).

  • x_request_id (string, Required) — Сквозной маркер трассировки логирования.

  • При передаче пустых обязательных полей сервер прерывает выполнение транзакции со статус-кодом codes.InvalidArgument.

• Исходящий контракт (pb.DirectNotificationResponse):

* success (bool) — Логический флаг факта доставки в целевой канал.
* channel (string) — Определенный рантаймом канал отправки (Строго фиксированные enum-строки: WEB_SOCKET или FIREBASE_CLOUD_MESSAGING).
* fcm_message_id (string, Optional) — Уникальный хэш транзакции, возвращенный облаком Google GCP (заполняется только при доставке через пуш).
* При отказе сетевой инфраструктуры Google или потере связи с кэшем, транспорт мапит ошибку в статус codes.Internal или codes.Unavailable.

4 Спецификация асинхронного интерфейса (Контракт Kafka)

Вместо синхронного gRPC-протокола, базовое межсервисное взаимодействие реализовано по событийной модели (Event-Driven). Сервис выступает исключительно в роли потребителя (Kafka Consumer) и не имеет встроенных продюсеров для публикации новых сообщений.

  • Целевой топик (Topic): user_notifications
  • Идентификатор группы (Group ID): bupar_notification_group

4.1 Структура входящего сообщения (Kafka Payload JSON)

Каждое событие в очереди представляет собой плоский JSON-объект со следующим набором обязательных полей:

  • user_id (string, Required) — Уникальный идентификатор получателя.
  • type (string, Required) — Инвариантный код инцидента для поиска шаблона (например, ‘SECURITY_BLOCKED’).
  • lang (string, Optional) — Двухсимвольный маркер локализации (‘ru’ / ‘en’). При отсутствии дублируется проверкой заголовков.
  • x_request_id (string, Required) — Сквозной UUID трассировки для распределенного логирования.

4.2 Чтение метаданных (Kafka Headers)

В целевой архитектуре на Go воркеры реализуют каскадный поиск токена языка. Если поле lang отсутствует внутри тела JSON, сервис в обязательном порядке выполняет чтение бинарных заголовков сообщения:

• Заголовок “lang” -> Декодируется в UTF-8 строку (Default fallback: “en”).

5 ER-диаграмма базы данных уведомлений (Связи сущностей)

Ниже представлена структура хранения локализованных шаблонов в реляционной СУБД и схема распределенного оперативного кэша FCM-токенов в оперативной памяти.

%%{init: {
  'theme': 'base', 
  'er': {
    'useMaxWidth': false
  },
  'themeVariables': {
    'mainBkg': '#FFF9C4',
    'lineColor': '#2C3E50',
    'borderClassName': '#FBC02D',
    'nodeBorder': '#FBC02D',
    'attributeBackend': '#FFF9C4'
  }
}}%%

erDiagram
    notification_templates {
        int id PK
        varchar system_code UK
        varchar lang UK
        varchar title
        text body
    }

    user_sessions_redis {
        string key_user_id PK
        string fcm_token
        boolean is_ws_active
    }

    notification_templates ||--o{ user_sessions_redis : "dispatched to"

6 Сценарий обработки события и логика ветвления каналов (Sequence Diagram)

Ниже представлена целевая схема обработки входящего события уведомления, включая валидацию локализации и каскадный выбор канала доставки (Web-Socket -> FCM Fallback).

%%{init: {
  'theme': 'base',
  'themeVariables': {
    'actorBkg': '#E3F2FD',
    'actorBorder': '#546E7A',
    'actorTextColor': '#0D47A1',
    'rectBkg': '#FFF9C4',
    'rectBorder': '#FBC02D',
    'noteBkgColor': '#F3E5F5',
    'noteBorderColor': '#7E57C2',
    'noteTextColor': '#311B92',
    'signalColor': '#2C3E50',
    'signalLineColor': '#2C3E50'
  }
}}%%
sequenceDiagram
    autonumber
    participant Kafka as Брокер Kafka
    participant Notif as Notification Service (Go)
    participant DB as База Данных и Кэш
    participant Google as Инфраструктура FCM Cloud

    Kafka->>Notif: Асинхронное сообщение из user-notifications (user_id, type, lang)
    activate Notif
    
    note over Notif: Шаг 2: Чтение параметров x-request-id и i18n-токена языка (из тела запроса или заголовков Kafka)
    
    Notif->>DB: Запрос текста шаблона по коду (type) и токена устройства по user_id
    activate DB
    DB-->>Notif: Локализованный текст (title, body) + fcm_token + статус ws_active
    deactivate DB
    
    alt Кейс А: Клиент удерживает активное Web-Socket соединение (ws_active == true)
        note over Notif: Шаг 5: Моментальная отправка сообщения внутрь открытого сокета клиента
        Notif-->>Kafka: Фиксация SUCCESS-статуса доставки через WS
    else Кейс Б: Клиент не в сети (ws_active == false) ИЛИ сбой Web-Socket доставки
        note over Notif: Шаг 7: Формирование структуры messaging.Message с payload данных трассировки
        Notif->>Google: Асинхронный HTTP-вызов пуш-нотификации messaging.send(fcm_message)
        activate Google
        Google-->>Notif: Ответ от Google Cloud Gateway (FCM-ID успеха или ошибка отклонения)
        deactivate Google
    end
    deactivate Notif

  • Примечание к схеме: Таблицы полностью изолированы на физическом уровне (разные хосты СУБД). Связь «один шаблон может быть отправлен на множество сессий пользователей» является сугубо логической и рассчитывается внутри горутин доменного слоя Go-сервиса на этапе склейки DTO.

6.1 Таблица расшифровки шагов сценария доставки уведомлений

Шаг Действие Параметры / Запросы / DTO Код ошибки (canonical_code) Статус
1 (Kafka -> Notif) Сервис вычитывает из топика user-notifications событие о необходимости отправки уведомления конкретному пользователю. Kafka Message Payload:
{ "user_id": "4512", "type": "ERR_OCR_FAILED", "lang": "ru", "x_request_id": "req-99" }
NOTIF_PARSE_FAILED WARN
2 (Notif -> Notif) Внутренняя логика: Разбор заголовков (Kafka Headers) для гарантированного извлечения целевого языка интерфейса и маркера трассировки. Внутренний парсер Go:
message.Headers.Get("lang")
Нет INFO
3 (Notif -> DB) Потоковый запрос во внутренний кэш и базу шаблонов для одновременного извлечения текста сообщения и токена адресата. SQL / Redis Query:
GET refresh:4512:fcm_token
SELECT title, body FROM templates WHERE code = $1 AND lang = $2;
DB_TEMPLATE_NOT_FOUND
DB_TOKEN_FETCH_ERROR
WARN
ERROR
4 (DB -> Notif) Система хранения отдает локализованную языковую пару (title, body), девайс-токен и флаг проверки активности сокета. Репозиторный DTO:
struct { Title string; Body string; Token string; IsWsActive bool }
Нет INFO
5 (Notif -> Notif) Условие (Кейс А): Пользователь находится в приложении. Сервис производит прямую отправку сообщения в дескриптор открытого веб-сокета. WebSocket JSON Payload:
{ "title": "Не удалось прочитать чек", "body": "Фотография размыта" }
WS_DELIVERY_FAILED (Каскадный переход к Шагу 7) WARN
6 (Notif -> Notif) Условие (Кейс Б): Пользователь вне сети. Сервис инициализирует сборку пуш-конверта и заполняет служебный блок data для Flutter-клиента. FCM Message DTO:
messaging.Message{ Notification: { Title: "..." }, Data: { "x_request_id": "req-99" }, Token: "..." }
Нет INFO
7 (Notif -> Google) Сервис выполняет асинхронный сетевой запрос к шлюзам Google Cloud Platform для доставки пуш-нотификации на устройство. Google API Call:
messaging.send(fcm_message)
FCM_API_CONNECTION_DOWN ERROR
8 (Google -> Notif) Облачная инфраструктура Firebase подтверждает успешный прием пакета и возвращает уникальный идентификатор транзакции. Google Cloud Response:
"projects/foodtracker/messages/109283012938"
FCM_PAYLOAD_REJECTED (Токен устройства устарел или отозван) WARN

7 Справочник Обсервабилити (Коды ошибок под выгрузку в ClickHouse)

Каждое отправленное уведомление или зафиксированный сбой маршрутизации регистрируется в виде плоской JSON-строки в stdout для последующей агрегации в аналитической базе ClickHouse:

1. Код (canonical_code) 2. Уровень (log_level) 3. Статусы (transport_statuses) 4. Получатель (error_target) 5. JSON для Фронтенда (ui_payload) 6. Метрики для ClickHouse (observability_json)
NOTIF_PARSE_FAILED WARN {“grpc”: null, “http”: 400} INTERNAL_SYSTEM null {“metric”: “kafka_consume_fail”, “labels”: {“reason”: “invalid_json_body”}}
DB_TEMPLATE_NOT_FOUND WARN {“grpc”: null, “http”: 404} INTERNAL_SYSTEM null {“metric”: “i18n_miss”, “labels”: {“service”: “notification_service”, “fallback”: “en_default”}}
FCM_PAYLOAD_REJECTED WARN {“grpc”: null, “http”: 422} GOOGLE_FCM_API null {“metric”: “push_rejected”, “labels”: {“reason”: “unregistered_device_token”}}
FCM_API_CONNECTION_DOWN ERROR {“grpc”: null, “http”: 503} GOOGLE_FCM_API null {“metric”: “push_network_lost”, “labels”: {“cloud_provider”: “google_gcp”}}