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{}, }) } }