From 5809f3e431d1a7d422fd8e14cdb28cc38f433b4c Mon Sep 17 00:00:00 2001 From: npeter83 Date: Sun, 26 Jul 2026 15:32:46 +0200 Subject: [PATCH] refactor(quota): own the quota session; Plex per-section commit (R6 S2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Transaction ownership (C-A5): - quota.record_usage / log_action write in their OWN SessionLocal(), committed immediately, and no longer take/commit the caller's db. The spend stays DURABLE regardless of the caller (read-only YouTube paths never commit; a failed job rolls back) — under-counting the shared daily budget would risk overspending the real YouTube quota — AND quota stops flushing a job's partial state mid-run behind its own `except: db.rollback()`. units_used_today switches to a scalar SELECT so measured()'s before/after diff sees the separate session's commits (READ COMMITTED). Call sites updated (youtube/client.py ×3, routes/search.py ×1). - plex sync commits PER SECTION (moved the end-of-loop commit inside), so a PlexError on a later library keeps the ones already mirrored (idempotent; next run reconciles). The sync jobs did NOT rely on record_usage's incidental commit for persistence — their per-item functions (apply_rss_feed, backfill_channel_recent, import_subscriptions, sync_user_playlists, apply_channel_autotags) all commit internally, so the in-loop rollback handlers already isolate a failed item. No loop restructuring needed. DB-lane tests (test_quota_ownership, 4): spend survives a caller rollback while the caller's uncommitted row does NOT leak; accumulation; measured-style mid-txn read. Gate: ruff clean; DB lane 180 green. --- backend/app/plex/sync.py | 5 ++- backend/app/quota.py | 46 ++++++++++++++------- backend/app/routes/search.py | 2 +- backend/app/youtube/client.py | 6 +-- backend/tests/test_quota_ownership.py | 57 +++++++++++++++++++++++++++ 5 files changed, 97 insertions(+), 19 deletions(-) create mode 100644 backend/tests/test_quota_ownership.py diff --git a/backend/app/plex/sync.py b/backend/app/plex/sync.py index 058349d..f1b9e82 100644 --- a/backend/app/plex/sync.py +++ b/backend/app/plex/sync.py @@ -129,7 +129,10 @@ def sync(db: Session) -> dict: else: _sync_shows(db, plex, lib, stats) _sync_collections(db, plex, lib, stats) - db.commit() + # Commit each section as it completes: a PlexError on a LATER section then keeps the + # libraries already mirrored (idempotent — the next run reconciles the rest), instead + # of the old single end-of-loop commit that a late failure rolled back entirely. + db.commit() # Rich-fetch a bounded batch of not-yet-enriched titles' FULL cast, so the whole library # becomes filterable over successive runs (the cheap listing above only stored the top ~3). try: diff --git a/backend/app/quota.py b/backend/app/quota.py index 0e44839..25772c7 100644 --- a/backend/app/quota.py +++ b/backend/app/quota.py @@ -9,10 +9,12 @@ import contextvars from datetime import datetime, timezone from zoneinfo import ZoneInfo +from sqlalchemy import func, select from sqlalchemy.dialects.postgresql import insert as pg_insert from sqlalchemy.orm import Session from app import sysconfig +from app.db import SessionLocal from app.models import ApiQuotaUsage, QuotaEvent _PACIFIC = ZoneInfo("America/Los_Angeles") @@ -99,8 +101,13 @@ def pacific_day_start_utc() -> datetime: def units_used_today(db: Session) -> int: - row = db.get(ApiQuotaUsage, pacific_today()) - return row.units_used if row else 0 + # A scalar SELECT (not db.get) so it bypasses the caller session's identity map — quota writes + # commit in a SEPARATE session (see record_usage), and a cached ORM row would otherwise hide + # their effect from a caller that already loaded today's row (measured()'s before/after diff). + # Under READ COMMITTED a fresh SELECT sees the other session's committed writes. + return db.scalar( + select(ApiQuotaUsage.units_used).where(ApiQuotaUsage.day == pacific_today()) + ) or 0 def remaining_today(db: Session) -> int: @@ -111,21 +118,24 @@ def can_spend(db: Session, units: int) -> bool: return remaining_today(db) >= units -def log_action(db: Session, user_id: int, action: str) -> None: +def log_action(user_id: int, action: str) -> None: """Record a zero-cost action event so per-user daily caps still count it. `record_usage` only logs when units are actually charged; a scrape-based search spends no - quota yet must still be rate-limited per user, so it logs its event here directly.""" - db.add(QuotaEvent(user_id=user_id, action=action, units=0)) - db.commit() + quota yet must still be rate-limited per user, so it logs its event here directly. + + Writes in its OWN session (like record_usage) — quota accounting is durable independent of the + caller's transaction and must never commit the caller's session (that would flush a job's + partial state mid-run, defeating its own rollback). See [[quota-owns-its-session]].""" + with SessionLocal() as qs: + qs.add(QuotaEvent(user_id=user_id, action=action, units=0)) + qs.commit() def actions_today(db: Session, user_id: int, action: str) -> int: """How many quota events of `action` this user has logged so far in the current Pacific day — for per-user, per-action daily caps (e.g. the live-search limit). Counts events, not units, so it only works for actions charged exactly once per user action.""" - from sqlalchemy import func, select # local import: keep the module's import head lean - return ( db.scalar( select(func.count()) @@ -140,8 +150,15 @@ def actions_today(db: Session, user_id: int, action: str) -> int: ) -def record_usage(db: Session, units: int) -> None: - """Atomically add `units` to today's counter (upsert) and log an attribution event.""" +def record_usage(units: int) -> None: + """Atomically add `units` to today's counter (upsert) and log an attribution event. + + Uses its OWN session, committed immediately, so the spend is DURABLE regardless of what the + calling job does next: the units were really consumed at YouTube, so a caller that later rolls + back (or is read-only and never commits) must NOT lose the record — under-counting the shared + daily budget risks overspending the real API quota. Equally, quota must never commit the + CALLER's session (the old behaviour) — that persisted a job's partial state mid-run behind its + own `except: db.rollback()`. See [[quota-owns-its-session]].""" if units <= 0: return day = pacific_today() @@ -153,7 +170,8 @@ def record_usage(db: Session, units: int) -> None: set_={"units_used": ApiQuotaUsage.units_used + units}, ) ) - db.execute(stmt) - # Audit detail: who/what spent it (per-user attribution; NULL actor = system). - db.add(QuotaEvent(user_id=_actor_id.get(), action=_action.get(), units=units)) - db.commit() + with SessionLocal() as qs: + qs.execute(stmt) + # Audit detail: who/what spent it (per-user attribution; NULL actor = system). + qs.add(QuotaEvent(user_id=_actor_id.get(), action=_action.get(), units=units)) + qs.commit() diff --git a/backend/app/routes/search.py b/backend/app/routes/search.py index 375664b..7f4809a 100644 --- a/backend/app/routes/search.py +++ b/backend/app/routes/search.py @@ -290,7 +290,7 @@ def search_youtube( # at least one page (a total failure raises 502 above and never reaches this), keeps the cap # counting user searches instead of internal continuation pages. if source != "api": - quota.log_action(db, user.id, quota.QuotaAction.VIDEOS_SEARCH) + quota.log_action(user.id, quota.QuotaAction.VIDEOS_SEARCH) ordered = collected[:limit] diff --git a/backend/app/youtube/client.py b/backend/app/youtube/client.py index c52dfb0..0ccf64a 100644 --- a/backend/app/youtube/client.py +++ b/backend/app/youtube/client.py @@ -113,7 +113,7 @@ class YouTubeClient: else: headers["Authorization"] = f"Bearer {self._access_token()}" resp = self._send("GET", f"{API_BASE}/{path}", params=p, headers=headers) - quota.record_usage(self.db, cost) + quota.record_usage(cost) if resp.status_code != 200: log.warning("YouTube API %s -> %s: %s", path, resp.status_code, resp.text[:200]) raise YouTubeError(f"GET {path} -> {resp.status_code}: {resp.text[:300]}") @@ -239,7 +239,7 @@ class YouTubeClient: params={"id": subscription_id}, headers={"Authorization": f"Bearer {self._access_token()}"}, ) - quota.record_usage(self.db, 50) + quota.record_usage(50) if resp.status_code not in (200, 204): log.warning( "YouTube subscriptions.delete -> %s: %s", @@ -282,7 +282,7 @@ class YouTubeClient: resp = self._send( method, f"{API_BASE}/{path}", params=params, json=json, headers=headers ) - quota.record_usage(self.db, 50) + quota.record_usage(50) if resp.status_code not in (200, 204): log.warning("YouTube %s %s -> %s: %s", method, path, resp.status_code, resp.text[:200]) raise YouTubeError(f"{method} {path} -> {resp.status_code}: {resp.text[:300]}") diff --git a/backend/tests/test_quota_ownership.py b/backend/tests/test_quota_ownership.py new file mode 100644 index 0000000..fd7138e --- /dev/null +++ b/backend/tests/test_quota_ownership.py @@ -0,0 +1,57 @@ +"""DB-backed tests for R6 S2 quota transaction ownership: quota accounting writes in its OWN +session, so a spend is DURABLE even when the calling job rolls back (the units were really consumed +at YouTube), and it NEVER commits the caller's session (the old behaviour flushed a job's partial +state mid-run behind its own rollback). Uses the DB lane (`db` fixture) — skipped without Postgres. +""" +from sqlalchemy import func, select + +from app import quota +from app.models import QuotaEvent, User + + +def test_record_usage_survives_caller_rollback_and_leaves_caller_uncommitted(db): + # The caller has uncommitted work of its own; then a YouTube call spends quota; then the caller + # rolls back (a failed job). + db.add(User(email="pending@example.com")) # caller's own, never committed + quota.record_usage(10) # dedicated session — commits independently + db.rollback() # the job fails and rolls back + + # Quota spend PERSISTED (real units were consumed) — the guard stays accurate ... + assert quota.units_used_today(db) == 10 + assert db.scalar(select(func.count()).select_from(QuotaEvent)) == 1 + # ... but the caller's own partial work did NOT leak (quota never committed the caller's session). + assert db.scalar(select(func.count()).select_from(User).where(User.email == "pending@example.com")) == 0 + + +def test_record_usage_accumulates(db): + quota.record_usage(10) + quota.record_usage(5) + assert quota.units_used_today(db) == 15 # atomic upsert: +10 then +5 + + +def test_reads_see_separate_session_spend_mid_transaction(db): + # Mimics measured()'s before/after diff: a read, a spend (in quota's SEPARATE session), a read. + # The second read must reflect the spend even though the caller's own transaction is still open — + # units_used_today is a scalar SELECT (READ COMMITTED sees other sessions' commits), NOT a cached + # ORM row that db.get would return stale. + before = quota.units_used_today(db) + quota.record_usage(7) + after = quota.units_used_today(db) + assert after - before == 7 + + +def test_log_action_survives_caller_rollback(db): + user = User(email="searcher@example.com") + db.add(user) + db.commit() # a committed user to attribute the (free) action to + db.add(User(email="pending@example.com")) # uncommitted caller work + quota.log_action(user.id, quota.QuotaAction.VIDEOS_SEARCH) + db.rollback() + + # The zero-cost action event survived (per-user daily caps must still count it) ... + events = db.execute(select(QuotaEvent)).scalars().all() + assert len(events) == 1 + assert events[0].user_id == user.id + assert events[0].units == 0 + # ... and the caller's uncommitted user is gone. + assert db.scalar(select(func.count()).select_from(User)) == 1