From 2ce19dc96edc233821519e2ba39202768ad776bd Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 17 Jul 2026 19:32:21 -0500 Subject: [PATCH 01/31] fix: harden API cost webhook outbox retries --- .../backend/services/api_cost_monitor.py | 465 ++++++++++++++---- 1 file changed, 378 insertions(+), 87 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index b5cbd56ac..e24309a40 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -8,9 +8,11 @@ """ import asyncio +import contextvars import json import logging import os +import random import threading import time from collections import defaultdict, deque @@ -25,6 +27,7 @@ Column, DateTime, Float, + Index, Integer, String, UniqueConstraint, @@ -32,6 +35,7 @@ create_engine, delete, func, + or_, text, ) from sqlalchemy.orm import declarative_base, sessionmaker @@ -49,6 +53,12 @@ # Configure logging logger = logging.getLogger(__name__) +# Keep the public webhook helper's one-argument signature for existing callers and +# tests while attaching an outbox event identifier to each delivery attempt. +_WEBHOOK_EVENT_ID: contextvars.ContextVar[Optional[str]] = contextvars.ContextVar( + "api_cost_webhook_event_id", default=None +) + # SQLAlchemy Base Declarative Base = declarative_base() @@ -92,12 +102,14 @@ class WebhookOutbox(Base): ) # pending, processing, sent, failed retry_count = Column(Integer, nullable=False, default=0) last_attempt = Column(DateTime, nullable=True) + next_attempt_at = Column(DateTime, nullable=True) error_message = Column(String, nullable=True) current_cost = Column(Float, nullable=False) payload = Column(String, nullable=True) __table_args__ = ( UniqueConstraint("utc_date", "alert_type", name="uq_utc_date_alert_type"), + Index("ix_webhook_outbox_due", "status", "next_attempt_at", "retry_count"), ) @@ -191,6 +203,22 @@ def __init__(self, db_path: Optional[str] = None): # Webhook notification settings self.webhook_url = os.getenv("API_COST_WEBHOOK_URL") + self.webhook_max_attempts = 5 + self.webhook_retry_base_seconds = max( + 0.0, float(os.getenv("API_COST_WEBHOOK_RETRY_BASE_SECONDS", "5")) + ) + self.webhook_retry_max_seconds = max( + self.webhook_retry_base_seconds, + float(os.getenv("API_COST_WEBHOOK_RETRY_MAX_SECONDS", "300")), + ) + self.webhook_poll_interval_seconds = max( + 0.01, float(os.getenv("API_COST_WEBHOOK_POLL_SECONDS", "1")) + ) + self.webhook_stale_timeout_seconds = max( + 1, int(os.getenv("API_COST_WEBHOOK_STALE_SECONDS", "30")) + ) + self._worker_task: Optional[asyncio.Task] = None + self._worker_wake_event: Optional[asyncio.Event] = None # Rate limiters for different services self.rate_limiters = { @@ -241,19 +269,53 @@ def _init_database(self): db_parent.mkdir(parents=True, exist_ok=True) Base.metadata.create_all(self.engine) - - # Trigger stale webhook recoveries asynchronously - try: - loop = asyncio.get_running_loop() - loop.create_task(self.recover_stale_deliveries()) - except RuntimeError: - threading.Thread( - target=lambda: asyncio.run(self.recover_stale_deliveries()), - daemon=True, - ).start() + self._upgrade_sqlite_outbox_schema() except Exception as e: logger.error(f"Failed to initialize cost monitoring database: {e}") + raise + + def _upgrade_sqlite_outbox_schema(self) -> None: + """Add scheduling state to existing SQLite databases without data loss.""" + if self.engine.dialect.name != "sqlite": + return + + # Serialize the inspect/ALTER/index sequence. A deferred transaction lets + # two processes both observe the missing column before either writes; + # BEGIN EXCLUSIVE makes the second process inspect only after the first + # migration commits. + with self.engine.connect() as connection: + connection.exec_driver_sql("BEGIN EXCLUSIVE") + try: + columns = { + row[1] + for row in connection.exec_driver_sql( + "PRAGMA table_info(webhook_outbox)" + ) + } + if "next_attempt_at" not in columns: + connection.exec_driver_sql( + "ALTER TABLE webhook_outbox ADD COLUMN next_attempt_at DATETIME" + ) + index_columns = [ + row[2] + for row in connection.exec_driver_sql( + "PRAGMA index_info(ix_webhook_outbox_due)" + ) + ] + expected_columns = ["status", "next_attempt_at", "retry_count"] + if index_columns and index_columns != expected_columns: + connection.exec_driver_sql( + "DROP INDEX IF EXISTS ix_webhook_outbox_due" + ) + connection.exec_driver_sql( + "CREATE INDEX IF NOT EXISTS ix_webhook_outbox_due " + "ON webhook_outbox (status, next_attempt_at, retry_count)" + ) + connection.commit() + except BaseException: + connection.rollback() + raise def check_rate_limit(self, service: str) -> tuple[bool, int]: """ @@ -512,37 +574,137 @@ def _claim_alert( session.close() def _trigger_delivery(self): - """Trigger non-blocking, asynchronous delivery of pending/failed alerts.""" + """Wake the explicitly managed worker without spawning per-alert tasks.""" + if self._worker_wake_event is not None: + self._worker_wake_event.set() + + async def start(self) -> asyncio.Task: + """Start the monitor's single managed outbox worker.""" + if self._worker_task is not None and not self._worker_task.done(): + return self._worker_task + + self._worker_wake_event = asyncio.Event() + self._worker_task = asyncio.create_task( + self._outbox_worker(), name="api-cost-webhook-outbox" + ) + return self._worker_task + + async def close(self) -> None: + """Stop the managed worker and wait for any claim cleanup to finish.""" + task = self._worker_task + if task is None: + return + + task.cancel() try: - loop = asyncio.get_running_loop() - loop.create_task(self.process_outbox()) - except RuntimeError: - threading.Thread( - target=lambda: asyncio.run(self.process_outbox()), daemon=True - ).start() - - async def recover_stale_deliveries(self, stale_timeout_seconds: int = 30): - """Recover items stuck in 'processing' status (due to crashes, task cancellation, etc.)""" + await task + except asyncio.CancelledError: + pass + finally: + if self._worker_task is task: + self._worker_task = None + self._worker_wake_event = None + + async def _outbox_worker(self) -> None: + """Continuously deliver due outbox items until explicitly closed.""" + while True: + wake_event = self._worker_wake_event + if wake_event is None: + return + wake_event.clear() + try: + await self.process_outbox() + except asyncio.CancelledError: + raise + except Exception: + logger.exception("Unhandled error in API-cost webhook outbox worker") + + try: + await asyncio.wait_for( + wake_event.wait(), timeout=self.webhook_poll_interval_seconds + ) + except asyncio.TimeoutError: + pass + + def _retry_at(self, attempt: int, now: datetime) -> datetime: + """Return a bounded exponential equal-jitter retry timestamp.""" + exponential_cap = min( + self.webhook_retry_max_seconds, + self.webhook_retry_base_seconds * (2 ** max(0, attempt - 1)), + ) + half_cap = exponential_cap / 2 + delay = half_cap + random.uniform(0, half_cap) + return now + timedelta(seconds=delay) + + def _retry_state( + self, attempt: int, now: datetime, error_message: str + ) -> tuple[Optional[datetime], str]: + """Return persisted scheduling and error state for a failed attempt.""" + if attempt >= self.webhook_max_attempts: + return ( + None, + f"Retry exhausted after {self.webhook_max_attempts} attempts: " + f"{error_message}", + ) + return self._retry_at(attempt, now), error_message + + async def recover_stale_deliveries( + self, stale_timeout_seconds: Optional[int] = None + ) -> None: + """Recover processing claims abandoned by a crash or cancellation.""" + if stale_timeout_seconds is None: + stale_timeout_seconds = self.webhook_stale_timeout_seconds + session = self.Session() try: - cutoff = datetime.now(timezone.utc) - timedelta( - seconds=stale_timeout_seconds - ) + now = datetime.now(timezone.utc) + cutoff = now - timedelta(seconds=stale_timeout_seconds) stale_items = ( session.query(WebhookOutbox) .filter( WebhookOutbox.status == "processing", - WebhookOutbox.last_attempt < cutoff, + or_( + WebhookOutbox.last_attempt.is_(None), + WebhookOutbox.last_attempt < cutoff, + ), ) .all() ) for item in stale_items: - item.status = "failed" - item.error_message = "Recovery: Stale/Crashed delivery task recovered" - logger.info( - f"Recovered stale webhook delivery {item.id} for {item.utc_date} ({item.alert_type})" + next_attempt_at, recovery_error = self._retry_state( + max(1, item.retry_count), + now, + "Recovery: Stale/Crashed delivery task recovered", ) + filters = [ + WebhookOutbox.id == item.id, + WebhookOutbox.status == "processing", + ] + if item.last_attempt is None: + filters.append(WebhookOutbox.last_attempt.is_(None)) + else: + filters.append(WebhookOutbox.last_attempt == item.last_attempt) + + recovered = ( + session.query(WebhookOutbox) + .filter(*filters) + .update( + { + WebhookOutbox.status: "failed", + WebhookOutbox.next_attempt_at: next_attempt_at, + WebhookOutbox.error_message: recovery_error, + }, + synchronize_session=False, + ) + ) + if recovered: + logger.info( + "Recovered stale webhook delivery %s for %s (%s)", + item.id, + item.utc_date, + item.alert_type, + ) session.commit() except Exception as e: @@ -554,75 +716,193 @@ async def recover_stale_deliveries(self, stale_timeout_seconds: int = 30): finally: session.close() - async def process_outbox(self): - """Process pending/failed outbox deliveries with bounded retry and crash recovery.""" - await self.recover_stale_deliveries() + def _try_claim_outbox_item( + self, + item_id: int, + claim_time: datetime, + respect_schedule: bool = True, + ) -> Optional[dict[str, Any]]: + """Claim one due item with a single compare-and-swap UPDATE.""" + session = self.Session() + try: + filters = [ + WebhookOutbox.id == item_id, + WebhookOutbox.status.in_(["pending", "failed"]), + WebhookOutbox.retry_count < self.webhook_max_attempts, + ] + if respect_schedule: + filters.append( + or_( + WebhookOutbox.next_attempt_at.is_(None), + WebhookOutbox.next_attempt_at <= claim_time, + ) + ) + + claimed = ( + session.query(WebhookOutbox) + .filter(*filters) + .update( + { + WebhookOutbox.status: "processing", + WebhookOutbox.retry_count: WebhookOutbox.retry_count + 1, + WebhookOutbox.last_attempt: claim_time, + WebhookOutbox.next_attempt_at: None, + }, + synchronize_session=False, + ) + ) + if claimed != 1: + session.rollback() + return None + + session.commit() + item = session.query(WebhookOutbox).filter_by(id=item_id).one() + return { + "id": item.id, + "payload": item.payload, + "utc_date": item.utc_date, + "alert_type": item.alert_type, + "retry_count": item.retry_count, + "last_attempt": item.last_attempt, + } + except Exception as e: + logger.debug("Could not claim webhook outbox item %s: %s", item_id, e) + try: + session.rollback() + except Exception: + pass + return None + finally: + session.close() + def _complete_outbox_claim( + self, + claim: dict[str, Any], + *, + success: bool, + error_message: Optional[str] = None, + ) -> bool: + """Conditionally complete exactly the attempt represented by ``claim``.""" session = self.Session() try: - items = ( + values: dict[Any, Any] + if success: + values = { + WebhookOutbox.status: "sent", + WebhookOutbox.next_attempt_at: None, + WebhookOutbox.error_message: None, + } + else: + next_attempt_at, persisted_error = self._retry_state( + claim["retry_count"], + datetime.now(timezone.utc), + error_message or "Delivery failed", + ) + values = { + WebhookOutbox.status: "failed", + WebhookOutbox.next_attempt_at: next_attempt_at, + WebhookOutbox.error_message: persisted_error, + } + + completed = ( session.query(WebhookOutbox) .filter( - WebhookOutbox.status.in_(["pending", "failed"]), - WebhookOutbox.retry_count < 5, # Bounded retry limit of 5 + WebhookOutbox.id == claim["id"], + WebhookOutbox.status == "processing", + WebhookOutbox.retry_count == claim["retry_count"], + WebhookOutbox.last_attempt == claim["last_attempt"], ) - .all() + .update(values, synchronize_session=False) ) + if completed != 1: + session.rollback() + return False + session.commit() + return True + except Exception as e: + logger.error("Error completing outbox item %s: %s", claim["id"], e) + try: + session.rollback() + except Exception: + pass + return False + finally: + session.close() - if not items: - return + async def process_outbox(self, *, force: bool = False) -> int: + """Deliver eligible items, honoring persisted due times by default. - for item in items: - item_id = item.id - payload = item.payload - - # Atomically claim this specific item - try: - db_item = session.query(WebhookOutbox).filter_by(id=item_id).first() - if ( - db_item - and db_item.status in ["pending", "failed"] - and db_item.retry_count < 5 - ): - db_item.status = "processing" - db_item.retry_count += 1 - db_item.last_attempt = datetime.now(timezone.utc) - session.commit() - else: - continue - except Exception as e: - logger.error(f"Error claiming outbox item {item_id}: {e}") - try: - session.rollback() - except Exception: - pass - continue - - # Run delivery without blocking the session or holding database locks - success = await self._send_webhook_notification(payload) - - # Update item status based on delivery outcome - try: - db_item = session.query(WebhookOutbox).filter_by(id=item_id).first() - if db_item: - if success: - db_item.status = "sent" - db_item.error_message = None - else: - db_item.status = "failed" - db_item.error_message = "Delivery failed" - session.commit() - except Exception as e: - logger.error( - f"Error updating status for outbox item {item_id}: {e}" + ``force=True`` is an explicit operational/test escape hatch that ignores + only the due timestamp; compare-and-swap claims and retry bounds remain. + """ + await self.recover_stale_deliveries() + + # Configuration absence is not a successful delivery and must not consume + # an attempt. The worker will revisit the pending row after configuration. + if not self.webhook_url: + return 0 + + session = self.Session() + try: + now = datetime.now(timezone.utc) + filters = [ + WebhookOutbox.status.in_(["pending", "failed"]), + WebhookOutbox.retry_count < self.webhook_max_attempts, + ] + if not force: + filters.append( + or_( + WebhookOutbox.next_attempt_at.is_(None), + WebhookOutbox.next_attempt_at <= now, ) - try: - session.rollback() - except Exception: - pass + ) + item_ids = [ + row[0] + for row in ( + session.query(WebhookOutbox.id) + .filter(*filters) + .order_by(WebhookOutbox.next_attempt_at, WebhookOutbox.id) + .all() + ) + ] + except Exception as e: + logger.error("Error selecting webhook outbox items: %s", e) + return 0 finally: session.close() + completed = 0 + for item_id in item_ids: + claim = self._try_claim_outbox_item( + item_id, + datetime.now(timezone.utc), + respect_schedule=not force, + ) + if claim is None: + continue + + event_id = f"api-cost:{claim['utc_date']}:{claim['alert_type']}" + token = _WEBHOOK_EVENT_ID.set(event_id) + try: + success = await self._send_webhook_notification(claim["payload"]) + except asyncio.CancelledError: + self._complete_outbox_claim( + claim, + success=False, + error_message="Delivery cancelled during worker shutdown", + ) + raise + except Exception as e: + logger.error("Webhook outbox delivery %s raised: %s", item_id, e) + success = False + finally: + _WEBHOOK_EVENT_ID.reset(token) + + if self._complete_outbox_claim(claim, success=success): + completed += 1 + + return completed + async def _send_webhook_notification(self, message: str) -> bool: """Send an async webhook notification if URL is configured. @@ -633,15 +913,26 @@ async def _send_webhook_notification(self, message: str) -> bool: True if the POST completed with a successful 2xx status; False otherwise. """ if not self.webhook_url: - return True # Behave as successful delivery if no webhook is configured + return False try: payload = {"text": message, "content": message} + event_id = _WEBHOOK_EVENT_ID.get() + headers = ( + {"Idempotency-Key": event_id, "X-Event-ID": event_id} + if event_id + else None + ) + request_kwargs: dict[str, Any] = { + "json": payload, + "timeout": aiohttp.ClientTimeout(total=5), + } + if headers is not None: + request_kwargs["headers"] = headers async with aiohttp.ClientSession() as session: async with session.post( self.webhook_url, - json=payload, - timeout=aiohttp.ClientTimeout(total=5), + **request_kwargs, ) as response: if response.status >= 200 and response.status < 300: return True From fb33cdbab48f133922e56ca61afe0e320e15b951 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 17 Jul 2026 19:32:23 -0500 Subject: [PATCH 02/31] fix: run API cost outbox with production app lifecycle --- src/youtube_extension/main.py | 25 +++++++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/src/youtube_extension/main.py b/src/youtube_extension/main.py index 9defa4778..ba53c18e4 100644 --- a/src/youtube_extension/main.py +++ b/src/youtube_extension/main.py @@ -6,6 +6,7 @@ import logging import os +from contextlib import asynccontextmanager from urllib.parse import urlparse import uvicorn @@ -17,6 +18,8 @@ from slowapi.util import get_remote_address from starlette.middleware.base import BaseHTTPMiddleware +from .backend.services.api_cost_monitor import cost_monitor + # Configure logging logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) @@ -43,6 +46,27 @@ except Exception as exc: logger.warning("Sentry init skipped: %s", exc) + +# Keep the outbox worker tied to the lifetime of the process that serves the +# production app. Failed startup still performs cleanup, and a cleanup failure +# in that path is logged without replacing the original startup exception. +@asynccontextmanager +async def _app_lifespan(_: FastAPI): + try: + await cost_monitor.start() + except BaseException: + try: + await cost_monitor.close() + except Exception: + logger.exception("API cost monitor cleanup failed after startup error") + raise + + try: + yield + finally: + await cost_monitor.close() + + # Create FastAPI application app = FastAPI( title="EventRelay API", @@ -50,6 +74,7 @@ version="1.1.0", docs_url="/docs", redoc_url="/redoc", + lifespan=_app_lifespan, ) # Configure CORS. From 61892a9e099eb5bf1aec65737ef74c5b6231e962 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 17 Jul 2026 19:32:30 -0500 Subject: [PATCH 03/31] test: preserve explicit forced retry coverage --- tests/unit/test_api_cost_monitor.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/unit/test_api_cost_monitor.py b/tests/unit/test_api_cost_monitor.py index f56055167..d9686bd35 100644 --- a/tests/unit/test_api_cost_monitor.py +++ b/tests/unit/test_api_cost_monitor.py @@ -740,7 +740,7 @@ async def fake_notification(message): # Attempt 2, 3, 4, 5 for expected_retry in [2, 3, 4, 5]: - await monitor.process_outbox() + await monitor.process_outbox(force=True) session = monitor.Session() try: item = ( @@ -754,7 +754,7 @@ async def fake_notification(message): session.close() # Attempt 6 (should not be retried because retry count reached 5) - await monitor.process_outbox() + await monitor.process_outbox(force=True) session = monitor.Session() try: item = ( From a93fe87c61939bfd9f80efcd9fe67893c83970c5 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 17 Jul 2026 19:32:32 -0500 Subject: [PATCH 04/31] test: cover durable API cost outbox worker --- tests/unit/test_api_cost_outbox_worker.py | 415 ++++++++++++++++++++++ 1 file changed, 415 insertions(+) create mode 100644 tests/unit/test_api_cost_outbox_worker.py diff --git a/tests/unit/test_api_cost_outbox_worker.py b/tests/unit/test_api_cost_outbox_worker.py new file mode 100644 index 000000000..93a858b15 --- /dev/null +++ b/tests/unit/test_api_cost_outbox_worker.py @@ -0,0 +1,415 @@ +"""Focused durability and lifecycle tests for the API-cost webhook outbox.""" + +from __future__ import annotations + +import asyncio +import sqlite3 +import sys +from datetime import datetime, timedelta, timezone +from pathlib import Path + +import pytest +from sqlalchemy import text + +sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "src")) + +from youtube_extension.backend.services import api_cost_monitor as monitor_module +from youtube_extension.backend.services.api_cost_monitor import ( + APICostMonitor, + WebhookOutbox, +) + + +def _get_item(monitor: APICostMonitor, utc_date: str) -> WebhookOutbox: + session = monitor.Session() + try: + item = ( + session.query(WebhookOutbox) + .filter_by(utc_date=utc_date, alert_type="threshold") + .one() + ) + session.expunge(item) + return item + finally: + session.close() + + +async def _wait_until(predicate, timeout: float = 1.0) -> None: + deadline = asyncio.get_running_loop().time() + timeout + while not predicate(): + if asyncio.get_running_loop().time() >= deadline: + raise AssertionError("condition was not reached before timeout") + await asyncio.sleep(0.005) + + +def test_additive_schema_upgrade_preserves_rows_and_adds_due_index(tmp_path): + db_path = tmp_path / "legacy.db" + connection = sqlite3.connect(db_path) + try: + connection.executescript(""" + CREATE TABLE webhook_outbox ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + utc_date VARCHAR NOT NULL, + alert_type VARCHAR NOT NULL, + status VARCHAR NOT NULL, + retry_count INTEGER NOT NULL, + last_attempt DATETIME, + error_message VARCHAR, + current_cost FLOAT NOT NULL, + payload VARCHAR, + CONSTRAINT uq_utc_date_alert_type UNIQUE (utc_date, alert_type) + ); + INSERT INTO webhook_outbox ( + utc_date, alert_type, status, retry_count, current_cost, payload + ) VALUES ( + '2026-07-17', 'threshold', 'pending', 0, 8.5, 'keep me' + ); + """) + connection.commit() + finally: + connection.close() + + monitor = APICostMonitor(db_path=str(db_path)) + + with monitor.engine.connect() as connection: + columns = { + row[1] + for row in connection.execute(text("PRAGMA table_info(webhook_outbox)")) + } + indexes = { + row[1] + for row in connection.execute(text("PRAGMA index_list(webhook_outbox)")) + } + due_index_columns = [ + row[2] + for row in connection.execute( + text("PRAGMA index_info(ix_webhook_outbox_due)") + ) + ] + + assert "next_attempt_at" in columns + assert "ix_webhook_outbox_due" in indexes + assert due_index_columns == ["status", "next_attempt_at", "retry_count"] + assert _get_item(monitor, "2026-07-17").payload == "keep me" + + +def test_schema_initialization_failure_is_not_suppressed(tmp_path, monkeypatch): + def fail_upgrade(self): + raise sqlite3.OperationalError("migration failed") + + monkeypatch.setattr(APICostMonitor, "_upgrade_sqlite_outbox_schema", fail_upgrade) + + with pytest.raises(sqlite3.OperationalError, match="migration failed"): + APICostMonitor(db_path=str(tmp_path / "broken.db")) + + +async def test_constructor_and_alert_do_not_spawn_background_work( + tmp_path, monkeypatch +): + recovery_started = asyncio.Event() + + async def blocking_recovery(self, stale_timeout_seconds=None): + recovery_started.set() + await asyncio.Event().wait() + + monkeypatch.setattr(APICostMonitor, "recover_stale_deliveries", blocking_recovery) + existing_tasks = asyncio.all_tasks() + monitor = APICostMonitor(db_path=str(tmp_path / "no_implicit_tasks.db")) + await asyncio.sleep(0) + spawned_tasks = asyncio.all_tasks() - existing_tasks + + try: + assert spawned_tasks == set() + assert not recovery_started.is_set() + + loop = asyncio.get_running_loop() + created = [] + original_create_task = loop.create_task + + def record_create_task(coro, *args, **kwargs): + created.append(coro) + coro.close() + return None + + with monkeypatch.context() as context: + context.setattr(loop, "create_task", record_create_task) + await monitor._send_budget_alert(8.5, "threshold") + + assert created == [] + assert loop.create_task == original_create_task + finally: + for task in spawned_tasks: + task.cancel() + if spawned_tasks: + await asyncio.gather(*spawned_tasks, return_exceptions=True) + + +async def test_start_is_idempotent_and_close_stops_the_single_worker(tmp_path): + monitor = APICostMonitor(db_path=str(tmp_path / "lifecycle.db")) + monitor.webhook_poll_interval_seconds = 60 + + first = await monitor.start() + second = await monitor.start() + + assert first is second + assert first is monitor._worker_task + assert not first.done() + + await monitor.close() + + assert first.done() + assert monitor._worker_task is None + await monitor.close() + + +async def test_missing_webhook_url_leaves_item_unattempted(tmp_path): + monitor = APICostMonitor(db_path=str(tmp_path / "missing_url.db")) + monitor.webhook_url = None + assert monitor._claim_alert("2026-07-18", "threshold", 8.5) + + await monitor.process_outbox() + + item = _get_item(monitor, "2026-07-18") + assert item.status == "pending" + assert item.retry_count == 0 + assert item.last_attempt is None + assert item.next_attempt_at is None + + +async def test_claim_is_compare_and_swap_across_monitor_instances(tmp_path): + db_path = str(tmp_path / "shared.db") + first = APICostMonitor(db_path=db_path) + second = APICostMonitor(db_path=db_path) + assert first._claim_alert("2026-07-19", "threshold", 8.5) + item_id = _get_item(first, "2026-07-19").id + claim_time = datetime.now(timezone.utc) + + claims = await asyncio.gather( + asyncio.to_thread(first._try_claim_outbox_item, item_id, claim_time, True), + asyncio.to_thread(second._try_claim_outbox_item, item_id, claim_time, True), + ) + + assert sum(claim is not None for claim in claims) == 1 + item = _get_item(first, "2026-07-19") + assert item.status == "processing" + assert item.retry_count == 1 + + +async def test_completion_is_conditional_on_the_original_claim(tmp_path): + db_path = str(tmp_path / "conditional-completion.db") + first = APICostMonitor(db_path=db_path) + second = APICostMonitor(db_path=db_path) + assert first._claim_alert("2026-07-25", "threshold", 8.5) + item_id = _get_item(first, "2026-07-25").id + + old_claim = first._try_claim_outbox_item(item_id, datetime.now(timezone.utc), False) + assert old_claim is not None + + session = second.Session() + try: + item = session.query(WebhookOutbox).filter_by(id=item_id).one() + item.status = "failed" + session.commit() + finally: + session.close() + + new_claim = second._try_claim_outbox_item( + item_id, datetime.now(timezone.utc) + timedelta(seconds=1), False + ) + assert new_claim is not None + + assert first._complete_outbox_claim(old_claim, success=True) is False + item = _get_item(first, "2026-07-25") + assert item.status == "processing" + assert item.retry_count == 2 + + assert second._complete_outbox_claim(new_claim, success=True) is True + assert _get_item(first, "2026-07-25").status == "sent" + + +async def test_failure_persists_equal_jitter_backoff_and_respects_due_time( + tmp_path, monkeypatch +): + monitor = APICostMonitor(db_path=str(tmp_path / "backoff.db")) + monitor.webhook_url = "https://example.test/hook" + monitor.webhook_retry_base_seconds = 10 + monitor.webhook_retry_max_seconds = 25 + monkeypatch.setattr(monitor_module.random, "uniform", lambda low, high: high) + + attempts = 0 + + async def fail(message): + nonlocal attempts + attempts += 1 + return False + + monkeypatch.setattr(monitor, "_send_webhook_notification", fail) + assert monitor._claim_alert("2026-07-20", "threshold", 8.5) + + expected_delays = [10, 20, 25, 25] + for expected_attempt, expected_delay in enumerate(expected_delays, start=1): + before = datetime.now(timezone.utc).replace(tzinfo=None) + await monitor.process_outbox(force=True) + item = _get_item(monitor, "2026-07-20") + assert item.retry_count == expected_attempt + assert item.status == "failed" + assert item.next_attempt_at is not None + actual_delay = (item.next_attempt_at - before).total_seconds() + assert expected_delay - 0.5 <= actual_delay <= expected_delay + 0.5 + + await monitor.process_outbox() + assert _get_item(monitor, "2026-07-20").retry_count == expected_attempt + + await monitor.process_outbox(force=True) + item = _get_item(monitor, "2026-07-20") + assert item.retry_count == 5 + assert item.next_attempt_at is None + assert item.error_message.startswith("Retry exhausted") + + await monitor.process_outbox(force=True) + assert _get_item(monitor, "2026-07-20").retry_count == 5 + assert attempts == 5 + + +async def test_worker_automatically_retries_due_delivery(tmp_path, monkeypatch): + monitor = APICostMonitor(db_path=str(tmp_path / "automatic.db")) + monitor.webhook_url = "https://example.test/hook" + monitor.webhook_retry_base_seconds = 0.01 + monitor.webhook_retry_max_seconds = 0.01 + monitor.webhook_poll_interval_seconds = 0.005 + monkeypatch.setattr(monitor_module.random, "uniform", lambda low, high: high) + attempts = 0 + + async def fail_once(message): + nonlocal attempts + attempts += 1 + return attempts > 1 + + monkeypatch.setattr(monitor, "_send_webhook_notification", fail_once) + assert monitor._claim_alert("2026-07-21", "threshold", 8.5) + + try: + await monitor.start() + await _wait_until(lambda: _get_item(monitor, "2026-07-21").status == "sent") + finally: + await monitor.close() + + assert attempts == 2 + assert _get_item(monitor, "2026-07-21").retry_count == 2 + + +@pytest.mark.parametrize("last_attempt", [None, datetime(2020, 1, 1)]) +async def test_stale_processing_recovery_handles_null_and_old_timestamps( + tmp_path, last_attempt +): + suffix = "null" if last_attempt is None else "old" + monitor = APICostMonitor(db_path=str(tmp_path / f"stale-{suffix}.db")) + assert monitor._claim_alert("2026-07-22", "threshold", 8.5) + session = monitor.Session() + try: + item = session.query(WebhookOutbox).one() + item.status = "processing" + item.retry_count = 1 + item.last_attempt = last_attempt + session.commit() + finally: + session.close() + + await monitor.recover_stale_deliveries(stale_timeout_seconds=30) + + item = _get_item(monitor, "2026-07-22") + assert item.status == "failed" + assert item.next_attempt_at is not None + assert "Recovery:" in item.error_message + + +async def test_stale_processing_at_max_attempts_is_terminal(tmp_path): + monitor = APICostMonitor(db_path=str(tmp_path / "stale-exhausted.db")) + assert monitor._claim_alert("2026-07-26", "threshold", 8.5) + session = monitor.Session() + try: + item = session.query(WebhookOutbox).one() + item.status = "processing" + item.retry_count = 5 + item.last_attempt = datetime(2020, 1, 1) + session.commit() + finally: + session.close() + + await monitor.recover_stale_deliveries(stale_timeout_seconds=30) + + item = _get_item(monitor, "2026-07-26") + assert item.status == "failed" + assert item.next_attempt_at is None + assert item.error_message.startswith("Retry exhausted") + + +async def test_cancellation_releases_claim_and_schedules_retry(tmp_path, monkeypatch): + monitor = APICostMonitor(db_path=str(tmp_path / "cancel.db")) + monitor.webhook_url = "https://example.test/hook" + monitor.webhook_retry_base_seconds = 0.01 + monitor.webhook_poll_interval_seconds = 60 + delivery_started = asyncio.Event() + + async def block(message): + delivery_started.set() + await asyncio.Event().wait() + + monkeypatch.setattr(monitor, "_send_webhook_notification", block) + assert monitor._claim_alert("2026-07-23", "threshold", 8.5) + + await monitor.start() + await asyncio.wait_for(delivery_started.wait(), timeout=1) + await monitor.close() + + item = _get_item(monitor, "2026-07-23") + assert item.status == "failed" + assert item.retry_count == 1 + assert item.next_attempt_at is not None + assert "cancel" in item.error_message.lower() + + +async def test_every_attempt_uses_stable_idempotency_headers_and_sent_is_terminal( + tmp_path, monkeypatch +): + monitor = APICostMonitor(db_path=str(tmp_path / "headers.db")) + monitor.webhook_url = "https://example.test/hook" + responses = iter([500, 204]) + captured_headers: list[dict[str, str]] = [] + + class FakeResponse: + def __init__(self, status): + self.status = status + + async def __aenter__(self): + return self + + async def __aexit__(self, *args): + return False + + class FakeSession: + async def __aenter__(self): + return self + + async def __aexit__(self, *args): + return False + + def post(self, url, json=None, timeout=None, headers=None): + captured_headers.append(headers) + return FakeResponse(next(responses)) + + monkeypatch.setattr(monitor_module.aiohttp, "ClientSession", FakeSession) + assert monitor._claim_alert("2026-07-24", "threshold", 8.5) + + await monitor.process_outbox(force=True) + await monitor.process_outbox(force=True) + await monitor.process_outbox() + + expected_event_id = "api-cost:2026-07-24:threshold" + assert captured_headers == [ + {"Idempotency-Key": expected_event_id, "X-Event-ID": expected_event_id}, + {"Idempotency-Key": expected_event_id, "X-Event-ID": expected_event_id}, + ] + item = _get_item(monitor, "2026-07-24") + assert item.status == "sent" + assert item.retry_count == 2 From 09d7891c89ea1f270b1e48090adacdf3c7ef69f5 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 17 Jul 2026 19:32:33 -0500 Subject: [PATCH 05/31] test: cover API cost monitor app lifecycle --- tests/unit/test_api_cost_monitor_lifecycle.py | 48 +++++++++++++++++++ 1 file changed, 48 insertions(+) create mode 100644 tests/unit/test_api_cost_monitor_lifecycle.py diff --git a/tests/unit/test_api_cost_monitor_lifecycle.py b/tests/unit/test_api_cost_monitor_lifecycle.py new file mode 100644 index 000000000..4ad354fb7 --- /dev/null +++ b/tests/unit/test_api_cost_monitor_lifecycle.py @@ -0,0 +1,48 @@ +"""Lifecycle coverage for the production FastAPI entrypoint.""" + +from unittest.mock import AsyncMock + +import pytest + +from youtube_extension import main as main_module + + +@pytest.fixture +def monitor(monkeypatch): + """Replace the process-global monitor without constructing another database.""" + fake = AsyncMock() + monkeypatch.setattr(main_module, "cost_monitor", fake, raising=False) + return fake + + +async def test_app_lifespan_starts_and_closes_cost_monitor(monitor) -> None: + async with main_module.app.router.lifespan_context(main_module.app): + monitor.start.assert_awaited_once_with() + monitor.close.assert_not_awaited() + + monitor.close.assert_awaited_once_with() + + +async def test_app_lifespan_closes_monitor_when_startup_fails(monitor) -> None: + startup_error = RuntimeError("monitor startup failed") + monitor.start.side_effect = startup_error + + with pytest.raises(RuntimeError, match="monitor startup failed") as raised: + async with main_module.app.router.lifespan_context(main_module.app): + pytest.fail("the application must not serve after failed startup") + + assert raised.value is startup_error + monitor.close.assert_awaited_once_with() + + +async def test_startup_error_is_not_masked_when_cleanup_also_fails(monitor) -> None: + startup_error = RuntimeError("monitor startup failed") + monitor.start.side_effect = startup_error + monitor.close.side_effect = RuntimeError("monitor cleanup failed") + + with pytest.raises(RuntimeError, match="monitor startup failed") as raised: + async with main_module.app.router.lifespan_context(main_module.app): + pytest.fail("the application must not serve after failed startup") + + assert raised.value is startup_error + monitor.close.assert_awaited_once_with() From cdea1aeca7537c660f13e1367a4e257d1972f471 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Sat, 18 Jul 2026 00:57:35 -0500 Subject: [PATCH 06/31] fix(cost): move outbox database work off event loop --- .../backend/services/api_cost_monitor.py | 63 ++++++++++++------- 1 file changed, 41 insertions(+), 22 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index e24309a40..bcee266b8 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -217,7 +217,7 @@ def __init__(self, db_path: Optional[str] = None): self.webhook_stale_timeout_seconds = max( 1, int(os.getenv("API_COST_WEBHOOK_STALE_SECONDS", "30")) ) - self._worker_task: Optional[asyncio.Task] = None + self._worker_task: Optional[asyncio.Task[None]] = None self._worker_wake_event: Optional[asyncio.Event] = None # Rate limiters for different services @@ -578,7 +578,7 @@ def _trigger_delivery(self): if self._worker_wake_event is not None: self._worker_wake_event.set() - async def start(self) -> asyncio.Task: + async def start(self) -> asyncio.Task[None]: """Start the monitor's single managed outbox worker.""" if self._worker_task is not None and not self._worker_task.done(): return self._worker_task @@ -651,7 +651,15 @@ def _retry_state( async def recover_stale_deliveries( self, stale_timeout_seconds: Optional[int] = None ) -> None: - """Recover processing claims abandoned by a crash or cancellation.""" + """Recover abandoned claims without blocking the application event loop.""" + await asyncio.to_thread( + self._recover_stale_deliveries_sync, stale_timeout_seconds + ) + + def _recover_stale_deliveries_sync( + self, stale_timeout_seconds: Optional[int] = None + ) -> None: + """Recover processing claims in a worker thread.""" if stale_timeout_seconds is None: stale_timeout_seconds = self.webhook_stale_timeout_seconds @@ -829,22 +837,10 @@ def _complete_outbox_claim( finally: session.close() - async def process_outbox(self, *, force: bool = False) -> int: - """Deliver eligible items, honoring persisted due times by default. - - ``force=True`` is an explicit operational/test escape hatch that ignores - only the due timestamp; compare-and-swap claims and retry bounds remain. - """ - await self.recover_stale_deliveries() - - # Configuration absence is not a successful delivery and must not consume - # an attempt. The worker will revisit the pending row after configuration. - if not self.webhook_url: - return 0 - + def _select_outbox_item_ids(self, *, now: datetime, force: bool) -> list[int]: + """Return due outbox IDs using a short worker-thread transaction.""" session = self.Session() try: - now = datetime.now(timezone.utc) filters = [ WebhookOutbox.status.in_(["pending", "failed"]), WebhookOutbox.retry_count < self.webhook_max_attempts, @@ -856,7 +852,7 @@ async def process_outbox(self, *, force: bool = False) -> int: WebhookOutbox.next_attempt_at <= now, ) ) - item_ids = [ + return [ row[0] for row in ( session.query(WebhookOutbox.id) @@ -867,13 +863,33 @@ async def process_outbox(self, *, force: bool = False) -> int: ] except Exception as e: logger.error("Error selecting webhook outbox items: %s", e) - return 0 + return [] finally: session.close() + async def process_outbox(self, *, force: bool = False) -> int: + """Deliver eligible items, honoring persisted due times by default. + + ``force=True`` is an explicit operational/test escape hatch that ignores + only the due timestamp; compare-and-swap claims and retry bounds remain. + """ + await self.recover_stale_deliveries() + + # Configuration absence is not a successful delivery and must not consume + # an attempt. The worker will revisit the pending row after configuration. + if not self.webhook_url: + return 0 + + item_ids = await asyncio.to_thread( + self._select_outbox_item_ids, + now=datetime.now(timezone.utc), + force=force, + ) + completed = 0 for item_id in item_ids: - claim = self._try_claim_outbox_item( + claim = await asyncio.to_thread( + self._try_claim_outbox_item, item_id, datetime.now(timezone.utc), respect_schedule=not force, @@ -886,7 +902,8 @@ async def process_outbox(self, *, force: bool = False) -> int: try: success = await self._send_webhook_notification(claim["payload"]) except asyncio.CancelledError: - self._complete_outbox_claim( + await asyncio.to_thread( + self._complete_outbox_claim, claim, success=False, error_message="Delivery cancelled during worker shutdown", @@ -898,7 +915,9 @@ async def process_outbox(self, *, force: bool = False) -> int: finally: _WEBHOOK_EVENT_ID.reset(token) - if self._complete_outbox_claim(claim, success=success): + if await asyncio.to_thread( + self._complete_outbox_claim, claim, success=success + ): completed += 1 return completed From 9810209ab37b40ec85bd57ce238fb74b77124bf1 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Sat, 18 Jul 2026 00:57:36 -0500 Subject: [PATCH 07/31] fix(api): harden cost monitor lifespan shutdown --- src/youtube_extension/main.py | 27 +++++++++++++++++++-------- 1 file changed, 19 insertions(+), 8 deletions(-) diff --git a/src/youtube_extension/main.py b/src/youtube_extension/main.py index ba53c18e4..d45e5fa54 100644 --- a/src/youtube_extension/main.py +++ b/src/youtube_extension/main.py @@ -6,6 +6,7 @@ import logging import os +from collections.abc import AsyncIterator from contextlib import asynccontextmanager from urllib.parse import urlparse @@ -33,7 +34,9 @@ sentry_sdk.init( dsn=_sentry_dsn, - environment=os.getenv("ENVIRONMENT", os.getenv("VERCEL_ENV", "development")), + environment=os.getenv( + "ENVIRONMENT", os.getenv("VERCEL_ENV", "development") + ), traces_sample_rate=float(os.getenv("SENTRY_TRACES_SAMPLE_RATE", "0.1")), integrations=[StarletteIntegration(), FastApiIntegration()], # Do not attach request headers/cookies/body/user IP to events. @@ -51,7 +54,7 @@ # production app. Failed startup still performs cleanup, and a cleanup failure # in that path is logged without replacing the original startup exception. @asynccontextmanager -async def _app_lifespan(_: FastAPI): +async def _app_lifespan(_: FastAPI) -> AsyncIterator[None]: try: await cost_monitor.start() except BaseException: @@ -64,7 +67,10 @@ async def _app_lifespan(_: FastAPI): try: yield finally: - await cost_monitor.close() + try: + await cost_monitor.close() + except Exception: + logger.exception("API cost monitor cleanup failed during shutdown") # Create FastAPI application @@ -136,14 +142,17 @@ def _is_loopback_origin(origin: str) -> bool: continue if _IS_PRODUCTION and _is_loopback_origin(_origin): logger.warning( - "Ignoring loopback origin %r from CORS_ALLOWED_ORIGINS in production", _origin + "Ignoring loopback origin %r from CORS_ALLOWED_ORIGINS in production", + _origin, ) continue _EXTRA_ORIGINS.append(_origin) -_allowed_origins = list(dict.fromkeys( - _PRODUCTION_ORIGINS + _EXTRA_ORIGINS + ([] if _IS_PRODUCTION else _DEV_ORIGINS) -)) +_allowed_origins = list( + dict.fromkeys( + _PRODUCTION_ORIGINS + _EXTRA_ORIGINS + ([] if _IS_PRODUCTION else _DEV_ORIGINS) + ) +) logger.info("CORS allow_origins configured for %s: %s", _ENVIRONMENT, _allowed_origins) app.add_middleware( @@ -177,7 +186,9 @@ async def dispatch(self, request, call_next): # API key auth middleware try: - from .backend.middleware.api_key_auth import APIKeyAuthMiddleware as APIKeyMiddleware + from .backend.middleware.api_key_auth import ( + APIKeyAuthMiddleware as APIKeyMiddleware, + ) app.add_middleware(APIKeyMiddleware) logger.info("API key auth middleware loaded") From 43ffe1090bd889d369dd2f2962d7f07080a5c054 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Sat, 18 Jul 2026 00:57:37 -0500 Subject: [PATCH 08/31] test(api): preserve application errors during cleanup --- tests/unit/test_api_cost_monitor_lifecycle.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/tests/unit/test_api_cost_monitor_lifecycle.py b/tests/unit/test_api_cost_monitor_lifecycle.py index 4ad354fb7..0ca116093 100644 --- a/tests/unit/test_api_cost_monitor_lifecycle.py +++ b/tests/unit/test_api_cost_monitor_lifecycle.py @@ -46,3 +46,17 @@ async def test_startup_error_is_not_masked_when_cleanup_also_fails(monitor) -> N assert raised.value is startup_error monitor.close.assert_awaited_once_with() + + +async def test_application_error_is_not_masked_when_shutdown_cleanup_fails( + monitor, +) -> None: + application_error = RuntimeError("application failed") + monitor.close.side_effect = RuntimeError("monitor cleanup failed") + + with pytest.raises(RuntimeError, match="application failed") as raised: + async with main_module.app.router.lifespan_context(main_module.app): + raise application_error + + assert raised.value is application_error + monitor.close.assert_awaited_once_with() From 3ae774c0aed4b9231d786ef66dc8e3bf230462b8 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Sat, 18 Jul 2026 00:57:39 -0500 Subject: [PATCH 09/31] test(cost): prove worker database work runs off loop --- tests/unit/test_api_cost_outbox_worker.py | 35 +++++++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/tests/unit/test_api_cost_outbox_worker.py b/tests/unit/test_api_cost_outbox_worker.py index 93a858b15..748ea2d94 100644 --- a/tests/unit/test_api_cost_outbox_worker.py +++ b/tests/unit/test_api_cost_outbox_worker.py @@ -5,6 +5,7 @@ import asyncio import sqlite3 import sys +import threading from datetime import datetime, timedelta, timezone from pathlib import Path @@ -413,3 +414,37 @@ def post(self, url, json=None, timeout=None, headers=None): item = _get_item(monitor, "2026-07-24") assert item.status == "sent" assert item.retry_count == 2 + + +async def test_worker_database_transactions_run_off_event_loop(tmp_path, monkeypatch): + monitor = APICostMonitor(db_path=str(tmp_path / "off-loop.db")) + monitor.webhook_url = "https://example.test/hook" + assert monitor._claim_alert("2026-07-27", "threshold", 8.5) + + event_loop_thread = threading.get_ident() + observed_threads: list[tuple[str, int]] = [] + helper_names = ( + "_recover_stale_deliveries_sync", + "_select_outbox_item_ids", + "_try_claim_outbox_item", + "_complete_outbox_claim", + ) + + for helper_name in helper_names: + original = getattr(monitor, helper_name) + + def record_thread(*args, _name=helper_name, _original=original, **kwargs): + observed_threads.append((_name, threading.get_ident())) + return _original(*args, **kwargs) + + monkeypatch.setattr(monitor, helper_name, record_thread) + + async def succeed(message): + return True + + monkeypatch.setattr(monitor, "_send_webhook_notification", succeed) + + assert await monitor.process_outbox(force=True) == 1 + assert _get_item(monitor, "2026-07-27").status == "sent" + assert {name for name, _ in observed_threads} == set(helper_names) + assert all(thread_id != event_loop_thread for _, thread_id in observed_threads) From 94c2005719d00541d6a523946e3eb62c1c50f1b1 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 20 Jul 2026 22:19:18 -0500 Subject: [PATCH 10/31] fix: count only successful outbox deliveries --- src/youtube_extension/backend/services/api_cost_monitor.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index a626c2bdd..db8d286ff 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -1782,9 +1782,10 @@ async def process_outbox( finally: _WEBHOOK_EVENT_ID.reset(token) - if await asyncio.to_thread( + claim_completed = await asyncio.to_thread( self._complete_outbox_claim, claim, success=success - ): + ) + if success and claim_completed: completed += 1 return completed From 768a392f0081b3f52a454fc1ab2109ed8bcc18cb Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 20 Jul 2026 22:19:31 -0500 Subject: [PATCH 11/31] test: cover in-memory and failed outbox delivery semantics --- tests/unit/test_api_cost_outbox_worker.py | 30 +++++++++++++++++++++++ 1 file changed, 30 insertions(+) diff --git a/tests/unit/test_api_cost_outbox_worker.py b/tests/unit/test_api_cost_outbox_worker.py index 748ea2d94..a6f45ae2e 100644 --- a/tests/unit/test_api_cost_outbox_worker.py +++ b/tests/unit/test_api_cost_outbox_worker.py @@ -177,6 +177,36 @@ async def test_missing_webhook_url_leaves_item_unattempted(tmp_path): assert item.next_attempt_at is None +async def test_in_memory_outbox_is_shared_with_worker_threads(monkeypatch): + monitor = APICostMonitor(db_path=":memory:") + monitor.webhook_url = "https://example.test/hook" + + async def succeed(message): + return True + + monkeypatch.setattr(monitor, "_send_webhook_notification", succeed) + assert monitor._claim_alert("2026-07-28", "threshold", 8.5) + + assert await monitor.process_outbox(force=True) == 1 + assert _get_item(monitor, "2026-07-28").status == "sent" + + +async def test_failed_delivery_is_not_counted_as_completed(tmp_path, monkeypatch): + monitor = APICostMonitor(db_path=str(tmp_path / "failed-count.db")) + monitor.webhook_url = "https://example.test/hook" + + async def fail(message): + return False + + monkeypatch.setattr(monitor, "_send_webhook_notification", fail) + assert monitor._claim_alert("2026-07-29", "threshold", 8.5) + + assert await monitor.process_outbox(force=True) == 0 + item = _get_item(monitor, "2026-07-29") + assert item.status == "failed" + assert item.retry_count == 1 + + async def test_claim_is_compare_and_swap_across_monitor_instances(tmp_path): db_path = str(tmp_path / "shared.db") first = APICostMonitor(db_path=db_path) From fb79912eba7819c802b5ec5472ef1b60a37399a5 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 20 Jul 2026 22:22:07 -0500 Subject: [PATCH 12/31] fix: keep API-cost delivery in the dedicated worker --- src/youtube_extension/main.py | 29 ----------------------------- 1 file changed, 29 deletions(-) diff --git a/src/youtube_extension/main.py b/src/youtube_extension/main.py index a08b7ef4a..86c12a26a 100644 --- a/src/youtube_extension/main.py +++ b/src/youtube_extension/main.py @@ -6,8 +6,6 @@ import logging import os -from collections.abc import AsyncIterator -from contextlib import asynccontextmanager from urllib.parse import urlparse import uvicorn @@ -19,8 +17,6 @@ from slowapi.util import get_remote_address from starlette.middleware.base import BaseHTTPMiddleware -from .backend.services.api_cost_monitor import cost_monitor - # Configure logging logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) @@ -49,30 +45,6 @@ except Exception as exc: logger.warning("Sentry init skipped: %s", exc) - -# Keep the outbox worker tied to the lifetime of the process that serves the -# production app. Failed startup still performs cleanup, and a cleanup failure -# in that path is logged without replacing the original startup exception. -@asynccontextmanager -async def _app_lifespan(_: FastAPI) -> AsyncIterator[None]: - try: - await cost_monitor.start() - except BaseException: - try: - await cost_monitor.close() - except Exception: - logger.exception("API cost monitor cleanup failed after startup error") - raise - - try: - yield - finally: - try: - await cost_monitor.close() - except Exception: - logger.exception("API cost monitor cleanup failed during shutdown") - - # Create FastAPI application app = FastAPI( title="EventRelay API", @@ -80,7 +52,6 @@ async def _app_lifespan(_: FastAPI) -> AsyncIterator[None]: version="1.1.0", docs_url="/docs", redoc_url="/redoc", - lifespan=_app_lifespan, ) # Configure CORS. From cd9964a05f200b5aa1d9fa71152a20ce44eaf23b Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 20 Jul 2026 22:22:16 -0500 Subject: [PATCH 13/31] test: drop obsolete FastAPI worker lifecycle coverage --- tests/unit/test_api_cost_monitor_lifecycle.py | 62 ------------------- 1 file changed, 62 deletions(-) delete mode 100644 tests/unit/test_api_cost_monitor_lifecycle.py diff --git a/tests/unit/test_api_cost_monitor_lifecycle.py b/tests/unit/test_api_cost_monitor_lifecycle.py deleted file mode 100644 index 0ca116093..000000000 --- a/tests/unit/test_api_cost_monitor_lifecycle.py +++ /dev/null @@ -1,62 +0,0 @@ -"""Lifecycle coverage for the production FastAPI entrypoint.""" - -from unittest.mock import AsyncMock - -import pytest - -from youtube_extension import main as main_module - - -@pytest.fixture -def monitor(monkeypatch): - """Replace the process-global monitor without constructing another database.""" - fake = AsyncMock() - monkeypatch.setattr(main_module, "cost_monitor", fake, raising=False) - return fake - - -async def test_app_lifespan_starts_and_closes_cost_monitor(monitor) -> None: - async with main_module.app.router.lifespan_context(main_module.app): - monitor.start.assert_awaited_once_with() - monitor.close.assert_not_awaited() - - monitor.close.assert_awaited_once_with() - - -async def test_app_lifespan_closes_monitor_when_startup_fails(monitor) -> None: - startup_error = RuntimeError("monitor startup failed") - monitor.start.side_effect = startup_error - - with pytest.raises(RuntimeError, match="monitor startup failed") as raised: - async with main_module.app.router.lifespan_context(main_module.app): - pytest.fail("the application must not serve after failed startup") - - assert raised.value is startup_error - monitor.close.assert_awaited_once_with() - - -async def test_startup_error_is_not_masked_when_cleanup_also_fails(monitor) -> None: - startup_error = RuntimeError("monitor startup failed") - monitor.start.side_effect = startup_error - monitor.close.side_effect = RuntimeError("monitor cleanup failed") - - with pytest.raises(RuntimeError, match="monitor startup failed") as raised: - async with main_module.app.router.lifespan_context(main_module.app): - pytest.fail("the application must not serve after failed startup") - - assert raised.value is startup_error - monitor.close.assert_awaited_once_with() - - -async def test_application_error_is_not_masked_when_shutdown_cleanup_fails( - monitor, -) -> None: - application_error = RuntimeError("application failed") - monitor.close.side_effect = RuntimeError("monitor cleanup failed") - - with pytest.raises(RuntimeError, match="application failed") as raised: - async with main_module.app.router.lifespan_context(main_module.app): - raise application_error - - assert raised.value is application_error - monitor.close.assert_awaited_once_with() From 11797e989837a26b1e13eccfc65353c9856faa64 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 20 Jul 2026 23:17:47 -0500 Subject: [PATCH 14/31] refactor(api-cost): manage outbox sessions through session scope --- .../backend/services/api_cost_monitor.py | 296 ++++++++---------- 1 file changed, 135 insertions(+), 161 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index db8d286ff..7b2c1eb5b 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -1514,6 +1514,7 @@ async def recover_stale_deliveries( self._recover_stale_deliveries_sync, stale_timeout_seconds ) + def _recover_stale_deliveries_sync( self, stale_timeout_seconds: Optional[int] = None ) -> None: @@ -1521,67 +1522,59 @@ def _recover_stale_deliveries_sync( if stale_timeout_seconds is None: stale_timeout_seconds = self.webhook_stale_timeout_seconds - session = self.Session() try: - now = datetime.now(timezone.utc) - cutoff = now - timedelta(seconds=stale_timeout_seconds) - stale_items = ( - session.query(WebhookOutbox) - .filter( - WebhookOutbox.status == "processing", - or_( - WebhookOutbox.last_attempt.is_(None), - WebhookOutbox.last_attempt < cutoff, - ), - ) - .all() - ) - - for item in stale_items: - next_attempt_at, recovery_error = self._retry_state( - max(1, item.retry_count), - now, - "Recovery: Stale/Crashed delivery task recovered", - ) - filters = [ - WebhookOutbox.id == item.id, - WebhookOutbox.status == "processing", - ] - if item.last_attempt is None: - filters.append(WebhookOutbox.last_attempt.is_(None)) - else: - filters.append(WebhookOutbox.last_attempt == item.last_attempt) - - recovered = ( + with self._session_scope(commit=True) as session: + now = datetime.now(timezone.utc) + cutoff = now - timedelta(seconds=stale_timeout_seconds) + stale_items = ( session.query(WebhookOutbox) - .filter(*filters) - .update( - { - WebhookOutbox.status: "failed", - WebhookOutbox.next_attempt_at: next_attempt_at, - WebhookOutbox.error_message: recovery_error, - WebhookOutbox.last_recovered_at: now, - }, - synchronize_session=False, + .filter( + WebhookOutbox.status == "processing", + or_( + WebhookOutbox.last_attempt.is_(None), + WebhookOutbox.last_attempt < cutoff, + ), ) + .all() ) - if recovered: - logger.info( - "Recovered stale webhook delivery %s for %s (%s)", - item.id, - item.utc_date, - item.alert_type, - ) - session.commit() + for item in stale_items: + next_attempt_at, recovery_error = self._retry_state( + max(1, item.retry_count), + now, + "Recovery: Stale/Crashed delivery task recovered", + ) + filters = [ + WebhookOutbox.id == item.id, + WebhookOutbox.status == "processing", + ] + if item.last_attempt is None: + filters.append(WebhookOutbox.last_attempt.is_(None)) + else: + filters.append(WebhookOutbox.last_attempt == item.last_attempt) + + recovered = ( + session.query(WebhookOutbox) + .filter(*filters) + .update( + { + WebhookOutbox.status: "failed", + WebhookOutbox.next_attempt_at: next_attempt_at, + WebhookOutbox.error_message: recovery_error, + WebhookOutbox.last_recovered_at: now, + }, + synchronize_session=False, + ) + ) + if recovered: + logger.info( + "Recovered stale webhook delivery %s for %s (%s)", + item.id, + item.utc_date, + item.alert_type, + ) except Exception as e: logger.error("Error during stale webhook delivery recovery: %s", e) - try: - session.rollback() - except Exception: - pass - finally: - session.close() def _try_claim_outbox_item( self, @@ -1590,58 +1583,50 @@ def _try_claim_outbox_item( respect_schedule: bool = True, ) -> Optional[dict[str, Any]]: """Claim one due item with a single compare-and-swap UPDATE.""" - session = self.Session() try: - filters = [ - WebhookOutbox.id == item_id, - WebhookOutbox.status.in_(["pending", "failed"]), - WebhookOutbox.retry_count < self.webhook_max_attempts, - ] - if respect_schedule: - filters.append( - or_( - WebhookOutbox.next_attempt_at.is_(None), - WebhookOutbox.next_attempt_at <= claim_time, + with self._session_scope(commit=True) as session: + filters = [ + WebhookOutbox.id == item_id, + WebhookOutbox.status.in_(["pending", "failed"]), + WebhookOutbox.retry_count < self.webhook_max_attempts, + ] + if respect_schedule: + filters.append( + or_( + WebhookOutbox.next_attempt_at.is_(None), + WebhookOutbox.next_attempt_at <= claim_time, + ) ) - ) - claimed = ( - session.query(WebhookOutbox) - .filter(*filters) - .update( - { - WebhookOutbox.status: "processing", - WebhookOutbox.retry_count: WebhookOutbox.retry_count + 1, - WebhookOutbox.last_attempt: claim_time, - WebhookOutbox.claimed_at: claim_time, - WebhookOutbox.next_attempt_at: None, - }, - synchronize_session=False, + claimed = ( + session.query(WebhookOutbox) + .filter(*filters) + .update( + { + WebhookOutbox.status: "processing", + WebhookOutbox.retry_count: WebhookOutbox.retry_count + 1, + WebhookOutbox.last_attempt: claim_time, + WebhookOutbox.claimed_at: claim_time, + WebhookOutbox.next_attempt_at: None, + }, + synchronize_session=False, + ) ) - ) - if claimed != 1: - session.rollback() - return None - - session.commit() - item = session.query(WebhookOutbox).filter_by(id=item_id).one() - return { - "id": item.id, - "payload": item.payload, - "utc_date": item.utc_date, - "alert_type": item.alert_type, - "retry_count": item.retry_count, - "last_attempt": item.last_attempt, - } + if claimed != 1: + return None + + item = session.query(WebhookOutbox).filter_by(id=item_id).one() + return { + "id": item.id, + "payload": item.payload, + "utc_date": item.utc_date, + "alert_type": item.alert_type, + "retry_count": item.retry_count, + "last_attempt": item.last_attempt, + } except Exception as e: logger.debug("Could not claim webhook outbox item %s: %s", item_id, e) - try: - session.rollback() - except Exception: - pass return None - finally: - session.close() def _complete_outbox_claim( self, @@ -1650,85 +1635,74 @@ def _complete_outbox_claim( success: bool, error_message: Optional[str] = None, ) -> bool: - """Conditionally complete exactly the attempt represented by ``claim``.""" - session = self.Session() + """Conditionally complete exactly the represented delivery attempt.""" try: - values: dict[Any, Any] - if success: - values = { - WebhookOutbox.status: "sent", - WebhookOutbox.next_attempt_at: None, - WebhookOutbox.error_message: None, - WebhookOutbox.sent_at: datetime.now(timezone.utc), - } - else: - next_attempt_at, persisted_error = self._retry_state( - claim["retry_count"], - datetime.now(timezone.utc), - error_message or "Delivery failed", - ) - values = { - WebhookOutbox.status: "failed", - WebhookOutbox.next_attempt_at: next_attempt_at, - WebhookOutbox.error_message: persisted_error, - } + with self._session_scope(commit=True) as session: + values: dict[Any, Any] + if success: + values = { + WebhookOutbox.status: "sent", + WebhookOutbox.next_attempt_at: None, + WebhookOutbox.error_message: None, + WebhookOutbox.sent_at: datetime.now(timezone.utc), + } + else: + next_attempt_at, persisted_error = self._retry_state( + claim["retry_count"], + datetime.now(timezone.utc), + error_message or "Delivery failed", + ) + values = { + WebhookOutbox.status: "failed", + WebhookOutbox.next_attempt_at: next_attempt_at, + WebhookOutbox.error_message: persisted_error, + } - completed = ( - session.query(WebhookOutbox) - .filter( - WebhookOutbox.id == claim["id"], - WebhookOutbox.status == "processing", - WebhookOutbox.retry_count == claim["retry_count"], - WebhookOutbox.last_attempt == claim["last_attempt"], + completed = ( + session.query(WebhookOutbox) + .filter( + WebhookOutbox.id == claim["id"], + WebhookOutbox.status == "processing", + WebhookOutbox.retry_count == claim["retry_count"], + WebhookOutbox.last_attempt == claim["last_attempt"], + ) + .update(values, synchronize_session=False) ) - .update(values, synchronize_session=False) - ) - if completed != 1: - session.rollback() - return False - session.commit() - return True + if completed != 1: + return False + return True except Exception as e: logger.error("Error completing outbox item %s: %s", claim["id"], e) - try: - session.rollback() - except Exception: - pass return False - finally: - session.close() def _select_outbox_item_ids( self, *, now: datetime, force: bool, max_items: Optional[int] ) -> list[int]: """Return due outbox IDs using a short worker-thread transaction.""" - session = self.Session() try: - filters = [ - WebhookOutbox.status.in_(["pending", "failed"]), - WebhookOutbox.retry_count < self.webhook_max_attempts, - ] - if not force: - filters.append( - or_( - WebhookOutbox.next_attempt_at.is_(None), - WebhookOutbox.next_attempt_at <= now, + with self._session_scope() as session: + filters = [ + WebhookOutbox.status.in_(["pending", "failed"]), + WebhookOutbox.retry_count < self.webhook_max_attempts, + ] + if not force: + filters.append( + or_( + WebhookOutbox.next_attempt_at.is_(None), + WebhookOutbox.next_attempt_at <= now, + ) ) + query = ( + session.query(WebhookOutbox.id) + .filter(*filters) + .order_by(WebhookOutbox.next_attempt_at, WebhookOutbox.id) ) - query = ( - session.query(WebhookOutbox.id) - .filter(*filters) - .order_by(WebhookOutbox.next_attempt_at, WebhookOutbox.id) - ) - if max_items is not None: - query = query.limit(max(0, max_items)) - return [row[0] for row in query.all()] + if max_items is not None: + query = query.limit(max(0, max_items)) + return [row[0] for row in query.all()] except Exception as e: logger.error("Error selecting webhook outbox items: %s", e) return [] - finally: - session.close() - async def process_outbox( self, max_items: Optional[int] = None, *, force: bool = False ) -> int: From 7eacfdbce77aae5a38298c5b722118fadd675486 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 20 Jul 2026 23:18:34 -0500 Subject: [PATCH 15/31] style(api-cost): normalize outbox method spacing --- src/youtube_extension/backend/services/api_cost_monitor.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index 7b2c1eb5b..3660a4c3f 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -1514,7 +1514,6 @@ async def recover_stale_deliveries( self._recover_stale_deliveries_sync, stale_timeout_seconds ) - def _recover_stale_deliveries_sync( self, stale_timeout_seconds: Optional[int] = None ) -> None: @@ -1703,6 +1702,7 @@ def _select_outbox_item_ids( except Exception as e: logger.error("Error selecting webhook outbox items: %s", e) return [] + async def process_outbox( self, max_items: Optional[int] = None, *, force: bool = False ) -> int: From 53eb7950f4de23a706cbb3a2e025bed1f59449b2 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Tue, 21 Jul 2026 00:31:16 -0500 Subject: [PATCH 16/31] fix(api-cost): enqueue alerts with usage transaction --- .../backend/services/api_cost_monitor.py | 98 +++++++++++++++++-- 1 file changed, 91 insertions(+), 7 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index 3660a4c3f..aa4c7a0d8 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -36,6 +36,8 @@ or_, text, ) +from sqlalchemy.dialects.postgresql import insert as postgresql_insert +from sqlalchemy.dialects.sqlite import insert as sqlite_insert from sqlalchemy.engine import URL, make_url from sqlalchemy.orm import Session, sessionmaker from sqlalchemy.pool import StaticPool @@ -1303,28 +1305,33 @@ async def record_usage( self.session_costs[service] += cost self.session_requests[service] += 1 - stored = False + claimed_alerts: list[tuple[str, float]] = [] if self.Session is None: logger.warning( "API usage was not persisted because persistence is disabled" ) else: try: - await asyncio.to_thread(self._record_usage_sync, record) - stored = True + claimed_alerts = await asyncio.to_thread( + self._record_usage_sync, record + ) except Exception as exc: # The provider operation has already completed. Telemetry is # best effort and must never make that paid result retry/fail. logger.error("Failed to record API usage: %s", exc) - if stored: - await self._check_budget_alerts() + # This only wakes the explicitly managed worker; network I/O remains + # outside the accounting path. + for alert_type, current_cost in claimed_alerts: + await self._send_budget_alert(current_cost, alert_type) logger.debug("API usage: %s - $%.4f (%s tokens)", service, cost, tokens_used) return record - def _record_usage_sync(self, record: APIUsageRecord) -> None: - """Persist one usage record on a worker thread.""" + def _record_usage_sync( + self, record: APIUsageRecord + ) -> list[tuple[str, float]]: + """Persist usage and any newly crossed alert in one transaction.""" with self._session_scope(commit=True) as session: session.add( APIUsage( @@ -1340,6 +1347,83 @@ def _record_usage_sync(self, record: APIUsageRecord) -> None: error_message=record.error_message, ) ) + session.flush() + return self._stage_budget_alerts(session, record.timestamp) + + def _stage_budget_alerts( + self, session: Session, timestamp: datetime + ) -> list[tuple[str, float]]: + """Aggregate the UTC day and enqueue crossed alerts transactionally.""" + utc_date = timestamp.astimezone(timezone.utc).date().isoformat() + start_at, end_at = self._utc_day_bounds(utc_date) + insert_factory = postgresql_insert if self._is_postgres else sqlite_insert + + session.execute( + insert_factory(DailyBudget) + .values( + date=utc_date, + total_cost=0.0, + alert_sent=False, + budget_exceeded=False, + ) + .on_conflict_do_nothing(index_elements=["date"]) + ) + budget = ( + session.query(DailyBudget) + .filter_by(date=utc_date) + .with_for_update() + .one() + ) + total = ( + session.query(func.sum(APIUsage.cost)) + .filter( + APIUsage.timestamp >= start_at, + APIUsage.timestamp < end_at, + ) + .scalar() + ) + current_cost = float(total) if total is not None else 0.0 + budget.total_cost = current_cost + + claimed: list[tuple[str, float]] = [] + alert_specs = ( + ("threshold", self.alert_threshold, "alert_sent"), + ("exceeded", self.daily_budget, "budget_exceeded"), + ) + for alert_type, limit, flag_name in alert_specs: + if current_cost < limit or getattr(budget, flag_name): + continue + + if alert_type == "threshold": + payload = ( + f"🚨 API Budget Alert: ${current_cost:.2f} " + f"(Alert threshold: ${self.alert_threshold})" + ) + else: + payload = ( + f"🚨 API Budget Alert: ${current_cost:.2f} " + f"EXCEEDED daily budget of ${self.daily_budget}" + ) + + inserted = session.execute( + insert_factory(WebhookOutbox) + .values( + utc_date=utc_date, + alert_type=alert_type, + status="pending", + retry_count=0, + current_cost=current_cost, + payload=payload, + ) + .on_conflict_do_nothing( + index_elements=["utc_date", "alert_type"] + ) + ) + setattr(budget, flag_name, True) + if inserted.rowcount == 1: + claimed.append((alert_type, current_cost)) + + return claimed async def _check_budget_alerts(self) -> None: """Check and enqueue budget alerts if thresholds are exceeded.""" From 8f944bd8b9827fbd26083ebbfd91c2cbb9c85dca Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Tue, 21 Jul 2026 00:31:42 -0500 Subject: [PATCH 17/31] test(api-cost): prove atomic alert staging --- tests/unit/test_api_cost_monitor.py | 50 +++++++++++++++++++++++++++++ 1 file changed, 50 insertions(+) diff --git a/tests/unit/test_api_cost_monitor.py b/tests/unit/test_api_cost_monitor.py index d9686bd35..97e02adf6 100644 --- a/tests/unit/test_api_cost_monitor.py +++ b/tests/unit/test_api_cost_monitor.py @@ -345,6 +345,56 @@ async def test_record_failure_usage(self, monitor): assert record.success is False assert record.error_message == "rate limited" + async def test_usage_and_crossed_alert_commit_atomically(self, monitor): + from youtube_extension.backend.models.api_cost import ( + APIUsage, + DailyBudget, + WebhookOutbox, + ) + + monitor.alert_threshold = 0.001 + monitor.daily_budget = 100.0 + record = await monitor.record_usage( + service="anthropic", + endpoint="/messages", + tokens_used=1000, + model="claude-opus-4-8", + ) + + with monitor._session_scope() as session: + assert session.query(APIUsage).count() == 1 + budget = session.query(DailyBudget).one() + alert = session.query(WebhookOutbox).one() + + assert budget.total_cost == pytest.approx(record.cost) + assert budget.alert_sent is True + assert alert.alert_type == "threshold" + assert alert.current_cost == pytest.approx(record.cost) + + async def test_alert_staging_failure_rolls_back_usage(self, monitor, monkeypatch): + from youtube_extension.backend.models.api_cost import ( + APIUsage, + DailyBudget, + WebhookOutbox, + ) + + def fail_staging(session, timestamp): + raise RuntimeError("simulated crash boundary") + + monkeypatch.setattr(monitor, "_stage_budget_alerts", fail_staging) + record = await monitor.record_usage( + service="anthropic", + endpoint="/messages", + tokens_used=1000, + model="claude-opus-4-8", + ) + + assert record is not None + with monitor._session_scope() as session: + assert session.query(APIUsage).count() == 0 + assert session.query(DailyBudget).count() == 0 + assert session.query(WebhookOutbox).count() == 0 + # =========================================================================== # APICostMonitor — get_daily_cost From 184f1acfa79420591e6e29139c35b6901999b0d6 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Tue, 21 Jul 2026 00:34:06 -0500 Subject: [PATCH 18/31] test(api-cost): enforce single atomic persistence call --- tests/unit/test_api_cost_database_substrate.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/unit/test_api_cost_database_substrate.py b/tests/unit/test_api_cost_database_substrate.py index 6baefd69c..48dd5510b 100644 --- a/tests/unit/test_api_cost_database_substrate.py +++ b/tests/unit/test_api_cost_database_substrate.py @@ -630,8 +630,9 @@ async def tracking_to_thread( await monitor.record_usage("openai", "/chat", 100, model="gpt-4o") - assert "_record_usage_sync" in calls - assert "_get_daily_cost_sync" in calls + # Usage persistence and UTC-day aggregation now share one worker-thread + # transaction; a second daily-cost query would reopen the crash boundary. + assert calls == ["_record_usage_sync"] async def test_telemetry_database_failure_does_not_fail_paid_api_result( From 3b4d66eba7d877fcce9351f2169f267853619bd9 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Tue, 21 Jul 2026 02:19:34 -0500 Subject: [PATCH 19/31] fix(api-cost): preserve Gemini usage metadata --- src/youtube_extension/services/ai/gemini_service.py | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/src/youtube_extension/services/ai/gemini_service.py b/src/youtube_extension/services/ai/gemini_service.py index 343311141..a2f0093ab 100644 --- a/src/youtube_extension/services/ai/gemini_service.py +++ b/src/youtube_extension/services/ai/gemini_service.py @@ -306,6 +306,7 @@ class GeminiResult: model_name: str backend: str # "api" or "vertex" error: Optional[str] = None + usage_metadata: Optional[Any] = None class GeminiService: @@ -589,7 +590,8 @@ async def process_image( response=response.text, latency=latency, model_name=self.config.model_name, - backend="vertex" if self._use_vertex else "api" + backend="vertex" if self._use_vertex else "api", + usage_metadata=getattr(response, "usage_metadata", None), ) except Exception as e: @@ -677,6 +679,7 @@ async def process_text( latency=time.time() - start_time, model_name=self.config.model_name, backend=self._backend_kind, + usage_metadata=getattr(response, "usage_metadata", None), ) except Exception as exc: @@ -826,7 +829,8 @@ async def process_video( response=response.text, latency=latency, model_name=self.config.model_name, - backend="vertex" if self._use_vertex else "api" + backend="vertex" if self._use_vertex else "api", + usage_metadata=getattr(response, "usage_metadata", None), ) except Exception as e: @@ -893,6 +897,7 @@ async def process_audio( latency=latency, model_name=self.config.model_name, backend="vertex" if self._use_vertex else "api", + usage_metadata=getattr(response, "usage_metadata", None), ) except Exception as e: @@ -1174,7 +1179,8 @@ async def process_youtube( response=response.text, latency=latency, model_name=self.config.model_name, - backend="api" + backend="api", + usage_metadata=getattr(response, "usage_metadata", None), ) except Exception as e: From 1c3ee57ccf36f11d6e6a71672862560ac35bd378 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Tue, 21 Jul 2026 02:20:00 -0500 Subject: [PATCH 20/31] fix(api-cost): track canonical Gemini usage --- .../services/ai/hybrid_processor_service.py | 51 +++++++++++++++++++ 1 file changed, 51 insertions(+) diff --git a/src/youtube_extension/services/ai/hybrid_processor_service.py b/src/youtube_extension/services/ai/hybrid_processor_service.py index 3d0b9b3bb..c1953e8e7 100644 --- a/src/youtube_extension/services/ai/hybrid_processor_service.py +++ b/src/youtube_extension/services/ai/hybrid_processor_service.py @@ -26,6 +26,13 @@ from .gemini_service import GeminiConfig, GeminiResult, GeminiService +async def _record_api_usage(*args: Any, **kwargs: Any) -> Any: + """Load cost tracking only when provider usage is actually available.""" + from youtube_extension.backend.services.api_cost_monitor import track_api_call + + return await track_api_call(*args, **kwargs) + + class ProcessingMode(Enum): """Processing mode roadmap retained for compatibility.""" @@ -261,6 +268,11 @@ async def process( **kwargs, ) + await self._track_gemini_usage( + cloud_result, + routing_decision.task_type, + ) + hybrid_result = HybridResult( success=cloud_result.success, response=cloud_result.response, @@ -290,6 +302,45 @@ async def process( error=str(exc), ) + async def _track_gemini_usage( + self, + result: GeminiResult, + task_type: TaskType, + ) -> None: + """Persist provider-reported usage without delaying a paid result.""" + if not result.success or result.backend not in {"api", "vertex", "gemini"}: + return + + usage = result.usage_metadata + if usage is None: + self.logger.warning( + "Gemini response omitted usage metadata; cost record skipped" + ) + return + + input_tokens = int(getattr(usage, "prompt_token_count", 0) or 0) + output_tokens = int(getattr(usage, "candidates_token_count", 0) or 0) + if input_tokens <= 0 and output_tokens <= 0: + self.logger.warning( + "Gemini usage metadata contained no billable token counts" + ) + return + + try: + await _record_api_usage( + "google", + "hybrid/process", + input_tokens, + model=result.model_name, + output_tokens=output_tokens, + request_type=task_type.value, + success=True, + ) + except Exception: + self.logger.exception( + "Gemini usage tracking failed after provider completion" + ) + async def _call_gemini( self, input_data: str | Path | Image.Image, From adbfb6b70ae86d358fe241432410030628a15c81 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Tue, 21 Jul 2026 02:20:13 -0500 Subject: [PATCH 21/31] test(api-cost): preserve Gemini usage metadata --- tests/unit/test_gemini_service.py | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/tests/unit/test_gemini_service.py b/tests/unit/test_gemini_service.py index 8e6d829b7..56fed1afe 100644 --- a/tests/unit/test_gemini_service.py +++ b/tests/unit/test_gemini_service.py @@ -53,6 +53,7 @@ def _make_service(api_key: str = "fake_key", model_name: str | None = None, **ex def _mock_response(text: str = "test response") -> MagicMock: resp = MagicMock() resp.text = text + resp.usage_metadata = None return resp @@ -691,6 +692,21 @@ async def test_process_text_success(self): assert result.success is True assert result.response == "text response here" + async def test_process_text_preserves_usage_metadata(self): + svc, mock_model, m = self._make_initialized_service() + usage = SimpleNamespace( + prompt_token_count=12, + candidates_token_count=7, + total_token_count=19, + ) + mock_response = _mock_response("tracked response") + mock_response.usage_metadata = usage + mock_model.generate_content.return_value = mock_response + + result = await svc.process_text("hello") + + assert result.usage_metadata is usage + async def test_process_text_with_input_text(self): svc, mock_model, m = self._make_initialized_service() mock_response = _mock_response("expanded response") From 45edc01037d72e7d2d9a56e18b2d5c2f6bb4ba76 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Tue, 21 Jul 2026 02:20:30 -0500 Subject: [PATCH 22/31] test(api-cost): prove canonical Gemini tracking --- tests/unit/test_hybrid_processor_service.py | 52 +++++++++++++++++++++ 1 file changed, 52 insertions(+) diff --git a/tests/unit/test_hybrid_processor_service.py b/tests/unit/test_hybrid_processor_service.py index 075d34284..2f34cbe57 100644 --- a/tests/unit/test_hybrid_processor_service.py +++ b/tests/unit/test_hybrid_processor_service.py @@ -6,6 +6,7 @@ import sys import types as _types from pathlib import Path +from types import SimpleNamespace from unittest.mock import AsyncMock, MagicMock, patch import pytest @@ -103,6 +104,7 @@ def _make_gemini_result( model_name: str = "gemini-2.0-flash", backend: str = "api", error: str | None = None, + usage_metadata: object | None = None, ) -> GeminiResult: return GeminiResult( success=success, @@ -111,6 +113,7 @@ def _make_gemini_result( model_name=model_name, backend=backend, error=error, + usage_metadata=usage_metadata, ) @@ -620,6 +623,55 @@ async def test_process_routes_youtube_url(self): await svc.process("https://www.youtube.com/watch?v=abc", "summarize") svc.gemini.process_youtube.assert_awaited_once() + async def test_process_tracks_provider_reported_usage(self): + usage = SimpleNamespace( + prompt_token_count=125, + candidates_token_count=40, + total_token_count=165, + ) + svc = self._svc( + _make_gemini_result(usage_metadata=usage) + ) + + with patch( + "youtube_extension.services.ai.hybrid_processor_service._record_api_usage", + new=AsyncMock(), + ) as track: + result = await svc.process( + "video.mp4", + "describe", + task_type=TaskType.VIDEO_UNDERSTANDING, + ) + + assert result.success is True + track.assert_awaited_once_with( + "google", + "hybrid/process", + 125, + model="gemini-2.0-flash", + output_tokens=40, + request_type="video_understanding", + success=True, + ) + + async def test_usage_tracking_failure_does_not_discard_paid_result(self): + usage = SimpleNamespace( + prompt_token_count=25, + candidates_token_count=10, + ) + svc = self._svc( + _make_gemini_result(usage_metadata=usage) + ) + + with patch( + "youtube_extension.services.ai.hybrid_processor_service._record_api_usage", + new=AsyncMock(side_effect=RuntimeError("database unavailable")), + ): + result = await svc.process("video.mp4", "describe") + + assert result.success is True + assert result.response == "ok" + async def test_process_routes_mp4_video(self): svc = self._svc() await svc.process("/data/video.mp4", "describe") From 2195c021551a2c4ded9cd6f211fca8e014d3b77a Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:26:45 -0500 Subject: [PATCH 23/31] fix(api-cost): fence claims and bill canonical Gemini usage --- .../backend/services/api_cost_monitor.py | 23 +++++++++++++++++-- 1 file changed, 21 insertions(+), 2 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index aa4c7a0d8..83cc8d627 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -16,6 +16,7 @@ import re import threading import time +import uuid from collections import defaultdict, deque from collections.abc import Iterator from contextlib import contextmanager @@ -576,6 +577,8 @@ class APICostMonitor: "claude-3-haiku-20240307": {"input": 0.00025, "output": 0.00125}, }, "google": { + # Standard paid-tier prices, normalized from per-million to per-1K USD. + "gemini-3.5-flash": {"input": 0.0015, "output": 0.009}, "gemini-3-pro": {"input": 0.000875, "output": 0.0035}, "gemini-3-flash": {"input": 0.000052, "output": 0.00021}, "gemini-1.5-pro": {"input": 0.00125, "output": 0.005}, @@ -1248,8 +1251,10 @@ def calculate_cost( service_costs = self.COST_MODELS[service] if model not in service_costs: - # Use average cost for unknown models - model = list(service_costs.keys())[0] + raise ValueError( + f"Unknown pricing model for {service}: {model!r}; " + "refusing to apply an unrelated fallback price" + ) if service == "youtube": # YouTube uses quota units, not token pricing @@ -1635,6 +1640,10 @@ def _recover_stale_deliveries_sync( filters.append(WebhookOutbox.last_attempt.is_(None)) else: filters.append(WebhookOutbox.last_attempt == item.last_attempt) + if item.claim_token is None: + filters.append(WebhookOutbox.claim_token.is_(None)) + else: + filters.append(WebhookOutbox.claim_token == item.claim_token) recovered = ( session.query(WebhookOutbox) @@ -1645,6 +1654,8 @@ def _recover_stale_deliveries_sync( WebhookOutbox.next_attempt_at: next_attempt_at, WebhookOutbox.error_message: recovery_error, WebhookOutbox.last_recovered_at: now, + WebhookOutbox.claimed_at: None, + WebhookOutbox.claim_token: None, }, synchronize_session=False, ) @@ -1681,6 +1692,7 @@ def _try_claim_outbox_item( ) ) + claim_token = uuid.uuid4().hex claimed = ( session.query(WebhookOutbox) .filter(*filters) @@ -1690,6 +1702,7 @@ def _try_claim_outbox_item( WebhookOutbox.retry_count: WebhookOutbox.retry_count + 1, WebhookOutbox.last_attempt: claim_time, WebhookOutbox.claimed_at: claim_time, + WebhookOutbox.claim_token: claim_token, WebhookOutbox.next_attempt_at: None, }, synchronize_session=False, @@ -1706,6 +1719,7 @@ def _try_claim_outbox_item( "alert_type": item.alert_type, "retry_count": item.retry_count, "last_attempt": item.last_attempt, + "claim_token": item.claim_token, } except Exception as e: logger.debug("Could not claim webhook outbox item %s: %s", item_id, e) @@ -1728,6 +1742,8 @@ def _complete_outbox_claim( WebhookOutbox.next_attempt_at: None, WebhookOutbox.error_message: None, WebhookOutbox.sent_at: datetime.now(timezone.utc), + WebhookOutbox.claimed_at: None, + WebhookOutbox.claim_token: None, } else: next_attempt_at, persisted_error = self._retry_state( @@ -1739,6 +1755,8 @@ def _complete_outbox_claim( WebhookOutbox.status: "failed", WebhookOutbox.next_attempt_at: next_attempt_at, WebhookOutbox.error_message: persisted_error, + WebhookOutbox.claimed_at: None, + WebhookOutbox.claim_token: None, } completed = ( @@ -1748,6 +1766,7 @@ def _complete_outbox_claim( WebhookOutbox.status == "processing", WebhookOutbox.retry_count == claim["retry_count"], WebhookOutbox.last_attempt == claim["last_attempt"], + WebhookOutbox.claim_token == claim["claim_token"], ) .update(values, synchronize_session=False) ) From 453e3835017a17830d368563c33149e9efb5f164 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:26:47 -0500 Subject: [PATCH 24/31] fix(api-cost): fence claims and bill canonical Gemini usage --- src/youtube_extension/services/ai/hybrid_processor_service.py | 1 + 1 file changed, 1 insertion(+) diff --git a/src/youtube_extension/services/ai/hybrid_processor_service.py b/src/youtube_extension/services/ai/hybrid_processor_service.py index c1953e8e7..3afed7606 100644 --- a/src/youtube_extension/services/ai/hybrid_processor_service.py +++ b/src/youtube_extension/services/ai/hybrid_processor_service.py @@ -320,6 +320,7 @@ async def _track_gemini_usage( input_tokens = int(getattr(usage, "prompt_token_count", 0) or 0) output_tokens = int(getattr(usage, "candidates_token_count", 0) or 0) + output_tokens += int(getattr(usage, "thoughts_token_count", 0) or 0) if input_tokens <= 0 and output_tokens <= 0: self.logger.warning( "Gemini usage metadata contained no billable token counts" From ebfdefdd309abaa75b31e62e8229b9665e96c428 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:26:49 -0500 Subject: [PATCH 25/31] fix(api-cost): fence claims and bill canonical Gemini usage --- tests/unit/test_api_cost_monitor.py | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/tests/unit/test_api_cost_monitor.py b/tests/unit/test_api_cost_monitor.py index 97e02adf6..21a3e78a2 100644 --- a/tests/unit/test_api_cost_monitor.py +++ b/tests/unit/test_api_cost_monitor.py @@ -89,11 +89,17 @@ def test_youtube_quota_cost(self, monitor): def test_unknown_service_returns_zero(self, monitor): assert monitor.calculate_cost("nonexistent", "model", 1000) == 0.0 - def test_unknown_model_falls_back_to_first_model(self, monitor): + def test_unknown_model_fails_closed(self, monitor): + with pytest.raises(ValueError, match="Unknown pricing model"): + monitor.calculate_cost( + "anthropic", "unknown-model", input_tokens=1000, output_tokens=0 + ) + + def test_google_gemini_35_flash_cost(self, monitor): cost = monitor.calculate_cost( - "anthropic", "unknown-model", input_tokens=1000, output_tokens=0 + "google", "gemini-3.5-flash", input_tokens=1000, output_tokens=1000 ) - assert cost > 0.0 + assert pytest.approx(cost, rel=1e-6) == 0.0015 + 0.009 def test_zero_tokens_returns_zero_cost(self, monitor): cost = monitor.calculate_cost( From ff806c1bfe4bcfdffbc772700d672d97ff0d9aaa Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:26:53 -0500 Subject: [PATCH 26/31] fix(api-cost): fence claims and bill canonical Gemini usage --- tests/unit/test_hybrid_processor_service.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/unit/test_hybrid_processor_service.py b/tests/unit/test_hybrid_processor_service.py index 2f34cbe57..3b3ed14e7 100644 --- a/tests/unit/test_hybrid_processor_service.py +++ b/tests/unit/test_hybrid_processor_service.py @@ -627,7 +627,8 @@ async def test_process_tracks_provider_reported_usage(self): usage = SimpleNamespace( prompt_token_count=125, candidates_token_count=40, - total_token_count=165, + thoughts_token_count=5, + total_token_count=170, ) svc = self._svc( _make_gemini_result(usage_metadata=usage) @@ -649,7 +650,7 @@ async def test_process_tracks_provider_reported_usage(self): "hybrid/process", 125, model="gemini-2.0-flash", - output_tokens=40, + output_tokens=45, request_type="video_understanding", success=True, ) From c6d123b1a81ec7abc10472e2bd17759f7d505d51 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:26:55 -0500 Subject: [PATCH 27/31] fix(api-cost): fence claims and bill canonical Gemini usage --- tests/unit/test_api_cost_outbox_worker.py | 24 +++++++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/tests/unit/test_api_cost_outbox_worker.py b/tests/unit/test_api_cost_outbox_worker.py index a6f45ae2e..8713256cc 100644 --- a/tests/unit/test_api_cost_outbox_worker.py +++ b/tests/unit/test_api_cost_outbox_worker.py @@ -478,3 +478,27 @@ async def succeed(message): assert _get_item(monitor, "2026-07-27").status == "sent" assert {name for name, _ in observed_threads} == set(helper_names) assert all(thread_id != event_loop_thread for _, thread_id in observed_threads) + + +def test_claim_token_fences_completion_and_is_cleared(tmp_path): + monitor = APICostMonitor(db_path=str(tmp_path / "claim-token.db")) + assert monitor._claim_alert("2026-07-28", "threshold", 8.5) + + item = _get_item(monitor, "2026-07-28") + claim = monitor._try_claim_outbox_item( + item.id, datetime.now(timezone.utc), respect_schedule=False + ) + + assert claim is not None + assert claim["claim_token"] + assert _get_item(monitor, "2026-07-28").claim_token == claim["claim_token"] + + stale_claim = dict(claim, claim_token="not-the-owner") + assert monitor._complete_outbox_claim(stale_claim, success=True) is False + assert _get_item(monitor, "2026-07-28").status == "processing" + + assert monitor._complete_outbox_claim(claim, success=True) is True + completed = _get_item(monitor, "2026-07-28") + assert completed.status == "sent" + assert completed.claim_token is None + assert completed.claimed_at is None From 0b3bbf2ce9bbe77a30f5c195224534210142618a Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:34:04 -0500 Subject: [PATCH 28/31] fix(cost-monitor): preserve legacy unspecified-model costing --- .../backend/services/api_cost_monitor.py | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index 83cc8d627..fd539b75f 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -1251,10 +1251,13 @@ def calculate_cost( service_costs = self.COST_MODELS[service] if model not in service_costs: - raise ValueError( - f"Unknown pricing model for {service}: {model!r}; " - "refusing to apply an unrelated fallback price" - ) + if model == "default": + model = next(iter(service_costs)) + else: + raise ValueError( + f"Unknown pricing model for {service}: {model!r}; " + "refusing to apply an unrelated fallback price" + ) if service == "youtube": # YouTube uses quota units, not token pricing From 9276cf11e129cacc30d879dd9c47ba90a2971f12 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Fri, 24 Jul 2026 19:34:23 -0500 Subject: [PATCH 29/31] test(cost-monitor): cover unspecified-model compatibility --- tests/unit/test_api_cost_monitor.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/tests/unit/test_api_cost_monitor.py b/tests/unit/test_api_cost_monitor.py index 21a3e78a2..adedcd82b 100644 --- a/tests/unit/test_api_cost_monitor.py +++ b/tests/unit/test_api_cost_monitor.py @@ -95,6 +95,12 @@ def test_unknown_model_fails_closed(self, monitor): "anthropic", "unknown-model", input_tokens=1000, output_tokens=0 ) + def test_default_model_preserves_legacy_service_costing(self, monitor): + cost = monitor.calculate_cost( + "openai", "default", input_tokens=1000, output_tokens=0 + ) + assert cost > 0.0 + def test_google_gemini_35_flash_cost(self, monitor): cost = monitor.calculate_cost( "google", "gemini-3.5-flash", input_tokens=1000, output_tokens=1000 From 3093abf885cbcb05d8d6703c33a856bca8f61119 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 27 Jul 2026 16:50:32 -0500 Subject: [PATCH 30/31] fix(outbox): align recovery queries with worker indexes --- .../backend/services/api_cost_monitor.py | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/src/youtube_extension/backend/services/api_cost_monitor.py b/src/youtube_extension/backend/services/api_cost_monitor.py index fd539b75f..7e69881d0 100644 --- a/src/youtube_extension/backend/services/api_cost_monitor.py +++ b/src/youtube_extension/backend/services/api_cost_monitor.py @@ -854,6 +854,7 @@ def _upgrade_sqlite_outbox_schema(self) -> None: "status", "next_attempt_at", "retry_count", + "id", ] if index_columns and index_columns != expected_due_index_columns: connection.exec_driver_sql( @@ -861,7 +862,7 @@ def _upgrade_sqlite_outbox_schema(self) -> None: ) connection.exec_driver_sql( "CREATE INDEX IF NOT EXISTS ix_webhook_outbox_due " - "ON webhook_outbox (status, next_attempt_at, retry_count)" + "ON webhook_outbox (status, next_attempt_at, retry_count, id)" ) connection.exec_driver_sql( "CREATE INDEX IF NOT EXISTS ix_webhook_outbox_stale_claims " @@ -1622,8 +1623,8 @@ def _recover_stale_deliveries_sync( .filter( WebhookOutbox.status == "processing", or_( - WebhookOutbox.last_attempt.is_(None), - WebhookOutbox.last_attempt < cutoff, + WebhookOutbox.claimed_at.is_(None), + WebhookOutbox.claimed_at < cutoff, ), ) .all() @@ -1639,10 +1640,10 @@ def _recover_stale_deliveries_sync( WebhookOutbox.id == item.id, WebhookOutbox.status == "processing", ] - if item.last_attempt is None: - filters.append(WebhookOutbox.last_attempt.is_(None)) + if item.claimed_at is None: + filters.append(WebhookOutbox.claimed_at.is_(None)) else: - filters.append(WebhookOutbox.last_attempt == item.last_attempt) + filters.append(WebhookOutbox.claimed_at == item.claimed_at) if item.claim_token is None: filters.append(WebhookOutbox.claim_token.is_(None)) else: From cc6fe44ab2d6092b2f550d8c1bec74e1a90e6db1 Mon Sep 17 00:00:00 2001 From: Hayden <154503486+groupthinking@users.noreply.github.com> Date: Mon, 27 Jul 2026 16:50:46 -0500 Subject: [PATCH 31/31] test(outbox): cover canonical index and claim timestamp --- tests/unit/test_api_cost_outbox_worker.py | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/tests/unit/test_api_cost_outbox_worker.py b/tests/unit/test_api_cost_outbox_worker.py index 8713256cc..f4603e8d6 100644 --- a/tests/unit/test_api_cost_outbox_worker.py +++ b/tests/unit/test_api_cost_outbox_worker.py @@ -90,7 +90,7 @@ def test_additive_schema_upgrade_preserves_rows_and_adds_due_index(tmp_path): assert "next_attempt_at" in columns assert "ix_webhook_outbox_due" in indexes - assert due_index_columns == ["status", "next_attempt_at", "retry_count"] + assert due_index_columns == ["status", "next_attempt_at", "retry_count", "id"] assert _get_item(monitor, "2026-07-17").payload == "keep me" @@ -329,11 +329,11 @@ async def fail_once(message): assert _get_item(monitor, "2026-07-21").retry_count == 2 -@pytest.mark.parametrize("last_attempt", [None, datetime(2020, 1, 1)]) +@pytest.mark.parametrize("claimed_at", [None, datetime(2020, 1, 1)]) async def test_stale_processing_recovery_handles_null_and_old_timestamps( - tmp_path, last_attempt + tmp_path, claimed_at ): - suffix = "null" if last_attempt is None else "old" + suffix = "null" if claimed_at is None else "old" monitor = APICostMonitor(db_path=str(tmp_path / f"stale-{suffix}.db")) assert monitor._claim_alert("2026-07-22", "threshold", 8.5) session = monitor.Session() @@ -341,7 +341,7 @@ async def test_stale_processing_recovery_handles_null_and_old_timestamps( item = session.query(WebhookOutbox).one() item.status = "processing" item.retry_count = 1 - item.last_attempt = last_attempt + item.claimed_at = claimed_at session.commit() finally: session.close() @@ -362,7 +362,7 @@ async def test_stale_processing_at_max_attempts_is_terminal(tmp_path): item = session.query(WebhookOutbox).one() item.status = "processing" item.retry_count = 5 - item.last_attempt = datetime(2020, 1, 1) + item.claimed_at = datetime(2020, 1, 1) session.commit() finally: session.close()