"""Daily queue-processing job (stages 3-6 of the pipeline) plus the daemon thread that runs it once a day. No cron/APScheduler dependency - a single background thread sleeping until the next configured run time is enough for a one-flight-a-day, single-container workload. process_one_flight() is also called synchronously, for a single row, by the review page's "re-query" button - same function, so a corrected field (wrong year, fixed airport code, etc.) re-runs exactly the same lookup. """ import logging import threading import time from datetime import date, datetime, timedelta import config import flightaware import models import notifier import reference_data log = logging.getLogger(__name__) def next_run_datetime(hhmm: str, now: datetime | None = None) -> datetime: now = now or datetime.now() try: hour, minute = (int(p) for p in hhmm.split(":")) except ValueError: hour, minute = 3, 0 candidate = now.replace(hour=hour, minute=minute, second=0, microsecond=0) if candidate <= now: candidate += timedelta(days=1) return candidate def process_one_flight(row: dict) -> None: flight_id = row["id"] from_icao = row.get("from_icao") or reference_data.iata_to_icao_airport(row.get("from_iata")) to_icao = row.get("to_icao") or reference_data.iata_to_icao_airport(row.get("to_iata")) carrier_icao = row.get("operating_carrier_icao") or reference_data.iata_to_icao_airline( row.get("operating_carrier_iata") ) updates = {"from_icao": from_icao, "to_icao": to_icao, "operating_carrier_icao": carrier_icao} if not from_icao or not to_icao or not carrier_icao: missing = [ name for name, val in (("from airport", from_icao), ("to airport", to_icao), ("carrier", carrier_icao)) if not val ] updates["status"] = "track_unavailable" updates["error_message"] = f"could not resolve ICAO code(s) for: {', '.join(missing)}" models.update_flight(flight_id, **updates) notifier.notify_flight_staged(models.get_flight(flight_id)) return try: flight_date = date.fromisoformat(row["flight_date"]) except (TypeError, ValueError): updates["status"] = "track_unavailable" updates["error_message"] = f"invalid flight_date: {row.get('flight_date')!r}" models.update_flight(flight_id, **updates) notifier.notify_flight_staged(models.get_flight(flight_id)) return icao_callsign = f"{carrier_icao}{row.get('flight_number') or ''}" try: history_url = flightaware.find_history_url(icao_callsign, flight_date, from_icao, to_icao) kml_bytes = flightaware.fetch_kml(history_url) if history_url else None track = flightaware.parse_kml(kml_bytes) if kml_bytes else None except flightaware.FlightAwareTransientError as e: attempts = row.get("fetch_attempts", 0) + 1 updates["fetch_attempts"] = attempts if attempts >= config.MAX_FETCH_ATTEMPTS: updates["status"] = "track_unavailable" updates["error_message"] = f"gave up after {attempts} attempts: {e}" models.update_flight(flight_id, **updates) notifier.notify_flight_staged(models.get_flight(flight_id)) else: updates["error_message"] = str(e) models.update_flight(flight_id, **updates) return if history_url: updates["flightaware_url"] = history_url if track is None: updates["status"] = "track_unavailable" updates["error_message"] = updates.get("error_message") or "no track found on FlightAware" else: coordinates, times = track updates["status"] = "pending_review" updates["track_coordinates"] = coordinates updates["track_times"] = times updates["track_point_count"] = len(coordinates) updates["error_message"] = None models.update_flight(flight_id, **updates) notifier.notify_flight_staged(models.get_flight(flight_id)) def run_daily_job() -> None: now_iso = datetime.now().isoformat() due = models.list_due_flights(now_iso) log.info("daily job: %d flight(s) due", len(due)) for row in due: try: process_one_flight(row) except Exception: log.exception("failed to process flight %s", row.get("id")) def scheduler_loop() -> None: while True: run_time = models.get_setting("daily_run_time", "03:00") run_at = next_run_datetime(run_time) sleep_seconds = max((run_at - datetime.now()).total_seconds(), 1) time.sleep(sleep_seconds) try: run_daily_job() except Exception: log.exception("daily job crashed") def start_background_thread() -> threading.Thread: thread = threading.Thread(target=scheduler_loop, name="scheduler", daemon=True) thread.start() return thread