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.
This commit is contained in:
@@ -2,9 +2,10 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import hmac
|
||||
import itertools
|
||||
import math
|
||||
import threading
|
||||
from typing import Any
|
||||
from typing import Any, Literal
|
||||
|
||||
from fastapi import FastAPI, Request
|
||||
from fastapi.concurrency import run_in_threadpool
|
||||
@@ -29,6 +30,7 @@ class Decision(BaseModel):
|
||||
id: str
|
||||
question: str
|
||||
options: list[Option]
|
||||
orderings: Literal["none", "rotations", "all"] = "none"
|
||||
|
||||
def row(self, state: State) -> dict:
|
||||
"""The SemIf row shape: exactly id, state, question, options."""
|
||||
@@ -64,6 +66,75 @@ def _first_error(exc: ValidationError) -> str:
|
||||
return f"{where}: {first.get('msg', 'invalid')}"
|
||||
|
||||
|
||||
MAX_OPTIONS_FOR_ALL = 4
|
||||
|
||||
|
||||
def ordering_perms(decision: Decision) -> list[tuple[int, ...]]:
|
||||
"""Index permutations of the caller's options, the caller's own order first."""
|
||||
n = len(decision.options)
|
||||
if decision.orderings == "rotations":
|
||||
return [tuple((start + k) % n for k in range(n)) for start in range(n)]
|
||||
if n > MAX_OPTIONS_FOR_ALL:
|
||||
raise ApiError(422, "invalid_request",
|
||||
f"orderings 'all' allows at most {MAX_OPTIONS_FOR_ALL} options ({n} given); use 'rotations'")
|
||||
return list(itertools.permutations(range(n)))
|
||||
|
||||
|
||||
def expanded_rows(decision: Decision, state: State, perms: list[tuple[int, ...]]) -> list[dict]:
|
||||
base = decision.row(state)
|
||||
return [{**base, "id": f"{decision.id}#o{k}", "options": [base["options"][i] for i in perm]}
|
||||
for k, perm in enumerate(perms)]
|
||||
|
||||
|
||||
def combine(decision: Decision, results: list[dict]) -> dict:
|
||||
"""Average per-ordering log-softmax by option id; the native results ride along unchanged."""
|
||||
option_ids = [o.id for o in decision.options]
|
||||
logp: dict[str, list[float]] = {i: [] for i in option_ids}
|
||||
probs: dict[str, list[float]] = {i: [] for i in option_ids}
|
||||
tops = []
|
||||
for result in results:
|
||||
logits = result["option_logits"]
|
||||
top = max(logits)
|
||||
lse = top + math.log(sum(math.exp(x - top) for x in logits))
|
||||
for oid, x, p in zip(result["option_ids"], logits, result["probabilities"]):
|
||||
logp[oid].append(x - lse)
|
||||
probs[oid].append(p)
|
||||
tops.append(result["option_ids"][logits.index(top)])
|
||||
means = [sum(logp[i]) / len(logp[i]) for i in option_ids]
|
||||
peak = max(means)
|
||||
weights = [math.exp(m - peak) for m in means]
|
||||
combined_p = [w / sum(weights) for w in weights]
|
||||
winner = option_ids[combined_p.index(max(combined_p))]
|
||||
return {
|
||||
"id": decision.id,
|
||||
"option_ids": option_ids,
|
||||
"combined": {
|
||||
"method": decision.orderings,
|
||||
"orderings": len(results),
|
||||
"probabilities": combined_p,
|
||||
"top": winner,
|
||||
"agreement": tops.count(winner) / len(tops),
|
||||
"spread": {i: [min(probs[i]), max(probs[i])] for i in option_ids},
|
||||
},
|
||||
"orderings": results,
|
||||
}
|
||||
|
||||
|
||||
async def read_limited(stream, declared: str | None, limit: int) -> bytes:
|
||||
"""Read a request body, refusing it once it would exceed `limit` bytes. The check runs BEFORE a
|
||||
chunk is kept, and nothing after the crossing chunk is read (bug hunt C2). A declared length is
|
||||
trusted only as ASCII digits: `"²".isdigit()` is True but `int("²")` raises (S3)."""
|
||||
too_large = ApiError(413, "request_too_large", f"request body exceeds {limit} bytes")
|
||||
if declared is not None and declared.isascii() and declared.isdigit() and int(declared) > limit:
|
||||
raise too_large
|
||||
body = bytearray()
|
||||
async for chunk in stream:
|
||||
if len(body) + len(chunk) > limit:
|
||||
raise too_large
|
||||
body.extend(chunk)
|
||||
return bytes(body)
|
||||
|
||||
|
||||
def calibrated_view(result: dict, workload: str, temperature: float) -> dict:
|
||||
"""softmax(option_logits / T): the native fields are left exactly as SemIf returned them (INV-1)."""
|
||||
scaled = [x / temperature for x in result["option_logits"]]
|
||||
@@ -77,6 +148,27 @@ def create_app(settings: Settings, engine: Any) -> FastAPI:
|
||||
app = FastAPI(title="semif-serve")
|
||||
expected = f"Bearer {settings.api_token}".encode()
|
||||
inference = threading.Lock() # INV-2: one scorer call at a time, off the event loop
|
||||
in_progress = 0 # POSTs admitted and not yet answered (bug hunt C6)
|
||||
|
||||
def admit():
|
||||
nonlocal in_progress
|
||||
if in_progress >= settings.max_queue:
|
||||
raise ApiError(429, "busy", f"{in_progress} requests already in progress (limit {settings.max_queue})")
|
||||
in_progress += 1
|
||||
|
||||
def leave():
|
||||
nonlocal in_progress
|
||||
in_progress -= 1
|
||||
|
||||
def build(fn, *args):
|
||||
"""Response construction from a scorer result (calibration, averaging) maps its failures to
|
||||
the 500 envelope too, instead of escaping as a bare 500 (bug hunt C3)."""
|
||||
try:
|
||||
return fn(*args)
|
||||
except ApiError:
|
||||
raise
|
||||
except Exception as exc: # noqa: BLE001
|
||||
raise ApiError(500, "scoring_failed", f"building the response failed: {type(exc).__name__}: {exc}") from exc
|
||||
|
||||
def locked(fn, *args):
|
||||
with inference:
|
||||
@@ -94,22 +186,10 @@ def create_app(settings: Settings, engine: Any) -> FastAPI:
|
||||
async def api_error(_request: Request, exc: ApiError):
|
||||
return error(exc.status, exc.code, exc.message)
|
||||
|
||||
async def read_limited(request: Request) -> bytes:
|
||||
limit = settings.max_body_bytes
|
||||
too_large = ApiError(413, "request_too_large", f"request body exceeds {limit} bytes")
|
||||
declared = request.headers.get("content-length")
|
||||
if declared is not None and declared.isdigit() and int(declared) > limit:
|
||||
raise too_large
|
||||
body = bytearray()
|
||||
async for chunk in request.stream(): # also caps bodies that declare no length
|
||||
body.extend(chunk)
|
||||
if len(body) > limit:
|
||||
raise too_large
|
||||
return bytes(body)
|
||||
|
||||
async def parse(request: Request, model: type[BaseModel]):
|
||||
try:
|
||||
return model.model_validate_json(await read_limited(request))
|
||||
return model.model_validate_json(
|
||||
await read_limited(request.stream(), request.headers.get("content-length"), settings.max_body_bytes))
|
||||
except ValidationError as exc:
|
||||
raise ApiError(422, "invalid_request", _first_error(exc)) from exc
|
||||
|
||||
@@ -142,20 +222,59 @@ def create_app(settings: Settings, engine: Any) -> FastAPI:
|
||||
"vram_cap_gib": settings.vram_cap_gib, "max_tokens": settings.max_tokens,
|
||||
"max_decisions": settings.max_decisions, "workloads": sorted(settings.calibration)}
|
||||
|
||||
async def score_batch(decisions: list[Decision], state: State, workload: str | None) -> tuple[list[dict], dict]:
|
||||
"""One engine.shared call for every row of every decision; results in request order."""
|
||||
plan = [] # (decision, perms or None, row count)
|
||||
rows: list[dict] = []
|
||||
for d in decisions:
|
||||
if d.orderings == "none":
|
||||
plan.append((d, None, 1))
|
||||
rows.append(d.row(state))
|
||||
else:
|
||||
if workload is not None:
|
||||
raise ApiError(422, "invalid_request",
|
||||
"workload calibration is not available together with orderings")
|
||||
perms = ordering_perms(d)
|
||||
plan.append((d, perms, len(perms)))
|
||||
rows.extend(expanded_rows(d, state, perms))
|
||||
if not 1 <= len(rows) <= settings.max_decisions:
|
||||
raise ApiError(422, "invalid_request",
|
||||
f"this request expands to {len(rows)} scored rows; the limit is 1..{settings.max_decisions}")
|
||||
temperature = temperature_for(workload)
|
||||
results, timing = await score(engine.shared, rows)
|
||||
out, cursor = [], 0
|
||||
for d, perms, count in plan:
|
||||
chunk = results[cursor:cursor + count]
|
||||
cursor += count
|
||||
out.append(build(with_calibration, chunk[0], workload, temperature) if perms is None
|
||||
else build(combine, d, chunk))
|
||||
return out, timing
|
||||
|
||||
@app.post("/decide")
|
||||
async def decide(request: Request):
|
||||
body = await parse(request, DecideBody)
|
||||
temperature = temperature_for(body.workload)
|
||||
return with_calibration(await score(engine.direct, body.row(body.state)), body.workload, temperature)
|
||||
admit()
|
||||
try:
|
||||
body = await parse(request, DecideBody)
|
||||
if body.orderings != "none":
|
||||
results, _timing = await score_batch([body], body.state, body.workload)
|
||||
return results[0]
|
||||
temperature = temperature_for(body.workload)
|
||||
result = await score(engine.direct, body.row(body.state))
|
||||
return build(with_calibration, result, body.workload, temperature)
|
||||
finally:
|
||||
leave()
|
||||
|
||||
@app.post("/decide/shared")
|
||||
async def decide_shared(request: Request):
|
||||
body = await parse(request, SharedBody)
|
||||
if not 1 <= len(body.decisions) <= settings.max_decisions:
|
||||
raise ApiError(422, "invalid_request",
|
||||
f"decisions must hold 1..{settings.max_decisions} entries, got {len(body.decisions)}")
|
||||
temperature = temperature_for(body.workload)
|
||||
results, timing = await score(engine.shared, [d.row(body.state) for d in body.decisions])
|
||||
return {"results": [with_calibration(r, body.workload, temperature) for r in results], "timing": timing}
|
||||
admit()
|
||||
try:
|
||||
body = await parse(request, SharedBody)
|
||||
if not 1 <= len(body.decisions) <= settings.max_decisions:
|
||||
raise ApiError(422, "invalid_request",
|
||||
f"decisions must hold 1..{settings.max_decisions} entries, got {len(body.decisions)}")
|
||||
results, timing = await score_batch(body.decisions, body.state, body.workload)
|
||||
return {"results": results, "timing": timing}
|
||||
finally:
|
||||
leave()
|
||||
|
||||
return app
|
||||
|
||||
@@ -1,4 +1,9 @@
|
||||
"""Settings for semif-serve. Contract: semif-serve.contract.md § Configuration."""
|
||||
"""Settings for semif-serve. Contract: semif-serve.contract.md § Configuration.
|
||||
|
||||
Every value is validated at startup and a bad one is refused with a ValueError naming the
|
||||
variable: a service that starts and then rejects every request (or runs uncapped) is worse
|
||||
than one that does not start (bug hunt 2026-09-27: C1, S1, S2, S8).
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
@@ -8,6 +13,7 @@ from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
|
||||
MIN_TOKEN_CHARS = 32
|
||||
MIN_TEMPERATURE, MAX_TEMPERATURE = 0.05, 20.0
|
||||
SEMIF_COMMIT = "23cf1f39fc9534fe81437200959b6dfc7106e45a"
|
||||
DEFAULT_MODEL = "Qwen/Qwen3.5-4B"
|
||||
DEFAULT_REVISION = "851bf6e806efd8d0a36b00ddf55e13ccb7b8cd0a"
|
||||
@@ -23,35 +29,71 @@ class Settings:
|
||||
max_tokens: int = 4096
|
||||
max_decisions: int = 64
|
||||
max_body_bytes: int = 1024 * 1024
|
||||
max_queue: int = 32
|
||||
calibration: dict[str, float] = field(default_factory=dict)
|
||||
|
||||
@classmethod
|
||||
def from_env(cls, env: Mapping[str, str]) -> "Settings":
|
||||
token = env.get("SEMIF_API_TOKEN", "")
|
||||
if len(token) < MIN_TOKEN_CHARS: # INV-6
|
||||
raise ValueError(f"SEMIF_API_TOKEN must be at least {MIN_TOKEN_CHARS} characters")
|
||||
cap = env.get("SEMIF_VRAM_CAP_GIB")
|
||||
# INV-6: visible ASCII only. A CR, LF or NUL can never arrive in a header, so a token
|
||||
# carrying one would lock every caller out while /health still said ok.
|
||||
if len(token) < MIN_TOKEN_CHARS or not all(33 <= ord(c) <= 126 for c in token):
|
||||
raise ValueError(f"SEMIF_API_TOKEN must be at least {MIN_TOKEN_CHARS} visible ASCII characters")
|
||||
return cls(
|
||||
api_token=token,
|
||||
model=env.get("SEMIF_MODEL", DEFAULT_MODEL),
|
||||
revision=env.get("SEMIF_REVISION", DEFAULT_REVISION),
|
||||
device=env.get("SEMIF_DEVICE", "cuda"),
|
||||
vram_cap_gib=float(cap) if cap else None,
|
||||
max_tokens=int(env.get("SEMIF_MAX_TOKENS", 4096)),
|
||||
max_decisions=int(env.get("SEMIF_MAX_DECISIONS", 64)),
|
||||
max_body_bytes=int(env.get("SEMIF_MAX_BODY_BYTES", 1024 * 1024)),
|
||||
vram_cap_gib=_positive_float(env, "SEMIF_VRAM_CAP_GIB"),
|
||||
max_tokens=_positive_int(env, "SEMIF_MAX_TOKENS", 4096),
|
||||
max_decisions=_positive_int(env, "SEMIF_MAX_DECISIONS", 64),
|
||||
max_body_bytes=_positive_int(env, "SEMIF_MAX_BODY_BYTES", 1024 * 1024),
|
||||
max_queue=_positive_int(env, "SEMIF_MAX_QUEUE", 32),
|
||||
calibration=_load_calibration(env.get("SEMIF_CALIBRATION")),
|
||||
)
|
||||
|
||||
|
||||
def _positive_int(env: Mapping[str, str], name: str, default: int) -> int:
|
||||
raw = env.get(name)
|
||||
if raw is None:
|
||||
return default
|
||||
try:
|
||||
value = int(raw)
|
||||
except ValueError:
|
||||
raise ValueError(f"{name} must be an integer, got {raw!r}") from None
|
||||
if value < 1:
|
||||
raise ValueError(f"{name} must be >= 1, got {value}")
|
||||
return value
|
||||
|
||||
|
||||
def _positive_float(env: Mapping[str, str], name: str) -> float | None:
|
||||
"""Unset means no cap. When set it must be finite and > 0: `0` used to slip through as 'no cap'."""
|
||||
raw = env.get(name)
|
||||
if raw is None or raw == "":
|
||||
return None
|
||||
try:
|
||||
value = float(raw)
|
||||
except ValueError:
|
||||
raise ValueError(f"{name} must be a number, got {raw!r}") from None
|
||||
if not math.isfinite(value) or value <= 0:
|
||||
raise ValueError(f"{name} must be a finite number > 0, got {raw!r}")
|
||||
return value
|
||||
|
||||
|
||||
def _load_calibration(path: str | None) -> dict[str, float]:
|
||||
"""{workload: T}, every T a finite number > 0 (T scales option logits before softmax)."""
|
||||
"""{workload: T}; T scales option logits before softmax, so it is kept in a sane range
|
||||
(a tiny T overflows to NaN and the response then fails to render)."""
|
||||
if not path:
|
||||
return {}
|
||||
table = json.loads(Path(path).read_text())
|
||||
try:
|
||||
table = json.loads(Path(path).read_text())
|
||||
except (OSError, ValueError) as exc:
|
||||
raise ValueError(f"SEMIF_CALIBRATION {path!r} could not be read as JSON: {exc}") from None
|
||||
if not isinstance(table, dict) or not all(
|
||||
isinstance(t, (int, float)) and not isinstance(t, bool) and math.isfinite(t) and t > 0
|
||||
isinstance(t, (int, float)) and not isinstance(t, bool) and math.isfinite(t)
|
||||
and MIN_TEMPERATURE <= t <= MAX_TEMPERATURE
|
||||
for t in table.values()
|
||||
):
|
||||
raise ValueError("SEMIF_CALIBRATION must be a JSON object of workload -> finite temperature > 0")
|
||||
raise ValueError(f"SEMIF_CALIBRATION must be a JSON object of workload -> temperature in "
|
||||
f"[{MIN_TEMPERATURE}, {MAX_TEMPERATURE}]")
|
||||
return {str(k): float(v) for k, v in table.items()}
|
||||
|
||||
@@ -7,10 +7,14 @@ unit-tested against a fake torch (tests/test_engine.py).
|
||||
from __future__ import annotations
|
||||
|
||||
import gc
|
||||
import logging
|
||||
import traceback
|
||||
from typing import Any, Callable
|
||||
|
||||
from .config import Settings
|
||||
from .errors import OutOfMemory
|
||||
from .errors import OutOfMemory, ScoringFailed
|
||||
|
||||
log = logging.getLogger("semif_serve.engine")
|
||||
|
||||
RELEASE_SLACK_BYTES = 512 * 2**20
|
||||
WARMUP_ROW = {
|
||||
@@ -24,6 +28,11 @@ WARMUP_ROW = {
|
||||
}
|
||||
|
||||
|
||||
def _first_line(exc: BaseException) -> str:
|
||||
lines = str(exc).splitlines()
|
||||
return lines[0] if lines else ""
|
||||
|
||||
|
||||
class TorchEngine:
|
||||
def __init__(self, torch: Any, model: Any, tokenizer: Any, metadata: dict, settings: Settings,
|
||||
direct_fn: Callable, shared_fn: Callable, release_above_bytes: int | None = None):
|
||||
@@ -46,7 +55,7 @@ class TorchEngine:
|
||||
arch = f"sm_{major}{minor}"
|
||||
if arch not in torch.cuda.get_arch_list(): # INV-3: no silent PTX/CPU fallback
|
||||
raise RuntimeError(f"torch {torch.__version__} has no kernels for {arch}: {torch.cuda.get_arch_list()}")
|
||||
if settings.vram_cap_gib: # INV-4: cap BEFORE the weights land
|
||||
if settings.vram_cap_gib is not None: # INV-4: cap BEFORE the weights land
|
||||
total = torch.cuda.get_device_properties(0).total_memory
|
||||
fraction = settings.vram_cap_gib * 2**30 / total
|
||||
if not 0 < fraction <= 1:
|
||||
@@ -82,17 +91,28 @@ class TorchEngine:
|
||||
def _guard(self, fn, *args):
|
||||
try:
|
||||
result = fn(*args)
|
||||
except ValueError:
|
||||
raise # validation: SemIf raises it before any GPU work
|
||||
except self._torch.cuda.OutOfMemoryError as exc:
|
||||
message = str(exc).splitlines()[0]
|
||||
failure, message = OutOfMemory, _first_line(exc) or "CUDA out of memory"
|
||||
except Exception as exc: # noqa: BLE001 — every other failure is released and reported below
|
||||
message = _first_line(exc)
|
||||
if "out of memory" in message.lower(): # cuBLAS/cuDNN allocation failures
|
||||
failure = OutOfMemory
|
||||
else:
|
||||
failure, message = ScoringFailed, f"{type(exc).__name__}: {message}"
|
||||
# Formatted text, not exc_info: a log record that keeps the traceback object alive
|
||||
# (pytest's capture handler does; so would any buffering handler) pins the tensors.
|
||||
log.error("scorer failed:\n%s", traceback.format_exc())
|
||||
else:
|
||||
self._release_burst()
|
||||
return result
|
||||
# INV-4, outside the except block on purpose: the torch exception's traceback holds the
|
||||
# failed scorer's frames, and with them its tensors (the replicated prefix cache). Raising
|
||||
# inside the block, or `from exc`, would chain to it and keep GiBs allocated after the 503.
|
||||
# INV-4, outside the except block on purpose: the exception's traceback holds the failed
|
||||
# scorer's frames, and with them its tensors (the replicated prefix cache). Raising inside
|
||||
# the block, or `from exc`, would chain to it and keep GiBs allocated after the response.
|
||||
gc.collect()
|
||||
self._torch.cuda.empty_cache()
|
||||
raise OutOfMemory(message)
|
||||
raise failure(message)
|
||||
|
||||
def direct(self, row: dict) -> dict:
|
||||
return self._guard(self._direct, self._model, self._tokenizer, row, self._metadata, self._settings.max_tokens)
|
||||
|
||||
@@ -3,3 +3,8 @@
|
||||
|
||||
class OutOfMemory(RuntimeError):
|
||||
"""The engine ran out of GPU memory during a request and has already released its cache (INV-4)."""
|
||||
|
||||
|
||||
class ScoringFailed(RuntimeError):
|
||||
"""A scorer call failed for a reason other than validation or OOM. Raised unchained, after the
|
||||
failed call's memory has been released; the original traceback is logged, not carried (INV-4)."""
|
||||
|
||||
@@ -11,6 +11,10 @@ from .config import Settings
|
||||
|
||||
def app_from_env() -> FastAPI:
|
||||
settings = Settings.from_env(os.environ)
|
||||
# INV-5: never download at runtime, inside the image or out of it (bug hunt S10). Set before
|
||||
# torch / transformers / huggingface_hub are imported, since they read it at import time.
|
||||
os.environ["HF_HUB_OFFLINE"] = "1"
|
||||
os.environ["TRANSFORMERS_OFFLINE"] = "1"
|
||||
from .engine import TorchEngine # torch loads only here, never in the unit tests
|
||||
|
||||
return create_app(settings, TorchEngine.load(settings))
|
||||
|
||||
Reference in New Issue
Block a user