diff --git a/.env.example b/.env.example index b8908d72..e699a644 100644 --- a/.env.example +++ b/.env.example @@ -59,6 +59,9 @@ POSTGRES_SHM_SIZE=8gb # docs/architecture/secret-encryption.md. # SECRET_KEY=replace-with-output-of-openssl-rand-base64-48 +# Optional public URL for Silo. If not set, the server will use the IP address of the container. +# SILO_PUBLIC_URL=https://silo.example.com + # Run from source / advanced overrides # Only DATABASE_URL is required when running Silo outside the default docker compose stack. # DATABASE_URL=postgres://silo:password@localhost:5432/silo diff --git a/cmd/silo/main.go b/cmd/silo/main.go index 507e6cdb..10b1310e 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" @@ -595,6 +596,16 @@ func main() { bootstrapSensitiveValues["redis.url"] = bc.RedisURL } + // Shared Redis client for components needing raw Redis beyond the event + // bus (websocket handshake tickets, session listing). Nil on Redis-less + // deployments; consumers fall back to in-process implementations. + apiRedisClient, apiRedisErr := cache.NewRedisClient(cfg.Redis) + if apiRedisErr != nil { + slog.Warn("redis client init failed; multi-node websocket tickets disabled", "error", apiRedisErr) + } else if apiRedisClient != nil { + defer func() { _ = apiRedisClient.Close() }() + } + deps := api.Dependencies{ Config: cfg, LiveConfig: configWatcher.Config, @@ -605,6 +616,7 @@ func main() { DB: pool, SecretCipher: dataCipher, EventBus: eventBus, + RedisClient: apiRedisClient, LogStreamHub: logStreamHub, RealtimeHub: realtimeHub, EventsHub: eventsHub, @@ -1208,8 +1220,10 @@ func main() { cfg.Scanner.MaxConcurrentLibraries, cfg.Scanner.MaxConcurrentScoped, ) - libraryScanQueue.Start() - defer libraryScanQueue.Stop() + // Started below, after the notification system has attached its + // availability detector to the executor: a scan resumed by the + // workers before that wiring would complete without recording + // episode availability, silently losing release notifications. deps.LibraryScanQueue = libraryScanQueue } if deps.DB != nil && deps.FileRepo != nil && metadataService != nil { @@ -1271,6 +1285,48 @@ func main() { } defer userStoreProvider.Close() } + + // User-facing release notifications. The system reads user state through + // the raw store provider; the provider handed to everything downstream is + // wrapped so every favorites/watchlist/progress mutation (REST handlers, + // jellycompat, imports, playback) feeds the interest index. + var notificationSystem *notifications.System + if deps.DB != nil && userStoreProvider != nil { + notificationScopes := access.NewResolver( + auth.NewUserRepository(deps.DB), + userStoreProvider, + access.NewProfileTokenService(cfg.Auth.JWTSecret, 0), + ) + notificationSystem = notifications.NewSystem( + deps.DB, + settingsRepo, + userStoreProvider, + notificationScopes, + auth.NewUserRepository(deps.DB), + deps.EventsHub, + deps.RedisClient, + deps.SecretCipher, + mail.NewSMTPSender(settingsRepo), + ) + userStoreProvider = notifications.WrapUserStoreProvider(userStoreProvider, notificationSystem) + deps.Notifications = notificationSystem + if libraryIngestExecutor != nil { + libraryIngestExecutor.SetAvailabilityDetector(notificationSystem.Detector) + } + if needsWorkers { + notificationSystem.Start(appCtx) + defer notificationSystem.Wait() + } + } + + // Start the scan queue only now that the availability detector (when + // notifications are enabled) is attached to the ingest executor, so scans + // resumed at startup cannot complete before the detector exists. + if libraryScanQueue != nil { + libraryScanQueue.Start() + defer libraryScanQueue.Stop() + } + if userStoreProvider != nil && pluginService != nil { deps.PluginUserConfig = plugins.NewUserConfigStore(userStoreProvider, pluginService) } @@ -1575,6 +1631,11 @@ func main() { } taskMgr.Register(tasks.NewActivityLogCleanupTask(deps.DB, settingsRepo, activityPM)) taskMgr.Register(tasks.NewOperationalLogCleanupTask(deps.DB, settingsRepo, opsPM)) + if notificationSystem != nil { + taskMgr.Register(tasks.NewSeedContentAvailabilityTask(notificationSystem)) + taskMgr.Register(tasks.NewRebuildReleaseInterestTask(notificationSystem)) + taskMgr.Register(tasks.NewNotificationsRetentionTask(notificationSystem)) + } if matchWorker != nil { taskMgr.Register(tasks.NewMatchMediaTask(matchWorker)) } @@ -1614,6 +1675,9 @@ func main() { ) requestReconcileSvc.SetEntitlementResolver(mediarequests.NewAccessEntitlements(reconcileResolver)) } + if notificationSystem != nil { + requestReconcileSvc.SetFulfillmentNotifier(notifications.NewRequestFulfillmentNotifier(notificationSystem)) + } taskMgr.Register(tasks.NewReconcileRequestsTask(requestReconcileSvc, 100)) if deps.FolderRepo != nil && deps.LibraryScanQueue != nil && pluginService != nil && pluginInstallationStore != nil { autoscanRepo := autoscan.NewRepository(deps.DB, deps.SecretCipher) diff --git a/docs/architecture/email.md b/docs/architecture/email.md new file mode 100644 index 00000000..baa1c8d3 --- /dev/null +++ b/docs/architecture/email.md @@ -0,0 +1,49 @@ +# Outbound Email (`internal/mail`) + +**Status:** Implemented 2026-06-11 + +Silo's shared outbound email facility. It is deliberately feature-agnostic: any +feature that sends mail (notification emails, account flows, invites) composes +a `mail.Message` and hands it to the shared `mail.Sender`, so SMTP +configuration, security policy, and diagnostics live in exactly one place. + +## Abstraction + +```go +type Sender interface { + Enabled(ctx context.Context) bool + Send(ctx context.Context, msg Message) error +} +``` + +`Message` carries recipients, subject, and text and/or HTML bodies (both set → +multipart/alternative). `Send` returns `mail.ErrNotConfigured` when email is +disabled or incomplete, so features treat email as an optional transport and +degrade gracefully. The SMTP implementation (`mail.NewSMTPSender`) is backed by +`github.com/wneessen/go-mail`. + +## Configuration + +Live server settings (no restart required; read on every send — volume is +low): + +| Key | Default | Notes | +|---|---|---| +| `email.enabled` | `false` | master switch | +| `email.smtp_host` | — | required | +| `email.smtp_port` | `587` | | +| `email.smtp_security` | `starttls` | `starttls` \| `tls` (implicit, port 465) \| `none` | +| `email.smtp_username` | — | empty = no auth | +| `email.smtp_password` | — | encrypted at rest (`SensitiveSettingKeys`) | +| `email.from_address` | — | required | +| `email.from_name` | `Silo` | | + +Admin UI: Admin Settings → Connections → Email, including a synchronous test +send (`POST /api/v1/admin/email/test`). + +## Adding a consumer + +Construct messages in the feature package and send through a `mail.Sender` +dependency. Do not read `email.*` settings from feature code, and check +`Enabled` (or branch on `ErrNotConfigured`) rather than treating a missing +SMTP configuration as an error. diff --git a/docs/superpowers/plans/notifications/05-web-push.md b/docs/superpowers/plans/notifications/05-web-push.md new file mode 100644 index 00000000..fb804c40 --- /dev/null +++ b/docs/superpowers/plans/notifications/05-web-push.md @@ -0,0 +1,54 @@ +# Web Push Spec + +**Date:** 2026-06-11 +**Status:** Implemented +**Scope:** Browser push notifications (Push API + VAPID) as a third push platform alongside the deferred APNs/FCM channels. +**Depends On:** +- [`00-architecture-overview.md`](./00-architecture-overview.md) +- [`01-release-events-and-inbox.md`](./01-release-events-and-inbox.md) + +## Why Web Push ships before APNs/FCM + +The architecture overview deferred mobile push because Apple and Google require pushes to official store builds to be signed by the publisher's credentials, forcing a Silo-operated relay. Web Push has neither problem: + +- **No accounts, no relay.** The server self-provisions a VAPID keypair on first use. Any standards-compliant browser push service (Chrome, Firefox, Edge, Safari 16+) accepts VAPID-signed requests from any origin. +- **Content-safe by protocol.** Payloads are encrypted end-to-end (RFC 8291, `aes128gcm`) to keys held only by the subscribed browser. The vendor push service relays ciphertext. Unlike the APNs/FCM design, payloads can therefore carry full display content (titles, episode numbers, poster URLs) without violating the self-hosted privacy model — there is no opaque-wake/fetch dance. + +The residual leak matches the relay threat model: the push service sees the user server's egress IP, delivery timing, and payload size. It never sees content or identity. + +## Data model + +- `web_push_subscriptions` — profile-scoped browser registrations: `endpoint` (unique; a resubscription from the same browser under a different profile reassigns the row), `p256dh`, `auth`, `device_name`, failure bookkeeping. No FK to profiles (per-user SQLite stores); profile deletion purges in code. +- `web_push_delivery_attempts` — the durable dispatch outbox, mirroring `webhook_delivery_attempts`: `pending` rows enqueued in the fanout transaction, claimed post-commit with a lease, swept by the retry loop after a crash. + +## VAPID identity + +Generated once and persisted in `server_settings`: + +- `notifications.web_push.vapid_public_key` — served to clients via the capability endpoint. +- `notifications.web_push.vapid_private_key` — encrypted at rest (`SensitiveSettingKeys`). + +The private key is persisted before the public key so a crash between writes regenerates the pair instead of stranding clients with an unusable public key. The pair must never be rotated casually: browsers bind subscriptions to it. + +## API surface (profile-scoped) + +- `GET /api/v1/notifications/capability` — `web_push: { available, public_key }`. +- `POST /api/v1/notifications/web-push/subscriptions` — body is `PushSubscription.toJSON()` plus `device_name`. +- `GET /api/v1/notifications/web-push/subscriptions` — for the settings UI device list. +- `DELETE /api/v1/notifications/web-push/subscriptions/{id}` +- `POST /api/v1/notifications/web-push/unsubscribe` — by endpoint (browsers don't know row IDs). + +Subscription endpoints are attacker-controllable URLs the server will POST to, so they pass the same HTTPS + private-destination guard as webhooks, both at registration and at connect time (guarded dialer). + +## Delivery semantics + +- Fanout enqueues one `pending` attempt per enabled subscription of each recipient profile, in the same transaction as the delivery rows. No per-reason filters: profile preferences already gate delivery creation. +- Retry schedule is short (30s/2m/10m/30m, 5 attempts): vendor push services queue messages for offline devices themselves (TTL 12h), so server-side retries only ride out transient push-service errors. +- `404`/`410` from the push service is the protocol's unsubscribe signal: the subscription row is deleted, not retried. +- `notifications.web_push_enabled` is the kill switch (default on). + +## Client + +- `web/public/sw.js` — displays notifications and routes clicks (episode deep link, or the inbox). +- `web/src/lib/webPush.ts` — permission + subscribe/unsubscribe flows. +- Settings → Notifications → "Browser Notifications" — this-browser toggle plus a revocable list of the profile's other subscribed devices. diff --git a/docs/superpowers/plans/notifications/06-v1.5-roadmap.md b/docs/superpowers/plans/notifications/06-v1.5-roadmap.md new file mode 100644 index 00000000..b423ca49 --- /dev/null +++ b/docs/superpowers/plans/notifications/06-v1.5-roadmap.md @@ -0,0 +1,169 @@ +# Notifications v1.5 Roadmap + +**Date:** 2026-06-11 +**Status:** Draft (work not started) +**Scope:** The remaining notification work between the shipped v1 and the deferred v2 push channels. Each item is independent and sized to land as its own PR. +**Depends On:** +- [`00-architecture-overview.md`](./00-architecture-overview.md) +- [`01-release-events-and-inbox.md`](./01-release-events-and-inbox.md) — implemented +- [`04-outbound-webhooks.md`](./04-outbound-webhooks.md) — implemented +- [`05-web-push.md`](./05-web-push.md) — implemented + +## Where v1 landed (context for this doc) + +Implemented 2026-06-11: the full foundation (availability seeding, release events, +interest index, fanout worker with burst caps, websocket channel with ticket +handshake, inbox/sync/preferences/capability APIs, web inbox + badge + settings), +outbound webhooks (Discord + generic HMAC, SSRF guard, durable outbox, retry + +auto-disable), Web Push (VAPID self-provisioned, E2E-encrypted payloads, service +worker), and the shared SMTP core (`internal/mail`, see +`docs/architecture/email.md`) with an admin Email settings page — but no feature +consuming email yet. + +**Deferred to v2 by explicit decision:** APNs (`02`) and FCM (`03`) — they +require Silo-operated relay infrastructure and developer accounts. Also v2 per +the original plans: movie availability, aggregated notifications ("3 new +episodes"), quiet hours, cross-profile views. + +--- + +## 1. Admin settings UI for notification controls + +**Why:** every `notifications.*` setting works today but is reachable only +through the raw admin settings API. Admins should not need `curl` to find the +kill switches. + +**What:** an admin settings page ("Notifications", next to the Email page added +in v1) exposing: + +| Group | Keys | +|---|---| +| Kill switches | `notifications.release_events_enabled`, `notifications.fanout_enabled`, `notifications.ui_enabled`, `notifications.webhooks_enabled`, `notifications.web_push_enabled` | +| Fanout tuning | `notifications.fanout.settle_seconds` (default 30), `notifications.fanout.max_series_burst` (default 3) | +| Webhook guards | `notifications.webhooks.max_per_profile` (10), `notifications.webhooks.allow_private_destinations` (false; dev only — label it loudly), `notifications.webhooks.deliveries_per_minute_per_profile` (60) | +| Retention | `notifications.retention.read_days` (90), `notifications.retention.unread_days` (180), `notifications.retention.event_days` (30) | + +**Files:** add `web/src/pages/admin-settings/NotificationsAdminSettings.tsx` +(follow `EmailSettings.tsx` / `useSettingsForm`), register in +`web/src/pages/admin-settings/AdminSettingsLayout.tsx`. No backend work — all +keys are live-read. + +**Effort:** small (one page, no migrations, no Go changes). + +--- + +## 2. Request-fulfilled notifications (`request.fulfilled`) + +**Why:** `00-architecture-overview.md` calls this "the most obvious next +notification type." Users who request media currently learn it arrived by +checking manually; every delivery channel they configured should tell them. + +**Design:** the `notification_deliveries.type` registry is extensible by +construction — no schema change. + +- New type `request.fulfilled`. `reason_flags` carries the operational shape + (like `webhook.auto_disabled` does), e.g. + `{"request_id": "...", "tmdb_id": 123, "media_type": "movie"}` — never the + four reason booleans. +- **Hook point:** the request reconciliation service (`internal/mediarequests`) + is where a request transitions to available/fulfilled. On that transition, + insert a delivery via `DeliveryRepository.InsertOperational` (the path the + webhook auto-disable notice already uses) and publish through the system's + dispatchers so websocket, web push, and webhooks all fire. +- **Recipient:** the requesting profile (requests are profile-attributed). No + `profile_series_interest` involvement — this is a direct, not fanned-out, + notification. +- **Webhook enqueue:** operational inserts bypass the fanout outbox, so either + (a) extend `InsertOperational` to optionally enqueue per-target attempt rows, + or (b) add a small shared "dispatch one delivery durably" helper used by both + this and the auto-disable notice. Prefer (b); the auto-disable notice + deliberately skips webhooks (loop guard) but request notices should not. +- **Per-reason preferences:** add nothing in v1.5. The profile master toggle + (`notification_preferences.enabled`) gates it; a dedicated + `notify_requests` flag can come later if users ask. +- **Clients:** the web inbox/toast/web-push renderers fall back to a generic + card for unknown types; add a `request.fulfilled` case with the media title, + poster, and a deep link to the item (or the request page until matched). + +**Effort:** medium-small. The delivery/dispatch machinery all exists. + +--- + +## 3. Email digest channel + +**Why:** first real consumer of `internal/mail`; reaches users who don't keep a +browser open and have no webhook. + +**Open design decisions (resolve before building):** + +- **Account-level, not profile-level.** Email addresses live on `users`; + profiles have none. A digest therefore aggregates across the account's + profiles (group by profile inside the email body). +- **Digest, not per-episode.** Per-episode email is spam at hundreds-of-users + scale and duplicates the realtime channels. Recommend: opt-in daily digest of + unread deliveries, sent by a taskmanager task (reuse the checkpointed + iteration pattern from the interest backfill), with a per-user + enable + cadence setting. +- **Unsubscribe / preference surface:** account settings, not profile + notification preferences. + +**Files (sketch):** `internal/notifications/email_digest.go` (compose from +`DeliveryRepository`, send via `mail.Sender`, branch on +`mail.ErrNotConfigured`), a `taskmanager` task, a small user-settings surface. + +**Effort:** medium. Blocked on the design decisions above, not on plumbing. + +--- + +## 4. Native client adoption (no push required) + +**Why:** the Android and Apple apps gain a full notification experience today — +APNs/FCM only add closed-app wake-ups later. + +Server surfaces ready for clients (`silo-android`, `silo-apple`): + +- `GET /api/v1/notifications` + `unread-count` + read endpoints — inbox UI. +- `GET /api/v1/notifications/sync` — opaque forward cursor for + reconnect/foreground catch-up (this is also the wake-fetch endpoint the v2 + push specs assume, so client work done now is reused). +- `POST /api/v1/events/ws-ticket` + `ticket` query param on `/api/v1/events/ws`, + `notifications` channel — realtime while the app is open. Snapshot on + subscribe hydrates recent unread. +- `GET /api/v1/notifications/capability` — drive setup UI from this, never from + admin settings. +- `GET/PUT /api/v1/notifications/preferences` — per-profile reason toggles. + +**Effort:** client-repo work; the server side is done. Coordinate per the +multi-repo guidance in the repo root `CLAUDE.md`. + +--- + +## 5. Hardening backlog (defer freely) + +- **DB-backed integration tests** from the `01` verification plan: idempotent + availability/event inserts, cross-library delivery dedupe, per-series burst + cap, outbox recovery (pending rows with no dispatch → retry worker sends), + multi-node claim safety. The behaviors shipped and were exercised manually on + dev; they are not yet pinned by automated tests because the repo has no + Postgres test harness for this package. +- **Metrics:** `01` names Prometheus-style counters + (`release_events_suppressed_total`, etc.); v1 ships them as structured log + fields. Revisit when the repo grows a metrics registry — keep the names. +- **Discord embed images** via the `media.discord-cdn-proxy` service + (see `04`, "v1.5 payload"). Requires a new Silo-operated repo/service plus + `webhook_image_signer.go`; v1 deliberately ships text-only embeds so the + user's server origin never reaches Discord. +- **Webhook delivery history endpoint:** `webhook_delivery_attempts` already + has the listing index; a `GET /api/v1/notifications/webhooks/{id}/attempts` + endpoint + UI table would make failures self-diagnosable beyond the + last-failure summary. + +--- + +## Suggested order + +1. Admin settings UI (#1) — smallest, completes operability. +2. Request-fulfilled (#2) — highest product value per effort. +3. Webhook history endpoint (#5, last bullet) — pairs naturally with #1. +4. Email digest (#3) — after its design decisions are made. +5. Native clients (#4) — parallel track in the client repos. 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/go.mod b/go.mod index bb9ec600..8f652015 100644 --- a/go.mod +++ b/go.mod @@ -20,6 +20,7 @@ require ( ) require ( + github.com/SherClockHolmes/webpush-go v1.4.0 github.com/abadojack/whatlanggo v1.0.1 github.com/go-chi/cors v1.2.2 github.com/gorilla/websocket v1.5.3 @@ -30,6 +31,7 @@ require ( github.com/oklog/ulid/v2 v2.1.0 github.com/pgvector/pgvector-go v0.3.0 github.com/pressly/goose/v3 v3.27.1 + github.com/wneessen/go-mail v0.7.3 github.com/zishang520/socket.io/v2 v2.5.0 go.n16f.net/thumbhash v1.1.0 golang.org/x/image v0.41.0 diff --git a/go.sum b/go.sum index 58843ee5..23fc5643 100644 --- a/go.sum +++ b/go.sum @@ -2,6 +2,8 @@ entgo.io/ent v0.14.3 h1:wokAV/kIlH9TeklJWGGS7AYJdVckr0DloWjIcO9iIIQ= entgo.io/ent v0.14.3/go.mod h1:aDPE/OziPEu8+OWbzy4UlvWmD2/kbRuWfK2A40hcxJM= github.com/PuerkitoBio/goquery v1.8.0 h1:PJTF7AmFCFKk1N6V6jmKfrNH9tV5pNE6lZMkG0gta/U= github.com/PuerkitoBio/goquery v1.8.0/go.mod h1:ypIiRMtY7COPGk+I/YbZLbxsxn9g5ejnI2HSMtkjZvI= +github.com/SherClockHolmes/webpush-go v1.4.0 h1:ocnzNKWN23T9nvHi6IfyrQjkIc0oJWv1B1pULsf9i3s= +github.com/SherClockHolmes/webpush-go v1.4.0/go.mod h1:XSq8pKX11vNV8MJEMwjrlTkxhAj1zKfxmyhdV7Pd6UA= github.com/Silo-Server/silo-plugin-sdk v0.6.0 h1:Gi9TdH9kt7b8X4xRXH493/nSYb9n0GO4VCWmlll0hKI= github.com/Silo-Server/silo-plugin-sdk v0.6.0/go.mod h1:etqmxLTwjxpFH9goAjBDfNDoqHMv2/sqUXu8yx3hNfA= github.com/abadojack/whatlanggo v1.0.1 h1:19N6YogDnf71CTHm3Mp2qhYfkRdyvbgwWdd2EPxJRG4= @@ -69,10 +71,12 @@ github.com/go-pg/pg/v10 v10.11.0 h1:CMKJqLgTrfpE/aOVeLdybezR2om071Vh38OLZjsyMI0= github.com/go-pg/pg/v10 v10.11.0/go.mod h1:4BpHRoxE61y4Onpof3x1a2SQvi9c+q1dJnrNdMjsroA= github.com/go-pg/zerochecker v0.2.0 h1:pp7f72c3DobMWOb2ErtZsnrPaSvHd2W4o9//8HtF4mU= github.com/go-pg/zerochecker v0.2.0/go.mod h1:NJZ4wKL0NmTtz0GKCoJ8kym6Xn/EQzXRl2OnAe7MmDo= +github.com/golang-jwt/jwt/v5 v5.2.1/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk= github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY= github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= @@ -209,12 +213,15 @@ github.com/vmihailenco/tagparser v0.1.2 h1:gnjoVuB/kljJ5wICEEOpx98oXMWPLj22G67Vb github.com/vmihailenco/tagparser v0.1.2/go.mod h1:OeAg3pn3UbLjkWt+rN9oFYB6u/cQgqMEUPoW2WPyhdI= github.com/vmihailenco/tagparser/v2 v2.0.0 h1:y09buUbR+b5aycVFQs/g70pqKVZNBmxwAhO7/IwNM9g= github.com/vmihailenco/tagparser/v2 v2.0.0/go.mod h1:Wri+At7QHww0WTrCBeu4J6bNtoV6mEfg5OIWRZA9qds= +github.com/wneessen/go-mail v0.7.3 h1:g3DravXC5SMlVdboFrQA8Jx95A8sOzoBeS5F+vzNRK0= +github.com/wneessen/go-mail v0.7.3/go.mod h1:QGhBX0yNbc1J+Mkjcu7z2rpj4B4l+BmDY8gYznPC9sk= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/xo/terminfo v0.0.0-20210125001918-ca9a967f8778 h1:QldyIu/L63oPpyvQmHgvgickp1Yw510KJOqX7H24mg8= github.com/xo/terminfo v0.0.0-20210125001918-ca9a967f8778/go.mod h1:2MuV+tbUrU1zIOPMxZ5EncGwgmMJsa+9ucAQZXxsObs= github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU= github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= +github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0= github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA= github.com/zishang520/engine.io-go-parser v1.3.2 h1:aEVrhQVhfk99Ct6htNffgHydUBC4dGclO/OXPz5CSy0= @@ -249,31 +256,90 @@ go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI= go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU= +golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= +golang.org/x/crypto v0.13.0/go.mod h1:y6Z2r+Rw4iayiXXAIxJIDAJ1zMW4yaTpebo8fPOliYc= +golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU= +golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8= +golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk= golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988= golang.org/x/crypto v0.52.0/go.mod h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc= golang.org/x/image v0.41.0 h1:8wS72eGJMJaBxK6okTzd4WaXumUlTVlb753MlsSvTCo= golang.org/x/image v0.41.0/go.mod h1:uIc348UZMSvS5Z65CVZ7iDPaNobNFEPeJ4kbqTOszmA= +golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= +golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= +golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= +golang.org/x/mod v0.15.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= +golang.org/x/mod v0.17.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= +golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20210916014120-12bc252f5db8/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= +golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= +golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= +golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg= +golang.org/x/net v0.15.0/go.mod h1:idbUs1IY1+zTqbi8yxTbhexhEEk5ur9LInksu6HrEpk= +golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44= +golang.org/x/net v0.25.0/go.mod h1:JkAGAh7GEvH74S6FOH42FLoXpXbE/aqXSrIQjXgsiwM= golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8= golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww= +golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.3.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y= +golang.org/x/sync v0.6.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20200116001909-b77594299b42/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200223170610-d5e6a3e2c0ae/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210630005230-0f9fa26af87c/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210927094055-39ccf1dd6fa6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220503163025-988cb79eb6c6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.28.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= +golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= +golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= +golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo= +golang.org/x/term v0.12.0/go.mod h1:owVbMEjm3cBLCHdkQu9b1opXd4ETQWc3BhuQGKgXgvU= +golang.org/x/term v0.17.0/go.mod h1:lLRBjIVuehSbZlaOtGMbcMncT+aqLLLmKrsjNrUguwk= +golang.org/x/term v0.20.0/go.mod h1:8UkIAJTvZgivsXaD6/pH6U9ecQzZ45awqEOzuCvwpFY= +golang.org/x/term v0.27.0/go.mod h1:iMsnZpn0cago0GOrHO2+Y7u7JPn5AylBrcoWkElMTSM= +golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= +golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= +golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= +golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= +golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8= +golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE= +golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= +golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= +golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ= golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= golang.org/x/time v0.14.0 h1:MRx4UaLrDotUKUdCIqzPC48t1Y9hANFKIRpNx+Te8PI= golang.org/x/time v0.14.0/go.mod h1:eL/Oa2bBBK0TkX57Fyni+NgnyQQN4LitPmob2Hjnqw4= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= +golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= +golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= +golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58= +golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk= +golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= google.golang.org/genproto/googleapis/rpc v0.0.0-20260420184626-e10c466a9529 h1:XF8+t6QQiS0o9ArVan/HW8Q7cycNPGsJf6GA2nXxYAg= diff --git a/internal/api/handlers/admin_server_channels.go b/internal/api/handlers/admin_server_channels.go new file mode 100644 index 00000000..788c7ba2 --- /dev/null +++ b/internal/api/handlers/admin_server_channels.go @@ -0,0 +1,232 @@ +package handlers + +import ( + "encoding/json" + "errors" + "net/http" + "time" + + apimw "github.com/Silo-Server/silo-server/internal/api/middleware" + "github.com/Silo-Server/silo-server/internal/notifications" + "github.com/go-chi/chi/v5" +) + +// AdminServerChannelsHandler exposes admin CRUD for server notification +// channels (community broadcast destinations). All routes are mounted inside +// the admin-only group. +type AdminServerChannelsHandler struct { + system *notifications.System +} + +// NewAdminServerChannelsHandler creates the handler. +func NewAdminServerChannelsHandler(system *notifications.System) *AdminServerChannelsHandler { + return &AdminServerChannelsHandler{system: system} +} + +func (h *AdminServerChannelsHandler) service() *notifications.ServerChannelService { + if h == nil || h.system == nil { + return nil + } + return h.system.ServerChannels +} + +// serverChannelResponse is the API view of a server channel. Like webhooks it +// never includes the destination URL (Discord webhook tokens are bearer +// credentials in the URL path) or the stored signing secret — only url_host. +type serverChannelResponse struct { + ID string `json:"id"` + Name string `json:"name"` + Type string `json:"type"` + URLHost string `json:"url_host"` + Enabled bool `json:"enabled"` + NotifyNewMovies bool `json:"notify_new_movies"` + NotifyNewEpisodes bool `json:"notify_new_episodes"` + NotifyRequestSubmitted bool `json:"notify_request_submitted"` + NotifyRequestApproved bool `json:"notify_request_approved"` + NotifyRequestDeclined bool `json:"notify_request_declined"` + NotifyRequestFulfilled bool `json:"notify_request_fulfilled"` + ConsecutiveFailures int `json:"consecutive_failures"` + DisabledReason *string `json:"disabled_reason"` + LastSuccessAt *time.Time `json:"last_success_at"` + LastFailureAt *time.Time `json:"last_failure_at"` + LastFailureStatus *int `json:"last_failure_status"` + LastFailureMessage *string `json:"last_failure_message"` + CreatedAt time.Time `json:"created_at"` + // SigningSecret is present only in create / rotate-secret responses. + SigningSecret string `json:"signing_secret,omitempty"` +} + +func serverChannelToResponse(ch notifications.ServerChannel) serverChannelResponse { + return serverChannelResponse{ + ID: ch.ID, + Name: ch.Name, + Type: ch.Type, + URLHost: ch.URLHost, + Enabled: ch.Enabled, + NotifyNewMovies: ch.NotifyNewMovies, + NotifyNewEpisodes: ch.NotifyNewEpisodes, + NotifyRequestSubmitted: ch.NotifyRequestSubmitted, + NotifyRequestApproved: ch.NotifyRequestApproved, + NotifyRequestDeclined: ch.NotifyRequestDeclined, + NotifyRequestFulfilled: ch.NotifyRequestFulfilled, + ConsecutiveFailures: ch.ConsecutiveFailures, + DisabledReason: ch.DisabledReason, + LastSuccessAt: ch.LastSuccessAt, + LastFailureAt: ch.LastFailureAt, + LastFailureStatus: ch.LastFailureStatus, + LastFailureMessage: ch.LastFailureMessage, + CreatedAt: ch.CreatedAt, + } +} + +type serverChannelRequest struct { + Name *string `json:"name"` + URL *string `json:"url"` + Type *string `json:"type"` + Enabled *bool `json:"enabled"` + NotifyNewMovies *bool `json:"notify_new_movies"` + NotifyNewEpisodes *bool `json:"notify_new_episodes"` + NotifyRequestSubmitted *bool `json:"notify_request_submitted"` + NotifyRequestApproved *bool `json:"notify_request_approved"` + NotifyRequestDeclined *bool `json:"notify_request_declined"` + NotifyRequestFulfilled *bool `json:"notify_request_fulfilled"` +} + +func (r serverChannelRequest) toInput() notifications.ServerChannelInput { + return notifications.ServerChannelInput{ + Name: r.Name, + URL: r.URL, + Type: r.Type, + Enabled: r.Enabled, + NotifyNewMovies: r.NotifyNewMovies, + NotifyNewEpisodes: r.NotifyNewEpisodes, + NotifyRequestSubmitted: r.NotifyRequestSubmitted, + NotifyRequestApproved: r.NotifyRequestApproved, + NotifyRequestDeclined: r.NotifyRequestDeclined, + NotifyRequestFulfilled: r.NotifyRequestFulfilled, + } +} + +func writeServerChannelError(w http.ResponseWriter, err error) { + switch { + case errors.Is(err, notifications.ErrServerChannelsDisabled): + writeError(w, http.StatusForbidden, "server_channels_disabled", "Server channels are disabled") + case errors.Is(err, notifications.ErrServerChannelNotFound): + writeError(w, http.StatusNotFound, "not_found", "Server channel not found") + case errors.Is(err, notifications.ErrServerChannelLimit): + writeError(w, http.StatusUnprocessableEntity, "limit_reached", "Server channel limit reached") + case errors.Is(err, notifications.ErrServerChannelInvalid): + writeError(w, http.StatusBadRequest, "bad_request", err.Error()) + default: + writeError(w, http.StatusInternalServerError, "internal_error", "Server channel operation failed") + } +} + +// HandleList handles GET /admin/notifications/server-channels. +func (h *AdminServerChannelsHandler) HandleList(w http.ResponseWriter, r *http.Request) { + service := h.service() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Server channels are not available") + return + } + channels, err := service.List(r.Context()) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to list server channels") + return + } + responses := make([]serverChannelResponse, 0, len(channels)) + for _, ch := range channels { + responses = append(responses, serverChannelToResponse(ch)) + } + writeJSON(w, http.StatusOK, map[string]any{"channels": responses}) +} + +// HandleCreate handles POST /admin/notifications/server-channels. For generic +// channels the response carries the signing secret exactly once. +func (h *AdminServerChannelsHandler) HandleCreate(w http.ResponseWriter, r *http.Request) { + service := h.service() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Server channels are not available") + return + } + var req serverChannelRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + ch, signingSecret, err := service.Create(r.Context(), apimw.GetUserID(r.Context()), req.toInput()) + if err != nil { + writeServerChannelError(w, err) + return + } + response := serverChannelToResponse(*ch) + response.SigningSecret = signingSecret + writeJSON(w, http.StatusCreated, response) +} + +// HandleUpdate handles PUT /admin/notifications/server-channels/{id}. +func (h *AdminServerChannelsHandler) HandleUpdate(w http.ResponseWriter, r *http.Request) { + service := h.service() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Server channels are not available") + return + } + var req serverChannelRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + ch, err := service.Update(r.Context(), chi.URLParam(r, "id"), req.toInput()) + if err != nil { + writeServerChannelError(w, err) + return + } + writeJSON(w, http.StatusOK, serverChannelToResponse(*ch)) +} + +// HandleDelete handles DELETE /admin/notifications/server-channels/{id}. +// Idempotent. +func (h *AdminServerChannelsHandler) HandleDelete(w http.ResponseWriter, r *http.Request) { + service := h.service() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Server channels are not available") + return + } + if err := service.Delete(r.Context(), chi.URLParam(r, "id")); err != nil { + writeServerChannelError(w, err) + return + } + w.WriteHeader(http.StatusNoContent) +} + +// HandleRotateSecret handles POST /admin/notifications/server-channels/{id}/rotate-secret. +func (h *AdminServerChannelsHandler) HandleRotateSecret(w http.ResponseWriter, r *http.Request) { + service := h.service() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Server channels are not available") + return + } + signingSecret, err := service.RotateSecret(r.Context(), chi.URLParam(r, "id")) + if err != nil { + writeServerChannelError(w, err) + return + } + writeJSON(w, http.StatusOK, map[string]string{"signing_secret": signingSecret}) +} + +// HandleTest handles POST /admin/notifications/server-channels/{id}/test. The +// test send is synchronous and never touches the watermark or failure +// counters. +func (h *AdminServerChannelsHandler) HandleTest(w http.ResponseWriter, r *http.Request) { + service := h.service() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Server channels are not available") + return + } + result, err := service.Test(r.Context(), chi.URLParam(r, "id")) + if err != nil { + writeServerChannelError(w, err) + return + } + writeJSON(w, http.StatusOK, result) +} diff --git a/internal/api/handlers/admin_server_channels_test.go b/internal/api/handlers/admin_server_channels_test.go new file mode 100644 index 00000000..9e75283c --- /dev/null +++ b/internal/api/handlers/admin_server_channels_test.go @@ -0,0 +1,82 @@ +package handlers + +import ( + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/Silo-Server/silo-server/internal/notifications" +) + +func TestWriteServerChannelErrorMapping(t *testing.T) { + cases := []struct { + err error + want int + }{ + {notifications.ErrServerChannelsDisabled, http.StatusForbidden}, + {notifications.ErrServerChannelNotFound, http.StatusNotFound}, + {notifications.ErrServerChannelLimit, http.StatusUnprocessableEntity}, + {fmt.Errorf("%w: name is required", notifications.ErrServerChannelInvalid), http.StatusBadRequest}, + {fmt.Errorf("boom"), http.StatusInternalServerError}, + } + for _, tc := range cases { + rec := httptest.NewRecorder() + writeServerChannelError(rec, tc.err) + if rec.Code != tc.want { + t.Errorf("writeServerChannelError(%v) = %d, want %d", tc.err, rec.Code, tc.want) + } + } +} + +// The response shape must never leak the destination URL (a bearer credential +// for Discord webhooks) or the stored signing secret; signing_secret appears +// only when explicitly set by the create/rotate paths. +func TestServerChannelResponseNeverLeaksSecrets(t *testing.T) { + secret := "ciphertext-secret" + ch := notifications.ServerChannel{ + ID: "ch-1", + Name: "Community", + Type: "discord", + URLCiphertext: "enc:v1:secret-url-material", + URLHost: "discord.com", + SigningSecretCiphertext: &secret, + Enabled: true, + } + body, err := json.Marshal(serverChannelToResponse(ch)) + if err != nil { + t.Fatal(err) + } + encoded := string(body) + if strings.Contains(encoded, "secret-url-material") || strings.Contains(encoded, "ciphertext-secret") { + t.Fatalf("response leaked ciphertext material: %s", encoded) + } + if strings.Contains(encoded, "signing_secret") { + t.Fatalf("signing_secret must be omitted unless explicitly set: %s", encoded) + } + if !strings.Contains(encoded, `"url_host":"discord.com"`) { + t.Fatalf("url_host missing from response: %s", encoded) + } + + response := serverChannelToResponse(ch) + response.SigningSecret = "shown-once" + body, err = json.Marshal(response) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(string(body), `"signing_secret":"shown-once"`) { + t.Fatalf("explicit signing_secret missing: %s", body) + } +} + +// A handler with no notifications system reports 503 instead of panicking. +func TestServerChannelHandlersUnavailableWithoutSystem(t *testing.T) { + h := NewAdminServerChannelsHandler(nil) + rec := httptest.NewRecorder() + h.HandleList(rec, httptest.NewRequest(http.MethodGet, "/admin/notifications/server-channels", nil)) + if rec.Code != http.StatusServiceUnavailable { + t.Fatalf("HandleList without system = %d, want 503", rec.Code) + } +} diff --git a/internal/api/handlers/email.go b/internal/api/handlers/email.go new file mode 100644 index 00000000..3c0a7373 --- /dev/null +++ b/internal/api/handlers/email.go @@ -0,0 +1,79 @@ +package handlers + +import ( + "encoding/json" + "errors" + "net/http" + "net/mail" + "time" + + silomail "github.com/Silo-Server/silo-server/internal/mail" +) + +// EmailHandler exposes admin operations for the shared outbound email +// facility (internal/mail). Feature-specific email content lives with the +// features; this handler only owns configuration verification. +type EmailHandler struct { + sender silomail.Sender +} + +// NewEmailHandler creates an EmailHandler. +func NewEmailHandler(sender silomail.Sender) *EmailHandler { + return &EmailHandler{sender: sender} +} + +type emailTestRequest struct { + To string `json:"to"` +} + +type emailTestResponse struct { + OK bool `json:"ok"` + DurationMS int64 `json:"duration_ms"` + Message string `json:"message,omitempty"` +} + +// HandleTest handles POST /admin/email/test: synchronously sends a test +// message so admins can verify SMTP settings before any feature depends on +// them. +func (h *EmailHandler) HandleTest(w http.ResponseWriter, r *http.Request) { + if h == nil || h.sender == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Email is not available") + return + } + var req emailTestRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + if _, err := mail.ParseAddress(req.To); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "A valid recipient address is required") + return + } + + started := time.Now() + err := h.sender.Send(r.Context(), silomail.Message{ + To: []string{req.To}, + Subject: "Silo test email", + TextBody: "This is a test email from your Silo server.\n\n" + + "If you received it, outbound email is configured correctly.", + HTMLBody: silomail.RenderLayout(silomail.LayoutOptions{ + Preheader: "Outbound email from your Silo server is configured correctly.", + Title: "Outbound email is working", + BodyHTML: silomail.EmailParagraph("This is a test email from your Silo server.") + + silomail.EmailParagraph("If you're reading it, the SMTP settings are correct and "+ + "notification emails will look like this one."), + }), + }) + response := emailTestResponse{ + OK: err == nil, + DurationMS: time.Since(started).Milliseconds(), + } + switch { + case err == nil: + case errors.Is(err, silomail.ErrNotConfigured): + response.Message = "Email is not configured. Set the SMTP host, from address, and enable email first." + default: + response.Message = err.Error() + } + writeJSON(w, http.StatusOK, response) +} diff --git a/internal/api/handlers/events_ws.go b/internal/api/handlers/events_ws.go index a6b668ff..cb390929 100644 --- a/internal/api/handlers/events_ws.go +++ b/internal/api/handlers/events_ws.go @@ -12,6 +12,7 @@ import ( "github.com/Silo-Server/silo-server/internal/auth" evt "github.com/Silo-Server/silo-server/internal/events" "github.com/Silo-Server/silo-server/internal/historyimport" + "github.com/Silo-Server/silo-server/internal/notifications" "github.com/Silo-Server/silo-server/internal/scanqueue" "github.com/Silo-Server/silo-server/internal/taskmanager" "github.com/gorilla/websocket" @@ -39,6 +40,15 @@ type EventsHandler struct { scans *evt.ScanRegistry persistedScans activeScanLister historyImports historyImportActiveLister + notifications *notifications.System +} + +// SetNotificationsSystem wires the user-notification system: websocket +// handshake tickets and the notifications channel snapshot. +func (h *EventsHandler) SetNotificationsSystem(system *notifications.System) { + if h != nil { + h.notifications = system + } } func NewEventsHandler( @@ -73,6 +83,27 @@ func (h *EventsHandler) HandleWebSocket(w http.ResponseWriter, r *http.Request) return } + // Browsers cannot set custom headers on websocket handshakes, so profile + // identity arrives as a short-lived single-use ticket minted via + // POST /events/ws-ticket. A connection without a ticket stays unbound and + // simply cannot subscribe to the profile-scoped notifications channel. + boundProfileID := "" + if ticket := r.URL.Query().Get("ticket"); ticket != "" && h.notifications != nil { + ticketUserID, ticketProfileID, ok := h.notifications.Tickets.Consume(r.Context(), ticket) + if ok && ticketUserID == claims.UserID { + boundProfileID = ticketProfileID + } else { + // Expired, reused, consumed on a different node, or minted for + // another user: degrade to an unbound connection instead of + // failing the handshake. The binding grants nothing on its own — + // the client retries it when its notifications subscription is + // rejected — whereas a hard 403 would take down every realtime + // channel over a notifications-only concern. + slog.Warn("events: websocket ticket rejected; connection unbound", + "user_id", claims.UserID) + } + } + conn, err := wsUpgrader.Upgrade(w, r, nil) if err != nil { return @@ -142,7 +173,7 @@ func (h *EventsHandler) HandleWebSocket(w http.ResponseWriter, r *http.Request) ) return case data := <-readMessages: - nextSubs, handled, ok := h.handleEventsClientMessage(conn, r, claims, data, allowedChannels) + nextSubs, handled, ok := h.handleEventsClientMessage(conn, r, claims, boundProfileID, data, allowedChannels) if !ok { return } @@ -164,10 +195,10 @@ func (h *EventsHandler) HandleWebSocket(w http.ResponseWriter, r *http.Request) if _, subscribed := subscriptions[env.Channel]; !subscribed { continue } - if !allowsEventForClaims(claims, env) { + if !allowsEventForClaims(claims, boundProfileID, env) { continue } - if err := h.writeEventFrame(conn, r, claims, env); err != nil { + if err := h.writeEventFrame(conn, r, claims, boundProfileID, env); err != nil { return } } @@ -178,6 +209,7 @@ func (h *EventsHandler) handleEventsClientMessage( conn *websocket.Conn, r *http.Request, claims *auth.Claims, + boundProfileID string, data []byte, allowed []evt.EventChannel, ) (map[evt.EventChannel]struct{}, bool, bool) { @@ -245,6 +277,16 @@ func (h *EventsHandler) handleEventsClientMessage( }) continue } + // The notifications channel is profile-scoped: it requires a + // connection bound to a profile via a websocket ticket. + if channel == evt.ChannelNotifications && boundProfileID == "" { + rejected = append(rejected, evt.EventsRejectedChannel{ + Channel: channel, + Code: "profile_required", + Message: "A profile-bound websocket ticket is required", + }) + continue + } if _, seen := nextSubs[channel]; seen { continue } @@ -262,7 +304,7 @@ func (h *EventsHandler) handleEventsClientMessage( } for _, channel := range accepted { - if err := h.writeSnapshotFrame(conn, r, claims, channel); err != nil { + if err := h.writeSnapshotFrame(conn, r, claims, boundProfileID, channel); err != nil { return nil, false, false } } @@ -275,6 +317,7 @@ func allowedChannelsForRole(role string) []evt.EventChannel { evt.ChannelCatalog, evt.ChannelHistoryImport, evt.ChannelUserState, + evt.ChannelNotifications, } if role == "admin" { channels = append(channels, @@ -287,13 +330,20 @@ func allowedChannelsForRole(role string) []evt.EventChannel { return channels } -func allowsEventForClaims(claims *auth.Claims, env evt.Envelope) bool { +func allowsEventForClaims(claims *auth.Claims, boundProfileID string, env evt.Envelope) bool { if claims == nil { return false } if env.AdminOnly && claims.Role != "admin" { return false } + if env.Channel == evt.ChannelNotifications { + // Notifications are personal: even admins only receive their own + // profile's deliveries, and only on a profile-bound connection. + return boundProfileID != "" && + env.UserID == claims.UserID && + env.ProfileID == boundProfileID + } if env.UserID > 0 && claims.Role != "admin" && env.UserID != claims.UserID { return false } @@ -314,11 +364,24 @@ func marshalJSON(value any) json.RawMessage { func (h *EventsHandler) snapshotForChannel( r *http.Request, claims *auth.Claims, + boundProfileID string, channel evt.EventChannel, ) (json.RawMessage, error) { switch channel { case evt.ChannelCatalog, evt.ChannelUserState: return json.RawMessage("null"), nil + case evt.ChannelNotifications: + // Recent unread deliveries for the bound profile so reconnecting + // clients hydrate without a separate REST call. Same row shape as the + // inbox list API. + if h == nil || h.notifications == nil || boundProfileID == "" { + return json.RawMessage("[]"), nil + } + rows, err := h.notifications.Deliveries.RecentUnread(r.Context(), boundProfileID, 25) + if err != nil { + return nil, err + } + return marshalJSON(h.notifications.PayloadsForRows(r.Context(), rows)), nil case evt.ChannelJobs: if h == nil || h.jobs == nil || h.jobs.repo == nil { return json.RawMessage("[]"), nil @@ -387,9 +450,10 @@ func (h *EventsHandler) writeSnapshotFrame( conn *websocket.Conn, r *http.Request, claims *auth.Claims, + boundProfileID string, channel evt.EventChannel, ) error { - data, err := h.snapshotForChannel(r, claims, channel) + data, err := h.snapshotForChannel(r, claims, boundProfileID, channel) if err != nil { slog.Error( "events: failed to build initial snapshot", @@ -415,14 +479,19 @@ func (h *EventsHandler) writeEventFrame( conn *websocket.Conn, r *http.Request, claims *auth.Claims, + boundProfileID string, env evt.Envelope, ) error { data := env.Data if len(data) == 0 || (env.Channel == evt.ChannelSessions && env.Event == "sessions.replaced") { - snapshot, err := h.snapshotForChannel(r, claims, env.Channel) + snapshot, err := h.snapshotForChannel(r, claims, boundProfileID, env.Channel) if err != nil { + // Drop the frame but keep the stream open (same contract as + // writeSnapshotFrame): durable state covers the gap on the next + // event or reconnect, while closing the socket tears down every + // channel the client subscribed to. slog.Error("events: failed to build event payload", "channel", env.Channel, "event", env.Event, "error", err) - return err + return nil } data = snapshot } diff --git a/internal/api/handlers/notifications.go b/internal/api/handlers/notifications.go new file mode 100644 index 00000000..7a4d351c --- /dev/null +++ b/internal/api/handlers/notifications.go @@ -0,0 +1,395 @@ +package handlers + +import ( + "encoding/json" + "net/http" + "strconv" + + apimw "github.com/Silo-Server/silo-server/internal/api/middleware" + evt "github.com/Silo-Server/silo-server/internal/events" + "github.com/Silo-Server/silo-server/internal/notifications" + "github.com/go-chi/chi/v5" +) + +const ( + notificationsDefaultLimit = 25 + notificationsMaxLimit = 100 + notificationsSyncLimit = 50 +) + +// NotificationsHandler serves the profile-scoped notification inbox, +// preferences, capability, and websocket-ticket endpoints. All routes are +// mounted behind RequireProfile. +type NotificationsHandler struct { + system *notifications.System + hub *evt.Hub +} + +// NewNotificationsHandler creates a NotificationsHandler. +func NewNotificationsHandler(system *notifications.System, hub *evt.Hub) *NotificationsHandler { + return &NotificationsHandler{system: system, hub: hub} +} + +type notificationListResponse struct { + Notifications []notifications.DeliveryRowPayload `json:"notifications"` + // NextCursor pages further into the past via the `before` query param. + // Empty when this page may be the last. + NextCursor string `json:"next_cursor,omitempty"` +} + +type notificationSyncResponse struct { + Notifications []notifications.DeliveryRowPayload `json:"notifications"` + NextCursor string `json:"next_cursor,omitempty"` + UnreadCount int `json:"unread_count"` +} + +type unreadCountResponse struct { + Count int `json:"count"` +} + +type wsTicketResponse struct { + Ticket string `json:"ticket"` + ExpiresIn int `json:"expires_in"` +} + +func parseNotificationsLimit(r *http.Request, fallback int) int { + raw := r.URL.Query().Get("limit") + if raw == "" { + return fallback + } + limit, err := strconv.Atoi(raw) + if err != nil || limit <= 0 { + return fallback + } + return min(limit, notificationsMaxLimit) +} + +// HandleList handles GET /notifications (newest-first inbox page). +func (h *NotificationsHandler) HandleList(w http.ResponseWriter, r *http.Request) { + profileID := apimw.GetProfileID(r.Context()) + unreadOnly := r.URL.Query().Get("status") == "unread" + limit := parseNotificationsLimit(r, notificationsDefaultLimit) + + var before *notifications.Cursor + if raw := r.URL.Query().Get("before"); raw != "" { + cursor, err := notifications.DecodeCursor(raw) + if err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid before cursor") + return + } + before = &cursor + } + + rows, err := h.system.Deliveries.ListInbox(r.Context(), profileID, unreadOnly, limit, before) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to list notifications") + return + } + response := notificationListResponse{Notifications: h.system.PayloadsForRows(r.Context(), rows)} + if len(rows) == limit { + last := rows[len(rows)-1] + response.NextCursor = notifications.Cursor{CreatedAt: last.CreatedAt, ID: last.ID}.Encode() + } + writeJSON(w, http.StatusOK, response) +} + +// HandleSync handles GET /notifications/sync — the forward (ascending) cursor +// sync used by clients waking from a push or reconnecting after a gap. +func (h *NotificationsHandler) HandleSync(w http.ResponseWriter, r *http.Request) { + profileID := apimw.GetProfileID(r.Context()) + limit := parseNotificationsLimit(r, notificationsSyncLimit) + + var since *notifications.Cursor + if raw := r.URL.Query().Get("since"); raw != "" { + cursor, err := notifications.DecodeCursor(raw) + if err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid since cursor") + return + } + since = &cursor + } + + rows, err := h.system.Deliveries.ListSync(r.Context(), profileID, since, limit) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to sync notifications") + return + } + unread, err := h.system.Deliveries.UnreadCount(r.Context(), profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to count unread notifications") + return + } + + response := notificationSyncResponse{ + Notifications: h.system.PayloadsForRows(r.Context(), rows), + UnreadCount: unread, + } + if len(rows) > 0 { + last := rows[len(rows)-1] + response.NextCursor = notifications.Cursor{CreatedAt: last.CreatedAt, ID: last.ID}.Encode() + } else if since != nil { + response.NextCursor = since.Encode() + } + writeJSON(w, http.StatusOK, response) +} + +// HandleGet handles GET /notifications/{id}; 404 for other profiles' rows. +func (h *NotificationsHandler) HandleGet(w http.ResponseWriter, r *http.Request) { + profileID := apimw.GetProfileID(r.Context()) + id := chi.URLParam(r, "id") + + row, err := h.system.Deliveries.GetByID(r.Context(), profileID, id) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to load notification") + return + } + if row == nil { + writeError(w, http.StatusNotFound, "not_found", "Notification not found") + return + } + writeJSON(w, http.StatusOK, h.system.PayloadForRow(r.Context(), *row)) +} + +// HandleUnreadCount handles GET /notifications/unread-count. +func (h *NotificationsHandler) HandleUnreadCount(w http.ResponseWriter, r *http.Request) { + profileID := apimw.GetProfileID(r.Context()) + count, err := h.system.Deliveries.UnreadCount(r.Context(), profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to count unread notifications") + return + } + writeJSON(w, http.StatusOK, unreadCountResponse{Count: count}) +} + +// HandleMarkRead handles POST /notifications/{id}/read. Idempotent. +func (h *NotificationsHandler) HandleMarkRead(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + id := chi.URLParam(r, "id") + + transitioned, err := h.system.Deliveries.MarkRead(r.Context(), profileID, id) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to mark notification read") + return + } + if !transitioned { + // Already read is fine (idempotent); unknown IDs are a 404. + exists, err := h.system.Deliveries.Exists(r.Context(), profileID, id) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to mark notification read") + return + } + if !exists { + writeError(w, http.StatusNotFound, "not_found", "Notification not found") + return + } + } + if transitioned { + h.publishReadEvent(r, userID, profileID, id) + } + w.WriteHeader(http.StatusNoContent) +} + +// HandleReadAll handles POST /notifications/read-all. +func (h *NotificationsHandler) HandleReadAll(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + + if _, err := h.system.Deliveries.MarkAllRead(r.Context(), profileID); err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to mark notifications read") + return + } + h.publishReadEvent(r, userID, profileID, "") + w.WriteHeader(http.StatusNoContent) +} + +// publishReadEvent lets other connected tabs of the same profile reconcile +// read state. An empty id means "all read". +func (h *NotificationsHandler) publishReadEvent(r *http.Request, userID int, profileID, id string) { + if h.hub == nil { + return + } + payload := map[string]any{"profile_id": profileID} + if id != "" { + payload["id"] = id + } else { + payload["all"] = true + } + _ = h.hub.PublishJSON(r.Context(), evt.ChannelNotifications, notifications.EventNotificationRead, + payload, evt.PublishOptions{UserID: userID, ProfileID: profileID}) +} + +// HandleGetPreferences handles GET /notifications/preferences. +func (h *NotificationsHandler) HandleGetPreferences(w http.ResponseWriter, r *http.Request) { + profileID := apimw.GetProfileID(r.Context()) + prefs, err := h.system.Preferences.Get(r.Context(), profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to load notification preferences") + return + } + writeJSON(w, http.StatusOK, prefs) +} + +type updatePreferencesRequest struct { + Enabled *bool `json:"enabled"` + NotifyFavorites *bool `json:"notify_favorites"` + NotifyWatchlist *bool `json:"notify_watchlist"` + NotifyContinueWatching *bool `json:"notify_continue_watching"` + NotifyNextUp *bool `json:"notify_next_up"` +} + +// HandleUpdatePreferences handles PUT /notifications/preferences. Fields are +// optional; omitted fields keep their current value. Idempotent. +func (h *NotificationsHandler) HandleUpdatePreferences(w http.ResponseWriter, r *http.Request) { + profileID := apimw.GetProfileID(r.Context()) + + var req updatePreferencesRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + + prefs, err := h.system.Preferences.Get(r.Context(), profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to load notification preferences") + return + } + if req.Enabled != nil { + prefs.Enabled = *req.Enabled + } + if req.NotifyFavorites != nil { + prefs.NotifyFavorites = *req.NotifyFavorites + } + if req.NotifyWatchlist != nil { + prefs.NotifyWatchlist = *req.NotifyWatchlist + } + if req.NotifyContinueWatching != nil { + prefs.NotifyContinueWatching = *req.NotifyContinueWatching + } + if req.NotifyNextUp != nil { + prefs.NotifyNextUp = *req.NotifyNextUp + } + if err := h.system.Preferences.Upsert(r.Context(), prefs); err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to save notification preferences") + return + } + writeJSON(w, http.StatusOK, prefs) +} + +type capabilityResponse struct { + InApp capabilityInApp `json:"in_app"` + ApplePush capabilityPush `json:"apple_push"` + AndroidPush capabilityPush `json:"android_push"` + WebPush capabilityWebPush `json:"web_push"` + Webhooks capabilityWebhooks `json:"webhooks"` + Email capabilityAccountChannel `json:"email"` + Discord capabilityAccountChannel `json:"discord"` +} + +// capabilityAccountChannel describes an account-level digest channel (email, +// Discord DMs). +type capabilityAccountChannel 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 { + Available bool `json:"available"` + PublicKey string `json:"public_key,omitempty"` +} + +type capabilityInApp struct { + Enabled bool `json:"enabled"` +} + +type capabilityPush struct { + Available bool `json:"available"` + Provider string `json:"provider"` + SupportedModes []string `json:"supported_modes"` +} + +type capabilityWebhooks struct { + Available bool `json:"available"` + MaxPerProfile int `json:"max_per_profile"` + SupportedTypes []string `json:"supported_types"` +} + +// HandleCapability handles GET /notifications/capability. Clients render +// setup UI from this response instead of introspecting admin settings. Push +// channels report unavailable until they ship +// (docs/superpowers/plans/notifications/02-03). +func (h *NotificationsHandler) HandleCapability(w http.ResponseWriter, r *http.Request) { + webhooks := capabilityWebhooks{Available: false, MaxPerProfile: 0, SupportedTypes: []string{}} + if h.system.Webhooks != nil && h.system.Settings.WebhooksEnabled(r.Context()) { + webhooks = capabilityWebhooks{ + Available: true, + MaxPerProfile: h.system.Settings.WebhooksMaxPerProfile(r.Context()), + SupportedTypes: []string{"discord", "generic"}, + } + } + webPush := capabilityWebPush{} + if h.system.WebPush != nil && h.system.Settings.WebPushEnabled(r.Context()) { + if publicKey, err := h.system.WebPush.PublicKey(r.Context()); err == nil && publicKey != "" { + webPush = capabilityWebPush{Available: true, PublicKey: publicKey} + } + } + email := capabilityAccountChannel{Modes: []string{}} + if h.system.EmailAvailable(r.Context()) { + modes := []string{notifications.ChannelModeDailyDigest} + if h.system.Settings.EmailAllowPerEpisode(r.Context()) { + modes = append(modes, + notifications.ChannelModePerEpisode, + notifications.ChannelModePerEpisodeAndDigest) + } + email = capabilityAccountChannel{ + Available: true, + Modes: modes, + DigestHour: h.system.Settings.EmailDigestHour(r.Context()), + } + } + discordCap := capabilityAccountChannel{Modes: []string{}} + if h.system.DiscordAvailable(r.Context()) { + modes := []string{notifications.ChannelModeDailyDigest} + if h.system.Settings.DiscordAllowPerEpisode(r.Context()) { + modes = append(modes, + notifications.ChannelModePerEpisode, + notifications.ChannelModePerEpisodeAndDigest) + } + discordCap = capabilityAccountChannel{ + Available: true, + Modes: modes, + DigestHour: h.system.Settings.DiscordDigestHour(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, + Discord: discordCap, + }) +} + +// HandleMintWSTicket handles POST /events/ws-ticket: mints a short-lived +// single-use websocket handshake ticket bound to (user, profile). Long-lived +// tokens must never ride the websocket query string — reverse-proxy access +// logs capture it. +func (h *NotificationsHandler) HandleMintWSTicket(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + + ticket, ttl, err := h.system.Tickets.Mint(r.Context(), userID, profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to mint websocket ticket") + return + } + writeJSON(w, http.StatusOK, wsTicketResponse{ + Ticket: ticket, + ExpiresIn: int(ttl.Seconds()), + }) +} diff --git a/internal/api/handlers/notifications_discord.go b/internal/api/handlers/notifications_discord.go new file mode 100644 index 00000000..c9a12333 --- /dev/null +++ b/internal/api/handlers/notifications_discord.go @@ -0,0 +1,230 @@ +package handlers + +import ( + "crypto/rand" + "encoding/hex" + "encoding/json" + "errors" + "net/http" + "net/url" + "strings" + "time" + + apimw "github.com/Silo-Server/silo-server/internal/api/middleware" + "github.com/Silo-Server/silo-server/internal/discord" + "github.com/Silo-Server/silo-server/internal/notifications" +) + +// discordSettingsPath is the SPA page the OAuth callback redirects back to, +// with ?discord_linked=1 or ?discord_error= for the page to toast. +const discordSettingsPath = "/settings/notifications" + +// DiscordNotificationsHandler serves the account-level Discord DM channel: +// preferences, the OAuth account-link flow, and the admin bot test. The +// preferences and link-init endpoints require auth; the link callback is +// public because Discord redirects the browser there without credentials — +// the one-time server-side state row authenticates it instead. +type DiscordNotificationsHandler struct { + system *notifications.System + publicURL string +} + +// NewDiscordNotificationsHandler creates a DiscordNotificationsHandler. +// publicURL may be empty, in which case linking is reported unavailable +// (Discord needs a stable redirect_uri origin). +func NewDiscordNotificationsHandler(system *notifications.System, publicURL string) *DiscordNotificationsHandler { + return &DiscordNotificationsHandler{system: system, publicURL: strings.TrimRight(publicURL, "/")} +} + +func (h *DiscordNotificationsHandler) redirectURI() string { + return h.publicURL + "/api/v1/notifications/discord/link/callback" +} + +// discordPreferencesResponse is the account-level Discord DM setting plus +// link state and health. +type discordPreferencesResponse struct { + Linked bool `json:"linked"` + DiscordUsername string `json:"discord_username,omitempty"` + Mode string `json:"mode"` + LinkFailure string `json:"link_failure,omitempty"` +} + +type updateDiscordPreferencesRequest struct { + Mode string `json:"mode"` +} + +func discordPrefsResponse(prefs notifications.DiscordPrefs) discordPreferencesResponse { + return discordPreferencesResponse{ + Linked: prefs.Linked(), + DiscordUsername: prefs.DiscordUsername, + Mode: prefs.Mode, + LinkFailure: prefs.LinkFailure, + } +} + +// HandleGetPreferences handles GET /notifications/discord-preferences. +func (h *DiscordNotificationsHandler) HandleGetPreferences(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + prefs, err := h.system.DiscordPrefsFor(r.Context(), userID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to load Discord preferences") + return + } + writeJSON(w, http.StatusOK, discordPrefsResponse(prefs)) +} + +// HandleUpdatePreferences handles PUT /notifications/discord-preferences. +func (h *DiscordNotificationsHandler) HandleUpdatePreferences(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + + var req updateDiscordPreferencesRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + + err := h.system.SetDiscordMode(r.Context(), userID, req.Mode) + switch { + case err == nil: + case errors.Is(err, notifications.ErrDiscordModeInvalid): + writeError(w, http.StatusBadRequest, "bad_request", "Unknown Discord notification mode") + return + case errors.Is(err, notifications.ErrDiscordModeNotAllowed): + writeError(w, http.StatusBadRequest, "not_allowed", "Per-episode Discord DMs are disabled by the administrator") + return + case errors.Is(err, notifications.ErrDiscordNotLinked): + writeError(w, http.StatusBadRequest, "not_linked", "Link a Discord account first") + return + default: + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to save Discord preferences") + return + } + prefs, err := h.system.DiscordPrefsFor(r.Context(), userID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to load Discord preferences") + return + } + writeJSON(w, http.StatusOK, discordPrefsResponse(prefs)) +} + +// HandleUnlink handles DELETE /notifications/discord-link: removes the +// linked identity and switches the channel off. +func (h *DiscordNotificationsHandler) HandleUnlink(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + if err := h.system.UnlinkDiscord(r.Context(), userID); err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to unlink Discord account") + return + } + writeJSON(w, http.StatusOK, discordPreferencesResponse{Mode: notifications.ChannelModeOff}) +} + +type discordLinkInitResponse struct { + URL string `json:"url"` +} + +// HandleLinkInit handles POST /notifications/discord/link/init: records a +// one-time state for the signed-in account and returns the Discord consent +// URL for the SPA to navigate to. +func (h *DiscordNotificationsHandler) HandleLinkInit(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + if !h.system.DiscordAvailable(r.Context()) { + writeError(w, http.StatusConflict, "not_configured", "Discord integration is not enabled by the administrator") + return + } + if h.publicURL == "" { + writeError(w, http.StatusConflict, "no_public_url", "Linking requires SILO_PUBLIC_URL to be configured") + return + } + + stateBytes := make([]byte, 32) + if _, err := rand.Read(stateBytes); err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to start Discord link") + return + } + state := hex.EncodeToString(stateBytes) + if err := h.system.BeginDiscordLink(r.Context(), state, userID); err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to start Discord link") + return + } + + query := url.Values{ + "client_id": {h.system.Settings.DiscordClientID(r.Context())}, + "response_type": {"code"}, + "scope": {"identify"}, + "redirect_uri": {h.redirectURI()}, + "state": {state}, + } + writeJSON(w, http.StatusOK, discordLinkInitResponse{ + URL: discord.AuthorizeURL + "?" + query.Encode(), + }) +} + +// HandleLinkCallback handles GET /notifications/discord/link/callback. +// Public route: Discord redirects the user's browser here. The one-time +// state row authenticates the request and recovers which account started +// the flow; the browser is then sent back to the settings page. +func (h *DiscordNotificationsHandler) HandleLinkCallback(w http.ResponseWriter, r *http.Request) { + redirectBack := func(params url.Values) { + http.Redirect(w, r, discordSettingsPath+"?"+params.Encode(), http.StatusFound) + } + + // The admin may have switched the integration off between init and + // callback; don't complete a link for a disabled channel. + if !h.system.DiscordAvailable(r.Context()) { + redirectBack(url.Values{"discord_error": {"disabled"}}) + return + } + + query := r.URL.Query() + if query.Get("error") != "" { + // The user declined on Discord's consent screen. + redirectBack(url.Values{"discord_error": {"denied"}}) + return + } + state := query.Get("state") + code := query.Get("code") + if state == "" || code == "" { + redirectBack(url.Values{"discord_error": {"invalid_callback"}}) + return + } + + userID, ok, err := h.system.ConsumeDiscordLinkState(r.Context(), state) + if err != nil || !ok { + redirectBack(url.Values{"discord_error": {"state_invalid"}}) + return + } + if _, err := h.system.CompleteDiscordLink(r.Context(), userID, code, h.redirectURI()); err != nil { + redirectBack(url.Values{"discord_error": {"exchange_failed"}}) + return + } + redirectBack(url.Values{"discord_linked": {"1"}}) +} + +type discordTestResponse struct { + OK bool `json:"ok"` + DurationMS int64 `json:"duration_ms"` + Message string `json:"message"` +} + +// HandleAdminTest handles POST /admin/notifications/discord/test: verifies +// the configured bot token by fetching the bot's own identity. +func (h *DiscordNotificationsHandler) HandleAdminTest(w http.ResponseWriter, r *http.Request) { + start := time.Now() + botUser, err := h.system.TestDiscordBot(r.Context()) + duration := time.Since(start).Milliseconds() + if errors.Is(err, notifications.ErrDiscordNotConfigured) { + writeJSON(w, http.StatusOK, discordTestResponse{ + OK: false, DurationMS: duration, Message: "Bot token is not configured", + }) + return + } + if err != nil { + writeJSON(w, http.StatusOK, discordTestResponse{ + OK: false, DurationMS: duration, Message: err.Error(), + }) + return + } + writeJSON(w, http.StatusOK, discordTestResponse{ + OK: true, DurationMS: duration, Message: "Connected as " + botUser.Username, + }) +} diff --git a/internal/api/handlers/notifications_email.go b/internal/api/handlers/notifications_email.go new file mode 100644 index 00000000..e9863965 --- /dev/null +++ b/internal/api/handlers/notifications_email.go @@ -0,0 +1,209 @@ +package handlers + +import ( + "encoding/json" + "errors" + "fmt" + "html" + "net/http" + + apimw "github.com/Silo-Server/silo-server/internal/api/middleware" + "github.com/Silo-Server/silo-server/internal/notifications" +) + +// emailPreferencesResponse is one profile's email notification state. The +// channel is profile-scoped: each profile verifies its own destination +// address and receives nothing until it has one — there is no account-email +// fallback. +type emailPreferencesResponse struct { + Mode string `json:"mode"` + // CustomEmail is the verified destination ('' = none; channel inert). + CustomEmail string `json:"custom_email"` + // PendingEmail is an address awaiting link-click verification. + PendingEmail string `json:"pending_email"` + // CanEditAddress is false for child profiles, which cannot set + // addresses (and so cannot receive email notifications). + CanEditAddress bool `json:"can_edit_address"` +} + +type updateEmailPreferencesRequest struct { + Mode string `json:"mode"` +} + +type updateEmailAddressRequest struct { + Email string `json:"email"` +} + +func emailPreferencesPayload(state notifications.EmailPreferencesState) emailPreferencesResponse { + return emailPreferencesResponse{ + Mode: state.Mode, + CustomEmail: state.CustomEmail, + PendingEmail: state.PendingEmail, + CanEditAddress: !state.IsChild, + } +} + +// respondEmailPreferences re-reads and writes the profile's full email state, +// so every mutation returns the same shape as GET. +func (h *NotificationsHandler) respondEmailPreferences(w http.ResponseWriter, r *http.Request, userID int, profileID string) { + state, err := h.system.EmailPreferences(r.Context(), userID, profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to load email preferences") + return + } + writeJSON(w, http.StatusOK, emailPreferencesPayload(state)) +} + +// HandleGetEmailPreferences handles GET /notifications/email-preferences. +func (h *NotificationsHandler) HandleGetEmailPreferences(w http.ResponseWriter, r *http.Request) { + h.respondEmailPreferences(w, r, apimw.GetUserID(r.Context()), apimw.GetProfileID(r.Context())) +} + +// HandleUpdateEmailPreferences handles PUT /notifications/email-preferences. +func (h *NotificationsHandler) HandleUpdateEmailPreferences(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(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, profileID, 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", "Verify an email address for this profile first") + return + default: + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to save email preferences") + return + } + h.respondEmailPreferences(w, r, userID, profileID) +} + +// HandleRequestEmailAddress handles PUT /notifications/email-preferences/address. +// It stores the candidate address and emails it a verification link; the +// address only becomes the destination once that link is clicked. +func (h *NotificationsHandler) HandleRequestEmailAddress(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + + var req updateEmailAddressRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + + err := h.system.RequestEmailAddress(r.Context(), userID, profileID, req.Email) + switch { + case err == nil: + case errors.Is(err, notifications.ErrEmailInvalidAddress): + writeError(w, http.StatusBadRequest, "bad_request", "Invalid email address") + return + case errors.Is(err, notifications.ErrEmailChildProfile): + writeError(w, http.StatusForbidden, "child_profile", "Child profiles cannot set a custom notification address") + return + case errors.Is(err, notifications.ErrEmailAddressInUse): + writeError(w, http.StatusConflict, "address_in_use", "That email address is already used by another profile or account") + return + case errors.Is(err, notifications.ErrEmailVerifyRateLimited): + writeError(w, http.StatusTooManyRequests, "rate_limited", "Too many verification emails; try again later") + return + case errors.Is(err, notifications.ErrEmailNoLinkBase): + writeError(w, http.StatusConflict, "no_external_url", "The server has no external URL configured for verification links") + return + default: + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to send the verification email") + return + } + h.respondEmailPreferences(w, r, userID, profileID) +} + +// HandleClearEmailAddress handles DELETE /notifications/email-preferences/address. +func (h *NotificationsHandler) HandleClearEmailAddress(w http.ResponseWriter, r *http.Request) { + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + + err := h.system.ClearEmailAddress(r.Context(), userID, profileID) + switch { + case err == nil: + case errors.Is(err, notifications.ErrEmailChildProfile): + writeError(w, http.StatusForbidden, "child_profile", "Child profiles cannot change the notification address") + return + default: + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to remove the custom address") + return + } + h.respondEmailPreferences(w, r, userID, profileID) +} + +// EmailLinkHandler serves the public tokenized email endpoints: address +// verification and unsubscribe. Both are clicked from email clients on +// devices that may have no Silo session, so they render minimal standalone +// HTML instead of redirecting into the authenticated app. +type EmailLinkHandler struct { + system *notifications.System +} + +// NewEmailLinkHandler creates an EmailLinkHandler. +func NewEmailLinkHandler(system *notifications.System) *EmailLinkHandler { + return &EmailLinkHandler{system: system} +} + +// HandleVerify handles GET /notifications/email/verify?token=... +func (h *EmailLinkHandler) HandleVerify(w http.ResponseWriter, r *http.Request) { + outcome, err := h.system.VerifyEmailToken(r.Context(), r.URL.Query().Get("token")) + switch { + case err != nil: + writeEmailLinkPage(w, http.StatusInternalServerError, "Something went wrong", + "The address could not be verified. Try the link again in a moment.") + case outcome == notifications.EmailVerifyConflict: + writeEmailLinkPage(w, http.StatusConflict, "Address already in use", + "This address now belongs to another profile or account. Choose a different address in Silo's notification settings.") + case outcome == notifications.EmailVerifyInvalid: + writeEmailLinkPage(w, http.StatusBadRequest, "Link expired or already used", + "Request a new verification email from Silo's notification settings.") + default: + writeEmailLinkPage(w, http.StatusOK, "Address verified", + "Silo notifications for this profile will now be delivered here. You can close this page.") + } +} + +// HandleUnsubscribe handles GET and POST /notifications/email/unsubscribe?token=... +// POST is the RFC 8058 one-click target mail clients call directly. +func (h *EmailLinkHandler) HandleUnsubscribe(w http.ResponseWriter, r *http.Request) { + ok, err := h.system.UnsubscribeEmail(r.Context(), r.URL.Query().Get("token")) + switch { + case err != nil: + writeEmailLinkPage(w, http.StatusInternalServerError, "Something went wrong", + "Could not unsubscribe. Try the link again in a moment.") + case !ok: + writeEmailLinkPage(w, http.StatusBadRequest, "Link invalid", + "This unsubscribe link is no longer valid. Manage notifications in Silo's settings.") + default: + writeEmailLinkPage(w, http.StatusOK, "Unsubscribed", + "This profile will no longer receive notification emails. Re-enable them any time in Silo's settings.") + } +} + +// writeEmailLinkPage renders the minimal standalone page behind tokenized +// email links. +func writeEmailLinkPage(w http.ResponseWriter, status int, title, detail string) { + w.Header().Set("Content-Type", "text/html; charset=utf-8") + w.WriteHeader(status) + fmt.Fprintf(w, `%s — Silo + +
+

%s

+

%s

+
`, + html.EscapeString(title), html.EscapeString(title), html.EscapeString(detail)) +} diff --git a/internal/api/handlers/notifications_webhooks.go b/internal/api/handlers/notifications_webhooks.go new file mode 100644 index 00000000..f584ea25 --- /dev/null +++ b/internal/api/handlers/notifications_webhooks.go @@ -0,0 +1,221 @@ +package handlers + +import ( + "encoding/json" + "errors" + "net/http" + "time" + + apimw "github.com/Silo-Server/silo-server/internal/api/middleware" + "github.com/Silo-Server/silo-server/internal/notifications" + "github.com/go-chi/chi/v5" +) + +// webhookResponse is the API view of a webhook. It never includes the +// destination URL (Discord webhook tokens are bearer credentials in the URL +// path) or the signing secret — only url_host for identification. +type webhookResponse struct { + ID string `json:"id"` + Name string `json:"name"` + Type string `json:"type"` + URLHost string `json:"url_host"` + Enabled bool `json:"enabled"` + NotifyFavorites bool `json:"notify_favorites"` + NotifyWatchlist bool `json:"notify_watchlist"` + NotifyContinueWatching bool `json:"notify_continue_watching"` + NotifyNextUp bool `json:"notify_next_up"` + NotifyRequests bool `json:"notify_requests"` + ConsecutiveFailures int `json:"consecutive_failures"` + DisabledReason *string `json:"disabled_reason"` + LastSuccessAt *time.Time `json:"last_success_at"` + LastFailureAt *time.Time `json:"last_failure_at"` + LastFailureStatus *int `json:"last_failure_status"` + LastFailureMessage *string `json:"last_failure_message"` + // SigningSecret is present only in create / rotate-secret responses. + SigningSecret string `json:"signing_secret,omitempty"` +} + +func webhookToResponse(hook notifications.Webhook) webhookResponse { + return webhookResponse{ + ID: hook.ID, + Name: hook.Name, + Type: hook.Type, + URLHost: hook.URLHost, + Enabled: hook.Enabled, + NotifyFavorites: hook.NotifyFavorites, + NotifyWatchlist: hook.NotifyWatchlist, + NotifyContinueWatching: hook.NotifyContinueWatching, + NotifyNextUp: hook.NotifyNextUp, + NotifyRequests: hook.NotifyRequests, + ConsecutiveFailures: hook.ConsecutiveFailures, + DisabledReason: hook.DisabledReason, + LastSuccessAt: hook.LastSuccessAt, + LastFailureAt: hook.LastFailureAt, + LastFailureStatus: hook.LastFailureStatus, + LastFailureMessage: hook.LastFailureMessage, + } +} + +type webhookRequest struct { + Name *string `json:"name"` + URL *string `json:"url"` + Type *string `json:"type"` + Enabled *bool `json:"enabled"` + NotifyFavorites *bool `json:"notify_favorites"` + NotifyWatchlist *bool `json:"notify_watchlist"` + NotifyContinueWatching *bool `json:"notify_continue_watching"` + NotifyNextUp *bool `json:"notify_next_up"` + NotifyRequests *bool `json:"notify_requests"` +} + +func (r webhookRequest) toInput() notifications.WebhookInput { + return notifications.WebhookInput{ + Name: r.Name, + URL: r.URL, + Type: r.Type, + Enabled: r.Enabled, + NotifyFavorites: r.NotifyFavorites, + NotifyWatchlist: r.NotifyWatchlist, + NotifyContinueWatching: r.NotifyContinueWatching, + NotifyNextUp: r.NotifyNextUp, + NotifyRequests: r.NotifyRequests, + } +} + +func (h *NotificationsHandler) webhooks() *notifications.WebhookService { + if h == nil || h.system == nil { + return nil + } + return h.system.Webhooks +} + +func writeWebhookError(w http.ResponseWriter, err error) { + switch { + case errors.Is(err, notifications.ErrWebhooksDisabled): + writeError(w, http.StatusForbidden, "webhooks_disabled", "Webhooks are disabled by the server administrator") + case errors.Is(err, notifications.ErrWebhookNotFound): + writeError(w, http.StatusNotFound, "not_found", "Webhook not found") + case errors.Is(err, notifications.ErrWebhookLimit): + writeError(w, http.StatusUnprocessableEntity, "limit_reached", "Webhook limit reached for this profile") + case errors.Is(err, notifications.ErrWebhookInvalid): + writeError(w, http.StatusBadRequest, "bad_request", err.Error()) + default: + writeError(w, http.StatusInternalServerError, "internal_error", "Webhook operation failed") + } +} + +// HandleListWebhooks handles GET /notifications/webhooks. +func (h *NotificationsHandler) HandleListWebhooks(w http.ResponseWriter, r *http.Request) { + service := h.webhooks() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Webhooks are not available") + return + } + profileID := apimw.GetProfileID(r.Context()) + hooks, err := service.List(r.Context(), profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to list webhooks") + return + } + responses := make([]webhookResponse, 0, len(hooks)) + for _, hook := range hooks { + responses = append(responses, webhookToResponse(hook)) + } + writeJSON(w, http.StatusOK, map[string]any{"webhooks": responses}) +} + +// HandleCreateWebhook handles POST /notifications/webhooks. For generic +// webhooks the response carries the signing secret exactly once. +func (h *NotificationsHandler) HandleCreateWebhook(w http.ResponseWriter, r *http.Request) { + service := h.webhooks() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Webhooks are not available") + return + } + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + + var req webhookRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + hook, signingSecret, err := service.Create(r.Context(), userID, profileID, req.toInput()) + if err != nil { + writeWebhookError(w, err) + return + } + response := webhookToResponse(*hook) + response.SigningSecret = signingSecret + writeJSON(w, http.StatusCreated, response) +} + +// HandleUpdateWebhook handles PUT /notifications/webhooks/{id}. +func (h *NotificationsHandler) HandleUpdateWebhook(w http.ResponseWriter, r *http.Request) { + service := h.webhooks() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Webhooks are not available") + return + } + profileID := apimw.GetProfileID(r.Context()) + + var req webhookRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + hook, err := service.Update(r.Context(), profileID, chi.URLParam(r, "id"), req.toInput()) + if err != nil { + writeWebhookError(w, err) + return + } + writeJSON(w, http.StatusOK, webhookToResponse(*hook)) +} + +// HandleDeleteWebhook handles DELETE /notifications/webhooks/{id}. Idempotent. +func (h *NotificationsHandler) HandleDeleteWebhook(w http.ResponseWriter, r *http.Request) { + service := h.webhooks() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Webhooks are not available") + return + } + profileID := apimw.GetProfileID(r.Context()) + if err := service.Delete(r.Context(), profileID, chi.URLParam(r, "id")); err != nil { + writeWebhookError(w, err) + return + } + w.WriteHeader(http.StatusNoContent) +} + +// HandleRotateWebhookSecret handles POST /notifications/webhooks/{id}/rotate-secret. +func (h *NotificationsHandler) HandleRotateWebhookSecret(w http.ResponseWriter, r *http.Request) { + service := h.webhooks() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Webhooks are not available") + return + } + profileID := apimw.GetProfileID(r.Context()) + signingSecret, err := service.RotateSecret(r.Context(), profileID, chi.URLParam(r, "id")) + if err != nil { + writeWebhookError(w, err) + return + } + writeJSON(w, http.StatusOK, map[string]string{"signing_secret": signingSecret}) +} + +// HandleTestWebhook handles POST /notifications/webhooks/{id}/test. The test +// send is synchronous and never touches the retry/auto-disable counters. +func (h *NotificationsHandler) HandleTestWebhook(w http.ResponseWriter, r *http.Request) { + service := h.webhooks() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Webhooks are not available") + return + } + profileID := apimw.GetProfileID(r.Context()) + result, err := service.Test(r.Context(), profileID, chi.URLParam(r, "id")) + if err != nil { + writeWebhookError(w, err) + return + } + writeJSON(w, http.StatusOK, result) +} diff --git a/internal/api/handlers/notifications_webpush.go b/internal/api/handlers/notifications_webpush.go new file mode 100644 index 00000000..c9de9963 --- /dev/null +++ b/internal/api/handlers/notifications_webpush.go @@ -0,0 +1,144 @@ +package handlers + +import ( + "encoding/json" + "errors" + "net/http" + "time" + + apimw "github.com/Silo-Server/silo-server/internal/api/middleware" + "github.com/Silo-Server/silo-server/internal/notifications" + "github.com/go-chi/chi/v5" +) + +// webPushSubscriptionResponse is the API view of a browser push +// registration. The keys are write-only: clients re-subscribe rather than +// read them back. +type webPushSubscriptionResponse struct { + ID string `json:"id"` + Endpoint string `json:"endpoint"` + DeviceName string `json:"device_name,omitempty"` + Enabled bool `json:"enabled"` + CreatedAt time.Time `json:"created_at"` + LastSuccessAt *time.Time `json:"last_success_at"` + LastFailureAt *time.Time `json:"last_failure_at"` +} + +func webPushToResponse(sub notifications.WebPushSubscription) webPushSubscriptionResponse { + return webPushSubscriptionResponse{ + ID: sub.ID, + Endpoint: sub.Endpoint, + DeviceName: sub.DeviceName, + Enabled: sub.Enabled, + CreatedAt: sub.CreatedAt, + LastSuccessAt: sub.LastSuccessAt, + LastFailureAt: sub.LastFailureAt, + } +} + +func (h *NotificationsHandler) webPush() *notifications.WebPushService { + if h == nil || h.system == nil { + return nil + } + return h.system.WebPush +} + +type webPushSubscribeRequest struct { + Endpoint string `json:"endpoint"` + Keys struct { + P256dh string `json:"p256dh"` + Auth string `json:"auth"` + } `json:"keys"` + DeviceName string `json:"device_name"` +} + +// HandleWebPushSubscribe handles POST /notifications/web-push/subscriptions. +// The body matches PushSubscription.toJSON() plus an optional device name. +func (h *NotificationsHandler) HandleWebPushSubscribe(w http.ResponseWriter, r *http.Request) { + service := h.webPush() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Web push is not available") + return + } + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + + var req webPushSubscribeRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + writeError(w, http.StatusBadRequest, "bad_request", "Invalid request body") + return + } + sub, err := service.Subscribe(r.Context(), userID, profileID, + req.Endpoint, req.Keys.P256dh, req.Keys.Auth, req.DeviceName) + if err != nil { + if errors.Is(err, notifications.ErrWebPushInvalid) { + writeError(w, http.StatusBadRequest, "bad_request", err.Error()) + return + } + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to register push subscription") + return + } + writeJSON(w, http.StatusCreated, webPushToResponse(*sub)) +} + +// HandleWebPushList handles GET /notifications/web-push/subscriptions. +func (h *NotificationsHandler) HandleWebPushList(w http.ResponseWriter, r *http.Request) { + service := h.webPush() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Web push is not available") + return + } + profileID := apimw.GetProfileID(r.Context()) + subs, err := service.List(r.Context(), profileID) + if err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to list push subscriptions") + return + } + responses := make([]webPushSubscriptionResponse, 0, len(subs)) + for _, sub := range subs { + responses = append(responses, webPushToResponse(sub)) + } + writeJSON(w, http.StatusOK, map[string]any{"subscriptions": responses}) +} + +// HandleWebPushDelete handles DELETE /notifications/web-push/subscriptions/{id}. +func (h *NotificationsHandler) HandleWebPushDelete(w http.ResponseWriter, r *http.Request) { + service := h.webPush() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Web push is not available") + return + } + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + if err := service.Unsubscribe(r.Context(), userID, profileID, chi.URLParam(r, "id"), ""); err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to remove push subscription") + return + } + w.WriteHeader(http.StatusNoContent) +} + +type webPushUnsubscribeRequest struct { + Endpoint string `json:"endpoint"` +} + +// HandleWebPushUnsubscribe handles POST /notifications/web-push/unsubscribe. +// Browsers only know their endpoint, not the server-side row ID. Idempotent. +func (h *NotificationsHandler) HandleWebPushUnsubscribe(w http.ResponseWriter, r *http.Request) { + service := h.webPush() + if service == nil { + writeError(w, http.StatusServiceUnavailable, "unavailable", "Web push is not available") + return + } + userID := apimw.GetUserID(r.Context()) + profileID := apimw.GetProfileID(r.Context()) + var req webPushUnsubscribeRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil || req.Endpoint == "" { + writeError(w, http.StatusBadRequest, "bad_request", "An endpoint is required") + return + } + if err := service.Unsubscribe(r.Context(), userID, profileID, "", req.Endpoint); err != nil { + writeError(w, http.StatusInternalServerError, "internal_error", "Failed to remove push subscription") + return + } + w.WriteHeader(http.StatusNoContent) +} diff --git a/internal/api/router.go b/internal/api/router.go index 3b026983..ae4ece0e 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -37,6 +37,7 @@ import ( "github.com/Silo-Server/silo-server/internal/intromarkers" "github.com/Silo-Server/silo-server/internal/libraryingest" "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" @@ -119,6 +120,7 @@ type Dependencies struct { NodeID string LogStreamHub *logstream.Hub RealtimeHub *notifications.Hub + Notifications *notifications.System // user-facing release notifications (may be nil) EventsHub *evt.Hub ScanRegistry *evt.ScanRegistry LibraryScanQueue *scanqueue.Service @@ -529,6 +531,12 @@ func NewRouter(deps Dependencies) chi.Router { if viewerResolver != nil { requestSvc.SetEntitlementResolver(mediarequests.NewAccessEntitlements(viewerResolver)) } + // Server-channel broadcast of request lifecycle events (submitted / + // approved / declined). Fulfilled rides the reconcile service's + // fulfillment notifier instead. + if lifecycle := notifications.NewServerChannelLifecycleNotifier(deps.Notifications); lifecycle != nil { + requestSvc.SetLifecycleNotifier(lifecycle) + } requestHandler = handlers.NewRequestsHandler(requestSvc) autoscanRepo := autoscan.NewRepository(deps.DB, deps.SecretCipher) @@ -1477,6 +1485,28 @@ func NewRouter(deps Dependencies) chi.Router { }) } + // Discord account-link OAuth callback: public — Discord redirects the + // browser here without credentials; the one-time link-state row + // authenticates the request and maps it back to the initiating + // account. The static path coexists with the authenticated + // /notifications subrouter below (static routes win in chi). + var discordNotificationsHandler *handlers.DiscordNotificationsHandler + if deps.Notifications != nil { + discordNotificationsHandler = handlers.NewDiscordNotificationsHandler(deps.Notifications, deps.PublicURL) + r.Get("/notifications/discord/link/callback", discordNotificationsHandler.HandleLinkCallback) + + // Tokenized email links: public — clicked from mail clients on + // devices without a Silo session; the single-use token (verify) + // or per-profile capability token (unsubscribe) authenticates the + // request. Static paths coexist with the authenticated + // /notifications subrouter below, same as the Discord callback. + deps.Notifications.SetPublicURL(deps.PublicURL) + emailLinkHandler := handlers.NewEmailLinkHandler(deps.Notifications) + r.Get("/notifications/email/verify", emailLinkHandler.HandleVerify) + r.Get("/notifications/email/unsubscribe", emailLinkHandler.HandleUnsubscribe) + r.Post("/notifications/email/unsubscribe", emailLinkHandler.HandleUnsubscribe) + } + // API key management routes (auth only, no viewer access needed). if apiKeyRepo != nil && authMiddleware != nil { r.Group(func(r chi.Router) { @@ -1522,9 +1552,60 @@ func NewRouter(deps Dependencies) chi.Router { deps.LibraryScanQueue, historyImportSvc, ) + eventsHandler.SetNotificationsSystem(deps.Notifications) r.Get("/events/ws", eventsHandler.HandleWebSocket) } + // User notifications: profile-scoped inbox, preferences, and + // the websocket handshake ticket. + if deps.Notifications != nil { + if detailSvc != nil { + deps.Notifications.SetImageResolver(detailSvc) + } + notificationsHandler := handlers.NewNotificationsHandler(deps.Notifications, deps.EventsHub) + r.With(apimw.RequireProfile).Post("/events/ws-ticket", notificationsHandler.HandleMintWSTicket) + // Discord DM channel: the linked identity and mode hang off + // the login account, not a profile, so these stay outside + // the RequireProfile subrouter below (static paths coexist + // with it, same as the public email-link routes above). + if discordNotificationsHandler != nil { + r.Get("/notifications/discord-preferences", discordNotificationsHandler.HandleGetPreferences) + r.Put("/notifications/discord-preferences", discordNotificationsHandler.HandleUpdatePreferences) + r.Delete("/notifications/discord-link", discordNotificationsHandler.HandleUnlink) + r.Post("/notifications/discord/link/init", discordNotificationsHandler.HandleLinkInit) + } + r.Route("/notifications", func(r chi.Router) { + r.Use(apimw.RequireProfile) + r.Get("/", notificationsHandler.HandleList) + r.Get("/sync", notificationsHandler.HandleSync) + r.Get("/unread-count", notificationsHandler.HandleUnreadCount) + 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.Put("/email-preferences/address", notificationsHandler.HandleRequestEmailAddress) + r.Delete("/email-preferences/address", notificationsHandler.HandleClearEmailAddress) + r.Post("/read-all", notificationsHandler.HandleReadAll) + r.Route("/webhooks", func(r chi.Router) { + r.Get("/", notificationsHandler.HandleListWebhooks) + r.Post("/", notificationsHandler.HandleCreateWebhook) + r.Put("/{id}", notificationsHandler.HandleUpdateWebhook) + r.Delete("/{id}", notificationsHandler.HandleDeleteWebhook) + r.Post("/{id}/rotate-secret", notificationsHandler.HandleRotateWebhookSecret) + r.Post("/{id}/test", notificationsHandler.HandleTestWebhook) + }) + r.Route("/web-push", func(r chi.Router) { + r.Get("/subscriptions", notificationsHandler.HandleWebPushList) + r.Post("/subscriptions", notificationsHandler.HandleWebPushSubscribe) + r.Delete("/subscriptions/{id}", notificationsHandler.HandleWebPushDelete) + r.Post("/unsubscribe", notificationsHandler.HandleWebPushUnsubscribe) + }) + r.Get("/{id}", notificationsHandler.HandleGet) + r.Post("/{id}/read", notificationsHandler.HandleMarkRead) + }) + } + // Marker read/write/clear for any authenticated viewer: users // fix and create intro/recap/credits/preview markers from the // player. Writes are stamped source="manual" and contributed to @@ -2118,6 +2199,24 @@ func NewRouter(deps Dependencies) chi.Router { r.Get("/settings/{key}", adminHandler.HandleGetSetting) r.Get("/settings", adminHandler.HandleGetSettings) r.Put("/settings/{key}", adminHandler.HandleUpdateSetting) + if settingsRepo != nil { + emailHandler := handlers.NewEmailHandler(mail.NewSMTPSender(settingsRepo)) + r.Post("/email/test", emailHandler.HandleTest) + } + if discordNotificationsHandler != nil { + r.Post("/notifications/discord/test", discordNotificationsHandler.HandleAdminTest) + } + if deps.Notifications != nil && deps.Notifications.ServerChannels != nil { + serverChannelsHandler := handlers.NewAdminServerChannelsHandler(deps.Notifications) + r.Route("/notifications/server-channels", func(r chi.Router) { + r.Get("/", serverChannelsHandler.HandleList) + r.Post("/", serverChannelsHandler.HandleCreate) + r.Put("/{id}", serverChannelsHandler.HandleUpdate) + r.Delete("/{id}", serverChannelsHandler.HandleDelete) + r.Post("/{id}/rotate-secret", serverChannelsHandler.HandleRotateSecret) + r.Post("/{id}/test", serverChannelsHandler.HandleTest) + }) + } if adminIntroHandler != nil { r.Post("/items/{id}/refresh-markers", adminIntroHandler.HandleRefreshEpisodeMarkers) r.Post("/items/{id}/redetect-intro", adminIntroHandler.HandleRedetectEpisodeIntro) diff --git a/internal/catalog/encrypted_settings_repo.go b/internal/catalog/encrypted_settings_repo.go index 30328d48..0c9e67a2 100644 --- a/internal/catalog/encrypted_settings_repo.go +++ b/internal/catalog/encrypted_settings_repo.go @@ -91,6 +91,19 @@ var SensitiveSettingKeys = map[string]bool{ // still referenced by older request_integrations rows until backfilled). "requests.radarr.api_key": true, "requests.sonarr.api_key": true, + + // Shared outbound email (internal/mail) SMTP credential. + "email.smtp_password": true, + + // Discord notification integration. The client_id is public in Discord's + // own UI, so only the secret and bot token are encrypted. + "discord.client_secret": true, + "discord.bot_token": true, + + // Web Push VAPID keypair JSON (generated + persisted atomically as one + // value by the notifications system; clients receive the public half via + // the capability endpoint, never from the settings store). + "notifications.web_push.vapid_keypair": true, } // EncryptedSettingsRepo decorates a raw settings store, transparently @@ -127,6 +140,30 @@ func (r *EncryptedSettingsRepo) Set(ctx context.Context, key, value string) erro return r.inner.Set(ctx, key, value) } +// settingsConditionalWriter is the optional conditional-write capability of a +// raw settings store (satisfied by *ServerSettingsRepo). +type settingsConditionalWriter interface { + SetIfAbsent(ctx context.Context, key, value string) (bool, error) +} + +// SetIfAbsent applies Set's encryption contract to a conditional write: the +// value lands only when the key currently has no value, so concurrent +// provisioners of generated secrets cannot overwrite each other. +func (r *EncryptedSettingsRepo) SetIfAbsent(ctx context.Context, key, value string) (bool, error) { + inner, ok := r.inner.(settingsConditionalWriter) + if !ok { + return false, fmt.Errorf("settings store does not support conditional writes") + } + if SensitiveSettingKeys[key] && value != "" { + ct, err := r.cipher.Encrypt(value, secret.SettingsAAD(key)) + if err != nil { + return false, fmt.Errorf("encrypt setting %q: %w", key, err) + } + value = ct + } + return inner.SetIfAbsent(ctx, key, value) +} + // Get reads a value and applies the read-path contract: legacy plaintext passes // through, an enc:v1: value is decrypted, and a corrupt ciphertext errors. func (r *EncryptedSettingsRepo) Get(ctx context.Context, key string) (string, error) { diff --git a/internal/catalog/encrypted_settings_repo_test.go b/internal/catalog/encrypted_settings_repo_test.go index 229f6934..44788228 100644 --- a/internal/catalog/encrypted_settings_repo_test.go +++ b/internal/catalog/encrypted_settings_repo_test.go @@ -154,6 +154,10 @@ func TestSensitiveSettingKeys_Audited(t *testing.T) { "requests.radarr.api_key", "requests.sonarr.api_key", "watchsync.trakt.client_secret", + "email.smtp_password", + "discord.client_secret", + "discord.bot_token", + "notifications.web_push.vapid_keypair", } for _, k := range mustHave { if !SensitiveSettingKeys[k] { diff --git a/internal/catalog/episode_catalog_source.go b/internal/catalog/episode_catalog_source.go index f4e34bca..758fcdc9 100644 --- a/internal/catalog/episode_catalog_source.go +++ b/internal/catalog/episode_catalog_source.go @@ -29,6 +29,7 @@ const episodeCatalogSelectBody = `( COALESCE(e.tmdb_id, '') AS tmdb_id, COALESCE(e.tvdb_id, '') AS tvdb_id, COALESCE(NULLIF(s.poster_path, ''), NULLIF(si.poster_path, ''), NULLIF(e.still_path, ''), '') AS poster_path, + ''::text AS poster_source_path, COALESCE(NULLIF(s.poster_thumbhash, ''), NULLIF(si.poster_thumbhash, ''), NULLIF(e.still_thumbhash, ''), '') AS poster_thumbhash, COALESCE(si.backdrop_path, '') AS backdrop_path, COALESCE(si.backdrop_thumbhash, '') AS backdrop_thumbhash, diff --git a/internal/catalog/item_repo.go b/internal/catalog/item_repo.go index 015d3be9..f2466ead 100644 --- a/internal/catalog/item_repo.go +++ b/internal/catalog/item_repo.go @@ -117,7 +117,7 @@ var itemColumnNames = []string{ "content_rating", "runtime", "overview", "tagline", "rating_imdb", "rating_tmdb", "rating_rt_critic", "rating_rt_audience", "imdb_id", "tmdb_id", "tvdb_id", - "poster_path", "poster_thumbhash", "backdrop_path", "backdrop_thumbhash", "logo_path", + "poster_path", "poster_source_path", "poster_thumbhash", "backdrop_path", "backdrop_thumbhash", "logo_path", "metadata_s3_path", "metadata_etag", "season_count", "studios", "networks", "countries", "keywords", "original_language", "release_date::text", "first_air_date", "last_air_date", "air_time", "air_timezone", "show_status", @@ -130,6 +130,7 @@ var itemColumnNames = []string{ // lists coalesce them to ”. var nullableStringItemColumns = map[string]bool{ "poster_path": true, + "poster_source_path": true, "poster_thumbhash": true, "backdrop_path": true, "backdrop_thumbhash": true, @@ -215,6 +216,7 @@ func scanItem(row pgx.Row) (*models.MediaItem, error) { &item.TmdbID, &item.TvdbID, &item.PosterPath, + &item.PosterSourcePath, &item.PosterThumbhash, &item.BackdropPath, &item.BackdropThumbhash, @@ -278,6 +280,7 @@ func scanItems(rows pgx.Rows) ([]*models.MediaItem, error) { &item.TmdbID, &item.TvdbID, &item.PosterPath, + &item.PosterSourcePath, &item.PosterThumbhash, &item.BackdropPath, &item.BackdropThumbhash, @@ -351,6 +354,7 @@ func scanItemsWithTotal(rows pgx.Rows) ([]*models.MediaItem, int, error) { &item.TmdbID, &item.TvdbID, &item.PosterPath, + &item.PosterSourcePath, &item.PosterThumbhash, &item.BackdropPath, &item.BackdropThumbhash, @@ -417,7 +421,7 @@ func (r *ItemRepository) upsert(ctx context.Context, execer itemExecer, item *mo content_rating, runtime, overview, tagline, rating_imdb, rating_tmdb, rating_rt_critic, rating_rt_audience, imdb_id, tmdb_id, tvdb_id, - poster_path, poster_thumbhash, backdrop_path, backdrop_thumbhash, logo_path, + poster_path, poster_source_path, poster_thumbhash, backdrop_path, backdrop_thumbhash, logo_path, metadata_s3_path, metadata_etag, season_count, studios, networks, countries, keywords, original_language, release_date, first_air_date, last_air_date, air_time, air_timezone, show_status, @@ -428,12 +432,12 @@ func (r *ItemRepository) upsert(ctx context.Context, execer itemExecer, item *mo $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, - $20, $21, $22, $23, $24, - $25, $26, $27, - $28, $29, $30, $31, $32, $33, $34, $35, $36, $37, - $38, - $39, $40, $41, - $42, $43, $44 + $20, $21, $22, $23, $24, $25, + $26, $27, $28, + $29, $30, $31, $32, $33, $34, $35, $36, $37, $38, + $39, + $40, $41, $42, + $43, $44, $45 ) ON CONFLICT (content_id) DO UPDATE SET type = EXCLUDED.type, @@ -455,6 +459,7 @@ func (r *ItemRepository) upsert(ctx context.Context, execer itemExecer, item *mo tmdb_id = EXCLUDED.tmdb_id, tvdb_id = EXCLUDED.tvdb_id, poster_path = EXCLUDED.poster_path, + poster_source_path = EXCLUDED.poster_source_path, poster_thumbhash = EXCLUDED.poster_thumbhash, backdrop_path = EXCLUDED.backdrop_path, backdrop_thumbhash = EXCLUDED.backdrop_thumbhash, @@ -502,6 +507,7 @@ func (r *ItemRepository) upsert(ctx context.Context, execer itemExecer, item *mo item.TmdbID, item.TvdbID, item.PosterPath, + item.PosterSourcePath, item.PosterThumbhash, item.BackdropPath, item.BackdropThumbhash, @@ -1423,6 +1429,12 @@ func (r *ItemRepository) UpdateMetadata(ctx context.Context, contentID string, u addString("tvdb_id", upd.TvdbID) addIntArray("locked_fields", upd.LockedFields) addString("poster_path", upd.PosterPath) + if upd.PosterPath != nil { + // An explicit poster override invalidates the provider-origin source + // path captured by image caching; outbound embeds must not keep + // rendering the replaced provider artwork. + setClauses = append(setClauses, "poster_source_path = NULL") + } addString("poster_thumbhash", upd.PosterThumbhash) addString("backdrop_path", upd.BackdropPath) addString("backdrop_thumbhash", upd.BackdropThumbhash) diff --git a/internal/catalog/server_settings_repo.go b/internal/catalog/server_settings_repo.go index 104032c7..0ad114f3 100644 --- a/internal/catalog/server_settings_repo.go +++ b/internal/catalog/server_settings_repo.go @@ -45,6 +45,23 @@ func (r *ServerSettingsRepo) Set(ctx context.Context, key, value string) error { return nil } +// SetIfAbsent inserts a setting only when the key has no value yet (absent or +// empty), reporting whether this call won the write. Generated credentials +// (e.g. the web push VAPID keypair) must be provisioned single-writer across +// concurrent nodes: exactly one generated value may ever land. +func (r *ServerSettingsRepo) SetIfAbsent(ctx context.Context, key, value string) (bool, error) { + tag, err := r.pool.Exec(ctx, + `INSERT INTO server_settings (key, value) VALUES ($1, $2) + ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value + WHERE server_settings.value = ''`, + key, value, + ) + if err != nil { + return false, fmt.Errorf("server_settings set-if-absent %q: %w", key, err) + } + return tag.RowsAffected() > 0, nil +} + // GetAll retrieves all settings as a map. func (r *ServerSettingsRepo) GetAll(ctx context.Context) (map[string]string, error) { rows, err := r.pool.Query(ctx, `SELECT key, value FROM server_settings`) diff --git a/internal/discord/client.go b/internal/discord/client.go new file mode 100644 index 00000000..a0850ae1 --- /dev/null +++ b/internal/discord/client.go @@ -0,0 +1,226 @@ +// Package discord is a minimal Discord REST client covering exactly what the +// notification channel needs: the OAuth2 code exchange used for account +// linking, identity lookups, and bot DM delivery. No Gateway connection is +// held; everything is short-lived REST against a fixed trusted host. +package discord + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "strings" + "time" + + "golang.org/x/time/rate" +) + +const ( + defaultAPIBase = "https://discord.com/api/v10" + // AuthorizeURL is the user-facing OAuth2 consent page. + AuthorizeURL = "https://discord.com/oauth2/authorize" + + requestTimeout = 10 * time.Second + // Discord asks bot user agents to identify themselves in this format. + userAgent = "DiscordBot (https://github.com/Silo-Server/silo-server, 1.0)" + + // errorBodyLimit bounds how much of an error response is read for + // diagnostics. + errorBodyLimit = 4 << 10 +) + +// Sentinel errors for the failure modes callers branch on. +var ( + // ErrDMBlocked is Discord error 50007: the bot cannot DM this user. The + // user does not share a guild with the bot or has server DMs disabled. + ErrDMBlocked = errors.New("discord: cannot send messages to this user") + // ErrUnauthorized means the bot token or OAuth credentials were rejected. + ErrUnauthorized = errors.New("discord: unauthorized") + // ErrRateLimited means Discord returned 429; retry later. + ErrRateLimited = errors.New("discord: rate limited") +) + +const dmBlockedCode = 50007 + +// User is the subset of a Discord user object the integration stores. +type User struct { + ID string `json:"id"` + Username string `json:"username"` +} + +// Client makes Discord REST calls. Tokens are passed per call so callers can +// read live settings; the client itself holds no credentials. +type Client struct { + httpClient *http.Client + apiBase string + // limiter paces outbound calls well under Discord's global rate limits + // (~50 req/s global, ~5 DMs/s); notification volume is far below this, + // but digest-hour bursts across many accounts need smoothing. + limiter *rate.Limiter +} + +// NewClient creates a Client. +func NewClient() *Client { + return &Client{ + httpClient: &http.Client{Timeout: requestTimeout}, + apiBase: defaultAPIBase, + limiter: rate.NewLimiter(rate.Every(250*time.Millisecond), 4), + } +} + +// ExchangeCode performs the OAuth2 authorization-code exchange and returns +// the user access token. Only the identify scope is requested at authorize +// time, so the token can do nothing beyond reading the user's own identity. +func (c *Client) ExchangeCode(ctx context.Context, clientID, clientSecret, code, redirectURI string) (string, error) { + form := url.Values{ + "grant_type": {"authorization_code"}, + "code": {code}, + "redirect_uri": {redirectURI}, + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, + c.apiBase+"/oauth2/token", strings.NewReader(form.Encode())) + if err != nil { + return "", err + } + req.Header.Set("Content-Type", "application/x-www-form-urlencoded") + req.SetBasicAuth(clientID, clientSecret) + + var token struct { + AccessToken string `json:"access_token"` + } + if err := c.do(req, &token); err != nil { + return "", fmt.Errorf("exchange oauth code: %w", err) + } + if token.AccessToken == "" { + return "", errors.New("discord: token response missing access_token") + } + return token.AccessToken, nil +} + +// GetUser returns the identity behind a user access token (GET /users/@me). +func (c *Client) GetUser(ctx context.Context, accessToken string) (User, error) { + return c.getMe(ctx, "Bearer "+accessToken) +} + +// GetBotUser returns the bot's own identity, verifying the bot token. Used by +// the admin "test" endpoint. +func (c *Client) GetBotUser(ctx context.Context, botToken string) (User, error) { + return c.getMe(ctx, "Bot "+botToken) +} + +func (c *Client) getMe(ctx context.Context, authorization string) (User, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.apiBase+"/users/@me", nil) + if err != nil { + return User{}, err + } + req.Header.Set("Authorization", authorization) + + var user User + if err := c.do(req, &user); err != nil { + return User{}, fmt.Errorf("get current user: %w", err) + } + if user.ID == "" { + return User{}, errors.New("discord: user response missing id") + } + return user, nil +} + +// OpenDMChannel opens (or returns the existing) DM channel with the user. +// Idempotent: Discord returns the same channel for repeated calls. +func (c *Client) OpenDMChannel(ctx context.Context, botToken, recipientDiscordUserID string) (string, error) { + body, err := json.Marshal(map[string]string{"recipient_id": recipientDiscordUserID}) + if err != nil { + return "", err + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, + c.apiBase+"/users/@me/channels", bytes.NewReader(body)) + if err != nil { + return "", err + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bot "+botToken) + + var channel struct { + ID string `json:"id"` + } + if err := c.do(req, &channel); err != nil { + return "", fmt.Errorf("open dm channel: %w", err) + } + if channel.ID == "" { + return "", errors.New("discord: channel response missing id") + } + return channel.ID, nil +} + +// SendDM posts a message payload (Discord message JSON, e.g. an embeds body) +// to a DM channel. +func (c *Client) SendDM(ctx context.Context, botToken, channelID string, payload []byte) error { + req, err := http.NewRequestWithContext(ctx, http.MethodPost, + c.apiBase+"/channels/"+url.PathEscape(channelID)+"/messages", bytes.NewReader(payload)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bot "+botToken) + + if err := c.do(req, nil); err != nil { + return fmt.Errorf("send dm: %w", err) + } + return nil +} + +// do executes the request under the rate limiter and decodes a 2xx JSON +// response into out (when non-nil). Non-2xx responses map to sentinel errors +// with the Discord error message attached. +func (c *Client) do(req *http.Request, out any) error { + if err := c.limiter.Wait(req.Context()); err != nil { + return err + } + req.Header.Set("User-Agent", userAgent) + resp, err := c.httpClient.Do(req) + if err != nil { + return err + } + defer func() { _ = resp.Body.Close() }() + + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return apiError(resp) + } + if out == nil { + _, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, errorBodyLimit)) + return nil + } + if err := json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(out); err != nil { + return fmt.Errorf("decode discord response: %w", err) + } + return nil +} + +// apiError maps a non-2xx response to a sentinel error, preserving Discord's +// message for logs and UI surfacing. +func apiError(resp *http.Response) error { + body, _ := io.ReadAll(io.LimitReader(resp.Body, errorBodyLimit)) + var payload struct { + Code int `json:"code"` + Message string `json:"message"` + } + _ = json.Unmarshal(body, &payload) + + switch { + case payload.Code == dmBlockedCode: + return ErrDMBlocked + case resp.StatusCode == http.StatusUnauthorized: + return fmt.Errorf("%w: %s", ErrUnauthorized, payload.Message) + case resp.StatusCode == http.StatusTooManyRequests: + return ErrRateLimited + } + message := payload.Message + if message == "" { + message = strings.TrimSpace(string(body)) + } + return fmt.Errorf("discord: HTTP %d: %s", resp.StatusCode, message) +} diff --git a/internal/discord/client_test.go b/internal/discord/client_test.go new file mode 100644 index 00000000..0d4c3c61 --- /dev/null +++ b/internal/discord/client_test.go @@ -0,0 +1,136 @@ +package discord + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/http/httptest" + "testing" + "time" + + "golang.org/x/time/rate" +) + +func testClient(t *testing.T, handler http.Handler) *Client { + t.Helper() + server := httptest.NewServer(handler) + t.Cleanup(server.Close) + return &Client{ + httpClient: server.Client(), + apiBase: server.URL, + limiter: rate.NewLimiter(rate.Inf, 1), + } +} + +func TestExchangeCodeAndGetUser(t *testing.T) { + client := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/oauth2/token": + user, pass, ok := r.BasicAuth() + if !ok || user != "client-id" || pass != "client-secret" { + t.Errorf("missing or wrong basic auth: %s/%s", user, pass) + } + if err := r.ParseForm(); err != nil { + t.Fatal(err) + } + if r.PostForm.Get("grant_type") != "authorization_code" || r.PostForm.Get("code") != "the-code" { + t.Errorf("unexpected form: %v", r.PostForm) + } + _ = json.NewEncoder(w).Encode(map[string]string{"access_token": "user-token"}) + case "/users/@me": + if got := r.Header.Get("Authorization"); got != "Bearer user-token" { + t.Errorf("unexpected authorization %q", got) + } + _ = json.NewEncoder(w).Encode(User{ID: "42", Username: "quick"}) + default: + t.Errorf("unexpected path %s", r.URL.Path) + w.WriteHeader(http.StatusNotFound) + } + })) + + token, err := client.ExchangeCode(context.Background(), "client-id", "client-secret", "the-code", "https://silo.example/cb") + if err != nil { + t.Fatalf("exchange: %v", err) + } + user, err := client.GetUser(context.Background(), token) + if err != nil { + t.Fatalf("get user: %v", err) + } + if user.ID != "42" || user.Username != "quick" { + t.Fatalf("unexpected user %+v", user) + } +} + +func TestOpenDMChannelAndSendDM(t *testing.T) { + client := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/users/@me/channels": + if got := r.Header.Get("Authorization"); got != "Bot bot-token" { + t.Errorf("unexpected authorization %q", got) + } + var body map[string]string + _ = json.NewDecoder(r.Body).Decode(&body) + if body["recipient_id"] != "42" { + t.Errorf("unexpected recipient %q", body["recipient_id"]) + } + _ = json.NewEncoder(w).Encode(map[string]string{"id": "dm-123"}) + case "/channels/dm-123/messages": + w.WriteHeader(http.StatusOK) + default: + t.Errorf("unexpected path %s", r.URL.Path) + w.WriteHeader(http.StatusNotFound) + } + })) + + channelID, err := client.OpenDMChannel(context.Background(), "bot-token", "42") + if err != nil { + t.Fatalf("open dm: %v", err) + } + if channelID != "dm-123" { + t.Fatalf("unexpected channel id %q", channelID) + } + if err := client.SendDM(context.Background(), "bot-token", channelID, []byte(`{"content":"hi"}`)); err != nil { + t.Fatalf("send dm: %v", err) + } +} + +func TestErrorMapping(t *testing.T) { + cases := []struct { + name string + status int + body string + want error + }{ + {"dm blocked", http.StatusForbidden, `{"code":50007,"message":"Cannot send messages to this user"}`, ErrDMBlocked}, + {"bad token", http.StatusUnauthorized, `{"message":"401: Unauthorized"}`, ErrUnauthorized}, + {"rate limited", http.StatusTooManyRequests, `{"message":"You are being rate limited."}`, ErrRateLimited}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + client := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(tc.status) + _, _ = w.Write([]byte(tc.body)) + })) + err := client.SendDM(context.Background(), "bot-token", "dm-123", []byte(`{}`)) + if !errors.Is(err, tc.want) { + t.Fatalf("got %v, want %v", err, tc.want) + } + }) + } +} + +func TestLimiterHonorsContextCancel(t *testing.T) { + client := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + })) + // Exhausted limiter with a long refill forces Wait to block on the context. + client.limiter = rate.NewLimiter(rate.Every(time.Hour), 1) + client.limiter.Allow() + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond) + defer cancel() + if err := client.SendDM(ctx, "bot-token", "dm-123", []byte(`{}`)); err == nil { + t.Fatal("expected context error from limiter wait") + } +} diff --git a/internal/events/types.go b/internal/events/types.go index 2f735c27..eb324e24 100644 --- a/internal/events/types.go +++ b/internal/events/types.go @@ -16,6 +16,10 @@ const ( ChannelHistoryImport EventChannel = "history_import" ChannelUserState EventChannel = "user_state" ChannelPlugins EventChannel = "plugins" + // ChannelNotifications carries profile-scoped user notifications + // (inbox deliveries). Subscriptions require a websocket ticket binding + // the connection to a (user, profile). + ChannelNotifications EventChannel = "notifications" ) var AllChannels = []EventChannel{ @@ -27,6 +31,7 @@ var AllChannels = []EventChannel{ ChannelHistoryImport, ChannelUserState, ChannelPlugins, + ChannelNotifications, } type Envelope struct { diff --git a/internal/libraryingest/executor.go b/internal/libraryingest/executor.go index bc53e73f..f1da04d1 100644 --- a/internal/libraryingest/executor.go +++ b/internal/libraryingest/executor.go @@ -88,6 +88,7 @@ type Executor struct { skippedRootRepo SkippedRootRepository events cache.EventBus realtime *notifications.Hub + availability *notifications.AvailabilityDetector now func() time.Time // tvDrainSettleWindow overrides scopedTVDrainSettleWindow when > 0. Kept @@ -119,6 +120,15 @@ func NewExecutor( } } +// SetAvailabilityDetector wires episode-availability detection for release +// notifications. Optional; runs after matching completes and never blocks or +// fails the ingest. +func (e *Executor) SetAvailabilityDetector(detector *notifications.AvailabilityDetector) { + if e != nil { + e.availability = detector + } +} + // IngestFolder runs the full ingest workflow for an entire library. func (e *Executor) IngestFolder(ctx context.Context, folder *models.MediaFolder) (*Result, error) { return e.ingest(ctx, folder, scopeModeLibrary, "") @@ -349,6 +359,21 @@ func (e *Executor) ingest(ctx context.Context, folder *models.MediaFolder, mode } } + // Content availability runs after matching/reconcile so releases are tied + // to resolved items. It runs detached: the detector is best-effort with + // its own deadline (it detaches from scanCtx internally, surviving its + // cancellation), and a slow pass must not delay scan completion or the + // serialized scan queue. + if e.availability != nil { + kinds := notifications.AvailabilityKinds{ + Episodes: isTVLibraryType(folder.Type) || isMixedLibraryType(folder.Type), + Movies: isMovieLibraryType(folder.Type) || isMixedLibraryType(folder.Type), + } + if kinds.Episodes || kinds.Movies { + go e.availability.HandleIngestCompleted(scanCtx, folder.ID, mode == scopeModeLibrary, matchScopes, kinds) + } + } + if shouldPublish(result) && e.events != nil { if err := e.events.Publish(scanCtx, cache.ChannelCatalog, cache.Event{ Type: cache.EventScanComplete, @@ -446,6 +471,17 @@ func isMixedLibraryType(libraryType string) bool { return strings.ToLower(strings.TrimSpace(libraryType)) == "mixed" } +// isMovieLibraryType mirrors the scanner's movie library naming +// (internal/scanner/scanner.go). +func isMovieLibraryType(libraryType string) bool { + switch strings.ToLower(strings.TrimSpace(libraryType)) { + case "movie", "movies": + return true + default: + return false + } +} + func (e *Executor) scan(ctx context.Context, folder *models.MediaFolder, mode scopeMode, scopePath string) ([]string, *scanner.ScanResult, error) { switch mode { case scopeModeLibrary: diff --git a/internal/mail/layout.go b/internal/mail/layout.go new file mode 100644 index 00000000..f692ccaa --- /dev/null +++ b/internal/mail/layout.go @@ -0,0 +1,135 @@ +package mail + +import ( + "html" + "strings" +) + +// Shared visual tokens for Silo's branded emails, mirroring the web UI's +// default "Midnight Cinema" theme (web/src/app.css): a near-black canvas, +// monochrome type, and a white primary action. Feature packages compose body +// fragments with these tokens and wrap them with RenderLayout so every email +// the server sends looks like it came from the same product. +// +// Email-client constraints shape everything here: styles must be inline, +// layout must be tables, and colors must be explicit on every element (no +// inheritance through client-rewritten DOM). Web fonts don't load in most +// clients, so the stacks lead with the brand font and degrade to common +// system faces. +const ( + EmailFont = "'Outfit','Avenir Next','Segoe UI',Helvetica,Arial,sans-serif" + EmailFontMono = "'SF Mono',SFMono-Regular,Menlo,Consolas,'Liberation Mono',monospace" + + EmailColorCanvas = "#141417" // page background + EmailColorCard = "#1c1c20" // content card surface + EmailColorBorder = "#2e2e35" // card outline + EmailColorText = "#e8e8ec" // primary text + EmailColorMuted = "#9696a0" // secondary text, badges, footer + EmailColorRule = "#26262c" // hairline row separators + EmailColorAction = "#e8e8ec" // primary button background (white-on-dark) + EmailColorOnAct = "#141417" // primary button label +) + +// LayoutOptions is the content RenderLayout places into the branded shell. +type LayoutOptions struct { + // Preheader is the hidden inbox-preview snippet shown next to the subject + // line. Plain text; optional. + Preheader string + // Title is the headline at the top of the card. Plain text; optional. + Title string + // BodyHTML is the card content below the title. Trusted HTML — callers + // must escape any user-controlled values before building it. + BodyHTML string + // FooterHTML is the fine print under the card. Trusted HTML; optional. + FooterHTML string +} + +// RenderLayout wraps content in Silo's dark branded email shell: wordmark, +// content card, and footer. It adds no links of its own, so an email whose +// options carry no hrefs renders fully link-free (some features require +// that when no external URL is configured). +func RenderLayout(opts LayoutOptions) string { + preheader := "" + if opts.Preheader != "" { + // The trailing zwnj/nbsp run pads the preview so clients don't pull + // body markup into the snippet after the real preheader text. + preheader = `
` + + html.EscapeString(opts.Preheader) + + strings.Repeat(" ‌", 40) + `
` + "\n" + } + title := "" + if opts.Title != "" { + title = `

` + html.EscapeString(opts.Title) + `

` + "\n" + } + footer := "" + if opts.FooterHTML != "" { + footer = `` + opts.FooterHTML + `` + "\n" + } + + return strings.NewReplacer( + "{{preheader}}", preheader, + "{{title}}", title, + "{{body}}", opts.BodyHTML, + "{{footer}}", footer, + "{{font}}", EmailFont, + "{{canvas}}", EmailColorCanvas, + "{{card}}", EmailColorCard, + "{{border}}", EmailColorBorder, + "{{text}}", EmailColorText, + ).Replace(emailShell) +} + +// EmailButton renders the primary call-to-action: a white pill on the dark +// card, matching the web UI's primary action style. Both arguments are +// escaped here. The wrapping table keeps the button shape in Outlook, which +// ignores padding on anchors. +func EmailButton(label, href string) string { + return `` + + `
` + + `` + html.EscapeString(label) + `` + + `
` +} + +// EmailParagraph renders one body paragraph in the standard text style, +// escaping the given plain text. +func EmailParagraph(text string) string { + return `

` + html.EscapeString(text) + `

` +} + +// emailShell is the document skeleton. The color-scheme meta plus explicit +// bgcolor attributes keep dark-mode-aware clients from inverting the design; +// the small stylesheet only tightens padding on narrow screens (supported by +// Gmail/Apple Mail, harmlessly ignored elsewhere). +const emailShell = ` + + + + + + + + + +{{preheader}} + +
+ + + +{{footer}}
▸︎  SILO
+{{title}}{{body}} +
+
+ +` diff --git a/internal/mail/layout_test.go b/internal/mail/layout_test.go new file mode 100644 index 00000000..46623e73 --- /dev/null +++ b/internal/mail/layout_test.go @@ -0,0 +1,52 @@ +package mail + +import ( + "strings" + "testing" +) + +func TestRenderLayoutEscapesAndPlacesContent(t *testing.T) { + out := RenderLayout(LayoutOptions{ + Preheader: `sneak `, + Title: `Title & bold`, + BodyHTML: `

trusted

`, + FooterHTML: `fine print`, + }) + if strings.Contains(out, "` + content := composeNotificationEmail(EmailModePerEpisode, []DeliveryRow{row}, emailComposeOptions{}) + if strings.Contains(content.HTML, "