RcmEmailAutomation / app /scheduler_main.py
Avinash Nalla
Phase 5: split scheduler into standalone process; enable uvicorn workers
b96346f
Raw
History Blame Contribute Delete
6.8 kB
"""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()