feat(semif): SemIf option-logit decisions on fv-ml1 GPU 1 (Prime)
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.
This commit is contained in:
@@ -0,0 +1,161 @@
|
||||
"""semif-serve HTTP layer. Contract: semif-serve.contract.md."""
|
||||
from __future__ import annotations
|
||||
|
||||
import hmac
|
||||
import math
|
||||
import threading
|
||||
from typing import Any
|
||||
|
||||
from fastapi import FastAPI, Request
|
||||
from fastapi.concurrency import run_in_threadpool
|
||||
from fastapi.responses import JSONResponse
|
||||
from pydantic import BaseModel, ConfigDict, ValidationError
|
||||
|
||||
from .config import SEMIF_COMMIT, Settings
|
||||
from .errors import OutOfMemory
|
||||
|
||||
OPEN_PATHS = frozenset({"/health"})
|
||||
State = str | dict[str, Any] | list[Any]
|
||||
|
||||
|
||||
class Option(BaseModel):
|
||||
model_config = ConfigDict(extra="ignore")
|
||||
id: str
|
||||
description: str
|
||||
|
||||
|
||||
class Decision(BaseModel):
|
||||
model_config = ConfigDict(extra="ignore")
|
||||
id: str
|
||||
question: str
|
||||
options: list[Option]
|
||||
|
||||
def row(self, state: State) -> dict:
|
||||
"""The SemIf row shape: exactly id, state, question, options."""
|
||||
return {"id": self.id, "state": state, "question": self.question,
|
||||
"options": [o.model_dump() for o in self.options]}
|
||||
|
||||
|
||||
class DecideBody(Decision):
|
||||
state: State
|
||||
workload: str | None = None
|
||||
|
||||
|
||||
class SharedBody(BaseModel):
|
||||
model_config = ConfigDict(extra="ignore")
|
||||
state: State
|
||||
decisions: list[Decision]
|
||||
workload: str | None = None
|
||||
|
||||
|
||||
class ApiError(Exception):
|
||||
def __init__(self, status: int, code: str, message: str):
|
||||
super().__init__(message)
|
||||
self.status, self.code, self.message = status, code, message
|
||||
|
||||
|
||||
def error(status: int, code: str, message: str) -> JSONResponse:
|
||||
return JSONResponse(status_code=status, content={"error": {"code": code, "message": message}})
|
||||
|
||||
|
||||
def _first_error(exc: ValidationError) -> str:
|
||||
first = exc.errors()[0]
|
||||
where = ".".join(str(p) for p in first.get("loc", ())) or "body"
|
||||
return f"{where}: {first.get('msg', 'invalid')}"
|
||||
|
||||
|
||||
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"]]
|
||||
top = max(scaled)
|
||||
weights = [math.exp(x - top) for x in scaled]
|
||||
total = sum(weights)
|
||||
return {"workload": workload, "temperature": temperature, "probabilities": [w / total for w in weights]}
|
||||
|
||||
|
||||
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
|
||||
|
||||
def locked(fn, *args):
|
||||
with inference:
|
||||
return fn(*args)
|
||||
|
||||
@app.middleware("http")
|
||||
async def require_bearer(request: Request, call_next):
|
||||
if request.url.path not in OPEN_PATHS:
|
||||
supplied = request.headers.get("authorization", "").encode()
|
||||
if not hmac.compare_digest(supplied, expected): # INV-6
|
||||
return error(401, "unauthorized", "missing or wrong bearer token")
|
||||
return await call_next(request)
|
||||
|
||||
@app.exception_handler(ApiError)
|
||||
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))
|
||||
except ValidationError as exc:
|
||||
raise ApiError(422, "invalid_request", _first_error(exc)) from exc
|
||||
|
||||
async def score(fn, *args):
|
||||
"""Run one scorer call in a worker thread under the lock; map its failures to contract codes."""
|
||||
try:
|
||||
return await run_in_threadpool(locked, fn, *args)
|
||||
except ValueError as exc: # SemIf validation, token limit, tokenisation
|
||||
raise ApiError(422, "invalid_request", str(exc)) from exc
|
||||
except OutOfMemory as exc:
|
||||
raise ApiError(503, "out_of_memory", str(exc)) from exc
|
||||
except Exception as exc: # noqa: BLE001 — any other scorer failure
|
||||
raise ApiError(500, "scoring_failed", f"{type(exc).__name__}: {exc}") from exc
|
||||
|
||||
def temperature_for(workload: str | None) -> float | None:
|
||||
if workload is None:
|
||||
return None
|
||||
if workload not in settings.calibration:
|
||||
raise ApiError(422, "invalid_request", f"unknown workload {workload!r}")
|
||||
return settings.calibration[workload]
|
||||
|
||||
def with_calibration(result: dict, workload: str | None, temperature: float | None) -> dict:
|
||||
if temperature is None:
|
||||
return result
|
||||
return {**result, "calibrated": calibrated_view(result, workload, temperature)}
|
||||
|
||||
@app.get("/health")
|
||||
async def health():
|
||||
return {"status": "ok", "semif_commit": SEMIF_COMMIT, "model": engine.health(),
|
||||
"vram_cap_gib": settings.vram_cap_gib, "max_tokens": settings.max_tokens,
|
||||
"max_decisions": settings.max_decisions, "workloads": sorted(settings.calibration)}
|
||||
|
||||
@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)
|
||||
|
||||
@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}
|
||||
|
||||
return app
|
||||
@@ -0,0 +1,57 @@
|
||||
"""Settings for semif-serve. Contract: semif-serve.contract.md § Configuration."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import math
|
||||
from collections.abc import Mapping
|
||||
from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
|
||||
MIN_TOKEN_CHARS = 32
|
||||
SEMIF_COMMIT = "23cf1f39fc9534fe81437200959b6dfc7106e45a"
|
||||
DEFAULT_MODEL = "Qwen/Qwen3.5-4B"
|
||||
DEFAULT_REVISION = "851bf6e806efd8d0a36b00ddf55e13ccb7b8cd0a"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Settings:
|
||||
api_token: str
|
||||
model: str = DEFAULT_MODEL
|
||||
revision: str = DEFAULT_REVISION
|
||||
device: str = "cuda"
|
||||
vram_cap_gib: float | None = None
|
||||
max_tokens: int = 4096
|
||||
max_decisions: int = 64
|
||||
max_body_bytes: int = 1024 * 1024
|
||||
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")
|
||||
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)),
|
||||
calibration=_load_calibration(env.get("SEMIF_CALIBRATION")),
|
||||
)
|
||||
|
||||
|
||||
def _load_calibration(path: str | None) -> dict[str, float]:
|
||||
"""{workload: T}, every T a finite number > 0 (T scales option logits before softmax)."""
|
||||
if not path:
|
||||
return {}
|
||||
table = json.loads(Path(path).read_text())
|
||||
if not isinstance(table, dict) or not all(
|
||||
isinstance(t, (int, float)) and not isinstance(t, bool) and math.isfinite(t) and t > 0
|
||||
for t in table.values()
|
||||
):
|
||||
raise ValueError("SEMIF_CALIBRATION must be a JSON object of workload -> finite temperature > 0")
|
||||
return {str(k): float(v) for k, v in table.items()}
|
||||
@@ -0,0 +1,101 @@
|
||||
"""The real engine: SemIf's torch scorers over one resident model. Needs the `model` extra.
|
||||
|
||||
Contract: semif-serve.contract.md, INV-3 (fail-closed startup), INV-4 (VRAM cap + OOM),
|
||||
INV-5 (offline weights). load() is checked on the card at acceptance; the OOM path is
|
||||
unit-tested against a fake torch (tests/test_engine.py).
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import gc
|
||||
from typing import Any, Callable
|
||||
|
||||
from .config import Settings
|
||||
from .errors import OutOfMemory
|
||||
|
||||
RELEASE_SLACK_BYTES = 512 * 2**20
|
||||
WARMUP_ROW = {
|
||||
"id": "semif-serve-warmup",
|
||||
"state": "The deployment completed at 14:02 UTC. Health checks passed in all three zones.",
|
||||
"question": "Is there evidence that the deployment succeeded?",
|
||||
"options": [
|
||||
{"id": "yes", "description": "The deployment succeeded."},
|
||||
{"id": "no", "description": "The deployment did not succeed."},
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
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):
|
||||
self._torch, self._model, self._tokenizer = torch, model, tokenizer
|
||||
self._metadata, self._settings = metadata, settings
|
||||
self._direct, self._shared = direct_fn, shared_fn
|
||||
self._release_above = release_above_bytes
|
||||
|
||||
@classmethod
|
||||
def load(cls, settings: Settings) -> "TorchEngine":
|
||||
import torch
|
||||
from semif_phase1.core import load_causal_model
|
||||
from semif_phase1.direct import score
|
||||
from semif_phase1.shared import score_shared
|
||||
|
||||
if settings.device == "cuda":
|
||||
if not torch.cuda.is_available():
|
||||
raise RuntimeError("SEMIF_DEVICE=cuda but torch sees no CUDA device")
|
||||
major, minor = torch.cuda.get_device_capability(0)
|
||||
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
|
||||
total = torch.cuda.get_device_properties(0).total_memory
|
||||
fraction = settings.vram_cap_gib * 2**30 / total
|
||||
if not 0 < fraction <= 1:
|
||||
raise ValueError(f"SEMIF_VRAM_CAP_GIB={settings.vram_cap_gib} does not fit a {total / 2**30:.1f} GiB card")
|
||||
torch.cuda.set_per_process_memory_fraction(fraction, 0)
|
||||
elif settings.device != "cpu":
|
||||
raise ValueError(f"SEMIF_DEVICE must be cuda or cpu, not {settings.device!r}")
|
||||
|
||||
model, tokenizer, metadata = load_causal_model(settings.model, settings.revision, settings.device, "bfloat16")
|
||||
placed = next(model.parameters()).device.type
|
||||
if placed != settings.device: # INV-3
|
||||
raise RuntimeError(f"model landed on {placed}, expected {settings.device}")
|
||||
engine = cls(torch, model, tokenizer, metadata, settings, direct_fn=score, shared_fn=score_shared)
|
||||
engine.direct(WARMUP_ROW) # INV-3: one decision must score
|
||||
if settings.device == "cuda": # INV-4: the resting footprint
|
||||
engine._release_above = torch.cuda.memory_reserved(0) + RELEASE_SLACK_BYTES
|
||||
return engine
|
||||
|
||||
def health(self) -> dict:
|
||||
info = dict(self._metadata)
|
||||
if self._settings.device == "cuda":
|
||||
info["device_name"] = self._torch.cuda.get_device_name(0)
|
||||
info["allocated_gib"] = round(self._torch.cuda.memory_allocated(0) / 2**30, 2)
|
||||
info["reserved_gib"] = round(self._torch.cuda.memory_reserved(0) / 2**30, 2)
|
||||
return info
|
||||
|
||||
def _release_burst(self) -> None:
|
||||
"""INV-4: hand a burst back to the driver so the card's shared headroom (scriberr, the
|
||||
vLLM seats) returns after a big request, instead of sitting in torch's cache."""
|
||||
if self._release_above is not None and self._torch.cuda.memory_reserved(0) > self._release_above:
|
||||
self._torch.cuda.empty_cache()
|
||||
|
||||
def _guard(self, fn, *args):
|
||||
try:
|
||||
result = fn(*args)
|
||||
except self._torch.cuda.OutOfMemoryError as exc:
|
||||
message = str(exc).splitlines()[0]
|
||||
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.
|
||||
gc.collect()
|
||||
self._torch.cuda.empty_cache()
|
||||
raise OutOfMemory(message)
|
||||
|
||||
def direct(self, row: dict) -> dict:
|
||||
return self._guard(self._direct, self._model, self._tokenizer, row, self._metadata, self._settings.max_tokens)
|
||||
|
||||
def shared(self, rows: list[dict]) -> tuple[list[dict], dict]:
|
||||
return self._guard(self._shared, self._model, self._tokenizer, rows, self._metadata, self._settings.max_tokens)
|
||||
@@ -0,0 +1,5 @@
|
||||
"""Torch-free exceptions shared by the HTTP layer and the engine."""
|
||||
|
||||
|
||||
class OutOfMemory(RuntimeError):
|
||||
"""The engine ran out of GPU memory during a request and has already released its cache (INV-4)."""
|
||||
@@ -0,0 +1,16 @@
|
||||
"""uvicorn entry point: `uvicorn semif_serve.main:app_from_env --factory --workers 1`."""
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
from fastapi import FastAPI
|
||||
|
||||
from .app import create_app
|
||||
from .config import Settings
|
||||
|
||||
|
||||
def app_from_env() -> FastAPI:
|
||||
settings = Settings.from_env(os.environ)
|
||||
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