diff --git a/app/__init__.py b/app/__init__.py index a2d54ec..13d99dc 100644 --- a/app/__init__.py +++ b/app/__init__.py @@ -183,6 +183,14 @@ def _evict_if_needed(self) -> None: pass +@dataclass +class GpkgLimiter: + """Per-app concurrency guard for the hydrofabric gpkg endpoint.""" + + semaphore: BoundedSemaphore + queue_timeout_s: float + + def get_catalog(request: Request) -> Catalog: """Gets the pyiceberg catalog reference from the app state diff --git a/app/routers/hydrofabric/router.py b/app/routers/hydrofabric/router.py index c989e32..b9f157b 100644 --- a/app/routers/hydrofabric/router.py +++ b/app/routers/hydrofabric/router.py @@ -3,6 +3,7 @@ import pathlib import sqlite3 import tempfile +import time import uuid import geopandas as gpd diff --git a/app/routers/streamflow_observations/router.py b/app/routers/streamflow_observations/router.py index bbb53f3..006f2f1 100644 --- a/app/routers/streamflow_observations/router.py +++ b/app/routers/streamflow_observations/router.py @@ -133,7 +133,6 @@ def validate_identifier(identifier: str, request: Request | None = None): @api_router.get("/{identifier}/info", tags=["Streamflow Observations"]) def get_identifier_info( - request: Request, identifier: str = Path( ..., description="Station/gauge ID", diff --git a/pyproject.toml b/pyproject.toml index cd03cfb..3a877ff 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -165,6 +165,7 @@ convention = "numpy" [tool.ruff.lint.per-file-ignores] "docs/*" = ["I"] "tests/*" = ["D"] +"scripts/*" = ["D", "BLE001"] "*/__init__.py" = ["F401"] [tool.mypy] diff --git a/scripts/load_test/analyze.py b/scripts/load_test/analyze.py new file mode 100644 index 0000000..a9d241b --- /dev/null +++ b/scripts/load_test/analyze.py @@ -0,0 +1,86 @@ +"""Summarise stats.csv from monitor.sh — look for memory creep, OOM proximity, CPU saturation.""" + +from __future__ import annotations + +import argparse +import csv +import re +from pathlib import Path + + +def to_bytes(s: str) -> float: + s = s.strip() + m = re.match(r"([0-9.]+)\s*([KMGT]?i?B)", s, re.I) + if not m: + return 0.0 + v = float(m.group(1)) + unit = m.group(2).lower() + mult = { + "b": 1, + "kb": 1e3, + "mb": 1e6, + "gb": 1e9, + "tb": 1e12, + "kib": 1024, + "mib": 1024**2, + "gib": 1024**3, + "tib": 1024**4, + }.get(unit, 1) + return v * mult + + +def main() -> None: + p = argparse.ArgumentParser() + p.add_argument("--stats", default="scripts/load_test/results/stats.csv") + args = p.parse_args() + path = Path(args.stats) + if not path.exists(): + print(f"no stats file at {path}") + return + + rows = list(csv.DictReader(path.open())) + if not rows: + print("no samples") + return + + cpu_vals = [float(r["cpu_pct"].rstrip("%")) for r in rows if r["cpu_pct"]] + mem_vals = [to_bytes(r["mem_usage"]) for r in rows if r["mem_usage"]] + mem_pct_vals = [float(r["mem_pct"].rstrip("%")) for r in rows if r["mem_pct"]] + + def q(xs, p): + if not xs: + return 0.0 + xs = sorted(xs) + return xs[int((p / 100.0) * (len(xs) - 1))] + + print(f"samples: {len(rows)}") + print( + f"cpu % min={min(cpu_vals):6.1f} mean={sum(cpu_vals) / len(cpu_vals):6.1f} " + f"p95={q(cpu_vals, 95):6.1f} max={max(cpu_vals):6.1f}" + ) + print( + f"mem GiB min={min(mem_vals) / 1024**3:6.2f} mean={sum(mem_vals) / len(mem_vals) / 1024**3:6.2f} " + f"p95={q(mem_vals, 95) / 1024**3:6.2f} max={max(mem_vals) / 1024**3:6.2f}" + ) + print( + f"mem % min={min(mem_pct_vals):6.1f} mean={sum(mem_pct_vals) / len(mem_pct_vals):6.1f} " + f"p95={q(mem_pct_vals, 95):6.1f} max={max(mem_pct_vals):6.1f}" + ) + + # Creep detection: compare first-third vs last-third mean memory + third = len(mem_vals) // 3 + if third >= 3: + head = sum(mem_vals[:third]) / third + tail = sum(mem_vals[-third:]) / third + creep = (tail - head) / max(head, 1) + print(f"mem creep (last-third vs first-third): {creep * 100:+.1f}%") + if creep > 0.25: + print(" ⚠️ > 25% growth — possible leak or unbounded caching") + elif creep > 0.10: + print(" ⚠️ > 10% growth — keep an eye on it over a longer run") + else: + print(" ✅ memory stable") + + +if __name__ == "__main__": + main() diff --git a/scripts/load_test/docker-compose.load.yml b/scripts/load_test/docker-compose.load.yml new file mode 100644 index 0000000..fbb3c38 --- /dev/null +++ b/scripts/load_test/docker-compose.load.yml @@ -0,0 +1,35 @@ +version: "2.4" +services: + api: + build: + context: ../.. + dockerfile: docker/Dockerfile.api + image: icefabric-api:loadtest + container_name: icefabric-api-loadtest + ports: + - "127.0.0.1:8000:8000" + env_file: + - ../../.env + environment: + - ICEFABRIC_DEPLOY_ENV=test + # Curb glibc per-thread arena fragmentation on heavy numpy/pandas use. + - MALLOC_ARENA_MAX=2 + # Stop numerical libs and polars from spawning 1 thread per core; + # we only have ~2 vCPU of budget and already fan out via asyncio + # threadpool + subset_nhf's ThreadPoolExecutor. + - OMP_NUM_THREADS=2 + - OPENBLAS_NUM_THREADS=2 + - MKL_NUM_THREADS=2 + - POLARS_MAX_THREADS=2 + # Emulating m6i.xlarge (4 vCPU, 16 GB) for this run. + cpus: 4.0 + mem_limit: 16g + memswap_limit: 16g + oom_kill_disable: false + restart: "no" + healthcheck: + test: ["CMD", "curl", "-f", "--head", "http://localhost:8000/health"] + interval: 15s + timeout: 10s + retries: 20 + start_period: 300s diff --git a/scripts/load_test/load_test.py b/scripts/load_test/load_test.py new file mode 100644 index 0000000..294b78d --- /dev/null +++ b/scripts/load_test/load_test.py @@ -0,0 +1,354 @@ +""" +Load test for the icefabric API. + +Targets 3 commonly-used endpoints at ~100 requests/minute total, with valid +CONUS NHF identifiers. Designed to run against a locally Dockerized API with +t3.large-equivalent resource caps. + +Usage +----- + python load_test.py --base-url http://localhost:8000 --rpm 100 --duration 300 +""" + +from __future__ import annotations + +import argparse +import asyncio +import csv +import json +import random +import statistics +import time +from collections import Counter, defaultdict +from dataclasses import dataclass, field +from pathlib import Path + +import httpx + +# --------------------------------------------------------------------------- +# Identifier pools (valid CONUS NHF) +# --------------------------------------------------------------------------- +# NHF CONUS VPU IDs. These are the standard National Hydrofabric VPU codes. +VPU_IDS: list[str] = [f"{i:02d}" for i in range(1, 19)] # "01".."18" + +# Known-valid NHF flowpath IDs (integers). Populated at runtime from the +# `/available` / hydrofabric responses when possible; defaults below are from +# the router's documented examples and a safe spread. +FLOWPATH_IDS: list[int] = [3490271] + +# Parameter-metadata modules +MODULES: list[str] = [ + "CFE-S", + "CFE-X", + "LASAM", + "LSTM", + "Noah-OWP-Modular", + "PET", + "Sac-SMA", + "SFT", + "SMP", + "Snow-17", + "T-Route", + "TopModel", + "Topoflow-Glacier", + "UEB", +] + +# Fallback gauge IDs (USGS) — the driver will try to fetch a fresh list from +# /streamflow_observations/available at start; these are the safety net. +FALLBACK_GAGES: list[str] = [ + "01010000", + "01031500", + "02GC002", + "08102730", + "01013500", + "01022500", + "01030500", + "01047000", +] + + +# --------------------------------------------------------------------------- +# Result tracking +# --------------------------------------------------------------------------- +@dataclass +class Sample: + endpoint: str + url: str + status: int + latency_ms: float + error: str | None = None + bytes: int = 0 + started_at: float = 0.0 + + +@dataclass +class Results: + samples: list[Sample] = field(default_factory=list) + + def add(self, s: Sample) -> None: + self.samples.append(s) + + def summary(self) -> dict: + by_ep: dict[str, list[Sample]] = defaultdict(list) + for s in self.samples: + by_ep[s.endpoint].append(s) + + out: dict = {"overall": self._stats(self.samples), "by_endpoint": {}} + for ep, xs in by_ep.items(): + out["by_endpoint"][ep] = self._stats(xs) + return out + + @staticmethod + def _stats(xs: list[Sample]) -> dict: + if not xs: + return {"count": 0} + lats = [s.latency_ms for s in xs] + codes = Counter(s.status for s in xs) + ok = sum(1 for s in xs if 200 <= s.status < 300) + errs = sum(1 for s in xs if s.status >= 500 or s.status == 0) + lats_sorted = sorted(lats) + + def pct(p): + if not lats_sorted: + return 0.0 + k = int(round((p / 100.0) * (len(lats_sorted) - 1))) + return lats_sorted[k] + + return { + "count": len(xs), + "success": ok, + "errors_5xx_or_conn": errs, + "status_codes": dict(codes), + "latency_ms": { + "min": round(min(lats), 1), + "mean": round(statistics.mean(lats), 1), + "p50": round(pct(50), 1), + "p95": round(pct(95), 1), + "p99": round(pct(99), 1), + "max": round(max(lats), 1), + }, + "total_bytes": sum(s.bytes for s in xs), + } + + +# --------------------------------------------------------------------------- +# Request builders +# --------------------------------------------------------------------------- +def build_hydrofabric_request(gages: list[str]) -> tuple[str, str]: + """Hydrofabric NHF subset — realistic mix of VPU and gage requests. + + VPU subsets are the largest (full region); gage subsets vary widely by + basin size. Random uniform sampling over the ~200 gauge pool gives the + natural cost distribution (headwater gauges are cheap, basin-outlet + gauges are expensive). + """ + # ~40% VPU (consistent heavy), ~60% gage (variable cost) + if random.random() < 0.4: + ident = random.choice(VPU_IDS) + id_type = "vpu_id" + else: + ident = random.choice(gages) + id_type = "gage_id" + url = f"/api/v1/hydrofabric/{ident}/gpkg?id_type={id_type}&source=nhf&domain=CONUS" + return "hydrofabric_gpkg", url + + +def build_parameter_metadata_request(gages: list[str]) -> tuple[str, str]: + """Parameter metadata with gage_id — always the heavy subset path.""" + mod = random.choice(MODULES) + gage = random.choice(gages) + url = f"/api/v1/modules/parameter_metadata/?modules={mod}&gage_id={gage}&domain=CONUS&source=hf" + return "parameter_metadata", url + + +def build_streamflow_request(gages: list[str]) -> tuple[str, str]: + gage = random.choice(gages) + url = f"/api/v1/streamflow_observations/{gage}/info" + return "streamflow_info", url + + +def pick_request(gages: list[str], weights: tuple[int, int, int]) -> tuple[str, str]: + """Pick a weighted random endpoint. weights = (hf, param, streamflow).""" + total = sum(weights) + r = random.randint(1, total) + if r <= weights[0]: + return build_hydrofabric_request(gages) + elif r <= weights[0] + weights[1]: + return build_parameter_metadata_request(gages) + else: + return build_streamflow_request(gages) + + +# --------------------------------------------------------------------------- +# Driver +# --------------------------------------------------------------------------- +async def fetch_available_gages(client: httpx.AsyncClient) -> list[str]: + try: + r = await client.get("/api/v1/streamflow_observations/available?limit=200", timeout=60) + if r.status_code == 200: + data = r.json() + # Endpoint may return dict or list; try common shapes + if isinstance(data, dict): + for key in ("identifiers", "ids", "available"): + if key in data and isinstance(data[key], list): + return [str(x) for x in data[key] if x][:200] + elif isinstance(data, list): + return [str(x) for x in data if x][:200] + except Exception as e: + print(f"[discover] couldn't fetch /available: {e}") + return FALLBACK_GAGES + + +async def one_request( + client: httpx.AsyncClient, + ep: str, + url: str, + results: Results, + timeout_s: float, +) -> None: + started = time.perf_counter() + wall = time.time() + status = 0 + err = None + nbytes = 0 + try: + # Stream so we don't buffer entire gpkg into memory on the client side + async with client.stream("GET", url, timeout=timeout_s) as r: + status = r.status_code + async for chunk in r.aiter_bytes(): + nbytes += len(chunk) + except httpx.TimeoutException: + err = "timeout" + except httpx.HTTPError as e: + err = f"http_error:{type(e).__name__}" + except Exception as e: + err = f"exc:{type(e).__name__}:{e}" + latency_ms = (time.perf_counter() - started) * 1000 + results.add( + Sample( + endpoint=ep, + url=url, + status=status, + latency_ms=latency_ms, + error=err, + bytes=nbytes, + started_at=wall, + ) + ) + + +async def run( + base_url: str, + rpm: int, + duration_s: int, + timeout_s: float, + weights: tuple[int, int, int], + out_dir: Path, +) -> Results: + interval = 60.0 / rpm + results = Results() + limits = httpx.Limits(max_connections=64, max_keepalive_connections=32) + async with httpx.AsyncClient(base_url=base_url, limits=limits) as client: + # Warm up / discover valid gages + print("[discover] fetching available gauge IDs...") + gages = await fetch_available_gages(client) + print(f"[discover] using {len(gages)} gauge IDs (sample: {gages[:5]})") + + start = time.time() + tasks: list[asyncio.Task] = [] + issued = 0 + next_fire = start + while time.time() - start < duration_s: + now = time.time() + if now >= next_fire: + ep, url = pick_request(gages, weights) + tasks.append(asyncio.create_task(one_request(client, ep, url, results, timeout_s))) + issued += 1 + next_fire += interval + if issued % 10 == 0: + elapsed = now - start + ok = sum(1 for s in results.samples if 200 <= s.status < 300) + print( + f"[{elapsed:5.0f}s] issued={issued} done={len(results.samples)} " + f"ok={ok} inflight={issued - len(results.samples)}" + ) + else: + await asyncio.sleep(min(0.05, next_fire - now)) + + # Drain + print(f"[drain] waiting for {issued - len(results.samples)} in-flight requests...") + if tasks: + await asyncio.gather(*tasks, return_exceptions=True) + + # Persist raw samples as CSV + out_dir.mkdir(parents=True, exist_ok=True) + csv_path = out_dir / "samples.csv" + with csv_path.open("w", newline="") as f: + w = csv.writer(f) + w.writerow(["started_at", "endpoint", "status", "latency_ms", "bytes", "error", "url"]) + for s in results.samples: + w.writerow( + [ + f"{s.started_at:.3f}", + s.endpoint, + s.status, + f"{s.latency_ms:.1f}", + s.bytes, + s.error or "", + s.url, + ] + ) + print(f"[io] wrote {csv_path}") + return results + + +def main() -> None: + p = argparse.ArgumentParser() + p.add_argument("--base-url", default="http://localhost:8000") + p.add_argument("--rpm", type=int, default=100, help="Requests per minute") + p.add_argument("--duration", type=int, default=300, help="Total test duration in seconds") + p.add_argument( + "--timeout", + type=float, + default=300, + help="Per-request timeout seconds (matches ICEFABRIC_HF_GPKG_QUEUE_TIMEOUT_S default).", + ) + p.add_argument( + "--weights", + default="1,2,2", + help=( + "Endpoint weights hf,param,stream. Hydrofabric gpkg is heavy (~130 MB / ~30s " + "per CONUS VPU) so we keep its share modest; a run at 100 rpm with 20% hf " + "weight still exercises ~20 concurrent heavy gpkg builds per minute." + ), + ) + p.add_argument("--out-dir", default="scripts/load_test/results") + args = p.parse_args() + weights = tuple(int(x) for x in args.weights.split(",")) + assert len(weights) == 3, "weights must be 3 ints" + + print(f"[cfg] base_url={args.base_url} rpm={args.rpm} duration={args.duration}s weights={weights}") + results = asyncio.run( + run( + base_url=args.base_url, + rpm=args.rpm, + duration_s=args.duration, + timeout_s=args.timeout, + weights=weights, # type: ignore[arg-type] + out_dir=Path(args.out_dir), + ) + ) + + summary = results.summary() + summary_path = Path(args.out_dir) / "summary.json" + summary_path.write_text(json.dumps(summary, indent=2)) + print("\n" + "=" * 60) + print("SUMMARY") + print("=" * 60) + print(json.dumps(summary, indent=2)) + print(f"\n[io] wrote {summary_path}") + + +if __name__ == "__main__": + main() diff --git a/scripts/load_test/monitor.sh b/scripts/load_test/monitor.sh new file mode 100755 index 0000000..9d5c026 --- /dev/null +++ b/scripts/load_test/monitor.sh @@ -0,0 +1,25 @@ +#!/usr/bin/env bash +# Stream `docker stats` into a CSV for later analysis. +# Usage: ./monitor.sh [interval_seconds] +set -euo pipefail + +CONTAINER="${1:-icefabric-api-loadtest}" +OUT="${2:-scripts/load_test/results/stats.csv}" +INTERVAL="${3:-2}" + +mkdir -p "$(dirname "$OUT")" +echo "timestamp,name,cpu_pct,mem_usage,mem_limit,mem_pct,net_io,block_io,pids" > "$OUT" + +while true; do + ts=$(date +%s) + # --no-stream gives one snapshot; we parse its non-header line + docker stats --no-stream --format '{{.Name}},{{.CPUPerc}},{{.MemUsage}},{{.MemPerc}},{{.NetIO}},{{.BlockIO}},{{.PIDs}}' \ + "$CONTAINER" 2>/dev/null | while IFS=, read -r name cpu mem mempct net block pids; do + # Split "123.4MiB / 8GiB" into usage,limit + usage=$(echo "$mem" | awk -F' / ' '{print $1}') + limit=$(echo "$mem" | awk -F' / ' '{print $2}') + printf '%s,%s,%s,%s,%s,%s,"%s","%s",%s\n' \ + "$ts" "$name" "$cpu" "$usage" "$limit" "$mempct" "$net" "$block" "$pids" >> "$OUT" + done || true + sleep "$INTERVAL" +done diff --git a/scripts/load_test/run.sh b/scripts/load_test/run.sh new file mode 100755 index 0000000..7540764 --- /dev/null +++ b/scripts/load_test/run.sh @@ -0,0 +1,95 @@ +#!/usr/bin/env bash +# Orchestrate: build → start container → wait healthy → monitor + load test → teardown. +set -euo pipefail + +HERE="$(cd "$(dirname "$0")" && pwd)" +cd "$HERE" + +RESULTS="results" +mkdir -p "$RESULTS" + +# DURATION, RPM (requests per minute), TIMEOUT, CPUS, and MEMORY can be set as environment variables, +# otherwise the defaults are used. +DURATION="${DURATION:-300}" # seconds +RPM="${RPM:-100}" +TIMEOUT="${TIMEOUT:-120}" # per-request seconds +CPUS="${CPUS:-2.0}" +MEMORY="${MEMORY:-8g}" + +echo "== config: rpm=$RPM duration=${DURATION}s per-req-timeout=${TIMEOUT}s cpus=$CPUS mem=$MEMORY ==" + +# Use docker compose v2 if available, else fall back to docker-compose v1. +if docker compose version >/dev/null 2>&1; then + DC="docker compose" +else + DC="docker-compose" +fi +echo "== using compose: $DC ==" + +# --- Build --- +echo "== building api image ==" +$DC -f docker-compose.load.yml build api + +# --- Start --- +echo "== starting container ==" +$DC -f docker-compose.load.yml up -d api + +# Apply runtime caps (compose v2 ignores top-level cpus/mem_limit without swarm; +# make it explicit with docker update after start). +docker update --cpus="$CPUS" --memory="$MEMORY" --memory-swap="$MEMORY" \ + icefabric-api-loadtest >/dev/null +echo "== applied --cpus=$CPUS --memory=$MEMORY ==" + +# --- Wait for healthy --- +echo "== waiting for /health (up to 10 min, cache build takes time) ==" +for i in $(seq 1 120); do + if curl -fsS --max-time 3 http://localhost:8000/health >/dev/null 2>&1; then + echo "== api is healthy after $((i*5))s ==" + break + fi + sleep 5 + if [[ $((i % 12)) -eq 0 ]]; then + echo " still waiting... ($((i*5))s elapsed)" + fi +done +if ! curl -fsS --max-time 3 http://localhost:8000/health >/dev/null 2>&1; then + echo "!! api never became healthy. Last 80 lines of container logs:" + docker logs --tail 80 icefabric-api-loadtest || true + exit 1 +fi + +# --- Monitor --- +echo "== starting docker-stats monitor ==" +bash monitor.sh icefabric-api-loadtest "$RESULTS/stats.csv" 2 & +MON_PID=$! +trap 'kill $MON_PID 2>/dev/null || true; $DC -f docker-compose.load.yml logs api > "$RESULTS/container.log" 2>&1 || true' EXIT + +# --- Load test --- +echo "== running load test ==" +python3 load_test.py \ + --base-url http://localhost:8000 \ + --rpm "$RPM" \ + --duration "$DURATION" \ + --timeout "$TIMEOUT" \ + --out-dir "$RESULTS" + +# --- Stop monitor, dump logs, analyze --- +kill $MON_PID 2>/dev/null || true +sleep 2 +$DC -f docker-compose.load.yml logs api > "$RESULTS/container.log" 2>&1 || true + +echo "" +echo "== resource analysis ==" +python3 analyze.py --stats "$RESULTS/stats.csv" + +echo "" +echo "== container state ==" +docker inspect icefabric-api-loadtest \ + --format '{{.State.Status}} OOMKilled={{.State.OOMKilled}} ExitCode={{.State.ExitCode}} RestartCount={{.RestartCount}}' + +echo "" +echo "== grep for errors / OOM / tracebacks in container log ==" +grep -i -E "oom|killed|memoryerror|traceback|error" "$RESULTS/container.log" | head -30 || true + +echo "" +echo "results in: $HERE/$RESULTS"