Files
esh-pfi-infrastructure/services/semif-serve/tests/test_app.py
T
vh 77b8cb449c feat(semif): 0.1.3 — order averaging, fast kernels, bug-hunt hardening (Prime)
Order averaging (Prime, after the 739aa03 spike):
- A decision may set orderings: rotations|all (all only for <= 4 options). Every
  ordering goes to the engine in one shared batch.
- The reply keeps each native result and adds combined {probabilities (log-mean),
  top, agreement, spread}.
- Through the service on SemIf's labelled sets (252 rows): 78.6% -> 88.1%
  (group-bootstrap 95% CI +5.1..+14.3). Unanimous agreement is 94.5% accurate.

Fast kernels: flash-linear-attention 0.5.2 and causal-conv1d 1.7.0 are now the
default build. A/B on the empty GPU 3:
- parity with upstream went from 142/144 to 144/144;
- a ~2k-token /decide went from 169 to 92 ms server-side;
- short 3-rotation batches cost ~3-6 ms more.
triton builds a C shim at runtime, so the image carries gcc. Without it the
warm-up failed and startup failed closed.

Heid bug-hunt panel (4/4 arms, thread 01M3H3F4RR7XBP90KQ3A39H4SX), folded:
- Startup validation: VRAM cap 0 no longer means uncapped (C1); limits must be
  >= 1 (S1); the token must be visible ASCII (S2); the calibration file must
  exist and parse, with T in [0.05, 20] (S8, and C3's NaN leg).
- The body limit is checked before a chunk is kept, and a Unicode-digit
  Content-Length no longer crashes (C2, S3).
- Failures while building the response now get the 500 envelope (C3).
- 429 busy past SEMIF_MAX_QUEUE requests in progress (C6).
- The engine releases memory on every non-validation failure, unchained after
  gc; an empty OOM message is handled; 'out of memory' RuntimeErrors map to 503
  (C4, C5, S9).
- The entry point forces HF_HUB_OFFLINE (S10). README wording fixed (S5, S6).
- New guard tests close the gaps the arms' mutation grids exposed: early stop of
  the body read, a shared-route lock, calibration pass-through, the gc cycle,
  the exact caps, TorchEngine.load's arch and device checks, and the offline
  entry point.
86 tests.

Deployed on fv-ml1 GPU 1: parity 144/144, OOM and burst release verified, shared
capacity 63/51/26/16 rows at ~140/520/1960/3900 prefix tokens.
2026-09-27 03:27:15 -07:00

337 lines
14 KiB
Python

"""semif-serve HTTP behaviour against a fake engine (no torch, no model).
Contract: services/semif-serve/semif-serve.contract.md"""
import json
import pytest
from fastapi.testclient import TestClient
from semif_serve.app import create_app
from semif_serve.config import Settings
from semif_serve.errors import OutOfMemory
TOKEN = "t" * 40
AUTH = {"Authorization": f"Bearer {TOKEN}"}
OPTIONS = [{"id": "yes", "description": "Yes."}, {"id": "no", "description": "No."}]
ROW = {"id": "r1", "state": "The deploy passed.", "question": "Did it pass?", "options": OPTIONS}
class FakeEngine:
"""Returns canned SemIf-shaped dicts and records what it was asked."""
def __init__(self, logits=(2.0, 0.0)):
self.logits = list(logits)
self.direct_calls = []
self.shared_calls = []
def health(self):
return {"source": "fake/model", "revision": "0" * 40}
def direct(self, row):
self.direct_calls.append(row)
return {
"id": row["id"],
"option_ids": [o["id"] for o in row["options"]],
"probabilities": [0.8807970779778823, 0.11920292202211769],
"option_logits": self.logits,
"prompt_sha256": "ab" * 32,
"prompt_version": "direct-options-v1",
"model": {"source": "fake/model"},
"probability_status": "conditional option score; uncalibrated as decision confidence",
}
def shared(self, rows):
self.shared_calls.append(rows)
return [self.direct(r) for r in rows], {"prefix_tokens": 7, "batch_size": len(rows)}
def make_client(engine=None, **overrides):
settings = Settings(api_token=TOKEN, **overrides)
return TestClient(create_app(settings, engine or FakeEngine()))
def test_decide_passes_the_row_through_and_returns_the_scorer_dict_unchanged():
engine = FakeEngine()
client = make_client(engine)
response = client.post("/decide", json=ROW, headers=AUTH)
assert response.status_code == 200
assert response.json() == engine.direct(ROW)
assert engine.direct_calls[0] == ROW
@pytest.mark.parametrize("headers", [{}, {"Authorization": "Bearer wrong"}, {"Authorization": TOKEN}])
def test_posts_without_the_right_bearer_are_401_and_never_reach_the_engine(headers):
engine = FakeEngine()
client = make_client(engine)
for path in ("/decide", "/decide/shared"):
response = client.post(path, json=ROW, headers=headers)
assert response.status_code == 401
assert response.json() == {"error": {"code": "unauthorized", "message": response.json()["error"]["message"]}}
assert engine.direct_calls == []
def test_health_needs_no_auth():
response = make_client().get("/health")
assert response.status_code == 200
assert response.json()["status"] == "ok"
def test_shared_builds_semif_rows_from_the_shared_state_and_returns_results_and_timing():
engine = FakeEngine()
body = {"state": {"deploy": "passed"},
"decisions": [{"id": "a", "question": "Q1?", "options": OPTIONS},
{"id": "b", "question": "Q2?", "options": OPTIONS}]}
response = make_client(engine).post("/decide/shared", json=body, headers=AUTH)
assert response.status_code == 200
assert engine.shared_calls == [[
{"id": "a", "state": {"deploy": "passed"}, "question": "Q1?", "options": OPTIONS},
{"id": "b", "state": {"deploy": "passed"}, "question": "Q2?", "options": OPTIONS},
]]
out = response.json()
assert [r["id"] for r in out["results"]] == ["a", "b"]
assert out["timing"] == {"prefix_tokens": 7, "batch_size": 2}
def test_a_known_workload_adds_a_calibrated_view_and_leaves_the_native_scores_alone():
import math
engine = FakeEngine(logits=(3.0, 1.0))
client = make_client(engine, calibration={"triage": 2.0})
native = engine.direct(ROW)
for path, body, pick in (
("/decide", {**ROW, "workload": "triage"}, lambda j: [j]),
("/decide/shared", {"state": ROW["state"], "workload": "triage",
"decisions": [{"id": "a", "question": "Q?", "options": OPTIONS}]}, lambda j: j["results"]),
):
for result in pick(client.post(path, json=body, headers=AUTH).json()):
cal = result.pop("calibrated")
assert (cal["workload"], cal["temperature"]) == ("triage", 2.0)
e = [math.exp(1.5), math.exp(0.5)]
assert cal["probabilities"] == pytest.approx([e[0] / sum(e), e[1] / sum(e)])
assert cal["probabilities"].index(max(cal["probabilities"])) == 0
assert {k: result[k] for k in ("probabilities", "option_logits", "probability_status")} == \
{k: native[k] for k in ("probabilities", "option_logits", "probability_status")}
def test_no_workload_means_no_calibrated_key():
client = make_client(calibration={"triage": 2.0})
assert "calibrated" not in client.post("/decide", json=ROW, headers=AUTH).json()
def test_an_unknown_workload_is_422_before_the_engine_runs():
engine = FakeEngine()
client = make_client(engine, calibration={"triage": 2.0})
response = client.post("/decide", json={**ROW, "workload": "nope"}, headers=AUTH)
assert response.status_code == 422
assert response.json()["error"]["code"] == "invalid_request"
assert engine.direct_calls == []
class RaisingEngine(FakeEngine):
def __init__(self, exc):
super().__init__()
self.exc = exc
def direct(self, row):
raise self.exc
def shared(self, rows):
raise self.exc
@pytest.mark.parametrize("exc, status, code", [
(ValueError("Row r1: 5000 input tokens exceed limit 4096; no truncation allowed"), 422, "invalid_request"),
(OutOfMemory("CUDA out of memory"), 503, "out_of_memory"),
(RuntimeError("Invalid native prefix cache"), 500, "scoring_failed"),
])
def test_scorer_failures_map_to_the_contract_status_codes(exc, status, code):
client = make_client(RaisingEngine(exc))
for path, body in (("/decide", ROW),
("/decide/shared", {"state": "s", "decisions": [{"id": "a", "question": "Q?", "options": OPTIONS}]})):
response = client.post(path, json=body, headers=AUTH)
assert response.status_code == status
assert response.json()["error"]["code"] == code
assert str(exc) in response.json()["error"]["message"]
@pytest.mark.parametrize("content", [b"{not json", b'{"id": "r1"}', b'{"state": "s", "decisions": "nope"}', b"[1, 2]"])
def test_malformed_json_or_a_wrong_shape_is_422(content):
client = make_client()
for path in ("/decide", "/decide/shared"):
response = client.post(path, content=content, headers={**AUTH, "content-type": "application/json"})
assert response.status_code == 422
assert response.json()["error"]["code"] == "invalid_request"
def test_more_decisions_than_the_cap_is_422_and_the_engine_never_runs():
engine = FakeEngine()
client = make_client(engine, max_decisions=2)
decisions = [{"id": str(i), "question": "Q?", "options": OPTIONS} for i in range(3)]
response = client.post("/decide/shared", json={"state": "s", "decisions": decisions}, headers=AUTH)
assert response.status_code == 422
assert response.json()["error"]["code"] == "invalid_request"
assert engine.shared_calls == []
def test_an_empty_decision_list_is_422():
response = make_client().post("/decide/shared", json={"state": "s", "decisions": []}, headers=AUTH)
assert response.status_code == 422
@pytest.mark.parametrize("chunked", [False, True])
def test_a_body_over_the_limit_is_413_whether_or_not_it_declares_its_length(chunked):
engine = FakeEngine()
client = make_client(engine, max_body_bytes=200)
body = ('{"id": "r1", "state": "' + "x" * 500 + '", "question": "Q?", "options": []}').encode()
content = (chunk for chunk in [body[:100], body[100:]]) if chunked else body
response = client.post("/decide", content=content, headers={**AUTH, "content-type": "application/json"})
assert response.status_code == 413
assert response.json()["error"]["code"] == "request_too_large"
assert engine.direct_calls == []
class SlowEngine(FakeEngine):
"""Holds each scorer call until released, and records the peak number of calls inside at once."""
def __init__(self):
super().__init__()
import threading
self.inside = 0
self.peak = 0
self.guard = threading.Lock()
self.release = threading.Event()
self.entered = threading.Event()
def direct(self, row):
with self.guard:
self.inside += 1
self.peak = max(self.peak, self.inside)
self.entered.set()
self.release.wait(5)
with self.guard:
self.inside -= 1
return super().direct(row)
def test_concurrent_requests_never_overlap_inside_the_scorer_and_health_still_answers():
from concurrent.futures import ThreadPoolExecutor
engine = SlowEngine()
with make_client(engine) as client, ThreadPoolExecutor(4) as pool:
futures = [pool.submit(client.post, "/decide", json={**ROW, "id": f"r{i}"}, headers=AUTH) for i in range(4)]
assert engine.entered.wait(5)
import time
started = time.monotonic()
assert client.get("/health").status_code == 200
assert time.monotonic() - started < 1.0 # answered while the scorer is still held
assert not engine.release.is_set() and engine.inside == 1
engine.release.set()
assert [f.result().status_code for f in futures] == [200] * 4
assert engine.peak == 1
def test_health_reports_the_pins_limits_and_workloads():
from semif_serve.config import SEMIF_COMMIT
client = make_client(vram_cap_gib=12.0, max_decisions=8, calibration={"triage": 2.0, "alerts": 1.3})
body = client.get("/health").json()
assert body == {"status": "ok", "semif_commit": SEMIF_COMMIT, "model": FakeEngine().health(),
"vram_cap_gib": 12.0, "max_tokens": 4096, "max_decisions": 8, "workloads": ["alerts", "triage"]}
def test_a_body_of_exactly_the_limit_is_accepted():
body = json.dumps(ROW).encode()
client = make_client(max_body_bytes=len(body))
assert client.post("/decide", content=body, headers={**AUTH, "content-type": "application/json"}).status_code == 200
def test_a_non_ascii_digit_content_length_is_ignored_not_a_crash():
"""HTTP clients cannot send one (httpx refuses; h11 rejects it), so check the reader directly."""
import asyncio
from semif_serve.app import read_limited
async def body():
yield b"{}"
assert asyncio.run(read_limited(body(), "²", 100)) == b"{}" # int("²") would raise
def test_exactly_max_decisions_is_accepted():
decisions = [{"id": str(i), "question": "Q?", "options": OPTIONS} for i in range(3)]
client = make_client(max_decisions=3)
assert client.post("/decide/shared", json={"state": "s", "decisions": decisions}, headers=AUTH).status_code == 200
def test_calibration_leaves_every_native_field_alone():
engine = FakeEngine(logits=(3.0, 1.0))
body = make_client(engine, calibration={"triage": 2.0}).post(
"/decide", json={**ROW, "workload": "triage"}, headers=AUTH).json()
body.pop("calibrated")
assert body == engine.direct(ROW)
class MalformedEngine(FakeEngine):
def direct(self, row):
return {"id": row["id"], "option_ids": ["yes", "no"], "probabilities": [0.5, 0.5]} # no option_logits
def shared(self, rows):
return [self.direct(r) for r in rows], {}
@pytest.mark.parametrize("path, body", [
("/decide", {**ROW, "workload": "triage"}),
("/decide", {**ROW, "orderings": "rotations"}),
])
def test_a_malformed_scorer_result_is_an_envelope_500_not_a_bare_one(path, body):
client = make_client(MalformedEngine(), calibration={"triage": 2.0})
response = client.post(path, json=body, headers=AUTH)
assert response.status_code == 500
assert response.json()["error"]["code"] == "scoring_failed"
class SharedSlowEngine(SlowEngine):
def shared(self, rows):
return [self.direct(r) for r in rows], {}
def test_shared_requests_are_serialised_too():
from concurrent.futures import ThreadPoolExecutor
engine = SharedSlowEngine()
body = {"state": "s", "decisions": [{"id": "a", "question": "Q?", "options": OPTIONS}]}
with make_client(engine) as client, ThreadPoolExecutor(3) as pool:
futures = [pool.submit(client.post, "/decide/shared", json=body, headers=AUTH) for _ in range(3)]
assert engine.entered.wait(5)
engine.release.set()
assert [f.result().status_code for f in futures] == [200] * 3
assert engine.peak == 1
def test_more_than_max_queue_requests_in_progress_get_429_busy():
from concurrent.futures import ThreadPoolExecutor
engine = SlowEngine()
with make_client(engine, max_queue=2) as client, ThreadPoolExecutor(3) as pool:
held = [pool.submit(client.post, "/decide", json={**ROW, "id": f"r{i}"}, headers=AUTH) for i in range(2)]
assert engine.entered.wait(5)
import time
deadline = time.monotonic() + 5
while engine.inside + 0 < 1 and time.monotonic() < deadline:
time.sleep(0.01)
time.sleep(0.2) # let the second request reach the queue
extra = client.post("/decide", json={**ROW, "id": "extra"}, headers=AUTH)
assert extra.status_code == 429 and extra.json()["error"]["code"] == "busy"
engine.release.set()
assert [f.result().status_code for f in held] == [200, 200]
assert make_client(FakeEngine(), max_queue=2).post("/decide", json=ROW, headers=AUTH).status_code == 200
def test_read_limited_stops_reading_at_the_crossing_chunk():
import asyncio
from semif_serve.app import ApiError, read_limited
consumed = []
async def chunks():
for i in range(10):
consumed.append(i)
yield b"x" * 100
with pytest.raises(ApiError) as info:
asyncio.run(read_limited(chunks(), None, 250))
assert info.value.status == 413
assert consumed == [0, 1, 2] # the third chunk crosses 250 and is never kept; nothing after is read