"""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, timezone import aerodatabox 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 _iso_from_unix(ts: int) -> str: return datetime.fromtimestamp(ts, tz=timezone.utc).isoformat(timespec="seconds") def _fetch_gate_times(flight_number: str, flight_date: date, updates: dict) -> None: """Best-effort AeroDataBox lookup for scheduled/actual GATE times (departure/arrival) - the data BCBP and the FlightAware track can't provide. Failures here are recorded (schedule_error) but never block the FlightAware track or push status to track_unavailable: worst case the gate fields stay blank exactly like they do today, and the existing fallback (AirTrail's own manual Search) still works.""" try: gate_times = aerodatabox.lookup_gate_times(flight_number, flight_date) except aerodatabox.AeroDataBoxError as e: updates["schedule_error"] = str(e) return updates["schedule_error"] = None if not gate_times: return updates["gate_departure_scheduled"] = gate_times.get("departureScheduled") updates["gate_departure_actual"] = gate_times.get("departure") updates["gate_arrival_scheduled"] = gate_times.get("arrivalScheduled") updates["gate_arrival_actual"] = gate_times.get("arrival") # AeroDataBox's runwayTime is an authoritative actual takeoff/landing # time - prefer it. Only set here if present; the FlightAware-track # fallback below fills in behind it (checks "not in updates") rather # than overwriting it unconditionally. if gate_times.get("takeoffActual"): updates["takeoff_actual"] = gate_times["takeoffActual"] if gate_times.get("landingActual"): updates["landing_actual"] = gate_times["landingActual"] 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} try: flight_date = date.fromisoformat(row["flight_date"]) except (TypeError, ValueError): flight_date = None # Gate times only need a flight number + date (no ICAO resolution), so # attempt this regardless of whether the airport/carrier lookups below # succeed - an unresolved airport shouldn't cost us data we could # otherwise have gotten. iata_flight_number = f"{row.get('operating_carrier_iata') or ''}{row.get('flight_number') or ''}" if iata_flight_number.strip() and flight_date: _fetch_gate_times(iata_flight_number, flight_date, updates) 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 if flight_date is None: 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 # The ADS-B track only exists while the aircraft is squawking - # roughly wheels-up to touchdown - so its first/last timestamps are # a decent proxy for actual takeoff/landing when AeroDataBox didn't # already give us its authoritative runwayTime above (no key # configured, lookup failed, or that particular flight had no # runwayTime in its response). Never overwrite a value AeroDataBox # already set - that one's more precise. if times: updates.setdefault("takeoff_actual", _iso_from_unix(times[0])) updates.setdefault("landing_actual", _iso_from_unix(times[-1])) 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