73 lines
2.4 KiB
Python
73 lines
2.4 KiB
Python
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 ImportAlreadyRunning, 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
|
|
try:
|
|
run = import_records(
|
|
session,
|
|
url=settings.official_records_url,
|
|
region=settings.official_records_region,
|
|
category=settings.official_records_category,
|
|
)
|
|
except ImportAlreadyRunning:
|
|
logger.info("official import skipped because it is already running", extra={"event": "official_import_locked"})
|
|
return False
|
|
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()
|