""" ================================================================ Auteur : Simon-Pierre Boucher Contact : contact@spboucher.ai Projet : Prévision de volatilité réalisée multi-actifs (HAR-RV vs GARCH vs Machine Learning) Fichier : api_client.py Description : Client réutilisable pour l'API HF Market Data — rate limiting, retry avec backoff exponentiel, pagination et cache local parquet. ================================================================ """ from __future__ import annotations import io import logging import time from pathlib import Path from typing import Any import pandas as pd import requests from . import config logger = logging.getLogger(__name__) _RETRYABLE = {429, 500, 502, 503, 504} class HFMarketDataClient: """Client for https://www.hfmarketdata.io with caching and retries. Parameters ---------- base_url : str, optional API root; defaults to the value in ``config.yaml``. cache_dir : pathlib.Path, optional Directory for the local parquet cache; defaults to ``data/raw``. Notes ----- Bars are fetched in ascending order and paginated by advancing the ``start`` parameter one minute past the last received timestamp until an empty page is returned — this is robust to any server-side row cap. """ def __init__(self, base_url: str | None = None, cache_dir: Path | None = None) -> None: cfg = config.load_config()["api"] self.base_url = (base_url or cfg["base_url"]).rstrip("/") self.cache_dir = Path(cache_dir) if cache_dir else config.path("raw") self.min_interval = 1.0 / float(cfg.get("rate_limit_per_sec", 4)) self.max_retries = int(cfg.get("max_retries", 5)) self.backoff_base = float(cfg.get("backoff_base_sec", 1.0)) self.timeout = float(cfg.get("timeout_sec", 60)) self.page_limit = int(cfg.get("page_limit", 50_000)) self._last_request_ts = 0.0 self._session = requests.Session() # ------------------------------------------------------------- low level def _throttle(self) -> None: """Sleep as needed to respect the configured request rate.""" wait = self._last_request_ts + self.min_interval - time.monotonic() if wait > 0: time.sleep(wait) self._last_request_ts = time.monotonic() def _get(self, path: str, params: dict[str, Any] | None = None) -> requests.Response: """Issue a GET with throttling and exponential-backoff retries. Parameters ---------- path : str Path relative to the API root (e.g. ``/v1/status``). params : dict, optional Query-string parameters. Returns ------- requests.Response The successful response (status code < 400). Raises ------ requests.HTTPError If the request keeps failing after ``max_retries`` attempts. """ url = f"{self.base_url}{path}" last_exc: Exception | None = None for attempt in range(self.max_retries + 1): self._throttle() try: resp = self._session.get(url, params=params, timeout=self.timeout) if resp.status_code in _RETRYABLE: raise requests.HTTPError(f"HTTP {resp.status_code}", response=resp) resp.raise_for_status() return resp except (requests.ConnectionError, requests.Timeout, requests.HTTPError) as exc: resp_obj = getattr(exc, "response", None) if resp_obj is not None and resp_obj.status_code not in _RETRYABLE: raise last_exc = exc delay = self.backoff_base * (2**attempt) logger.warning( "GET %s failed (%s), retry %d/%d in %.1fs", path, exc, attempt + 1, self.max_retries, delay, ) time.sleep(delay) raise requests.HTTPError(f"GET {url} failed after {self.max_retries} retries") from last_exc def get_json(self, path: str, params: dict[str, Any] | None = None) -> Any: """GET a JSON endpoint and return the decoded payload.""" return self._get(path, params).json() # ------------------------------------------------------------- endpoints def status(self) -> dict[str, Any]: """Return the ``/v1/status`` dataset inventory.""" return self.get_json("/v1/status") def tickers(self, asset: str, search: str | None = None, limit: int = 100) -> list[str]: """List tickers available for an asset class. Parameters ---------- asset : str API asset class (``stock``, ``etf``, ``futures``, ``crypto``, ``index``, ``fx``). search : str, optional Substring filter applied server-side. limit : int Maximum number of tickers returned. """ params: dict[str, Any] = {"limit": limit} if search: params["search"] = search payload = self.get_json(f"/v1/{asset}/tickers", params) return payload.get("tickers", payload) def fetch_bars( self, asset: str, ticker: str, timeframe: str = "1min", adjustment: str | None = None, start: str | None = None, end: str | None = None, ) -> pd.DataFrame: """Download OHLCV bars, transparently handling pagination. Parameters ---------- asset, ticker : str Instrument identification. timeframe : str ``1min``, ``5min``, ``30min``, ``1hour`` or ``1day``. adjustment : str, optional Price-adjustment scheme; defaults to the class setting in ``config.yaml``. start, end : str, optional Inclusive date bounds (``YYYY-MM-DD`` or full datetime). Returns ------- pandas.DataFrame Columns ``datetime`` (naive exchange-local timestamps), ``open``, ``high``, ``low``, ``close``, ``volume``; sorted, duplicate timestamps dropped. """ adjustment = adjustment or config.adjustment_for(asset) frames: list[pd.DataFrame] = [] cursor = start n_pages = 0 while True: params: dict[str, Any] = { "timeframe": timeframe, "adjustment": adjustment, "order": "asc", "limit": self.page_limit, "format": "csv", } if cursor: params["start"] = cursor if end: params["end"] = end resp = self._get(f"/v1/bars/{asset}/{ticker}", params) page = pd.read_csv(io.StringIO(resp.text)) if resp.text.strip() else pd.DataFrame() if page.empty: break page["datetime"] = pd.to_datetime(page["datetime"]) frames.append(page) n_pages += 1 last_dt = page["datetime"].iloc[-1] logger.debug("%s/%s %s page %d: %d rows, up to %s", asset, ticker, timeframe, n_pages, len(page), last_dt) if len(page) < self.page_limit: break cursor = (last_dt + pd.Timedelta(minutes=1)).strftime("%Y-%m-%d %H:%M:%S") if not frames: logger.warning("%s/%s %s: no data returned", asset, ticker, timeframe) return pd.DataFrame(columns=["datetime", "open", "high", "low", "close", "volume"]) df = pd.concat(frames, ignore_index=True) df = ( df.drop(columns=["ticker"], errors="ignore") .drop_duplicates(subset="datetime") .sort_values("datetime") .reset_index(drop=True) ) if "volume" not in df.columns: # some index series carry no volume df["volume"] = float("nan") logger.info("%s/%s %s: %d bars (%s -> %s)", asset, ticker, timeframe, len(df), df["datetime"].iloc[0], df["datetime"].iloc[-1]) return df # ------------------------------------------------------------- cache def _cache_file(self, asset: str, ticker: str, timeframe: str) -> Path: safe = ticker.replace("/", "-") return self.cache_dir / f"{asset}_{safe}_{timeframe}.parquet" def get_bars( self, asset: str, ticker: str, timeframe: str = "1min", start: str | None = None, end: str | None = None, refresh: bool = False, ) -> pd.DataFrame: """Return bars from the local parquet cache, downloading if absent. Parameters ---------- asset, ticker, timeframe : str Instrument identification and bar frequency. start, end : str, optional Date bounds used when the series must be downloaded; when the cache is hit the full cached range is returned (filter at the call site if needed). refresh : bool Force a re-download even when a cache file exists. Returns ------- pandas.DataFrame Same layout as :meth:`fetch_bars`. """ f = self._cache_file(asset, ticker, timeframe) if f.exists() and not refresh: logger.debug("cache hit: %s", f.name) return pd.read_parquet(f) df = self.fetch_bars(asset, ticker, timeframe, start=start, end=end) if not df.empty: f.parent.mkdir(parents=True, exist_ok=True) df.to_parquet(f, index=False) logger.info("cached %s (%.1f MB)", f.name, f.stat().st_size / 1e6) return df