248 lines
11 KiB
Python
248 lines
11 KiB
Python
"""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 _tz_from_offset(offset: str | None):
|
|
"""'+08:00' -> a datetime.timezone. None if missing/unparseable."""
|
|
if not offset:
|
|
return None
|
|
try:
|
|
return datetime.strptime(
|
|
f"2000-01-01T00:00:00{offset}", "%Y-%m-%dT%H:%M:%S%z"
|
|
).tzinfo
|
|
except ValueError:
|
|
return None
|
|
|
|
|
|
def _iso_local_from_unix(ts: int, offset: str | None) -> str | None:
|
|
"""A FlightAware track timestamp (UTC epoch seconds) as an ISO string in
|
|
the *airport's* local wall clock.
|
|
|
|
This conversion is not cosmetic. AirTrail reads only the literal Y-M-D
|
|
out of a datetime field and the HH:MM out of its companion *Time field,
|
|
then interprets that pair in the airport's own timezone
|
|
(mergeTimeWithDate - confirmed in AirTrail's source). A UTC timestamp
|
|
sent as-is is therefore silently stored wrong by the whole offset, and
|
|
when takeoff comes from AeroDataBox (already airport-local) while
|
|
landing falls back to this track (UTC), the mismatch can invert their
|
|
order and get the save rejected outright with "Actual landing must be
|
|
after actual takeoff" - which is exactly what happened on VA556.
|
|
|
|
With no offset known (AeroDataBox unavailable, and airports.csv carries
|
|
no timezone data) this returns None so the caller omits the field
|
|
entirely - a missing takeoff time is recoverable, a wrong one is not.
|
|
"""
|
|
tzinfo = _tz_from_offset(offset)
|
|
if tzinfo is None:
|
|
return None
|
|
return datetime.fromtimestamp(ts, tz=tzinfo).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")
|
|
# Each leg's UTC offset on the day - what lets the FlightAware track's
|
|
# UTC timestamps be expressed in the same airport-local terms as
|
|
# everything else (see _iso_local_from_unix).
|
|
updates["departure_utc_offset"] = gate_times.get("departureUtcOffset")
|
|
updates["arrival_utc_offset"] = gate_times.get("arrivalUtcOffset")
|
|
# AeroDataBox's runwayTime is an authoritative actual takeoff/landing
|
|
# time - prefer it. Assigned unconditionally (None included) so a
|
|
# re-query always rewrites these rather than leaving a stale value
|
|
# behind; the FlightAware-track fallback below fills in only where
|
|
# AeroDataBox came up empty.
|
|
updates["takeoff_actual"] = gate_times.get("takeoffActual")
|
|
updates["landing_actual"] = gate_times.get("landingActual")
|
|
updates["departure_terminal"] = gate_times.get("departureTerminal")
|
|
updates["departure_gate"] = gate_times.get("departureGate")
|
|
updates["arrival_terminal"] = gate_times.get("arrivalTerminal")
|
|
updates["arrival_gate"] = gate_times.get("arrivalGate")
|
|
|
|
registration = gate_times.get("aircraftReg")
|
|
updates["aircraft_reg"] = registration
|
|
updates["aircraft_icao"] = None
|
|
if registration:
|
|
# One extra API unit, and only worth spending once we actually have a
|
|
# registration to look up. A failure here is not worth losing the rest
|
|
# of the lookup over - the registration alone still reaches AirTrail.
|
|
try:
|
|
updates["aircraft_icao"] = aerodatabox.lookup_aircraft_icao(registration)
|
|
except aerodatabox.AeroDataBoxError as e:
|
|
log.warning("aircraft type lookup failed for %s: %s", registration, e)
|
|
|
|
|
|
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 - and convert into the
|
|
# departure/arrival airport's local wall clock, since that's the
|
|
# only form AirTrail interprets correctly.
|
|
if times:
|
|
dep_offset = updates.get("departure_utc_offset") or row.get("departure_utc_offset")
|
|
arr_offset = updates.get("arrival_utc_offset") or row.get("arrival_utc_offset")
|
|
if not updates.get("takeoff_actual"):
|
|
updates["takeoff_actual"] = _iso_local_from_unix(times[0], dep_offset)
|
|
if not updates.get("landing_actual"):
|
|
updates["landing_actual"] = _iso_local_from_unix(times[-1], arr_offset)
|
|
|
|
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
|