services/semif-serve is a FastAPI wrapper around SemIf's direct and shared torch scorers (SemIf-OpenJev @ 23cf1f39, MIT). Upstream ships only a batch CLI. The wrapper loads the pinned Qwen3.5-4B (851bf6e8, BF16) once from the offline HF cache and returns SemIf's result dicts unchanged, with an optional per-workload temperature-calibrated view. Contract: semif-serve.contract.md. Built with a short contract, TDD (39 tests, fake engine and fake torch, no GPU) and a heid bug-hunt panel (pending). On the card: - torch 2.10.0+cu128 with sm_120 kernels, which is SemIf's own stack; - a hard 12 GiB VRAM cap. Two defects surfaced only on the card, and each fix is covered by a test: - 0.1.1: an OOM raised as a chained exception kept the failed request's tensors alive (11.9 GiB after the 503). It is now raised unchained, after gc. - 0.1.2: a large request left 12.6 GB reserved on the shared card. After each call, reserved memory over the baseline + 512 MiB is now released. Acceptance against SemIf's committed torch predictions (authored144): - 142/144 same top choice; both misses are exact bf16 ties; - 144/144 identical prompt hashes; - deterministic A-vs-A; - negative control 14/144; - shared vs direct 72/72. 21 binary criteria over one state take 159 ms. The shared-mode capacity table under the cap is in stacks/semif/README.md. The Dockerfile installs dependencies from a manifest with the project version blanked, so a version bump reuses the ~4 GB torch layer. Verified: 41 s rebuild, dependency layer CACHED. DNS: semif.fv.internal. Token: vault semif/api-token.
234 lines
9.9 KiB
Python
234 lines
9.9 KiB
Python
"""semif-serve HTTP behaviour against a fake engine (no torch, no model).
|
|
Contract: services/semif-serve/semif-serve.contract.md"""
|
|
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"]}
|