feat(notifications): the inbox store, one writer, coalescing and retention (#5338)
release / govulncheck (push) Successful in 42s
release / web (push) Successful in 1m34s
release / go (push) Successful in 1m51s
release / integration (push) Successful in 4m59s
release / android (push) Successful in 5m24s
release / Build signed APK (releases and dev) (push) Successful in 5m32s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 1m20s
release / Verify release artifacts (tag releases only) (push) Skipped

M489 step 1. The event bus is fire-and-forget, so a client that isn't
connected never hears that a request completed or that tracks went missing.
user_notifications is the durable record; the bus only nudges.

- Migration 0073: user_notifications (kind CHECK-gated, payload jsonb,
  read_at, coalesce_key, emailed_at) and user_notification_prefs (per user,
  per kind: inbox, phone, email). A missing pref row means the kind's
  defaults, so nothing is seeded.
- internal/notifications.Notifier is the only writer. It resolves recipients
  (admin kinds reach admins only, and never the excepted user), honours the
  inbox pref (phone and email ride on it), writes, and publishes a
  contentless notification.created nudge per recipient.
- Burst-prone admin kinds coalesce into one unread row: tracks_missing and
  scan_failed add up their counts, duplicates_found and playback_errors take
  the latest total. Once read, the next event is a new row.
- Retention: read rows go after 90 days, anything after a year, on the
  library_changes compactor's daily shape.

Tests: unit (channel rules, kind table) and integration (recipients, nudge,
coalescing both ways, prefs, owner-scoped idempotent mark-read, every kind
against both schema CHECKs, retention cut-offs).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-10-08 07:06:39 -04:00
co-authored by Claude Opus 5.5
parent 5dffe51b95
commit 8bf333e748
12 changed files with 1266 additions and 0 deletions
+20
View File
@@ -831,6 +831,26 @@ type UserNormalizationPref struct {
UpdatedAt pgtype.Timestamptz
}
type UserNotification struct {
ID pgtype.UUID
UserID pgtype.UUID
Kind string
Payload []byte
CreatedAt pgtype.Timestamptz
ReadAt pgtype.Timestamptz
CoalesceKey *string
EmailedAt pgtype.Timestamptz
}
type UserNotificationPref struct {
UserID pgtype.UUID
Kind string
Inbox bool
Phone bool
Email bool
UpdatedAt pgtype.Timestamptz
}
type YouMightLikeAlbum struct {
UserID pgtype.UUID
AlbumID pgtype.UUID
+346
View File
@@ -0,0 +1,346 @@
// Code generated by sqlc. DO NOT EDIT.
// versions:
// sqlc v1.31.1
// source: notifications.sql
package dbq
import (
"context"
"github.com/jackc/pgx/v5/pgtype"
)
const countUnreadNotifications = `-- name: CountUnreadNotifications :one
SELECT count(*) FROM user_notifications
WHERE user_id = $1 AND read_at IS NULL
`
func (q *Queries) CountUnreadNotifications(ctx context.Context, userID pgtype.UUID) (int64, error) {
row := q.db.QueryRow(ctx, countUnreadNotifications, userID)
var count int64
err := row.Scan(&count)
return count, err
}
const insertNotification = `-- name: InsertNotification :one
INSERT INTO user_notifications (user_id, kind, payload)
VALUES ($1, $2, $3)
RETURNING id
`
type InsertNotificationParams struct {
UserID pgtype.UUID
Kind string
Payload []byte
}
// M489 notifications inbox (#726). Rows are written only through
// internal/notifications.Notifier, which decides recipients and honours each
// recipient's inbox preference before it reaches these.
func (q *Queries) InsertNotification(ctx context.Context, arg InsertNotificationParams) (pgtype.UUID, error) {
row := q.db.QueryRow(ctx, insertNotification, arg.UserID, arg.Kind, arg.Payload)
var id pgtype.UUID
err := row.Scan(&id)
return id, err
}
const listAdminUserIDs = `-- name: ListAdminUserIDs :many
SELECT id FROM users WHERE is_admin = true ORDER BY created_at, id
`
func (q *Queries) ListAdminUserIDs(ctx context.Context) ([]pgtype.UUID, error) {
rows, err := q.db.Query(ctx, listAdminUserIDs)
if err != nil {
return nil, err
}
defer rows.Close()
var items []pgtype.UUID
for rows.Next() {
var id pgtype.UUID
if err := rows.Scan(&id); err != nil {
return nil, err
}
items = append(items, id)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listNotificationPrefsForKind = `-- name: ListNotificationPrefsForKind :many
SELECT user_id, inbox, phone, email
FROM user_notification_prefs
WHERE kind = $1 AND user_id = ANY($2::uuid[])
`
type ListNotificationPrefsForKindParams struct {
Kind string
UserIds []pgtype.UUID
}
type ListNotificationPrefsForKindRow struct {
UserID pgtype.UUID
Inbox bool
Phone bool
Email bool
}
// The stored prefs of one kind for a set of recipients. A recipient with no
// row has the kind's defaults.
func (q *Queries) ListNotificationPrefsForKind(ctx context.Context, arg ListNotificationPrefsForKindParams) ([]ListNotificationPrefsForKindRow, error) {
rows, err := q.db.Query(ctx, listNotificationPrefsForKind, arg.Kind, arg.UserIds)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListNotificationPrefsForKindRow
for rows.Next() {
var i ListNotificationPrefsForKindRow
if err := rows.Scan(
&i.UserID,
&i.Inbox,
&i.Phone,
&i.Email,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listNotificationPrefsForUser = `-- name: ListNotificationPrefsForUser :many
SELECT kind, inbox, phone, email
FROM user_notification_prefs
WHERE user_id = $1
`
type ListNotificationPrefsForUserRow struct {
Kind string
Inbox bool
Phone bool
Email bool
}
func (q *Queries) ListNotificationPrefsForUser(ctx context.Context, userID pgtype.UUID) ([]ListNotificationPrefsForUserRow, error) {
rows, err := q.db.Query(ctx, listNotificationPrefsForUser, userID)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListNotificationPrefsForUserRow
for rows.Next() {
var i ListNotificationPrefsForUserRow
if err := rows.Scan(
&i.Kind,
&i.Inbox,
&i.Phone,
&i.Email,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listNotifications = `-- name: ListNotifications :many
SELECT id, kind, payload, created_at, read_at
FROM user_notifications
WHERE user_id = $1
AND ($2::timestamptz IS NULL
OR (created_at, id) < ($2::timestamptz, $3::uuid))
ORDER BY created_at DESC, id DESC
LIMIT $4
`
type ListNotificationsParams struct {
UserID pgtype.UUID
BeforeCreatedAt pgtype.Timestamptz
BeforeID pgtype.UUID
PageLimit int32
}
type ListNotificationsRow struct {
ID pgtype.UUID
Kind string
Payload []byte
CreatedAt pgtype.Timestamptz
ReadAt pgtype.Timestamptz
}
// Newest first, keyset-paged on (created_at, id). Pass both cursor halves
// from the last row of the previous page, or neither for the first page.
func (q *Queries) ListNotifications(ctx context.Context, arg ListNotificationsParams) ([]ListNotificationsRow, error) {
rows, err := q.db.Query(ctx, listNotifications,
arg.UserID,
arg.BeforeCreatedAt,
arg.BeforeID,
arg.PageLimit,
)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListNotificationsRow
for rows.Next() {
var i ListNotificationsRow
if err := rows.Scan(
&i.ID,
&i.Kind,
&i.Payload,
&i.CreatedAt,
&i.ReadAt,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const markAllNotificationsRead = `-- name: MarkAllNotificationsRead :execrows
UPDATE user_notifications
SET read_at = now()
WHERE user_id = $1 AND read_at IS NULL
`
func (q *Queries) MarkAllNotificationsRead(ctx context.Context, userID pgtype.UUID) (int64, error) {
result, err := q.db.Exec(ctx, markAllNotificationsRead, userID)
if err != nil {
return 0, err
}
return result.RowsAffected(), nil
}
const markNotificationRead = `-- name: MarkNotificationRead :execrows
UPDATE user_notifications
SET read_at = COALESCE(read_at, now())
WHERE id = $1 AND user_id = $2
`
type MarkNotificationReadParams struct {
ID pgtype.UUID
UserID pgtype.UUID
}
// Scoped to the owner: another user's id matches no row. Marking an already
// read row keeps its original read_at and still matches, so a repeat is not
// mistaken for "not yours".
func (q *Queries) MarkNotificationRead(ctx context.Context, arg MarkNotificationReadParams) (int64, error) {
result, err := q.db.Exec(ctx, markNotificationRead, arg.ID, arg.UserID)
if err != nil {
return 0, err
}
return result.RowsAffected(), nil
}
const trimNotifications = `-- name: TrimNotifications :execrows
DELETE FROM user_notifications
WHERE (read_at IS NOT NULL AND read_at < $1)
OR created_at < $2
`
type TrimNotificationsParams struct {
ReadCutoff pgtype.Timestamptz
AnyCutoff pgtype.Timestamptz
}
// Retention: read rows go after read_cutoff, and anything at all after
// any_cutoff, so an inbox nobody opens doesn't grow without bound either.
func (q *Queries) TrimNotifications(ctx context.Context, arg TrimNotificationsParams) (int64, error) {
result, err := q.db.Exec(ctx, trimNotifications, arg.ReadCutoff, arg.AnyCutoff)
if err != nil {
return 0, err
}
return result.RowsAffected(), nil
}
const upsertCoalescedNotification = `-- name: UpsertCoalescedNotification :one
INSERT INTO user_notifications (user_id, kind, payload, coalesce_key)
VALUES ($1, $2, $3, $4)
ON CONFLICT (user_id, coalesce_key) WHERE read_at IS NULL AND coalesce_key IS NOT NULL
DO UPDATE SET
payload = CASE
WHEN $5::boolean THEN
EXCLUDED.payload || jsonb_build_object('count',
COALESCE((user_notifications.payload->>'count')::bigint, 0)
+ COALESCE((EXCLUDED.payload->>'count')::bigint, 0))
ELSE EXCLUDED.payload
END,
created_at = now(),
emailed_at = NULL
RETURNING id
`
type UpsertCoalescedNotificationParams struct {
UserID pgtype.UUID
Kind string
Payload []byte
CoalesceKey *string
SumCount bool
}
// One unread row per (user, coalesce_key). A new event while that row is
// unread updates it in place and moves it back to the top of the inbox.
//
// sum_count: the payload's `count` adds to the unread row's count rather than
// replacing it. For kinds whose event is "N more happened" (tracks marked
// missing). Kinds whose event states the whole current total (pending
// duplicate groups) pass false and the payload simply replaces.
//
// emailed_at is cleared so the newer state goes out in the next batch; an
// emailed-but-unread row that keeps growing is news the user hasn't seen.
func (q *Queries) UpsertCoalescedNotification(ctx context.Context, arg UpsertCoalescedNotificationParams) (pgtype.UUID, error) {
row := q.db.QueryRow(ctx, upsertCoalescedNotification,
arg.UserID,
arg.Kind,
arg.Payload,
arg.CoalesceKey,
arg.SumCount,
)
var id pgtype.UUID
err := row.Scan(&id)
return id, err
}
const upsertNotificationPref = `-- name: UpsertNotificationPref :exec
INSERT INTO user_notification_prefs (user_id, kind, inbox, phone, email)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (user_id, kind) DO UPDATE SET
inbox = EXCLUDED.inbox,
phone = EXCLUDED.phone,
email = EXCLUDED.email,
updated_at = now()
`
type UpsertNotificationPrefParams struct {
UserID pgtype.UUID
Kind string
Inbox bool
Phone bool
Email bool
}
func (q *Queries) UpsertNotificationPref(ctx context.Context, arg UpsertNotificationPrefParams) error {
_, err := q.db.Exec(ctx, upsertNotificationPref,
arg.UserID,
arg.Kind,
arg.Inbox,
arg.Phone,
arg.Email,
)
return err
}
@@ -0,0 +1,2 @@
DROP TABLE IF EXISTS user_notification_prefs;
DROP TABLE IF EXISTS user_notifications;
@@ -0,0 +1,74 @@
-- M489: a persistent, per-user notifications inbox (#726).
--
-- The event bus is fire-and-forget: a client that is not connected when a
-- request completes, or when the scanner marks tracks missing, never hears of
-- it. A row here is the durable record; the bus only nudges open clients to
-- come and read it.
--
-- `kind` is CHECK-gated (rule 36). A new kind adds its value here, in the
-- same migration as the code that produces it. The list is mirrored in
-- internal/notifications/kinds.go, and both CHECKs below must agree with it.
CREATE TABLE user_notifications (
id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
user_id uuid NOT NULL REFERENCES users(id) ON DELETE CASCADE,
kind text NOT NULL CHECK (kind IN (
'request_approved',
'request_rejected',
'request_completed',
'request_pending',
'quarantine_flagged',
'scan_failed',
'tracks_missing',
'duplicates_found',
'playback_errors'
)),
payload jsonb NOT NULL DEFAULT '{}'::jsonb,
created_at timestamptz NOT NULL DEFAULT now(),
read_at timestamptz,
-- Burst-prone kinds share one key per user. While a row with that key is
-- unread, a new event updates it in place rather than adding another, so
-- fourteen missing tracks are one notification with a count.
coalesce_key text,
-- Stamped once the row has gone out in an email (or been judged not to
-- need one), so the digest never sends the same item twice.
emailed_at timestamptz
);
-- The inbox list, newest first.
CREATE INDEX user_notifications_user_created_idx
ON user_notifications (user_id, created_at DESC);
-- The unread badge, and the digest's selection of unread rows.
CREATE INDEX user_notifications_unread_idx
ON user_notifications (user_id, created_at DESC)
WHERE read_at IS NULL;
-- At most one unread row per coalesce key per user. The upsert in
-- notifications.sql targets this index.
CREATE UNIQUE INDEX user_notifications_coalesce_idx
ON user_notifications (user_id, coalesce_key)
WHERE read_at IS NULL AND coalesce_key IS NOT NULL;
-- Per user, per kind: which channels a kind reaches. A missing row means the
-- kind's defaults (internal/notifications/kinds.go), so nothing is seeded and
-- a new kind needs no backfill.
CREATE TABLE user_notification_prefs (
user_id uuid NOT NULL REFERENCES users(id) ON DELETE CASCADE,
kind text NOT NULL CHECK (kind IN (
'request_approved',
'request_rejected',
'request_completed',
'request_pending',
'quarantine_flagged',
'scan_failed',
'tracks_missing',
'duplicates_found',
'playback_errors'
)),
inbox boolean NOT NULL,
phone boolean NOT NULL,
email boolean NOT NULL,
updated_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (user_id, kind)
);
+93
View File
@@ -0,0 +1,93 @@
-- M489 notifications inbox (#726). Rows are written only through
-- internal/notifications.Notifier, which decides recipients and honours each
-- recipient's inbox preference before it reaches these.
-- name: InsertNotification :one
INSERT INTO user_notifications (user_id, kind, payload)
VALUES (sqlc.arg(user_id), sqlc.arg(kind), sqlc.arg(payload))
RETURNING id;
-- name: UpsertCoalescedNotification :one
-- One unread row per (user, coalesce_key). A new event while that row is
-- unread updates it in place and moves it back to the top of the inbox.
--
-- sum_count: the payload's `count` adds to the unread row's count rather than
-- replacing it. For kinds whose event is "N more happened" (tracks marked
-- missing). Kinds whose event states the whole current total (pending
-- duplicate groups) pass false and the payload simply replaces.
--
-- emailed_at is cleared so the newer state goes out in the next batch; an
-- emailed-but-unread row that keeps growing is news the user hasn't seen.
INSERT INTO user_notifications (user_id, kind, payload, coalesce_key)
VALUES (sqlc.arg(user_id), sqlc.arg(kind), sqlc.arg(payload), sqlc.arg(coalesce_key))
ON CONFLICT (user_id, coalesce_key) WHERE read_at IS NULL AND coalesce_key IS NOT NULL
DO UPDATE SET
payload = CASE
WHEN sqlc.arg(sum_count)::boolean THEN
EXCLUDED.payload || jsonb_build_object('count',
COALESCE((user_notifications.payload->>'count')::bigint, 0)
+ COALESCE((EXCLUDED.payload->>'count')::bigint, 0))
ELSE EXCLUDED.payload
END,
created_at = now(),
emailed_at = NULL
RETURNING id;
-- name: ListNotifications :many
-- Newest first, keyset-paged on (created_at, id). Pass both cursor halves
-- from the last row of the previous page, or neither for the first page.
SELECT id, kind, payload, created_at, read_at
FROM user_notifications
WHERE user_id = sqlc.arg(user_id)
AND (sqlc.narg(before_created_at)::timestamptz IS NULL
OR (created_at, id) < (sqlc.narg(before_created_at)::timestamptz, sqlc.narg(before_id)::uuid))
ORDER BY created_at DESC, id DESC
LIMIT sqlc.arg(page_limit);
-- name: CountUnreadNotifications :one
SELECT count(*) FROM user_notifications
WHERE user_id = $1 AND read_at IS NULL;
-- name: MarkNotificationRead :execrows
-- Scoped to the owner: another user's id matches no row. Marking an already
-- read row keeps its original read_at and still matches, so a repeat is not
-- mistaken for "not yours".
UPDATE user_notifications
SET read_at = COALESCE(read_at, now())
WHERE id = sqlc.arg(id) AND user_id = sqlc.arg(user_id);
-- name: MarkAllNotificationsRead :execrows
UPDATE user_notifications
SET read_at = now()
WHERE user_id = $1 AND read_at IS NULL;
-- name: TrimNotifications :execrows
-- Retention: read rows go after read_cutoff, and anything at all after
-- any_cutoff, so an inbox nobody opens doesn't grow without bound either.
DELETE FROM user_notifications
WHERE (read_at IS NOT NULL AND read_at < sqlc.arg(read_cutoff))
OR created_at < sqlc.arg(any_cutoff);
-- name: ListAdminUserIDs :many
SELECT id FROM users WHERE is_admin = true ORDER BY created_at, id;
-- name: ListNotificationPrefsForUser :many
SELECT kind, inbox, phone, email
FROM user_notification_prefs
WHERE user_id = $1;
-- name: ListNotificationPrefsForKind :many
-- The stored prefs of one kind for a set of recipients. A recipient with no
-- row has the kind's defaults.
SELECT user_id, inbox, phone, email
FROM user_notification_prefs
WHERE kind = sqlc.arg(kind) AND user_id = ANY(sqlc.arg(user_ids)::uuid[]);
-- name: UpsertNotificationPref :exec
INSERT INTO user_notification_prefs (user_id, kind, inbox, phone, email)
VALUES (sqlc.arg(user_id), sqlc.arg(kind), sqlc.arg(inbox), sqlc.arg(phone), sqlc.arg(email))
ON CONFLICT (user_id, kind) DO UPDATE SET
inbox = EXCLUDED.inbox,
phone = EXCLUDED.phone,
email = EXCLUDED.email,
updated_at = now();
+4
View File
@@ -97,6 +97,10 @@ var dataTables = []string{
"track_loudness", // M464
"album_loudness", // M464
"user_normalization_prefs", // M464
// M489. Both cascade from users, but the operator's own admin row is
// never deleted, so rows a test wrote for it would otherwise survive.
"user_notifications",
"user_notification_prefs",
"tracks",
"albums",
"artists",
+130
View File
@@ -0,0 +1,130 @@
// Package notifications is the one writer of the per-user notifications inbox
// (M489, #726). Producers across the server call Notifier.Notify; nothing else
// inserts into user_notifications.
//
// The event bus alone was fire-and-forget: a client that was not connected
// when a request completed, or when tracks went missing, never heard of it. A
// row here is the durable record. The bus only nudges open clients to come
// and read it.
package notifications
// Kind names one sort of notification. The set is CHECK-gated in migration
// 0073 (rule 36): a new kind adds its value there, in the same change, and
// TestKindsMatchMigrationCheck fails until it does.
type Kind string
const (
KindRequestApproved Kind = "request_approved"
KindRequestRejected Kind = "request_rejected"
KindRequestCompleted Kind = "request_completed"
KindRequestPending Kind = "request_pending"
KindQuarantineFlagged Kind = "quarantine_flagged"
KindScanFailed Kind = "scan_failed"
KindTracksMissing Kind = "tracks_missing"
KindDuplicatesFound Kind = "duplicates_found"
KindPlaybackErrors Kind = "playback_errors"
)
// Audience is who a kind can reach. Admin kinds are never offered to, or
// stored for, a non-admin.
type Audience int
const (
AudienceRequester Audience = iota
AudienceAdmin
)
// EmailGroup is how a kind's emails are grouped. Nothing is emailed per event.
type EmailGroup int
const (
// EmailBatch: one email per batch window, holding everything that
// accumulated since the first un-emailed item.
EmailBatch EmailGroup = iota
// EmailSummary: at most one summary a day, at a set local hour. For new
// music arriving, which comes in bursts and is not urgent.
EmailSummary
)
// Channels are where a kind reaches a user. Inbox is the master switch: with
// it off nothing is stored, so there is nothing for the phone or email to
// deliver. Effective applies that.
type Channels struct {
Inbox bool
Phone bool
Email bool
}
// Effective is what will actually be delivered: phone and email ride on the
// inbox row, so they cannot be on without it.
func (c Channels) Effective() Channels {
return Channels{Inbox: c.Inbox, Phone: c.Inbox && c.Phone, Email: c.Inbox && c.Email}
}
type spec struct {
audience Audience
// coalesce: while one of these is unread, a new event updates it in place
// rather than adding a row. For the burst-prone admin kinds.
coalesce bool
// sumCount: a coalesced update adds the payload's `count` to the unread
// row's, because each event is "N more". Without it the newer payload
// replaces the older, because each event states the whole current total.
sumCount bool
group EmailGroup
defaults Channels
}
var (
allOn = Channels{Inbox: true, Phone: true, Email: true}
healthAlert = Channels{Inbox: true, Phone: true, Email: false}
)
var specs = map[Kind]spec{
KindRequestApproved: {audience: AudienceRequester, group: EmailBatch, defaults: allOn},
KindRequestRejected: {audience: AudienceRequester, group: EmailBatch, defaults: allOn},
KindRequestCompleted: {audience: AudienceRequester, group: EmailSummary, defaults: allOn},
KindRequestPending: {audience: AudienceAdmin, group: EmailBatch, defaults: allOn},
KindQuarantineFlagged: {audience: AudienceAdmin, group: EmailBatch, defaults: allOn},
KindScanFailed: {audience: AudienceAdmin, coalesce: true, sumCount: true, group: EmailBatch, defaults: healthAlert},
KindTracksMissing: {audience: AudienceAdmin, coalesce: true, sumCount: true, group: EmailBatch, defaults: healthAlert},
KindDuplicatesFound: {audience: AudienceAdmin, coalesce: true, group: EmailBatch, defaults: healthAlert},
KindPlaybackErrors: {audience: AudienceAdmin, coalesce: true, group: EmailBatch, defaults: healthAlert},
}
// order is the display order for settings screens: the requester's own kinds
// first, then the admin ones.
var order = []Kind{
KindRequestApproved,
KindRequestRejected,
KindRequestCompleted,
KindRequestPending,
KindQuarantineFlagged,
KindScanFailed,
KindTracksMissing,
KindDuplicatesFound,
KindPlaybackErrors,
}
// Kinds returns every kind in display order.
func Kinds() []Kind { return append([]Kind(nil), order...) }
// Valid reports whether k is a known kind.
func (k Kind) Valid() bool { _, ok := specs[k]; return ok }
// AdminOnly reports whether k reaches only admins.
func (k Kind) AdminOnly() bool { return specs[k].audience == AudienceAdmin }
// Defaults are the channels a user has for k until they change them.
func (k Kind) Defaults() Channels { return specs[k].defaults }
// EmailGroup is how k's emails are grouped.
func (k Kind) EmailGroup() EmailGroup { return specs[k].group }
// coalesceKey is the key k's rows coalesce on, or "" when k never coalesces.
// One key per kind: two unread "tracks missing" rows would only split a count.
func (k Kind) coalesceKey() string {
if specs[k].coalesce {
return string(k)
}
return ""
}
+61
View File
@@ -0,0 +1,61 @@
package notifications
import (
"context"
"testing"
"github.com/stretchr/testify/require"
)
func TestEffective_PhoneAndEmailRideOnInbox(t *testing.T) {
cases := []struct {
in, want Channels
}{
{Channels{true, true, true}, Channels{true, true, true}},
{Channels{true, false, true}, Channels{true, false, true}},
{Channels{false, true, true}, Channels{false, false, false}},
{Channels{false, false, false}, Channels{false, false, false}},
}
for _, c := range cases {
require.Equal(t, c.want, c.in.Effective(), "in=%+v", c.in)
}
}
// The display order and the spec table must name the same kinds, or a kind
// is either unreachable from settings or has no behaviour.
func TestKinds_OrderAndSpecsAgree(t *testing.T) {
require.Len(t, order, len(specs))
seen := map[Kind]bool{}
for _, k := range Kinds() {
require.True(t, k.Valid(), "%s in order but not in specs", k)
require.False(t, seen[k], "%s listed twice", k)
seen[k] = true
}
}
func TestKinds_AdminKindsDefaultEmailOffExceptRequestQueue(t *testing.T) {
// The burst-prone health kinds stay out of email unless an admin opts
// in; the request queue and quarantine flags are things an admin acts on.
for _, k := range []Kind{KindScanFailed, KindTracksMissing, KindDuplicatesFound, KindPlaybackErrors} {
require.True(t, k.AdminOnly(), k)
require.False(t, k.Defaults().Email, k)
require.NotEmpty(t, k.coalesceKey(), "%s should coalesce", k)
}
for _, k := range []Kind{KindRequestApproved, KindRequestRejected, KindRequestCompleted} {
require.False(t, k.AdminOnly(), k)
require.Equal(t, Channels{true, true, true}, k.Defaults(), k)
require.Empty(t, k.coalesceKey(), "%s must never coalesce: each request is its own news", k)
}
require.Equal(t, EmailSummary, KindRequestCompleted.EmailGroup())
require.Equal(t, EmailBatch, KindRequestApproved.EmailGroup())
}
func TestNotify_NilNotifierIsANoOp(t *testing.T) {
var n *Notifier
require.NoError(t, n.Notify(context.Background(), KindRequestApproved, Recipients{}, nil))
}
func TestNotify_UnknownKindIsRefusedBeforeTouchingTheDB(t *testing.T) {
n := New(nil, nil, nil)
require.Error(t, n.Notify(context.Background(), Kind("bogus"), Recipients{}, nil))
}
+198
View File
@@ -0,0 +1,198 @@
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); 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
}
func write(ctx context.Context, q *dbq.Queries, kind Kind, userID pgtype.UUID, body []byte) error {
key := kind.coalesceKey()
if key == "" {
_, err := q.InsertNotification(ctx, dbq.InsertNotificationParams{
UserID: userID, Kind: string(kind), Payload: body,
})
return err
}
_, err := q.UpsertCoalescedNotification(ctx, dbq.UpsertCoalescedNotificationParams{
UserID: userID,
Kind: string(kind),
Payload: body,
CoalesceKey: &key,
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{},
})
}
}
+267
View File
@@ -0,0 +1,267 @@
package notifications_test
import (
"context"
"encoding/json"
"io"
"log/slog"
"os"
"testing"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/stretchr/testify/require"
"git.fabledsword.com/bvandeusen/minstrel/internal/db"
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
"git.fabledsword.com/bvandeusen/minstrel/internal/dbtest"
"git.fabledsword.com/bvandeusen/minstrel/internal/eventbus"
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
)
func testPool(t *testing.T) *pgxpool.Pool {
t.Helper()
if testing.Short() {
t.Skip("skipping integration test in -short mode")
}
dsn := os.Getenv("MINSTREL_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("MINSTREL_TEST_DATABASE_URL not set")
}
if err := db.Migrate(dsn, slog.New(slog.NewTextHandler(io.Discard, nil))); err != nil {
t.Fatalf("migrate: %v", err)
}
pool, err := pgxpool.New(context.Background(), dsn)
if err != nil {
t.Fatalf("pool: %v", err)
}
t.Cleanup(pool.Close)
dbtest.ResetDB(t, pool)
return pool
}
func mkUser(t *testing.T, q *dbq.Queries, name string, admin bool) pgtype.UUID {
t.Helper()
u, err := q.CreateUser(context.Background(), dbq.CreateUserParams{
Username: dbtest.TestUserPrefix + name,
PasswordHash: "x",
ApiTokenHash: "tok-" + name, // unique per user
IsAdmin: admin,
})
require.NoError(t, err)
return u.ID
}
func rowsFor(t *testing.T, q *dbq.Queries, user pgtype.UUID) []dbq.ListNotificationsRow {
t.Helper()
rows, err := q.ListNotifications(context.Background(), dbq.ListNotificationsParams{UserID: user, PageLimit: 100})
require.NoError(t, err)
return rows
}
func count(t *testing.T, payload []byte) int64 {
t.Helper()
var p struct {
Count int64 `json:"count"`
}
require.NoError(t, json.Unmarshal(payload, &p))
return p.Count
}
func TestNotify_WritesARowAndNudgesOnlyThatUser(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
alice := mkUser(t, q, "alice", false)
bob := mkUser(t, q, "bob", false)
bus := eventbus.New()
events, unsub := bus.Subscribe(8)
defer unsub()
n := notifications.New(pool, bus, nil)
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved,
notifications.ToUser(alice), map[string]any{"request_id": "r1", "title": "WWW"}))
rows := rowsFor(t, q, alice)
require.Len(t, rows, 1)
require.Equal(t, "request_approved", rows[0].Kind)
require.False(t, rows[0].ReadAt.Valid)
require.JSONEq(t, `{"request_id":"r1","title":"WWW"}`, string(rows[0].Payload))
require.Empty(t, rowsFor(t, q, bob))
select {
case e := <-events:
require.Equal(t, notifications.EventCreated, e.Kind)
require.Equal(t, uuid.UUID(alice.Bytes).String(), e.UserID)
require.Empty(t, e.Data, "the nudge carries no content; the client fetches")
case <-time.After(time.Second):
t.Fatal("no nudge published")
}
}
func TestNotify_AdminKindReachesAdminsOnly_AndNeverTheExcepted(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
requester := mkUser(t, q, "admin-requester", true)
other := mkUser(t, q, "admin-other", true)
plain := mkUser(t, q, "plain", false)
n := notifications.New(pool, nil, nil)
// A plain user named directly must still get nothing: the kind is admin-only.
to := notifications.ToAdmins(requester)
to.Users = []pgtype.UUID{plain}
require.NoError(t, n.Notify(ctx, notifications.KindRequestPending, to, nil))
require.Empty(t, rowsFor(t, q, requester), "the admin who filed it needs no notice")
require.Len(t, rowsFor(t, q, other), 1)
require.Empty(t, rowsFor(t, q, plain))
}
func TestNotify_CoalescedCountAddsUpWhileUnread_AndStartsAfreshOnceRead(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
admin := mkUser(t, q, "admin", true)
n := notifications.New(pool, nil, nil)
to := notifications.ToUser(admin)
require.NoError(t, n.Notify(ctx, notifications.KindTracksMissing, to, map[string]any{"count": 3}))
require.NoError(t, n.Notify(ctx, notifications.KindTracksMissing, to, map[string]any{"count": 4}))
rows := rowsFor(t, q, admin)
require.Len(t, rows, 1, "a burst is one notification")
require.Equal(t, int64(7), count(t, rows[0].Payload))
_, err := q.MarkNotificationRead(ctx, dbq.MarkNotificationReadParams{ID: rows[0].ID, UserID: admin})
require.NoError(t, err)
require.NoError(t, n.Notify(ctx, notifications.KindTracksMissing, to, map[string]any{"count": 2}))
rows = rowsFor(t, q, admin)
require.Len(t, rows, 2, "after it is read, the next event is new news")
require.Equal(t, int64(2), count(t, rows[0].Payload))
require.False(t, rows[0].ReadAt.Valid)
require.True(t, rows[1].ReadAt.Valid)
}
func TestNotify_CoalescedTotalReplacesWhenEventsStateTheWhole(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
admin := mkUser(t, q, "admin", true)
n := notifications.New(pool, nil, nil)
// duplicates_found states the pending total each time; it must not sum.
require.NoError(t, n.Notify(ctx, notifications.KindDuplicatesFound, notifications.ToUser(admin), map[string]any{"count": 5}))
require.NoError(t, n.Notify(ctx, notifications.KindDuplicatesFound, notifications.ToUser(admin), map[string]any{"count": 3}))
rows := rowsFor(t, q, admin)
require.Len(t, rows, 1)
require.Equal(t, int64(3), count(t, rows[0].Payload))
}
func TestNotify_RequestKindsNeverCoalesce(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
alice := mkUser(t, q, "alice", false)
n := notifications.New(pool, nil, nil)
for _, id := range []string{"r1", "r2"} {
require.NoError(t, n.Notify(ctx, notifications.KindRequestCompleted, notifications.ToUser(alice), map[string]any{"request_id": id}))
}
require.Len(t, rowsFor(t, q, alice), 2)
}
func TestNotify_InboxOffStoresNothing_AndDefaultsApplyWithoutAPref(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
muted := mkUser(t, q, "muted", false)
fresh := mkUser(t, q, "fresh", false)
require.NoError(t, q.UpsertNotificationPref(ctx, dbq.UpsertNotificationPrefParams{
UserID: muted, Kind: "request_rejected", Inbox: false, Phone: true, Email: true,
}))
n := notifications.New(pool, nil, nil)
to := notifications.Recipients{Users: []pgtype.UUID{muted, fresh}}
require.NoError(t, n.Notify(ctx, notifications.KindRequestRejected, to, nil))
require.Empty(t, rowsFor(t, q, muted), "inbox off is off, whatever phone and email say")
require.Len(t, rowsFor(t, q, fresh), 1)
}
func TestMarkRead_ScopedToTheOwner_AndIdempotent(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
alice := mkUser(t, q, "alice", false)
mallory := mkUser(t, q, "mallory", false)
n := notifications.New(pool, nil, nil)
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved, notifications.ToUser(alice), nil))
id := rowsFor(t, q, alice)[0].ID
got, err := q.MarkNotificationRead(ctx, dbq.MarkNotificationReadParams{ID: id, UserID: mallory})
require.NoError(t, err)
require.Zero(t, got, "another user's id matches nothing")
require.False(t, rowsFor(t, q, alice)[0].ReadAt.Valid)
for i := 0; i < 2; i++ {
got, err = q.MarkNotificationRead(ctx, dbq.MarkNotificationReadParams{ID: id, UserID: alice})
require.NoError(t, err)
require.Equal(t, int64(1), got, "a repeat still matches, so it is not mistaken for 'not yours'")
}
unread, err := q.CountUnreadNotifications(ctx, alice)
require.NoError(t, err)
require.Zero(t, unread)
}
// Every kind the code knows must pass both CHECKs in the real schema. This is
// rule 36's guard: a kind added in Go without the migration fails here, not in
// production at INSERT time.
func TestEveryKindPassesTheSchemaChecks(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
user := mkUser(t, q, "kinds", true)
for _, k := range notifications.Kinds() {
_, err := q.InsertNotification(ctx, dbq.InsertNotificationParams{UserID: user, Kind: string(k), Payload: []byte(`{}`)})
require.NoError(t, err, "user_notifications rejects %s", k)
require.NoError(t, q.UpsertNotificationPref(ctx, dbq.UpsertNotificationPrefParams{
UserID: user, Kind: string(k), Inbox: true, Phone: true, Email: true,
}), "user_notification_prefs rejects %s", k)
}
// And the check is real: an unknown kind is refused.
_, err := q.InsertNotification(ctx, dbq.InsertNotificationParams{UserID: user, Kind: "bogus", Payload: []byte(`{}`)})
require.Error(t, err)
}
func TestRetention_TrimsOldReadRowsAndAncientUnreadOnes(t *testing.T) {
pool := testPool(t)
q := dbq.New(pool)
ctx := context.Background()
user := mkUser(t, q, "retention", false)
now := time.Now()
insert := func(created time.Time, read *time.Time) {
var readAt any
if read != nil {
readAt = *read
}
_, err := pool.Exec(ctx,
`INSERT INTO user_notifications (user_id, kind, created_at, read_at) VALUES ($1, 'request_approved', $2, $3)`,
user, created, readAt)
require.NoError(t, err)
}
day := 24 * time.Hour
oldRead := now.Add(-100 * day)
recentRead := now.Add(-10 * day)
insert(now.Add(-120*day), &oldRead) // read 100 days ago: goes
insert(now.Add(-20*day), &recentRead) // read 10 days ago: stays
insert(now.Add(-400*day), nil) // unread but over a year old: goes
insert(now.Add(-200*day), nil) // unread, 200 days: stays
r := notifications.NewRetention(pool, slog.New(slog.NewTextHandler(io.Discard, nil)))
require.Equal(t, int64(2), r.TrimOnce(ctx, now))
require.Len(t, rowsFor(t, q, user), 2)
}
+66
View File
@@ -0,0 +1,66 @@
package notifications
import (
"context"
"log/slog"
"time"
"github.com/jackc/pgx/v5/pgtype"
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
)
const (
// ReadRetention is how long a read notification is kept. An inbox is a
// record of what happened lately, not an archive.
ReadRetention = 90 * 24 * time.Hour
// AnyRetention bounds an inbox nobody opens: past it, even unread rows go.
AnyRetention = 365 * 24 * time.Hour
// retentionInterval matches the library_changes compactor's daily tick.
retentionInterval = 24 * time.Hour
)
// Retention trims old notifications on a daily tick, in the shape of the
// library_changes compactor (internal/sync/compactor.go).
type Retention struct {
db dbq.DBTX
logger *slog.Logger
}
// NewRetention returns a Retention trimming through db.
func NewRetention(db dbq.DBTX, logger *slog.Logger) *Retention {
return &Retention{db: db, logger: logger}
}
// Run blocks until ctx is cancelled. It trims once at startup, so a process
// that has been down a while catches up, then daily.
func (r *Retention) Run(ctx context.Context) {
r.TrimOnce(ctx, time.Now())
t := time.NewTicker(retentionInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case now := <-t.C:
r.TrimOnce(ctx, now)
}
}
}
// TrimOnce deletes what has outlived its retention as of now. Errors are
// logged, never fatal: the next tick tries again.
func (r *Retention) TrimOnce(ctx context.Context, now time.Time) int64 {
n, err := dbq.New(r.db).TrimNotifications(ctx, dbq.TrimNotificationsParams{
ReadCutoff: pgtype.Timestamptz{Time: now.Add(-ReadRetention), Valid: true},
AnyCutoff: pgtype.Timestamptz{Time: now.Add(-AnyRetention), Valid: true},
})
if err != nil {
r.logger.Warn("notifications retention: trim failed", "err", err)
return 0
}
if n > 0 {
r.logger.Info("notifications retention: trimmed rows", "count", n)
}
return n
}