diff --git a/.env.example b/.env.example index 7791da1..439a8da 100644 --- a/.env.example +++ b/.env.example @@ -5,3 +5,8 @@ DEBUG=false # Database # Format: postgresql+asyncpg://:@:/ DATABASE_URL=postgresql+asyncpg://user:password@localhost:5432/dbname + +# External APIs +XRAPIDAPI_KEY= +OPENSKY_CLIENT_ID= +OPENSKY_CLIENT_SECRET= diff --git a/backend/foresight/config.py b/backend/foresight/config.py index 24eb3da..d073b79 100644 --- a/backend/foresight/config.py +++ b/backend/foresight/config.py @@ -3,7 +3,7 @@ class Settings(BaseSettings): model_config = SettingsConfigDict( - env_file=".env", + env_file=(".env", "../.env"), # repo-root .env when run from backend/ env_file_encoding="utf-8", extra="ignore", ) @@ -17,5 +17,10 @@ class Settings(BaseSettings): "postgresql+asyncpg://foresight:foresight@localhost:5432/foresight_db" ) + # External APIs + xrapidapi_key: str = "" # RapidAPI key for AeroDataBox (flight arrivals) + opensky_client_id: str = "" # OpenSky API client (historical flights) + opensky_client_secret: str = "" + settings = Settings() diff --git a/backend/foresight/services/flight_data_service.py b/backend/foresight/services/flight_data_service.py new file mode 100644 index 0000000..1c2379a --- /dev/null +++ b/backend/foresight/services/flight_data_service.py @@ -0,0 +1,249 @@ +r"""Fetch flight data from the CSO, OpenSky and AeroDataBox and save it as CSV. + +Dates and times are in Dublin local time (Europe/Dublin). + +Example: + Run from ``backend/``:: + + uv run python foresight/services/flight_data_service.py \ + --start 2026-06-01 --end 2026-06-07 + uv run python foresight/services/flight_data_service.py \ + --api cso --start 2026-01-01 --end 2026-06-30 +""" + +import argparse +import io +from datetime import date, datetime, time, timedelta +from pathlib import Path +from zoneinfo import ZoneInfo + +import httpx +import pandas as pd + +from foresight.config import settings + +OUTPUT_DIR = Path(__file__).resolve().parents[2] / "data" / "flights" + +CSO_URL = "https://ws.cso.ie/public/api.restful/PxStat.Data.Cube_API.ReadDataset/TAM07/CSV/1.0/en" +OPENSKY_URL = "https://opensky-network.org/api/flights/arrival" +OPENSKY_TOKEN_URL = "https://auth.opensky-network.org/auth/realms/opensky-network/protocol/openid-connect/token" +AERODATABOX_HOST = "aerodatabox.p.rapidapi.com" +DUBLIN = ZoneInfo("Europe/Dublin") + + +def _days(start: date, end: date) -> list[date]: + """List every date in a range. + + Args: + start (date): First date. + end (date): Last date, inclusive. + + Returns: + list[date]: Dates from ``start`` to ``end``. + """ + return [start + timedelta(days=i) for i in range((end - start).days + 1)] + + +def fetch_cso(start: date, end: date, airport: str = "Dublin") -> pd.DataFrame: + """Fetch monthly passenger numbers for an Irish airport from CSO table TAM07. + + Args: + start (date): First day of the range. Its whole month is included. + end (date): Last day of the range, inclusive. Its whole month is included. + airport (str): CSO airport name, e.g. "Dublin" or "Cork". Defaults to "Dublin". + + Returns: + pandas.DataFrame: One row per month, country, direction and flight type, with + columns ``year``, ``month``, ``airport``, ``country``, ``direction``, + ``flight_type`` and ``passengers_thousands``. + + Raises: + httpx.HTTPStatusError: If the CSO request fails. + """ + response = httpx.get(CSO_URL, timeout=60, follow_redirects=True) + response.raise_for_status() + df = pd.read_csv(io.StringIO(response.content.decode("utf-8-sig"))) + + df = df[ + (df["Statistic Label"] == "Passengers") & (df["Airports in Ireland"] == airport) + ] + dates = pd.to_datetime(df["Month"], format="%Y %B") + df = pd.DataFrame( + { + "year": dates.dt.year, + "month": dates.dt.month, + "airport": df["Airports in Ireland"], + "country": df["Country"], + "direction": df["Direction"], + "flight_type": df["Flight Type"], + "passengers_thousands": df["VALUE"], + } + ) + first_month = pd.Timestamp(start).replace(day=1) + return df[(dates >= first_month) & (dates <= pd.Timestamp(end))] + + +def fetch_opensky(start: date, end: date, airport: str = "EIDW") -> pd.DataFrame: + """Fetch individual flights arriving at an airport from OpenSky. + + Args: + start (date): First day of the range (Dublin time). + end (date): Last day of the range, inclusive (Dublin time). + airport (str): ICAO airport code. Defaults to "EIDW" (Dublin). + + Returns: + pandas.DataFrame: One row per flight with OpenSky's fields; ``firstSeen``, + ``lastSeen`` and ``arrival_time`` are in Dublin time. Empty if there were no + flights. + + Raises: + httpx.HTTPStatusError: If authentication or a request fails. + """ + token = httpx.post( + OPENSKY_TOKEN_URL, + data={ + "grant_type": "client_credentials", + "client_id": settings.opensky_client_id, + "client_secret": settings.opensky_client_secret, + }, + timeout=30, + ) + token.raise_for_status() + headers = {"Authorization": f"Bearer {token.json()['access_token']}"} + + rows = [] + for day in _days(start, end): + # Midnight to midnight in Dublin; 23 or 25 hours on clock-change days. + begin = datetime.combine(day, time.min, DUBLIN) + end_of_day = datetime.combine(day + timedelta(days=1), time.min, DUBLIN) + response = httpx.get( + OPENSKY_URL, + params={ + "airport": airport, + "begin": int(begin.timestamp()), + "end": int(end_of_day.timestamp()), + }, + headers=headers, + timeout=60, + ) + if response.status_code == 404: # no flights that day + continue + response.raise_for_status() + rows.extend(response.json()) + + df = pd.DataFrame(rows) + if df.empty: + return df + for col in ("firstSeen", "lastSeen"): + df[col] = pd.to_datetime(df[col], unit="s", utc=True).dt.tz_convert(DUBLIN) + df["arrival_time"] = df["lastSeen"] + return df + + +def fetch_aerodatabox(start: date, end: date, airport: str = "DUB") -> pd.DataFrame: + """Fetch scheduled flights arriving at an airport from AeroDataBox (via RapidAPI). + + Each day costs about 4 API units; the free plan has 400 a month. + + Args: + start (date): First day of the range (airport local time). + end (date): Last day of the range, inclusive (airport local time). + airport (str): IATA airport code. Defaults to "DUB" (Dublin). + + Returns: + pandas.DataFrame: One row per flight, with columns ``flight_number``, + ``airline``, ``aircraft_model``, ``origin_iata``, ``origin_name``, + ``origin_country``, ``status`` and ``arrival_time`` (scheduled, in Dublin + time). Empty if there were no flights. + + Raises: + httpx.HTTPStatusError: If a request fails. + """ + headers = { + "X-RapidAPI-Key": settings.xrapidapi_key, + "X-RapidAPI-Host": AERODATABOX_HOST, + } + params = { + "direction": "Arrival", + "withCancelled": "false", + "withCodeshared": "false", + "withCargo": "false", + "withPrivate": "false", + "withLocation": "false", + } + + rows = [] + for day in _days(start, end): + # The API allows at most 12 hours per request. + midnight = datetime.combine(day, time.min) + for window_start in (midnight, midnight + timedelta(hours=12)): + window_end = window_start + timedelta(hours=11, minutes=59) + response = httpx.get( + f"https://{AERODATABOX_HOST}/flights/airports/iata/{airport}" + f"/{window_start:%Y-%m-%dT%H:%M}/{window_end:%Y-%m-%dT%H:%M}", + params=params, + headers=headers, + timeout=30, + ) + response.raise_for_status() + if response.status_code == 204: # no flights in this window + continue + for item in response.json().get("arrivals", []): + move = item.get("movement") or {} + origin = move.get("airport") or {} + rows.append( + { + "flight_number": item.get("number"), + "airline": (item.get("airline") or {}).get("name"), + "aircraft_model": (item.get("aircraft") or {}).get("model"), + "origin_iata": origin.get("iata"), + "origin_name": origin.get("name"), + "origin_country": origin.get("countryCode"), + "status": item.get("status"), + "arrival_time": (move.get("scheduledTime") or {}).get("utc"), + } + ) + + df = pd.DataFrame(rows) + if df.empty: + return df + # A delayed flight near a window edge can be returned twice. + df = df.drop_duplicates(["flight_number", "arrival_time", "origin_iata"]) + df["arrival_time"] = pd.to_datetime(df["arrival_time"], utc=True).dt.tz_convert( + DUBLIN + ) + return df + + +FETCHERS = { + "cso": fetch_cso, + "opensky": fetch_opensky, + "aerodatabox": fetch_aerodatabox, +} + + +def main() -> None: + """Fetch the chosen APIs for a date range and save a CSV for each.""" + parser = argparse.ArgumentParser( + description="Fetch Dublin flight data and save CSVs." + ) + parser.add_argument("--api", nargs="+", choices=FETCHERS, default=list(FETCHERS)) + parser.add_argument( + "--start", type=date.fromisoformat, required=True, help="YYYY-MM-DD" + ) + parser.add_argument( + "--end", type=date.fromisoformat, help="YYYY-MM-DD, defaults to --start" + ) + args = parser.parse_args() + end = args.end or args.start + + OUTPUT_DIR.mkdir(parents=True, exist_ok=True) + for api in args.api: + df = FETCHERS[api](args.start, end) + path = OUTPUT_DIR / f"{api}_{args.start}_{end}.csv" + df.to_csv(path, index=False) + print(f"[{api}] {len(df)} rows saved to {path}") + + +if __name__ == "__main__": + main() diff --git a/backend/pyproject.toml b/backend/pyproject.toml index 14bcc7f..1acadc0 100644 --- a/backend/pyproject.toml +++ b/backend/pyproject.toml @@ -18,6 +18,7 @@ dependencies = [ "pandas>=2.2", "psycopg2-binary>=2.9", "scikit-learn>=1.9.1", + "httpx>=0.28", ] [dependency-groups] @@ -29,7 +30,6 @@ dev = [ "pre-commit>=4.2", "pytest>=8.0", "pytest-asyncio>=0.25", - "httpx>=0.28", ] ml = [ "matplotlib>=3.11.2", diff --git a/backend/scripts/fetch_opensky.py b/backend/scripts/fetch_opensky.py new file mode 100644 index 0000000..d03c60f --- /dev/null +++ b/backend/scripts/fetch_opensky.py @@ -0,0 +1,117 @@ +r"""Download historical airport arrivals from OpenSky in 2-day chunks. + +Each chunk is saved as its own CSV in ``data/flights/opensky/`` and chunks that already +exist are skipped, so re-running the same command resumes where it stopped. The script +stops when OpenSky's daily credit limit is reached. + +OpenSky allows at most 2 days per arrivals request. A request covering up to 2 UTC days +costs 30 credits and a standard account gets 4,000 credits a day, so a day's credits +cover about 133 chunks (about 266 days of flights). Dates are UTC days so that every +request stays within 2 UTC days. Only flights up to yesterday (UTC) are available. + +Example: + Run from ``backend/``:: + + uv run python scripts/fetch_opensky.py --start 2022-01-01 --end 2025-12-31 +""" + +import argparse +from datetime import UTC, date, datetime, time, timedelta +from pathlib import Path + +import httpx +import pandas as pd + +from foresight.config import settings + +OUTPUT_DIR = Path(__file__).resolve().parents[1] / "data" / "flights" / "opensky" +ARRIVAL_URL = "https://opensky-network.org/api/flights/arrival" +TOKEN_URL = "https://auth.opensky-network.org/auth/realms/opensky-network/protocol/openid-connect/token" + + +def get_token() -> str: + """Get an OpenSky access token, which expires after 30 minutes. + + Returns: + str: The bearer token. + + Raises: + httpx.HTTPStatusError: If authentication fails. + """ + response = httpx.post( + TOKEN_URL, + data={ + "grant_type": "client_credentials", + "client_id": settings.opensky_client_id, + "client_secret": settings.opensky_client_secret, + }, + timeout=30, + ) + response.raise_for_status() + return response.json()["access_token"] + + +def main() -> None: + """Download arrivals for a date range, one CSV per 2-day chunk.""" + parser = argparse.ArgumentParser( + description="Download OpenSky airport arrivals in 2-day chunks." + ) + parser.add_argument( + "--start", type=date.fromisoformat, required=True, help="YYYY-MM-DD (UTC)" + ) + parser.add_argument( + "--end", + type=date.fromisoformat, + help="YYYY-MM-DD (UTC), defaults to yesterday", + ) + parser.add_argument( + "--airport", default="EIDW", help="ICAO code, defaults to Dublin (EIDW)" + ) + args = parser.parse_args() + + yesterday = datetime.now(UTC).date() - timedelta(days=1) + end = min(args.end or yesterday, yesterday) + OUTPUT_DIR.mkdir(parents=True, exist_ok=True) + token = get_token() + + for offset in range(0, (end - args.start).days + 1, 2): + first = args.start + timedelta(days=offset) + last = min(first + timedelta(days=1), end) + path = OUTPUT_DIR / f"{args.airport}_{first}_{last}.csv" + if path.exists(): + continue + + # Midnight UTC to the last second of ``last``: never more than 2 UTC days. + begin = datetime.combine(first, time.min, UTC) + until = datetime.combine(last + timedelta(days=1), time.min, UTC) + params = { + "airport": args.airport, + "begin": int(begin.timestamp()), + "end": int(until.timestamp()) - 1, + } + for _ in range(2): + response = httpx.get( + ARRIVAL_URL, + params=params, + headers={"Authorization": f"Bearer {token}"}, + timeout=60, + ) + if response.status_code != 401: + break + token = get_token() # the token expired, so get a new one and retry + + if response.status_code == 429: + wait = response.headers.get("X-Rate-Limit-Retry-After-Seconds") + print(f"Credit limit reached, retry in {wait}s. Re-run to resume.") + return + if response.status_code == 404: # no flights in this chunk + rows = [] + else: + response.raise_for_status() + rows = response.json() + pd.DataFrame(rows).to_csv(path, index=False) + print(f"{first} to {last}: {len(rows)} flights") + + +if __name__ == "__main__": + main() diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index a341be7..38c6ca4 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -2,7 +2,6 @@ import pandas as pd import pytest - from generate_mock_data import write_csv # Known only after the night, so never used to predict price. diff --git a/backend/uv.lock b/backend/uv.lock index 87cadef..add21bc 100644 --- a/backend/uv.lock +++ b/backend/uv.lock @@ -331,6 +331,7 @@ dependencies = [ { name = "alembic" }, { name = "asyncpg" }, { name = "fastapi" }, + { name = "httpx" }, { name = "pandas" }, { name = "psycopg2-binary" }, { name = "pydantic-settings" }, @@ -342,7 +343,6 @@ dependencies = [ [package.dev-dependencies] dev = [ - { name = "httpx" }, { name = "pre-commit" }, { name = "pytest" }, { name = "pytest-asyncio" }, @@ -362,6 +362,7 @@ requires-dist = [ { name = "alembic", specifier = ">=1.15" }, { name = "asyncpg", specifier = ">=0.30" }, { name = "fastapi", specifier = ">=0.115" }, + { name = "httpx", specifier = ">=0.28" }, { name = "pandas", specifier = ">=2.2" }, { name = "psycopg2-binary", specifier = ">=2.9" }, { name = "pydantic-settings", specifier = ">=2.9" }, @@ -373,7 +374,6 @@ requires-dist = [ [package.metadata.requires-dev] dev = [ - { name = "httpx", specifier = ">=0.28" }, { name = "pre-commit", specifier = ">=4.2" }, { name = "pytest", specifier = ">=8.0" }, { name = "pytest-asyncio", specifier = ">=0.25" },