Files
stroblmeandClaude Opus 5 51464941ac
Docs / docs (push) Successful in 35s
Playwright Tests / test-playwright (1, 2) (push) Successful in 3m11s
Playwright Tests / test-playwright (2, 2) (push) Successful in 2m17s
pre-commit / pre-commit (push) Failing after 2m44s
Test Backend / test-backend (push) Successful in 3m0s
Compose Smoke Test / test-compose (push) Successful in 41s
Playwright Tests / merge-reports (push) Successful in 8m14s
Export runs and their curves as tables an analysis reads
`fluksio export metrics` is the long table — a row per run, metric and step —
and `fluksio export runs` the wide one, a row per run with the inputs that
*vary* across the selection as columns beside its final numbers, status,
duration and the commit and digest of the code it ran. Both carry the run id
on every row, which is the join back to the run page and what makes an
exported file auditable. `Client.export_metrics`/`export_runs` answer the same
rows to a notebook.

The engine streams csv or jsonl from two routes declared above `/{run_id}`;
parquet is a client-side conversion behind the new `fluksio[parquet]` extra,
so nobody pays for pyarrow who does not want dtypes kept. The long export
reads each run through `_series`, so a cached node's curve comes with it, and
`--stride` thins each series rather than the concatenation of all of them.

Two things they needed on the way: `GET /runs` takes `?since=` and `?before=`,
so a long history pages by the last row's own timestamp instead of an offset
that shifts under it; and a read that reaches no engine now says so in half a
second rather than seven, because `runs`, `flavors`, `export` and an unwatched
`status` pass `retries=0`. Everything that submits keeps them.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01A9Hdrmf2cwNABCnE5x9UJa
2026-08-27 17:43:30 +02:00

155 lines
5.0 KiB
Python

"""What the client does when the engine is slow, busy, or briefly gone."""
import httpx
import pytest
from fluksio.sdk.client import WAIT_TOLERANCE, ApiError, Client, RunHandle
def a_client(handler, monkeypatch, **kwargs):
"""A client whose transport is a function, and whose backoff costs nothing."""
monkeypatch.setattr("fluksio.sdk.client.time.sleep", lambda _seconds: None)
transport = httpx.MockTransport(handler)
http = httpx.Client(transport=transport, base_url="http://engine")
return Client(url="http://engine", token="t", http=http, **kwargs)
def test_idempotent_get_is_tried_again(monkeypatch):
calls = []
def handler(request):
calls.append(request)
if len(calls) < 3:
raise httpx.ReadTimeout("too slow", request=request)
return httpx.Response(200, json={"id": "r1", "status": "ok"})
client = a_client(handler, monkeypatch)
assert client.run("r1")["status"] == "ok"
assert len(calls) == 3
def test_a_busy_engine_is_asked_again(monkeypatch):
codes = iter([503, 503, 200])
def handler(request):
code = next(codes)
return httpx.Response(code, json={"id": "r1"} if code == 200 else {})
client = a_client(handler, monkeypatch)
assert client.run("r1")["id"] == "r1"
def test_giving_up_raises_what_it_last_saw(monkeypatch):
def handler(request):
raise httpx.ReadTimeout("too slow", request=request)
client = a_client(handler, monkeypatch, retries=2)
with pytest.raises(httpx.ReadTimeout):
client.run("r1")
def test_a_write_is_not_repeated(monkeypatch):
calls = []
def handler(request):
calls.append(request)
raise httpx.ReadTimeout("too slow", request=request)
client = a_client(handler, monkeypatch)
with pytest.raises(httpx.ReadTimeout):
client.put_source("train", "fit", "code")
assert len(calls) == 1, "a source write means something different twice"
def test_submit_carries_one_key_across_its_retries(monkeypatch):
bodies = []
def handler(request):
bodies.append(httpx.Response(200, content=request.content).json())
if len(bodies) < 3:
raise httpx.ConnectError("no route", request=request)
return httpx.Response(202, json={"id": "r1", "status": "queued"})
client = a_client(handler, monkeypatch)
assert client.submit("train", {"lr": 0.1}).id == "r1"
keys = {body["idempotency_key"] for body in bodies}
assert len(keys) == 1, "a retry must not read as a second run"
assert len(next(iter(keys))) == 32
def test_a_sweep_keys_every_entry(monkeypatch):
seen = {}
def handler(request):
body = httpx.Response(200, content=request.content).json()
seen["runs"] = body["runs"]
return httpx.Response(202, json=[{"id": "r1"}, {"id": "r2"}])
client = a_client(handler, monkeypatch)
client.sweep("train", [{"params": {"lr": 0.1}}, {"params": {"lr": 0.2}}])
keys = [entry["idempotency_key"] for entry in seen["runs"]]
assert len(set(keys)) == 2
def test_waiting_survives_a_few_bad_answers(monkeypatch):
answers = iter(
[503] * (WAIT_TOLERANCE - 1) + [200] # then the run is finished
)
def handler(request):
code = next(answers)
if code != 200:
return httpx.Response(code, json={})
return httpx.Response(200, json={"id": "r1", "status": "ok"})
client = a_client(handler, monkeypatch, retries=0)
handle = RunHandle(client, "r1", {"status": "running"})
assert handle.wait(poll=0).status == "ok"
def test_waiting_gives_up_eventually(monkeypatch):
def handler(request):
return httpx.Response(503, json={})
client = a_client(handler, monkeypatch, retries=0)
handle = RunHandle(client, "r1", {"status": "running"})
with pytest.raises(ApiError):
handle.wait(poll=0)
def test_a_run_that_is_gone_stops_the_wait_at_once(monkeypatch):
calls = []
def handler(request):
calls.append(request)
return httpx.Response(404, json={"detail": "no such run"})
client = a_client(handler, monkeypatch, retries=0)
handle = RunHandle(client, "r1", {"status": "running"})
with pytest.raises(ApiError) as caught:
handle.wait(poll=0)
assert caught.value.status == 404
assert len(calls) == 1, "a 404 is an answer, not a blip"
def test_an_export_is_read_as_rows(monkeypatch):
"""jsonl on the wire: a stream of documents, not one."""
seen = {}
def handler(request):
seen.update(request.url.params)
return httpx.Response(
200,
text='{"run": "r1", "step": 0}\n{"run": "r1", "step": 1}\n',
headers={"content-type": "application/x-ndjson"},
)
client = a_client(handler, monkeypatch)
rows = client.export_metrics(flow="train", ids=["r1", "r2"], status="ok")
assert seen["format"] == "jsonl"
assert seen["ids"] == "r1,r2"
assert seen["flow"] == "train"
assert seen["status"] == "ok"
assert [row["step"] for row in rows] == [0, 1]