%%{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"
[DRAFT] Microservices: Сервис уведомлений
Спецификация асинхронного конвейера i18n-локализации и подсистемы доставки PUSH / Web-Socket
В открытом доступе представлена демонстрационная версия метода. В настоящей публичной документации отображены не все шаги, технические сценарии и приватные эндпоинты для системы цифровых симуляторов бизнес-процессов.
- Полная спецификация метода: Будет доступна только во внутреннем контуре разработки (Confluence / Swagger Enterprise).
1 Назначение и целевая бизнес-логика компонента
Микросервис Notification Service является асинхронным конвейером, отвечающим за гарантированную доставку системных, транзакционных и маркетинговых уведомлений пользователям. Сервис считывает события из брокера очередей Kafka, определяет целевой язык интерфейса пользователя, извлекает шаблоны текстов и осуществляет отправку через два независимых канала: Web-Socket соединение (основной канал быстрого информирования внутри приложения) или Firebase Cloud Messaging (FCM) (резервный фоновый канал доставки пушей).
1.1 Архитектурные требования для миграции на Go:
- Слой хранения (База данных и Кэш): В целевой архитектуре локализованные шаблоны уведомлений выносятся из кода в СУБД (таблица notification_templates). Соответствие user_id -> fcm_token и признак его текущего активного Web-Socket соединения кэшируются в Redis для мгновенной проверки рантайм-статуса «клиент в сети».
- Многопоточность и Highload-пайплайн: На языке Go обработка сообщений из топика организуется через пул воркеров (Worker Pool). Каждый воркер обрабатывает событие в отдельной горутине (Goroutine), выполняя неблокирующие вызовы к сетевой инфраструктуре Google.
- Оптимизация массовых рассылок (Маркетинговые кампании): Для исключения перегрузки брокера и 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-токенов в оперативной памяти.
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_tokenSELECT 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/bupar/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”}} |