feat(notifications): grouped email digest, new music as a daily summary (#5346)
release / govulncheck (push) Successful in 21s
release / web (push) Successful in 1m19s
release / go (push) Successful in 1m39s
release / integration (push) Successful in 5m27s
release / android (push) Successful in 5m47s
release / Build signed APK (releases and dev) (push) Successful in 5m34s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 26s
release / Verify release artifacts (tag releases only) (push) Skipped
release / govulncheck (push) Successful in 21s
release / web (push) Successful in 1m19s
release / go (push) Successful in 1m39s
release / integration (push) Successful in 5m27s
release / android (push) Successful in 5m47s
release / Build signed APK (releases and dev) (push) Successful in 5m34s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 26s
release / Verify release artifacts (tag releases only) (push) Skipped
Nothing is emailed per event. New music (request_completed) goes out at most once a day, at the summary hour in each user's own timezone, grouped by artist. Everything else is batched: one email a window after the first un-emailed item, holding whatever accumulated. - Migration 0074: notification_email_settings (summary hour, batch window, admin-configurable) and user_notification_email_state (batch start, last sent, failures and retry_after per user and group). Existing rows are stamped emailed so the upgrade sends no backlog. - The Notifier stamps emailed_at at write time when the recipient's email channel is off, so turning email on later doesn't send old items. - Read rows are never selected. A row is stamped only after the mailer accepts, in one transaction with the state, against the read's clock, so a coalesced row updated mid-send stays pending. - A failed send backs off 5m doubling to 6h; SMTP not configured just waits. - Links come from the public address; without one the email has none. - The mailer now RFC 2047-encodes subjects and strips line breaks from them. - Admin → Integrations gains a Notification emails card. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,55 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/apierror"
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
|
||||
)
|
||||
|
||||
// notificationEmailBody is the wire shape for GET and PUT
|
||||
// /api/admin/notification-email (M489 #5346): when the grouped emails go out.
|
||||
type notificationEmailBody struct {
|
||||
SummaryHour int32 `json:"summary_hour"`
|
||||
BatchWindowMinutes int32 `json:"batch_window_minutes"`
|
||||
}
|
||||
|
||||
func notificationEmailBodyOf(s notifications.EmailSettings) notificationEmailBody {
|
||||
return notificationEmailBody{SummaryHour: s.SummaryHour, BatchWindowMinutes: s.BatchWindowMinutes}
|
||||
}
|
||||
|
||||
// handleGetNotificationEmail implements GET /api/admin/notification-email.
|
||||
func (h *handlers) handleGetNotificationEmail(w http.ResponseWriter, r *http.Request) {
|
||||
s, err := notifications.LoadEmailSettings(r.Context(), dbq.New(h.pool))
|
||||
if err != nil {
|
||||
writeErrWithLog(w, h.logger, "admin notification email: load failed", apierror.Internal(err))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, notificationEmailBodyOf(s))
|
||||
}
|
||||
|
||||
// handleUpdateNotificationEmail implements PUT /api/admin/notification-email.
|
||||
// A whole-row write; the digest reads it on its next tick.
|
||||
func (h *handlers) handleUpdateNotificationEmail(w http.ResponseWriter, r *http.Request) {
|
||||
var req notificationEmailBody
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeErr(w, apierror.BadRequest("invalid_body", "malformed JSON"))
|
||||
return
|
||||
}
|
||||
saved, err := notifications.SaveEmailSettings(r.Context(), dbq.New(h.pool), notifications.EmailSettings{
|
||||
SummaryHour: req.SummaryHour,
|
||||
BatchWindowMinutes: req.BatchWindowMinutes,
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, notifications.ErrEmailSettingOutOfRange) {
|
||||
writeErr(w, apierror.BadRequest("invalid_setting", err.Error()))
|
||||
return
|
||||
}
|
||||
writeErrWithLog(w, h.logger, "admin notification email: update failed", apierror.Internal(err))
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, notificationEmailBodyOf(saved))
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestUpdateNotificationEmail_Rejects(t *testing.T) {
|
||||
// No pool: every case is refused before the database is reached.
|
||||
h := &handlers{logger: slog.New(slog.NewTextHandler(io.Discard, nil))}
|
||||
for name, tc := range map[string]struct {
|
||||
body string
|
||||
code string
|
||||
mentions string
|
||||
}{
|
||||
"an hour past 23": {body: `{"summary_hour":24,"batch_window_minutes":60}`, code: "invalid_setting", mentions: "summary_hour"},
|
||||
"a window under 15m": {body: `{"summary_hour":9,"batch_window_minutes":5}`, code: "invalid_setting", mentions: "batch_window_minutes"},
|
||||
"a body missing both": {body: `{}`, code: "invalid_setting", mentions: "batch_window_minutes"},
|
||||
"malformed JSON": {body: `{"summary_hour":`, code: "invalid_body"},
|
||||
} {
|
||||
rec := httptest.NewRecorder()
|
||||
h.handleUpdateNotificationEmail(rec, httptest.NewRequest(
|
||||
http.MethodPut, "/api/admin/notification-email", strings.NewReader(tc.body)))
|
||||
body := rec.Body.String()
|
||||
if rec.Code != http.StatusBadRequest || !strings.Contains(body, `"`+tc.code+`"`) || !strings.Contains(body, tc.mentions) {
|
||||
t.Errorf("%s: status %d body %s; want 400 %s mentioning %q", name, rec.Code, body, tc.code, tc.mentions)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -300,6 +300,8 @@ func Mount(r chi.Router, pool *pgxpool.Pool, logger *slog.Logger, events *playev
|
||||
admin.Get("/smtp-config", h.handleGetSMTPConfig)
|
||||
admin.Put("/smtp-config", h.handleUpdateSMTPConfig)
|
||||
admin.Post("/smtp-config/test", h.handleTestSMTPConfig)
|
||||
admin.Get("/notification-email", h.handleGetNotificationEmail)
|
||||
admin.Put("/notification-email", h.handleUpdateNotificationEmail)
|
||||
|
||||
// Recommendation tuning lab (#1250): scoring-weight
|
||||
// profiles + taste-build knobs, DB-backed, live effect.
|
||||
|
||||
@@ -472,6 +472,13 @@ type NetworkSetting struct {
|
||||
PublicUrl string
|
||||
}
|
||||
|
||||
type NotificationEmailSetting struct {
|
||||
ID bool
|
||||
SummaryHour int32
|
||||
BatchWindowMinutes int32
|
||||
UpdatedAt pgtype.Timestamptz
|
||||
}
|
||||
|
||||
type PasswordReset struct {
|
||||
Token string
|
||||
UserID pgtype.UUID
|
||||
@@ -842,6 +849,15 @@ type UserNotification struct {
|
||||
EmailedAt pgtype.Timestamptz
|
||||
}
|
||||
|
||||
type UserNotificationEmailState struct {
|
||||
UserID pgtype.UUID
|
||||
EmailGroup string
|
||||
BatchOpenedAt pgtype.Timestamptz
|
||||
LastSentAt pgtype.Timestamptz
|
||||
Failures int32
|
||||
RetryAfter pgtype.Timestamptz
|
||||
}
|
||||
|
||||
type UserNotificationPref struct {
|
||||
UserID pgtype.UUID
|
||||
Kind string
|
||||
|
||||
@@ -23,24 +23,50 @@ func (q *Queries) CountUnreadNotifications(ctx context.Context, userID pgtype.UU
|
||||
return count, err
|
||||
}
|
||||
|
||||
const getNotificationEmailSettings = `-- name: GetNotificationEmailSettings :one
|
||||
SELECT id, summary_hour, batch_window_minutes, updated_at FROM notification_email_settings WHERE id = true
|
||||
`
|
||||
|
||||
func (q *Queries) GetNotificationEmailSettings(ctx context.Context) (NotificationEmailSetting, error) {
|
||||
row := q.db.QueryRow(ctx, getNotificationEmailSettings)
|
||||
var i NotificationEmailSetting
|
||||
err := row.Scan(
|
||||
&i.ID,
|
||||
&i.SummaryHour,
|
||||
&i.BatchWindowMinutes,
|
||||
&i.UpdatedAt,
|
||||
)
|
||||
return i, err
|
||||
}
|
||||
|
||||
const insertNotification = `-- name: InsertNotification :one
|
||||
|
||||
INSERT INTO user_notifications (user_id, kind, payload)
|
||||
VALUES ($1, $2, $3)
|
||||
INSERT INTO user_notifications (user_id, kind, payload, emailed_at)
|
||||
VALUES ($1, $2, $3,
|
||||
CASE WHEN $4::boolean THEN NULL ELSE now() END)
|
||||
RETURNING id
|
||||
`
|
||||
|
||||
type InsertNotificationParams struct {
|
||||
UserID pgtype.UUID
|
||||
Kind string
|
||||
Payload []byte
|
||||
UserID pgtype.UUID
|
||||
Kind string
|
||||
Payload []byte
|
||||
EmailWanted bool
|
||||
}
|
||||
|
||||
// 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.
|
||||
// email_wanted is the recipient's email channel for this kind as of now. A
|
||||
// row nobody wants emailed is stamped at once, so the digest never has to
|
||||
// judge it and a later "email on" doesn't send an old backlog.
|
||||
func (q *Queries) InsertNotification(ctx context.Context, arg InsertNotificationParams) (pgtype.UUID, error) {
|
||||
row := q.db.QueryRow(ctx, insertNotification, arg.UserID, arg.Kind, arg.Payload)
|
||||
row := q.db.QueryRow(ctx, insertNotification,
|
||||
arg.UserID,
|
||||
arg.Kind,
|
||||
arg.Payload,
|
||||
arg.EmailWanted,
|
||||
)
|
||||
var id pgtype.UUID
|
||||
err := row.Scan(&id)
|
||||
return id, err
|
||||
@@ -70,6 +96,100 @@ func (q *Queries) ListAdminUserIDs(ctx context.Context) ([]pgtype.UUID, error) {
|
||||
return items, nil
|
||||
}
|
||||
|
||||
const listEmailPendingNotifications = `-- name: ListEmailPendingNotifications :many
|
||||
|
||||
SELECT n.id, n.user_id, n.kind, n.payload, n.created_at,
|
||||
u.email::text AS email, u.username, u.display_name, u.timezone,
|
||||
now()::timestamptz AS read_as_of
|
||||
FROM user_notifications n
|
||||
JOIN users u ON u.id = n.user_id
|
||||
WHERE n.read_at IS NULL
|
||||
AND n.emailed_at IS NULL
|
||||
AND u.email IS NOT NULL AND u.email <> ''
|
||||
ORDER BY n.user_id, n.created_at, n.id
|
||||
`
|
||||
|
||||
type ListEmailPendingNotificationsRow struct {
|
||||
ID pgtype.UUID
|
||||
UserID pgtype.UUID
|
||||
Kind string
|
||||
Payload []byte
|
||||
CreatedAt pgtype.Timestamptz
|
||||
Email string
|
||||
Username string
|
||||
DisplayName *string
|
||||
Timezone string
|
||||
ReadAsOf pgtype.Timestamptz
|
||||
}
|
||||
|
||||
// Email digest (#5346) ------------------------------------------------------
|
||||
// Every unread, un-emailed row of a user who has an address, oldest first.
|
||||
// Read rows are never selected: the user has seen them, so they are not news.
|
||||
// read_as_of is the database's clock at the read, handed back to
|
||||
// MarkNotificationsEmailed.
|
||||
func (q *Queries) ListEmailPendingNotifications(ctx context.Context) ([]ListEmailPendingNotificationsRow, error) {
|
||||
rows, err := q.db.Query(ctx, listEmailPendingNotifications)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var items []ListEmailPendingNotificationsRow
|
||||
for rows.Next() {
|
||||
var i ListEmailPendingNotificationsRow
|
||||
if err := rows.Scan(
|
||||
&i.ID,
|
||||
&i.UserID,
|
||||
&i.Kind,
|
||||
&i.Payload,
|
||||
&i.CreatedAt,
|
||||
&i.Email,
|
||||
&i.Username,
|
||||
&i.DisplayName,
|
||||
&i.Timezone,
|
||||
&i.ReadAsOf,
|
||||
); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items = append(items, i)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
const listNotificationEmailState = `-- name: ListNotificationEmailState :many
|
||||
SELECT user_id, email_group, batch_opened_at, last_sent_at, failures, retry_after
|
||||
FROM user_notification_email_state
|
||||
`
|
||||
|
||||
func (q *Queries) ListNotificationEmailState(ctx context.Context) ([]UserNotificationEmailState, error) {
|
||||
rows, err := q.db.Query(ctx, listNotificationEmailState)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var items []UserNotificationEmailState
|
||||
for rows.Next() {
|
||||
var i UserNotificationEmailState
|
||||
if err := rows.Scan(
|
||||
&i.UserID,
|
||||
&i.EmailGroup,
|
||||
&i.BatchOpenedAt,
|
||||
&i.LastSentAt,
|
||||
&i.Failures,
|
||||
&i.RetryAfter,
|
||||
); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items = append(items, i)
|
||||
}
|
||||
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
|
||||
@@ -258,6 +378,31 @@ func (q *Queries) MarkNotificationRead(ctx context.Context, arg MarkNotification
|
||||
return result.RowsAffected(), nil
|
||||
}
|
||||
|
||||
const markNotificationsEmailed = `-- name: MarkNotificationsEmailed :execrows
|
||||
UPDATE user_notifications
|
||||
SET emailed_at = now()
|
||||
WHERE id = ANY($1::uuid[])
|
||||
AND created_at <= $2::timestamptz
|
||||
AND emailed_at IS NULL
|
||||
`
|
||||
|
||||
type MarkNotificationsEmailedParams struct {
|
||||
Ids []pgtype.UUID
|
||||
ReadAsOf pgtype.Timestamptz
|
||||
}
|
||||
|
||||
// Stamps the rows an email carried, or that were judged not to need one.
|
||||
// read_as_of is from ListEmailPendingNotifications: a coalesced row updated
|
||||
// since that read carries a later created_at, holds newer news, and stays
|
||||
// pending for the next email.
|
||||
func (q *Queries) MarkNotificationsEmailed(ctx context.Context, arg MarkNotificationsEmailedParams) (int64, error) {
|
||||
result, err := q.db.Exec(ctx, markNotificationsEmailed, arg.Ids, arg.ReadAsOf)
|
||||
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)
|
||||
@@ -279,20 +424,48 @@ func (q *Queries) TrimNotifications(ctx context.Context, arg TrimNotificationsPa
|
||||
return result.RowsAffected(), nil
|
||||
}
|
||||
|
||||
const updateNotificationEmailSettings = `-- name: UpdateNotificationEmailSettings :one
|
||||
UPDATE notification_email_settings
|
||||
SET summary_hour = $1,
|
||||
batch_window_minutes = $2,
|
||||
updated_at = now()
|
||||
WHERE id = true
|
||||
RETURNING id, summary_hour, batch_window_minutes, updated_at
|
||||
`
|
||||
|
||||
type UpdateNotificationEmailSettingsParams struct {
|
||||
SummaryHour int32
|
||||
BatchWindowMinutes int32
|
||||
}
|
||||
|
||||
// Migration 0074's CHECKs are the backstop behind the service's validation.
|
||||
func (q *Queries) UpdateNotificationEmailSettings(ctx context.Context, arg UpdateNotificationEmailSettingsParams) (NotificationEmailSetting, error) {
|
||||
row := q.db.QueryRow(ctx, updateNotificationEmailSettings, arg.SummaryHour, arg.BatchWindowMinutes)
|
||||
var i NotificationEmailSetting
|
||||
err := row.Scan(
|
||||
&i.ID,
|
||||
&i.SummaryHour,
|
||||
&i.BatchWindowMinutes,
|
||||
&i.UpdatedAt,
|
||||
)
|
||||
return i, err
|
||||
}
|
||||
|
||||
const upsertCoalescedNotification = `-- name: UpsertCoalescedNotification :one
|
||||
INSERT INTO user_notifications (user_id, kind, payload, coalesce_key)
|
||||
VALUES ($1, $2, $3, $4)
|
||||
INSERT INTO user_notifications (user_id, kind, payload, coalesce_key, emailed_at)
|
||||
VALUES ($1, $2, $3, $4,
|
||||
CASE WHEN $5::boolean THEN NULL ELSE now() END)
|
||||
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
|
||||
WHEN $6::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
|
||||
emailed_at = EXCLUDED.emailed_at
|
||||
RETURNING id
|
||||
`
|
||||
|
||||
@@ -301,6 +474,7 @@ type UpsertCoalescedNotificationParams struct {
|
||||
Kind string
|
||||
Payload []byte
|
||||
CoalesceKey *string
|
||||
EmailWanted bool
|
||||
SumCount bool
|
||||
}
|
||||
|
||||
@@ -314,12 +488,14 @@ type UpsertCoalescedNotificationParams struct {
|
||||
//
|
||||
// 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.
|
||||
// Unless email_wanted is false, as in InsertNotification.
|
||||
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.EmailWanted,
|
||||
arg.SumCount,
|
||||
)
|
||||
var id pgtype.UUID
|
||||
@@ -327,6 +503,39 @@ func (q *Queries) UpsertCoalescedNotification(ctx context.Context, arg UpsertCoa
|
||||
return id, err
|
||||
}
|
||||
|
||||
const upsertNotificationEmailState = `-- name: UpsertNotificationEmailState :exec
|
||||
INSERT INTO user_notification_email_state
|
||||
(user_id, email_group, batch_opened_at, last_sent_at, failures, retry_after)
|
||||
VALUES ($1, $2, $3,
|
||||
$4, $5, $6)
|
||||
ON CONFLICT (user_id, email_group) DO UPDATE SET
|
||||
batch_opened_at = EXCLUDED.batch_opened_at,
|
||||
last_sent_at = EXCLUDED.last_sent_at,
|
||||
failures = EXCLUDED.failures,
|
||||
retry_after = EXCLUDED.retry_after
|
||||
`
|
||||
|
||||
type UpsertNotificationEmailStateParams struct {
|
||||
UserID pgtype.UUID
|
||||
EmailGroup string
|
||||
BatchOpenedAt pgtype.Timestamptz
|
||||
LastSentAt pgtype.Timestamptz
|
||||
Failures int32
|
||||
RetryAfter pgtype.Timestamptz
|
||||
}
|
||||
|
||||
func (q *Queries) UpsertNotificationEmailState(ctx context.Context, arg UpsertNotificationEmailStateParams) error {
|
||||
_, err := q.db.Exec(ctx, upsertNotificationEmailState,
|
||||
arg.UserID,
|
||||
arg.EmailGroup,
|
||||
arg.BatchOpenedAt,
|
||||
arg.LastSentAt,
|
||||
arg.Failures,
|
||||
arg.RetryAfter,
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
const upsertNotificationPref = `-- name: UpsertNotificationPref :exec
|
||||
INSERT INTO user_notification_prefs (user_id, kind, inbox, phone, email)
|
||||
VALUES ($1, $2, $3, $4, $5)
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
DROP INDEX IF EXISTS user_notifications_email_pending_idx;
|
||||
DROP TABLE IF EXISTS user_notification_email_state;
|
||||
DROP TABLE IF EXISTS notification_email_settings;
|
||||
@@ -0,0 +1,54 @@
|
||||
-- M489 #5346: notifications by email, grouped. Nothing is emailed per event.
|
||||
--
|
||||
-- New music (request_completed) goes out as at most one summary a day, at a set
|
||||
-- hour in each user's own timezone. Everything else is batched: a batch opens
|
||||
-- at the first un-emailed item and one email goes out a window later, holding
|
||||
-- whatever accumulated. internal/notifications/digest.go is the one sender.
|
||||
|
||||
-- The two knobs, in admin Settings (rule 25). Singleton in the style of
|
||||
-- fingerprint_settings (0061).
|
||||
CREATE TABLE notification_email_settings (
|
||||
id boolean PRIMARY KEY DEFAULT true,
|
||||
-- The local hour (0-23, in each user's timezone) the daily new-music
|
||||
-- summary goes out.
|
||||
summary_hour integer NOT NULL DEFAULT 9,
|
||||
-- How long a batch stays open after its first item before it is sent.
|
||||
batch_window_minutes integer NOT NULL DEFAULT 60,
|
||||
updated_at timestamptz NOT NULL DEFAULT now(),
|
||||
|
||||
CONSTRAINT notification_email_settings_singleton CHECK (id = true),
|
||||
CONSTRAINT notification_email_settings_hour_range
|
||||
CHECK (summary_hour >= 0 AND summary_hour <= 23),
|
||||
CONSTRAINT notification_email_settings_window_range
|
||||
CHECK (batch_window_minutes >= 15 AND batch_window_minutes <= 1440)
|
||||
);
|
||||
INSERT INTO notification_email_settings (id) VALUES (true) ON CONFLICT (id) DO NOTHING;
|
||||
|
||||
-- Per user, per email group: where the digest stands. A missing row is a user
|
||||
-- who has never had a batch open or a summary sent.
|
||||
CREATE TABLE user_notification_email_state (
|
||||
user_id uuid NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
email_group text NOT NULL CHECK (email_group IN ('batch', 'summary')),
|
||||
-- When the open batch started. Held here rather than read off the items,
|
||||
-- because a coalesced item moves its created_at forward on every update
|
||||
-- and would otherwise keep a batch from ever coming due.
|
||||
batch_opened_at timestamptz,
|
||||
-- The last email of this group the mailer accepted. A summary is due once
|
||||
-- per local day, after the summary hour, if this is before it.
|
||||
last_sent_at timestamptz,
|
||||
-- A failed send backs off: nothing is tried for this group before
|
||||
-- retry_after, and the gap doubles with each failure in a row.
|
||||
failures integer NOT NULL DEFAULT 0,
|
||||
retry_after timestamptz,
|
||||
PRIMARY KEY (user_id, email_group)
|
||||
);
|
||||
|
||||
-- The digest's selection: unread rows not yet emailed.
|
||||
CREATE INDEX user_notifications_email_pending_idx
|
||||
ON user_notifications (user_id, created_at)
|
||||
WHERE read_at IS NULL AND emailed_at IS NULL;
|
||||
|
||||
-- Rows written before this migration were never judged for email. Treat them
|
||||
-- as already handled, so the first digest after the upgrade doesn't send a
|
||||
-- backlog nobody asked for.
|
||||
UPDATE user_notifications SET emailed_at = now() WHERE emailed_at IS NULL;
|
||||
@@ -3,8 +3,12 @@
|
||||
-- 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))
|
||||
-- email_wanted is the recipient's email channel for this kind as of now. A
|
||||
-- row nobody wants emailed is stamped at once, so the digest never has to
|
||||
-- judge it and a later "email on" doesn't send an old backlog.
|
||||
INSERT INTO user_notifications (user_id, kind, payload, emailed_at)
|
||||
VALUES (sqlc.arg(user_id), sqlc.arg(kind), sqlc.arg(payload),
|
||||
CASE WHEN sqlc.arg(email_wanted)::boolean THEN NULL ELSE now() END)
|
||||
RETURNING id;
|
||||
|
||||
-- name: UpsertCoalescedNotification :one
|
||||
@@ -18,8 +22,10 @@ RETURNING id;
|
||||
--
|
||||
-- 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))
|
||||
-- Unless email_wanted is false, as in InsertNotification.
|
||||
INSERT INTO user_notifications (user_id, kind, payload, coalesce_key, emailed_at)
|
||||
VALUES (sqlc.arg(user_id), sqlc.arg(kind), sqlc.arg(payload), sqlc.arg(coalesce_key),
|
||||
CASE WHEN sqlc.arg(email_wanted)::boolean THEN NULL ELSE now() END)
|
||||
ON CONFLICT (user_id, coalesce_key) WHERE read_at IS NULL AND coalesce_key IS NOT NULL
|
||||
DO UPDATE SET
|
||||
payload = CASE
|
||||
@@ -30,7 +36,7 @@ DO UPDATE SET
|
||||
ELSE EXCLUDED.payload
|
||||
END,
|
||||
created_at = now(),
|
||||
emailed_at = NULL
|
||||
emailed_at = EXCLUDED.emailed_at
|
||||
RETURNING id;
|
||||
|
||||
-- name: ListNotifications :many
|
||||
@@ -97,3 +103,58 @@ ON CONFLICT (user_id, kind) DO UPDATE SET
|
||||
phone = EXCLUDED.phone,
|
||||
email = EXCLUDED.email,
|
||||
updated_at = now();
|
||||
|
||||
-- Email digest (#5346) ------------------------------------------------------
|
||||
|
||||
-- name: ListEmailPendingNotifications :many
|
||||
-- Every unread, un-emailed row of a user who has an address, oldest first.
|
||||
-- Read rows are never selected: the user has seen them, so they are not news.
|
||||
-- read_as_of is the database's clock at the read, handed back to
|
||||
-- MarkNotificationsEmailed.
|
||||
SELECT n.id, n.user_id, n.kind, n.payload, n.created_at,
|
||||
u.email::text AS email, u.username, u.display_name, u.timezone,
|
||||
now()::timestamptz AS read_as_of
|
||||
FROM user_notifications n
|
||||
JOIN users u ON u.id = n.user_id
|
||||
WHERE n.read_at IS NULL
|
||||
AND n.emailed_at IS NULL
|
||||
AND u.email IS NOT NULL AND u.email <> ''
|
||||
ORDER BY n.user_id, n.created_at, n.id;
|
||||
|
||||
-- name: MarkNotificationsEmailed :execrows
|
||||
-- Stamps the rows an email carried, or that were judged not to need one.
|
||||
-- read_as_of is from ListEmailPendingNotifications: a coalesced row updated
|
||||
-- since that read carries a later created_at, holds newer news, and stays
|
||||
-- pending for the next email.
|
||||
UPDATE user_notifications
|
||||
SET emailed_at = now()
|
||||
WHERE id = ANY(sqlc.arg(ids)::uuid[])
|
||||
AND created_at <= sqlc.arg(read_as_of)::timestamptz
|
||||
AND emailed_at IS NULL;
|
||||
|
||||
-- name: ListNotificationEmailState :many
|
||||
SELECT user_id, email_group, batch_opened_at, last_sent_at, failures, retry_after
|
||||
FROM user_notification_email_state;
|
||||
|
||||
-- name: UpsertNotificationEmailState :exec
|
||||
INSERT INTO user_notification_email_state
|
||||
(user_id, email_group, batch_opened_at, last_sent_at, failures, retry_after)
|
||||
VALUES (sqlc.arg(user_id), sqlc.arg(email_group), sqlc.narg(batch_opened_at),
|
||||
sqlc.narg(last_sent_at), sqlc.arg(failures), sqlc.narg(retry_after))
|
||||
ON CONFLICT (user_id, email_group) DO UPDATE SET
|
||||
batch_opened_at = EXCLUDED.batch_opened_at,
|
||||
last_sent_at = EXCLUDED.last_sent_at,
|
||||
failures = EXCLUDED.failures,
|
||||
retry_after = EXCLUDED.retry_after;
|
||||
|
||||
-- name: GetNotificationEmailSettings :one
|
||||
SELECT * FROM notification_email_settings WHERE id = true;
|
||||
|
||||
-- name: UpdateNotificationEmailSettings :one
|
||||
-- Migration 0074's CHECKs are the backstop behind the service's validation.
|
||||
UPDATE notification_email_settings
|
||||
SET summary_hour = sqlc.arg(summary_hour),
|
||||
batch_window_minutes = sqlc.arg(batch_window_minutes),
|
||||
updated_at = now()
|
||||
WHERE id = true
|
||||
RETURNING *;
|
||||
|
||||
@@ -101,6 +101,7 @@ var dataTables = []string{
|
||||
// never deleted, so rows a test wrote for it would otherwise survive.
|
||||
"user_notifications",
|
||||
"user_notification_prefs",
|
||||
"user_notification_email_state", // #5346
|
||||
"tracks",
|
||||
"albums",
|
||||
"artists",
|
||||
@@ -158,4 +159,11 @@ func ResetDB(t *testing.T, pool *pgxpool.Pool) {
|
||||
); err != nil {
|
||||
t.Fatalf("dbtest.ResetDB reset loudness settings: %v", err)
|
||||
}
|
||||
// Notification email settings (M489 #5346), reset the same way.
|
||||
if _, err := pool.Exec(ctx, `
|
||||
UPDATE notification_email_settings
|
||||
SET summary_hour = DEFAULT, batch_window_minutes = DEFAULT, updated_at = DEFAULT`,
|
||||
); err != nil {
|
||||
t.Fatalf("dbtest.ResetDB reset notification email settings: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -75,6 +75,8 @@ func (r *Reconciler) publishCompleted(ctx context.Context, row dbq.LidarrRequest
|
||||
RequestID: formatUUIDForBus(row.ID),
|
||||
RequestKind: string(row.Kind),
|
||||
Name: DisplayName(row),
|
||||
Artist: row.ArtistName,
|
||||
Title: requestedTitle(row),
|
||||
}
|
||||
if row.MatchedAlbumID.Valid {
|
||||
p.AlbumID = formatUUIDForBus(row.MatchedAlbumID)
|
||||
|
||||
@@ -141,6 +141,18 @@ func DisplayName(row dbq.LidarrRequest) string {
|
||||
return row.ArtistName
|
||||
}
|
||||
|
||||
// requestedTitle is the album or track a request names, or "" for an
|
||||
// artist request. DisplayName is the artist and this, joined.
|
||||
func requestedTitle(row dbq.LidarrRequest) string {
|
||||
switch {
|
||||
case row.Kind == dbq.LidarrRequestKindTrack && row.TrackTitle != nil:
|
||||
return *row.TrackTitle
|
||||
case row.Kind == dbq.LidarrRequestKindAlbum && row.AlbumTitle != nil:
|
||||
return *row.AlbumTitle
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func validateKindFields(p CreateParams) error {
|
||||
if p.LidarrArtistMBID == "" || p.ArtistName == "" {
|
||||
return fmt.Errorf("%w: artist_mbid and artist_name are always required", ErrInvalidKindFields)
|
||||
|
||||
@@ -12,9 +12,11 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"mime"
|
||||
"net"
|
||||
"net/smtp"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -151,7 +153,7 @@ func composeMessage(cfg dbq.SmtpConfig, to, subject, textBody, htmlBody string)
|
||||
}
|
||||
fmt.Fprintf(&buf, "From: %s\r\n", from)
|
||||
fmt.Fprintf(&buf, "To: %s\r\n", to)
|
||||
fmt.Fprintf(&buf, "Subject: %s\r\n", subject)
|
||||
fmt.Fprintf(&buf, "Subject: %s\r\n", encodeSubject(subject))
|
||||
fmt.Fprintf(&buf, "MIME-Version: 1.0\r\n")
|
||||
fmt.Fprintf(&buf, "Content-Type: multipart/alternative; boundary=\"%s\"\r\n\r\n", boundary)
|
||||
|
||||
@@ -169,6 +171,15 @@ func composeMessage(cfg dbq.SmtpConfig, to, subject, textBody, htmlBody string)
|
||||
return buf.Bytes()
|
||||
}
|
||||
|
||||
// encodeSubject makes subject safe as a header value. Line breaks are
|
||||
// replaced, so a name carried into a subject cannot start a header of its
|
||||
// own, and non-ASCII text ("Moe Shop – WWW") is RFC 2047 encoded, which
|
||||
// mime.QEncoding leaves alone when there is nothing to encode.
|
||||
func encodeSubject(subject string) string {
|
||||
subject = strings.NewReplacer("\r\n", " ", "\r", " ", "\n", " ").Replace(subject)
|
||||
return mime.QEncoding.Encode("utf-8", subject)
|
||||
}
|
||||
|
||||
// SentEmail is the in-memory record FakeSender keeps. Tests assert
|
||||
// against these.
|
||||
type SentEmail struct {
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
package mailer
|
||||
|
||||
import (
|
||||
"mime"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestEncodeSubject(t *testing.T) {
|
||||
if got := encodeSubject("Minstrel: 2 updates"); got != "Minstrel: 2 updates" {
|
||||
t.Errorf("plain ASCII changed: %q", got)
|
||||
}
|
||||
|
||||
enc := encodeSubject("Minstrel: Moe Shop – WWW")
|
||||
if !strings.HasPrefix(enc, "=?utf-8?q?") {
|
||||
t.Errorf("non-ASCII not RFC 2047 encoded: %q", enc)
|
||||
}
|
||||
if dec, err := new(mime.WordDecoder).DecodeHeader(enc); err != nil || dec != "Minstrel: Moe Shop – WWW" {
|
||||
t.Errorf("round trip = %q, %v", dec, err)
|
||||
}
|
||||
|
||||
if got := encodeSubject("Hi\r\nBcc: everyone@example.com"); strings.ContainsAny(got, "\r\n") {
|
||||
t.Errorf("a line break survived into the header: %q", got)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,390 @@
|
||||
package notifications
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/mailer"
|
||||
)
|
||||
|
||||
// The email digest (#5346). Nothing is emailed per event. Operator,
|
||||
// 2026-10-08: "music and stuff coming in should land as a summary email and
|
||||
// not everytime same for a approvals the emails should be grouped by a
|
||||
// reasonable amount of time like an hour or so."
|
||||
//
|
||||
// - EmailSummary kinds (new music) go out at most once a day, at the summary
|
||||
// hour in the user's own timezone, listing everything since the last one.
|
||||
// - EmailBatch kinds go out a batch window after the first un-emailed item,
|
||||
// holding whatever accumulated.
|
||||
// - A row the user has read is never emailed: they have seen it.
|
||||
|
||||
const (
|
||||
// digestInterval is how often the digest looks for due groups. It bounds
|
||||
// how late an email can be, not how often one is sent.
|
||||
digestInterval = 5 * time.Minute
|
||||
// maxEmailAge: an item older than this when the digest first gets to it is
|
||||
// not news worth an email. It covers a user who adds an address, or an
|
||||
// SMTP server set up, after weeks of unread notifications.
|
||||
maxEmailAge = 7 * 24 * time.Hour
|
||||
// A failed send waits retryBase, doubling with each failure in a row, up
|
||||
// to retryMax.
|
||||
retryBase = 5 * time.Minute
|
||||
retryMax = 6 * time.Hour
|
||||
)
|
||||
|
||||
const (
|
||||
groupBatch = "batch"
|
||||
groupSummary = "summary"
|
||||
)
|
||||
|
||||
func (g EmailGroup) key() string {
|
||||
if g == EmailSummary {
|
||||
return groupSummary
|
||||
}
|
||||
return groupBatch
|
||||
}
|
||||
|
||||
// groupState is where one user's group stands; user_notification_email_state.
|
||||
type groupState struct {
|
||||
BatchOpenedAt time.Time
|
||||
LastSentAt time.Time
|
||||
Failures int32
|
||||
RetryAfter time.Time
|
||||
}
|
||||
|
||||
// pendingItem is one unread, un-emailed notification.
|
||||
type pendingItem struct {
|
||||
ID pgtype.UUID
|
||||
Kind Kind
|
||||
Payload []byte
|
||||
CreatedAt time.Time
|
||||
}
|
||||
|
||||
// userPlan is what the digest does for one user on one tick.
|
||||
type userPlan struct {
|
||||
// Skip is stamped emailed without sending: email is off for its kind now,
|
||||
// or it is too old to be news.
|
||||
Skip []pendingItem
|
||||
|
||||
Batch []pendingItem
|
||||
BatchDue bool
|
||||
// BatchOpenedAt is the open batch's start, kept until it is sent. Zero
|
||||
// when nothing is pending in the batch group.
|
||||
BatchOpenedAt time.Time
|
||||
|
||||
Summary []pendingItem
|
||||
SummaryDue bool
|
||||
}
|
||||
|
||||
// planUser decides, from the clock and stored state alone, what goes out for
|
||||
// one user. items are oldest first. emailOn is the user's effective email
|
||||
// channel per kind.
|
||||
func planUser(now time.Time, cfg EmailSettings, loc *time.Location, items []pendingItem,
|
||||
emailOn map[Kind]bool, batch, summary groupState) userPlan {
|
||||
var p userPlan
|
||||
for _, it := range items {
|
||||
switch {
|
||||
case !emailOn[it.Kind] || now.Sub(it.CreatedAt) > maxEmailAge:
|
||||
p.Skip = append(p.Skip, it)
|
||||
case it.Kind.EmailGroup() == EmailSummary:
|
||||
p.Summary = append(p.Summary, it)
|
||||
default:
|
||||
p.Batch = append(p.Batch, it)
|
||||
}
|
||||
}
|
||||
if len(p.Batch) > 0 {
|
||||
p.BatchOpenedAt = batch.BatchOpenedAt
|
||||
if p.BatchOpenedAt.IsZero() {
|
||||
p.BatchOpenedAt = p.Batch[0].CreatedAt
|
||||
}
|
||||
p.BatchDue = !now.Before(batch.RetryAfter) && !now.Before(p.BatchOpenedAt.Add(cfg.BatchWindow()))
|
||||
}
|
||||
if len(p.Summary) > 0 {
|
||||
slot := summarySlot(now, loc, int(cfg.SummaryHour))
|
||||
p.SummaryDue = !now.Before(summary.RetryAfter) && !now.Before(slot) && summary.LastSentAt.Before(slot)
|
||||
}
|
||||
return p
|
||||
}
|
||||
|
||||
// summarySlot is today's summary time in loc: the summary hour on the user's
|
||||
// local calendar day. Across a DST change the hour keeps its local meaning,
|
||||
// so the instant moves by the shift. An hour that does not exist that day
|
||||
// (inside a spring-forward gap) is normalised by time.Date to one that does.
|
||||
func summarySlot(now time.Time, loc *time.Location, hour int) time.Time {
|
||||
l := now.In(loc)
|
||||
return time.Date(l.Year(), l.Month(), l.Day(), hour, 0, 0, 0, loc)
|
||||
}
|
||||
|
||||
// retryDelay is the wait after the failures-th failure in a row.
|
||||
func retryDelay(failures int32) time.Duration {
|
||||
d := retryBase
|
||||
for i := int32(1); i < failures && d < retryMax; i++ {
|
||||
d *= 2
|
||||
}
|
||||
return min(d, retryMax)
|
||||
}
|
||||
|
||||
// Digest sends the grouped emails. One per process, on a ticker.
|
||||
type Digest struct {
|
||||
pool *pgxpool.Pool
|
||||
sender mailer.Sender
|
||||
logger *slog.Logger
|
||||
publicURL func(context.Context) string
|
||||
}
|
||||
|
||||
// NewDigest returns a Digest sending through sender. Links in its emails are
|
||||
// built from the public address in admin Settings, never from a request.
|
||||
func NewDigest(pool *pgxpool.Pool, sender mailer.Sender, logger *slog.Logger) *Digest {
|
||||
if logger == nil {
|
||||
logger = slog.Default()
|
||||
}
|
||||
d := &Digest{pool: pool, sender: sender, logger: logger}
|
||||
d.publicURL = func(ctx context.Context) string {
|
||||
row, err := dbq.New(pool).GetNetworkSettings(ctx)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
return row.PublicUrl
|
||||
}
|
||||
return d
|
||||
}
|
||||
|
||||
// Run blocks until ctx is cancelled. It ticks once at startup, so a process
|
||||
// that has been down past a batch or a summary hour catches up, then every
|
||||
// few minutes.
|
||||
func (d *Digest) Run(ctx context.Context) {
|
||||
d.Tick(ctx, time.Now())
|
||||
t := time.NewTicker(digestInterval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case now := <-t.C:
|
||||
d.Tick(ctx, now)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Tick sends every group that is due as of now. Failures are logged and
|
||||
// retried on a later tick; nothing here is fatal.
|
||||
func (d *Digest) Tick(ctx context.Context, now time.Time) {
|
||||
q := dbq.New(d.pool)
|
||||
cfg, err := LoadEmailSettings(ctx, q)
|
||||
if err != nil {
|
||||
d.logger.Warn("notification digest: using default settings", "err", err)
|
||||
}
|
||||
rows, err := q.ListEmailPendingNotifications(ctx)
|
||||
if err != nil {
|
||||
d.logger.Warn("notification digest: list pending", "err", err)
|
||||
return
|
||||
}
|
||||
states, err := d.loadStates(ctx, q)
|
||||
if err != nil {
|
||||
d.logger.Warn("notification digest: load state", "err", err)
|
||||
return
|
||||
}
|
||||
base := d.publicURL(ctx)
|
||||
|
||||
handled := map[pgtype.UUID]bool{}
|
||||
for start := 0; start < len(rows); {
|
||||
end := start + 1
|
||||
for end < len(rows) && rows[end].UserID == rows[start].UserID {
|
||||
end++
|
||||
}
|
||||
u := rows[start]
|
||||
handled[u.UserID] = true
|
||||
d.digestUser(ctx, q, now, cfg, base, rows[start:end], states[u.UserID])
|
||||
start = end
|
||||
}
|
||||
|
||||
// A user with nothing pending at all whose batch is still open read
|
||||
// everything before it came due. Close it, so their next item opens a
|
||||
// fresh batch rather than going out at once.
|
||||
for userID, st := range states {
|
||||
if !handled[userID] && !st[groupBatch].BatchOpenedAt.IsZero() {
|
||||
closed := st[groupBatch]
|
||||
closed.BatchOpenedAt = time.Time{}
|
||||
d.saveState(ctx, q, userID, groupBatch, closed)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// digestUser handles one user's pending rows.
|
||||
func (d *Digest) digestUser(ctx context.Context, q *dbq.Queries, now time.Time, cfg EmailSettings, base string,
|
||||
rows []dbq.ListEmailPendingNotificationsRow, st map[string]groupState) {
|
||||
u := rows[0]
|
||||
emailOn, err := emailChannels(ctx, q, u.UserID)
|
||||
if err != nil {
|
||||
d.logger.Warn("notification digest: read prefs", "err", err)
|
||||
return
|
||||
}
|
||||
items := make([]pendingItem, len(rows))
|
||||
for i, r := range rows {
|
||||
items[i] = pendingItem{ID: r.ID, Kind: Kind(r.Kind), Payload: r.Payload, CreatedAt: r.CreatedAt.Time}
|
||||
}
|
||||
plan := planUser(now, cfg, userLocation(u.Timezone), items, emailOn, st[groupBatch], st[groupSummary])
|
||||
readAsOf := u.ReadAsOf.Time
|
||||
|
||||
if len(plan.Skip) > 0 {
|
||||
if _, err := q.MarkNotificationsEmailed(ctx, dbq.MarkNotificationsEmailedParams{
|
||||
Ids: itemIDs(plan.Skip), ReadAsOf: ts(readAsOf),
|
||||
}); err != nil {
|
||||
d.logger.Warn("notification digest: stamp skipped", "err", err)
|
||||
}
|
||||
}
|
||||
|
||||
// The batch's start is recorded before anything is sent, so a failed send
|
||||
// below keeps it along with its retry. With nothing left in the batch
|
||||
// group (all read, or skipped), an open batch is closed.
|
||||
batchSt := st[groupBatch]
|
||||
if !plan.BatchOpenedAt.Equal(batchSt.BatchOpenedAt) {
|
||||
batchSt.BatchOpenedAt = plan.BatchOpenedAt
|
||||
d.saveState(ctx, q, u.UserID, groupBatch, batchSt)
|
||||
}
|
||||
|
||||
to := recipient{email: u.Email, name: displayName(u)}
|
||||
if plan.SummaryDue {
|
||||
d.send(ctx, q, now, u.UserID, to, EmailSummary, plan.Summary, st[groupSummary], base, readAsOf)
|
||||
}
|
||||
if plan.BatchDue {
|
||||
d.send(ctx, q, now, u.UserID, to, EmailBatch, plan.Batch, batchSt, base, readAsOf)
|
||||
}
|
||||
}
|
||||
|
||||
type recipient struct{ email, name string }
|
||||
|
||||
// send renders and sends one group's email. Only once the mailer has accepted
|
||||
// it are its rows stamped and the group's state reset, in one transaction, so
|
||||
// a failure sends again later and a success never sends twice.
|
||||
func (d *Digest) send(ctx context.Context, q *dbq.Queries, now time.Time, userID pgtype.UUID, to recipient,
|
||||
group EmailGroup, items []pendingItem, st groupState, base string, readAsOf time.Time) {
|
||||
msg, err := renderDigest(group, to.name, items, base)
|
||||
if err != nil {
|
||||
d.logger.Error("notification digest: render", "group", group.key(), "err", err)
|
||||
return
|
||||
}
|
||||
err = d.sender.Send(ctx, to.email, msg.Subject, msg.Text, msg.HTML)
|
||||
if errors.Is(err, mailer.ErrNotConfigured) {
|
||||
// Not a failure to back off from: nothing can go until SMTP is set
|
||||
// up, and the rows wait, unread and un-emailed, until it is.
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
st.Failures++
|
||||
st.RetryAfter = now.Add(retryDelay(st.Failures))
|
||||
d.logger.Warn("notification digest: send failed",
|
||||
"group", group.key(), "failures", st.Failures, "retry_after", st.RetryAfter, "err", err)
|
||||
d.saveState(ctx, q, userID, group.key(), st)
|
||||
return
|
||||
}
|
||||
|
||||
sent := groupState{LastSentAt: now}
|
||||
if err := pgx.BeginFunc(ctx, d.pool, func(tx pgx.Tx) error {
|
||||
tq := dbq.New(tx)
|
||||
if _, err := tq.MarkNotificationsEmailed(ctx, dbq.MarkNotificationsEmailedParams{
|
||||
Ids: itemIDs(items), ReadAsOf: ts(readAsOf),
|
||||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
return tq.UpsertNotificationEmailState(ctx, stateParams(userID, group.key(), sent))
|
||||
}); err != nil {
|
||||
// The email is out but its rows are not stamped, so the next due
|
||||
// group would repeat them. Logged at ERROR: it needs a look.
|
||||
d.logger.Error("notification digest: sent but not recorded", "group", group.key(), "err", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (d *Digest) loadStates(ctx context.Context, q *dbq.Queries) (map[pgtype.UUID]map[string]groupState, error) {
|
||||
rows, err := q.ListNotificationEmailState(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make(map[pgtype.UUID]map[string]groupState, len(rows))
|
||||
for _, r := range rows {
|
||||
if out[r.UserID] == nil {
|
||||
out[r.UserID] = map[string]groupState{}
|
||||
}
|
||||
out[r.UserID][r.EmailGroup] = groupState{
|
||||
BatchOpenedAt: tsTime(r.BatchOpenedAt),
|
||||
LastSentAt: tsTime(r.LastSentAt),
|
||||
Failures: r.Failures,
|
||||
RetryAfter: tsTime(r.RetryAfter),
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (d *Digest) saveState(ctx context.Context, q *dbq.Queries, userID pgtype.UUID, group string, st groupState) {
|
||||
if err := q.UpsertNotificationEmailState(ctx, stateParams(userID, group, st)); err != nil {
|
||||
d.logger.Warn("notification digest: save state", "group", group, "err", err)
|
||||
}
|
||||
}
|
||||
|
||||
// emailChannels is the user's effective email channel for every kind: their
|
||||
// stored preference, or the kind's default.
|
||||
func emailChannels(ctx context.Context, q *dbq.Queries, userID pgtype.UUID) (map[Kind]bool, error) {
|
||||
rows, err := q.ListNotificationPrefsForUser(ctx, userID)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("prefs: %w", err)
|
||||
}
|
||||
out := make(map[Kind]bool, len(order))
|
||||
for _, k := range order {
|
||||
out[k] = k.Defaults().Effective().Email
|
||||
}
|
||||
for _, r := range rows {
|
||||
out[Kind(r.Kind)] = Channels{Inbox: r.Inbox, Phone: r.Phone, Email: r.Email}.Effective().Email
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// userLocation is the user's timezone, or UTC when it does not parse.
|
||||
func userLocation(tz string) *time.Location {
|
||||
if loc, err := time.LoadLocation(tz); err == nil && tz != "" {
|
||||
return loc
|
||||
}
|
||||
return time.UTC
|
||||
}
|
||||
|
||||
func displayName(u dbq.ListEmailPendingNotificationsRow) string {
|
||||
if u.DisplayName != nil && *u.DisplayName != "" {
|
||||
return *u.DisplayName
|
||||
}
|
||||
return u.Username
|
||||
}
|
||||
|
||||
func itemIDs(items []pendingItem) []pgtype.UUID {
|
||||
out := make([]pgtype.UUID, len(items))
|
||||
for i, it := range items {
|
||||
out[i] = it.ID
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func stateParams(userID pgtype.UUID, group string, st groupState) dbq.UpsertNotificationEmailStateParams {
|
||||
return dbq.UpsertNotificationEmailStateParams{
|
||||
UserID: userID,
|
||||
EmailGroup: group,
|
||||
BatchOpenedAt: ts(st.BatchOpenedAt),
|
||||
LastSentAt: ts(st.LastSentAt),
|
||||
Failures: st.Failures,
|
||||
RetryAfter: ts(st.RetryAfter),
|
||||
}
|
||||
}
|
||||
|
||||
func ts(t time.Time) pgtype.Timestamptz { return pgtype.Timestamptz{Time: t, Valid: !t.IsZero()} }
|
||||
|
||||
func tsTime(t pgtype.Timestamptz) time.Time {
|
||||
if !t.Valid {
|
||||
return time.Time{}
|
||||
}
|
||||
return t.Time
|
||||
}
|
||||
@@ -0,0 +1,223 @@
|
||||
package notifications_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/mailer"
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
|
||||
)
|
||||
|
||||
// emailUser is a user with an address, so the digest can reach them.
|
||||
func emailUser(t *testing.T, pool *pgxpool.Pool, name string, admin bool) pgtype.UUID {
|
||||
t.Helper()
|
||||
id := mkUser(t, dbq.New(pool), name, admin)
|
||||
_, err := pool.Exec(context.Background(), "UPDATE users SET email = $2 WHERE id = $1", id, name+"@example.com")
|
||||
require.NoError(t, err)
|
||||
return id
|
||||
}
|
||||
|
||||
func digestWith(pool *pgxpool.Pool) (*notifications.Digest, *mailer.FakeSender) {
|
||||
fake := &mailer.FakeSender{}
|
||||
return notifications.NewDigest(pool, fake, nil), fake
|
||||
}
|
||||
|
||||
func TestDigest_BatchGoesOutOnceAWindowAfterTheFirstItem(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
alice := emailUser(t, pool, "alice", false)
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
start := time.Now()
|
||||
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved, notifications.ToUser(alice),
|
||||
notifications.Payload{Name: "Moe Shop – WWW"}.Map()))
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestRejected, notifications.ToUser(alice),
|
||||
notifications.Payload{Name: "Tycho – Awake", Reason: "already have it"}.Map()))
|
||||
|
||||
d.Tick(ctx, start.Add(30*time.Minute))
|
||||
require.Empty(t, fake.Sent, "inside the window nothing goes")
|
||||
|
||||
d.Tick(ctx, start.Add(61*time.Minute))
|
||||
require.Len(t, fake.Sent, 1, "both items, one email")
|
||||
require.Equal(t, "alice@example.com", fake.Sent[0].To)
|
||||
require.Equal(t, "Minstrel: 2 updates", fake.Sent[0].Subject)
|
||||
require.Contains(t, fake.Sent[0].TextBody, "Moe Shop – WWW is on its way.")
|
||||
require.Contains(t, fake.Sent[0].TextBody, "Tycho – Awake was declined: already have it")
|
||||
|
||||
d.Tick(ctx, start.Add(3*time.Hour))
|
||||
require.Len(t, fake.Sent, 1, "never sent twice")
|
||||
}
|
||||
|
||||
func TestDigest_ReadBeforeTheEmailIsNotEmailed(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
q := dbq.New(pool)
|
||||
alice := emailUser(t, pool, "alice", false)
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
start := time.Now()
|
||||
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved, notifications.ToUser(alice), nil))
|
||||
_, err := q.MarkAllNotificationsRead(ctx, dbq.MarkAllNotificationsReadParams{UserID: alice})
|
||||
require.NoError(t, err)
|
||||
|
||||
d.Tick(ctx, start.Add(2*time.Hour))
|
||||
require.Empty(t, fake.Sent, "everything was read: nothing to send")
|
||||
}
|
||||
|
||||
func TestDigest_EmailOffAtNotifyTimeIsNeverSent(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
q := dbq.New(pool)
|
||||
alice := emailUser(t, pool, "alice", false)
|
||||
off := false
|
||||
_, err := notifications.SaveSettings(ctx, q, alice, false, []notifications.SettingChange{
|
||||
{Kind: notifications.KindRequestApproved, Email: &off},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
start := time.Now()
|
||||
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved, notifications.ToUser(alice), nil))
|
||||
on := true
|
||||
_, err = notifications.SaveSettings(ctx, q, alice, false, []notifications.SettingChange{
|
||||
{Kind: notifications.KindRequestApproved, Email: &on},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
d.Tick(ctx, start.Add(2*time.Hour))
|
||||
require.Empty(t, fake.Sent, "turning email on later does not send what came before")
|
||||
}
|
||||
|
||||
func TestDigest_NoAddressNoEmail(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
bob := mkUser(t, dbq.New(pool), "bob", false)
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved, notifications.ToUser(bob), nil))
|
||||
d.Tick(ctx, time.Now().Add(2*time.Hour))
|
||||
require.Empty(t, fake.Sent)
|
||||
}
|
||||
|
||||
func TestDigest_CoalescedItemAppearsOnceWithItsLatestCount(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
q := dbq.New(pool)
|
||||
root := emailUser(t, pool, "root", true)
|
||||
on := true
|
||||
_, err := notifications.SaveSettings(ctx, q, root, true, []notifications.SettingChange{
|
||||
{Kind: notifications.KindTracksMissing, Email: &on},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
start := time.Now()
|
||||
|
||||
for _, c := range []int64{3, 2} {
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindTracksMissing,
|
||||
notifications.ToAdmins(pgtype.UUID{}), notifications.Payload{Count: c}.Map()))
|
||||
}
|
||||
d.Tick(ctx, start.Add(2*time.Hour))
|
||||
require.Len(t, fake.Sent, 1)
|
||||
require.Equal(t, "Minstrel: 5 tracks went missing", fake.Sent[0].Subject)
|
||||
}
|
||||
|
||||
func TestDigest_MailerFailureRetriesWithBackoffAndSendsOnce(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
alice := emailUser(t, pool, "alice", false)
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
start := time.Now()
|
||||
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved, notifications.ToUser(alice), nil))
|
||||
|
||||
fake.FailNext = errors.New("smtp: 421 try later")
|
||||
d.Tick(ctx, start.Add(61*time.Minute))
|
||||
require.Empty(t, fake.Sent)
|
||||
|
||||
d.Tick(ctx, start.Add(63*time.Minute))
|
||||
require.Empty(t, fake.Sent, "waits out the first retry gap")
|
||||
|
||||
d.Tick(ctx, start.Add(67*time.Minute))
|
||||
require.Len(t, fake.Sent, 1, "then goes")
|
||||
|
||||
d.Tick(ctx, start.Add(80*time.Minute))
|
||||
require.Len(t, fake.Sent, 1, "and is not sent again")
|
||||
}
|
||||
|
||||
func TestDigest_NewMusicWaitsForTheSummaryHour(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
alice := emailUser(t, pool, "alice", false) // timezone defaults to UTC
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
|
||||
now := time.Now().UTC()
|
||||
slot := time.Date(now.Year(), now.Month(), now.Day(), 9, 0, 0, 0, time.UTC)
|
||||
if !slot.After(now) {
|
||||
slot = slot.AddDate(0, 0, 1)
|
||||
}
|
||||
|
||||
for _, p := range []notifications.Payload{
|
||||
{Name: "Moe Shop – WWW", Artist: "Moe Shop", Title: "WWW"},
|
||||
{Name: "Moe Shop – Pure", Artist: "Moe Shop", Title: "Pure"},
|
||||
} {
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestCompleted, notifications.ToUser(alice), p.Map()))
|
||||
}
|
||||
|
||||
d.Tick(ctx, slot.Add(-time.Minute))
|
||||
require.Empty(t, fake.Sent, "new music never goes out in an hourly batch")
|
||||
|
||||
d.Tick(ctx, slot)
|
||||
require.Len(t, fake.Sent, 1)
|
||||
require.Equal(t, "Minstrel: new music in your library", fake.Sent[0].Subject)
|
||||
require.Contains(t, fake.Sent[0].TextBody, "Moe Shop\n - WWW")
|
||||
require.Contains(t, fake.Sent[0].TextBody, " - Pure")
|
||||
|
||||
d.Tick(ctx, slot.Add(5*time.Hour))
|
||||
require.Len(t, fake.Sent, 1, "one summary a day")
|
||||
}
|
||||
|
||||
func TestEmailSettings_DefaultsMatchTheMigrationAndRoundTrip(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
q := dbq.New(pool)
|
||||
|
||||
got, err := notifications.LoadEmailSettings(ctx, q)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, notifications.DefaultEmailSettings, got)
|
||||
|
||||
saved, err := notifications.SaveEmailSettings(ctx, q, notifications.EmailSettings{SummaryHour: 7, BatchWindowMinutes: 30})
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, notifications.EmailSettings{SummaryHour: 7, BatchWindowMinutes: 30}, saved)
|
||||
got, err = notifications.LoadEmailSettings(ctx, q)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, saved, got)
|
||||
}
|
||||
|
||||
func TestDigest_UsesTheConfiguredBatchWindow(t *testing.T) {
|
||||
pool := testPool(t)
|
||||
ctx := context.Background()
|
||||
alice := emailUser(t, pool, "alice", false)
|
||||
_, err := notifications.SaveEmailSettings(ctx, dbq.New(pool), notifications.EmailSettings{SummaryHour: 9, BatchWindowMinutes: 15})
|
||||
require.NoError(t, err)
|
||||
n := notifications.New(pool, nil, nil)
|
||||
d, fake := digestWith(pool)
|
||||
start := time.Now()
|
||||
|
||||
require.NoError(t, n.Notify(ctx, notifications.KindRequestApproved, notifications.ToUser(alice), nil))
|
||||
d.Tick(ctx, start.Add(16*time.Minute))
|
||||
require.Len(t, fake.Sent, 1)
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
package notifications
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"embed"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
htmltemplate "html/template"
|
||||
"strings"
|
||||
texttemplate "text/template"
|
||||
)
|
||||
|
||||
//go:embed templates/digest.txt templates/digest.html
|
||||
var digestFS embed.FS
|
||||
|
||||
var (
|
||||
digestText = texttemplate.Must(texttemplate.ParseFS(digestFS, "templates/digest.txt"))
|
||||
digestHTML = htmltemplate.Must(htmltemplate.ParseFS(digestFS, "templates/digest.html"))
|
||||
)
|
||||
|
||||
// digestEmail is one rendered email.
|
||||
type digestEmail struct {
|
||||
Subject string
|
||||
Text string
|
||||
HTML string
|
||||
}
|
||||
|
||||
type digestLine struct {
|
||||
Title string
|
||||
Body string
|
||||
URL string
|
||||
}
|
||||
|
||||
type digestArtist struct {
|
||||
Name string
|
||||
URL string
|
||||
Titles []digestLine
|
||||
}
|
||||
|
||||
type digestVars struct {
|
||||
Name string
|
||||
Intro string
|
||||
Items []digestLine
|
||||
Artists []digestArtist
|
||||
SettingsURL string
|
||||
}
|
||||
|
||||
// renderDigest renders one group's email. base is the public address; with
|
||||
// none set, the email still goes out, without links.
|
||||
func renderDigest(group EmailGroup, name string, items []pendingItem, base string) (digestEmail, error) {
|
||||
vars := digestVars{Name: name, SettingsURL: joinURL(base, "/settings#notifications")}
|
||||
var subject string
|
||||
if group == EmailSummary {
|
||||
vars.Artists = summaryArtists(items, base)
|
||||
vars.Intro = "New in your library since the last summary:"
|
||||
subject = "Minstrel: new music in your library"
|
||||
} else {
|
||||
for _, it := range items {
|
||||
r := Render(it.Kind, it.Payload)
|
||||
vars.Items = append(vars.Items, digestLine{Title: r.Title, Body: r.Body, URL: joinURL(base, r.Link)})
|
||||
}
|
||||
vars.Intro = "Here's what happened on Minstrel:"
|
||||
subject = fmt.Sprintf("Minstrel: %d updates", len(items))
|
||||
if len(items) == 1 {
|
||||
subject = "Minstrel: " + vars.Items[0].Title
|
||||
}
|
||||
}
|
||||
|
||||
var text, html bytes.Buffer
|
||||
if err := digestText.Execute(&text, vars); err != nil {
|
||||
return digestEmail{}, fmt.Errorf("render text: %w", err)
|
||||
}
|
||||
if err := digestHTML.Execute(&html, vars); err != nil {
|
||||
return digestEmail{}, fmt.Errorf("render html: %w", err)
|
||||
}
|
||||
return digestEmail{Subject: subject, Text: text.String(), HTML: html.String()}, nil
|
||||
}
|
||||
|
||||
// summaryArtists groups new-music arrivals by artist, in the order each
|
||||
// artist first arrived, with each artist's albums and tracks under them. An
|
||||
// artist request has no title of its own; the artist links to its page.
|
||||
func summaryArtists(items []pendingItem, base string) []digestArtist {
|
||||
var out []digestArtist
|
||||
index := map[string]int{}
|
||||
for _, it := range items {
|
||||
var p Payload
|
||||
_ = json.Unmarshal(it.Payload, &p)
|
||||
artist, title := p.Artist, p.Title
|
||||
if artist == "" {
|
||||
// Rows from before Artist was recorded carry only "Artist – Title".
|
||||
artist, title, _ = strings.Cut(p.Name, " – ")
|
||||
}
|
||||
if artist == "" {
|
||||
artist = "Something you asked for"
|
||||
}
|
||||
i, ok := index[artist]
|
||||
if !ok {
|
||||
i = len(out)
|
||||
index[artist] = i
|
||||
out = append(out, digestArtist{Name: artist})
|
||||
}
|
||||
link := Render(it.Kind, it.Payload).Link
|
||||
if title == "" {
|
||||
if p.ArtistID != "" {
|
||||
out[i].URL = joinURL(base, link)
|
||||
}
|
||||
continue
|
||||
}
|
||||
out[i].Titles = append(out[i].Titles, digestLine{Title: title, URL: joinURL(base, link)})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// joinURL is base followed by path, or "" when no public address is set.
|
||||
func joinURL(base, path string) string {
|
||||
if base == "" || path == "" {
|
||||
return ""
|
||||
}
|
||||
return strings.TrimRight(base, "/") + path
|
||||
}
|
||||
@@ -0,0 +1,178 @@
|
||||
package notifications
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
var (
|
||||
digestCfg = EmailSettings{SummaryHour: 9, BatchWindowMinutes: 60}
|
||||
t0 = time.Date(2026, 10, 8, 12, 0, 0, 0, time.UTC)
|
||||
)
|
||||
|
||||
func item(k Kind, at time.Time) pendingItem {
|
||||
return pendingItem{ID: pgtype.UUID{Bytes: [16]byte{byte(at.Minute()), byte(len(k))}, Valid: true}, Kind: k, CreatedAt: at}
|
||||
}
|
||||
|
||||
func allEmail() map[Kind]bool {
|
||||
m := map[Kind]bool{}
|
||||
for _, k := range Kinds() {
|
||||
m[k] = true
|
||||
}
|
||||
return m
|
||||
}
|
||||
|
||||
func TestPlanUser_BatchWindow(t *testing.T) {
|
||||
first := item(KindRequestApproved, t0)
|
||||
later := item(KindRequestRejected, t0.Add(40*time.Minute))
|
||||
cases := []struct {
|
||||
name string
|
||||
now time.Time
|
||||
state groupState
|
||||
items []pendingItem
|
||||
wantDue bool
|
||||
wantOpen time.Time
|
||||
}{
|
||||
{"opens at the first item and is not due inside the window", t0.Add(59 * time.Minute), groupState{}, []pendingItem{first, later}, false, t0},
|
||||
{"due a window after the first item, carrying everything since", t0.Add(60 * time.Minute), groupState{}, []pendingItem{first, later}, true, t0},
|
||||
{"a recorded start wins over an item that moved later", t0.Add(61 * time.Minute), groupState{BatchOpenedAt: t0}, []pendingItem{later}, true, t0},
|
||||
{"a failed send waits out its retry", t0.Add(90 * time.Minute), groupState{BatchOpenedAt: t0, RetryAfter: t0.Add(95 * time.Minute)}, []pendingItem{first}, false, t0},
|
||||
{"and goes once the retry has passed", t0.Add(96 * time.Minute), groupState{BatchOpenedAt: t0, RetryAfter: t0.Add(95 * time.Minute)}, []pendingItem{first}, true, t0},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
p := planUser(c.now, digestCfg, time.UTC, c.items, allEmail(), c.state, groupState{})
|
||||
require.Equal(t, c.wantDue, p.BatchDue)
|
||||
require.Equal(t, c.wantOpen, p.BatchOpenedAt)
|
||||
require.Len(t, p.Batch, len(c.items))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanUser_QuietBatchSendsNothing(t *testing.T) {
|
||||
p := planUser(t0.Add(3*time.Hour), digestCfg, time.UTC, nil, allEmail(), groupState{BatchOpenedAt: t0}, groupState{})
|
||||
require.False(t, p.BatchDue)
|
||||
require.True(t, p.BatchOpenedAt.IsZero(), "nothing pending leaves no batch open")
|
||||
}
|
||||
|
||||
func TestPlanUser_SummaryHourInUserTimezone(t *testing.T) {
|
||||
ny, err := time.LoadLocation("America/New_York")
|
||||
require.NoError(t, err)
|
||||
arrived := []pendingItem{item(KindRequestCompleted, time.Date(2026, 7, 1, 3, 0, 0, 0, time.UTC))}
|
||||
// 09:00 in New York in July is 13:00 UTC (EDT, UTC-4).
|
||||
cases := []struct {
|
||||
name string
|
||||
now time.Time
|
||||
lastSent time.Time
|
||||
want bool
|
||||
}{
|
||||
{"before the local hour", time.Date(2026, 7, 1, 12, 59, 0, 0, time.UTC), time.Time{}, false},
|
||||
{"at the local hour", time.Date(2026, 7, 1, 13, 0, 0, 0, time.UTC), time.Time{}, true},
|
||||
{"once a day: already sent after today's hour", time.Date(2026, 7, 1, 18, 0, 0, 0, time.UTC), time.Date(2026, 7, 1, 13, 1, 0, 0, time.UTC), false},
|
||||
{"yesterday's summary does not hold today's", time.Date(2026, 7, 2, 13, 5, 0, 0, time.UTC), time.Date(2026, 7, 1, 13, 1, 0, 0, time.UTC), true},
|
||||
// 03:30 UTC on 2 July is 23:30 on 1 July in New York: still the 1st there.
|
||||
{"the local day, not the UTC day", time.Date(2026, 7, 2, 3, 30, 0, 0, time.UTC), time.Date(2026, 7, 1, 13, 1, 0, 0, time.UTC), false},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
p := planUser(c.now, digestCfg, ny, arrived, allEmail(), groupState{}, groupState{LastSentAt: c.lastSent})
|
||||
require.Equal(t, c.want, p.SummaryDue)
|
||||
require.False(t, p.BatchDue, "new music never goes out in a batch")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestPlanUser_SummaryAcrossDST(t *testing.T) {
|
||||
ny, err := time.LoadLocation("America/New_York")
|
||||
require.NoError(t, err)
|
||||
arrived := []pendingItem{item(KindRequestCompleted, time.Date(2026, 11, 1, 0, 0, 0, 0, time.UTC))}
|
||||
// Clocks go back on 1 November 2026: 09:00 the day before is 13:00 UTC
|
||||
// (EDT), and 09:00 that day is 14:00 UTC (EST).
|
||||
before := groupState{LastSentAt: time.Date(2026, 10, 31, 13, 0, 0, 0, time.UTC)}
|
||||
p := planUser(time.Date(2026, 11, 1, 13, 30, 0, 0, time.UTC), digestCfg, ny, arrived, allEmail(), groupState{}, before)
|
||||
require.False(t, p.SummaryDue, "08:30 local after the change is not yet the hour")
|
||||
p = planUser(time.Date(2026, 11, 1, 14, 0, 0, 0, time.UTC), digestCfg, ny, arrived, allEmail(), groupState{}, before)
|
||||
require.True(t, p.SummaryDue, "09:00 local after the change")
|
||||
|
||||
// Spring forward, 8 March 2026: 02:00 does not exist. A summary hour of
|
||||
// 2 still goes out that day, once.
|
||||
gap := EmailSettings{SummaryHour: 2, BatchWindowMinutes: 60}
|
||||
spring := []pendingItem{item(KindRequestCompleted, time.Date(2026, 3, 8, 0, 0, 0, 0, time.UTC))}
|
||||
at := time.Date(2026, 3, 8, 9, 0, 0, 0, time.UTC) // 05:00 EDT
|
||||
p = planUser(at, gap, ny, spring, allEmail(), groupState{}, groupState{LastSentAt: time.Date(2026, 3, 7, 7, 0, 0, 0, time.UTC)})
|
||||
require.True(t, p.SummaryDue)
|
||||
p = planUser(at.Add(time.Hour), gap, ny, spring, allEmail(), groupState{}, groupState{LastSentAt: at})
|
||||
require.False(t, p.SummaryDue)
|
||||
}
|
||||
|
||||
func TestPlanUser_EmailOffAndStaleItemsAreSkippedNotSent(t *testing.T) {
|
||||
off := allEmail()
|
||||
off[KindRequestRejected] = false
|
||||
items := []pendingItem{
|
||||
item(KindRequestRejected, t0),
|
||||
item(KindRequestApproved, t0.Add(-8*24*time.Hour)),
|
||||
}
|
||||
p := planUser(t0.Add(2*time.Hour), digestCfg, time.UTC, items, off, groupState{}, groupState{})
|
||||
require.Len(t, p.Skip, 2)
|
||||
require.Empty(t, p.Batch)
|
||||
require.False(t, p.BatchDue, "nothing left after filtering: no email")
|
||||
require.False(t, p.SummaryDue)
|
||||
}
|
||||
|
||||
func TestRetryDelay_DoublesToACap(t *testing.T) {
|
||||
require.Equal(t, 5*time.Minute, retryDelay(1))
|
||||
require.Equal(t, 10*time.Minute, retryDelay(2))
|
||||
require.Equal(t, 40*time.Minute, retryDelay(4))
|
||||
require.Equal(t, 6*time.Hour, retryDelay(20))
|
||||
}
|
||||
|
||||
func TestRenderDigest_SummaryGroupsByArtist(t *testing.T) {
|
||||
arrival := func(p Payload) pendingItem {
|
||||
b, _ := json.Marshal(p.Map())
|
||||
return pendingItem{Kind: KindRequestCompleted, Payload: b}
|
||||
}
|
||||
items := []pendingItem{
|
||||
arrival(Payload{Name: "Moe Shop – WWW", Artist: "Moe Shop", Title: "WWW", AlbumID: "al-1"}),
|
||||
arrival(Payload{Name: "Boards of Canada", Artist: "Boards of Canada", ArtistID: "ar-2"}),
|
||||
arrival(Payload{Name: "Moe Shop – Pure", Artist: "Moe Shop", Title: "Pure", AlbumID: "al-3"}),
|
||||
arrival(Payload{Name: "Tycho – Awake"}), // stored before Artist was recorded
|
||||
}
|
||||
e, err := renderDigest(EmailSummary, "alice", items, "https://music.example/")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "Minstrel: new music in your library", e.Subject)
|
||||
|
||||
moe := strings.Index(e.Text, "Moe Shop")
|
||||
require.GreaterOrEqual(t, moe, 0)
|
||||
require.Equal(t, moe, strings.LastIndex(e.Text, "Moe Shop\n"), "one heading per artist")
|
||||
require.Less(t, moe, strings.Index(e.Text, "Boards of Canada"), "artists in the order they arrived")
|
||||
require.Contains(t, e.Text, " - WWW\n https://music.example/albums/al-1")
|
||||
require.Contains(t, e.Text, " - Pure\n https://music.example/albums/al-3")
|
||||
require.Contains(t, e.Text, "Boards of Canada\n https://music.example/artists/ar-2")
|
||||
require.Contains(t, e.Text, "Tycho\n - Awake")
|
||||
require.Contains(t, e.Text, "https://music.example/settings#notifications")
|
||||
require.Contains(t, e.HTML, `<a href="https://music.example/albums/al-1"`)
|
||||
}
|
||||
|
||||
func TestRenderDigest_BatchListsEachItemAndEscapesHTML(t *testing.T) {
|
||||
b, _ := json.Marshal(Payload{Name: "<script>x</script>", Reason: "dup"}.Map())
|
||||
items := []pendingItem{
|
||||
{Kind: KindRequestApproved, Payload: []byte(`{"name":"WWW"}`)},
|
||||
{Kind: KindRequestRejected, Payload: b},
|
||||
}
|
||||
e, err := renderDigest(EmailBatch, "alice", items, "")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "Minstrel: 2 updates", e.Subject)
|
||||
require.Contains(t, e.Text, "- Request approved\n WWW is on its way.")
|
||||
require.NotContains(t, e.Text, "http", "no public address: no links")
|
||||
require.Contains(t, e.Text, "Settings → Notifications")
|
||||
require.NotContains(t, e.HTML, "<script>")
|
||||
|
||||
one, err := renderDigest(EmailBatch, "alice", items[:1], "")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "Minstrel: Request approved", one.Subject)
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package notifications
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
|
||||
)
|
||||
|
||||
// EmailSettings are the digest's two knobs, in admin Settings (rule 25).
|
||||
type EmailSettings struct {
|
||||
// SummaryHour is the local hour (0-23, in each user's own timezone) the
|
||||
// daily new-music summary goes out.
|
||||
SummaryHour int32
|
||||
// BatchWindowMinutes is how long a batch stays open after its first item.
|
||||
BatchWindowMinutes int32
|
||||
}
|
||||
|
||||
const (
|
||||
minBatchWindowMinutes = 15
|
||||
maxBatchWindowMinutes = 1440
|
||||
)
|
||||
|
||||
// DefaultEmailSettings mirrors migration 0074's column defaults.
|
||||
var DefaultEmailSettings = EmailSettings{SummaryHour: 9, BatchWindowMinutes: 60}
|
||||
|
||||
// ErrEmailSettingOutOfRange is a value migration 0074's CHECKs would reject,
|
||||
// so the API answers 400 naming the field rather than a constraint violation.
|
||||
var ErrEmailSettingOutOfRange = errors.New("notification email setting out of range")
|
||||
|
||||
// BatchWindow is BatchWindowMinutes as a duration.
|
||||
func (s EmailSettings) BatchWindow() time.Duration {
|
||||
return time.Duration(s.BatchWindowMinutes) * time.Minute
|
||||
}
|
||||
|
||||
// LoadEmailSettings reads the settings. A failed read returns the defaults
|
||||
// with the error, so the digest keeps its shipped cadence rather than stopping.
|
||||
func LoadEmailSettings(ctx context.Context, q *dbq.Queries) (EmailSettings, error) {
|
||||
row, err := q.GetNotificationEmailSettings(ctx)
|
||||
if err != nil {
|
||||
return DefaultEmailSettings, fmt.Errorf("notification email settings: load: %w", err)
|
||||
}
|
||||
return EmailSettings{SummaryHour: row.SummaryHour, BatchWindowMinutes: row.BatchWindowMinutes}, nil
|
||||
}
|
||||
|
||||
// SaveEmailSettings validates and writes the settings, returning them as saved.
|
||||
// The digest reads them on its next tick; nothing is cached.
|
||||
func SaveEmailSettings(ctx context.Context, q *dbq.Queries, in EmailSettings) (EmailSettings, error) {
|
||||
switch {
|
||||
case in.SummaryHour < 0 || in.SummaryHour > 23:
|
||||
return EmailSettings{}, fmt.Errorf("%w: summary_hour must be 0-23", ErrEmailSettingOutOfRange)
|
||||
case in.BatchWindowMinutes < minBatchWindowMinutes || in.BatchWindowMinutes > maxBatchWindowMinutes:
|
||||
return EmailSettings{}, fmt.Errorf("%w: batch_window_minutes must be %d-%d",
|
||||
ErrEmailSettingOutOfRange, minBatchWindowMinutes, maxBatchWindowMinutes)
|
||||
}
|
||||
row, err := q.UpdateNotificationEmailSettings(ctx, dbq.UpdateNotificationEmailSettingsParams{
|
||||
SummaryHour: in.SummaryHour,
|
||||
BatchWindowMinutes: in.BatchWindowMinutes,
|
||||
})
|
||||
if err != nil {
|
||||
return EmailSettings{}, fmt.Errorf("notification email settings: save: %w", err)
|
||||
}
|
||||
return EmailSettings{SummaryHour: row.SummaryHour, BatchWindowMinutes: row.BatchWindowMinutes}, nil
|
||||
}
|
||||
@@ -92,7 +92,7 @@ func (n *Notifier) Notify(ctx context.Context, kind Kind, to Recipients, payload
|
||||
if !channels[id].Inbox {
|
||||
continue
|
||||
}
|
||||
if err := write(ctx, q, kind, id, body); err != nil {
|
||||
if err := write(ctx, q, kind, id, body, channels[id].Email); err != nil {
|
||||
return fmt.Errorf("notifications: write %s: %w", kind, err)
|
||||
}
|
||||
delivered = append(delivered, id)
|
||||
@@ -166,11 +166,14 @@ func channelsFor(ctx context.Context, q *dbq.Queries, kind Kind, ids []pgtype.UU
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func write(ctx context.Context, q *dbq.Queries, kind Kind, userID pgtype.UUID, body []byte) error {
|
||||
// write stores one row. emailWanted is the recipient's email channel for the
|
||||
// kind: a row nobody wants emailed is stamped emailed_at at once, so the
|
||||
// digest never sends it, even if email is switched on later.
|
||||
func write(ctx context.Context, q *dbq.Queries, kind Kind, userID pgtype.UUID, body []byte, emailWanted bool) error {
|
||||
key := kind.coalesceKey()
|
||||
if key == "" {
|
||||
_, err := q.InsertNotification(ctx, dbq.InsertNotificationParams{
|
||||
UserID: userID, Kind: string(kind), Payload: body,
|
||||
UserID: userID, Kind: string(kind), Payload: body, EmailWanted: emailWanted,
|
||||
})
|
||||
return err
|
||||
}
|
||||
@@ -179,6 +182,7 @@ func write(ctx context.Context, q *dbq.Queries, kind Kind, userID pgtype.UUID, b
|
||||
Kind: string(kind),
|
||||
Payload: body,
|
||||
CoalesceKey: &key,
|
||||
EmailWanted: emailWanted,
|
||||
SumCount: specs[kind].sumCount,
|
||||
})
|
||||
return err
|
||||
|
||||
@@ -24,7 +24,12 @@ type Payload struct {
|
||||
RequestKind string `json:"request_kind,omitempty"` // artist | album | track
|
||||
// Name is what was requested or flagged, as the user would say it:
|
||||
// "WWW", "Moe Shop", "Moe Shop – WWW".
|
||||
Name string `json:"name,omitempty"`
|
||||
Name string `json:"name,omitempty"`
|
||||
// Artist and Title are Name's parts, for request_completed: the new-music
|
||||
// summary email groups arrivals by artist. Title is empty for an artist
|
||||
// request.
|
||||
Artist string `json:"artist,omitempty"`
|
||||
Title string `json:"title,omitempty"`
|
||||
ArtistID string `json:"artist_id,omitempty"`
|
||||
AlbumID string `json:"album_id,omitempty"`
|
||||
// Actor is the other person involved: who asked, who flagged.
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
<!doctype html>
|
||||
<html>
|
||||
<body style="font-family: -apple-system, system-ui, sans-serif; max-width: 560px; margin: 0 auto; padding: 24px; color: #2c2c2c;">
|
||||
<p>Hello <strong>{{.Name}}</strong>,</p>
|
||||
|
||||
<p>{{.Intro}}</p>
|
||||
{{range .Artists}}
|
||||
<p style="margin: 20px 0 4px; font-weight: 600;">{{if .URL}}<a href="{{.URL}}" style="color: #2c2c2c;">{{.Name}}</a>{{else}}{{.Name}}{{end}}</p>
|
||||
{{if .Titles}}<ul style="margin: 0; padding-left: 20px;">
|
||||
{{range .Titles}}<li style="margin: 2px 0;">{{if .URL}}<a href="{{.URL}}" style="color: #4a6b5c;">{{.Title}}</a>{{else}}{{.Title}}{{end}}</li>
|
||||
{{end}}</ul>{{end}}
|
||||
{{end}}{{if .Items}}
|
||||
<ul style="margin: 16px 0; padding-left: 20px;">
|
||||
{{range .Items}}<li style="margin: 8px 0;">
|
||||
{{if .URL}}<a href="{{.URL}}" style="color: #2c2c2c; font-weight: 600;">{{.Title}}</a>{{else}}<strong>{{.Title}}</strong>{{end}}
|
||||
{{if .Body}}<br><span style="color: #555;">{{.Body}}</span>{{end}}
|
||||
</li>
|
||||
{{end}}</ul>
|
||||
{{end}}
|
||||
<p style="font-size: 12px; color: #999; margin-top: 32px; border-top: 1px solid #eee; padding-top: 16px;">
|
||||
{{if .SettingsURL}}<a href="{{.SettingsURL}}" style="color: #999;">Choose which notifications are emailed to you</a>{{else}}Choose which notifications are emailed to you in Minstrel, under Settings → Notifications.{{end}}<br>
|
||||
— Minstrel
|
||||
</p>
|
||||
</body>
|
||||
</html>
|
||||
@@ -0,0 +1,16 @@
|
||||
Hello {{.Name}},
|
||||
|
||||
{{.Intro}}
|
||||
{{range .Artists}}
|
||||
{{.Name}}{{if .URL}}
|
||||
{{.URL}}{{end}}{{range .Titles}}
|
||||
- {{.Title}}{{if .URL}}
|
||||
{{.URL}}{{end}}{{end}}
|
||||
{{end}}{{range .Items}}
|
||||
- {{.Title}}{{if .Body}}
|
||||
{{.Body}}{{end}}{{if .URL}}
|
||||
{{.URL}}{{end}}
|
||||
{{end}}
|
||||
{{if .SettingsURL}}Choose which notifications are emailed to you: {{.SettingsURL}}{{else}}Choose which notifications are emailed to you in Minstrel, under Settings → Notifications.{{end}}
|
||||
|
||||
— Minstrel
|
||||
Reference in New Issue
Block a user