feat(intern-decision-serve): 0.1.1 adds POST /v1/systemone (Jev wire shape)
Straight passthrough to the checkpoint's own DecisionEngine.predict — never the semif mapping, whose different prompt would change the answers. Reuses the one inference thread, bearer auth, MAX_QUEUE, VRAM cap and error envelope; no new concurrency. 1..16 questions in ONE call (never chunked: Jev questions share a prompt); images 422; over MAX_TOKENS 422 before the forward. /health advertises the surface. Response 'model' is a string name@revision (JevBench's runner hashes it; a dict broke its manifest step). Acceptance on the live service (see acceptance/systemone-2026-09-30/): JevBench v1.2.16 typesafe adapter over the 231 public items scores all 202/231, hard 83/111, with 0 changed answers across all 924 rows of the bench's own r1..r4; controls 401/422x3 (token boundary proven at 7168 pass / 7169 refuse); GPU 1 per-process peak 9,866 MiB under the largest accepted request (budget 9,876); /decide/shared positive control unchanged. 120 tests green.
This commit is contained in:
@@ -22,6 +22,32 @@ def score_by_description(field: str, question: dict) -> list[float]:
|
||||
return [len(d) + (0.5 if i == 0 else 0.0) for i, d in enumerate(descs)]
|
||||
|
||||
|
||||
def answer_for(field: str, question: dict, scorer) -> dict:
|
||||
"""One Jev answer in the shape DecisionEngine.predict returns, for any of the three types."""
|
||||
kind = question["type"]
|
||||
if kind == "noul":
|
||||
probs = {"yes": 0.75, "no": 0.25}
|
||||
return {"type": "noul", "probabilities": probs, "confidence": 0.75, "noul": 0.75,
|
||||
"source": "local", "decision": "yes"}
|
||||
criteria = question["criteria"]
|
||||
if isinstance(criteria, dict):
|
||||
keys = list(criteria)
|
||||
descriptions = list(criteria.values())
|
||||
else: # score: a list of levels
|
||||
keys = [str(i) for i in range(len(criteria))]
|
||||
descriptions = [str(level) for level in criteria]
|
||||
probs = dict(zip(keys, softmax(scorer(field, {**question, "criteria": dict(zip(keys, descriptions))}))))
|
||||
best = min(keys, key=lambda k: (-probs[k], k))
|
||||
answer = {"type": kind, "probabilities": probs, "confidence": probs[best], "source": "local", "decision": best}
|
||||
if kind == "choice":
|
||||
answer["choice"] = best
|
||||
else:
|
||||
answer["score"] = sum(float(k) * p for k, p in probs.items())
|
||||
answer["legend"] = {k: (criteria[i] if isinstance(criteria, list) else v)
|
||||
for i, (k, v) in enumerate(zip(keys, criteria.values() if isinstance(criteria, dict) else criteria))}
|
||||
return answer
|
||||
|
||||
|
||||
class FakeEngine:
|
||||
def __init__(self, scorer=score_by_description, tokens_per_question: int = 100):
|
||||
self.scorer = scorer
|
||||
@@ -35,13 +61,8 @@ class FakeEngine:
|
||||
|
||||
def predict(self, request: dict) -> tuple[dict, str]:
|
||||
self.calls.append(copy.deepcopy(request))
|
||||
answers = {}
|
||||
for field, question in request["questions"].items():
|
||||
ids = list(question["criteria"])
|
||||
probs = dict(zip(ids, softmax(self.scorer(field, question))))
|
||||
best = min(ids, key=lambda i: (-probs[i], i))
|
||||
answers[field] = {"type": "choice", "probabilities": probs, "confidence": probs[best],
|
||||
"choice": best, "source": "local", "decision": best}
|
||||
answers = {field: answer_for(field, question, self.scorer)
|
||||
for field, question in request["questions"].items()}
|
||||
response = {"answers": answers,
|
||||
"usage": {"input_tokens": self.tokens_per_question * len(answers),
|
||||
"output_tokens": len(answers), "decision_count": len(answers)},
|
||||
|
||||
@@ -0,0 +1,224 @@
|
||||
"""POST /v1/systemone: Jev wire shape, straight passthrough to DecisionEngine.predict.
|
||||
Contract: intern-decision-serve.contract.md § POST /v1/systemone"""
|
||||
import threading
|
||||
import time
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
|
||||
import pytest
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from fake_engine import FakeEngine
|
||||
from intern_decision_serve.app import create_app
|
||||
from intern_decision_serve.config import Settings
|
||||
from intern_decision_serve.errors import OutOfMemory
|
||||
|
||||
TOKEN = "t" * 40
|
||||
AUTH = {"Authorization": f"Bearer {TOKEN}"}
|
||||
|
||||
CHOICE = {"type": "choice", "instructions": "Did it pass?",
|
||||
"criteria": {"yes": "It passed.", "no": "It did not pass at all."}}
|
||||
NOUL = {"type": "noul", "instructions": "Is the deploy healthy?"}
|
||||
SCORE = {"type": "score", "instructions": "How risky?", "criteria": ["low", "medium", "high"]}
|
||||
|
||||
|
||||
def make_client(engine=None, **overrides):
|
||||
return TestClient(create_app(Settings(api_token=TOKEN, **overrides), engine or FakeEngine()))
|
||||
|
||||
|
||||
def body(**over):
|
||||
b = {"state": "The deploy finished at 14:00.", "model": "jev-latest", "questions": {"decision": CHOICE}}
|
||||
b.update(over)
|
||||
return b
|
||||
|
||||
|
||||
# ---------------------------------------------------------------- happy paths
|
||||
|
||||
def test_choice_request_is_passed_straight_through_and_the_engine_answer_is_returned():
|
||||
engine = FakeEngine()
|
||||
response = make_client(engine).post("/v1/systemone", json=body(), headers=AUTH)
|
||||
assert response.status_code == 200
|
||||
native, _ = FakeEngine().predict(engine.calls[0])
|
||||
# straight through: only state+questions reach predict, model is not forwarded, the question
|
||||
# dict is untouched (no semif mapping: no A=id:description conversion, no field rename)
|
||||
assert engine.calls == [{"state": "The deploy finished at 14:00.", "questions": {"decision": CHOICE}}]
|
||||
out = response.json()
|
||||
assert out["answers"] == native["answers"]
|
||||
assert out["answers"]["decision"]["type"] == "choice"
|
||||
assert out["answers"]["decision"]["choice"] in ("yes", "no")
|
||||
assert out["usage"] == native["usage"]
|
||||
# a STRING: JevBench's runner hashes the value into its manifest (a dict broke it on the
|
||||
# 0.1.1 first run) and real clients compare it; the string names the real model and revision
|
||||
assert out["model"] == "Intern-Decision-4B@" + "0" * 40
|
||||
|
||||
|
||||
def test_noul_question_needs_no_criteria_and_answers_with_a_probability():
|
||||
engine = FakeEngine()
|
||||
response = make_client(engine).post("/v1/systemone", json=body(questions={"q": NOUL}), headers=AUTH)
|
||||
assert response.status_code == 200
|
||||
assert engine.calls == [{"state": "The deploy finished at 14:00.", "questions": {"q": NOUL}}]
|
||||
answer = response.json()["answers"]["q"]
|
||||
assert answer["type"] == "noul" and answer["noul"] == 0.75
|
||||
|
||||
|
||||
def test_score_question_with_a_list_of_levels_returns_probabilities_over_the_levels():
|
||||
engine = FakeEngine()
|
||||
response = make_client(engine).post("/v1/systemone", json=body(questions={"q": SCORE}), headers=AUTH)
|
||||
assert response.status_code == 200
|
||||
answer = response.json()["answers"]["q"]
|
||||
assert answer["type"] == "score"
|
||||
assert set(answer["probabilities"]) == {"0", "1", "2"}
|
||||
assert abs(sum(answer["probabilities"].values()) - 1.0) < 1e-9
|
||||
|
||||
|
||||
def test_any_number_of_questions_up_to_16_is_ONE_call_so_they_share_one_prompt():
|
||||
engine = FakeEngine()
|
||||
questions = {f"q{i}": NOUL for i in range(16)}
|
||||
response = make_client(engine).post("/v1/systemone", json=body(questions=questions), headers=AUTH)
|
||||
assert response.status_code == 200
|
||||
assert len(engine.calls) == 1 and len(engine.calls[0]["questions"]) == 16
|
||||
|
||||
|
||||
def test_the_model_field_is_accepted_with_any_value_and_ignored():
|
||||
engine = FakeEngine()
|
||||
for model in ("jev-latest", "anything-else", 7, None):
|
||||
response = make_client(engine).post("/v1/systemone", json=body(model=model), headers=AUTH)
|
||||
assert response.status_code == 200
|
||||
|
||||
|
||||
class NoUsageEngine(FakeEngine):
|
||||
def predict(self, request):
|
||||
response, sha = super().predict(request)
|
||||
response.pop("usage")
|
||||
return response, sha
|
||||
|
||||
|
||||
def test_an_engine_that_reports_no_usage_yields_an_empty_usage_object():
|
||||
response = make_client(NoUsageEngine()).post("/v1/systemone", json=body(), headers=AUTH)
|
||||
assert response.status_code == 200
|
||||
assert response.json()["usage"] == {}
|
||||
|
||||
|
||||
# ---------------------------------------------------------------- refusals
|
||||
|
||||
@pytest.mark.parametrize("headers", [{}, {"Authorization": "Bearer wrong"}, {"Authorization": TOKEN}])
|
||||
def test_without_the_right_bearer_a_systemone_post_is_401_and_never_reaches_the_engine(headers):
|
||||
engine = FakeEngine()
|
||||
response = make_client(engine).post("/v1/systemone", json=body(), headers=headers)
|
||||
assert response.status_code == 401
|
||||
assert response.json()["error"]["code"] == "unauthorized"
|
||||
assert engine.calls == []
|
||||
|
||||
|
||||
@pytest.mark.parametrize("images", [["a.png"], [], "a.png"])
|
||||
def test_images_any_value_at_all_is_422_before_the_engine_runs(images):
|
||||
engine = FakeEngine()
|
||||
response = make_client(engine).post("/v1/systemone", json=body(images=images), headers=AUTH)
|
||||
assert response.status_code == 422
|
||||
assert response.json()["error"]["code"] == "invalid_request"
|
||||
assert "images not supported" in response.json()["error"]["message"]
|
||||
assert engine.calls == []
|
||||
|
||||
|
||||
def test_17_questions_is_422_not_a_chunked_2_call_request():
|
||||
engine = FakeEngine()
|
||||
questions = {f"q{i}": NOUL for i in range(17)}
|
||||
response = make_client(engine).post("/v1/systemone", json=body(questions=questions), headers=AUTH)
|
||||
assert response.status_code == 422
|
||||
assert "16" in response.json()["error"]["message"]
|
||||
assert engine.calls == []
|
||||
|
||||
|
||||
def test_zero_questions_is_422():
|
||||
response = make_client().post("/v1/systemone", json=body(questions={}), headers=AUTH)
|
||||
assert response.status_code == 422
|
||||
|
||||
|
||||
def test_a_body_that_is_not_json_or_has_no_questions_is_422():
|
||||
client = make_client()
|
||||
assert client.post("/v1/systemone", content=b"{nope", headers={**AUTH, "content-type": "application/json"}
|
||||
).status_code == 422
|
||||
assert client.post("/v1/systemone", json={"state": "s"}, headers=AUTH).status_code == 422
|
||||
|
||||
|
||||
class TokenLimitEngine(FakeEngine):
|
||||
"""Stands in for the model's own pre-forward check: predict raises ValueError, never a tensor."""
|
||||
def predict(self, request):
|
||||
self.calls.append(request)
|
||||
raise ValueError("Example has 7169 tokens, above 7168; truncation is forbidden")
|
||||
|
||||
|
||||
def test_over_max_tokens_is_422_from_the_models_own_before_the_forward_check():
|
||||
response = make_client(TokenLimitEngine(), max_tokens=7168).post("/v1/systemone", json=body(), headers=AUTH)
|
||||
assert response.status_code == 422
|
||||
assert "7169" in response.json()["error"]["message"]
|
||||
|
||||
|
||||
# ---------------------------------------------------------------- queue and memory
|
||||
|
||||
class SlowEngine(FakeEngine):
|
||||
def __init__(self):
|
||||
super().__init__()
|
||||
self.entered, self.release = threading.Event(), threading.Event()
|
||||
|
||||
def predict(self, request):
|
||||
self.entered.set()
|
||||
assert self.release.wait(5)
|
||||
return super().predict(request)
|
||||
|
||||
|
||||
def test_more_than_max_queue_systemone_posts_get_429_busy():
|
||||
engine = SlowEngine()
|
||||
with make_client(engine, max_queue=1) as client, ThreadPoolExecutor(2) as pool:
|
||||
held = pool.submit(client.post, "/v1/systemone", json=body(), headers=AUTH)
|
||||
assert engine.entered.wait(5)
|
||||
time.sleep(0.2)
|
||||
extra = client.post("/v1/systemone", json=body(), headers=AUTH)
|
||||
assert extra.status_code == 429 and extra.json()["error"]["code"] == "busy"
|
||||
engine.release.set()
|
||||
assert held.result().status_code == 200
|
||||
|
||||
|
||||
class OomOnceEngine(FakeEngine):
|
||||
def __init__(self):
|
||||
super().__init__()
|
||||
self.raised = False
|
||||
|
||||
def predict(self, request):
|
||||
if not self.raised:
|
||||
self.raised, _ = True, self.calls.append(request)
|
||||
raise OutOfMemory("CUDA out of memory.")
|
||||
return super().predict(request)
|
||||
|
||||
|
||||
def test_an_oom_is_503_with_the_contract_envelope_and_the_next_request_succeeds():
|
||||
engine = OomOnceEngine()
|
||||
client = make_client(engine)
|
||||
first = client.post("/v1/systemone", json=body(), headers=AUTH)
|
||||
assert first.status_code == 503 and first.json()["error"]["code"] == "out_of_memory"
|
||||
second = client.post("/v1/systemone", json=body(), headers=AUTH)
|
||||
assert second.status_code == 200
|
||||
|
||||
|
||||
def test_systemone_runs_on_the_same_single_inference_thread_as_decide():
|
||||
seen = []
|
||||
|
||||
class ThreadSpy(FakeEngine):
|
||||
def predict(self, request):
|
||||
seen.append(threading.current_thread().name)
|
||||
return super().predict(request)
|
||||
|
||||
client = make_client(ThreadSpy())
|
||||
client.post("/v1/systemone", json=body(), headers=AUTH)
|
||||
client.post("/decide", json={"id": "r", "state": "s", "question": "q?", "options":
|
||||
[{"id": "a", "description": "x"}, {"id": "b", "description": "yy"}]},
|
||||
headers=AUTH)
|
||||
assert len(set(seen)) == 1 and seen[0].startswith("inference")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------- /health
|
||||
|
||||
def test_health_advertises_the_systemone_surface_and_its_limits():
|
||||
out = make_client(max_tokens=7168).get("/health").json()
|
||||
assert set(out["endpoints"]) == {"/decide", "/decide/shared", "/v1/systemone"}
|
||||
assert out["systemone"] == {"max_questions": 16, "chunking": "none",
|
||||
"images": "not supported", "max_tokens": 7168}
|
||||
Reference in New Issue
Block a user