from __future__ import annotations import logging import time from datetime import datetime, timedelta, timezone from sqlalchemy import select from sqlalchemy.orm import Session from .config import settings from .database import SessionLocal from .importer import import_records from .logging_config import configure_logging from .models import OfficialRecordImport logger = logging.getLogger("rf4.import_scheduler") def import_is_due(session: Session, *, now: datetime | None = None) -> bool: current = now or datetime.now(timezone.utc) latest = session.scalar( select(OfficialRecordImport.started_at) .where(OfficialRecordImport.source_url == settings.official_records_url) .order_by(OfficialRecordImport.started_at.desc()) .limit(1) ) if latest is None: return True if latest.tzinfo is None: latest = latest.replace(tzinfo=timezone.utc) return latest <= current - timedelta(seconds=settings.import_interval_seconds) def run_due_import() -> bool: with SessionLocal() as session: if not import_is_due(session): return False run = import_records( session, url=settings.official_records_url, region=settings.official_records_region, category=settings.official_records_category, ) logger.info( "official import completed", extra={ "event": "official_import_completed", "result": run.status.value, "rows_seen": run.rows_seen, "rows_created": run.rows_created, "rows_updated": run.rows_updated, "not_modified": run.not_modified, }, ) return True def main() -> None: configure_logging(settings.log_level) logger.info("scheduler started", extra={"event": "scheduler_started", "interval_seconds": settings.import_interval_seconds}) while True: try: run_due_import() except Exception: logger.exception("scheduled official import failed", extra={"event": "official_import_failed"}) time.sleep(settings.import_interval_seconds) if __name__ == "__main__": main()