From fa078a57609432298bd723b1158e09e352a8ebe2 Mon Sep 17 00:00:00 2001 From: Mikko Kotila Date: Fri, 2 Oct 2026 12:24:19 +0300 Subject: [PATCH] Add Hyperliquid market data and reusable weighted IP routing --- README.md | 25 ++- dedomena/sources/__init__.py | 6 +- dedomena/sources/__main__.py | 74 ++++++- dedomena/sources/core.py | 151 ++++++++++--- dedomena/sources/egress.py | 315 +++++++++++++++++++++++++++ dedomena/sources/hyperliquid.py | 332 +++++++++++++++++++++++++++++ dedomena/sources/openalex.py | 5 + docs/FINANCE.md | 3 +- docs/HYPERLIQUID.md | 123 +++++++++++ docs/IP_ROUTING.md | 121 +++++++++++ setup.py | 5 +- tests/test_egress.py | 362 ++++++++++++++++++++++++++++++++ tests/test_hyperliquid.py | 297 ++++++++++++++++++++++++++ tests/test_hyperliquid_cli.py | 125 +++++++++++ tests/test_ip_transport.py | 214 +++++++++++++++++++ 15 files changed, 2109 insertions(+), 49 deletions(-) create mode 100644 dedomena/sources/egress.py create mode 100644 dedomena/sources/hyperliquid.py create mode 100644 docs/HYPERLIQUID.md create mode 100644 docs/IP_ROUTING.md create mode 100644 tests/test_egress.py create mode 100644 tests/test_hyperliquid.py create mode 100644 tests/test_hyperliquid_cli.py create mode 100644 tests/test_ip_transport.py diff --git a/README.md b/README.md index bc5cf6e..5660817 100644 --- a/README.md +++ b/README.md @@ -1,7 +1,7 @@ # Dedomena Consistent research data access for agents and scientists across biology and finance. -OpenAlex, Europe PMC, EPO, SEC EDGAR, FRED/ALFRED, ECB and World Bank share a +OpenAlex, Europe PMC, EPO, SEC EDGAR, FRED/ALFRED, ECB, World Bank and Hyperliquid share a streaming page and provenance contract. Queries retain their native source semantics. ## Install @@ -38,7 +38,7 @@ source identifiers, exact request provenance, SHA-256 hashes, a saved-response snapshot ID, retrieval times, completeness, cost and transfer size. Partial enumeration fails explicitly. Native records retain source-specific fields. -Default caches persist for 24 hours at `~/.cache/dedomena/sources.sqlite3`. +Research and traditional finance caches persist for 24 hours at `~/.cache/dedomena/sources.sqlite3`. `refresh=True` fetches new data while keeping previous snapshots for offline replay. `DEDOMENA_SOURCE_STORE` selects a shared store, including across Canary workers. @@ -73,6 +73,25 @@ FRED supports ALFRED knowledge dates; SEC retains amendments; ECB exposes revisions; World Bank serves current revised indicators. See [finance access, throughput and semantics](docs/FINANCE.md). +~~~python +from dedomena.sources import Hyperliquid + +with Hyperliquid() as source: # Public market data, no key or wallet + markets = source.markets() # Native metadata and asset contexts + book = source.order_book("BTC") + consume(book.records, book.provenance.to_dict()) +~~~ + +Hyperliquid adds mids, spot/perpetual metadata, books, candles and funding history. +Market snapshots default to a zero cache TTL. Candle retention is explicitly +incomplete; funding supports timestamp checkpoints. +See [Hyperliquid retrieval and weight](docs/HYPERLIQUID.md). + +`IPPool` routes the same source clients through owned local IPs or proxies, with +shared weighted per-IP admission, key budgets and cooldowns. Reuse one pool for +Hyperliquid and OpenAlex; extra IPs do not multiply OpenAlex's key allowance. +See [shared IP routing examples](docs/IP_ROUTING.md). + ## Agent CLI ~~~sh @@ -80,6 +99,8 @@ python -m dedomena.sources benchmark openalex 'CRISPR' python -m dedomena.sources benchmark europepmc 'TITLE:CRISPR' python -m dedomena.sources search epo 'ta="CRISPR"' --max-pages 1 python -m dedomena.sources quota openalex +python -m dedomena.sources markets hyperliquid +python -m dedomena.sources book hyperliquid BTC python -m dedomena.sources fetch sec 320193 python -m dedomena.sources benchmark fred 53 --operation release python -m dedomena.sources observations fred GDP --as-of 2020-01-01 diff --git a/dedomena/sources/__init__.py b/dedomena/sources/__init__.py index e9e5c54..a2e5ab7 100644 --- a/dedomena/sources/__init__.py +++ b/dedomena/sources/__init__.py @@ -1,6 +1,7 @@ """Agent-ready research sources; streaming pages share one provenance contract.""" from .core import (BudgetExceeded, InvalidResponse, Page, Provenance, SearchLimitExceeded, SourceError, Store, Throttled) +from .egress import IPPool, IPRoute from .openalex import OpenAlex from .europepmc import EuropePMC from .epo import EPO @@ -8,8 +9,9 @@ from .fred import FRED from .ecb import ECB from .worldbank import WorldBank +from .hyperliquid import Hyperliquid __all__ = [ - "OpenAlex", "EuropePMC", "EPO", "SEC", "FRED", "ECB", "WorldBank", "Page", "Provenance", "Store", - "SourceError", "BudgetExceeded", "Throttled", "InvalidResponse", "SearchLimitExceeded", + "OpenAlex", "EuropePMC", "EPO", "SEC", "FRED", "ECB", "WorldBank", "Hyperliquid", "Page", "Provenance", "Store", + "IPPool", "IPRoute", "SourceError", "BudgetExceeded", "Throttled", "InvalidResponse", "SearchLimitExceeded", ] diff --git a/dedomena/sources/__main__.py b/dedomena/sources/__main__.py index eeb9c14..2ed78a3 100644 --- a/dedomena/sources/__main__.py +++ b/dedomena/sources/__main__.py @@ -1,5 +1,6 @@ """Agent-friendly JSON CLI. Credentials are read exclusively from the environment.""" import argparse +from contextlib import ExitStack from datetime import date import json import os @@ -7,16 +8,41 @@ import sys import time -from . import ECB, EPO, FRED, SEC, EuropePMC, OpenAlex, SourceError, Store, WorldBank +from . import (ECB, EPO, FRED, SEC, EuropePMC, Hyperliquid, IPPool, IPRoute, + OpenAlex, SourceError, Store, WorldBank) + + +def _ip_pool(path): + try: + with Path(path).expanduser().open("rb") as handle: + data = handle.read(65537) + if len(data) > 65536: + raise ValueError + routes = json.loads(data) + except (OSError, ValueError, UnicodeDecodeError): + raise ValueError("IP routes require a readable JSON file of at most 64 KiB") from None + if (not isinstance(routes, list) or not 1 <= len(routes) <= 1000 + or any(not isinstance(route, dict) or "public_ip" not in route + or not set(route).issubset({"public_ip", "local_address", "proxy"}) + for route in routes)): + raise ValueError("IP routes require public_ip and one local_address or proxy per entry") + return IPPool(IPRoute(**route) for route in routes) def main(argv=None): parser = argparse.ArgumentParser(description=__doc__) - operations = ("search", "observations", "release", "catalogue") + market_operations = ("markets", "mids", "book", "candles", "funding") + operations = ("search", "observations", "release", "catalogue", *market_operations) parser.add_argument("action", choices=(*operations, "fetch", "quota", "usage", "replay", "benchmark")) - parser.add_argument("source", choices=("openalex", "europepmc", "epo", "sec", "fred", "ecb", "worldbank")) + parser.add_argument("source", choices=("openalex", "europepmc", "epo", "sec", "fred", "ecb", "worldbank", "hyperliquid")) parser.add_argument("query", nargs="?") parser.add_argument("--store", help="Shared SQLite snapshot/allowance store") + parser.add_argument("--ip-routes", help="JSON file of owned public IP routes and local bindings/proxies") + parser.add_argument("--start-time", type=int, help="Hyperliquid start time in epoch milliseconds") + parser.add_argument("--end-time", type=int, help="Hyperliquid end time in epoch milliseconds") + parser.add_argument("--interval", help="Hyperliquid candle interval, e.g. 1h") + parser.add_argument("--spot", action="store_true", help="Hyperliquid spot market metadata") + parser.add_argument("--dex", help="Hyperliquid perpetual DEX name") parser.add_argument("--profile", choices=("discovery", "evidence"), default="discovery") parser.add_argument("--filter") parser.add_argument("--refresh", action="store_true") @@ -45,24 +71,54 @@ def main(argv=None): parser.error("SDMX period/delta/last-n options require ECB search or observations") if args.action != "benchmark" and args.operation != "search": parser.error("--operation selects the benchmark operation only") + if operation in market_operations and args.source != "hyperliquid": + parser.error("market operations require hyperliquid") + if args.source == "hyperliquid" and operation in ("search", "observations", "release"): + parser.error("hyperliquid supports markets, mids, book, candles, funding, and fetch") + if (args.start_time is not None or args.end_time is not None or args.interval) and ( + args.source != "hyperliquid" or operation not in ("candles", "funding")): + parser.error("epoch time and interval options require Hyperliquid candles or funding") + if args.interval and operation != "candles": + parser.error("--interval requires candles") + if operation in ("candles", "funding") and args.start_time is None: + parser.error("historical market operations require --start-time") + if operation == "candles" and (args.end_time is None or not args.interval): + parser.error("candles require --end-time and --interval") + if args.spot and (args.source != "hyperliquid" or operation not in ("markets", "catalogue")): + parser.error("--spot requires Hyperliquid markets or catalogue") + if args.dex is not None and (args.source != "hyperliquid" or operation not in ("markets", "catalogue", "mids")): + parser.error("--dex requires Hyperliquid markets or mids") if operation == "release" and args.source != "fred": parser.error("release requires fred") - if operation == "catalogue" and args.source not in ("ecb", "worldbank"): - parser.error("catalogue requires ecb or worldbank") + if operation == "catalogue" and args.source not in ("ecb", "worldbank", "hyperliquid"): + parser.error("catalogue requires ecb, worldbank or hyperliquid") if operation == "observations" and args.source not in ("fred", "ecb", "worldbank"): parser.error("observations requires fred, ecb, or worldbank") - if operation in ("fetch", "replay", "release", "observations") and not args.query: + if operation in ("fetch", "replay", "release", "observations", "book", "candles", "funding") and not args.query: parser.error("this action needs an identifier") if operation == "search" and not args.query and not (args.source == "openalex" and args.filter): parser.error("search needs a query, or an OpenAlex filter") if bool(args.start_date) != bool(args.end_date) or (args.start_date and args.source != "epo"): parser.error("date partitions require both dates and the epo source") store = None + resources = ExitStack() def emit(value): print(json.dumps(value, ensure_ascii=False, allow_nan=False)) def pages(source): + if args.source == "hyperliquid": + if operation in ("markets", "catalogue"): + return iter((source.markets(spot=args.spot, dex=args.dex or "", refresh=args.refresh),)) + if operation == "mids": + return iter((source.all_mids(dex=args.dex or "", refresh=args.refresh),)) + if operation == "book": + return iter((source.order_book(args.query, refresh=args.refresh),)) + if operation == "candles": + return iter((source.candles(args.query, args.interval, start_time=args.start_time, + end_time=args.end_time, refresh=args.refresh),)) + return source.funding_history(args.query, start_time=args.start_time, + end_time=args.end_time, refresh=args.refresh) if operation == "release": if not args.query.isascii() or not args.query.isdecimal(): raise ValueError("release identifier must be a positive integer") @@ -105,8 +161,11 @@ def pages(source): return 0 store = Store(args.store) if args.store else None constructors = {"openalex": OpenAlex, "europepmc": EuropePMC, "epo": EPO, - "sec": SEC, "fred": FRED, "ecb": ECB, "worldbank": WorldBank} + "sec": SEC, "fred": FRED, "ecb": ECB, "worldbank": WorldBank, + "hyperliquid": Hyperliquid} kwargs = {"store": store} if store else {} + if args.ip_routes: + kwargs["ip_pool"] = resources.enter_context(_ip_pool(args.ip_routes)) if args.source == "openalex": kwargs["profile"] = args.profile elif args.source == "europepmc": @@ -154,6 +213,7 @@ def pages(source): print(json.dumps(emit_error), file=sys.stderr) return 2 finally: + resources.close() if store: store.close() return 0 diff --git a/dedomena/sources/core.py b/dedomena/sources/core.py index 990c6e0..5d403e8 100644 --- a/dedomena/sources/core.py +++ b/dedomena/sources/core.py @@ -21,7 +21,7 @@ import httpx -SCHEMA_VERSION = "0.5.0" +SCHEMA_VERSION = "0.6.0" PUBLIC_HEADERS = { "content-type", "content-encoding", "x-ratelimit-limit", "x-ratelimit-remaining", "x-ratelimit-credits-used", "x-ratelimit-reset", "retry-after", @@ -95,6 +95,7 @@ class Provenance: request_body: str | None = None completed_at_utc: str | None = None http_status: int | None = None + egress_ip_sha256: str | None = None def to_dict(self) -> dict: return asdict(self) @@ -279,13 +280,18 @@ def charge(self, scope: str, unit: str, amount: int, limit: int | None, period: def reserve(self, scope: str, unit: str, amount: int, limit: int | None, period: str, *, window_key: str) -> str: """Reserve atomically without mixing refundable estimates with provider usage.""" - reservation = uuid.uuid4().hex with self.transaction() as db: - baseline, pending = self._budget(db, scope, unit, window_key) - if limit is not None and baseline + pending + amount > limit: - raise BudgetExceeded(f"{scope.split(':')[0]}: {period} {unit} allowance exhausted") - db.execute("INSERT INTO budget_reservations VALUES(?,?,?,?,?)", - (reservation, scope, unit, window_key, amount)) + return self._reserve_in(db, scope, unit, amount, limit, period, window_key=window_key) + + def _reserve_in(self, db, scope: str, unit: str, amount: int, limit: int | None, + period: str, *, window_key: str) -> str: + """Reserve inside the caller's final dispatch transaction.""" + reservation = uuid.uuid4().hex + baseline, pending = self._budget(db, scope, unit, window_key) + if limit is not None and baseline + pending + amount > limit: + raise BudgetExceeded(f"{scope.split(':')[0]}: {period} {unit} allowance exhausted") + db.execute("INSERT INTO budget_reservations VALUES(?,?,?,?,?)", + (reservation, scope, unit, window_key, amount)) return reservation def settle(self, reservation: str, actual: int): @@ -342,7 +348,7 @@ def pace(self, scope: str, rate: float, sleep: Callable[[float], None], max_wait self.pace_many(((scope, rate),), sleep, max_wait) def pace_many(self, limits: tuple[tuple[str, float], ...], - sleep: Callable[[float], None], max_wait: float): + sleep: Callable[[float], None], max_wait: float, before_admit=None): """Admit all rate scopes together; waiters acquire no stale future dispatch slots.""" began = self.clock() while True: @@ -358,9 +364,16 @@ def pace_many(self, limits: tuple[tuple[str, float], ...], due = max(due, row[0] if row else now, cooldown[0] if cooldown else now) rates.append((scope, effective)) delay = due - now - if delay > max(0, max_wait - (now - began)): + if (delay > max(0, max_wait - (now - began)) + or (max_wait > 0 and now - began > max_wait)): raise Throttled(limits[0][0].split(":")[0], delay) if not delay: + if before_admit is not None: + before_admit(db, now) + # No separately locking work may intervene between admission and send. + now = self.clock() + if max_wait > 0 and now - began > max_wait: + raise Throttled(limits[0][0].split(":")[0], 0) for scope, rate in rates: db.execute("INSERT OR REPLACE INTO pacing VALUES(?,?)", (scope, now + 1 / rate)) return @@ -382,7 +395,8 @@ def __init__(self, source: str, base_url: str, *, credential: str = "", requests_per_second: float = 10, cache_ttl: float = 86400, max_response_bytes: int = 16_000_000, max_attempts: int = 3, max_wait: float = 60, sleep: Callable[[float], None] = time.sleep, - license: str = ""): + license: str = "", ip_pool=None, ip_limit: tuple[int, float] | None = None, + ip_throttle_only: bool = False): if not base_url.startswith("https://"): raise ValueError("source URL must use HTTPS") if requests_per_second <= 0 or not math.isfinite(requests_per_second): @@ -392,13 +406,29 @@ def __init__(self, source: str, base_url: str, *, credential: str = "", or type(max_attempts) is not int or max_attempts < 1 or not math.isfinite(max_wait) or max_wait < 0): raise ValueError("invalid transport limits") + from .egress import IPPool, IPRoute + if ip_pool is not None and (not isinstance(ip_pool, IPPool) or client is not None): + raise ValueError("ip_pool must be an IPPool and cannot be combined with client") + if ip_limit is not None and ( + not isinstance(ip_limit, tuple) or len(ip_limit) != 2 + or type(ip_limit[0]) is not int or ip_limit[0] < 1 + or type(ip_limit[1]) not in (int, float) + or not math.isfinite(ip_limit[1]) or ip_limit[1] <= 0): + raise ValueError("ip_limit must contain positive integer capacity and window seconds") + if type(ip_throttle_only) is not bool: + raise ValueError("ip_throttle_only must be boolean") self.source, self.base_url = source, base_url.rstrip("/") self.scope = source + ":" + digest(credential.encode()) self._secrets = tuple(value for value in {credential, *credential.split(":")} if len(value) >= 8) - self.client = client or httpx.Client(timeout=30, follow_redirects=False, - limits=httpx.Limits(max_connections=100, - max_keepalive_connections=100)) - self._owns_client = client is None + self.client = None if ip_pool is not None else client or httpx.Client( + timeout=30, follow_redirects=False, + limits=httpx.Limits(max_connections=100, max_keepalive_connections=100)) + self._owns_client = client is None and ip_pool is None + self.ip_pool = ip_pool + self._owns_pool = ip_pool is None and ip_limit is not None + if self._owns_pool: + self.ip_pool = IPPool([IPRoute("0.0.0.0", client=self.client)]) + self.ip_limit, self.ip_throttle_only = ip_limit, ip_throttle_only self.store = store or Store(os.environ.get("DEDOMENA_SOURCE_STORE", str(Path.home() / ".cache" / "dedomena" / "sources.sqlite3"))) self._owns_store = store is None @@ -407,6 +437,8 @@ def __init__(self, source: str, base_url: str, *, credential: str = "", self.max_wait, self.sleep, self.license = max_wait, sleep, license def close(self): + if self._owns_pool: + self.ip_pool.close() if self._owns_client: self.client.close() if self._owns_store: @@ -452,10 +484,12 @@ def _request(self, method: str, path: str, *, params: Mapping | None = None, refresh: bool = False, sensitive: bool = False, observe: Callable[..., None] | None = None, before_request: Callable[[], None] | None = None, - headers_factory: Callable[[], Mapping] | None = None) -> Response: + headers_factory: Callable[[], Mapping] | None = None, + rate_weight: int = 1, actual_weight: Callable[[bytes], int] | None = None) -> Response: if method not in ("GET", "POST") or not path.startswith("/") or path.startswith("//") or "?" in path: raise ValueError("invalid source request") routes = { + "hyperliquid": (("POST", r"/info"),), "fred": (("GET", r"/(?:series(?:/search|/observations|/vintagedates)?|v2/release/observations)"),), "worldbank": (("GET", r"/country/[A-Za-z0-9;]+/indicator/[A-Za-z0-9._;\-]+"), ("GET", r"/indicator(?:/[A-Za-z0-9._;\-]+)?")), @@ -477,6 +511,13 @@ def _request(self, method: str, path: str, *, params: Mapping | None = None, method == verb and re.fullmatch(pattern, path) for verb, pattern in routes[self.source]): raise ValueError("source route is outside the read-only adapter") + if self.source == "hyperliquid": + from .hyperliquid import info_policy + if params or private_params or form is not None or auth is not None: + raise ValueError("Hyperliquid info requests use a public JSON body only") + rate_weight, actual_weight = info_policy(content) + if type(rate_weight) is not int or rate_weight < 1: + raise ValueError("rate_weight must be a positive integer") params = {str(k): str(v) for k, v in (params or {}).items()} if any(k.lower() in ("api_key", "access_token", "authorization") for k in params): raise ValueError("credentials must not appear in source parameters") @@ -531,23 +572,44 @@ def _request(self, method: str, path: str, *, params: Mapping | None = None, limits = ((self.scope + ":all", self.rate),) if service_rate is not None: limits += ((pace_scope, service_rate),) - self.store.pace_many(limits, self.sleep, self.max_wait) - credit_window = window("day", self.store.clock()) - byte_window = window("week", self.store.clock()) - credit_reservation = self.store.reserve(self.scope, "credits", estimated_credits, - credit_limit, "day", window_key=credit_window) - byte_reservation = None - byte_reserved = self.max_response_bytes if byte_limit is not None else 0 - try: - if byte_reserved: - byte_reservation = self.store.reserve(self.scope, "bytes", byte_reserved, - byte_limit, "week", window_key=byte_window) - except BudgetExceeded: - self.store.settle(credit_reservation, 0) - raise + admitted = {} + + def reserve_budgets(db, now): + # IP capacity, pacing and these reservations commit together. No + # additional database lock can leave a queued dispatch lease stale. + while True: + day, week = window("day", now), window("week", now) + credit = self.store._reserve_in(db, self.scope, "credits", estimated_credits, + credit_limit, "day", window_key=day) + byte = (self.store._reserve_in(db, self.scope, "bytes", self.max_response_bytes, + byte_limit, "week", window_key=week) + if byte_limit is not None else None) + fresh = self.store.clock() + if day == window("day", fresh) and week == window("week", fresh): + admitted.update(credit=credit, byte=byte, day=day, week=week) + return + # Reanchor if even the local transaction spans a quota reset. + db.execute("DELETE FROM budget_reservations WHERE id=?", (credit,)) + if byte: + db.execute("DELETE FROM budget_reservations WHERE id=?", (byte,)) + now = fresh + + lease = None + if self.ip_pool is not None: + capacity, period = self.ip_limit or (None, 60.0) + lease = self.ip_pool.acquire(self.store, self.source, rate_weight, capacity, + period, self.sleep, self.max_wait, pacing=limits, + before_admit=reserve_budgets) + else: + self.store.pace_many(limits, self.sleep, self.max_wait, reserve_budgets) + active_client = lease.client if lease is not None else self.client + weight_stats = ({"ip_weight_reserved": rate_weight, "ip_weight_used": rate_weight} + if lease is not None else {}) + credit_window, byte_window = admitted["day"], admitted["week"] + credit_reservation, byte_reservation = admitted["credit"], admitted["byte"] started = self.store.clock() try: - with self.client.stream( + with active_client.stream( method, url, params={**params, **private_params}, headers=attempt_headers, content=request_body, auth=auth, timeout=30, follow_redirects=False) as upstream: @@ -565,7 +627,7 @@ def _request(self, method: str, path: str, *, params: Mapping | None = None, if k.lower() in PUBLIC_HEADERS} except httpx.HTTPError: self.store.record(self.scope, requests=1, failures=1, credits_used=estimated_credits, - **{service + "_requests": 1}) + **{service + "_requests": 1}, **weight_stats) # An uncertain request may have been charged. Keep its reservations. if attempt + 1 < self.max_attempts: self.sleep(min(2 ** attempt, self.max_wait)) @@ -573,8 +635,20 @@ def _request(self, method: str, path: str, *, params: Mapping | None = None, raise SourceError(f"{self.source}: transport failure") from None except InvalidResponse: self.store.record(self.scope, requests=1, failures=1, credits_used=estimated_credits, - **{service + "_requests": 1}) + **{service + "_requests": 1}, **weight_stats) raise + if lease is not None and 200 <= status < 300 and actual_weight is not None: + try: + used_weight = actual_weight(body) + if type(used_weight) is not int or not 1 <= used_weight <= rate_weight: + raise ValueError + except (ValueError, TypeError, UnicodeDecodeError): + self.store.record(self.scope, requests=1, failures=1, + credits_used=estimated_credits, + **{service + "_requests": 1}, **weight_stats) + raise InvalidResponse(f"{self.source}: invalid response weight") from None + self.ip_pool.settle(self.store, lease, used_weight) + weight_stats["ip_weight_used"] = used_weight try: charged = int(safe_headers.get("x-ratelimit-credits-used", estimated_credits)) if charged < 0: @@ -585,11 +659,16 @@ def _request(self, method: str, path: str, *, params: Mapping | None = None, if byte_reservation is not None: self.store.settle(byte_reservation, wire) self.store.record(self.scope, requests=1, wire_bytes=wire, decoded_bytes=len(body), - credits_used=charged, failures=int(status >= 400), **{service + "_requests": 1}) + credits_used=charged, failures=int(status >= 400), + **{service + "_requests": 1}, **weight_stats) # Persist provider cooldowns even when the caller cannot wait or retry. # Budget rejection takes precedence in the error contract, never in pacing. retry_delay = max(2 ** attempt, self._retry_after(safe_headers.get("retry-after"))) - if status == 429 or 500 <= status <= 599: + if status == 429 and lease is not None: + self.ip_pool.defer(self.store, self.source, lease.route, retry_delay) + if (500 <= status <= 599 or (status == 429 and ( + lease is None or not self.ip_throttle_only + or safe_headers.get("x-ratelimit-remaining") == "0"))): self.store.defer(self.scope + ":all", retry_delay) if observe: observe(safe_headers, request_windows={"day": credit_window, "week": byte_window}) @@ -622,7 +701,9 @@ def _request(self, method: str, path: str, *, params: Mapping | None = None, acquired, snapshot_id, license=self.license, public_request_headers=identity["headers"], request_body=request_body.decode("utf-8") if request_body is not None else None, - completed_at_utc=utc_stamp(self.store.clock()), http_status=status) + completed_at_utc=utc_stamp(self.store.clock()), http_status=status, + egress_ip_sha256=(digest(lease.public_ip.encode()) + if lease and lease.public_ip != "0.0.0.0" else None)) response = Response(body, safe_headers, p, False, charged, wire) if not sensitive: self.store.put(cache_key, response) diff --git a/dedomena/sources/egress.py b/dedomena/sources/egress.py new file mode 100644 index 0000000..8779a56 --- /dev/null +++ b/dedomena/sources/egress.py @@ -0,0 +1,315 @@ +"""Owned egress routes and durable, weighted per-IP admission. + +The declared public IP identifies a quota; it does not configure routing. Each +route therefore needs a local interface, a proxy, or an already routed HTTP +client. No address is discovered automatically. Different local addresses +behind the same NAT must declare the same public IP and cannot form a pool. +""" +from __future__ import annotations + +from dataclasses import dataclass, field +import hashlib +import ipaddress +import math +import time +import uuid +from typing import Callable, Iterable, TYPE_CHECKING + +import httpx + +if TYPE_CHECKING: + from .core import Store + + +def _address(value: str) -> str: + if not isinstance(value, str) or "%" in value: + raise ValueError("route IP must be an IPv4 or IPv6 address") + try: + address = ipaddress.ip_address(value) + except ValueError: + raise ValueError("route IP must be an IPv4 or IPv6 address") from None + if isinstance(address, ipaddress.IPv6Address) and address.ipv4_mapped: + address = address.ipv4_mapped + return str(address) + + +@dataclass(frozen=True) +class IPRoute: + """An explicitly routed connection and its operator-declared public IP. + + Injected clients must already route through the declared address. They + cannot be combined with local_address or proxy, which configure a client + owned by the pool. Proxy credentials are deliberately absent from repr. + """ + public_ip: str + local_address: str | None = None + proxy: str | None = field(default=None, repr=False) + client: httpx.Client | None = field(default=None, repr=False, compare=False) + + def __post_init__(self): + object.__setattr__(self, "public_ip", _address(self.public_ip)) + if sum(value is not None for value in (self.local_address, self.proxy, self.client)) != 1: + raise ValueError("route requires exactly one local address, proxy or routed client") + if self.local_address is not None: + object.__setattr__(self, "local_address", _address(self.local_address)) + if self.proxy is not None: + if not isinstance(self.proxy, str): + raise ValueError("route proxy must be an HTTP or HTTPS URL") + try: + url = httpx.URL(self.proxy) + valid = (url.scheme in ("http", "https") and bool(url.host) + and not url.query and not url.fragment and url.path in ("", "/")) + except (httpx.InvalidURL, ValueError): + valid = False + if not valid: + raise ValueError("route proxy must be an HTTP or HTTPS origin") from None + if self.client is not None and not isinstance(self.client, httpx.Client): + raise ValueError("route client must be an httpx.Client") + + +@dataclass(frozen=True) +class IPLease: + """One admitted request, with a uniquely owned refundable IP estimate.""" + route: IPRoute + reservation_id: str | None + weight: int + scope: str + + @property + def client(self) -> httpx.Client: + return self.route.client + + @property + def public_ip(self) -> str: + return self.route.public_ip + + +class IPPool: + """Share a route pool across sources while keeping each provider's IP quota. + + Admission uses rolling weighted windows in Store, shared by threads, + processes and restarts using that database. Provider/key limits supplied + in limits are admitted in the same transaction; adding routes cannot + multiply those allowances. Quota charges represent dispatched attempts, + including failed requests. Only a verified response may reduce an estimate. + """ + def __init__(self, routes: Iterable[IPRoute]): + routes = tuple(routes) + if not routes or any(not isinstance(route, IPRoute) for route in routes): + raise ValueError("pool requires at least one IPRoute") + if len({route.public_ip for route in routes}) != len(routes): + raise ValueError("pool public IP addresses must be distinct, including NAT aliases") + self._owned = [] + self._closed = False + try: + configured = [] + for route in routes: + if route.client is None: + options = dict(timeout=30, follow_redirects=False, trust_env=False, + limits=httpx.Limits(max_connections=100, max_keepalive_connections=100)) + if route.local_address is not None: + options["transport"] = httpx.HTTPTransport(local_address=route.local_address, + trust_env=False, limits=options["limits"]) + else: + options["proxy"] = route.proxy + client = httpx.Client(**options) + self._owned.append(client) + configured.append(IPRoute(route.public_ip, client=client)) + else: + configured.append(route) + self.routes = tuple(configured) + except BaseException: + self.close() + raise + # The scheduling cursor is shared by pools with the same ordered routes. + self._identity = hashlib.sha256("\0".join(r.public_ip for r in self.routes).encode()).hexdigest() + + def close(self): + if not self._closed: + self._closed = True + for client in self._owned: + client.close() + + def __enter__(self): + if self._closed: + raise ValueError("IP pool is closed") + return self + + def __exit__(self, *args): + self.close() + + @staticmethod + def scope(source: str, route: IPRoute) -> str: + """Stable public-IP quota identity; excludes proxy and client credentials.""" + if not isinstance(source, str) or not source: + raise ValueError("source must be a nonempty string") + return source + ":ip:" + route.public_ip + + @staticmethod + def _schema(db): + db.execute("CREATE TABLE IF NOT EXISTS ip_events(" + "id TEXT PRIMARY KEY, scope TEXT NOT NULL, charged REAL NOT NULL, weight INTEGER NOT NULL)") + db.execute("CREATE INDEX IF NOT EXISTS ip_events_by_scope ON ip_events(scope,charged)") + db.execute("CREATE TABLE IF NOT EXISTS ip_windows(scope TEXT PRIMARY KEY, period REAL NOT NULL)") + db.execute("CREATE TABLE IF NOT EXISTS ip_cursors(scope TEXT PRIMARY KEY, position INTEGER NOT NULL)") + # Store owns cooldowns so admission and existing pacing see the same deferral. + db.execute("CREATE TABLE IF NOT EXISTS cooldowns(scope TEXT PRIMARY KEY, until REAL NOT NULL)") + + @staticmethod + def _limit(scope, weight, capacity, period): + if (not isinstance(scope, str) or not scope + or type(weight) is not int or weight < 1 + or type(capacity) is not int or capacity < 1 + or isinstance(period, bool) or not isinstance(period, (int, float)) + or not math.isfinite(period) or period <= 0): + raise ValueError("weighted admission requires a scope, positive integer cost and capacity, and finite period") + if weight > capacity: + from .core import BudgetExceeded + raise BudgetExceeded("request weight exceeds the configured window allowance") + return scope, weight, capacity, float(period) + + @staticmethod + def _ready(db, limit, now): + scope, weight, capacity, period = limit + row = db.execute("SELECT period FROM ip_windows WHERE scope=?", (scope,)).fetchone() + if row and row[0] != period: + # Changing a scope's period after pruning history could hide consumption. + raise ValueError("quota scope cannot change its rolling window period") + db.execute("INSERT OR IGNORE INTO ip_windows VALUES(?,?)", (scope, period)) + db.execute("DELETE FROM ip_events WHERE scope=? AND charged<=?", (scope, now - period)) + rows = db.execute("SELECT charged,weight FROM ip_events WHERE scope=? ORDER BY charged", + (scope,)).fetchall() + used = sum(amount for _, amount in rows) + due = now + if used + weight > capacity: + for charged, amount in rows: + used -= amount + due = charged + period + if used + weight <= capacity: + break + cooldown = db.execute("SELECT until FROM cooldowns WHERE scope=?", (scope,)).fetchone() + return max(due, cooldown[0] if cooldown else now) + + def acquire(self, store: Store, source: str, weight: int, capacity: int | None, period: float, + sleep: Callable[[float], None] = time.sleep, max_wait: float = 60, + *, limits: Iterable[tuple[str, int, int, float]] = (), + pacing: Iterable[tuple[str, float]] = (), before_admit=None) -> IPLease: + """Admit the next ready IP without reserving future dispatch slots. + + limits contains (scope, weight, capacity, period) for provider/key + ceilings. Scopes must be distinct from the IP scopes and each other. + Reuse a stable scope and period across every worker and pool. Existing + Store pacing scopes are admitted atomically too. capacity=None skips + weighted IP admission while retaining routing and cooldowns. Optional + before_admit(db, now) reserves provider budgets in the same transaction + immediately before admission; an exception rolls back every reservation. + """ + if self._closed: + raise ValueError("IP pool is closed") + if (isinstance(max_wait, bool) or not isinstance(max_wait, (int, float)) + or not math.isfinite(max_wait) or max_wait < 0): + raise ValueError("max_wait must be finite and nonnegative") + common = tuple(self._limit(*limit) for limit in limits) + rate_limits = tuple(pacing) + if any(not isinstance(scope, str) or not scope + or isinstance(rate, bool) or not isinstance(rate, (int, float)) + or not math.isfinite(rate) or rate <= 0 for scope, rate in rate_limits): + raise ValueError("pacing requires a scope and positive finite rate") + if len({scope for scope, _ in rate_limits}) != len(rate_limits): + raise ValueError("pacing scopes must be distinct") + # Validate weight and period even when this provider has no IP quota. + self._limit("validation", weight, capacity if capacity is not None else weight, period) + route_limits = tuple(self._limit(self.scope(source, route), weight, + capacity if capacity is not None else weight, period) + for route in self.routes) + common_scopes = [limit[0] for limit in common] + if (len(set(common_scopes)) != len(common_scopes) + or set(common_scopes).intersection(limit[0] for limit in route_limits)): + raise ValueError("admission scopes must be distinct") + began = store.clock() + scheduling = source + ":pool:" + self._identity + waited = False + while True: + if self._closed: + raise ValueError("IP pool is closed") + with store.transaction() as db: + self._schema(db) + now = store.clock() + db.execute("DELETE FROM throttles WHERE expires<=?", (now,)) + common_due = max([now] + [self._ready(db, limit, now) for limit in common]) + effective_rates = [] + for scope, rate in rate_limits: + throttled = db.execute("SELECT min(rate) FROM throttles WHERE scope=?", (scope,)).fetchone()[0] + effective_rates.append((scope, min(rate, throttled) if throttled else rate)) + next_at = db.execute("SELECT next_at FROM pacing WHERE scope=?", (scope,)).fetchone() + cooldown = db.execute("SELECT until FROM cooldowns WHERE scope=?", (scope,)).fetchone() + common_due = max(common_due, next_at[0] if next_at else now, + cooldown[0] if cooldown else now) + due = [] + for limit in route_limits: + if capacity is not None: + route_due = self._ready(db, limit, now) + else: + cooldown = db.execute("SELECT until FROM cooldowns WHERE scope=?", (limit[0],)).fetchone() + route_due = cooldown[0] if cooldown else now + due.append(max(common_due, route_due)) + row = db.execute("SELECT position FROM ip_cursors WHERE scope=?", (scheduling,)).fetchone() + first = row[0] % len(self.routes) if row else 0 + ordering = [(first + offset) % len(self.routes) for offset in range(len(self.routes))] + chosen = min(ordering, key=lambda index: due[index]) + delay = max(0, due[chosen] - now) + remaining = max_wait - max(0, now - began) + if delay > max(0, remaining) or (remaining < 0 and (waited or max_wait > 0)): + from .core import Throttled + raise Throttled(source, delay) + if delay == 0: + if before_admit is not None: + before_admit(db, now) + now = store.clock() + if max_wait > 0 and now - began > max_wait: + from .core import Throttled + raise Throttled(source, 0) + reservation = uuid.uuid4().hex if capacity is not None else None + for scope, amount, _, _ in common: + db.execute("INSERT INTO ip_events VALUES(?,?,?,?)", (uuid.uuid4().hex, scope, now, amount)) + if reservation: + db.execute("INSERT INTO ip_events VALUES(?,?,?,?)", + (reservation, route_limits[chosen][0], now, weight)) + for scope, rate in effective_rates: + db.execute("INSERT OR REPLACE INTO pacing VALUES(?,?)", (scope, now + 1 / rate)) + db.execute("INSERT OR REPLACE INTO ip_cursors VALUES(?,?)", + (scheduling, (chosen + 1) % len(self.routes))) + return IPLease(self.routes[chosen], reservation, weight, route_limits[chosen][0]) + waited = True + sleep(delay) + + @staticmethod + def settle(store: Store, lease: IPLease, actual_weight: int): + """Reduce only this request's known estimate; expired charges stay expired. + + Settle after a trustworthy response has established actual cost. An + uncertain send or malformed response must retain its full reservation. + """ + if not isinstance(lease, IPLease) or type(actual_weight) is not int or not 0 <= actual_weight <= lease.weight: + raise ValueError("actual weight must be a nonnegative integer within the reserved estimate") + if lease.reservation_id is None: + return + with store.transaction() as db: + IPPool._schema(db) + row = db.execute("SELECT weight FROM ip_events WHERE id=? AND scope=?", + (lease.reservation_id, lease.scope)).fetchone() + if row and row[0] != lease.weight and actual_weight != row[0]: + raise ValueError("settlement cannot revise an already reduced reservation") + db.execute("UPDATE ip_events SET weight=? WHERE id=? AND scope=?", + (actual_weight, lease.reservation_id, lease.scope)) + + def defer(self, store: Store, source: str, route: IPRoute | IPLease, seconds: float): + """Persist a 429 cooldown only for this provider and public IP.""" + if (isinstance(seconds, bool) or not isinstance(seconds, (int, float)) + or not math.isfinite(seconds) or seconds < 0): + raise ValueError("cooldown must be finite and nonnegative") + if isinstance(route, IPLease): + route = route.route + if route not in self.routes: + raise ValueError("cooldown route must belong to this pool") + store.defer(self.scope(source, route), seconds) diff --git a/dedomena/sources/hyperliquid.py b/dedomena/sources/hyperliquid.py new file mode 100644 index 0000000..4cc3d5f --- /dev/null +++ b/dedomena/sources/hyperliquid.py @@ -0,0 +1,332 @@ +"""Hyperliquid public market data with native precision and replayable receipts.""" +from __future__ import annotations + +from decimal import Decimal, InvalidOperation +import json +import math +import re +from typing import Callable, Iterator + +from .core import InvalidResponse, Page, SearchLimitExceeded, Transport, canonical, digest + + +REQUEST_TYPES = frozenset(("allMids", "metaAndAssetCtxs", "spotMetaAndAssetCtxs", "l2Book", + "candleSnapshot", "fundingHistory")) + + +def info_policy(content: bytes) -> tuple[int, Callable[[bytes], int] | None]: + """Validate the fixed read-only info schema and derive its dispatch weight. + + Transport invokes this policy even for direct _request callers. Unsupported + actions and unexpected fields fail before cache lookup or network dispatch. + Response-sized weights reserve a conservative bound and settle by item count. + """ + try: + body = json.loads(content) + except (ValueError, TypeError, UnicodeDecodeError): + raise ValueError("Hyperliquid info body must be JSON") from None + if (not isinstance(body, dict) or not isinstance(body.get("type"), str) + or body["type"] not in REQUEST_TYPES): + raise ValueError("Hyperliquid request type is outside the read-only adapter") + kind = body["type"] + required = { + "allMids": {"type"}, "metaAndAssetCtxs": {"type"}, "spotMetaAndAssetCtxs": {"type"}, + "l2Book": {"type", "coin"}, + "candleSnapshot": {"type", "req"}, "fundingHistory": {"type", "coin", "startTime", "endTime"}, + }[kind] + optional = {"allMids": {"dex"}, "metaAndAssetCtxs": {"dex"}, + "l2Book": {"nSigFigs", "mantissa"}}.get(kind, set()) + if not required <= set(body) <= required | optional: + raise ValueError("Hyperliquid info body has unsupported or missing fields") + if "dex" in body: + Hyperliquid._dex(body["dex"]) + if "coin" in body: + Hyperliquid._coin(body["coin"]) + if kind == "l2Book": + Hyperliquid._aggregation(body.get("nSigFigs"), body.get("mantissa")) + if kind == "fundingHistory": + Hyperliquid._range(body["startTime"], body["endTime"]) + if kind == "candleSnapshot": + request = body["req"] + if not isinstance(request, dict) or set(request) != {"coin", "interval", "startTime", "endTime"}: + raise ValueError("Hyperliquid candle request has invalid fields") + Hyperliquid._coin(request["coin"]) + Hyperliquid._interval(request["interval"]) + Hyperliquid._range(request["startTime"], request["endTime"]) + if kind in ("allMids", "l2Book"): + return 2, None + if kind in ("metaAndAssetCtxs", "spotMetaAndAssetCtxs"): + return 20, None + maximum, per_unit = {"fundingHistory": (500, 20), + "candleSnapshot": (5000, 60)}[kind] + def actual_weight(content: bytes) -> int: + rows = json.loads(content) + if (not isinstance(rows, list) or len(rows) > maximum + or any(not isinstance(row, dict) for row in rows)): + raise ValueError("invalid Hyperliquid item-count response") + return 20 + math.ceil(len(rows) / per_unit) + return 20 + math.ceil(maximum / per_unit), actual_weight + + +class Hyperliquid(Transport): + """Restricted public /info queries; no account, signing, or exchange actions. + + Snapshot records preserve the native response envelope. Historical funding + pages remove exact inclusive-boundary repeats; candles retain the provider's + recent-history limitation rather than claiming complete historical coverage. + """ + + MAX_FUNDING_PAGE = 500 + MAX_CANDLES = 5000 + INTERVALS = frozenset(("1m", "3m", "5m", "15m", "30m", "1h", "2h", "4h", + "8h", "12h", "1d", "3d", "1w", "1M")) + + def __init__(self, **kwargs): + kwargs.setdefault("requests_per_second", 1000) + kwargs.setdefault("ip_limit", (1200, 60.0)) + limit = kwargs["ip_limit"] + if (not isinstance(limit, tuple) or len(limit) != 2 or type(limit[0]) is not int + or not 0 < limit[0] <= 1200 or isinstance(limit[1], bool) + or not isinstance(limit[1], (int, float)) or not math.isfinite(limit[1]) or limit[1] < 60): + raise ValueError("Hyperliquid IP policy cannot exceed 1,200 weight per minute") + kwargs["ip_throttle_only"] = True + kwargs.setdefault("cache_ttl", 0) + super().__init__("hyperliquid", "https://api.hyperliquid.xyz", + license="Hyperliquid public API; provider terms apply", **kwargs) + + @staticmethod + def _coin(coin: str) -> str: + if (not isinstance(coin, str) or len(coin) > 128 or not re.fullmatch( + r"(?:[A-Za-z0-9_]+:)?[A-Za-z0-9_@#]+(?:/[A-Za-z0-9_]+)?", coin)): + raise ValueError("coin must be a native Hyperliquid asset name") + return coin + + @staticmethod + def _dex(dex: str) -> str: + if not isinstance(dex, str) or not re.fullmatch(r"[A-Za-z0-9_-]{0,64}", dex): + raise ValueError("dex must be a Hyperliquid perpetual dex name") + return dex + + @staticmethod + def _time(value: int, name: str) -> int: + if type(value) is not int or not 0 <= value <= 9_007_199_254_740_991: + raise ValueError(f"{name} must be a nonnegative integer timestamp in milliseconds") + return value + + @classmethod + def _range(cls, start_time: int, end_time: int) -> tuple[int, int]: + start, end = cls._time(start_time, "start_time"), cls._time(end_time, "end_time") + if start > end: + raise ValueError("start_time must not follow end_time") + return start, end + + @classmethod + def _interval(cls, interval: str) -> str: + if not isinstance(interval, str) or interval not in cls.INTERVALS: + raise ValueError("interval must be a supported Hyperliquid candle interval") + return interval + + @staticmethod + def _aggregation(n_sig_figs, mantissa): + if n_sig_figs is not None and (type(n_sig_figs) is not int or n_sig_figs not in (2, 3, 4, 5)): + raise ValueError("n_sig_figs must be 2, 3, 4, 5 or None") + if mantissa is not None and (type(mantissa) is not int or mantissa not in (1, 2, 5) or n_sig_figs != 5): + raise ValueError("mantissa must be 1, 2 or 5 and requires n_sig_figs=5") + + @staticmethod + def _decimal(value, *, nonnegative: bool = False) -> bool: + if not isinstance(value, str) or not value or len(value) > 256: + return False + try: + number = Decimal(value) + return number.is_finite() and (not nonnegative or number >= 0) + except InvalidOperation: + return False + + def _info(self, body: dict, *, refresh: bool): + content = canonical(body).encode() + weight, actual = info_policy(content) + response = self._request("POST", "/info", content=content, + headers={"Content-Type": "application/json"}, + rate_weight=weight, actual_weight=actual, refresh=refresh) + try: + data = json.loads(response.body) + except (ValueError, UnicodeDecodeError): + raise InvalidResponse("hyperliquid: invalid JSON response") from None + return response, data + + def all_mids(self, *, dex: str = "", refresh: bool = False) -> Page: + """One native coin-to-price map; spot mids belong to the default dex.""" + dex = self._dex(dex) + response, data = self._info({"type": "allMids", "dex": dex}, refresh=refresh) + if (not isinstance(data, dict) or any(not isinstance(coin, str) or not coin + or not self._decimal(price, nonnegative=True) for coin, price in data.items())): + raise InvalidResponse("hyperliquid: invalid mids response") + return self.page(response, [data], total=1, complete=True, records_seen=1, + warnings=("Point-in-time mids; an empty book can fall back to the last trade price.",)) + + def markets(self, *, spot: bool = False, dex: str = "", refresh: bool = False) -> Page: + """One native metadata/context envelope; spot contexts can outnumber pairs.""" + dex = self._dex(dex) + if type(spot) is not bool or (spot and dex): + raise ValueError("spot must be boolean; spot metadata does not accept a dex") + body = {"type": "spotMetaAndAssetCtxs" if spot else "metaAndAssetCtxs"} + if not spot: + body["dex"] = dex + response, data = self._info(body, refresh=refresh) + if (not isinstance(data, list) or len(data) != 2 or not isinstance(data[0], dict) + or not isinstance(data[1], list) or not isinstance(data[0].get("universe"), list)): + raise InvalidResponse("hyperliquid: invalid market metadata envelope") + meta, contexts = data + universe = meta["universe"] + if ((not spot and len(universe) != len(contexts)) or any(not isinstance(row, dict) + or not isinstance(row.get("name"), str) or not row["name"] for row in universe) + or len({row["name"] for row in universe}) != len(universe) + or any(not isinstance(row, dict) for row in contexts) + or (spot and (not isinstance(meta.get("tokens"), list) + or any(not isinstance(row, dict) for row in meta["tokens"])))): + raise InvalidResponse("hyperliquid: inconsistent market metadata or asset contexts") + return self.page(response, [{"meta": meta, "assetCtxs": contexts}], total=1, + complete=True, records_seen=1, + warnings=(("Spot contexts and pair metadata have independent lengths; native coin fields and token indices are retained." + if spot else "Perpetual asset contexts match universe positions; native decimal strings are retained."),)) + + def order_book(self, coin: str, *, n_sig_figs: int | None = None, + mantissa: int | None = None, refresh: bool = False) -> Page: + """Native L2 snapshot, at most twenty levels on each side.""" + coin = self._coin(coin) + self._aggregation(n_sig_figs, mantissa) + body = {"type": "l2Book", "coin": coin} + if n_sig_figs is not None: + body["nSigFigs"] = n_sig_figs + if mantissa is not None: + body["mantissa"] = mantissa + response, data = self._info(body, refresh=refresh) + if (not isinstance(data, dict) or data.get("coin") != coin + or type(data.get("time")) is not int or data["time"] < 0 + or not isinstance(data.get("levels"), list) or len(data["levels"]) != 2 + or any(not isinstance(side, list) or len(side) > 20 for side in data["levels"])): + raise InvalidResponse("hyperliquid: invalid book identity or levels") + for side in data["levels"]: + if any(not isinstance(row, dict) or not self._decimal(row.get("px"), nonnegative=True) + or not self._decimal(row.get("sz"), nonnegative=True) + or type(row.get("n")) is not int or row["n"] < 0 for row in side): + raise InvalidResponse("hyperliquid: invalid book level") + return self.page(response, [data], total=1, complete=True, records_seen=1, + warnings=("Point-in-time L2 snapshot; the provider returns at most 20 levels per side.",)) + + def candles(self, coin: str, interval: str, start_time: int, end_time: int, + *, refresh: bool = False) -> Page: + """Return available candles; only the provider's latest 5,000 exist here.""" + coin, interval = self._coin(coin), self._interval(interval) + start, end = self._range(start_time, end_time) + response, rows = self._info({"type": "candleSnapshot", "req": { + "coin": coin, "interval": interval, "startTime": start, "endTime": end}}, refresh=refresh) + if not isinstance(rows, list) or len(rows) > self.MAX_CANDLES: + raise InvalidResponse("hyperliquid: invalid candle response") + previous = -1 + for row in rows: + if (not isinstance(row, dict) or row.get("s") != coin or row.get("i") != interval + or type(row.get("t")) is not int or type(row.get("T")) is not int + or row["t"] < 0 or row["T"] < row["t"] or row["t"] > end or row["T"] < start + or row["t"] <= previous or type(row.get("n")) is not int or row["n"] < 0 + or any(not self._decimal(row.get(key), nonnegative=True) for key in ("o", "h", "l", "c", "v"))): + raise InvalidResponse("hyperliquid: invalid candle identity, range or values") + previous = row["t"] + return self.page(response, rows, complete=False, records_seen=len(rows), + warnings=("Only the latest 5,000 candles are retained; a short or empty response does not prove historical coverage.", + "The current interval can be unfinished; OHLCV values remain native decimal strings.")) + + def funding_history(self, coin: str, start_time: int, end_time: int | None = None, + *, cursor: str | None = None, refresh: bool = False) -> Iterator[Page]: + """Enumerate a pinned inclusive range with safe timestamp overlap. + + Resume with the same coin/start/end and the preceding Page.next_cursor. + If end_time is omitted, the cursor retains the initial retrieval cutoff. + The opaque cursor contains the count and hashes of already emitted rows + at the inclusive boundary. A saturated timestamp cannot be enumerated + safely and raises SearchLimitExceeded before that page is delivered. + """ + coin = self._coin(coin) + start = self._time(start_time, "start_time") + if end_time is not None: + self._time(end_time, "end_time") + current, seen, boundary = start, 0, set() + if cursor is not None: + if not isinstance(cursor, str) or len(cursor) > 50_000: + raise ValueError("cursor must be a funding checkpoint for this exact query") + try: + state = json.loads(cursor) + if (not isinstance(state, dict) or set(state) != {"coin", "start", "end", "after", "seen", "boundary"} + or state["coin"] != coin or state["start"] != start + or (end_time is not None and state["end"] != end_time) + or type(state["seen"]) is not int or state["seen"] <= 0 + or not isinstance(state["boundary"], list) or not state["boundary"] + or len(state["boundary"]) > self.MAX_FUNDING_PAGE + or any(not isinstance(value, str) or not re.fullmatch(r"[0-9a-f]{64}", value) + for value in state["boundary"])): + raise ValueError + end = self._time(state["end"], "cursor end") + current = self._time(state["after"], "cursor after") + if not start < current <= end: + raise ValueError + seen, boundary = state["seen"], set(state["boundary"]) + except (ValueError, TypeError, KeyError): + raise ValueError("cursor must be a funding checkpoint for this exact query") from None + else: + end = self._time(end_time if end_time is not None else int(self.store.clock() * 1000), "end_time") + self._range(start, end) + while True: + response, rows = self._info({"type": "fundingHistory", "coin": coin, + "startTime": current, "endTime": end}, refresh=refresh) + if not isinstance(rows, list) or len(rows) > self.MAX_FUNDING_PAGE: + raise InvalidResponse("hyperliquid: invalid funding history response") + previous, hashes = current, set() + emitted = [] + for row in rows: + if (not isinstance(row, dict) or row.get("coin") != coin + or type(row.get("time")) is not int or not current <= row["time"] <= end + or row["time"] < previous or not self._decimal(row.get("fundingRate")) + or not self._decimal(row.get("premium"))): + raise InvalidResponse("hyperliquid: invalid funding identity, order or range") + ident = digest(canonical(row).encode()) + if ident in hashes: + raise InvalidResponse("hyperliquid: duplicate funding row in response") + hashes.add(ident) + previous = row["time"] + if not (row["time"] == current and ident in boundary): + emitted.append(row) + more = len(rows) == self.MAX_FUNDING_PAGE + if more and previous == current: + raise SearchLimitExceeded("hyperliquid: funding timestamp reaches the page cap; narrow the range") + seen += len(emitted) + following = None + if more: + boundary = {digest(canonical(row).encode()) for row in rows if row["time"] == previous} + following = canonical({"coin": coin, "start": start, "end": end, + "after": previous, "seen": seen, "boundary": sorted(boundary)}) + if len(rows) > len(emitted): + self.store.record(self.scope, upstream_records=0 if response.cache_hit else len(rows) - len(emitted), + boundary_repeats=len(rows) - len(emitted)) + yield self.page(response, emitted, complete=not more, next_cursor=following, + records_seen=seen, + warnings=("Inclusive funding timestamp overlap is deduplicated by exact row hash; this is provider-available history.",)) + if not more: + return + current = previous + + def fetch(self, identifier: str, *, refresh: bool = False) -> Page: + return self.order_book(identifier, refresh=refresh) + + def search(self, query: str, *, refresh: bool = False) -> Iterator[Page]: + yield self.fetch(query, refresh=refresh) + + def quota(self) -> dict: + return {"api_key_required": False, "daily_provider_quota": None, + "ip_weight_per_minute": 1200, "ip_weight_per_second_average": 20, + "local_requests_per_second": self.rate, + "reserved_request_weights": {"all_mids": 2, "order_book": 2, "markets": 20, + "funding_history": 45, "candles": 104}, + "max_candles_retained": self.MAX_CANDLES, + "funding_page_cap": self.MAX_FUNDING_PAGE, + "note": "Response-sized requests reserve maximum weight, then use rounded-up returned counts for local accounting; provider billing is not reported. IP pools do not expand address limits."} diff --git a/dedomena/sources/openalex.py b/dedomena/sources/openalex.py index 08d7b28..88ff4c5 100644 --- a/dedomena/sources/openalex.py +++ b/dedomena/sources/openalex.py @@ -45,6 +45,11 @@ def __init__(self, api_key: str | None = None, *, daily_credit_limit: int = 10_0 kwargs.setdefault("requests_per_second", 100) if kwargs["requests_per_second"] > 100: raise ValueError("OpenAlex permits at most 100 requests per second") + ip_limit = kwargs.setdefault("ip_limit", (100, 1.0)) + if (not isinstance(ip_limit, tuple) or len(ip_limit) != 2 + or type(ip_limit[0]) is not int or not 1 <= ip_limit[0] <= 100 + or type(ip_limit[1]) not in (int, float) or ip_limit[1] != 1): + raise ValueError("OpenAlex IP policy permits at most 100 requests per second") super().__init__("openalex", "https://api.openalex.org", credential=self.api_key, license="CC0 metadata; article text has separate rights", **kwargs) diff --git a/docs/FINANCE.md b/docs/FINANCE.md index 360b374..be97d03 100644 --- a/docs/FINANCE.md +++ b/docs/FINANCE.md @@ -3,7 +3,8 @@ Four official APIs cover corporate fundamentals, dated macroeconomic data, monetary statistics, reference exchange rates and international indicators. Every adapter shares the research-source transport, native records, durable -snapshots, verified replay and JSON CLI. +snapshots, verified replay and JSON CLI. [Hyperliquid](HYPERLIQUID.md) adds public +crypto market data with [weighted IP routing](IP_ROUTING.md). ## Access and retrieval weight diff --git a/docs/HYPERLIQUID.md b/docs/HYPERLIQUID.md new file mode 100644 index 0000000..5d6841b --- /dev/null +++ b/docs/HYPERLIQUID.md @@ -0,0 +1,123 @@ +# Hyperliquid market data + +`Hyperliquid` retrieves public market data from the official +[POST `/info` endpoint](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/info-endpoint). +No key or wallet is required. The transport allows only the supported public +query schemas; account queries, signing and `/exchange` are excluded. + +## Choose the retrieval unit + +| SDK method | Native result | Request weight | +| --- | --- | --- | +| `all_mids(dex="")` | One coin-to-price map | 2 | +| `markets(spot=False, dex="")` | One metadata/context envelope | 20 | +| `order_book(coin)`; `fetch(coin)` alias | One L2 snapshot, at most 20 levels/side | 2 | +| `candles(coin, interval, start_time, end_time)` | Available OHLCV rows | Reserve 104; settle 20 + ceil(rows/60) | +| `funding_history(coin, start_time, end_time=None, cursor=None)` | Funding rows in timestamp pages | Reserve 45; settle 20 + ceil(rows/20) | + +Metadata retains native universe, token indices, margin tables and asset +contexts. Perpetual contexts follow universe positions. Spot contexts and pair +metadata can have different lengths; both remain intact without an inferred join. +Decimals remain strings; units, nulls and additional native fields are unchanged. +Native schemas: [perpetuals](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/info-endpoint/perpetuals), +[spot](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/info-endpoint/spot). + +Use provider asset names: `BTC` for a perpetual, `xyz:XYZ100` for a builder DEX, +`PURR/USDC` or `@` for spot. Resolve spot indices from the native universe; +UI labels can differ from API names. `dex` selects a perpetual DEX for mids or +metadata; spot metadata does not accept it. +[Asset naming](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/info-endpoint#perpetuals-vs-spot). + +~~~python +import time +from dedomena.sources import Hyperliquid + +end = int(time.time() * 1000) +start = end - 86_400_000 +with Hyperliquid() as api: + mids = api.all_mids() + perpetuals = api.markets() + spot = api.markets(spot=True) + book = api.order_book("BTC") + candles = api.candles("BTC", "1h", start, end) + for page in api.funding_history("BTC", start, end): + consume(page.records, page.provenance.to_dict()) +~~~ + +Times are inclusive epoch milliseconds. Candle intervals are `1m`, `3m`, `5m`, +`15m`, `30m`, `1h`, `2h`, `4h`, `8h`, `12h`, `1d`, `3d`, `1w`, `1M`. +Only the latest 5,000 candles are available. Every candle Page has +`complete=False`: short or empty results cannot prove historical coverage. +The current candle can be unfinished. [Candle contract](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/info-endpoint#candle-snapshot). + +## Funding checkpoints and receipts + +Funding uses the provider's 500-item time-range ceiling. Each following request +starts at the preceding page's final timestamp; exact boundary-row hashes remove +inclusive repeats. A full page that cannot advance past its starting timestamp +raises `SearchLimitExceeded` before delivering that page. Malformed identities, +reordered timestamps and duplicate rows fail explicitly. +[Provider pagination](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/info-endpoint#pagination). + +The opaque `next_cursor` contains the exact coin, start/end bounds, emitted +count and boundary hashes. Resume with the same coin/start and cursor; supply +the same end or omit it to use the cutoff stored in the checkpoint. An omitted +initial end is pinned once to retrieval time. `complete=True` finishes enumeration +of provider-available funding for that range. + +~~~python +with Hyperliquid() as api: + start = end - 30 * 86_400_000 + stream = api.funding_history("BTC", start, end) + checkpoint = next(stream) + consume(checkpoint.records, checkpoint.provenance.to_dict()) + stream.close() + if checkpoint.next_cursor is not None: + for page in api.funding_history("BTC", start, cursor=checkpoint.next_cursor): + consume(page.records, page.provenance.to_dict()) +~~~ + +Every Page carries the canonical public JSON request body, hashes, retrieval +and completion times, HTTP status and durable snapshot ID. Replay verifies the +stored response hash. Mutable snapshots default to `cache_ttl=0`; explicit cache +TTL or `refresh=True` follows the [shared source contract](SOURCES.md). + +## CLI + +~~~sh +python -m dedomena.sources markets hyperliquid +python -m dedomena.sources markets hyperliquid --spot +python -m dedomena.sources mids hyperliquid --dex xyz +python -m dedomena.sources book hyperliquid BTC +python -m dedomena.sources fetch hyperliquid BTC +python -m dedomena.sources candles hyperliquid BTC --interval 1h --start-time 1790812800000 --end-time 1790899200000 +python -m dedomena.sources funding hyperliquid BTC --start-time 1790812800000 --end-time 1790899200000 +python -m dedomena.sources quota hyperliquid +~~~ + +The CLI emits JSON Pages followed by retrieval totals and checkpoint fields. +`--max-pages` bounds funding acquisition; checkpoint continuation uses the SDK. + +## IP capacity and measured validation + +Hyperliquid publishes an aggregate REST limit of **1,200 weight/minute per IP**. +The shared limiter reserves each request's conservative maximum, then settles +valid list responses by returned item count. Failed or uncertain sends retain +the reservation; 429 cooldowns apply to the affected IP. The local 1,000 +requests/s dispatch ceiling is not a provider allowance. Other clients outside +the shared store still consume the provider's quota. +[Official rate limits](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/rate-limits-and-user-limits). + +Use the same [generic IPPool](IP_ROUTING.md) with Hyperliquid and OpenAlex. +Ten distinct owned egress IPs have a theoretical ceiling of 6,000 weight-2 +requests/minute, subject to other traffic, latency and provider availability. +This is capacity arithmetic, not a sustained throughput measurement or daily +entitlement. + +Bounded live adapter probes on **2026-10-02**: five successful queries covered +mids, perpetual metadata, spot metadata, candles and funding. All returned HTTP +200; all five snapshots passed verified replay. Combined transfer was 399,814 +decoded bytes, with 191 weight reserved and 85 locally accounted weight. +Item-count costs are rounded up; provider billing was not reported. Candles +returned three rows with `complete=False`; funding returned 24 rows with +`complete=True`. No sustained high-volume or multi-IP benchmark was performed. diff --git a/docs/IP_ROUTING.md b/docs/IP_ROUTING.md new file mode 100644 index 0000000..c92d44f --- /dev/null +++ b/docs/IP_ROUTING.md @@ -0,0 +1,121 @@ +# Shared IP routing + +`IPPool` connects any source adapter to the same owned egress routes. Each route +pairs its **actual public IP** with a local interface address, an HTTP(S) proxy, +or an already routed `httpx.Client`. Declaring an IP identifies its allowance; +it does not configure or verify network routing. Dedomena never discovers IPs +through an external service. + +Two local addresses behind the same NAT share one public IP and one allowance. +Duplicate public IPs, including IPv4-mapped IPv6 aliases, are rejected within a +pool. Configure every worker with the same declared IP for the same egress and +use the same SQLite `Store` to coordinate across threads, processes and restarts. +Other applications using that IP still consume the provider's allowance. + +## One pool, multiple sources + +Replace all example addresses with interfaces and public IPs you own. The +`203.0.113.*` addresses below are reserved documentation examples; local addresses +must exist on the machine running Dedomena. + +~~~python +from dedomena.sources import Hyperliquid, IPPool, IPRoute, OpenAlex, Store + +store = Store("sources.sqlite3") +try: + with IPPool([ + IPRoute("203.0.113.10", local_address="10.0.0.10"), + IPRoute("203.0.113.11", local_address="10.0.0.11"), + ]) as routes: + with OpenAlex(store=store, ip_pool=routes) as papers, \ + Hyperliquid(store=store, ip_pool=routes) as markets: + paper = papers.fetch("W2741809807") + book = markets.order_book("BTC") +finally: + store.close() +~~~ + +`OPENALEX_API_KEY` supplies the paper source's key. Hyperliquid public market +reads require no key. Closing a source leaves the caller's pool and store open; +closing the pool closes only clients it created. Local bindings and proxies use +pooled connections with `trust_env=False`, so environment proxy settings cannot +replace the configured route. An injected client must already implement that +route and cannot be combined with `local_address` or `proxy`. + +For a proxy, keep its URL and credentials in the caller's environment: + +~~~python +import os + +route = IPRoute("203.0.113.10", proxy=os.environ["RESEARCH_PROXY_1"]) +# RESEARCH_PROXY_1 is an http:// or https:// proxy origin, optionally with auth. +~~~ + +## Allowances and retrieval weight + +The pool admits weighted rolling windows, shared provider/key pacing and +cooldowns atomically. It chooses a ready route in deterministic round-robin +order, or waits for the earliest eligible route within `max_wait`. Waiting does +not reserve a future dispatch slot. A positive `max_wait` bounds admission, +including budget work and lock wait; zero disables scheduled quota sleeps. A rolling window's duration remains fixed +for its quota scope so changing a worker's configuration cannot hide prior use. + +- **Hyperliquid:** 1,200 aggregate REST weight per minute per public IP. + Mid-price and order-book requests cost 2; market metadata costs 20. + Ten distinct IPs therefore have a theoretical combined ceiling of 6,000 + cheap requests/minute, versus 600 for one IP. This is allowance arithmetic, + subject to RTT, local pacing and other traffic; it is not a throughput benchmark. + [Official rate policy](https://hyperliquid.gitbook.io/hyperliquid-docs/for-developers/api/rate-limits-and-user-limits). +- **OpenAlex:** Dedomena retains the 100 requests/second ceiling across the key, + plus a conservative IP guard and the shared daily credit budget. The current + help page does not establish an IP-only scope for that rate, so adding IPs + does not increase the key's request or credit allowance. + [Official authentication and limits](https://help.openalex.org/api/authentication/). + +Response-sized Hyperliquid reads reserve their maximum estimated weight before +sending, then reduce only that request's reservation after a verified response +establishes its returned item count. Uncertain sends and failed responses retain +the reservation. Budget rejection before a send consumes no IP allowance. +A 429 persists a cooldown on the affected provider/IP; shared provider or key +exhaustion still applies across the whole pool. Provider errors and retries never +supply extra allowances. + +Without explicit routes, adapters with IP limits use the conservative unknown-IP +scope `0.0.0.0`. All such instances of that provider share its IP guard in the +same store. No extra IP allowance is inferred from separate clients or keys. + +For another adapter, pass `ip_pool=routes` and `ip_limit=(capacity, seconds)` through +its transport configuration. The provider's own key, daily, service and global +limits remain in effect. `ip_throttle_only=True` is appropriate only when that +provider's 429 is known to describe an IP allowance; otherwise the transport also +retains its shared cooldown. A pool with `ip_limit=None` provides routing and IP +cooldowns while keeping the adapter's existing shared pacing. + +Request hashes and cache identity exclude routing, so a query can reuse the same +snapshot across IPs. Network receipts carry optional `egress_ip_sha256` for the +chosen declared public IP; unknown-IP and older receipts leave it unset. Public +provenance, pool representations and admission state exclude proxy credentials. + +## CLI route file + +Save the same routes as a JSON array in `ip-routes.json`: + +~~~json +[ + {"public_ip": "203.0.113.10", "local_address": "10.0.0.10"}, + {"public_ip": "203.0.113.11", "local_address": "10.0.0.11"} +] +~~~ + +~~~sh +python -m dedomena.sources fetch hyperliquid BTC --store sources.sqlite3 --ip-routes ip-routes.json +python -m dedomena.sources fetch openalex W2741809807 --store sources.sqlite3 --ip-routes ip-routes.json +~~~ + +Each JSON route accepts `public_ip` and exactly one of `local_address` or `proxy`; +Python client objects cannot appear in the file. Use the same route identities +and store path across CLI workers. + +Routing admission does not schedule parallel work. A generic bounded executor, +with stable ordering, shared admission and durable checkpoints, is tracked in +[issue #18](https://github.com/autonomio/dedomena/issues/18). diff --git a/setup.py b/setup.py index 24c26c7..3e5f0dc 100755 --- a/setup.py +++ b/setup.py @@ -5,8 +5,9 @@ DESCRIPTION = "Reproducible research data sources for agents and scientists" LONG_DESCRIPTION = """\ Dedomena provides compact, high-throughput access to research and finance data: -OpenAlex, Europe PMC, EPO, SEC EDGAR, FRED/ALFRED, ECB and World Bank. +OpenAlex, Europe PMC, EPO, SEC EDGAR, FRED/ALFRED, ECB, World Bank and Hyperliquid. Every source shares consistent provenance, durable caching and replay. +Reusable IP pools coordinate weighted egress limits across sources. Legacy dataset and API functions remain available through the legacy extra. """ @@ -16,7 +17,7 @@ URL = 'http://autonom.io' LICENSE = 'MIT' DOWNLOAD_URL = 'https://github.com/autonomio/dedomena/' -VERSION = '0.5.0' +VERSION = '0.6.0' try: from setuptools import setup diff --git a/tests/test_egress.py b/tests/test_egress.py new file mode 100644 index 0000000..17148e4 --- /dev/null +++ b/tests/test_egress.py @@ -0,0 +1,362 @@ +from concurrent.futures import ThreadPoolExecutor +import math + +import httpx +import pytest + +from dedomena.sources.core import BudgetExceeded, Store, Throttled +from dedomena.sources.egress import IPLease, IPPool, IPRoute + + +class Clock: + def __init__(self): + self.now = 1000.0 + self.waits = [] + + def __call__(self): + return self.now + + def sleep(self, seconds): + self.waits.append(seconds) + self.now += seconds + + +def pool(*ips): + return IPPool(IPRoute(ip, client=httpx.Client(transport=httpx.MockTransport( + lambda _: httpx.Response(200)))) for ip in ips) + + +def acquire(routes, store, **kwargs): + return routes.acquire(store, "provider", 1, 2, 60, max_wait=0, **kwargs) + + +def test_canonical_nat_aliases_do_not_multiply_quota(): + with pool("2001:db8::1") as routes: + assert routes.routes[0].public_ip == "2001:db8::1" + client = httpx.Client(transport=httpx.MockTransport(lambda _: httpx.Response(200))) + assert IPRoute("::ffff:192.0.2.1", client=client).public_ip == "192.0.2.1" + with pytest.raises(ValueError, match="distinct"): + IPPool([IPRoute("192.0.2.1", client=client), IPRoute("::ffff:c000:201", client=client)]) + with pytest.raises(ValueError, match="distinct"): + IPPool([IPRoute("2001:db8::1", client=client), IPRoute("2001:0DB8:0:0:0:0:0:1", client=client)]) + client.close() + + +@pytest.mark.parametrize("ip", ["host.example", "", "192.0.2.999", "fe80::1%en0", None]) +def test_invalid_ip_is_not_disclosed_in_error(ip): + with pytest.raises(ValueError, match="IPv4 or IPv6"): + IPRoute(ip, local_address="127.0.0.1") + + +def test_routing_is_real_and_cannot_be_overridden_by_injected_client(monkeypatch): + clients, transports = [], [] + original_client, original_transport = httpx.Client, httpx.HTTPTransport + + def transport(**kwargs): + transports.append(kwargs) + return original_transport(**kwargs) + + # IPRoute validates injected clients against the real class. + route = IPRoute("192.0.2.1", local_address="127.0.0.1") + monkeypatch.setattr(httpx, "HTTPTransport", transport) + # The constructor creates a bound client, then validates its immutable route. + # Restore the class while constructing the post-bind route. + def make_client(**kwargs): + clients.append(kwargs) + result = original_client(**kwargs) + monkeypatch.setattr(httpx, "Client", original_client) + return result + monkeypatch.setattr(httpx, "Client", make_client) + routes = IPPool([route]) + owned = routes.routes[0].client + assert transports[0]["local_address"] == "127.0.0.1" + assert transports[0]["trust_env"] is False + assert transports[0]["limits"].max_keepalive_connections == 100 + assert clients[0]["trust_env"] is False + assert clients[0]["follow_redirects"] is False + routes.close() + assert owned.is_closed + with pytest.raises(ValueError, match="exactly one"): + IPRoute("192.0.2.1", local_address="127.0.0.1", client=original_client()) + with pytest.raises(ValueError, match="exactly one"): + IPRoute("192.0.2.1") + + +def test_proxy_construction_and_secret_free_representations(monkeypatch): + calls = [] + original_client = httpx.Client + secret = "proxy-password-token" + route = IPRoute("192.0.2.1", proxy=f"http://user:{secret}@proxy.example:8080") + + def client(**kwargs): + calls.append(kwargs) + result = original_client(**kwargs) + monkeypatch.setattr(httpx, "Client", original_client) + return result + + monkeypatch.setattr(httpx, "Client", client) + with IPPool([route]) as routes: + assert calls[0]["proxy"] == route.proxy + assert calls[0]["trust_env"] is False + store = Store() + lease = acquire(routes, store) + assert secret not in repr(route) + repr(routes.routes[0]) + repr(lease) + with store.transaction() as db: + assert secret not in "\n".join(db.iterdump()) + with pytest.raises(ValueError) as error: + IPRoute("192.0.2.1", proxy=f"file://user:{secret}@host") + assert secret not in str(error.value) + + +def test_pool_does_not_close_injected_clients(): + routes = pool("192.0.2.1") + client = routes.routes[0].client + routes.close() + routes.close() + assert not client.is_closed + with pytest.raises(ValueError, match="closed"): + acquire(routes, Store()) + client.close() + + +def test_round_robin_is_deterministic_and_shared_across_pool_instances(): + store = Store() + with pool("192.0.2.1", "192.0.2.2") as first, pool("192.0.2.1", "192.0.2.2") as second: + assert [acquire(routes, store).public_ip for routes in [first, second, first, second]] == [ + "192.0.2.1", "192.0.2.2", "192.0.2.1", "192.0.2.2"] + with pytest.raises(Throttled): + acquire(first, store) + + +def test_ip_consumption_shared_across_restarts_and_subsets(tmp_path): + path = tmp_path / "shared.sqlite3" + clock = Clock() + with pool("192.0.2.1", "192.0.2.2") as first: + store = Store(path, clock=clock) + first.acquire(store, "provider", 2, 2, 60, max_wait=0) + store.close() + with pool("192.0.2.1") as restarted: + store = Store(path, clock=clock) + with pytest.raises(Throttled): + acquire(restarted, store) + # The same IP's unrelated provider remains independent. + assert restarted.acquire(store, "other", 2, 2, 60, max_wait=0).public_ip == "192.0.2.1" + store.close() + + +def test_weighted_rolling_window_waits_for_exact_capacity(): + clock = Clock() + store = Store(clock=clock) + with pool("192.0.2.1") as routes: + routes.acquire(store, "provider", 2, 3, 10, clock.sleep, 0) + clock.now += 3 + routes.acquire(store, "provider", 1, 3, 10, clock.sleep, 0) + lease = routes.acquire(store, "provider", 2, 3, 10, clock.sleep, 7) + assert lease.public_ip == "192.0.2.1" + assert clock.waits == [7] + with pytest.raises(Throttled) as error: + routes.acquire(store, "provider", 2, 3, 10, clock.sleep, 0) + assert error.value.retry_after == 10 + + +def test_waiters_charge_no_future_slot_and_recheck_cooldown(): + clock = Clock() + store = Store(clock=clock) + with pool("192.0.2.1", "192.0.2.2") as routes: + first = routes.acquire(store, "provider", 1, 1, 10, max_wait=0) + second = routes.acquire(store, "provider", 1, 1, 10, max_wait=0) + + def sleep(seconds): + clock.sleep(seconds) + if len(clock.waits) == 1: + routes.defer(store, "provider", first, 15) + # Another worker acquires the newly available route before this waiter. + assert routes.acquire(store, "provider", 1, 1, 10, max_wait=0).public_ip == second.public_ip + + admitted = routes.acquire(store, "provider", 1, 1, 10, sleep, 20) + assert admitted.public_ip == second.public_ip + assert clock.waits == [10, 10] + + +def test_one_ip_cooldown_does_not_poison_other_routes_or_providers(tmp_path): + clock = Clock() + first = Store(tmp_path / "shared.sqlite3", clock=clock) + second = Store(tmp_path / "shared.sqlite3", clock=clock) + with pool("192.0.2.1", "192.0.2.2") as routes, pool("192.0.2.1") as subset: + lease = acquire(routes, first) + routes.defer(first, "provider", lease, 25) + routes.defer(first, "provider", lease, 10) # cannot shorten an existing cooldown + assert acquire(routes, second).public_ip == "192.0.2.2" + with pytest.raises(Throttled) as error: + acquire(subset, second) + assert error.value.retry_after == 25 + assert subset.acquire(second, "other", 1, 2, 60, max_wait=0).public_ip == "192.0.2.1" + first.close() + second.close() + + +def test_shared_provider_limit_cannot_be_multiplied_by_more_ips(): + clock = Clock() + store = Store(clock=clock) + limits = (("provider:key:hashed", 1, 2, 60),) + with pool("192.0.2.1", "192.0.2.2", "192.0.2.3") as routes: + acquire(routes, store, limits=limits) + acquire(routes, store, limits=limits) + with pytest.raises(Throttled): + acquire(routes, store, limits=limits) + with store.transaction() as db: + assert db.execute("SELECT sum(weight) FROM ip_events WHERE scope=?", (limits[0][0],)).fetchone()[0] == 2 + assert db.execute("SELECT sum(weight) FROM ip_events WHERE scope LIKE 'provider:ip:%'").fetchone()[0] == 2 + + +def test_existing_core_pacing_is_atomic_with_ip_admission(): + clock = Clock() + store = Store(clock=clock) + pacing = (("provider:key", 2),) + with pool("192.0.2.1", "192.0.2.2") as routes: + store.pace("provider:key", 2, clock.sleep, 0) + with pytest.raises(Throttled) as error: + acquire(routes, store, pacing=pacing) + assert error.value.retry_after == 0.5 + with store.transaction() as db: + exists = db.execute("SELECT 1 FROM sqlite_master WHERE name='ip_events'").fetchone() + assert not exists or db.execute("SELECT count(*) FROM ip_events").fetchone()[0] == 0 + lease = routes.acquire(store, "provider", 1, 2, 60, clock.sleep, 1, pacing=pacing) + assert clock.waits == [0.5] + assert lease.public_ip == "192.0.2.1" + with pytest.raises(Throttled): + store.pace("provider:key", 2, clock.sleep, 0) + + +def test_core_throttle_and_global_cooldown_survive_route_selection(): + clock = Clock() + store = Store(clock=clock) + store.throttle("provider:key", 1) + store.defer("provider:key", 3) + with pool("192.0.2.1", "192.0.2.2") as routes: + routes.acquire(store, "provider", 1, 10, 60, clock.sleep, 3, pacing=(("provider:key", 10),)) + assert clock.waits == [3] + with pytest.raises(Throttled) as error: + routes.acquire(store, "provider", 1, 10, 60, clock.sleep, 0, pacing=(("provider:key", 10),)) + assert error.value.retry_after == 1 + + +def test_route_only_mode_still_observes_cooldown_and_pacing(): + clock = Clock() + store = Store(clock=clock) + with pool("192.0.2.1", "192.0.2.2") as routes: + lease = routes.acquire(store, "provider", 1, None, 60, max_wait=0) + assert lease.reservation_id is None + routes.defer(store, "provider", lease, 10) + assert routes.acquire(store, "provider", 1, None, 60, max_wait=0).public_ip == "192.0.2.2" + with store.transaction() as db: + assert db.execute("SELECT count(*) FROM ip_events").fetchone()[0] == 0 + + +def test_settlement_releases_only_this_lease_and_never_reinserts_expired_events(): + clock = Clock() + store = Store(clock=clock) + with pool("192.0.2.1") as routes: + first = routes.acquire(store, "provider", 4, 8, 10, max_wait=0) + second = routes.acquire(store, "provider", 4, 8, 10, max_wait=0) + routes.settle(store, second, 1) + assert routes.acquire(store, "provider", 3, 8, 10, max_wait=0).public_ip == first.public_ip + with pytest.raises(Throttled): + routes.acquire(store, "provider", 1, 8, 10, max_wait=0) + routes.settle(store, second, 1) # same-cost settlement is idempotent + with pytest.raises(ValueError): + routes.settle(store, second, 0) + with pytest.raises(ValueError): + routes.settle(store, second, 4) + with pytest.raises(ValueError): + routes.settle(store, first, 5) + clock.now += 10 + routes.acquire(store, "provider", 1, 8, 10, max_wait=0) + routes.settle(store, first, 0) + with store.transaction() as db: + assert db.execute("SELECT sum(weight) FROM ip_events").fetchone()[0] == 1 + + +def test_oversleep_cannot_dispatch_after_deadline(): + clock = Clock() + store = Store(clock=clock) + with pool("192.0.2.1") as routes: + routes.acquire(store, "provider", 1, 1, 10, max_wait=0) + with pytest.raises(Throttled): + routes.acquire(store, "provider", 1, 1, 10, lambda seconds: clock.sleep(seconds + 1), 10) + with store.transaction() as db: + assert db.execute("SELECT count(*) FROM ip_events").fetchone()[0] == 1 + + +def test_atomic_admission_across_concurrent_stores(tmp_path): + path = tmp_path / "shared.sqlite3" + clock = Clock() + stores = [Store(path, clock=clock) for _ in range(8)] + with pool("192.0.2.1", "192.0.2.2") as routes: + def request(index): + try: + return routes.acquire(stores[index % len(stores)], "provider", 1, 3, 60, max_wait=0) + except Throttled: + return None + with ThreadPoolExecutor(max_workers=8) as threads: + admitted = list(threads.map(request, range(40))) + assert sum(isinstance(lease, IPLease) for lease in admitted) == 6 + assert sum(lease.public_ip == "192.0.2.1" for lease in admitted if lease) == 3 + assert sum(lease.public_ip == "192.0.2.2" for lease in admitted if lease) == 3 + for store in stores: + store.close() + + +@pytest.mark.parametrize("kwargs", [ + {"weight": True}, {"weight": 0}, {"weight": -1}, {"capacity": 0}, + {"period": 0}, {"period": math.inf}, {"max_wait": -1}, {"max_wait": math.nan}, +]) +def test_invalid_admission_configuration_has_no_side_effects(kwargs): + defaults = dict(weight=1, capacity=2, period=60, max_wait=0) + defaults.update(kwargs) + store = Store() + with pool("192.0.2.1") as routes: + with pytest.raises(ValueError): + routes.acquire(store, "provider", **defaults) + with store.transaction() as db: + assert db.execute("SELECT count(*) FROM pacing").fetchone()[0] == 0 + + +def test_excessive_weight_fails_without_wait_or_charge(): + with pool("192.0.2.1") as routes: + with pytest.raises(BudgetExceeded): + routes.acquire(Store(), "provider", 3, 2, 60, lambda _: pytest.fail("must not sleep"), 60) + + +def test_window_cannot_be_shortened_after_history_has_been_pruned(): + store = Store() + with pool("192.0.2.1") as routes: + acquire(routes, store) + with pytest.raises(ValueError, match="period"): + routes.acquire(store, "provider", 1, 2, 10, max_wait=0) + + +def test_budget_callback_runs_only_on_final_admission_and_rolls_back_failures(): + clock = Clock() + store = Store(clock=clock) + callbacks = [] + with pool("192.0.2.1") as routes: + routes.acquire(store, "provider", 1, 1, 10, max_wait=0) + + def deny(db, now): + callbacks.append(now) + db.execute("INSERT INTO budgets VALUES('provider:key','credits','day',1)") + raise BudgetExceeded("provider: day allowance exhausted") + + with pytest.raises(BudgetExceeded): + routes.acquire(store, "provider", 1, 1, 10, clock.sleep, 10, + pacing=(("provider:key", 2),), before_admit=deny) + assert callbacks == [1010] + assert clock.waits == [10] + with store.transaction() as db: + assert db.execute("SELECT count(*) FROM budgets").fetchone()[0] == 0 + assert db.execute("SELECT count(*) FROM pacing").fetchone()[0] == 0 + lease = routes.acquire(store, "provider", 1, 1, 10, max_wait=0, + before_admit=lambda db, now: callbacks.append(now)) + assert lease.public_ip == "192.0.2.1" + assert callbacks == [1010, 1010] diff --git a/tests/test_hyperliquid.py b/tests/test_hyperliquid.py new file mode 100644 index 0000000..05b4fe8 --- /dev/null +++ b/tests/test_hyperliquid.py @@ -0,0 +1,297 @@ +import json + +import httpx +import pytest + +from dedomena.sources.core import InvalidResponse, SearchLimitExceeded, Store, canonical +from dedomena.sources.hyperliquid import Hyperliquid, info_policy + + +BOOK = {"coin": "BTC", "time": 1000, + "levels": [[{"px": "113377.0123456789", "sz": "7.6699", "n": 17}], []]} +META = [{"universe": [{"name": "BTC", "szDecimals": 5}], + "marginTables": [[50, {"marginTiers": [{"lowerBound": "0.0", "maxLeverage": 50}]}]], + "collateralToken": 0}, + [{"markPx": "113377.0123456789", "premium": None, "impactPxs": None, "midPx": None}]] +CANDLE = {"s": "BTC", "i": "1m", "t": 1000, "T": 60999, "n": 1, + "o": "100.0123456789", "h": "100.3", "l": "99.9", "c": "100.1", "v": "0.00123"} + + +def source(handler, **kwargs): + return Hyperliquid(client=httpx.Client(transport=httpx.MockTransport(handler)), + store=kwargs.pop("store", Store()), sleep=lambda _: None, **kwargs) + + +def funding(time, rate="0.0000125"): + return {"coin": "BTC", "time": time, "fundingRate": rate, "premium": "-0.00031774"} + + +def test_book_preserves_native_precision_receipt_and_replay(): + calls = [] + def handler(request): + calls.append(request) + assert request.method == "POST" + assert request.url == "https://api.hyperliquid.xyz/info" + assert json.loads(request.content) == {"type": "l2Book", "coin": "BTC", "nSigFigs": 5, "mantissa": 2} + return httpx.Response(200, json=BOOK) + with source(handler, cache_ttl=86400) as liquid: + page = liquid.order_book("BTC", n_sig_figs=5, mantissa=2) + assert page.complete and page.total == 1 + assert page.records == (BOOK,) + assert page.records[0]["levels"][0][0]["px"] == "113377.0123456789" + assert json.loads(page.provenance.request_body)["mantissa"] == 2 + assert json.loads(liquid.store.replay(page.provenance.snapshot_id).body) == BOOK + again = liquid.order_book("BTC", n_sig_figs=5, mantissa=2) + assert again.cache_hit and len(calls) == 1 + + +def test_metadata_retains_margin_tables_nulls_and_context_positions(): + with source(lambda _: httpx.Response(200, json=META)) as liquid: + page = liquid.markets() + assert page.records == ({"meta": META[0], "assetCtxs": META[1]},) + assert page.records[0]["assetCtxs"][0]["premium"] is None + assert page.complete and page.warnings + + +def test_spot_contexts_are_not_zipped_or_truncated_to_pair_universe(): + data = [{"universe": [{"name": "PURR/USDC", "tokens": [1, 0], "index": 0}], + "tokens": [{"name": "USDC", "index": 0}, {"name": "PURR", "index": 1}]}, + [{"coin": "PURR/USDC", "midPx": "0.15"}, {"coin": "@999", "midPx": None}]] + def handler(request): + assert json.loads(request.content) == {"type": "spotMetaAndAssetCtxs"} + return httpx.Response(200, json=data) + with source(handler) as liquid: + page = liquid.markets(spot=True) + assert len(page.records[0]["assetCtxs"]) == 2 + assert page.records[0]["meta"]["tokens"] == data[0]["tokens"] + + +def test_all_mids_is_one_native_map_including_spot_and_prediction_asset_ids(): + data = {"BTC": "113377.0123456789", "@107": "36.456", "#14720": "0.385"} + with source(lambda _: httpx.Response(200, json=data)) as liquid: + page = liquid.all_mids() + assert page.records == (data,) + assert page.complete and page.records_seen == 1 + + +def test_candles_preserve_precision_and_never_claim_complete_history(): + def handler(request): + assert json.loads(request.content) == {"type": "candleSnapshot", "req": { + "coin": "BTC", "interval": "1m", "startTime": 1000, "endTime": 2000}} + return httpx.Response(200, json=[CANDLE]) + with source(handler) as liquid: + page = liquid.candles("BTC", "1m", 1000, 2000) + assert page.records == (CANDLE,) + assert not page.complete and page.next_cursor is None + assert any("5,000" in warning for warning in page.warnings) + + +def test_empty_candles_do_not_prove_historical_coverage(): + with source(lambda _: httpx.Response(200, json=[])) as liquid: + assert not liquid.candles("BTC", "1d", 0, 1000).complete + + +def test_funding_pagination_inclusive_overlap_and_count_are_exact(): + calls = [] + first = [funding(time) for time in range(1, 501)] + second = [funding(500), funding(501)] + def handler(request): + body = json.loads(request.content) + calls.append(body) + assert body["endTime"] == 1000 + return httpx.Response(200, json=first if body["startTime"] == 1 else second) + with source(handler) as liquid: + pages = list(liquid.funding_history("BTC", 1, 1000)) + assert len(pages) == 2 + assert len(pages[0].records) == 500 and not pages[0].complete + assert pages[1].records == (funding(501),) + assert pages[1].complete and pages[1].records_seen == 501 + assert pages[1].next_cursor is None + assert [body["startTime"] for body in calls] == [1, 500] + assert liquid.usage()["upstream_records"] == 502 + assert liquid.usage()["records_delivered"] == 501 + assert liquid.usage()["boundary_repeats"] == 1 + + +def test_funding_checkpoint_pins_default_end_and_deduplicates_resume(): + store = Store(clock=lambda: 1.0) + bodies = [] + def handler(request): + body = json.loads(request.content) + bodies.append(body) + return httpx.Response(200, json=[funding(time) for time in range(1, 501)] + if body["startTime"] == 1 else [funding(500), funding(501)]) + with source(handler, store=store) as liquid: + page = next(liquid.funding_history("BTC", 1)) + store.clock = lambda: 2.0 + resumed = list(liquid.funding_history("BTC", 1, cursor=page.next_cursor)) + assert resumed[0].records == (funding(501),) + assert resumed[0].records_seen == 501 + assert resumed[0].complete + assert bodies[-1]["endTime"] == 1000 + with pytest.raises(ValueError): + list(liquid.funding_history("ETH", 1, cursor=page.next_cursor)) + with pytest.raises(ValueError): + list(liquid.funding_history("BTC", 2, cursor=page.next_cursor)) + with pytest.raises(ValueError): + list(liquid.funding_history("BTC", 1, 2000, cursor=page.next_cursor)) + + +def test_full_same_timestamp_funding_fails_before_partial_page_is_yielded(): + rows = [funding(1, str(index)) for index in range(500)] + with source(lambda _: httpx.Response(200, json=rows)) as liquid: + with pytest.raises(SearchLimitExceeded): + next(liquid.funding_history("BTC", 1, 1000)) + assert liquid.usage().get("records_delivered", 0) == 0 + + +def test_multiple_distinct_boundary_rows_are_deduplicated_without_losing_new_row(): + first = [funding(time) for time in range(1, 499)] + [funding(499, "0.1"), funding(499, "0.2")] + second = [funding(499, "0.1"), funding(499, "0.2"), funding(499, "0.3"), funding(500)] + with source(lambda request: httpx.Response(200, json=first + if json.loads(request.content)["startTime"] == 1 else second)) as liquid: + pages = list(liquid.funding_history("BTC", 1, 1000)) + assert pages[-1].records == (funding(499, "0.3"), funding(500)) + assert pages[-1].records_seen == 502 + + +@pytest.mark.parametrize("method,args,kwargs", [ + ("order_book", ("../BTC",), {}), + ("order_book", ("BTC",), {"n_sig_figs": True}), + ("order_book", ("BTC",), {"mantissa": 2}), + ("order_book", ("BTC",), {"n_sig_figs": 5, "mantissa": 3}), + ("all_mids", (), {"dex": "../test"}), + ("markets", (), {"spot": "true"}), + ("markets", (), {"spot": True, "dex": "xyz"}), + ("candles", ("BTC", "2m", 1, 2), {}), + ("candles", ("BTC", "1m", True, 2), {}), + ("candles", ("BTC", "1m", 3, 2), {}), + ("funding_history", ("BTC", 1, 2), {"cursor": "{}"}), + ("funding_history", ("BTC", -1, 2), {}), + ("funding_history", ("BTC", 1, True), {}), +]) +def test_invalid_queries_fail_before_network(method, args, kwargs): + with source(lambda _: pytest.fail("unexpected network")) as liquid: + with pytest.raises(ValueError): + result = getattr(liquid, method)(*args, **kwargs) + if method == "funding_history": + list(result) + + +@pytest.mark.parametrize("method,body", [ + ("all_mids", {"BTC": 123.4}), + ("all_mids", {"BTC": "NaN"}), + ("order_book", {**BOOK, "coin": "ETH"}), + ("order_book", {**BOOK, "time": True}), + ("order_book", {**BOOK, "levels": [[{"px": "1", "sz": "-1", "n": 1}], []]}), + ("markets", [{"universe": [{"name": "BTC"}]}, []]), + ("markets", [{"universe": [{"name": "BTC"}, {"name": "BTC"}]}, [{}, {}]]), +]) +def test_malformed_snapshot_responses_fail_closed(method, body): + with source(lambda _: httpx.Response(200, json=body)) as liquid: + with pytest.raises(InvalidResponse): + getattr(liquid, method)("BTC") if method == "order_book" else getattr(liquid, method)() + + +@pytest.mark.parametrize("row", [ + {**CANDLE, "s": "ETH"}, {**CANDLE, "i": "1d"}, {**CANDLE, "t": True}, + {**CANDLE, "t": 3000}, {**CANDLE, "c": 100.1}, {**CANDLE, "v": "Infinity"}, +]) +def test_malformed_candle_identity_and_values_fail_closed(row): + with source(lambda _: httpx.Response(200, json=[row])) as liquid: + with pytest.raises(InvalidResponse): + liquid.candles("BTC", "1m", 1000, 2000) + + +@pytest.mark.parametrize("rows", [ + [funding(2), funding(1)], [funding(1), funding(1)], + [{**funding(1), "coin": "ETH"}], [funding(3000)], + [{**funding(1), "fundingRate": 0.1}], [{**funding(1), "time": True}], +]) +def test_malformed_funding_order_identity_and_values_fail_closed(rows): + with source(lambda _: httpx.Response(200, json=rows)) as liquid: + with pytest.raises(InvalidResponse): + list(liquid.funding_history("BTC", 1, 2000)) + + +@pytest.mark.parametrize("body", [ + {"type": "order", "coin": "BTC"}, {"type": "userFills", "user": "0xabc"}, + {"type": "recentTrades", "coin": "BTC"}, {"type": []}, + {"type": "allMids", "user": "0xabc"}, {"type": "l2Book", "coin": "BTC", "action": {}}, + {"type": "candleSnapshot", "req": {"coin": "BTC", "interval": "1m", "startTime": 0}}, + {"type": "fundingHistory", "coin": "BTC", "startTime": False, "endTime": 1000}, +]) +def test_transport_policy_rejects_account_actions_unknown_types_and_fields(body): + with pytest.raises(ValueError): + info_policy(canonical(body).encode()) + with source(lambda _: pytest.fail("unexpected network")) as liquid: + with pytest.raises(ValueError): + liquid._request("POST", "/info", content=canonical(body).encode()) + + +@pytest.mark.parametrize("path", ["/exchange", "/info/../exchange", "//elsewhere/info"]) +def test_transport_cannot_reach_exchange_even_directly(path): + with source(lambda _: pytest.fail("unexpected network")) as liquid: + with pytest.raises(ValueError): + liquid._request("POST", path, content=b'{"type":"allMids"}') + + +@pytest.mark.parametrize("body,maximum,actual_count,actual", [ + ({"type": "allMids"}, 2, None, 2), + ({"type": "metaAndAssetCtxs"}, 20, None, 20), + ({"type": "fundingHistory", "coin": "BTC", "startTime": 0, "endTime": 1000}, 45, 21, 22), + ({"type": "candleSnapshot", "req": {"coin": "BTC", "interval": "1m", "startTime": 0, "endTime": 1000}}, 104, 61, 22), +]) +def test_provider_policy_reserves_maximum_and_settles_by_count(body, maximum, actual_count, actual): + reserved, count = info_policy(canonical(body).encode()) + assert reserved == maximum + if actual_count is not None: + assert count(canonical([{}] * actual_count).encode()) == actual + with pytest.raises(ValueError): + count(b'{}') + else: + assert count is None + + +@pytest.mark.parametrize("limit", [(1201, 60), (1200, 59), None, (True, 60), (1200, float("nan"))]) +def test_provider_ip_policy_cannot_be_weakened(limit): + with pytest.raises(ValueError): + source(lambda _: pytest.fail("unexpected network"), ip_limit=limit) + + +def test_quota_and_fetch_alias_are_agent_accessible(): + with source(lambda _: httpx.Response(200, json=BOOK)) as liquid: + assert liquid.quota()["ip_weight_per_minute"] == 1200 + assert liquid.quota()["api_key_required"] is False + assert liquid.ip_throttle_only + assert liquid.fetch("BTC").records == (BOOK,) + + + +def test_direct_transport_cannot_undercharge_a_count_weighted_endpoint(): + body = {"type": "candleSnapshot", "req": {"coin": "BTC", "interval": "1m", + "startTime": 0, "endTime": 1000}} + with source(lambda _: httpx.Response(200, json=[])) as liquid: + liquid._request("POST", "/info", content=canonical(body).encode(), + rate_weight=1, actual_weight=lambda _: 1) + usage = liquid.usage() + assert usage["ip_weight_reserved"] == 104 + assert usage["ip_weight_used"] == 20 + + +@pytest.mark.parametrize("kwargs", [{"params": {"type": "allMids"}}, + {"form": {"type": "allMids"}}, + {"auth": ("user", "pass")}, + {"private_params": {"api_key": "secret"}}]) +def test_hyperliquid_transport_only_accepts_public_json_body(kwargs): + with source(lambda _: pytest.fail("unexpected network")) as liquid: + with pytest.raises(ValueError): + liquid._request("POST", "/info", content=b'{"type":"allMids"}', **kwargs) + + +def test_malformed_response_keeps_conservative_weight_reservation(): + with source(lambda _: httpx.Response(200, json={"unexpected": "schema"})) as liquid: + with pytest.raises(InvalidResponse): + liquid.candles("BTC", "1m", 0, 1000) + usage = liquid.usage() + assert usage["ip_weight_reserved"] == usage["ip_weight_used"] == 104 diff --git a/tests/test_hyperliquid_cli.py b/tests/test_hyperliquid_cli.py new file mode 100644 index 0000000..323b878 --- /dev/null +++ b/tests/test_hyperliquid_cli.py @@ -0,0 +1,125 @@ +"""Public market operations, route configuration and offline replay in JSON CLI.""" +import json + +import httpx +import pytest + +from dedomena.sources import Hyperliquid, IPPool, IPRoute, Store +from dedomena.sources import __main__ as cli + + +@pytest.fixture +def market_cli(monkeypatch): + store, requests = Store(), [] + def handler(request): + body = json.loads(request.content) + requests.append(body) + kind = body["type"] + if kind == "allMids": + data = {"BTC": "100.00000001"} + elif kind in ("metaAndAssetCtxs", "spotMetaAndAssetCtxs"): + data = [{"universe": [{"name": "BTC"}], "tokens": []}, [{"markPx": "100.1"}]] + elif kind == "l2Book": + data = {"coin": body["coin"], "time": 1234, "levels": [[], []]} + elif kind == "candleSnapshot": + data = [{"s": "BTC", "i": "1h", "t": 0, "T": 3599999, "n": 2, + "o": "100.01", "h": "101", "l": "99", "c": "100.02", "v": "1.1"}] + else: + assert kind == "fundingHistory" + data = [{"coin": "BTC", "time": 1234, "fundingRate": "0.0001", "premium": "0.001"}] + return httpx.Response(200, json=data) + client = httpx.Client(transport=httpx.MockTransport(handler)) + def constructor(**options): + options.pop("store", None) + if "ip_pool" not in options: + options["client"] = client + return Hyperliquid(store=store, **options) + monkeypatch.setattr(cli, "Hyperliquid", constructor) + yield requests, client, store + client.close() + store.close() + + +@pytest.mark.parametrize("args,kind,complete", [ + (["mids", "hyperliquid", "--dex", "xyz"], "allMids", True), + (["markets", "hyperliquid"], "metaAndAssetCtxs", True), + (["catalogue", "hyperliquid", "--spot"], "spotMetaAndAssetCtxs", True), + (["book", "hyperliquid", "BTC"], "l2Book", True), + (["candles", "hyperliquid", "BTC", "--interval", "1h", "--start-time", "0", + "--end-time", "3600000"], "candleSnapshot", False), + (["funding", "hyperliquid", "BTC", "--start-time", "0", "--end-time", "3600000"], + "fundingHistory", True), +]) +def test_market_operations_emit_native_page_and_truthful_receipt(market_cli, capsys, args, kind, complete): + assert cli.main(args) == 0 + requests, _, _ = market_cli + page, receipt = [json.loads(line) for line in capsys.readouterr().out.splitlines()] + assert requests[0]["type"] == kind + assert page["provenance"]["method"] == "POST" + assert json.loads(page["provenance"]["request_body"]) == requests[0] + assert page["complete"] == receipt["complete"] == complete + assert receipt["usage"]["ip_weight_used"] == (2 if kind in ("allMids", "l2Book") else + 21 if kind in ("candleSnapshot", "fundingHistory") else 20) + + +def test_fetch_alias_and_benchmark_mids(market_cli, capsys): + assert cli.main(["fetch", "hyperliquid", "BTC"]) == 0 + page = json.loads(capsys.readouterr().out) + assert page["records"][0]["coin"] == "BTC" + assert cli.main(["benchmark", "hyperliquid", "--operation", "mids"]) == 0 + receipt = json.loads(capsys.readouterr().out) + assert receipt["operation"] == "mids" and receipt["complete"] + assert receipt["records"] == 1 + + +def test_route_file_configures_shared_pool_and_receipt(tmp_path, monkeypatch, market_cli, capsys): + _, client, _ = market_cli + configured, pools = [], [] + def pool_factory(routes): + routes = list(routes) + configured.extend(routes) + pool = IPPool([IPRoute(route.public_ip, client=client) for route in routes]) + pools.append(pool) + return pool + monkeypatch.setattr(cli, "IPPool", pool_factory) + path = tmp_path / "routes.json" + path.write_text(json.dumps([{"public_ip": "203.0.113.1", "local_address": "10.0.0.1"}])) + assert cli.main(["mids", "hyperliquid", "--ip-routes", str(path)]) == 0 + page = json.loads(capsys.readouterr().out.splitlines()[0]) + assert configured[0].local_address == "10.0.0.1" + assert page["provenance"]["egress_ip_sha256"] is not None + assert pools[0]._closed and not client.is_closed + + +@pytest.mark.parametrize("contents", ["bad JSON", "{}", "[]", + '[{"public_ip":"203.0.113.1","client":"not a client"}]', + '[{"public_ip":"203.0.113.1","proxy":"http://user:secret@proxy.invalid/?bad"}]', +]) +def test_bad_route_file_is_structured_and_never_dispatches(tmp_path, market_cli, capsys, contents): + path = tmp_path / "routes.json" + path.write_text(contents) + assert cli.main(["mids", "hyperliquid", "--ip-routes", str(path)]) == 2 + output = capsys.readouterr() + assert not output.out and not market_cli[0] + assert json.loads(output.err)["error"] == "ValueError" + assert "secret" not in output.err + + +def test_missing_route_file_is_structured(tmp_path, market_cli, capsys): + assert cli.main(["mids", "hyperliquid", "--ip-routes", str(tmp_path / "absent")]) == 2 + assert not market_cli[0] + assert json.loads(capsys.readouterr().err)["error"] == "ValueError" + + +def test_hyperliquid_replay_needs_no_routes(tmp_path, capsys): + path = tmp_path / "saved.sqlite3" + store = Store(path) + client = httpx.Client(transport=httpx.MockTransport(lambda _: httpx.Response(200, json={"BTC": "100.00001"}))) + with Hyperliquid(store=store, client=client) as source: + page = source.all_mids() + store.close() + client.close() + assert cli.main(["replay", "hyperliquid", page.provenance.snapshot_id, + "--store", str(path), "--ip-routes", "not-needed.json"]) == 0 + replay = json.loads(capsys.readouterr().out) + assert json.loads(replay["body"]) == {"BTC": "100.00001"} diff --git a/tests/test_ip_transport.py b/tests/test_ip_transport.py new file mode 100644 index 0000000..73fbedb --- /dev/null +++ b/tests/test_ip_transport.py @@ -0,0 +1,214 @@ +"""Quota, cache, dispatch and retry contracts across routed source clients.""" +from datetime import datetime, timezone + +import httpx +import pytest + +from dedomena.sources import BudgetExceeded, IPPool, IPRoute, OpenAlex, Store, Throttled +from dedomena.sources.core import SourceError, Transport, digest, window + + +class Clock: + def __init__(self, now=1000.0): + self.now = now + + def __call__(self): + return self.now + + def sleep(self, seconds): + self.now += seconds + + +def routed(*handlers): + return IPPool([IPRoute(f"203.0.113.{i + 1}", client=httpx.Client( + transport=httpx.MockTransport(handler))) for i, handler in enumerate(handlers)]) + + +def source(store, clock, pool, **kwargs): + kwargs.setdefault("requests_per_second", 1000) + return Transport("test", "https://example.org", store=store, ip_pool=pool, + sleep=clock.sleep, **kwargs) + + +def test_shared_cache_identity_and_replay_preserve_acquisition_route(): + clock, sent = Clock(), [] + def handler(label): + return lambda request: sent.append(label) or httpx.Response(200, json={"value": "1.00001"}) + with routed(handler("first"), handler("second")) as pool: + store = Store(clock=clock) + client = source(store, clock, pool, ip_limit=(10, 60)) + first = client._request("GET", "/same") + cached = client._request("GET", "/same") + refreshed = client._request("GET", "/same", refresh=True) + assert sent == ["first", "second"] + assert cached.cache_hit and cached.provenance == first.provenance + assert first.provenance.request_sha256 == refreshed.provenance.request_sha256 + assert first.provenance.egress_ip_sha256 == digest(b"203.0.113.1") + assert refreshed.provenance.egress_ip_sha256 == digest(b"203.0.113.2") + assert store.replay(first.provenance.snapshot_id).provenance == first.provenance + assert client.usage()["ip_weight_used"] == 2 + client.close() + assert not pool._closed + store.close() + + +def test_openalex_pool_cannot_multiply_key_budget_or_charge_unsent_requests(): + clock, sent = Clock(), [] + def handler(request): + sent.append(request.url.path) + if request.url.path == "/rate-limit": + return httpx.Response(200, json={"rate_limit": { + "credits_limit": 10000, "credits_remaining": 10000, + "credit_costs": {"singleton": 0, "list": 1, "search": 10}}}) + return httpx.Response(200, json={"results": [{"id": "https://openalex.org/W1"}], + "meta": {"count": 1, "next_cursor": None}}, + headers={"X-RateLimit-Credits-Used": "1"}) + store = Store(clock=clock) + with routed(handler, handler) as pool: + client = OpenAlex("fixture-secret", ip_pool=pool, store=store, + daily_credit_limit=1, sleep=clock.sleep) + assert next(client.search(filter="type:article")).complete + with pytest.raises(BudgetExceeded): + next(client.search(filter="type:article", refresh=True)) + assert sent == ["/rate-limit", "/works"] + with store.transaction() as db: + assert db.execute("SELECT sum(weight) FROM ip_events").fetchone()[0] == 2 + assert store.used(client.scope, "credits", "day") == 1 + store.close() + + +def test_verified_weight_releases_capacity_for_a_known_cheap_request(): + clock = Clock() + store = Store(clock=clock) + with routed(lambda _: httpx.Response(200, json={})) as pool: + client = source(store, clock, pool, ip_limit=(4, 60), max_wait=0) + client._request("GET", "/large-estimate", rate_weight=4, actual_weight=lambda _: 2) + clock.sleep(0.01) + client._request("GET", "/cheap", rate_weight=2) + clock.sleep(0.01) + with pytest.raises(Throttled): + client._request("GET", "/exhausted", rate_weight=1) + assert client.usage()["ip_weight_reserved"] == 6 + assert client.usage()["ip_weight_used"] == 4 + store.close() + + +def test_uncertain_send_retains_weight_and_error_does_not_expose_credentials(): + clock = Clock() + store = Store(clock=clock) + def handler(request): + raise httpx.ProxyError("http://user:private-password@example.net", request=request) + with routed(handler) as pool: + client = source(store, clock, pool, ip_limit=(4, 60), max_wait=0, max_attempts=1) + with pytest.raises(SourceError, match="transport failure") as error: + client._request("GET", "/uncertain", rate_weight=4) + assert "private-password" not in str(error.value) + clock.sleep(0.01) + with pytest.raises(Throttled): + client._request("GET", "/next") + assert client.usage()["ip_weight_used"] == 4 + store.close() + + +def test_budget_work_precedes_fresh_dispatch_lease(): + clock, sent = Clock(), [] + store = Store(clock=clock) + original = store._reserve_in + delayed = False + def reserve(*args, **kwargs): + nonlocal delayed + if not delayed: + delayed = True + clock.sleep(60) + return original(*args, **kwargs) + store._reserve_in = reserve + with routed(lambda _: sent.append(clock()) or httpx.Response(200, json={})) as pool: + client = source(store, clock, pool, ip_limit=(1, 60)) + client._request("GET", "/first") + client._request("GET", "/second") + assert sent == [1060.0, 1120.0] + store.close() + + +def test_slow_budget_admission_reanchors_both_calendar_windows(): + clock = Clock(datetime(2026, 10, 4, 23, 59, 59, tzinfo=timezone.utc).timestamp()) + store = Store(clock=clock) + old_day, old_week = window("day", clock()), window("week", clock()) + original = store._reserve_in + delayed = False + def reserve(*args, **kwargs): + nonlocal delayed + result = original(*args, **kwargs) + if not delayed: + delayed = True + clock.sleep(2) + return result + store._reserve_in = reserve + with routed(lambda _: httpx.Response(200, content=b"{}")) as pool: + client = source(store, clock, pool, ip_limit=(10, 60), max_response_bytes=10) + response = client._request("GET", "/data", estimated_credits=1, + credit_limit=100, byte_limit=100) + assert response.provenance.retrieved_at_utc.startswith("2026-10-05") + assert store.used(client.scope, "credits", "day") == 1 + assert store.used(client.scope, "bytes", "week") == 2 + with store.transaction() as db: + rows = db.execute("SELECT window FROM budgets UNION SELECT window FROM budget_reservations").fetchall() + assert all(key not in (old_day, old_week) for (key,) in rows) + store.close() + + +@pytest.mark.parametrize("ip_only", [False, True]) +def test_429_preserves_scope_while_an_independent_ip_can_continue(ip_only): + clock, sent = Clock(), [] + store = Store(clock=clock) + def first(request): + sent.append(("first", clock())) + return httpx.Response(429, headers={"Retry-After": "30"}) + def second(request): + sent.append(("second", clock())) + return httpx.Response(200, json={}) + with routed(first, second) as pool: + client = source(store, clock, pool, ip_limit=(10, 60), + ip_throttle_only=ip_only, max_attempts=1) + with pytest.raises(Throttled): + client._request("GET", "/first") + assert client._request("GET", "/second").provenance.http_status == 200 + assert [label for label, _ in sent] == ["first", "second"] + assert (sent[1][1] - sent[0][1] < 30) == ip_only + store.close() + + +@pytest.mark.parametrize("limit", [None, (101, 1), (100, 2), (True, 1)]) +def test_openalex_provider_ip_policy_cannot_be_disabled_or_expanded(limit): + with pytest.raises(ValueError, match="OpenAlex IP policy"): + OpenAlex("fixture-secret", ip_limit=limit) + + +@pytest.mark.parametrize("use_pool", [False, True]) +def test_positive_admission_deadline_includes_slow_budget_work(use_pool): + clock, sent = Clock(), [] + store = Store(clock=clock) + original = store._reserve_in + delayed = False + def reserve(*args, **kwargs): + nonlocal delayed + if not delayed: + delayed = True + clock.sleep(2) + return original(*args, **kwargs) + store._reserve_in = reserve + handler = lambda _: sent.append(clock()) or httpx.Response(200, json={}) + with routed(handler) as pool: + options = {"ip_pool": pool, "ip_limit": (10, 60)} if use_pool else { + "client": httpx.Client(transport=httpx.MockTransport(handler))} + client = Transport("test", "https://example.org", store=store, + sleep=clock.sleep, max_wait=1, **options) + with pytest.raises(Throttled): + client._request("GET", "/data", estimated_credits=1, credit_limit=100) + assert not sent + with store.transaction() as db: + assert db.execute("SELECT count(*) FROM budget_reservations").fetchone()[0] == 0 + assert db.execute("SELECT count(*) FROM pacing").fetchone()[0] == 0 + if not use_pool: + options["client"].close() + store.close()