Files
minstrel/internal/notifications/digest.go
T
bvandeusenandClaude Opus 5.5 63709a433d
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
feat(notifications): grouped email digest, new music as a daily summary (#5346)
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>
2026-10-08 07:48:22 -04:00

391 lines
12 KiB
Go

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
}