Files
minstrel/internal/db/dbq/notifications.sql.go
T
bvandeusenandClaude Opus 5.5 8bf333e748
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
feat(notifications): the inbox store, one writer, coalescing and retention (#5338)
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>
2026-10-08 07:06:39 -04:00

347 lines
9.1 KiB
Go

// 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
}