From 1979e02cde24542a187242c724d938e4c950462d Mon Sep 17 00:00:00 2001 From: Zak Bearman Date: Wed, 9 Sep 2026 22:39:22 +1200 Subject: [PATCH] Add opt-in monthly email recaps with scheduling and delivery history --- backend/app/db.py | 2 + backend/app/main.py | 4 + backend/app/routers/recaps.py | 133 +++++ backend/app/services/email_recaps.py | 243 +++++++++ backend/app/services/recap_email.py | 172 ++++++ backend/app/services/recap_store.py | 235 +++++++++ backend/tests/test_email_recaps.py | 493 ++++++++++++++++++ docs/jellystat-integration.md | 25 +- frontend/app/admin/configNavigation.ts | 1 + frontend/app/admin/recaps/page.tsx | 131 +++++ frontend/app/email-recaps/page.tsx | 66 +++ frontend/app/email-recaps/recaps.css | 78 +++ frontend/app/insights/reports/page.tsx | 17 +- .../app/profile/MonthlyRecapPreference.tsx | 77 +++ frontend/app/profile/page.tsx | 2 + frontend/app/ui/ApplicationChrome.tsx | 2 +- scripts/review_email_recaps_ui.cjs | 146 ++++++ scripts/review_monthly_reports_ui.cjs | 12 +- 18 files changed, 1832 insertions(+), 7 deletions(-) create mode 100644 backend/app/routers/recaps.py create mode 100644 backend/app/services/email_recaps.py create mode 100644 backend/app/services/recap_email.py create mode 100644 backend/app/services/recap_store.py create mode 100644 backend/tests/test_email_recaps.py create mode 100644 frontend/app/admin/recaps/page.tsx create mode 100644 frontend/app/email-recaps/page.tsx create mode 100644 frontend/app/email-recaps/recaps.css create mode 100644 frontend/app/profile/MonthlyRecapPreference.tsx create mode 100644 scripts/review_email_recaps_ui.cjs diff --git a/backend/app/db.py b/backend/app/db.py index 77e1990..0b6a584 100644 --- a/backend/app/db.py +++ b/backend/app/db.py @@ -733,6 +733,8 @@ def init_db() -> None: conn.execute("PRAGMA optimize") except sqlite3.OperationalError: pass + from .services.recap_store import init_schema as init_recap_schema + init_recap_schema(conn) _backfill_auth_providers() ensure_admin_user() _backfill_request_repairs() diff --git a/backend/app/main.py b/backend/app/main.py index 887d6be..022552f 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -29,8 +29,10 @@ from .routers.portal import router as portal_router from .routers.operations import router as operations_router from .routers.insights import router as insights_router from .routers.identities import router as identities_router +from .routers.recaps import router as recaps_router from .services.jellyfin_sync import run_daily_jellyfin_sync from .services.issue_resolution import run_issue_confirmation_loop +from .services.email_recaps import run_email_recap_loop from .services.operation_progress import ( begin_operation, finish_operation, @@ -267,6 +269,7 @@ async def startup() -> None: _launch_background_task("requests-full-sync", run_daily_requests_full_sync) _launch_background_task("db-cleanup", run_daily_db_cleanup) _launch_background_task("issue-confirmation", run_issue_confirmation_loop) + _launch_background_task("email-recaps", run_email_recap_loop) logger.info("startup complete") @@ -284,3 +287,4 @@ app.include_router(portal_router) app.include_router(operations_router) app.include_router(insights_router) app.include_router(identities_router) +app.include_router(recaps_router) diff --git a/backend/app/routers/recaps.py b/backend/app/routers/recaps.py new file mode 100644 index 0000000..b6fbf56 --- /dev/null +++ b/backend/app/routers/recaps.py @@ -0,0 +1,133 @@ +from datetime import datetime, timezone +from typing import Literal +from urllib.parse import urlsplit +from uuid import UUID + +from fastapi import APIRouter, Depends, HTTPException, Query, Response +from pydantic import BaseModel, ConfigDict, Field, field_validator + +from ..auth import get_current_user, require_admin +from ..services import email_recaps as recaps, recap_store as store + + +def no_cache(response: Response): + response.headers["Cache-Control"] = "no-store" + + +router = APIRouter(tags=["email-recaps"], dependencies=[Depends(no_cache)]) + + +class StrictPayload(BaseModel): + model_config = ConfigDict(extra="forbid") + + +class Preference(StrictPayload): + enabled: bool + + +class RecapSettings(StrictPayload): + enabled: bool + day: int = Field(ge=1, le=28) + hour: int = Field(ge=0, le=23) + public_url: str = Field(max_length=500) + + @field_validator("public_url") + @classmethod + def origin_only(cls, value: str) -> str: + value = value.strip().rstrip('/') + if not value: + return value + try: + url = urlsplit(value) + port = url.port + except ValueError as exc: + raise ValueError("Enter the public Magent address, such as https://magent.example.com.") from exc + if (url.scheme not in {"http", "https"} or not url.hostname or url.username or url.password + or url.path or url.query or url.fragment or any(char.isspace() or ord(char) < 33 for char in value) + or any(char in value for char in '<>"\\') or (port is not None and port < 1)): + raise ValueError("Enter a http(s) Magent address without a path, credentials or query.") + return value + + +class TestEmail(StrictPayload): + month: str | None = Field(default=None, pattern=r"^[0-9]{4}-[0-9]{2}$") + request_id: UUID + + +class TokenAction(StrictPayload): + token: str = Field(min_length=40, max_length=100, pattern=r"^[A-Za-z0-9_-]+$") + action: Literal["confirm", "unsubscribe"] + + +def error(exc: recaps.RecapError): + raise HTTPException(exc.status, exc.detail) from exc + + +@router.get("/profile/email-recaps") +def preferences(user: dict = Depends(get_current_user)) -> dict: + try: + return recaps.preferences(user) + except recaps.RecapError as exc: + error(exc) + + +@router.put("/profile/email-recaps") +async def preference(payload: Preference, user: dict = Depends(get_current_user)) -> dict: + try: + if payload.enabled: + return await recaps.subscribe(user) + store.disable(recaps.current_account(user)["id"]) + return recaps.preferences(user) + except recaps.RecapError as exc: + error(exc) + + +@router.post("/email-recaps/check") +def check_token(payload: TokenAction) -> dict: + try: + return recaps.token_action(payload.token, payload.action) + except recaps.RecapError as exc: + error(exc) + + +@router.post("/email-recaps/confirm") +def apply_token(payload: TokenAction) -> dict: + try: + return recaps.token_action(payload.token, payload.action, apply=True) + except recaps.RecapError as exc: + error(exc) + + +@router.get("/admin/email-recaps") +def overview(offset: int = Query(default=0, ge=0), user: dict = Depends(require_admin)) -> dict: + ready, detail = recaps.delivery_ready() + months = recaps.month_periods(None, datetime.now(timezone.utc))["available_months"][1:] + return {"settings": store.settings(), "ready": ready, "detail": detail, "months": months, + "worker_enabled": recaps.worker_enabled(), **store.history(offset=offset)} + + +@router.put("/admin/email-recaps") +def settings(payload: RecapSettings, user: dict = Depends(require_admin)) -> dict: + if payload.enabled: + # Validate against the proposed URL without writing any partial settings. + ready, detail = recaps.smtp_email_config_ready() + runtime = recaps.get_runtime_settings() + if not payload.public_url or not ready or not recaps.worker_enabled() or not runtime.jellystat_base_url or not runtime.jellystat_api_key: + raise HTTPException(409, "Set the public address, enable SMTP email and connect Jellystat before starting the schedule." if ready else detail) + return store.save_settings(payload.model_dump(), datetime.now(timezone.utc)) + + +@router.get("/admin/email-recaps/preview") +async def preview(month: str | None = Query(default=None, max_length=7, pattern=r"^[0-9]{4}-[0-9]{2}$"), user: dict = Depends(require_admin)) -> dict: + try: + return await recaps.preview(user, month) + except recaps.RecapError as exc: + error(exc) + + +@router.post("/admin/email-recaps/test", status_code=202) +def test_email(payload: TestEmail, user: dict = Depends(require_admin)) -> dict: + try: + return recaps.queue_test(user, payload.month, str(payload.request_id)) + except recaps.RecapError as exc: + error(exc) diff --git a/backend/app/services/email_recaps.py b/backend/app/services/email_recaps.py new file mode 100644 index 0000000..6aeb35b --- /dev/null +++ b/backend/app/services/email_recaps.py @@ -0,0 +1,243 @@ +"""Opt-in monthly recaps. Scheduling and delivery are safe to run in multiple workers.""" + +import asyncio +import logging +import os +import time +import uuid +from datetime import datetime, timezone +from urllib.parse import urlencode + +from .. import db +from ..clients.jellystat import HistoryLimitError, JellystatError +from ..runtime import get_runtime_settings +from . import recap_email as mail, recap_store as store +from .invite_email import smtp_email_config_ready +from .jellyfin_identity import linked_user_id, source_key +from .monthly_reports import get_monthly_report, month_periods + +logger = logging.getLogger(__name__) + + +class RecapError(Exception): + def __init__(self, detail: str, status: int = 409): + self.detail, self.status = detail, status + super().__init__(detail) + + +def worker_enabled() -> bool: + return os.environ.get("BACKGROUND_TASKS_ENABLED", "true").lower() != "false" + + +def delivery_ready() -> tuple[bool, str]: + config = store.settings() + if not config["public_url"]: + return False, "Set the public Magent address for email links." + ready, detail = smtp_email_config_ready() + if not ready: + return False, detail + runtime = get_runtime_settings() + if not runtime.jellystat_base_url or not runtime.jellystat_api_key: + return False, "Connect Jellystat to generate viewing recaps." + if not worker_enabled(): + return False, "Background automation is paused on this server." + return True, "Email delivery is configured." + + +def current_account(user: dict) -> dict: + account = db.get_user_by_username(user.get("username", "")) + if not account or account.get("is_blocked") or account.get("is_expired"): + raise RecapError("This account cannot receive viewing recaps.", 403) + return account + + +def binding_matches(sub: dict, account: dict) -> bool: + runtime = get_runtime_settings() + return bool(account and not account.get("is_blocked") and not account.get("is_expired") + and mail.valid_email(account.get("email")) + and account["email"].strip().casefold() == sub["email"].strip().casefold() + and source_key(runtime.jellyfin_base_url) == sub["identity_source"] + and linked_user_id(account["username"], runtime.jellyfin_base_url) == sub["identity_id"]) + + +def active_subscription(account: dict) -> dict | None: + sub = store.subscription(account["id"]) + if sub and sub["state"] != "off" and not binding_matches(sub, account): + store.disable(account["id"]) + sub = store.subscription(account["id"]) + return sub + + +def preferences(user: dict) -> dict: + account = current_account(user) + sub = active_subscription(account) + config = store.settings() + ready, detail = delivery_ready() + runtime = get_runtime_settings() + linked = bool(linked_user_id(account["username"], runtime.jellyfin_base_url)) + email = mail.valid_email(account.get("email")) + state = sub["state"] if sub else "off" + if state == "pending" and sub["confirmation_expires"] <= time.time(): + state = "expired" + return {"state": state, "email": account.get("email"), "can_subscribe": ready and linked and bool(email), + "detail": detail if not ready else "Save a valid email address in your profile." if not email else + "Your Jellyfin account needs a saved identity link." if not linked else "Your monthly story, in your inbox.", + "schedule_enabled": config["enabled"], "next_send_at": config["next_send_at"], + "day": config["day"], "hour": config["hour"], "timezone": "UTC", + "resend_after": (sub["requested_at"] + 300) if sub else None} + + +async def subscribe(user: dict) -> dict: + account = current_account(user) + preference = preferences(user) + if preference["state"] == "enabled": + return preference + if not preference["can_subscribe"]: + raise RecapError(preference["detail"]) + config = store.settings() + runtime = get_runtime_settings() + try: + token = store.request_confirmation(account, source_key(runtime.jellyfin_base_url), + linked_user_id(account["username"], runtime.jellyfin_base_url), time.time()) + except ValueError as exc: + raise RecapError(str(exc), 429) from exc + url = f"{config['public_url']}/email-recaps#" + urlencode({"action": "confirm", "token": token}) + rendered = mail.render_confirmation(account["username"], url) + try: + await asyncio.to_thread(mail.send_email, account["email"].strip(), rendered, + mail.message_id(uuid.uuid4().hex, config["public_url"])) + except mail.DeliveryError as exc: + raise RecapError("Could not confirm delivery of the verification email. Check your inbox; you can request another in five minutes.", 502) from exc + return {**preferences(user), "message": "Check your inbox and confirm within 24 hours to turn on monthly recaps."} + + +def token_action(token: str, action: str, *, apply: bool = False) -> dict: + sub = store.token_subscription(token, action) + if not sub: + raise RecapError("This email link is invalid or has already been used. Open Profile to manage your recaps.", 410) + if action == "unsubscribe": + if apply: + store.disable(sub["user_id"]) + return {"action": action, "state": "off" if apply or sub["state"] == "off" else "ready"} + account = db.get_user_by_id(sub["user_id"]) + if (sub["state"] != "pending" or sub["confirmation_expires"] <= time.time() + or not binding_matches(sub, account)): + raise RecapError("This confirmation has expired or your account details changed. Request a new link from Profile.", 410) + if apply and not store.confirm(sub, time.time()): + raise RecapError("This confirmation is no longer available. Request a new link from Profile.", 410) + return {"action": action, "state": "enabled" if apply else "ready"} + + +def completed_month(month: str | None) -> str: + try: + period = month_periods(month, datetime.now(timezone.utc)) + except ValueError as exc: + raise RecapError(str(exc), 422) from exc + if period["is_partial"]: + raise RecapError("Choose a completed month for an email recap.", 422) + return period["month"] + + +async def preview(user: dict, month: str | None) -> dict: + account = current_account(user) + selected = completed_month(month) + config = store.settings() + if not config["public_url"]: + raise RecapError("Save the public Magent address before previewing an email.") + try: + report = await asyncio.wait_for(get_monthly_report(account, selected), timeout=180) + except HistoryLimitError as exc: + raise RecapError("This report exceeds Jellystat's history limit. No partial recap was generated.", 422) from exc + except (JellystatError, TimeoutError) as exc: + raise RecapError("Your report is temporarily unavailable. Please try again shortly.", 502) from exc + if report["state"] != "ready": + raise RecapError("Connect Jellystat and link your Jellyfin account to preview your recap.") + return {"month": selected, "email": account.get("email"), **mail.render_recap( + report, account["username"], config["public_url"], config["public_url"] + "/profile#monthly-recaps")} + + +def queue_test(user: dict, month: str | None, request_id: str) -> dict: + account = current_account(user) + ready, detail = delivery_ready() + if not ready: + raise RecapError(detail) + sub = active_subscription(account) + if not sub or sub["state"] != "enabled": + raise RecapError("Turn on email recaps and confirm your email in Profile before sending a personal test.") + selected = completed_month(month) + try: + delivery_id = store.enqueue_test(sub, selected, request_id, store.settings()["public_url"], time.time()) + except ValueError as exc: + raise RecapError(str(exc), 429) from exc + return {"id": delivery_id, "message": "Test queued for your confirmed email. Check delivery history for the result."} + + +def eligible_delivery(delivery: dict) -> tuple[dict, dict]: + account = db.get_user_by_id(delivery["user_id"]) + sub = active_subscription(account) if account else None + config = store.settings() + ready, _ = delivery_ready() + if (not ready or not sub or sub["state"] != "enabled" or sub["version"] != delivery["subscription_version"] + or sub["email"] != delivery["email"] or not binding_matches(sub, account) + or config["public_url"] != delivery["public_url"] + or (delivery["kind"] == "scheduled" and not config["enabled"])): + raise mail.DeliveryCancelled() + return account, sub + + +async def process_delivery(delivery: dict) -> None: + state, detail, delay = "failed", "Could not prepare the recap. Check the report and email settings.", 0 + try: + account, sub = eligible_delivery(delivery) + report = await asyncio.wait_for(get_monthly_report(account, delivery["month"]), timeout=180) + if report["state"] != "ready" or report["is_partial"]: + raise mail.DeliveryError("failed", "A complete personal report is not available.") + unsubscribe = f"{delivery['public_url']}/email-recaps#" + urlencode({"action": "unsubscribe", "token": sub["unsubscribe_token"]}) + rendered = mail.render_recap(report, account["username"], delivery["public_url"], unsubscribe, test=delivery["kind"] == "test") + + def before_data(): + eligible_delivery(delivery) + if not store.begin_sending(delivery, time.time()): + raise mail.DeliveryCancelled() + + await asyncio.to_thread(mail.send_email, delivery["email"], rendered, + mail.message_id(delivery["id"], delivery["public_url"]), before_data) + state, detail = "sent", "Accepted by the mail server." + except mail.DeliveryCancelled: + state, detail = "cancelled", "Consent, account details or email configuration changed." + except HistoryLimitError: + state, detail = "failed", "Jellystat's history limit was reached. No partial recap was sent." + except (JellystatError, TimeoutError): + state, detail = "retry", "Viewing history is temporarily unavailable." + except mail.DeliveryError as exc: + state, detail = exc.state, exc.detail + except Exception as exc: + # Do not expose provider errors or private report content in history/logs. + logger.error("recap delivery error id=%s type=%s", delivery["id"], type(exc).__name__) + row = store.read_one("SELECT state FROM email_recap_deliveries WHERE id=?", (delivery["id"],)) + if row and row["state"] == "sending": + state, detail = "unknown", "Delivery outcome is unknown; check the mail server." + if state == "retry": + if delivery["attempts"] >= 3: + state, detail = "failed", detail + " Stopped after three attempts." + else: + delay = 300 if delivery["attempts"] == 1 else 1800 + store.finish(delivery, state, detail, time.time(), delay) + + +async def run_once() -> None: + store.enqueue_due(datetime.now(timezone.utc)) + for _ in range(10): + delivery = store.claim_delivery(time.time()) + if not delivery: + break + await process_delivery(delivery) + + +async def run_email_recap_loop() -> None: + while True: + try: + await run_once() + except Exception as exc: + logger.error("email recap worker failed type=%s", type(exc).__name__) + await asyncio.sleep(30) diff --git a/backend/app/services/recap_email.py b/backend/app/services/recap_email.py new file mode 100644 index 0000000..d2af859 --- /dev/null +++ b/backend/app/services/recap_email.py @@ -0,0 +1,172 @@ +"""Personal recap email rendering and SMTP delivery with explicit acceptance tracking.""" + +import html +import re +import smtplib +import ssl +from contextlib import suppress +from datetime import datetime +from email.message import EmailMessage +from email.policy import SMTP as SMTP_POLICY +from email.utils import formataddr, formatdate +from urllib.parse import urlsplit + +from ..runtime import get_runtime_settings + + +class DeliveryError(Exception): + def __init__(self, state: str, detail: str): + self.state, self.detail = state, detail + super().__init__(detail) + + +class DeliveryCancelled(Exception): + pass + + +def valid_email(value: str | None) -> str | None: + value = str(value or "").strip() + if (len(value) <= 254 and re.fullmatch(r"[^@\s<>;,\"\\]+@[^@\s<>;,\"\\]+\.[^@\s<>;,\"\\]+", value) + and all(32 < ord(char) < 127 for char in value)): + return value + return None + + +def month_label(value: str) -> str: + return datetime.strptime(value, "%Y-%m").strftime("%B %Y") + + +def number(value: float) -> str: + return f"{value:,.0f}" + + +def document(*, title: str, intro: str, content: str, action: str, url: str, footer: str) -> str: + esc = html.escape + return f'''{esc(title)} + +
+ + + + + +
MAGENT / YOUR MONTH IN VIEWING

{esc(title)}

{esc(intro)}

{content}
{esc(action)} ↗
{footer}
+
''' + + +def render_confirmation(username: str, url: str) -> dict: + title = "Your month, delivered." + intro = f"Hi {username}, confirm this email address to receive your personal monthly viewing recap from Magent." + text = f"{intro}\n\nConfirm email recaps: {url}\n\nThis link expires in 24 hours. If you did not request this, ignore this email. No viewing history will be emailed until you confirm." + body = document(title=title, intro=intro, + content='

Minutes watched, movies, episodes, your longest run and requests — with a link to your full monthly report.

', + action="Confirm email recaps", url=url, + footer="This link expires in 24 hours. If you did not request this, ignore this email.
No viewing history will be emailed until you confirm.") + return {"subject": "Confirm your Magent email recaps", "body_text": text, "body_html": body} + + +def render_recap(report: dict, username: str, public_url: str, unsubscribe_url: str, *, test: bool = False) -> dict: + esc = html.escape + month = month_label(report["month"]) + previous = month_label(report["comparison_month"]) + summary = report["summary"] + metrics = (("Minutes watched", "minutes", summary["minutes"]), ("Movies played", "movies", summary["movies"]), + ("Episodes played", "episodes", summary["episodes"]), ("Requests made", "requests", report["requests"]["total"])) + cells, lines = [], [] + for label, key, value in metrics: + change = report["changes"][key] + difference = change["difference"] + comparison = ("No change" if difference == 0 else f"{'+' if difference > 0 else '−'}{number(abs(difference))}") + if change["percent"] is not None and difference: + comparison += f" ({'+' if difference > 0 else '−'}{abs(change['percent']):g}%)" + comparison += f" from {previous}" + lines.append(f"{label}: {number(value)}. {comparison}.") + cells.append(f'{label}
{number(value)}{esc(comparison)}') + content = '' + ''.join(cells[:2]) + '' + ''.join(cells[2:]) + '' + habit = f"{number(summary['active_days'])} days watched · {number(summary['longest_streak'])}-day longest run" + content += f'

{esc(habit)}

' + top = report.get("top_titles", [])[:3] + if top: + content += '

Your most watched

' + for item in top: + content += f'

{esc(item["title"])}
{number(item["minutes"])} minutes · {number(item["plays"])} plays

' + else: + content += '

No viewing was recorded this month. Your requests are still included.

' + report_url = f"{public_url}/insights/reports?month={report['month']}" + intro = f"Hi {username}, here’s your {month} in viewing. A little look back at the stories you spent time with." + footer = f'You opted in to personal monthly recaps from Magent.
Based on retained Jellystat history. Calendar months use UTC; request statuses are current.
Unsubscribe from recaps · Email preferences' + if test: + intro = "This is your test recap. " + intro + body = document(title=month, intro=intro, content=content, action="Explore your full report", url=report_url, footer=footer) + text = '\n'.join([intro, '', *lines, '', habit, '', 'Most watched:', + *(f"{item['title']}: {number(item['minutes'])} minutes" for item in top), '', + f"Your full report: {report_url}", '', 'Based on retained Jellystat history. Calendar months use UTC; request statuses are current.', + f"Unsubscribe from recaps: {unsubscribe_url}", f"Email preferences: {public_url}/profile#monthly-recaps"]) + return {"subject": f"{'[Test] ' if test else ''}Your {month} in viewing · Magent", "body_text": text, "body_html": body} + + +def send_email(recipient: str, rendered: dict, message_id: str, before_data=lambda: None) -> None: + """Return only after SMTP accepts DATA. Never retry an ambiguous DATA disconnect. + + A stable Message-ID aids diagnosis; it is not an SMTP deduplication guarantee. + See RFC 5321 §4.5.3.2.6 and Python's smtplib exception definitions. + """ + runtime = get_runtime_settings() + sender = valid_email(runtime.magent_notify_email_from_address) + if not sender or not valid_email(recipient): + raise DeliveryError("failed", "A valid sender and recipient email are required.") + message = EmailMessage(policy=SMTP_POLICY) + message["From"] = formataddr((str(runtime.magent_notify_email_from_name or "Magent").replace('\r', '').replace('\n', ''), sender)) + message["To"], message["Subject"] = recipient, rendered["subject"] + message["Date"], message["Message-ID"] = formatdate(localtime=False), message_id + message["Auto-Submitted"], message["X-Auto-Response-Suppress"] = "auto-generated", "All" + message.set_content(rendered["body_text"]) + message.add_alternative(rendered["body_html"], subtype="html") + payload = message.as_bytes() + smtp, stage = None, "connect" + try: + kwargs = {"timeout": 30, "local_hostname": sender.split('@', 1)[1]} + if runtime.magent_notify_email_use_ssl: + smtp = smtplib.SMTP_SSL(runtime.magent_notify_email_smtp_host, runtime.magent_notify_email_smtp_port, + context=ssl.create_default_context(), **kwargs) + else: + smtp = smtplib.SMTP(runtime.magent_notify_email_smtp_host, runtime.magent_notify_email_smtp_port, **kwargs) + smtp.ehlo_or_helo_if_needed() + if runtime.magent_notify_email_use_tls and not runtime.magent_notify_email_use_ssl: + smtp.starttls(context=ssl.create_default_context()) + smtp.ehlo() + if runtime.magent_notify_email_smtp_username: + smtp.login(runtime.magent_notify_email_smtp_username, runtime.magent_notify_email_smtp_password) + code, reply = smtp.mail(sender) + if code != 250: + raise smtplib.SMTPResponseException(code, reply) + code, reply = smtp.rcpt(recipient) + if code not in (250, 251): + raise smtplib.SMTPResponseException(code, reply) + before_data() + stage = "data" + code, reply = smtp.data(payload) + if code != 250: + raise smtplib.SMTPDataError(code, reply) + stage = "accepted" + except smtplib.SMTPResponseException as exc: + state = "retry" if 400 <= exc.smtp_code < 500 else "failed" + raise DeliveryError(state, f"Mail server returned SMTP {exc.smtp_code}.") from exc + except (ssl.SSLError, smtplib.SMTPNotSupportedError, UnicodeError, ValueError) as exc: + raise DeliveryError("failed", "Check the SMTP security and sender settings.") from exc + except (OSError, smtplib.SMTPException) as exc: + state = "unknown" if stage == "data" else "retry" + detail = "Mail server acceptance is unknown; check its logs before taking further action." if state == "unknown" else "Could not reach or finish connecting to the mail server." + raise DeliveryError(state, detail) from exc + finally: + if smtp: + # A failed QUIT after a 250 DATA response must not turn an accepted email into a retry. + with suppress(Exception): + smtp.quit() + with suppress(Exception): + smtp.close() + + +def message_id(delivery_id: str, public_url: str) -> str: + host = urlsplit(public_url).hostname or "magent.local" + return f"" diff --git a/backend/app/services/recap_store.py b/backend/app/services/recap_store.py new file mode 100644 index 0000000..9e916a2 --- /dev/null +++ b/backend/app/services/recap_store.py @@ -0,0 +1,235 @@ +"""Durable consent, schedule and delivery records for personal email recaps.""" + +import hashlib +import secrets +import sqlite3 +import uuid +from contextlib import closing, contextmanager +from datetime import datetime + +from .. import db +from .monthly_reports import shift_month + + +def init_schema(conn: sqlite3.Connection) -> None: + for statement in ( + """CREATE TABLE IF NOT EXISTS email_recap_settings ( + id INTEGER PRIMARY KEY CHECK (id = 1), enabled INTEGER NOT NULL DEFAULT 0, + day INTEGER NOT NULL DEFAULT 2, hour INTEGER NOT NULL DEFAULT 9, + public_url TEXT NOT NULL DEFAULT '', next_send_at REAL)""", + "INSERT OR IGNORE INTO email_recap_settings (id) VALUES (1)", + """CREATE TABLE IF NOT EXISTS email_recap_subscriptions ( + user_id INTEGER PRIMARY KEY, state TEXT NOT NULL, email TEXT NOT NULL, + identity_source TEXT NOT NULL, identity_id TEXT NOT NULL, version TEXT NOT NULL, + confirmation_hash TEXT UNIQUE, confirmation_expires REAL, requested_at REAL NOT NULL, + confirmed_at REAL, unsubscribe_token TEXT NOT NULL UNIQUE)""", + """CREATE TABLE IF NOT EXISTS email_recap_deliveries ( + id TEXT PRIMARY KEY, dedupe_key TEXT NOT NULL UNIQUE, user_id INTEGER NOT NULL, + month TEXT NOT NULL, kind TEXT NOT NULL, email TEXT NOT NULL, + subscription_version TEXT NOT NULL, public_url TEXT NOT NULL, + state TEXT NOT NULL DEFAULT 'queued', attempts INTEGER NOT NULL DEFAULT 0, + created_at REAL NOT NULL, updated_at REAL NOT NULL, next_attempt_at REAL NOT NULL, + claim TEXT, lease_until REAL, detail TEXT NOT NULL DEFAULT '')""", + "CREATE INDEX IF NOT EXISTS idx_email_recap_queue ON email_recap_deliveries (state, next_attempt_at)", + """CREATE TRIGGER IF NOT EXISTS email_recap_account_changed AFTER UPDATE OF email, is_blocked ON users + WHEN LOWER(TRIM(COALESCE(NEW.email, ''))) != LOWER(TRIM(COALESCE(OLD.email, ''))) + OR NEW.is_blocked = 1 + BEGIN UPDATE email_recap_subscriptions SET state = 'off', confirmation_hash = NULL, + confirmed_at = NULL WHERE user_id = NEW.id; END""", + """CREATE TRIGGER IF NOT EXISTS email_recap_account_deleted AFTER DELETE ON users + BEGIN DELETE FROM email_recap_subscriptions WHERE user_id = OLD.id; + UPDATE email_recap_deliveries SET state = 'cancelled', detail = 'Account removed.' + WHERE user_id = OLD.id AND state IN ('queued', 'retry', 'preparing'); END""", + """CREATE TRIGGER IF NOT EXISTS email_recap_identity_changed AFTER UPDATE ON jellyfin_user_links + WHEN NEW.jellyfin_user_id != OLD.jellyfin_user_id OR NEW.source != OLD.source + OR NEW.local_user_id != OLD.local_user_id + BEGIN UPDATE email_recap_subscriptions SET state = 'off', confirmation_hash = NULL, + confirmed_at = NULL WHERE user_id = OLD.local_user_id; END""", + """CREATE TRIGGER IF NOT EXISTS email_recap_identity_deleted AFTER DELETE ON jellyfin_user_links + BEGIN UPDATE email_recap_subscriptions SET state = 'off', confirmation_hash = NULL, + confirmed_at = NULL WHERE user_id = OLD.local_user_id; END""", + ): + conn.execute(statement) + + +@contextmanager +def transaction(): + with closing(db._connect()) as conn, conn: + conn.row_factory = sqlite3.Row + conn.execute("BEGIN IMMEDIATE") + yield conn + + +def read_one(sql: str, args=()) -> dict | None: + with closing(db._connect()) as conn: + conn.row_factory = sqlite3.Row + row = conn.execute(sql, args).fetchone() + return dict(row) if row else None + + +def settings() -> dict: + row = read_one("SELECT * FROM email_recap_settings WHERE id = 1") + return {key: (bool(value) if key == "enabled" else value) for key, value in row.items() if key != "id"} + + +def next_due(now: datetime, day: int, hour: int) -> datetime: + due = shift_month(now, 0).replace(day=day, hour=hour) + return due if due > now else shift_month(now, 1).replace(day=day, hour=hour) + + +def save_settings(values: dict, now: datetime) -> dict: + with transaction() as conn: + old = dict(conn.execute("SELECT * FROM email_recap_settings WHERE id = 1").fetchone()) + changed = any(old[key] != values[key] for key in ("day", "hour", "public_url")) + due = old["next_send_at"] + if not values["enabled"]: + due = None + elif not old["enabled"] or changed: + due = next_due(now, values["day"], values["hour"]).timestamp() + conn.execute("UPDATE email_recap_settings SET enabled=?, day=?, hour=?, public_url=?, next_send_at=? WHERE id=1", + (values["enabled"], values["day"], values["hour"], values["public_url"], due)) + if not values["enabled"] or changed: + conn.execute("""UPDATE email_recap_deliveries SET state='cancelled', detail='Schedule paused or changed.', updated_at=? + WHERE kind='scheduled' AND state IN ('queued', 'retry', 'preparing')""", (now.timestamp(),)) + return settings() + + +def subscription(user_id: int) -> dict | None: + return read_one("SELECT * FROM email_recap_subscriptions WHERE user_id=?", (user_id,)) + + +def disable(user_id: int) -> None: + with transaction() as conn: + conn.execute("UPDATE email_recap_subscriptions SET state='off', confirmation_hash=NULL, confirmed_at=NULL WHERE user_id=?", (user_id,)) + conn.execute("""UPDATE email_recap_deliveries SET state='cancelled', detail='Email recaps turned off.' + WHERE user_id=? AND state IN ('queued', 'retry', 'preparing')""", (user_id,)) + + +def request_confirmation(user: dict, source: str, identity: str, now: float) -> str: + token = secrets.token_urlsafe(32) + with transaction() as conn: + old = conn.execute("SELECT * FROM email_recap_subscriptions WHERE user_id=?", (user["id"],)).fetchone() + if old and old["requested_at"] > now - 300: + raise ValueError("Please wait five minutes before requesting another confirmation email.") + conn.execute("""INSERT INTO email_recap_subscriptions + (user_id, state, email, identity_source, identity_id, version, confirmation_hash, + confirmation_expires, requested_at, confirmed_at, unsubscribe_token) + VALUES (?, 'pending', ?, ?, ?, ?, ?, ?, ?, NULL, ?) + ON CONFLICT(user_id) DO UPDATE SET state='pending', email=excluded.email, + identity_source=excluded.identity_source, identity_id=excluded.identity_id, version=excluded.version, + confirmation_hash=excluded.confirmation_hash, confirmation_expires=excluded.confirmation_expires, + requested_at=excluded.requested_at, confirmed_at=NULL, unsubscribe_token=excluded.unsubscribe_token""", + (user["id"], user["email"].strip(), source, identity, uuid.uuid4().hex, + hashlib.sha256(token.encode()).hexdigest(), now + 86400, now, secrets.token_urlsafe(32))) + return token + + +def token_subscription(token: str, action: str) -> dict | None: + if action == "confirm": + return read_one("SELECT * FROM email_recap_subscriptions WHERE confirmation_hash=?", + (hashlib.sha256(token.encode()).hexdigest(),)) + return read_one("SELECT * FROM email_recap_subscriptions WHERE unsubscribe_token=?", (token,)) + + +def confirm(sub: dict, now: float) -> bool: + with transaction() as conn: + # Recheck address and blocked state in the same transaction as the consent write. + result = conn.execute("""UPDATE email_recap_subscriptions SET state='enabled', confirmed_at=?, confirmation_hash=NULL + WHERE user_id=? AND version=? AND state='pending' AND confirmation_expires>? + AND EXISTS (SELECT 1 FROM users WHERE users.id=user_id AND is_blocked=0 + AND LOWER(TRIM(users.email))=LOWER(TRIM(email_recap_subscriptions.email)))""", + (now, sub["user_id"], sub["version"], now)) + return result.rowcount == 1 + + +def _enqueue(conn, sub: dict, month: str, kind: str, key: str, public_url: str, now: float) -> str: + delivery_id = uuid.uuid4().hex + conn.execute("""INSERT OR IGNORE INTO email_recap_deliveries + (id, dedupe_key, user_id, month, kind, email, subscription_version, public_url, + created_at, updated_at, next_attempt_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", + (delivery_id, key, sub["user_id"], month, kind, sub["email"], sub["version"], public_url, now, now, now)) + return conn.execute("SELECT id FROM email_recap_deliveries WHERE dedupe_key=?", (key,)).fetchone()[0] + + +def enqueue_test(sub: dict, month: str, request_id: str, public_url: str, now: float) -> str: + key = f"test:{sub['user_id']}:{request_id}" + with transaction() as conn: + existing = conn.execute("SELECT id FROM email_recap_deliveries WHERE dedupe_key=?", (key,)).fetchone() + if existing: + return existing[0] + recent = conn.execute("SELECT 1 FROM email_recap_deliveries WHERE user_id=? AND kind='test' AND created_at>?", + (sub["user_id"], now - 300)).fetchone() + if recent: + raise ValueError("Please wait five minutes between test emails.") + return _enqueue(conn, sub, month, "test", key, public_url, now) + + +def enqueue_due(now: datetime) -> int: + with transaction() as conn: + config = dict(conn.execute("SELECT * FROM email_recap_settings WHERE id=1").fetchone()) + if not config["enabled"] or not config["next_send_at"] or config["next_send_at"] > now.timestamp(): + return 0 + # After long downtime, send only the latest due recap; never backfill a pile of old emails. + due = shift_month(now, 0).replace(day=config["day"], hour=config["hour"]) + if due > now: + due = shift_month(now, -1).replace(day=config["day"], hour=config["hour"]) + month = shift_month(due, -1).strftime("%Y-%m") + subs = conn.execute("SELECT * FROM email_recap_subscriptions WHERE state='enabled' AND confirmed_at<=?", (due.timestamp(),)).fetchall() + before = conn.total_changes + for sub in subs: + _enqueue(conn, dict(sub), month, "scheduled", f"scheduled:{sub['user_id']}:{month}", config["public_url"], now.timestamp()) + count = conn.total_changes - before + conn.execute("UPDATE email_recap_settings SET next_send_at=? WHERE id=1", + (next_due(now, config["day"], config["hour"]).timestamp(),)) + return count + + +def claim_delivery(now: float) -> dict | None: + with transaction() as conn: + # A crashed worker could already have handed DATA to SMTP. Do not resend it automatically. + conn.execute("""UPDATE email_recap_deliveries SET state='unknown', detail='Delivery interrupted after sending began; check the mail server.', updated_at=? + WHERE state='sending' AND lease_until=3 THEN 'failed' ELSE 'retry' END, + next_attempt_at=?, updated_at=?, detail='Report preparation interrupted.' + WHERE state='preparing' AND lease_until bool: + with transaction() as conn: + # Consent may have changed while the report or SMTP connection was being prepared. + result = conn.execute("""UPDATE email_recap_deliveries SET state='sending', updated_at=?, lease_until=? + WHERE id=? AND claim=? AND state='preparing' + AND EXISTS (SELECT 1 FROM email_recap_subscriptions s JOIN users u ON u.id=s.user_id + JOIN jellyfin_user_links j ON j.local_user_id=u.id AND j.source=s.identity_source + WHERE s.user_id=email_recap_deliveries.user_id AND s.state='enabled' + AND s.version=email_recap_deliveries.subscription_version AND u.is_blocked=0 + AND LOWER(TRIM(u.email))=LOWER(TRIM(s.email)) AND j.jellyfin_user_id=s.identity_id) + AND EXISTS (SELECT 1 FROM email_recap_settings c WHERE c.id=1 AND c.public_url=email_recap_deliveries.public_url + AND (email_recap_deliveries.kind='test' OR c.enabled=1))""", (now, now + 1800, delivery["id"], delivery["claim"])) + return result.rowcount == 1 + + +def finish(delivery: dict, state: str, detail: str, now: float, delay: int = 0) -> None: + with transaction() as conn: + conn.execute("""UPDATE email_recap_deliveries SET state=?, detail=?, updated_at=?, next_attempt_at=?, lease_until=NULL + WHERE id=? AND claim=? AND state IN ('preparing', 'sending')""", + (state, detail, now, now + delay, delivery["id"], delivery["claim"])) + + +def history(limit: int = 50, offset: int = 0) -> dict: + with closing(db._connect()) as conn: + conn.row_factory = sqlite3.Row + rows = conn.execute("""SELECT d.id, d.month, d.kind, d.email, d.state, d.attempts, d.created_at, d.updated_at, + d.next_attempt_at, d.detail, u.username FROM email_recap_deliveries d LEFT JOIN users u ON u.id=d.user_id + ORDER BY d.created_at DESC, d.id LIMIT ? OFFSET ?""", (limit, offset)).fetchall() + total = conn.execute("SELECT COUNT(*) FROM email_recap_deliveries").fetchone()[0] + subscribers = conn.execute("SELECT COUNT(*) FROM email_recap_subscriptions WHERE state='enabled'").fetchone()[0] + return {"deliveries": [dict(row) for row in rows], "total": total, "subscribers": subscribers} diff --git a/backend/tests/test_email_recaps.py b/backend/tests/test_email_recaps.py new file mode 100644 index 0000000..0bdab22 --- /dev/null +++ b/backend/tests/test_email_recaps.py @@ -0,0 +1,493 @@ +import asyncio +import json +import re +import smtplib +import socketserver +import threading +import time +import unittest +from concurrent.futures import ThreadPoolExecutor +from datetime import datetime, timedelta, timezone +from email import policy +from email.parser import BytesParser +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock, patch +from urllib.parse import parse_qs, urlsplit + +from fastapi import FastAPI +from fastapi.testclient import TestClient + +from backend.app import db +from backend.app.auth import get_current_user +from backend.app.clients.jellystat import HistoryLimitError, JellystatError +from backend.app.routers import recaps as router +from backend.app.services import email_recaps as recaps, recap_email as mail, recap_store as store +from backend.app.services.jellyfin_identity import link_user, source_key +from backend.app.services.monthly_reports import change, month_periods, shift_month +from backend.tests.test_backend_quality import TempDatabaseMixin + + +def fixture_report(): + periods = month_periods(None, datetime.now(timezone.utc)) + summary = dict(minutes=1500, movies=8, episodes=24, plays=35, active_days=20, longest_streak=6) + changes = {key: change(value, round(value / 2)) for key, value in summary.items()} + changes['requests'] = change(3, 2) + return {**periods, 'state': 'ready', 'summary': summary, 'changes': changes, 'requests': {'total': 3}, + 'top_titles': [{'title': 'Severance', 'type': 'series', 'minutes': 460, 'plays': 10}, + {'title': 'Arrival', 'type': 'movie', 'minutes': 116, 'plays': 1}], + 'recent': [{'artwork_url': '/insights/artwork/SECRET?token=PRIVATE-TOKEN'}]} + + +def runtime(): + return SimpleNamespace(jellyfin_base_url='http://jellyfin', jellystat_base_url='http://jellystat', + jellystat_api_key='PRIVATE-STATS-KEY', magent_notify_enabled=True, magent_notify_email_enabled=True, + magent_notify_email_smtp_host='127.0.0.1', magent_notify_email_smtp_port=1, + magent_notify_email_smtp_username='', magent_notify_email_smtp_password='', + magent_notify_email_from_address='magent@example.test', magent_notify_email_from_name='Magent', + magent_notify_email_use_tls=False, magent_notify_email_use_ssl=False) + + +class RecapFixture(TempDatabaseMixin): + def setUp(self): + super().setUp() + db.create_user('viewer', 'Example-Password123!', role='admin', email='viewer@example.test') + link_user('viewer', 'jf-viewer', 'http://jellyfin') + self.user = db.get_user_by_username('viewer') + self.runtime = runtime() + for target, name, value in [(recaps, 'get_runtime_settings', self.runtime), (mail, 'get_runtime_settings', self.runtime), + (recaps, 'smtp_email_config_ready', (True, 'ok'))]: + mocked = patch.object(target, name, return_value=value) + mocked.start(); self.addCleanup(mocked.stop) + env = patch.dict('os.environ', {'BACKGROUND_TASKS_ENABLED': 'true'}) + env.start(); self.addCleanup(env.stop) + self.config = dict(enabled=False, day=2, hour=9, public_url='https://beta.example.test') + store.save_settings(self.config, datetime.now(timezone.utc)) + self.report = fixture_report() + + def subscribe(self, timestamp=None): + now = time.time() if timestamp is None else timestamp + token = store.request_confirmation(self.user, source_key('http://jellyfin'), 'jf-viewer', now) + sub = store.subscription(self.user['id']) + self.assertTrue(store.confirm(sub, now + 1)) + return store.subscription(self.user['id']), token + + def queue(self, sub=None, request_id='request-1'): + if sub is None: + sub, _ = self.subscribe() + return store.enqueue_test(sub, self.report['month'], request_id, self.config['public_url'], time.time()) + + def delivery(self, delivery_id): + return store.read_one('SELECT * FROM email_recap_deliveries WHERE id=?', (delivery_id,)) + + +class RecapConsentTests(RecapFixture, unittest.IsolatedAsyncioTestCase): + async def test_opt_in_only_emails_confirmation_and_check_link_does_not_confirm(self): + with patch.object(mail, 'send_email') as sender, patch.object(recaps, 'get_monthly_report') as report: + result = await recaps.subscribe(self.user) + self.assertEqual(result['state'], 'pending') + report.assert_not_called() + recipient, rendered, _ = sender.call_args.args + self.assertEqual(recipient, 'viewer@example.test') + self.assertNotIn('Severance', rendered['body_html']) + url = re.search(r'https://[^\s]+', rendered['body_text']).group(0) + token = parse_qs(urlsplit(url).fragment)['token'][0] + self.assertNotIn(token, store.subscription(self.user['id'])['confirmation_hash']) + self.assertEqual(recaps.token_action(token, 'confirm')['state'], 'ready') + self.assertEqual(store.subscription(self.user['id'])['state'], 'pending') + self.assertEqual(recaps.token_action(token, 'confirm', apply=True)['state'], 'enabled') + with self.assertRaises(recaps.RecapError): + recaps.token_action(token, 'confirm', apply=True) + with self.assertRaises(recaps.RecapError): + recaps.token_action(token, 'unsubscribe', apply=True) + + async def test_confirmation_failure_is_pending_and_resend_is_rate_limited(self): + with patch.object(mail, 'send_email', side_effect=mail.DeliveryError('unknown', 'unknown')): + with self.assertRaises(recaps.RecapError) as exc: + await recaps.subscribe(self.user) + self.assertEqual(exc.exception.status, 502) + self.assertEqual(recaps.preferences(self.user)['state'], 'pending') + with patch.object(mail, 'send_email') as sender: + with self.assertRaises(recaps.RecapError) as exc: + await recaps.subscribe(self.user) + self.assertEqual(exc.exception.status, 429) + sender.assert_not_called() + + def test_unsubscribe_is_public_idempotent_and_cancels_queued_email(self): + sub, _ = self.subscribe() + delivery_id = self.queue(sub) + token = sub['unsubscribe_token'] + self.assertEqual(recaps.token_action(token, 'unsubscribe')['state'], 'ready') + self.assertEqual(self.delivery(delivery_id)['state'], 'queued') + recaps.token_action(token, 'unsubscribe', apply=True) + self.assertEqual(recaps.token_action(token, 'unsubscribe', apply=True)['state'], 'off') + self.assertEqual(self.delivery(delivery_id)['state'], 'cancelled') + + def test_expired_confirmation_does_not_subscribe(self): + token = store.request_confirmation(self.user, source_key('http://jellyfin'), 'jf-viewer', time.time() - 90000) + self.assertEqual(recaps.preferences(self.user)['state'], 'expired') + with self.assertRaises(recaps.RecapError): + recaps.token_action(token, 'confirm', apply=True) + + def test_email_change_back_does_not_restore_consent(self): + self.subscribe() + db.set_user_email('viewer', 'changed@example.test') + db.set_user_email('viewer', 'viewer@example.test') + self.assertEqual(recaps.preferences(self.user)['state'], 'off') + + def test_changed_link_or_source_requires_new_consent(self): + self.subscribe() + with store.transaction() as conn: + conn.execute("UPDATE jellyfin_user_links SET jellyfin_user_id='new-identity' WHERE local_user_id=?", (self.user['id'],)) + self.assertEqual(recaps.preferences(self.user)['state'], 'off') + with store.transaction() as conn: + conn.execute("UPDATE email_recap_subscriptions SET state='enabled'") + self.runtime.jellyfin_base_url = 'http://other-jellyfin' + self.assertEqual(recaps.preferences(self.user)['state'], 'off') + + def test_missing_email_or_stored_identity_cannot_subscribe(self): + db.set_user_email('viewer', None) + self.assertFalse(recaps.preferences(self.user)['can_subscribe']) + db.set_user_email('viewer', 'viewer@example.test') + with store.transaction() as conn: + conn.execute('DELETE FROM jellyfin_user_links') + self.assertFalse(recaps.preferences(self.user)['can_subscribe']) + + def test_confirmation_rechecks_email_atomically(self): + store.request_confirmation(self.user, source_key('http://jellyfin'), 'jf-viewer', time.time()) + old = store.subscription(self.user['id']) + db.set_user_email('viewer', 'different@example.test') + self.assertFalse(store.confirm(old, time.time())) + + +class RecapScheduleTests(RecapFixture, unittest.TestCase): + def test_defaults_are_paused_and_no_users_are_opted_in(self): + self.assertFalse(store.settings()['enabled']) + self.assertEqual(store.history()['subscribers'], 0) + self.assertEqual(store.enqueue_due(datetime.now(timezone.utc)), 0) + + def test_utc_next_send_month_end_leap_year_and_new_year(self): + for now, expected in [ + (datetime(2026, 12, 31, tzinfo=timezone.utc), '2027-01-02T09:00:00+00:00'), + (datetime(2024, 2, 29, tzinfo=timezone.utc), '2024-03-02T09:00:00+00:00'), + (datetime(2026, 9, 2, 8, tzinfo=timezone.utc), '2026-09-02T09:00:00+00:00'), + (datetime(2026, 9, 2, 9, tzinfo=timezone.utc), '2026-10-02T09:00:00+00:00')]: + self.assertEqual(store.next_due(now, 2, 9).isoformat(), expected) + + def test_schedule_catches_up_once_and_excludes_late_subscribers(self): + before = datetime(2026, 8, 30, tzinfo=timezone.utc) + self.subscribe(before.timestamp()) + config = store.save_settings({**self.config, 'enabled': True}, before) + self.assertEqual(config['next_send_at'], datetime(2026, 9, 2, 9, tzinfo=timezone.utc).timestamp()) + db.create_user('late', 'Example-Password123!', email='late@example.test') + late = db.get_user_by_username('late') + store.request_confirmation(late, 'source', 'late-id', datetime(2026, 9, 2, 10, tzinfo=timezone.utc).timestamp()) + store.confirm(store.subscription(late['id']), datetime(2026, 9, 2, 11, tzinfo=timezone.utc).timestamp()) + now = datetime(2026, 9, 5, tzinfo=timezone.utc) + with ThreadPoolExecutor(max_workers=4) as pool: + counts = list(pool.map(store.enqueue_due, [now] * 4)) + self.assertEqual(sum(counts), 1) + rows = store.history()['deliveries'] + self.assertEqual(len(rows), 1) + self.assertEqual(rows[0]['month'], '2026-08') + self.assertEqual(rows[0]['email'], 'viewer@example.test') + # Revisit the same due date after a restart: the durable unique key still wins. + with store.transaction() as conn: + conn.execute('UPDATE email_recap_settings SET next_send_at=?', (config['next_send_at'],)) + self.assertEqual(store.enqueue_due(now), 0) + + def test_long_downtime_does_not_backfill_multiple_months(self): + before = datetime(2026, 5, 1, tzinfo=timezone.utc) + self.subscribe(before.timestamp()) + store.save_settings({**self.config, 'enabled': True}, before) + self.assertEqual(store.enqueue_due(datetime(2026, 9, 9, tzinfo=timezone.utc)), 1) + self.assertEqual(store.history()['deliveries'][0]['month'], '2026-08') + + def test_enable_after_due_date_waits_and_pause_cancels_pending_monthlies(self): + now = datetime(2026, 9, 9, tzinfo=timezone.utc) + self.subscribe(now.timestamp()) + result = store.save_settings({**self.config, 'enabled': True}, now) + self.assertEqual(result['next_send_at'], datetime(2026, 10, 2, 9, tzinfo=timezone.utc).timestamp()) + self.assertEqual(store.enqueue_due(now), 0) + store.enqueue_due(datetime(2026, 10, 3, tzinfo=timezone.utc)) + store.save_settings(self.config, now) + self.assertEqual(store.history()['deliveries'][0]['state'], 'cancelled') + self.assertIsNone(store.settings()['next_send_at']) + + +class RecapDeliveryTests(RecapFixture, unittest.IsolatedAsyncioTestCase): + async def run_claim(self): + delivery = store.claim_delivery(time.time()) + self.assertIsNotNone(delivery) + await recaps.process_delivery(delivery) + + async def test_private_report_is_delivered_once_using_confirmed_account(self): + delivery_id = self.queue() + sent = [] + def capture(recipient, rendered, message_id, before_data): + before_data() + self.assertEqual(self.delivery(delivery_id)['state'], 'sending') + sent.append((recipient, rendered, message_id)) + with patch.object(recaps, 'get_monthly_report', new=AsyncMock(return_value=self.report)) as report, patch.object(mail, 'send_email', side_effect=capture): + await recaps.run_once() + await recaps.run_once() + self.assertEqual(len(sent), 1) + self.assertEqual(sent[0][0], 'viewer@example.test') + self.assertIn(f'?month={self.report["month"]}', sent[0][1]['body_html']) + self.assertNotIn('PRIVATE-TOKEN', json.dumps(sent)) + self.assertEqual(report.await_args.args[0]['id'], self.user['id']) + self.assertEqual(self.delivery(delivery_id)['state'], 'sent') + self.assertNotIn('unsubscribe_token', json.dumps(store.history())) + + def test_concurrent_claim_and_test_deduplication(self): + sub, _ = self.subscribe() + with ThreadPoolExecutor(max_workers=4) as pool: + ids = list(pool.map(lambda _: self.queue(sub), range(4))) + rows = list(pool.map(lambda _: store.claim_delivery(time.time()), range(4))) + self.assertEqual(len(set(ids)), 1) + self.assertEqual(sum(row is not None for row in rows), 1) + with self.assertRaises(ValueError): + self.queue(sub, 'another-click') + + async def test_unsubscribe_or_email_change_during_report_prevents_sending(self): + delivery_id = self.queue() + async def report(*args): + db.set_user_email('viewer', 'other@example.test') + return self.report + def transport(recipient, rendered, message_id, before_data): + before_data() + self.fail('Private data must not reach SMTP DATA after an address change') + with patch.object(recaps, 'get_monthly_report', side_effect=report), patch.object(mail, 'send_email', side_effect=transport): + await self.run_claim() + self.assertEqual(self.delivery(delivery_id)['state'], 'cancelled') + + async def test_blocked_expired_and_deleted_accounts_are_not_sent(self): + for kind in ['blocked', 'expired', 'deleted']: + with self.subTest(kind=kind): + # Each subcase starts with a fresh account and confirmed subscription. + db.create_user(kind, 'Example-Password123!', email=f'{kind}@example.test') + account = db.get_user_by_username(kind) + link_user(kind, f'jf-{kind}', 'http://jellyfin') + store.request_confirmation(account, source_key('http://jellyfin'), f'jf-{kind}', time.time()) + store.confirm(store.subscription(account['id']), time.time()) + delivery_id = self.queue(store.subscription(account['id']), kind) + with store.transaction() as conn: + if kind == 'blocked': conn.execute('UPDATE users SET is_blocked=1 WHERE id=?', (account['id'],)) + elif kind == 'expired': conn.execute("UPDATE users SET expires_at='2000-01-01T00:00:00+00:00' WHERE id=?", (account['id'],)) + else: conn.execute('DELETE FROM users WHERE id=?', (account['id'],)) + with patch.object(mail, 'send_email') as sender, patch.object(recaps, 'get_monthly_report') as report: + await recaps.run_once() + sender.assert_not_called(); report.assert_not_called() + self.assertEqual(self.delivery(delivery_id)['state'], 'cancelled') + + async def test_known_temporary_failure_retries_three_times_with_stable_id(self): + delivery_id = self.queue() + with patch.object(recaps, 'get_monthly_report', new=AsyncMock(return_value=self.report)), patch.object(mail, 'send_email', side_effect=mail.DeliveryError('retry', 'SMTP 451')) as sender: + for attempt in range(1, 4): + await self.run_claim() + row = self.delivery(delivery_id) + self.assertEqual(row['attempts'], attempt) + self.assertEqual(row['state'], 'failed' if attempt == 3 else 'retry') + if attempt < 3: + self.assertGreater(row['next_attempt_at'], time.time() + 250) + with store.transaction() as conn: + conn.execute('UPDATE email_recap_deliveries SET next_attempt_at=0 WHERE id=?', (delivery_id,)) + self.assertEqual(len(set(call.args[2] for call in sender.call_args_list)), 1) + self.assertIsNone(store.claim_delivery(time.time())) + + async def test_ambiguous_smtp_failure_never_automatically_retries(self): + delivery_id = self.queue() + with patch.object(recaps, 'get_monthly_report', new=AsyncMock(return_value=self.report)), patch.object(mail, 'send_email', side_effect=mail.DeliveryError('unknown', 'Check mail logs')): + await self.run_claim() + self.assertEqual(self.delivery(delivery_id)['state'], 'unknown') + self.assertIsNone(store.claim_delivery(time.time() + 86400)) + + def test_stale_worker_claims_are_recovered_without_resending_uncertain_mail(self): + delivery_id = self.queue() + first = store.claim_delivery(time.time()) + second = store.claim_delivery(time.time() + 1801) + self.assertNotEqual(first['claim'], second['claim']) + self.assertFalse(store.begin_sending(first, time.time())) + self.assertTrue(store.begin_sending(second, time.time())) + store.claim_delivery(time.time() + 1801) + self.assertEqual(self.delivery(delivery_id)['state'], 'unknown') + store.finish(first, 'sent', 'Old worker', time.time()) + self.assertEqual(self.delivery(delivery_id)['state'], 'unknown') + + async def test_partial_or_over_limit_report_is_not_emailed(self): + delivery_id = self.queue() + with patch.object(recaps, 'get_monthly_report', new=AsyncMock(side_effect=HistoryLimitError('limit'))), patch.object(mail, 'send_email') as sender: + await self.run_claim() + sender.assert_not_called() + self.assertEqual(self.delivery(delivery_id)['state'], 'failed') + + +class RecapApiTests(RecapFixture, unittest.TestCase): + def setUp(self): + super().setUp() + app = FastAPI() + app.include_router(router.router) + self.app = app + self.client = TestClient(app) + self.addCleanup(self.client.close) + + def login(self, role='admin'): + self.app.dependency_overrides[get_current_user] = lambda: {**self.user, 'role': role} + + def test_authentication_roles_and_recipient_override(self): + self.assertEqual(self.client.get('/admin/email-recaps').status_code, 401) + self.assertEqual(self.client.get('/profile/email-recaps').status_code, 401) + self.login('user') + self.assertEqual(self.client.get('/admin/email-recaps').status_code, 403) + self.assertEqual(self.client.get('/admin/email-recaps/preview').status_code, 403) + self.assertEqual(self.client.post('/admin/email-recaps/test', json={}).status_code, 403) + self.login() + result = self.client.get('/admin/email-recaps') + self.assertEqual(result.status_code, 200) + self.assertEqual(result.headers['cache-control'], 'no-store') + self.assertNotIn('PRIVATE-STATS-KEY', result.text) + result = self.client.post('/admin/email-recaps/test', json={'request_id': 'c49b0c52-4528-4c1d-8c78-57aafeb24f58', 'recipient_email': 'other@example.test'}) + self.assertEqual(result.status_code, 422) + result = self.client.put('/profile/email-recaps', json={'enabled': False, 'user_id': 5}) + self.assertEqual(result.status_code, 422) + + def test_url_and_schedule_validation_do_not_write_partial_settings(self): + self.login() + for value in ['javascript:alert(1)', 'https://user:secret@example.test', 'https://example.test/path', 'https://example.test?token=secret', 'https://example.test#token', 'https://example.test:0', 'https://example.test\\evil']: + result = self.client.put('/admin/email-recaps', json={**self.config, 'public_url': value}) + self.assertEqual(result.status_code, 422, value) + for field, value in [('day', 0), ('day', 29), ('hour', 24)]: + self.assertEqual(self.client.put('/admin/email-recaps', json={**self.config, field: value}).status_code, 422) + with patch.object(recaps, 'smtp_email_config_ready', return_value=(False, 'Email is disabled.')): + self.assertEqual(self.client.put('/admin/email-recaps', json={**self.config, 'enabled': True}).status_code, 409) + self.assertEqual(store.settings()['public_url'], self.config['public_url']) + self.assertFalse(store.settings()['enabled']) + + def test_preview_uses_own_report_and_test_requires_confirmed_email(self): + self.login() + with patch.object(recaps, 'get_monthly_report', new=AsyncMock(return_value=self.report)) as report, patch.object(mail, 'send_email') as sender: + result = self.client.get('/admin/email-recaps/preview') + self.assertEqual(result.status_code, 200) + self.assertEqual(report.await_args.args[0]['id'], self.user['id']) + self.assertNotIn('PRIVATE-TOKEN', result.text) + sender.assert_not_called() + payload = {'request_id': 'c49b0c52-4528-4c1d-8c78-57aafeb24f58', 'month': self.report['month']} + self.assertEqual(self.client.post('/admin/email-recaps/test', json=payload).status_code, 409) + self.subscribe() + with patch.object(mail, 'send_email') as sender: + first = self.client.post('/admin/email-recaps/test', json=payload) + second = self.client.post('/admin/email-recaps/test', json=payload) + self.assertEqual(first.status_code, 202) + self.assertEqual(first.json()['id'], second.json()['id']) + sender.assert_not_called() + + def test_partial_month_test_rejected_and_public_get_does_not_mutate(self): + self.login(); sub, token = self.subscribe() + result = self.client.post('/admin/email-recaps/test', json={'request_id': 'c49b0c52-4528-4c1d-8c78-57aafeb24f58', 'month': datetime.now(timezone.utc).strftime('%Y-%m')}) + self.assertEqual(result.status_code, 422) + self.assertEqual(self.client.get('/email-recaps/confirm').status_code, 405) + result = self.client.post('/email-recaps/check', json={'action': 'unsubscribe', 'token': sub['unsubscribe_token']}) + self.assertEqual(result.status_code, 200) + self.assertEqual(store.subscription(self.user['id'])['state'], 'enabled') + + +class RecapEmailTests(unittest.TestCase): + def setUp(self): + self.runtime = runtime() + patched = patch.object(mail, 'get_runtime_settings', return_value=self.runtime) + patched.start(); self.addCleanup(patched.stop) + self.rendered = mail.render_recap(fixture_report(), 'Viewer', 'https://beta.example.test', 'https://beta.example.test/email-recaps#action=unsubscribe&token=fixture') + + def fake_smtp(self): + smtp = MagicMock() + smtp.mail.return_value = (250, b'OK') + smtp.rcpt.return_value = (250, b'OK') + smtp.data.return_value = (250, b'Accepted') + return smtp + + def test_render_escapes_names_and_titles_and_includes_no_artwork_credentials(self): + report = fixture_report() + report['top_titles'][0]['title'] = '' + rendered = mail.render_recap(report, '', 'https://beta.example.test', 'https://beta.example.test/email-recaps#token=example') + self.assertNotIn('