Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions app/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
1 change: 1 addition & 0 deletions app/routers/hydrofabric/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
import pathlib
import sqlite3
import tempfile
import time
import uuid

import geopandas as gpd
Expand Down
1 change: 0 additions & 1 deletion app/routers/streamflow_observations/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,7 @@ convention = "numpy"
[tool.ruff.lint.per-file-ignores]
"docs/*" = ["I"]
"tests/*" = ["D"]
"scripts/*" = ["D", "BLE001"]
"*/__init__.py" = ["F401"]

[tool.mypy]
Expand Down
86 changes: 86 additions & 0 deletions scripts/load_test/analyze.py
Original file line number Diff line number Diff line change
@@ -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()
35 changes: 35 additions & 0 deletions scripts/load_test/docker-compose.load.yml
Original file line number Diff line number Diff line change
@@ -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
Loading
Loading