diff --git a/cmd/minstrel/main.go b/cmd/minstrel/main.go index 92a7ac43..a276ee53 100644 --- a/cmd/minstrel/main.go +++ b/cmd/minstrel/main.go @@ -26,6 +26,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrconfig" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrrequests" "git.fabledsword.com/bvandeusen/minstrel/internal/logging" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" "git.fabledsword.com/bvandeusen/minstrel/internal/playlists" "git.fabledsword.com/bvandeusen/minstrel/internal/reacquisition" "git.fabledsword.com/bvandeusen/minstrel/internal/recsettings" @@ -345,6 +346,10 @@ func run() error { libraryChangesCompactor := syncpkg.NewCompactor(pool, logger.With("component", "library_changes_compactor")) go libraryChangesCompactor.Run(ctx) + // Notifications inbox retention (M489). Daily: read rows go after 90 + // days, anything at all after a year. + go notifications.NewRetention(pool, logger.With("component", "notifications_retention")).Run(ctx) + // Per-user system-playlist scheduler (#392 Half B). Fires each // active user's daily build at 03:00 in their stored timezone. // Replaces the 24h-anchored cron loop (removed in the next commit diff --git a/internal/db/dbq/models.go b/internal/db/dbq/models.go index 3fb3d24f..6e0fd731 100644 --- a/internal/db/dbq/models.go +++ b/internal/db/dbq/models.go @@ -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 diff --git a/internal/db/dbq/notifications.sql.go b/internal/db/dbq/notifications.sql.go new file mode 100644 index 00000000..82e43f42 --- /dev/null +++ b/internal/db/dbq/notifications.sql.go @@ -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 +} diff --git a/internal/db/migrations/0073_user_notifications.down.sql b/internal/db/migrations/0073_user_notifications.down.sql new file mode 100644 index 00000000..df94a29d --- /dev/null +++ b/internal/db/migrations/0073_user_notifications.down.sql @@ -0,0 +1,2 @@ +DROP TABLE IF EXISTS user_notification_prefs; +DROP TABLE IF EXISTS user_notifications; diff --git a/internal/db/migrations/0073_user_notifications.up.sql b/internal/db/migrations/0073_user_notifications.up.sql new file mode 100644 index 00000000..945bb0d1 --- /dev/null +++ b/internal/db/migrations/0073_user_notifications.up.sql @@ -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) +); diff --git a/internal/db/queries/notifications.sql b/internal/db/queries/notifications.sql new file mode 100644 index 00000000..8abc10fe --- /dev/null +++ b/internal/db/queries/notifications.sql @@ -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(); diff --git a/internal/dbtest/reset.go b/internal/dbtest/reset.go index 4dc79322..beca6a84 100644 --- a/internal/dbtest/reset.go +++ b/internal/dbtest/reset.go @@ -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", diff --git a/internal/notifications/kinds.go b/internal/notifications/kinds.go new file mode 100644 index 00000000..31585a54 --- /dev/null +++ b/internal/notifications/kinds.go @@ -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 "" +} diff --git a/internal/notifications/kinds_test.go b/internal/notifications/kinds_test.go new file mode 100644 index 00000000..0d9d8856 --- /dev/null +++ b/internal/notifications/kinds_test.go @@ -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)) +} diff --git a/internal/notifications/notifier.go b/internal/notifications/notifier.go new file mode 100644 index 00000000..3401a031 --- /dev/null +++ b/internal/notifications/notifier.go @@ -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{}, + }) + } +} diff --git a/internal/notifications/notifier_db_test.go b/internal/notifications/notifier_db_test.go new file mode 100644 index 00000000..a551e135 --- /dev/null +++ b/internal/notifications/notifier_db_test.go @@ -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) +} diff --git a/internal/notifications/retention.go b/internal/notifications/retention.go new file mode 100644 index 00000000..a329cbc1 --- /dev/null +++ b/internal/notifications/retention.go @@ -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 +}