feat(notifications): email notification channel

Adds email as a notification channel built on the shared SMTP core
(mail.Sender). Email mode is a per-account preference (off, daily
digest, or per-episode) stored in notification_email_prefs; delivery is
an account-watermark sweep over notification_deliveries that dedupes
cross-profile duplicates, advancing the watermark only after a
successful send. Admin controls cover the channel kill switch, the
per-episode allowance (off coerces those accounts to the digest),
digest hour, and an external URL for deep links inside emails.
Availability is advertised through /notifications/capability and the
user settings page gains an Email section for opt-in.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Quick
2026-06-11 18:09:37 -04:00
co-authored by Claude Fable 5
parent e5b210589d
commit df95e3cb95
18 changed files with 1540 additions and 1 deletions
+2
View File
@@ -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
@@ -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`).
+22
View File
@@ -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,
})
}
@@ -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})
}
+2
View File
@@ -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)
+16
View File
@@ -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,
+304
View File
@@ -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(
`<li style="margin:2px 0;"><a href="%s" style="color:#6d6df7;text-decoration:none;">%s</a></li>`,
html.EscapeString(href), html.EscapeString(plain)))
} else {
body.WriteString(fmt.Sprintf(`<li style="margin:2px 0;">%s</li>`, html.EscapeString(plain)))
}
}
writeHeading := func(title, href string) {
text.WriteString(title + "\n")
label := html.EscapeString(title)
if href != "" {
label = fmt.Sprintf(`<a href="%s" style="color:inherit;text-decoration:none;">%s</a>`,
html.EscapeString(href), label)
}
body.WriteString(fmt.Sprintf(
`<h3 style="margin:14px 0 4px;font-size:15px;">%s</h3>`, label))
}
openList := func() { body.WriteString(`<ul style="margin:4px 0;padding-left:20px;">`) }
closeList := func() { body.WriteString(`</ul>`) }
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(
`<p style="margin:8px 0;color:#888;">%s</p>`, 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(`<a href="%s" style="color:#888;">Settings → Notifications</a>`, settingsURL), 1)
}
htmlBody := fmt.Sprintf(`<div style="font-family:-apple-system,Segoe UI,Roboto,Helvetica,Arial,sans-serif;font-size:14px;line-height:1.5;color:#1a1a1a;max-width:560px;">
<p style="margin:0 0 8px;">%s</p>
%s
<hr style="border:none;border-top:1px solid #e5e5e5;margin:16px 0 8px;">
<p style="margin:0;font-size:12px;color:#888;">%s</p>
</div>`,
html.EscapeString(intro), body.String(), footerHTML)
return emailContent{
Subject: emailSubject(mode, items),
Text: intro + "\n\n" + text.String() + "\n" + footer + "\n",
HTML: htmlBody,
}
}
+360
View File
@@ -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
}
+191
View File
@@ -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 = `<script>alert("x")</script>`
content := composeNotificationEmail(EmailModePerEpisode, []DeliveryRow{row}, "")
if strings.Contains(content.HTML, "<script>") {
t.Fatalf("series title not escaped:\n%s", content.HTML)
}
}
func TestComposeNotificationEmailCapsRenderedItems(t *testing.T) {
rows := make([]DeliveryRow, 0, emailMaxItemsRendered+10)
for i := range emailMaxItemsRendered + 10 {
rows = append(rows, emailEpisodeRow(
fmt.Sprintf("01%03d", i), "p1", fmt.Sprintf("ep-%d", i), 1, i+1))
}
content := composeNotificationEmail(EmailModeDailyDigest, rows, "")
if !strings.Contains(content.Text, "and 10 more in your Silo inbox") {
t.Fatalf("overflow line missing:\n%s", content.Text)
}
if count := strings.Count(content.HTML, "<li"); count != emailMaxItemsRendered {
t.Fatalf("expected %d rendered items, got %d", emailMaxItemsRendered, count)
}
}
+210
View File
@@ -0,0 +1,210 @@
package notifications
import (
"context"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// Email notification modes. The channel is account-level: email addresses
// live on users, not profiles, so one setting covers every profile on the
// account and the worker collapses cross-profile duplicates.
const (
EmailModeOff = "off"
EmailModePerEpisode = "per_episode"
EmailModeDailyDigest = "daily_digest"
)
// ValidEmailMode reports whether mode is a recognized email mode value.
func ValidEmailMode(mode string) bool {
switch mode {
case EmailModeOff, EmailModePerEpisode, EmailModeDailyDigest:
return true
}
return false
}
// EmailPrefs is one account's email notification state: the user-chosen mode
// plus the worker's dispatch watermark and failure backoff counters.
type EmailPrefs struct {
UserID int
Mode string
WatermarkCreatedAt time.Time
WatermarkID string
LastDigestAt *time.Time
LastAttemptAt *time.Time
ConsecutiveFailures int
}
// emailRecipient pairs an active prefs row with the account's email address.
type emailRecipient struct {
EmailPrefs
Email string
}
// EmailPrefsRepository owns notification_email_prefs.
type EmailPrefsRepository struct {
pool *pgxpool.Pool
}
// NewEmailPrefsRepository creates an EmailPrefsRepository.
func NewEmailPrefsRepository(pool *pgxpool.Pool) *EmailPrefsRepository {
return &EmailPrefsRepository{pool: pool}
}
// Get returns the account's email prefs; missing rows default to mode off.
func (r *EmailPrefsRepository) Get(ctx context.Context, userID int) (EmailPrefs, error) {
prefs := EmailPrefs{UserID: userID, Mode: EmailModeOff}
err := r.pool.QueryRow(ctx, `
SELECT mode, watermark_created_at, watermark_id, last_digest_at,
last_attempt_at, consecutive_failures
FROM notification_email_prefs WHERE user_id = $1`,
userID,
).Scan(&prefs.Mode, &prefs.WatermarkCreatedAt, &prefs.WatermarkID,
&prefs.LastDigestAt, &prefs.LastAttemptAt, &prefs.ConsecutiveFailures)
if errors.Is(err, pgx.ErrNoRows) {
return prefs, nil
}
if err != nil {
return prefs, fmt.Errorf("get email prefs: %w", err)
}
return prefs, nil
}
// SetMode upserts the account's email mode. Enabling from off (or creating
// the row) resets the watermark to now so the backlog never floods a fresh
// opt-in, and clears failure backoff so the first send happens promptly.
func (r *EmailPrefsRepository) SetMode(ctx context.Context, userID int, mode string) error {
if !ValidEmailMode(mode) {
return fmt.Errorf("invalid email mode %q", mode)
}
_, err := r.pool.Exec(ctx, `
INSERT INTO notification_email_prefs (user_id, mode)
VALUES ($1, $2)
ON CONFLICT (user_id) DO UPDATE SET
mode = EXCLUDED.mode,
watermark_created_at = CASE
WHEN notification_email_prefs.mode = 'off' THEN now()
ELSE notification_email_prefs.watermark_created_at
END,
watermark_id = CASE
WHEN notification_email_prefs.mode = 'off' THEN ''
ELSE notification_email_prefs.watermark_id
END,
last_attempt_at = NULL,
consecutive_failures = 0,
updated_at = now()`,
userID, mode)
if err != nil {
return fmt.Errorf("set email mode: %w", err)
}
return nil
}
// ListActiveRecipients returns every account with email notifications on and
// a usable address. Disabled or deleted accounts drop out of the join.
func (r *EmailPrefsRepository) ListActiveRecipients(ctx context.Context) ([]emailRecipient, error) {
rows, err := r.pool.Query(ctx, `
SELECT p.user_id, p.mode, p.watermark_created_at, p.watermark_id,
p.last_digest_at, p.last_attempt_at, p.consecutive_failures,
u.email
FROM notification_email_prefs p
JOIN users u ON u.id = p.user_id AND u.enabled AND COALESCE(u.email, '') <> ''
WHERE p.mode <> 'off'
ORDER BY p.user_id`)
if err != nil {
return nil, fmt.Errorf("list email recipients: %w", err)
}
defer rows.Close()
out := make([]emailRecipient, 0, 8)
for rows.Next() {
var rec emailRecipient
if err := rows.Scan(&rec.UserID, &rec.Mode, &rec.WatermarkCreatedAt, &rec.WatermarkID,
&rec.LastDigestAt, &rec.LastAttemptAt, &rec.ConsecutiveFailures, &rec.Email); err != nil {
return nil, fmt.Errorf("scan email recipient: %w", err)
}
out = append(out, rec)
}
return out, rows.Err()
}
// claimForUpdate locks the account's prefs row for one dispatch attempt.
// SKIP LOCKED makes concurrent nodes pass over each other's in-flight users
// instead of double-sending; (nil, nil) means another node holds the row.
func (r *EmailPrefsRepository) claimForUpdate(ctx context.Context, tx pgx.Tx, userID int) (*EmailPrefs, error) {
prefs := EmailPrefs{UserID: userID}
err := tx.QueryRow(ctx, `
SELECT mode, watermark_created_at, watermark_id, last_digest_at,
last_attempt_at, consecutive_failures
FROM notification_email_prefs
WHERE user_id = $1
FOR UPDATE SKIP LOCKED`,
userID,
).Scan(&prefs.Mode, &prefs.WatermarkCreatedAt, &prefs.WatermarkID,
&prefs.LastDigestAt, &prefs.LastAttemptAt, &prefs.ConsecutiveFailures)
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
if err != nil {
return nil, fmt.Errorf("claim email prefs: %w", err)
}
return &prefs, nil
}
// markSent advances the watermark past everything the email covered and
// resets failure backoff. digestAt is non-nil for digest sends (including
// empty digests, so eligibility stops re-checking until tomorrow).
func (r *EmailPrefsRepository) markSent(ctx context.Context, tx pgx.Tx, userID int, watermark Cursor, digestAt *time.Time) error {
_, err := tx.Exec(ctx, `
UPDATE notification_email_prefs SET
watermark_created_at = $2,
watermark_id = $3,
last_digest_at = COALESCE($4, last_digest_at),
last_attempt_at = now(),
consecutive_failures = 0,
updated_at = now()
WHERE user_id = $1`,
userID, watermark.CreatedAt, watermark.ID, digestAt)
if err != nil {
return fmt.Errorf("mark email sent: %w", err)
}
return nil
}
// markFailure records a failed send for backoff; the watermark stays put so
// the next eligible pass retries the same items.
func (r *EmailPrefsRepository) markFailure(ctx context.Context, tx pgx.Tx, userID int) error {
_, err := tx.Exec(ctx, `
UPDATE notification_email_prefs SET
last_attempt_at = now(),
consecutive_failures = consecutive_failures + 1,
updated_at = now()
WHERE user_id = $1`,
userID)
if err != nil {
return fmt.Errorf("mark email failure: %w", err)
}
return nil
}
// HasDeliveriesSince reports whether the account has any delivery newer than
// the watermark. Cheap pre-check (index-only) so the per-episode sweep does
// not open a claim transaction for idle accounts every pass.
func (r *EmailPrefsRepository) HasDeliveriesSince(ctx context.Context, userID int, since Cursor) (bool, error) {
var exists bool
err := r.pool.QueryRow(ctx, `
SELECT EXISTS (
SELECT 1 FROM notification_deliveries
WHERE user_id = $1 AND (created_at, id) > ($2, $3)
)`,
userID, since.CreatedAt, since.ID,
).Scan(&exists)
if err != nil {
return false, fmt.Errorf("check deliveries since watermark: %w", err)
}
return exists, nil
}
+32
View File
@@ -34,6 +34,11 @@ const (
SettingWebhooksMaxPerProfile = "notifications.webhooks.max_per_profile"
SettingWebhooksAllowPrivate = "notifications.webhooks.allow_private_destinations"
SettingWebhooksRatePerMinute = "notifications.webhooks.deliveries_per_minute_per_profile"
SettingEmailEnabled = "notifications.email_enabled"
SettingEmailAllowPerEpisode = "notifications.email.allow_per_episode"
SettingEmailDigestHour = "notifications.email.digest_hour"
SettingEmailExternalURL = "notifications.email.external_url"
)
const (
@@ -43,6 +48,7 @@ const (
defaultRetentionReadDays = 90
defaultRetentionUnread = 180
defaultRetentionEventDays = 30
defaultEmailDigestHour = 8
settingsCacheTTL = 15 * time.Second
)
@@ -197,3 +203,29 @@ func (s *Settings) WebhooksAllowPrivateDestinations(ctx context.Context) bool {
func (s *Settings) WebhooksDeliveriesPerMinute(ctx context.Context) int {
return s.intSetting(ctx, SettingWebhooksRatePerMinute, 60, 1, 10000)
}
// EmailEnabled gates the email notification channel (kill switch). Actual
// availability additionally requires a configured SMTP sender (mail.Sender).
func (s *Settings) EmailEnabled(ctx context.Context) bool {
return s.boolSetting(ctx, SettingEmailEnabled, true)
}
// EmailAllowPerEpisode controls whether users may choose per-episode email
// alerts. When off, accounts set to per-episode are coerced to the daily
// digest instead of going silent.
func (s *Settings) EmailAllowPerEpisode(ctx context.Context) bool {
return s.boolSetting(ctx, SettingEmailAllowPerEpisode, true)
}
// EmailDigestHour is the hour of day (0-23, server-local time) at which daily
// digest emails go out.
func (s *Settings) EmailDigestHour(ctx context.Context) int {
return s.intSetting(ctx, SettingEmailDigestHour, defaultEmailDigestHour, 0, 23)
}
// EmailExternalURL is the externally reachable base URL of this server, used
// for deep links inside notification emails. Empty renders emails without
// links (webhooks deliberately never leak the origin; email is opt-in here).
func (s *Settings) EmailExternalURL(ctx context.Context) string {
return strings.TrimRight(strings.TrimSpace(s.raw(ctx, SettingEmailExternalURL)), "/")
}
+29 -1
View File
@@ -10,6 +10,7 @@ import (
"time"
evt "github.com/Silo-Server/silo-server/internal/events"
"github.com/Silo-Server/silo-server/internal/mail"
"github.com/Silo-Server/silo-server/internal/models"
"github.com/Silo-Server/silo-server/internal/secret"
"github.com/Silo-Server/silo-server/internal/userstore"
@@ -49,6 +50,11 @@ type System struct {
// WebPush is nil when the settings store is not writable (VAPID keys
// could not be provisioned).
WebPush *WebPushService
// EmailPrefs is nil when no mail sender was provided.
EmailPrefs *EmailPrefsRepository
mailSender mail.Sender
emailWorker *EmailWorker
webhookRepo *WebhookRepository
webhookDispatcher *WebhookDispatcher
@@ -69,7 +75,8 @@ type System struct {
}
// NewSystem wires the notification system. hub may be nil (no realtime
// publishing); redisClient may be nil (in-memory websocket tickets).
// publishing); redisClient may be nil (in-memory websocket tickets);
// mailSender may be nil (no email channel).
func NewSystem(
pool *pgxpool.Pool,
settingsReader SettingReader,
@@ -79,6 +86,7 @@ func NewSystem(
hub *evt.Hub,
redisClient *redis.Client,
cipher *secret.Cipher,
mailSender mail.Sender,
) *System {
settings := NewSettings(settingsReader)
releases := NewReleaseRepository(pool)
@@ -120,6 +128,16 @@ func NewSystem(
dispatchers = append(dispatchers, webPushDispatcher)
}
// Email rides the shared SMTP core. Unlike the per-target channels it
// keeps no outbox: its dispatcher only nudges the watermark sweep.
var emailPrefs *EmailPrefsRepository
var emailWorker *EmailWorker
if mailSender != nil {
emailPrefs = NewEmailPrefsRepository(pool)
emailWorker = newEmailWorker(pool, deliveries, emailPrefs, settings, mailSender)
dispatchers = append(dispatchers, newEmailDispatcher(emailWorker))
}
multiDispatcher := NewMultiDispatcher(dispatchers...)
fanout := NewFanoutWorker(pool, releases, interests, deliveries, preferences, settings, multiDispatcher)
if webhookRepo != nil {
@@ -144,6 +162,9 @@ func NewSystem(
Tickets: NewTicketStore(redisClient),
Webhooks: webhookService,
WebPush: webPushService,
EmailPrefs: emailPrefs,
mailSender: mailSender,
emailWorker: emailWorker,
webhookRepo: webhookRepo,
webhookDispatcher: webhookDispatcher,
webhookRetry: webhookRetry,
@@ -218,6 +239,13 @@ func (s *System) Start(ctx context.Context) {
s.webhookRetry.Run(ctx)
}()
}
if s.emailWorker != nil {
s.wg.Add(1)
go func() {
defer s.wg.Done()
s.emailWorker.Run(ctx)
}()
}
if s.webPushDispatcher != nil {
s.wg.Add(1)
go func() {
@@ -0,0 +1,34 @@
-- +goose Up
-- +goose StatementBegin
-- Email notification channel (docs/superpowers/plans/notifications/06,
-- item 3). Email addresses live on login accounts (users), not profiles, so
-- the mode and dispatch state are account-level. The email worker sweeps
-- notification_deliveries per user and advances the watermark only after a
-- successful SMTP send, so a crash or SMTP outage re-sends instead of
-- dropping; the watermark is initialized to now() whenever the channel is
-- enabled so history never floods a fresh opt-in.
CREATE TABLE public.notification_email_prefs (
user_id integer PRIMARY KEY,
mode text NOT NULL DEFAULT 'off'
CHECK (mode IN ('off', 'per_episode', 'daily_digest')),
watermark_created_at timestamptz NOT NULL DEFAULT now(),
watermark_id text NOT NULL DEFAULT '',
last_digest_at timestamptz,
last_attempt_at timestamptz,
consecutive_failures integer NOT NULL DEFAULT 0,
updated_at timestamptz NOT NULL DEFAULT now()
);
-- Deliberately no FK to users: notification tables stay FK-free toward
-- account/profile storage (see 20260611100000). Rows for deleted accounts
-- drop out of the recipient join and are inert.
-- The email sweep reads deliveries by account, not profile.
CREATE INDEX notification_deliveries_user_created_idx
ON public.notification_deliveries (user_id, created_at, id);
-- +goose StatementEnd
-- +goose Down
-- +goose StatementBegin
DROP INDEX IF EXISTS public.notification_deliveries_user_created_idx;
DROP TABLE IF EXISTS public.notification_email_prefs;
-- +goose StatementEnd
+8
View File
@@ -2333,6 +2333,14 @@ export interface NotificationCapability {
android_push: { available: boolean; provider: string; supported_modes: string[] };
web_push: { available: boolean; public_key?: string };
webhooks: { available: boolean; max_per_profile: number; supported_types: string[] };
email: { available: boolean; modes: string[]; digest_hour: number };
}
export type NotificationEmailMode = "off" | "per_episode" | "daily_digest";
/** Account-level (not per-profile): one mode covers all profiles. */
export interface NotificationEmailPreferences {
mode: NotificationEmailMode;
}
export interface WebPushSubscriptionView {
+1
View File
@@ -226,6 +226,7 @@ export const notificationKeys = {
list: (status: "all" | "unread" = "all") => ["notifications", "list", status] as const,
unreadCount: () => ["notifications", "unread-count"] as const,
preferences: () => ["notifications", "preferences"] as const,
emailPreferences: () => ["notifications", "email-preferences"] as const,
capability: () => ["notifications", "capability"] as const,
webhooks: () => ["notifications", "webhooks"] as const,
webPushSubscriptions: () => ["notifications", "web-push-subscriptions"] as const,
+26
View File
@@ -3,6 +3,7 @@ import { useInfiniteQuery, useMutation, useQuery, useQueryClient } from "@tansta
import { api } from "@/api/client";
import type {
AppNotification,
NotificationEmailPreferences,
NotificationListResponse,
NotificationPreferences,
NotificationReadEventPayload,
@@ -93,6 +94,31 @@ export function useUpdateNotificationPreferences() {
});
}
export function useEmailNotificationPreferences(enabled = true) {
return useQuery({
queryKey: notificationKeys.emailPreferences(),
queryFn: () => api<NotificationEmailPreferences>("/notifications/email-preferences"),
enabled,
});
}
export function useUpdateEmailNotificationPreferences() {
const queryClient = useQueryClient();
return useMutation({
mutationFn: (update: NotificationEmailPreferences) =>
api<NotificationEmailPreferences>("/notifications/email-preferences", {
method: "PUT",
body: JSON.stringify(update),
}),
onSuccess: (prefs) => {
queryClient.setQueryData(notificationKeys.emailPreferences(), prefs);
},
onError: (error) => {
toast.error(error instanceof Error ? error.message : "Failed to save email preferences");
},
});
}
// --- Realtime cache reducers (used by RealtimeEventsProvider) ---
type NotificationsInfiniteData = {
@@ -17,6 +17,10 @@ const KEYS = [
"notifications.webhooks.max_per_profile",
"notifications.webhooks.allow_private_destinations",
"notifications.webhooks.deliveries_per_minute_per_profile",
"notifications.email_enabled",
"notifications.email.allow_per_episode",
"notifications.email.digest_hour",
"notifications.email.external_url",
"notifications.retention.read_days",
"notifications.retention.unread_days",
"notifications.retention.event_days",
@@ -145,6 +149,37 @@ export default function NotificationsAdminSettings() {
)}
</FieldGroup>
<FieldGroup label="Email">
<SettingField
label="Email Notifications"
hint="Deliver notifications by email to accounts that opt in. Requires SMTP to be configured on the Email page."
type="toggle"
value={toggleValue("notifications.email_enabled")}
onChange={(v) => form.setValue("notifications.email_enabled", v)}
/>
<SettingField
label="Allow Per-Episode Email"
hint="Let users choose an email per episode instead of the daily digest. Off coerces those accounts to the digest."
type="toggle"
value={toggleValue("notifications.email.allow_per_episode")}
onChange={(v) => form.setValue("notifications.email.allow_per_episode", v)}
/>
<SettingField
label="Digest Hour"
hint="Hour of day (0-23, server time) when daily digest emails go out (default 8)"
type="number"
value={numberValue("notifications.email.digest_hour", "8")}
onChange={(v) => form.setValue("notifications.email.digest_hour", v)}
/>
<SettingField
label="External URL"
hint="Public base URL of this server (e.g. https://silo.example.com) used for links inside emails. Empty sends emails without links."
type="text"
value={form.getValue("notifications.email.external_url")}
onChange={(v) => form.setValue("notifications.email.external_url", v)}
/>
</FieldGroup>
<FieldGroup label="Retention">
<SettingField
label="Read Notifications (days)"
@@ -16,6 +16,7 @@ import {
} from "lucide-react";
import { toast } from "sonner";
import type {
NotificationEmailMode,
NotificationPreferences,
NotificationWebhook,
NotificationWebhookInput,
@@ -36,10 +37,20 @@ import {
} from "@/components/ui/dialog";
import { Input } from "@/components/ui/input";
import { Label } from "@/components/ui/label";
import {
Select,
SelectContent,
SelectItem,
SelectTrigger,
SelectValue,
} from "@/components/ui/select";
import { Skeleton } from "@/components/ui/skeleton";
import { Switch } from "@/components/ui/switch";
import { useAuth } from "@/hooks/useAuth";
import {
useEmailNotificationPreferences,
useNotificationPreferences,
useUpdateEmailNotificationPreferences,
useUpdateNotificationPreferences,
} from "@/hooks/queries/notifications";
import {
@@ -141,6 +152,90 @@ function PreferencesSection() {
);
}
function EmailSection() {
const { user } = useAuth();
const capability = useNotificationCapability();
const emailCap = capability.data?.email;
const available = emailCap?.available ?? false;
const { data: prefs, isLoading } = useEmailNotificationPreferences(available);
const updatePrefs = useUpdateEmailNotificationPreferences();
if (!available) {
return null;
}
const mode = prefs?.mode ?? "off";
const enabled = mode !== "off";
const allowPerEpisode = emailCap?.modes.includes("per_episode") ?? false;
const digestHour = String(emailCap?.digest_hour ?? 8).padStart(2, "0");
if (isLoading) {
return (
<SettingsGroup title="Email Notifications">
<Skeleton className="h-16 w-full" />
</SettingsGroup>
);
}
return (
<SettingsGroup
title="Email Notifications"
description="Account-wide: one email covers every profile on this account."
>
<div className="flex items-center justify-between gap-3">
<div>
<div className="text-sm font-medium">Send to {user?.email || "your account email"}</div>
<div className="text-muted-foreground text-xs">
Notifications you'd see in the inbox, delivered by email
</div>
</div>
<Switch
checked={enabled}
disabled={updatePrefs.isPending}
onCheckedChange={(checked) =>
updatePrefs.mutate({ mode: checked ? "daily_digest" : "off" })
}
/>
</div>
{enabled && (
<div className="flex items-center justify-between gap-3">
<div>
<div className="text-sm">Frequency</div>
<div className="text-muted-foreground text-xs">
{mode === "per_episode"
? "An email as soon as each notification arrives"
: `One summary per day, around ${digestHour}:00 server time`}
</div>
</div>
<Select
value={mode}
disabled={updatePrefs.isPending}
onValueChange={(value) => updatePrefs.mutate({ mode: value as NotificationEmailMode })}
>
<SelectTrigger className="w-[180px]">
<SelectValue />
</SelectTrigger>
<SelectContent>
<SelectItem value="daily_digest">Daily digest</SelectItem>
{(allowPerEpisode || mode === "per_episode") && (
<SelectItem value="per_episode" disabled={!allowPerEpisode}>
Every episode
</SelectItem>
)}
</SelectContent>
</Select>
</div>
)}
{enabled && mode === "per_episode" && !allowPerEpisode && (
<div className="text-xs text-amber-500">
Per-episode email is disabled by the administrator; you'll receive the daily digest
instead.
</div>
)}
</SettingsGroup>
);
}
/**
* Delivery health for one push subscription, derived the same way as webhook
* health: a failure newer than the last success means the device is failing.
@@ -656,6 +751,8 @@ export default function NotificationsSettings() {
<WebPushSection />
<EmailSection />
<SettingsGroup
title="Webhooks"
description="Send this profile's notifications to a webhook URL. Discord URLs render as native embeds; other URLs receive signed JSON."