diff --git a/cmd/silo/main.go b/cmd/silo/main.go index 3169f7b7..cafcc1b9 100644 --- a/cmd/silo/main.go +++ b/cmd/silo/main.go @@ -56,6 +56,7 @@ import ( "github.com/Silo-Server/silo-server/internal/libraryingest" "github.com/Silo-Server/silo-server/internal/logfilter" "github.com/Silo-Server/silo-server/internal/logstream" + "github.com/Silo-Server/silo-server/internal/mail" "github.com/Silo-Server/silo-server/internal/markers" "github.com/Silo-Server/silo-server/internal/mdblist" "github.com/Silo-Server/silo-server/internal/metadata" @@ -1305,6 +1306,7 @@ func main() { deps.EventsHub, deps.RedisClient, deps.SecretCipher, + mail.NewSMTPSender(settingsRepo), ) userStoreProvider = notifications.WrapUserStoreProvider(userStoreProvider, notificationSystem) deps.Notifications = notificationSystem diff --git a/docs/superpowers/plans/notifications/07-email-channel.md b/docs/superpowers/plans/notifications/07-email-channel.md new file mode 100644 index 00000000..42372399 --- /dev/null +++ b/docs/superpowers/plans/notifications/07-email-channel.md @@ -0,0 +1,109 @@ +# Notifications: Email Channel + +**Date:** 2026-06-11 +**Status:** Implemented (written post-implementation) +**Scope:** Item 3 of [`06-v1.5-roadmap.md`](./06-v1.5-roadmap.md) — the first real consumer of the shared SMTP core (`internal/mail`, `docs/architecture/email.md`). +**Depends On:** [`00-architecture-overview.md`](./00-architecture-overview.md), [`01-release-events-and-inbox.md`](./01-release-events-and-inbox.md) + +## Decisions (resolving the roadmap's open questions) + +- **Account-level, as planned.** Email addresses live on `users`; one mode + covers every profile on the account and one email aggregates across them. +- **Per-episode AND digest, not digest-only.** The roadmap recommended + digest-only; product direction chose to offer per-episode alerts too, gated + by an admin allowance (`notifications.email.allow_per_episode`). When the + admin disallows it, accounts set to per-episode are **coerced to the daily + digest** rather than silenced. +- **Interest-scoped by construction.** Email consumes existing + `notification_deliveries` rows, which the fanout only creates for profiles + with series interest (favorites, watchlist, continue-watching, next-up) and + for direct notices (`request.fulfilled`, `webhook.auto_disabled`). The + channel adds no targeting of its own — it is never "all new content". +- **Opt-in, default off.** Users enable it per account in Settings → + Notifications; enabling initializes the watermark to now so history never + floods a fresh opt-in. + +## Architecture: watermark sweep, not a third outbox + +Webhooks and web push use per-target outbox attempt rows. Email deliberately +does not: + +- Deliveries already carry `user_id`, and an account whose profiles follow the + same series gets one row per profile — a per-row outbox would email the same + episode several times. The sweep collapses them (dedupe by `episode_id`, by + `request_id` for requests). +- A per-account watermark over `(created_at, id)` that advances **only after a + successful SMTP send** gives durability for free: a crash or SMTP outage + re-sends on the next pass instead of dropping. +- Both cadences are the same mechanism: per-episode sweeps every minute (and + is nudged by the dispatcher seconds after fanout commits); the digest is the + same sweep gated on "today's send hour passed and not yet stamped today". + +State lives in `notification_email_prefs` (`migrations/sql/`, +`email_notification_channel`): mode, watermark, `last_digest_at`, and failure +backoff counters (`last_attempt_at`, `consecutive_failures`; 1m doubling, +capped at 6h). No FK to `users` per the notification-tables rule; deleted or +disabled accounts drop out of the recipient join. A supporting index +`notification_deliveries_user_created_idx (user_id, created_at, id)` serves +the sweep. + +**Multi-node safety:** each account is processed inside one transaction that +claims the prefs row `FOR UPDATE SKIP LOCKED`, re-derives eligibility from the +locked row (mode flips and another node's digest stamp are both re-checked), +sends, then commits the watermark/stamp. Failed sends commit only the backoff +counters. `mail.ErrNotConfigured` aborts the whole pass; three consecutive +send failures end it early (SMTP trouble is global, not per-recipient). + +**Flood bounds:** one email renders at most 30 lines (`…and N more in your +Silo inbox`), one pass fetches at most 200 rows per account, and upstream the +per-series burst cap already limits fanout volume. Digest emails include only +rows still unread at compose time; the watermark passes read rows silently. + +## Files + +| Piece | Location | +|---|---| +| Modes, prefs repo (`notification_email_prefs`) | `internal/notifications/email_prefs_repo.go` | +| Worker, dispatcher nudge, System service methods | `internal/notifications/email_digest.go` | +| Subject/text/HTML rendering | `internal/notifications/email_compose.go` | +| Account sweep query | `DeliveryRepository.ListForUserSince` (`internal/notifications/delivery_repo.go`) | +| Settings accessors | `internal/notifications/settings.go` | +| API handlers | `internal/api/handlers/notifications_email.go` (+ capability in `notifications.go`) | +| Web UI (user) | `EmailSection` in `web/src/pages/settings/NotificationsSettings.tsx` | +| Web UI (admin) | Email group in `web/src/pages/admin-settings/NotificationsAdminSettings.tsx` | +| Logic tests | `internal/notifications/email_logic_test.go` | + +The worker is wired in `notifications.NewSystem` (new `mail.Sender` parameter, +passed from `cmd/silo/main.go`); its dispatcher joins the `MultiDispatcher`, +so operational deliveries (`request.fulfilled`, `webhook.auto_disabled`) nudge +it exactly like fanout rows do. + +## Settings + +| Key | Default | Meaning | +|---|---|---| +| `notifications.email_enabled` | `true` | Channel kill switch (availability still requires SMTP configured via `email.*`) | +| `notifications.email.allow_per_episode` | `true` | Admin allowance for the per-episode cadence | +| `notifications.email.digest_hour` | `8` | Hour (0–23, server-local) daily digests go out | +| `notifications.email.external_url` | empty | Public base URL for deep links in emails; empty sends link-free emails (the server origin is never leaked implicitly) | + +## API + +- `GET /api/v1/notifications/email-preferences` → `{"mode": "off" | "per_episode" | "daily_digest"}` +- `PUT /api/v1/notifications/email-preferences` `{"mode": ...}` — 400 codes: + `bad_request` (unknown mode), `not_allowed` (per-episode disallowed), + `no_email` (account has no address). Any profile on the account may set it. +- `GET /api/v1/notifications/capability` gained + `"email": {"available", "modes", "digest_hour"}`; clients gate setup UI on + it as usual. `available` requires the kill switch on **and** + `mail.Sender.Enabled()` — never read `email.*` settings directly. + +## Deliberately not in v1 + +- Posters/images in emails (would require externally reachable presigned URLs). +- `List-Unsubscribe` headers / tokenized unsubscribe endpoint (self-hosted, + opt-in; revisit if servers grow beyond household scale). +- Per-user digest hour (admin-global for now). +- DB-backed integration tests for the sweep (same Postgres-harness gap as the + rest of `01`'s verification backlog; pure logic is covered by + `email_logic_test.go`). diff --git a/internal/api/handlers/notifications.go b/internal/api/handlers/notifications.go index 545ba7c9..6c9279e3 100644 --- a/internal/api/handlers/notifications.go +++ b/internal/api/handlers/notifications.go @@ -282,6 +282,15 @@ type capabilityResponse struct { AndroidPush capabilityPush `json:"android_push"` WebPush capabilityWebPush `json:"web_push"` Webhooks capabilityWebhooks `json:"webhooks"` + Email capabilityEmail `json:"email"` +} + +type capabilityEmail struct { + Available bool `json:"available"` + // Modes lists the cadences users may pick (per-episode is an admin + // allowance); DigestHour tells the UI when daily digests go out. + Modes []string `json:"modes"` + DigestHour int `json:"digest_hour"` } type capabilityWebPush struct { @@ -324,12 +333,25 @@ func (h *NotificationsHandler) HandleCapability(w http.ResponseWriter, r *http.R webPush = capabilityWebPush{Available: true, PublicKey: publicKey} } } + email := capabilityEmail{Modes: []string{}} + if h.system.EmailAvailable(r.Context()) { + modes := []string{notifications.EmailModeDailyDigest} + if h.system.Settings.EmailAllowPerEpisode(r.Context()) { + modes = append(modes, notifications.EmailModePerEpisode) + } + email = capabilityEmail{ + Available: true, + Modes: modes, + DigestHour: h.system.Settings.EmailDigestHour(r.Context()), + } + } writeJSON(w, http.StatusOK, capabilityResponse{ InApp: capabilityInApp{Enabled: h.system.Settings.UIEnabled(r.Context())}, ApplePush: capabilityPush{Available: false, Provider: "off", SupportedModes: []string{"in_app_only"}}, AndroidPush: capabilityPush{Available: false, Provider: "off", SupportedModes: []string{"in_app_only"}}, WebPush: webPush, Webhooks: webhooks, + Email: email, }) } diff --git a/internal/api/handlers/notifications_email.go b/internal/api/handlers/notifications_email.go new file mode 100644 index 00000000..64dcb9e3 --- /dev/null +++ b/internal/api/handlers/notifications_email.go @@ -0,0 +1,62 @@ +package handlers + +import ( + "encoding/json" + "errors" + "net/http" + + apimw "github.com/Silo-Server/silo-server/internal/api/middleware" + "github.com/Silo-Server/silo-server/internal/notifications" +) + +// emailPreferencesResponse is the account-level email notification setting. +// Unlike the per-profile preferences, one mode covers every profile on the +// login account: email addresses live on accounts, and the emails themselves +// aggregate across profiles. +type emailPreferencesResponse struct { + Mode string `json:"mode"` +} + +type updateEmailPreferencesRequest struct { + Mode string `json:"mode"` +} + +// HandleGetEmailPreferences handles GET /notifications/email-preferences. +func (h *NotificationsHandler) HandleGetEmailPreferences(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + mode, err := h.system.EmailMode(r.Context(), userID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to load email preferences") + return + } + writeJSON(w, http.StatusOK, emailPreferencesResponse{Mode: mode}) +} + +// HandleUpdateEmailPreferences handles PUT /notifications/email-preferences. +func (h *NotificationsHandler) HandleUpdateEmailPreferences(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + + var req updateEmailPreferencesRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + + err := h.system.SetEmailMode(r.Context(), userID, req.Mode) + switch { + case err == nil: + case errors.Is(err, notifications.ErrEmailModeInvalid): + writeError(w, http.StatusBadRequest, "bad_request", "Unknown email notification mode") + return + case errors.Is(err, notifications.ErrEmailModeNotAllowed): + writeError(w, http.StatusBadRequest, "not_allowed", "Per-episode email is disabled by the administrator") + return + case errors.Is(err, notifications.ErrEmailNoAddress): + writeError(w, http.StatusBadRequest, "no_email", "Your account has no email address") + return + default: + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to save email preferences") + return + } + writeJSON(w, http.StatusOK, emailPreferencesResponse{Mode: req.Mode}) +} diff --git a/internal/api/router.go b/internal/api/router.go index b6ccff10..d263213c 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -1544,6 +1544,8 @@ func NewRouter(deps Dependencies) chi.Router { r.Get("/capability", notificationsHandler.HandleCapability) r.Get("/preferences", notificationsHandler.HandleGetPreferences) r.Put("/preferences", notificationsHandler.HandleUpdatePreferences) + r.Get("/email-preferences", notificationsHandler.HandleGetEmailPreferences) + r.Put("/email-preferences", notificationsHandler.HandleUpdateEmailPreferences) r.Post("/read-all", notificationsHandler.HandleReadAll) r.Route("/webhooks", func(r chi.Router) { r.Get("/", notificationsHandler.HandleListWebhooks) diff --git a/internal/notifications/delivery_repo.go b/internal/notifications/delivery_repo.go index e114a688..cc545d9e 100644 --- a/internal/notifications/delivery_repo.go +++ b/internal/notifications/delivery_repo.go @@ -228,6 +228,22 @@ func (r *DeliveryRepository) GetRowByID(ctx context.Context, id string) (*Delive return &out[0], nil } +// ListForUserSince returns the account's deliveries newer than the watermark, +// ascending, across all of its profiles. Runs inside the email worker's claim +// transaction so the rows read are the rows the advanced watermark covers. +func (r *DeliveryRepository) ListForUserSince(ctx context.Context, tx pgx.Tx, userID int, since Cursor, limit int) ([]DeliveryRow, error) { + rows, err := tx.Query(ctx, + deliveryRowSelect+` + WHERE d.user_id = $1 AND (d.created_at, d.id) > ($2, $3) + ORDER BY d.created_at ASC, d.id ASC + LIMIT $4`, + userID, since.CreatedAt, since.ID, limit) + if err != nil { + return nil, fmt.Errorf("list user deliveries since: %w", err) + } + return scanDeliveryRows(rows) +} + // RecentUnread returns the newest unread rows for the websocket snapshot. func (r *DeliveryRepository) RecentUnread(ctx context.Context, profileID string, limit int) ([]DeliveryRow, error) { rows, err := r.pool.Query(ctx, diff --git a/internal/notifications/email_compose.go b/internal/notifications/email_compose.go new file mode 100644 index 00000000..d3cd18e3 --- /dev/null +++ b/internal/notifications/email_compose.go @@ -0,0 +1,304 @@ +package notifications + +import ( + "fmt" + "html" + "sort" + "strings" +) + +// emailMaxItemsRendered caps how many lines one email renders; the remainder +// collapses into a "+N more" line (the inbox always has everything). +const emailMaxItemsRendered = 30 + +// untitledSeriesGroup is the group heading for episode rows whose series +// metadata is missing or was deleted. +const untitledSeriesGroup = "New episodes" + +// emailContent is one rendered notification email. +type emailContent struct { + Subject string + Text string + HTML string +} + +// emailSeriesGroup is one series' new episodes, in first-appearance order. +type emailSeriesGroup struct { + seriesID string + title string + episodes []DeliveryRow +} + +// emailItems is the collated, deduplicated content of one email. Accounts +// with several profiles following the same series have one delivery row per +// profile; an email reports the episode once. +type emailItems struct { + series []emailSeriesGroup + episodes int + requests []DeliveryRow + others []DeliveryRow +} + +// collateEmailItems groups and dedupes delivery rows for rendering. +func collateEmailItems(rows []DeliveryRow) emailItems { + var items emailItems + seenEpisodes := make(map[string]struct{}, len(rows)) + seenRequests := make(map[string]struct{}, 4) + groupIndex := make(map[string]int, 4) + + for _, row := range rows { + switch row.Type { + case DeliveryTypeEpisodeAvailable: + key := row.ID + if row.EpisodeID != nil && *row.EpisodeID != "" { + key = *row.EpisodeID + } + if _, ok := seenEpisodes[key]; ok { + continue + } + seenEpisodes[key] = struct{}{} + seriesID := "" + if row.SeriesID != nil { + seriesID = *row.SeriesID + } + idx, ok := groupIndex[seriesID] + if !ok { + idx = len(items.series) + groupIndex[seriesID] = idx + title := row.SeriesTitle + if title == "" { + title = untitledSeriesGroup + } + items.series = append(items.series, emailSeriesGroup{seriesID: seriesID, title: title}) + } + items.series[idx].episodes = append(items.series[idx].episodes, row) + items.episodes++ + case DeliveryTypeRequestFulfilled: + key := parseRequestFulfilledFlags(row.ReasonFlags).RequestID + if key == "" { + key = row.ID + } + if _, ok := seenRequests[key]; ok { + continue + } + seenRequests[key] = struct{}{} + items.requests = append(items.requests, row) + default: + items.others = append(items.others, row) + } + } + + for i := range items.series { + sort.SliceStable(items.series[i].episodes, func(a, b int) bool { + ea, eb := items.series[i].episodes[a], items.series[i].episodes[b] + if ea.SeasonNumber == nil || eb.SeasonNumber == nil || + ea.EpisodeNumber == nil || eb.EpisodeNumber == nil { + return false + } + if *ea.SeasonNumber != *eb.SeasonNumber { + return *ea.SeasonNumber < *eb.SeasonNumber + } + return *ea.EpisodeNumber < *eb.EpisodeNumber + }) + } + return items +} + +// episodeCode renders "S02E03"; empty when numbering is unknown. +func episodeCode(row DeliveryRow) string { + if row.SeasonNumber == nil || row.EpisodeNumber == nil { + return "" + } + return fmt.Sprintf("S%02dE%02d", *row.SeasonNumber, *row.EpisodeNumber) +} + +// episodeLine renders one episode entry: "S02E03 — Title", degrading to +// whichever part exists. +func episodeLine(row DeliveryRow) string { + code := episodeCode(row) + switch { + case code != "" && row.EpisodeTitle != "": + return code + " — " + row.EpisodeTitle + case code != "": + return code + case row.EpisodeTitle != "": + return row.EpisodeTitle + default: + return genericEpisodeTitle + } +} + +// requestLine renders one fulfilled-request entry. +func requestLine(row DeliveryRow) string { + if row.SeriesTitle != "" { + return row.SeriesTitle + " is now available" + } + return "Your media request is now available" +} + +// otherLine renders operational and unknown delivery types generically. +func otherLine(row DeliveryRow) string { + if row.Type == DeliveryTypeWebhookAutoDisabled { + return "A webhook stopped working — open notification settings to fix it" + } + return genericNotificationTitle +} + +// countPart pluralizes "3 new episodes" style subject fragments. +func countPart(count int, singular, plural string) string { + if count == 1 { + return "1 " + singular + } + return fmt.Sprintf("%d %s", count, plural) +} + +// emailSubject builds the subject line from the collated items. +func emailSubject(mode string, items emailItems) string { + parts := make([]string, 0, 3) + if items.episodes > 0 { + parts = append(parts, countPart(items.episodes, "new episode", "new episodes")) + } + if len(items.requests) > 0 { + parts = append(parts, countPart(len(items.requests), "request ready", "requests ready")) + } + if len(items.others) > 0 { + parts = append(parts, countPart(len(items.others), "update", "updates")) + } + summary := strings.Join(parts, ", ") + + if mode == EmailModeDailyDigest { + return "Silo daily digest: " + summary + } + if items.episodes == 1 && len(items.requests) == 0 && len(items.others) == 0 { + row := items.series[0].episodes[0] + subject := genericEpisodeTitle + if items.series[0].title != untitledSeriesGroup { + subject += " of " + items.series[0].title + } + if line := episodeLine(row); line != genericEpisodeTitle { + subject += ": " + line + } + return subject + } + if items.episodes == 0 && len(items.requests) == 1 && len(items.others) == 0 { + return requestLine(items.requests[0]) + } + return "Silo: " + summary +} + +// itemURL builds a deep link; empty when no external URL is configured. +func itemURL(baseURL, itemID string) string { + if baseURL == "" || itemID == "" { + return "" + } + return baseURL + "/item/" + itemID +} + +// composeNotificationEmail renders one email (text + HTML) for the given +// delivery rows. baseURL is the admin-configured external URL; empty renders +// without links. +func composeNotificationEmail(mode string, rows []DeliveryRow, baseURL string) emailContent { + items := collateEmailItems(rows) + + var text strings.Builder + var body strings.Builder + rendered := 0 + total := items.episodes + len(items.requests) + len(items.others) + + writeLine := func(plain, href string) { + rendered++ + if rendered > emailMaxItemsRendered { + return + } + text.WriteString(" " + plain + "\n") + if href != "" { + body.WriteString(fmt.Sprintf( + `
  • %s
  • `, + html.EscapeString(href), html.EscapeString(plain))) + } else { + body.WriteString(fmt.Sprintf(`
  • %s
  • `, html.EscapeString(plain))) + } + } + writeHeading := func(title, href string) { + text.WriteString(title + "\n") + label := html.EscapeString(title) + if href != "" { + label = fmt.Sprintf(`%s`, + html.EscapeString(href), label) + } + body.WriteString(fmt.Sprintf( + `

    %s

    `, label)) + } + openList := func() { body.WriteString(``) } + + for _, group := range items.series { + if rendered >= emailMaxItemsRendered { + break + } + writeHeading(group.title, itemURL(baseURL, group.seriesID)) + openList() + for _, row := range group.episodes { + episodeID := "" + if row.EpisodeID != nil { + episodeID = *row.EpisodeID + } + writeLine(episodeLine(row), itemURL(baseURL, episodeID)) + } + closeList() + } + if len(items.requests) > 0 && rendered < emailMaxItemsRendered { + writeHeading("Requests ready", "") + openList() + for _, row := range items.requests { + seriesID := "" + if row.SeriesID != nil { + seriesID = *row.SeriesID + } + writeLine(requestLine(row), itemURL(baseURL, seriesID)) + } + closeList() + } + if len(items.others) > 0 && rendered < emailMaxItemsRendered { + writeHeading("Other updates", "") + openList() + for _, row := range items.others { + writeLine(otherLine(row), "") + } + closeList() + } + if remainder := total - emailMaxItemsRendered; remainder > 0 { + more := fmt.Sprintf("…and %d more in your Silo inbox.", remainder) + text.WriteString(more + "\n") + body.WriteString(fmt.Sprintf( + `

    %s

    `, html.EscapeString(more))) + } + + intro := "New in your library:" + if mode == EmailModeDailyDigest { + intro = "Here's what's new since your last digest:" + } + footer := "You're receiving this because email notifications are enabled for your Silo account. " + + "Manage them in Settings → Notifications." + footerHTML := html.EscapeString(footer) + if baseURL != "" { + settingsURL := html.EscapeString(baseURL + "/settings/notifications") + footerHTML = strings.Replace(footerHTML, + "Settings → Notifications", + fmt.Sprintf(`Settings → Notifications`, settingsURL), 1) + } + + htmlBody := fmt.Sprintf(`
    +

    %s

    +%s +
    +

    %s

    +
    `, + html.EscapeString(intro), body.String(), footerHTML) + + return emailContent{ + Subject: emailSubject(mode, items), + Text: intro + "\n\n" + text.String() + "\n" + footer + "\n", + HTML: htmlBody, + } +} diff --git a/internal/notifications/email_digest.go b/internal/notifications/email_digest.go new file mode 100644 index 00000000..110cd817 --- /dev/null +++ b/internal/notifications/email_digest.go @@ -0,0 +1,360 @@ +package notifications + +import ( + "context" + "errors" + "fmt" + "log/slog" + "time" + + "github.com/Silo-Server/silo-server/internal/mail" + "github.com/jackc/pgx/v5/pgxpool" +) + +const ( + emailPollInterval = time.Minute + // emailNudgeDelay coalesces the per-row dispatch nudges of one fanout + // batch (all rows commit before the first nudge fires) into one pass. + emailNudgeDelay = 2 * time.Second + // emailFetchLimit bounds one email's worth of watermark progress; the + // next pass drains the remainder. + emailFetchLimit = 200 + // emailMaxFailuresPerPass stops a pass early when sends keep failing — + // SMTP trouble is almost always global, not per-recipient. + emailMaxFailuresPerPass = 3 + + emailFailureBackoffBase = time.Minute + emailFailureBackoffMax = 6 * time.Hour +) + +// effectiveEmailMode coerces per-episode to the daily digest when the admin +// has disallowed per-episode email, instead of silencing those accounts. +func effectiveEmailMode(mode string, allowPerEpisode bool) string { + if mode == EmailModePerEpisode && !allowPerEpisode { + return EmailModeDailyDigest + } + return mode +} + +// emailDigestDue reports whether a daily digest should go out: today's send +// time (digestHour, local) has passed and no digest was stamped since. +func emailDigestDue(now time.Time, digestHour int, lastDigestAt *time.Time) bool { + todaySend := time.Date(now.Year(), now.Month(), now.Day(), digestHour, 0, 0, 0, now.Location()) + if now.Before(todaySend) { + return false + } + return lastDigestAt == nil || lastDigestAt.Before(todaySend) +} + +// emailRetryEligible applies exponential backoff after failed sends: +// 1m, 2m, 4m, ... capped at emailFailureBackoffMax. +func emailRetryEligible(now time.Time, lastAttemptAt *time.Time, consecutiveFailures int) bool { + if consecutiveFailures <= 0 || lastAttemptAt == nil { + return true + } + backoff := emailFailureBackoffBase << min(consecutiveFailures-1, 30) + if backoff > emailFailureBackoffMax || backoff <= 0 { + backoff = emailFailureBackoffMax + } + return !now.Before(lastAttemptAt.Add(backoff)) +} + +// EmailWorker delivers notification emails. Unlike webhooks and web push it +// keeps no per-target outbox: deliveries already carry user_id, so a per-user +// watermark over notification_deliveries is the durable dispatch state. The +// watermark advances only after a successful SMTP send, and one email covers +// everything since the last one — which also collapses the duplicate rows an +// account gets when several of its profiles follow the same series. +type EmailWorker struct { + pool *pgxpool.Pool + deliveries *DeliveryRepository + prefs *EmailPrefsRepository + settings *Settings + sender mail.Sender + logger *slog.Logger + nudge chan struct{} + now func() time.Time +} + +func newEmailWorker( + pool *pgxpool.Pool, + deliveries *DeliveryRepository, + prefs *EmailPrefsRepository, + settings *Settings, + sender mail.Sender, +) *EmailWorker { + return &EmailWorker{ + pool: pool, + deliveries: deliveries, + prefs: prefs, + settings: settings, + sender: sender, + logger: slog.Default().With("component", "notifications.email"), + nudge: make(chan struct{}, 1), + now: time.Now, + } +} + +// Nudge schedules a near-term pass so per-episode emails follow fanout within +// seconds instead of waiting for the next poll. Non-blocking. +func (w *EmailWorker) Nudge() { + if w == nil { + return + } + select { + case w.nudge <- struct{}{}: + default: + } +} + +// Run sweeps eligible accounts until ctx is canceled. +func (w *EmailWorker) Run(ctx context.Context) { + ticker := time.NewTicker(emailPollInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + case <-w.nudge: + select { + case <-ctx.Done(): + return + case <-time.After(emailNudgeDelay): + } + } + if !w.settings.EmailEnabled(ctx) || !w.sender.Enabled(ctx) { + continue + } + w.runPass(ctx) + } +} + +// runPass attempts one send per eligible account. Failures back off per +// account; the pass aborts entirely on ErrNotConfigured or after a few +// consecutive failures, since both indicate a global SMTP problem. +func (w *EmailWorker) runPass(ctx context.Context) { + recipients, err := w.prefs.ListActiveRecipients(ctx) + if err != nil { + w.logger.Error("email pass: list recipients failed", "error", err) + return + } + if len(recipients) == 0 { + return + } + allowPerEpisode := w.settings.EmailAllowPerEpisode(ctx) + digestHour := w.settings.EmailDigestHour(ctx) + now := w.now() + + failures := 0 + for _, rec := range recipients { + if ctx.Err() != nil || failures >= emailMaxFailuresPerPass { + return + } + if !emailRetryEligible(now, rec.LastAttemptAt, rec.ConsecutiveFailures) { + continue + } + mode := effectiveEmailMode(rec.Mode, allowPerEpisode) + switch mode { + case EmailModePerEpisode: + // Cheap pre-check so idle accounts don't open a claim + // transaction every pass. A stale watermark only ever + // produces a harmless extra claim. + pending, err := w.prefs.HasDeliveriesSince(ctx, rec.UserID, + Cursor{CreatedAt: rec.WatermarkCreatedAt, ID: rec.WatermarkID}) + if err != nil { + w.logger.Warn("email pass: pending check failed", "user_id", rec.UserID, "error", err) + continue + } + if !pending { + continue + } + case EmailModeDailyDigest: + if !emailDigestDue(now, digestHour, rec.LastDigestAt) { + continue + } + default: + continue + } + if err := w.processAccount(ctx, rec); err != nil { + if errors.Is(err, mail.ErrNotConfigured) { + return // email turned off mid-pass; nothing else will send either + } + failures++ + w.logger.Warn("email send failed", "user_id", rec.UserID, "mode", mode, "error", err) + } + } +} + +// processAccount sends one account's pending notifications under the prefs +// row lock. The SMTP send happens inside the claim transaction: the row lock +// is per-account and only contends with other nodes, and committing the +// watermark only after a successful send is what makes the channel durable. +func (w *EmailWorker) processAccount(ctx context.Context, rec emailRecipient) error { + tx, err := w.pool.Begin(ctx) + if err != nil { + return fmt.Errorf("begin email dispatch tx: %w", err) + } + defer func() { _ = tx.Rollback(ctx) }() + + claimed, err := w.prefs.claimForUpdate(ctx, tx, rec.UserID) + if err != nil { + return err + } + if claimed == nil { + return nil // another node is handling this account + } + + // Re-derive eligibility from the locked row: the pre-scan snapshot may + // predate a user mode flip or another node's digest stamp. + mode := effectiveEmailMode(claimed.Mode, w.settings.EmailAllowPerEpisode(ctx)) + switch mode { + case EmailModePerEpisode: + case EmailModeDailyDigest: + if !emailDigestDue(w.now(), w.settings.EmailDigestHour(ctx), claimed.LastDigestAt) { + return nil + } + default: + return nil + } + + since := Cursor{CreatedAt: claimed.WatermarkCreatedAt, ID: claimed.WatermarkID} + rows, err := w.deliveries.ListForUserSince(ctx, tx, rec.UserID, since, emailFetchLimit) + if err != nil { + return err + } + + var digestAt *time.Time + if mode == EmailModeDailyDigest { + now := w.now() + digestAt = &now + } + if len(rows) == 0 { + // Nothing new. Digests still stamp so eligibility stops re-checking + // until tomorrow; the watermark needs no update. + if digestAt != nil { + if err := w.prefs.markSent(ctx, tx, rec.UserID, since, digestAt); err != nil { + return err + } + return tx.Commit(ctx) + } + return nil + } + + items := rows + if mode == EmailModeDailyDigest { + // The digest reports what the user hasn't seen; rows already read in + // another client are skipped but the watermark still passes them. + items = make([]DeliveryRow, 0, len(rows)) + for _, row := range rows { + if row.ReadAt == nil { + items = append(items, row) + } + } + } + + last := rows[len(rows)-1] + watermark := Cursor{CreatedAt: last.CreatedAt, ID: last.ID} + + if len(items) > 0 { + content := composeNotificationEmail(mode, items, w.settings.EmailExternalURL(ctx)) + err = w.sender.Send(ctx, mail.Message{ + To: []string{rec.Email}, + Subject: content.Subject, + TextBody: content.Text, + HTMLBody: content.HTML, + }) + if err != nil { + if errors.Is(err, mail.ErrNotConfigured) { + return err + } + if markErr := w.prefs.markFailure(ctx, tx, rec.UserID); markErr != nil { + return errors.Join(err, markErr) + } + if commitErr := tx.Commit(ctx); commitErr != nil { + return errors.Join(err, commitErr) + } + return err + } + w.logger.Info("notification email sent", + "user_id", rec.UserID, "mode", mode, "items", len(items)) + } + + if err := w.prefs.markSent(ctx, tx, rec.UserID, watermark, digestAt); err != nil { + return err + } + return tx.Commit(ctx) +} + +// Errors surfaced by SetEmailMode for the API layer to map to 4xx responses. +var ( + ErrEmailModeInvalid = errors.New("invalid email notification mode") + ErrEmailModeNotAllowed = errors.New("per-episode email is disabled by the administrator") + ErrEmailNoAddress = errors.New("account has no email address") +) + +// EmailAvailable reports whether the email channel can deliver right now: +// a sender is wired, the kill switch is on, and SMTP is configured. +func (s *System) EmailAvailable(ctx context.Context) bool { + return s != nil && s.emailWorker != nil && + s.Settings.EmailEnabled(ctx) && s.mailSender.Enabled(ctx) +} + +// EmailMode returns the account's chosen email mode (off when never set). +func (s *System) EmailMode(ctx context.Context, userID int) (string, error) { + if s == nil || s.EmailPrefs == nil { + return EmailModeOff, nil + } + prefs, err := s.EmailPrefs.Get(ctx, userID) + if err != nil { + return "", err + } + return prefs.Mode, nil +} + +// SetEmailMode validates and stores the account's email mode. Enabling +// requires an email address on the account and, for per-episode, the admin +// allowance. +func (s *System) SetEmailMode(ctx context.Context, userID int, mode string) error { + if s == nil || s.EmailPrefs == nil { + return ErrEmailModeInvalid + } + if !ValidEmailMode(mode) { + return ErrEmailModeInvalid + } + if mode == EmailModePerEpisode && !s.Settings.EmailAllowPerEpisode(ctx) { + return ErrEmailModeNotAllowed + } + if mode != EmailModeOff { + var email string + err := s.pool.QueryRow(ctx, + `SELECT COALESCE(email, '') FROM users WHERE id = $1`, userID, + ).Scan(&email) + if err != nil { + return fmt.Errorf("look up account email: %w", err) + } + if email == "" { + return ErrEmailNoAddress + } + } + return s.EmailPrefs.SetMode(ctx, userID, mode) +} + +// EmailDispatcher plugs the email worker into the MultiDispatcher: a new +// delivery just nudges the sweep, which reads everything since the watermark. +// No per-delivery state is kept, so dropped nudges cost only poll latency. +type EmailDispatcher struct { + worker *EmailWorker +} + +func newEmailDispatcher(worker *EmailWorker) *EmailDispatcher { + return &EmailDispatcher{worker: worker} +} + +// Dispatch implements Dispatcher. +func (d *EmailDispatcher) Dispatch(_ context.Context, _ DeliveryRow) error { + if d != nil { + d.worker.Nudge() + } + return nil +} diff --git a/internal/notifications/email_logic_test.go b/internal/notifications/email_logic_test.go new file mode 100644 index 00000000..064d693e --- /dev/null +++ b/internal/notifications/email_logic_test.go @@ -0,0 +1,191 @@ +package notifications + +import ( + "fmt" + "strings" + "testing" + "time" +) + +func TestEffectiveEmailMode(t *testing.T) { + if got := effectiveEmailMode(EmailModePerEpisode, true); got != EmailModePerEpisode { + t.Fatalf("allowed per-episode coerced to %q", got) + } + if got := effectiveEmailMode(EmailModePerEpisode, false); got != EmailModeDailyDigest { + t.Fatalf("disallowed per-episode should coerce to digest, got %q", got) + } + if got := effectiveEmailMode(EmailModeDailyDigest, false); got != EmailModeDailyDigest { + t.Fatalf("digest mode changed to %q", got) + } + if got := effectiveEmailMode(EmailModeOff, true); got != EmailModeOff { + t.Fatalf("off mode changed to %q", got) + } +} + +func TestEmailDigestDue(t *testing.T) { + loc := time.UTC + morning := time.Date(2026, 6, 11, 7, 30, 0, 0, loc) + afternoon := time.Date(2026, 6, 11, 14, 0, 0, 0, loc) + yesterday := time.Date(2026, 6, 10, 9, 0, 0, 0, loc) + today := time.Date(2026, 6, 11, 8, 5, 0, 0, loc) + + if emailDigestDue(morning, 8, nil) { + t.Fatal("digest due before today's send hour") + } + if !emailDigestDue(afternoon, 8, nil) { + t.Fatal("first-ever digest not due after send hour") + } + if !emailDigestDue(afternoon, 8, &yesterday) { + t.Fatal("digest not due when last one was yesterday") + } + if emailDigestDue(afternoon, 8, &today) { + t.Fatal("digest due twice in one day") + } +} + +func TestEmailRetryEligible(t *testing.T) { + now := time.Date(2026, 6, 11, 12, 0, 0, 0, time.UTC) + recent := now.Add(-30 * time.Second) + stale := now.Add(-10 * time.Minute) + + if !emailRetryEligible(now, nil, 0) { + t.Fatal("clean account not eligible") + } + if !emailRetryEligible(now, &recent, 0) { + t.Fatal("successful account not eligible") + } + if emailRetryEligible(now, &recent, 1) { + t.Fatal("eligible 30s after first failure (backoff is 1m)") + } + if !emailRetryEligible(now, &stale, 3) { + t.Fatal("not eligible 10m after third failure (backoff is 4m)") + } + // Large failure counts must not overflow the shift; cap applies. + old := now.Add(-7 * time.Hour) + if !emailRetryEligible(now, &old, 60) { + t.Fatal("not eligible past the 6h backoff cap") + } + if emailRetryEligible(now, &recent, 60) { + t.Fatal("eligible 30s after many failures") + } +} + +// emailEpisodeRow builds an episode.available row for one profile. +func emailEpisodeRow(id, profileID, episodeID string, season, episode int) DeliveryRow { + seriesID := "series-123" + return DeliveryRow{ + Delivery: Delivery{ + ID: id, + ProfileID: profileID, + SeriesID: &seriesID, + EpisodeID: &episodeID, + Type: DeliveryTypeEpisodeAvailable, + ReasonFlags: []byte(`{"favorite":true}`), + CreatedAt: time.Date(2026, 6, 11, 12, 0, 0, 0, time.UTC), + }, + SeriesTitle: "Severance", + EpisodeTitle: fmt.Sprintf("Episode %d", episode), + SeasonNumber: &season, + EpisodeNumber: &episode, + } +} + +func TestCollateEmailItemsDedupesAcrossProfiles(t *testing.T) { + rows := []DeliveryRow{ + emailEpisodeRow("01A", "profile-1", "ep-1", 2, 3), + emailEpisodeRow("01B", "profile-2", "ep-1", 2, 3), // same episode, second profile + emailEpisodeRow("01C", "profile-1", "ep-2", 2, 4), + requestFulfilledTestRow(), + {Delivery: Delivery{ID: "01D", Type: DeliveryTypeWebhookAutoDisabled}}, + } + items := collateEmailItems(rows) + if items.episodes != 2 { + t.Fatalf("expected 2 deduped episodes, got %d", items.episodes) + } + if len(items.series) != 1 || items.series[0].title != "Severance" { + t.Fatalf("unexpected series groups: %+v", items.series) + } + if len(items.requests) != 1 || len(items.others) != 1 { + t.Fatalf("unexpected request/other split: %d/%d", len(items.requests), len(items.others)) + } +} + +func TestCollateEmailItemsSortsEpisodesWithinSeries(t *testing.T) { + rows := []DeliveryRow{ + emailEpisodeRow("01A", "profile-1", "ep-2", 2, 4), + emailEpisodeRow("01B", "profile-1", "ep-1", 2, 3), + } + items := collateEmailItems(rows) + first := items.series[0].episodes[0] + if *first.EpisodeNumber != 3 { + t.Fatalf("episodes not sorted by number: got E%d first", *first.EpisodeNumber) + } +} + +func TestEmailSubject(t *testing.T) { + single := collateEmailItems([]DeliveryRow{emailEpisodeRow("01A", "p1", "ep-1", 2, 3)}) + if got := emailSubject(EmailModePerEpisode, single); got != "New episode of Severance: S02E03 — Episode 3" { + t.Fatalf("unexpected single-episode subject %q", got) + } + + request := collateEmailItems([]DeliveryRow{requestFulfilledTestRow()}) + if got := emailSubject(EmailModePerEpisode, request); got != "Dune is now available" { + t.Fatalf("unexpected single-request subject %q", got) + } + + mixed := collateEmailItems([]DeliveryRow{ + emailEpisodeRow("01A", "p1", "ep-1", 2, 3), + emailEpisodeRow("01B", "p1", "ep-2", 2, 4), + requestFulfilledTestRow(), + }) + if got := emailSubject(EmailModePerEpisode, mixed); got != "Silo: 2 new episodes, 1 request ready" { + t.Fatalf("unexpected mixed subject %q", got) + } + if got := emailSubject(EmailModeDailyDigest, mixed); got != "Silo daily digest: 2 new episodes, 1 request ready" { + t.Fatalf("unexpected digest subject %q", got) + } +} + +func TestComposeNotificationEmailLinks(t *testing.T) { + rows := []DeliveryRow{emailEpisodeRow("01A", "p1", "ep-1", 2, 3)} + + withLinks := composeNotificationEmail(EmailModePerEpisode, rows, "https://silo.example.com") + if !strings.Contains(withLinks.HTML, `href="https://silo.example.com/item/ep-1"`) { + t.Fatalf("episode link missing from HTML:\n%s", withLinks.HTML) + } + if !strings.Contains(withLinks.HTML, `href="https://silo.example.com/settings/notifications"`) { + t.Fatalf("settings link missing from HTML footer:\n%s", withLinks.HTML) + } + + withoutLinks := composeNotificationEmail(EmailModePerEpisode, rows, "") + if strings.Contains(withoutLinks.HTML, "href=") { + t.Fatalf("HTML contains links with no external URL configured:\n%s", withoutLinks.HTML) + } + if !strings.Contains(withoutLinks.Text, "S02E03 — Episode 3") { + t.Fatalf("text body missing episode line:\n%s", withoutLinks.Text) + } +} + +func TestComposeNotificationEmailEscapesHTML(t *testing.T) { + row := emailEpisodeRow("01A", "p1", "ep-1", 2, 3) + row.SeriesTitle = `` + content := composeNotificationEmail(EmailModePerEpisode, []DeliveryRow{row}, "") + if strings.Contains(content.HTML, "