from __future__ import annotations import json from datetime import date, timedelta import httpx import pandas as pd from loguru import logger from ..config import API_BASE_URL # Keep chunks small to avoid truncated responses from the API. _MAX_DAYS_PER_REQUEST = 90 _NDJSON_HEADERS = {"Accept": "application/x-ndjson"} def _build_url(start: date, end: date, base: str) -> str: return f"{API_BASE_URL}/rates?from={start.isoformat()}&to={end.isoformat()}&base={base}" def _parse_ndjson(text: str) -> list[dict]: """Parse NDJSON (one JSON object per line) into our column schema.""" rows: list[dict] = [] for line in text.splitlines(): line = line.strip() if not line: continue row = json.loads(line) rows.append( { "rate_date": row["date"], "base_currency": row["base"], "quote_currency": row["quote"], "rate": row["rate"], } ) return rows def fetch_rates(start_date: date, end_date: date, base: str = "USD") -> pd.DataFrame: """Fetch daily spot rates from Frankfurter API for a date range. Large ranges are automatically chunked into 90-day requests. Uses NDJSON streaming to avoid JSON truncation on large responses. Returns a DataFrame with columns: rate_date, base_currency, quote_currency, rate. """ all_rows: list[dict] = [] chunk_start = start_date with httpx.Client(timeout=120.0) as client: while chunk_start <= end_date: chunk_end = min(chunk_start + timedelta(days=_MAX_DAYS_PER_REQUEST), end_date) url = _build_url(chunk_start, chunk_end, base) logger.info(f"Fetching rates {chunk_start} -> {chunk_end}") response = client.get(url, headers=_NDJSON_HEADERS) response.raise_for_status() rows = _parse_ndjson(response.text) all_rows.extend(rows) logger.debug(f" Got {len(rows)} rate rows") chunk_start = chunk_end + timedelta(days=1) df = pd.DataFrame(all_rows) if not df.empty: df["rate_date"] = pd.to_datetime(df["rate_date"]).dt.date logger.info(f"Total: {len(df)} rate rows fetched") return df