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
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>
203 lines
6.3 KiB
Go
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{},
|
|
})
|
|
}
|
|
}
|