feat(push): Web Push — VAPID, подписки устройств и отправка (AGENT.md 7.16)
Уведомления в браузере и на телефоне: сервер сам решает, кому их слать, и подписывает запрос VAPID-ключом, поэтому уведомление приходит, даже когда клиент закрыт. Сервер: - миграция 00021: `push_subscriptions` (эндпоинт уникален, ключи, счётчик неудач, время последней доставки); - `internal/push` — VAPID-ключи (приватный PKCS#8 из конфига, публичный выводится из него), правила уведомлений повторяют `web/src/lib/desktopNotifications.ts` (упоминания и личные беседы, без своих и системных сообщений), очередь доставки, TTL 12 часов, Topic по комнате; - 404/410 от push-сервиса удаляют подписку сразу, 5 неудач подряд — тоже, иначе копились бы мёртвые эндпоинты; retention чистит «молчащие» подписки; - `internal/httpx/safeurl.go` — общий запрет внутренних адресов с проверкой адреса в момент подключения (DNS rebinding): эндпоинт подписки приходит от клиента, и без проверки сервер сам себе организует SSRF; - `internal/gateway/presence.go` — `IsUserOnline`: если получатель в клиенте, уведомление покажет клиент, дублировать на телефон не нужно; - ручки `GET /push/config`, `GET|POST|DELETE /push/subscriptions`, лимиты 10 подписок и 20 уведомлений в минуту; ключ p256dh проверяется как настоящая точка P-256; - установщик генерирует `VAPID_PRIVATE_KEY` (openssl, PKCS#8 DER в base64) и `VAPID_SUBJECT`, ключ переиспользуется при переустановке; для ручной установки есть `glchat vapid-keys`; `features.web_push_enabled` виден в `/meta`. Тесты: правила и VAPID, доставка с поддельным push-эндпоинтом (шифрование, заголовки VAPID/TTL/Topic, очистка мёртвых подписок), ручки и лимиты, отправка при личном сообщении и упоминании, пропуск онлайн-получателей, миграция на чистой БД.
This commit is contained in:
@@ -0,0 +1,343 @@
|
||||
// Package push отправляет Web Push уведомления браузерам и телефонам
|
||||
// (AGENT.md 7.16, Фаза 7): VAPID-подпись, шифрование полезной нагрузки по
|
||||
// RFC 8291 и очередь доставки, чтобы HTTP-ручки не ждали push-сервисы.
|
||||
//
|
||||
// Решение «уведомлять или нет» принимает сервер (internal/push/rules.go) по
|
||||
// тем же правилам, что нативные уведомления desktop-обёртки
|
||||
// (web/src/lib/desktopNotifications.ts), а показывает уведомление service
|
||||
// worker клиента (web/public/sw.js).
|
||||
package push
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/ecdsa"
|
||||
"crypto/elliptic"
|
||||
"crypto/x509"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
webpush "github.com/SherClockHolmes/webpush-go"
|
||||
|
||||
"glchat/internal/config"
|
||||
"glchat/internal/httpx"
|
||||
"glchat/internal/store"
|
||||
)
|
||||
|
||||
const (
|
||||
// DefaultTTL — сколько push-сервис хранит недоставленное уведомление.
|
||||
DefaultTTL = 12 * time.Hour
|
||||
// MaxPayload — предел полезной нагрузки: запись RFC 8291 — 4096 байт,
|
||||
// оставляем запас на шифрование и служебные поля.
|
||||
MaxPayload = 3000
|
||||
// DeliverTimeout — сколько ждём push-сервис на одну доставку.
|
||||
DeliverTimeout = 10 * time.Second
|
||||
// MaxFailures — после стольких неудач подряд подписка удаляется.
|
||||
MaxFailures = 5
|
||||
// Воркеры и очередь: push не должен ни блокировать запросы, ни плодить
|
||||
// горутины на каждого получателя.
|
||||
workerCount = 4
|
||||
queueSize = 256
|
||||
)
|
||||
|
||||
// VAPIDSubject по умолчанию: push-сервисы требуют контакт администратора.
|
||||
const defaultSubjectPrefix = "mailto:admin@"
|
||||
|
||||
// Payload — то, что получает service worker (web/public/sw.js).
|
||||
type Payload struct {
|
||||
Title string `json:"title"`
|
||||
Body string `json:"body"`
|
||||
// Kind — "mention" или "direct": клиент по нему выбирает иконку и текст.
|
||||
Kind string `json:"kind"`
|
||||
GuildID string `json:"guild_id,omitempty"`
|
||||
GuildName string `json:"guild_name,omitempty"`
|
||||
ChannelID string `json:"channel_id"`
|
||||
ChannelName string `json:"channel_name,omitempty"`
|
||||
AuthorID string `json:"author_id,omitempty"`
|
||||
AuthorName string `json:"author_name,omitempty"`
|
||||
MessageID string `json:"message_id,omitempty"`
|
||||
// URL — путь, который открывает клик по уведомлению.
|
||||
URL string `json:"url"`
|
||||
// Tag — ключ схлопывания: уведомления одной комнаты заменяют друг друга.
|
||||
Tag string `json:"tag"`
|
||||
}
|
||||
|
||||
// Service — отправитель Web Push. Нулевое значение безопасно: методы
|
||||
// проверяют nil, поэтому «push выключен» не требует ветвлений у вызывающего.
|
||||
type Service struct {
|
||||
logger *slog.Logger
|
||||
store *store.Store
|
||||
publicKey string
|
||||
privateKey string
|
||||
subject string
|
||||
ttl int
|
||||
client webpush.HTTPClient
|
||||
|
||||
startOnce sync.Once
|
||||
queue chan delivery
|
||||
closed chan struct{}
|
||||
closeOnce sync.Once
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
type delivery struct {
|
||||
subscription store.PushSubscription
|
||||
payload Payload
|
||||
}
|
||||
|
||||
// New собирает сервис из конфига. Если VAPID-ключ не задан — возвращает nil
|
||||
// без ошибки: инстанс без push работает как обычно. Некорректный ключ —
|
||||
// ошибка: молча выключенный push выглядел бы как «уведомления не приходят».
|
||||
func New(cfg config.Config, st *store.Store, logger *slog.Logger) (*Service, error) {
|
||||
if logger == nil {
|
||||
logger = slog.New(slog.DiscardHandler)
|
||||
}
|
||||
if strings.TrimSpace(cfg.VAPIDPrivateKey) == "" {
|
||||
return nil, nil
|
||||
}
|
||||
privateKey, publicKey, err := ParseVAPIDKeys(cfg.VAPIDPrivateKey)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("parse VAPID key: %w", err)
|
||||
}
|
||||
subject := strings.TrimSpace(cfg.VAPIDSubject)
|
||||
if subject == "" {
|
||||
subject = defaultSubjectPrefix + cfg.Domain
|
||||
}
|
||||
return &Service{
|
||||
logger: logger,
|
||||
store: st,
|
||||
publicKey: publicKey,
|
||||
privateKey: privateKey,
|
||||
subject: subject,
|
||||
ttl: int(DefaultTTL.Seconds()),
|
||||
client: httpx.SafeClient(DeliverTimeout),
|
||||
queue: make(chan delivery, queueSize),
|
||||
closed: make(chan struct{}),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Enabled сообщает, настроен ли Web Push на инстансе.
|
||||
func (s *Service) Enabled() bool { return s != nil && s.privateKey != "" }
|
||||
|
||||
// PublicKey — публичный VAPID-ключ (base64url) для `applicationServerKey`.
|
||||
func (s *Service) PublicKey() string {
|
||||
if s == nil {
|
||||
return ""
|
||||
}
|
||||
return s.publicKey
|
||||
}
|
||||
|
||||
// Enqueue ставит доставку в очередь. Возвращает false, если push выключен или
|
||||
// очередь переполнена: уведомление теряется осознанно, но запрос не ждёт.
|
||||
func (s *Service) Enqueue(subscription store.PushSubscription, payload Payload) bool {
|
||||
if !s.Enabled() {
|
||||
return false
|
||||
}
|
||||
select {
|
||||
case <-s.closed:
|
||||
return false
|
||||
default:
|
||||
}
|
||||
s.start()
|
||||
select {
|
||||
case s.queue <- delivery{subscription: subscription, payload: payload}:
|
||||
return true
|
||||
default:
|
||||
s.logger.Warn("push queue is full, notification dropped",
|
||||
slog.String("user_id", formatID(subscription.UserID)))
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// Close останавливает воркеры: вызывается при завершении процесса и в тестах.
|
||||
func (s *Service) Close() {
|
||||
if s == nil {
|
||||
return
|
||||
}
|
||||
s.closeOnce.Do(func() { close(s.closed) })
|
||||
s.wg.Wait()
|
||||
}
|
||||
|
||||
func (s *Service) start() {
|
||||
s.startOnce.Do(func() {
|
||||
for i := 0; i < workerCount; i++ {
|
||||
s.wg.Add(1)
|
||||
go s.work()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Service) work() {
|
||||
defer s.wg.Done()
|
||||
for {
|
||||
select {
|
||||
case <-s.closed:
|
||||
return
|
||||
case job := <-s.queue:
|
||||
s.deliver(context.Background(), job)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// deliver отправляет одно уведомление и обслуживает ответ push-сервиса:
|
||||
// успех отмечается, 404/410 удаляют подписку, остальные ошибки копят счётчик.
|
||||
func (s *Service) deliver(ctx context.Context, job delivery) {
|
||||
payload, err := json.Marshal(job.payload)
|
||||
if err != nil {
|
||||
s.logger.WarnContext(ctx, "push payload marshal failed", slog.Any("error", err))
|
||||
return
|
||||
}
|
||||
if len(payload) > MaxPayload {
|
||||
s.logger.WarnContext(ctx, "push payload is too large",
|
||||
slog.Int("bytes", len(payload)), slog.String("user_id", formatID(job.subscription.UserID)))
|
||||
return
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(ctx, DeliverTimeout)
|
||||
defer cancel()
|
||||
|
||||
response, err := webpush.SendNotificationWithContext(ctx, payload, &webpush.Subscription{
|
||||
Endpoint: job.subscription.Endpoint,
|
||||
Keys: webpush.Keys{
|
||||
P256dh: job.subscription.P256dh,
|
||||
Auth: job.subscription.Auth,
|
||||
},
|
||||
}, &webpush.Options{
|
||||
HTTPClient: s.client,
|
||||
Subscriber: s.subject,
|
||||
TTL: s.ttl,
|
||||
Urgency: webpush.UrgencyNormal,
|
||||
Topic: topicFor(job.payload),
|
||||
VAPIDPublicKey: s.publicKey,
|
||||
VAPIDPrivateKey: s.privateKey,
|
||||
})
|
||||
if err != nil {
|
||||
s.noteFailure(ctx, job.subscription, 0, err)
|
||||
return
|
||||
}
|
||||
defer func() { _ = response.Body.Close() }()
|
||||
|
||||
switch {
|
||||
case response.StatusCode >= 200 && response.StatusCode < 300:
|
||||
if err := s.store.TouchPushSubscription(ctx, job.subscription.ID); err != nil {
|
||||
s.logger.WarnContext(ctx, "push subscription touch failed", slog.Any("error", err))
|
||||
}
|
||||
case response.StatusCode == http.StatusNotFound || response.StatusCode == http.StatusGone:
|
||||
// Подписка отозвана браузером или push-сервисом — хранить нечего.
|
||||
if err := s.store.DeletePushSubscriptionByID(ctx, job.subscription.ID); err != nil {
|
||||
s.logger.WarnContext(ctx, "push subscription delete failed", slog.Any("error", err))
|
||||
}
|
||||
default:
|
||||
s.noteFailure(ctx, job.subscription, response.StatusCode, nil)
|
||||
}
|
||||
}
|
||||
|
||||
// noteFailure копит неудачи и убирает подписку, которая не работает подряд.
|
||||
func (s *Service) noteFailure(ctx context.Context, subscription store.PushSubscription, status int, cause error) {
|
||||
count, err := s.store.FailPushSubscription(ctx, subscription.ID)
|
||||
if err != nil {
|
||||
s.logger.WarnContext(ctx, "push subscription failure note failed", slog.Any("error", err))
|
||||
return
|
||||
}
|
||||
attributes := []any{
|
||||
slog.String("user_id", formatID(subscription.UserID)),
|
||||
slog.Int("status", status),
|
||||
slog.Int("failures", count),
|
||||
}
|
||||
if cause != nil {
|
||||
attributes = append(attributes, slog.Any("error", cause))
|
||||
}
|
||||
s.logger.WarnContext(ctx, "web push delivery failed", attributes...)
|
||||
if count >= MaxFailures {
|
||||
if err := s.store.DeletePushSubscriptionByID(ctx, subscription.ID); err != nil {
|
||||
s.logger.WarnContext(ctx, "push subscription delete failed", slog.Any("error", err))
|
||||
return
|
||||
}
|
||||
s.logger.InfoContext(ctx, "stale push subscription removed",
|
||||
slog.String("user_id", formatID(subscription.UserID)))
|
||||
}
|
||||
}
|
||||
|
||||
// topicFor схлопывает уведомления одной комнаты: свежее заменяет предыдущее,
|
||||
// иначе упоминания в активной переписке заваливают телефон.
|
||||
func topicFor(payload Payload) string {
|
||||
if payload.ChannelID == "" {
|
||||
return ""
|
||||
}
|
||||
return "channel-" + payload.ChannelID
|
||||
}
|
||||
|
||||
// ParseVAPIDKeys принимает приватный ключ в двух видах и возвращает пару в
|
||||
// формате, которого ждёт web-push (base64url без выравнивания):
|
||||
// - PKCS#8 DER в base64 — так ключ генерирует установщик (`openssl pkcs8`);
|
||||
// - «сырой» 32-байтовый скаляр в base64url — вывод GenerateVAPIDKeys().
|
||||
//
|
||||
// Публичный ключ всегда выводится из приватного: хранить его отдельно
|
||||
// незачем, а рассинхронизация пары ломала бы подписку молча.
|
||||
func ParseVAPIDKeys(encoded string) (privateKey, publicKey string, err error) {
|
||||
trimmed := strings.TrimSpace(encoded)
|
||||
raw, decodeErr := decodeBase64(trimmed)
|
||||
if decodeErr != nil {
|
||||
return "", "", fmt.Errorf("decode: %w", decodeErr)
|
||||
}
|
||||
var key *ecdsa.PrivateKey
|
||||
switch {
|
||||
case len(raw) == 32:
|
||||
// «Сырой» скаляр: так ключ печатает GenerateVAPIDKeys().
|
||||
parsed, parseErr := ecdsa.ParseRawPrivateKey(elliptic.P256(), raw)
|
||||
if parseErr != nil {
|
||||
return "", "", fmt.Errorf("parse raw private key: %w", parseErr)
|
||||
}
|
||||
key = parsed
|
||||
default:
|
||||
parsed, parseErr := x509.ParsePKCS8PrivateKey(raw)
|
||||
if parseErr != nil {
|
||||
return "", "", fmt.Errorf("parse PKCS#8: %w", parseErr)
|
||||
}
|
||||
ecdsaKey, ok := parsed.(*ecdsa.PrivateKey)
|
||||
if !ok {
|
||||
return "", "", errors.New("key is not an ECDSA private key")
|
||||
}
|
||||
if ecdsaKey.Curve != elliptic.P256() {
|
||||
return "", "", fmt.Errorf("curve %s is not P-256", ecdsaKey.Curve.Params().Name)
|
||||
}
|
||||
key = ecdsaKey
|
||||
}
|
||||
// Bytes() у ecdsa отдаёт фиксированную длину и несжатую точку (SEC 1):
|
||||
// именно в таком виде ключи ждёт web-push, и никаких big.Int снаружи.
|
||||
scalar, err := key.Bytes()
|
||||
if err != nil {
|
||||
return "", "", fmt.Errorf("encode private key: %w", err)
|
||||
}
|
||||
point, err := key.PublicKey.Bytes()
|
||||
if err != nil {
|
||||
return "", "", fmt.Errorf("encode public key: %w", err)
|
||||
}
|
||||
return base64.RawURLEncoding.EncodeToString(scalar),
|
||||
base64.RawURLEncoding.EncodeToString(point), nil
|
||||
}
|
||||
|
||||
// decodeBase64 принимает base64 (standard/raw/url) — установщик пишет
|
||||
// стандартный base64, а ключи из библиотек приходят в base64url.
|
||||
func decodeBase64(value string) ([]byte, error) {
|
||||
encodings := []*base64.Encoding{
|
||||
base64.StdEncoding, base64.RawStdEncoding,
|
||||
base64.URLEncoding, base64.RawURLEncoding,
|
||||
}
|
||||
var lastErr error
|
||||
for _, encoding := range encodings {
|
||||
decoded, err := encoding.DecodeString(value)
|
||||
if err == nil {
|
||||
return decoded, nil
|
||||
}
|
||||
lastErr = err
|
||||
}
|
||||
return nil, lastErr
|
||||
}
|
||||
|
||||
func formatID(id uint64) string { return fmt.Sprintf("%d", id) }
|
||||
Reference in New Issue
Block a user