Initial commit: boarding pass to AirTrail pipeline
Photograph/screenshot a boarding pass, decode its BCBP barcode (PDF417/ Aztec/QR/DataMatrix), fetch the historical flight track from FlightAware the day after the flight, stage it for review, and write it into AirTrail via its REST API on approval. ntfy notifications carry signed Approve/Deny/Investigate actions.
This commit is contained in:
+132
@@ -0,0 +1,132 @@
|
||||
"""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
|
||||
Reference in New Issue
Block a user