Promote dev to main (deploy beta): ERA5 lake stack #26
3 changed files with 32 additions and 7 deletions
|
|
@ -26,8 +26,8 @@ RUN chmod +x /app/deploy/entrypoint.sh
|
||||||
# Docker seeds them from this image content -- including this ownership --
|
# Docker seeds them from this image content -- including this ownership --
|
||||||
# so they stay writable.
|
# so they stay writable.
|
||||||
RUN useradd --system --create-home --uid 10001 thermograph \
|
RUN useradd --system --create-home --uid 10001 thermograph \
|
||||||
&& mkdir -p /app/data /app/logs \
|
&& mkdir -p /app/data /app/logs /state \
|
||||||
&& chown -R thermograph:thermograph /app
|
&& chown -R thermograph:thermograph /app /state
|
||||||
|
|
||||||
USER thermograph
|
USER thermograph
|
||||||
WORKDIR /app
|
WORKDIR /app
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,7 @@ Scale-out safe: replicas share nothing — each keeps its own cache under
|
||||||
THERMOGRAPH_LAKE_CACHE (a per-task local volume), so 1..2 replicas behind the
|
THERMOGRAPH_LAKE_CACHE (a per-task local volume), so 1..2 replicas behind the
|
||||||
Swarm VIP need no coordination.
|
Swarm VIP need no coordination.
|
||||||
"""
|
"""
|
||||||
|
import io
|
||||||
import os
|
import os
|
||||||
import re
|
import re
|
||||||
import threading
|
import threading
|
||||||
|
|
@ -31,7 +32,7 @@ import time
|
||||||
|
|
||||||
import polars as pl
|
import polars as pl
|
||||||
from fastapi import FastAPI, HTTPException
|
from fastapi import FastAPI, HTTPException
|
||||||
from fastapi.responses import FileResponse
|
from fastapi.responses import FileResponse, Response
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
|
|
||||||
from data import era5lake
|
from data import era5lake
|
||||||
|
|
@ -106,10 +107,20 @@ def history(lat: float, lon: float):
|
||||||
df = era5lake._read_bucket(i, j, cfg)
|
df = era5lake._read_bucket(i, j, cfg)
|
||||||
except era5lake.LakeUnavailable as e:
|
except era5lake.LakeUnavailable as e:
|
||||||
raise HTTPException(status_code=404, detail=str(e)) from e
|
raise HTTPException(status_code=404, detail=str(e)) from e
|
||||||
os.makedirs(CACHE_DIR, exist_ok=True)
|
try:
|
||||||
tmp = f"{path}.{os.getpid()}.partial"
|
os.makedirs(CACHE_DIR, exist_ok=True)
|
||||||
df.write_parquet(tmp)
|
tmp = f"{path}.{os.getpid()}.partial"
|
||||||
os.replace(tmp, path) # atomic: concurrent requests see whole files only
|
df.write_parquet(tmp)
|
||||||
|
os.replace(tmp, path) # atomic: concurrent readers see whole files only
|
||||||
|
except OSError:
|
||||||
|
# An unwritable cache (bad volume ownership, full disk) must never
|
||||||
|
# fail the hot path — serve from memory; the cache is an
|
||||||
|
# accelerator, not a dependency.
|
||||||
|
buf = io.BytesIO()
|
||||||
|
df.write_parquet(buf)
|
||||||
|
return Response(buf.getvalue(),
|
||||||
|
media_type="application/vnd.apache.parquet",
|
||||||
|
headers={"X-Lake-Point": f"{i},{j}"})
|
||||||
return FileResponse(path, media_type="application/vnd.apache.parquet",
|
return FileResponse(path, media_type="application/vnd.apache.parquet",
|
||||||
headers={"X-Lake-Point": f"{i},{j}"})
|
headers={"X-Lake-Point": f"{i},{j}"})
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -69,6 +69,20 @@ def test_history_serves_parquet_and_caches(lake, tmp_path):
|
||||||
params={"lat": 51.5074, "lon": -0.1278}).content == r.content
|
params={"lat": 51.5074, "lon": -0.1278}).content == r.content
|
||||||
|
|
||||||
|
|
||||||
|
def test_history_serves_from_memory_when_cache_unwritable(lake, tmp_path, monkeypatch):
|
||||||
|
"""An unwritable cache dir (bad volume ownership — seen live) degrades to
|
||||||
|
serving from memory, never a 500."""
|
||||||
|
import lake_app
|
||||||
|
ro = tmp_path / "ro-cache" / "nested" # parent made read-only below
|
||||||
|
(tmp_path / "ro-cache").mkdir()
|
||||||
|
(tmp_path / "ro-cache").chmod(0o555)
|
||||||
|
monkeypatch.setattr(lake_app, "CACHE_DIR", str(ro))
|
||||||
|
r = lake.get("/history", params={"lat": 51.5074, "lon": -0.1278})
|
||||||
|
(tmp_path / "ro-cache").chmod(0o755)
|
||||||
|
assert r.status_code == 200
|
||||||
|
assert r.headers["x-lake-point"] == "154,1439"
|
||||||
|
|
||||||
|
|
||||||
def test_history_404_for_point_outside_lake(lake):
|
def test_history_404_for_point_outside_lake(lake):
|
||||||
r = lake.get("/history", params={"lat": -45.0, "lon": 170.0})
|
r = lake.get("/history", params={"lat": -45.0, "lon": 170.0})
|
||||||
assert r.status_code == 404
|
assert r.status_code == 404
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue