Files
minstrel/internal/notifications/notifier.go
T
bvandeusenandClaude Opus 5.5 63709a433d
release / govulncheck (push) Successful in 21s
release / web (push) Successful in 1m19s
release / go (push) Successful in 1m39s
release / integration (push) Successful in 5m27s
release / android (push) Successful in 5m47s
release / Build signed APK (releases and dev) (push) Successful in 5m34s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 26s
release / Verify release artifacts (tag releases only) (push) Skipped
feat(notifications): grouped email digest, new music as a daily summary (#5346)
Nothing is emailed per event. New music (request_completed) goes out at
most once a day, at the summary hour in each user's own timezone, grouped
by artist. Everything else is batched: one email a window after the first
un-emailed item, holding whatever accumulated.

- Migration 0074: notification_email_settings (summary hour, batch window,
  admin-configurable) and user_notification_email_state (batch start, last
  sent, failures and retry_after per user and group). Existing rows are
  stamped emailed so the upgrade sends no backlog.
- The Notifier stamps emailed_at at write time when the recipient's email
  channel is off, so turning email on later doesn't send old items.
- Read rows are never selected. A row is stamped only after the mailer
  accepts, in one transaction with the state, against the read's clock, so
  a coalesced row updated mid-send stays pending.
- A failed send backs off 5m doubling to 6h; SMTP not configured just waits.
- Links come from the public address; without one the email has none.
- The mailer now RFC 2047-encodes subjects and strips line breaks from them.
- Admin → Integrations gains a Notification emails card.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-08 07:48:22 -04:00

203 lines
6.3 KiB
Go

package notifications
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
"git.fabledsword.com/bvandeusen/minstrel/internal/eventbus"
)
// EventCreated is the bus event a recipient's open clients receive when a
// row lands for them. It is a nudge with no content: the client fetches the
// inbox. A frame that carried the notification could not be told apart from
// a repeat after a reconnect, and a frame nobody received is then one nobody
// needed, because the client pulls on reconnect anyway (Roundtable #2535).
const EventCreated = "notification.created"
// Recipients names who an event is for.
type Recipients struct {
Users []pgtype.UUID
// Admins adds every admin.
Admins bool
// Except is never notified, even when it is in Users or is an admin. The
// admin who filed a request needs no notice that it is pending.
Except pgtype.UUID
}
// ToUser addresses one user.
func ToUser(id pgtype.UUID) Recipients { return Recipients{Users: []pgtype.UUID{id}} }
// ToAdmins addresses every admin except one (pass the zero UUID for none).
func ToAdmins(except pgtype.UUID) Recipients { return Recipients{Admins: true, Except: except} }
// Notifier writes notifications and nudges recipients' clients. A nil
// *Notifier is valid and does nothing, so producers built without one (most
// tests) need no special case.
type Notifier struct {
db dbq.DBTX
bus *eventbus.Bus
logger *slog.Logger
}
// New returns a Notifier writing through db. bus may be nil.
func New(db dbq.DBTX, bus *eventbus.Bus, logger *slog.Logger) *Notifier {
if logger == nil {
logger = slog.Default()
}
return &Notifier{db: db, bus: bus, logger: logger}
}
// Notify records kind for each recipient whose inbox preference allows it,
// then nudges their open clients.
//
// For a coalescing kind, an unread row of that kind is updated in place
// rather than a new one added. payload["count"] then adds to the unread
// row's count for kinds whose events each mean "N more".
//
// The error is for logging. A notification must never fail the action that
// caused it, so producers log and carry on; NotifyLogged does exactly that.
func (n *Notifier) Notify(ctx context.Context, kind Kind, to Recipients, payload map[string]any) error {
if n == nil {
return nil
}
if !kind.Valid() {
return fmt.Errorf("notifications: unknown kind %q", kind)
}
q := dbq.New(n.db)
ids, err := n.resolve(ctx, q, kind, to)
if err != nil || len(ids) == 0 {
return err
}
channels, err := channelsFor(ctx, q, kind, ids)
if err != nil {
return err
}
if payload == nil {
payload = map[string]any{}
}
body, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("notifications: encode %s payload: %w", kind, err)
}
var delivered []pgtype.UUID
for _, id := range ids {
if !channels[id].Inbox {
continue
}
if err := write(ctx, q, kind, id, body, channels[id].Email); err != nil {
return fmt.Errorf("notifications: write %s: %w", kind, err)
}
delivered = append(delivered, id)
}
n.nudge(delivered)
return nil
}
// NotifyLogged is Notify for producers: a failure is logged at WARN and
// otherwise ignored.
func (n *Notifier) NotifyLogged(ctx context.Context, kind Kind, to Recipients, payload map[string]any) {
if err := n.Notify(ctx, kind, to, payload); err != nil {
n.logger.Warn("notifications: notify failed", "kind", string(kind), "err", err)
}
}
// resolve turns Recipients into a deduplicated list of user ids. Admin kinds
// are delivered only to admins, whatever Users says, so a producer that
// addresses a non-admin by mistake stores nothing for them.
func (n *Notifier) resolve(ctx context.Context, q *dbq.Queries, kind Kind, to Recipients) ([]pgtype.UUID, error) {
var admins []pgtype.UUID
if to.Admins || kind.AdminOnly() {
var err error
if admins, err = q.ListAdminUserIDs(ctx); err != nil {
return nil, fmt.Errorf("notifications: list admins: %w", err)
}
}
isAdmin := make(map[pgtype.UUID]bool, len(admins))
for _, a := range admins {
isAdmin[a] = true
}
var candidates []pgtype.UUID
candidates = append(candidates, to.Users...)
if to.Admins {
candidates = append(candidates, admins...)
}
seen := make(map[pgtype.UUID]bool, len(candidates))
out := make([]pgtype.UUID, 0, len(candidates))
for _, id := range candidates {
if !id.Valid || seen[id] || (to.Except.Valid && id == to.Except) {
continue
}
if kind.AdminOnly() && !isAdmin[id] {
continue
}
seen[id] = true
out = append(out, id)
}
return out, nil
}
// channelsFor returns each recipient's effective channels for kind: their
// stored preference, or the kind's defaults when they have none.
func channelsFor(ctx context.Context, q *dbq.Queries, kind Kind, ids []pgtype.UUID) (map[pgtype.UUID]Channels, error) {
rows, err := q.ListNotificationPrefsForKind(ctx, dbq.ListNotificationPrefsForKindParams{
Kind: string(kind),
UserIds: ids,
})
if err != nil {
return nil, fmt.Errorf("notifications: read prefs: %w", err)
}
out := make(map[pgtype.UUID]Channels, len(ids))
for _, id := range ids {
out[id] = kind.Defaults().Effective()
}
for _, r := range rows {
out[r.UserID] = Channels{Inbox: r.Inbox, Phone: r.Phone, Email: r.Email}.Effective()
}
return out, nil
}
// write stores one row. emailWanted is the recipient's email channel for the
// kind: a row nobody wants emailed is stamped emailed_at at once, so the
// digest never sends it, even if email is switched on later.
func write(ctx context.Context, q *dbq.Queries, kind Kind, userID pgtype.UUID, body []byte, emailWanted bool) error {
key := kind.coalesceKey()
if key == "" {
_, err := q.InsertNotification(ctx, dbq.InsertNotificationParams{
UserID: userID, Kind: string(kind), Payload: body, EmailWanted: emailWanted,
})
return err
}
_, err := q.UpsertCoalescedNotification(ctx, dbq.UpsertCoalescedNotificationParams{
UserID: userID,
Kind: string(kind),
Payload: body,
CoalesceKey: &key,
EmailWanted: emailWanted,
SumCount: specs[kind].sumCount,
})
return err
}
func (n *Notifier) nudge(ids []pgtype.UUID) {
if n.bus == nil {
return
}
for _, id := range ids {
n.bus.Publish(eventbus.Event{
Kind: EventCreated,
UserID: uuid.UUID(id.Bytes).String(),
Data: map[string]any{},
})
}
}