Spaces:
Sleeping
Sleeping
| """Standalone scheduler process — runs the cron jobs in its own OS process | |
| so the FastAPI server can run with multiple uvicorn workers. | |
| Responsibilities (taken over from the FastAPI lifespan): | |
| - Resolve the monitored mail folder ID via Microsoft Graph. | |
| - Bootstrap the Graph subscription: reuse an active one if it has more | |
| than 60 minutes of life left; otherwise create a new one. | |
| - Start the APScheduler-driven RenewalScheduler with three jobs: | |
| * subscription_renewal — keeps the Graph subscription alive | |
| * daily_brief — 18:00 America/Chicago email | |
| * retention_cleanup — 02:00 UTC daily DB + log file purge. | |
| - On signal, shut everything down cleanly. | |
| Why a separate process? | |
| The previous design ran APScheduler inside the FastAPI lifespan. That | |
| only works with `--workers 1`; with N workers every cron job would fire | |
| N times. Moving the scheduler out of the API process lets us scale the | |
| webhook receiver horizontally without duplicating cron work. | |
| Startup timing note: | |
| Graph's subscription-creation handshake requires our /webhook/notify | |
| endpoint to be online (Graph POSTs a validation token to it). So this | |
| process retries subscription creation with exponential backoff for up | |
| to ~60 seconds, giving uvicorn time to come up. | |
| Run with: python -m app.scheduler_main | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import logging | |
| import logging.handlers | |
| import signal | |
| import traceback | |
| from datetime import datetime, timedelta, timezone | |
| from app.config import settings | |
| from app.database import get_active_subscription, get_db, init_db | |
| from app.lib.graph.auth import GraphAuthProvider | |
| from app.lib.graph.client import GraphClient | |
| from app.lib.graph.folder_resolver import resolve_folder_id | |
| from app.lib.graph.subscription import SubscriptionManager | |
| from app.lib.utils.notifier import send_developer_alert | |
| from app.scheduler import RenewalScheduler | |
| logging.basicConfig( | |
| level=logging.INFO, | |
| format="%(asctime)s [scheduler] %(levelname)s %(message)s", | |
| ) | |
| logger = logging.getLogger(__name__) | |
| _shutdown = asyncio.Event() | |
| def _setup_log_file() -> None: | |
| settings.log_path.mkdir(parents=True, exist_ok=True) | |
| handler = logging.handlers.TimedRotatingFileHandler( | |
| filename=settings.log_path / "scheduler.log", | |
| when="midnight", | |
| interval=1, | |
| backupCount=0, | |
| encoding="utf-8", | |
| ) | |
| handler.setFormatter( | |
| logging.Formatter("%(asctime)s [scheduler] %(levelname)s %(message)s") | |
| ) | |
| logging.getLogger().addHandler(handler) | |
| def _handle_signal(*_: object) -> None: | |
| logger.info("Shutdown signal received") | |
| _shutdown.set() | |
| async def _bootstrap_subscription( | |
| sub_manager: SubscriptionManager, | |
| folder_id: str, | |
| ) -> str: | |
| """Reuse the active subscription if it has >60 minutes of life left; | |
| otherwise create a new one. Retries creation with exponential backoff | |
| so we don't crash if uvicorn isn't fully up yet to receive Graph's | |
| validation POST.""" | |
| # Try to reuse first — cheap, no Graph round trip. | |
| async with get_db(settings.database_path) as conn: | |
| existing = await get_active_subscription(conn) | |
| if existing: | |
| expiry = datetime.fromisoformat( | |
| existing["expiry_datetime"].replace("Z", "+00:00") | |
| ) | |
| cutoff = datetime.now(timezone.utc) + timedelta(minutes=60) | |
| if expiry > cutoff: | |
| logger.info("Reusing existing subscription %s", existing["id"]) | |
| return existing["id"] | |
| # Need to create a new one. Retry with backoff to give uvicorn time | |
| # to come up so Graph's validation POST will succeed. | |
| delays = [2, 4, 8, 16, 30] | |
| last_exc: Exception | None = None | |
| for i, delay in enumerate(delays): | |
| try: | |
| sub_id = await sub_manager.ensure_subscription(folder_id) | |
| logger.info("Registered new subscription %s (attempt %d)", sub_id, i + 1) | |
| return sub_id | |
| except Exception as e: | |
| last_exc = e | |
| logger.warning( | |
| "Subscription creation attempt %d failed (%s); retrying in %ds", | |
| i + 1, type(e).__name__, delay, | |
| ) | |
| try: | |
| await asyncio.wait_for(_shutdown.wait(), timeout=delay) | |
| # If we get here, shutdown was requested — abort. | |
| raise asyncio.CancelledError("shutdown during subscription bootstrap") | |
| except asyncio.TimeoutError: | |
| pass | |
| # All attempts exhausted. | |
| tb = "".join(traceback.format_exception(last_exc)) if last_exc else "(no traceback)" | |
| await send_developer_alert( | |
| settings, | |
| subject="[RCM] Scheduler: subscription bootstrap failed after retries", | |
| body=( | |
| "The scheduler process could not create a Graph subscription after " | |
| f"{len(delays)} retries. Email notifications will NOT be received " | |
| f"until this is resolved (container restart or manual intervention).\n\n" | |
| f"{tb}" | |
| ), | |
| ) | |
| raise RuntimeError("Subscription bootstrap failed") from last_exc | |
| async def run() -> None: | |
| # DB schema is normally applied by `alembic upgrade head` in entrypoint.sh | |
| # before this process starts. init_db is idempotent (CREATE TABLE IF NOT | |
| # EXISTS), so we run it again here for local-dev startups that skip | |
| # alembic. | |
| await init_db(settings.database_path) | |
| auth = GraphAuthProvider(settings) | |
| client = GraphClient(auth) | |
| folder_id = await resolve_folder_id( | |
| client, settings.mail_folder_name, settings.mailbox_user | |
| ) | |
| logger.info( | |
| "Monitoring folder '%s' (id=%s)", settings.mail_folder_name, folder_id | |
| ) | |
| sub_manager = SubscriptionManager(client, settings, settings.database_path) | |
| scheduler = RenewalScheduler(sub_manager, settings) | |
| try: | |
| sub_id = await _bootstrap_subscription(sub_manager, folder_id) | |
| except Exception: | |
| # _bootstrap_subscription already sent the developer alert. | |
| logger.exception("Scheduler aborting due to subscription bootstrap failure") | |
| await client.aclose() | |
| return | |
| scheduler.start(sub_id) | |
| logger.info("Scheduler running. Waiting for shutdown signal...") | |
| try: | |
| await _shutdown.wait() | |
| finally: | |
| logger.info("Shutting down scheduler") | |
| scheduler.shutdown() | |
| await client.aclose() | |
| logger.info("Scheduler stopped") | |
| def main() -> None: | |
| _setup_log_file() | |
| loop = asyncio.new_event_loop() | |
| asyncio.set_event_loop(loop) | |
| for sig in (signal.SIGTERM, signal.SIGINT): | |
| loop.add_signal_handler(sig, _handle_signal) | |
| try: | |
| loop.run_until_complete(run()) | |
| finally: | |
| loop.close() | |
| if __name__ == "__main__": | |
| main() | |