Spaces:
Running
Running
File size: 21,243 Bytes
8827155 f8a2d1d 8827155 59ba696 8827155 59ba696 8827155 59ba696 8827155 59ba696 8827155 d58afdb 8827155 59ba696 8827155 d58afdb 59ba696 d58afdb 59ba696 d58afdb 59ba696 d58afdb 8827155 e9ab6b6 8827155 59ba696 8827155 d58afdb f8a2d1d d58afdb 7222857 d58afdb f8a2d1d 8827155 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 | """The outbound HTTP worker carries attested, strict results through the queue."""
import asyncio
import io
import json
import math
import os
import time
import unittest
from concurrent.futures import ThreadPoolExecutor
from unittest.mock import patch
import httpx
from fastapi.testclient import TestClient
from app import create_app
from contract import MODEL
from direct_gateway import ARTIFACT_RESPONSE_HEADERS
from http_pull_worker import HTTPRuntime, HTTPWorker, runtime_from_environment
from model_registry import MODEL_ORDER, PROFILES
from pull_worker import Gateway, GatewayError
from relay import Relay
TOKEN = "t" * 32
MANIFEST = "a" * 64
REVISION = "b" * 40
CONTENT = "c" * 64
CANONICAL = PROFILES[MODEL]["repo_id"]
# The deployed Decision 1.0 runtime attests its pre-move Hub ID.
ATTESTED = PROFILES[MODEL]["runtime_model"]
ORIGIN = "http://127.0.0.1:18401"
ARTIFACT = {
"repo_id": CANONICAL,
"manifest_sha256": MANIFEST,
"revision": REVISION,
"content_sha256": CONTENT,
}
class NoTetris:
def public_config(self):
return {"competitors": []}
class LocalGateway:
def __init__(self, client):
self.client = client
def post(self, action, payload):
response = self.client.post(
"/internal/worker/" + action,
json=payload,
headers={"Authorization": "Bearer " + TOKEN},
)
if response.status_code >= 400:
raise GatewayError(response.status_code)
return response.json()
def payload(*, batch=False, model=MODEL):
request = {
"model": model,
"questions": {"route": {
"type": "choice",
"instructions": "Choose a route.",
"criteria": {"left": None, "right": None},
}},
}
if batch:
request["states"] = [
{"id": "one", "state": "First request"},
{"id": "two", "state": "Second request"},
]
else:
request["state"] = "One request"
return request
def answer():
return {
"type": "choice",
"choice": "left",
"confidence": 0.4,
"probabilities": {"left": 0.7, "right": 0.3},
}
def strict_response(request):
usage = {"input_tokens": 12, "output_tokens": 0}
if "states" in request:
return {
"model": request["model"],
"results": [
{"id": row["id"], "answers": {"route": answer()}, "usage": usage}
for row in request["states"]
],
"usage": {"input_tokens": 12 * len(request["states"]), "output_tokens": 0},
}
return {"model": request["model"], "answers": {"route": answer()}, "usage": usage}
class HTTPPullWorkerTests(unittest.TestCase):
def setUp(self):
self.calls = []
self.bad_header = False
self.bad_response = False
self.offline = False
def handler(request):
if request.method == "GET" and request.url.path == "/api/status":
return httpx.Response(200, json={
"status": "offline" if self.offline else "ready",
"artifact": {
"model": ATTESTED,
"revision": REVISION,
"manifest_sha256": MANIFEST,
"content_sha256": CONTENT,
},
})
body = json.loads(request.content)
self.calls.append((request.url.path, body))
response = strict_response(body)
if self.bad_response:
response["model"] = "another/model"
headers = {
header: {
"model": ATTESTED,
"revision": REVISION,
"manifest_sha256": "d" * 64 if self.bad_header else MANIFEST,
"content_sha256": CONTENT,
}[field]
for field, header in ARTIFACT_RESPONSE_HEADERS.items()
}
return httpx.Response(200, json=response, headers=headers)
self.runtime_client = httpx.AsyncClient(transport=httpx.MockTransport(handler))
self.runtime = HTTPRuntime(MODEL, ORIGIN, ARTIFACT, client=self.runtime_client)
self.app_context = TestClient(create_app(
mode="pull_queue",
relay=Relay(TOKEN, MANIFEST, model=MODEL),
registry=[{
"id": MODEL,
"label": "Kai",
"version": "1.0",
"manifest_sha256": MANIFEST,
}],
tetris_manager=NoTetris(),
))
self.client = self.app_context.__enter__()
self.worker = HTTPWorker(LocalGateway(self.client), self.runtime)
self.runtime._load()
self.worker.heartbeat()
def tearDown(self):
self.runtime.close()
asyncio.run(self.runtime_client.aclose())
self.app_context.__exit__(None, None, None)
def process_one(self):
claimed = self.worker.post("claim", {
"worker_id": self.worker.worker_id,
"wait_seconds": 0,
})["job"]
self.assertIsNotNone(claimed)
self.worker.execute(claimed)
return claimed
def test_single_and_batch_jobs_preserve_canonical_runtime_response(self):
for batch, expected_path in (
(False, "/v1/systemone"),
(True, "/v1/systemone/batches"),
):
with self.subTest(batch=batch):
submitted = self.client.post("/api/jobs", json=payload(batch=batch))
self.assertEqual(submitted.status_code, 202)
self.process_one()
completed = self.client.get("/api/jobs/" + submitted.json()["id"])
self.assertEqual(completed.status_code, 200)
self.assertEqual(completed.json()["status"], "succeeded")
self.assertEqual(completed.json()["result"],
strict_response(payload(batch=batch, model=CANONICAL)))
self.assertEqual(self.calls[-1],
(expected_path, payload(batch=batch, model=ATTESTED)))
def test_public_synchronous_route_returns_strict_result(self):
with ThreadPoolExecutor(max_workers=1) as executor:
pending = executor.submit(
self.client.post, "/v1/systemone", json=payload(model=CANONICAL)
)
for _ in range(200):
if self.client.get("/api/status").json()["queued"]:
break
time.sleep(0.01)
else:
self.fail("The public request was not queued")
self.process_one()
response = pending.result(timeout=5)
self.assertEqual(response.status_code, 200)
self.assertEqual(response.json(), strict_response(payload(model=CANONICAL)))
def test_studio_batch_routes_admit_over_256_kib_but_single_stays_bounded(self):
large_state = "x" * (300 * 1024)
batch_request = payload(batch=True)
batch_request["states"] = [{"id": "large", "state": large_state}]
single_request = payload()
single_request["state"] = large_state
for route in ("/api/jobs", "/api/evaluate"):
with self.subTest(route=route):
self.assertEqual(
self.client.post(route, json=single_request).status_code, 413
)
submitted = self.client.post("/api/jobs", json=batch_request)
self.assertEqual(submitted.status_code, 202)
self.process_one()
completed = self.client.get("/api/jobs/" + submitted.json()["id"]).json()
self.assertEqual(completed["status"], "succeeded")
with ThreadPoolExecutor(max_workers=1) as executor:
pending = executor.submit(
self.client.post, "/api/evaluate", json=batch_request
)
for _ in range(200):
if self.client.get("/api/status").json()["queued"]:
break
time.sleep(0.01)
else:
self.fail("The Studio batch was not queued")
self.process_one()
response = pending.result(timeout=5)
self.assertEqual(response.status_code, 200)
self.assertEqual(
response.json(), strict_response(dict(batch_request, model=CANONICAL))
)
def test_bad_artifact_or_response_fails_without_prediction(self):
for field in ("bad_header", "bad_response"):
with self.subTest(field=field):
setattr(self, field, True)
submitted = self.client.post("/api/jobs", json=payload())
self.assertEqual(submitted.status_code, 202)
self.process_one()
completed = self.client.get("/api/jobs/" + submitted.json()["id"]).json()
self.assertEqual(completed["status"], "failed")
self.assertNotIn("result", completed)
setattr(self, field, False)
def test_offline_runtime_does_not_send_a_fresh_heartbeat(self):
before = self.client.app.state.relay.worker["seen"]
self.offline = True
with self.assertRaises(RuntimeError):
self.worker.heartbeat()
self.assertEqual(self.client.app.state.relay.worker["seen"], before)
def test_space_rejects_mismatched_strict_result(self):
submitted = self.client.post("/api/jobs", json=payload()).json()
claimed = self.worker.post("claim", {
"worker_id": self.worker.worker_id,
"wait_seconds": 0,
})["job"]
wrong = strict_response(payload(model=ATTESTED))
wrong["model"] = "another/model"
with self.assertRaises(GatewayError) as rejected:
self.worker.post("result", {
"worker_id": self.worker.worker_id,
"id": claimed["id"],
"lease_token": claimed["lease_token"],
"result": {"kind": "http_runtime_v1", "response": wrong},
})
self.assertEqual(rejected.exception.code, 422)
self.worker.execute(claimed)
completed = self.client.get("/api/jobs/" + submitted["id"]).json()
self.assertEqual(completed["status"], "succeeded")
def test_strict_results_accept_reordered_question_maps(self):
questions = {
"beta": {"type": "noul", "instructions": "Is beta true?"},
"alpha": {"type": "noul", "instructions": "Is alpha true?"},
}
answers = {
"alpha": {"type": "noul", "noul": 0.2},
"beta": {"type": "noul", "noul": 0.8},
}
for batch in (False, True):
with self.subTest(batch=batch):
request = {"model": MODEL, "questions": questions}
if batch:
request["states"] = [
{"id": "first", "state": "First"},
{"id": "second", "state": "Second"},
]
strict = {
"model": ATTESTED,
"results": [
{"id": row["id"], "answers": answers,
"usage": {"input_tokens": 12, "output_tokens": 0}}
for row in request["states"]
],
"usage": {"input_tokens": 24, "output_tokens": 0},
}
else:
request["state"] = "One request"
strict = {
"model": ATTESTED,
"answers": answers,
"usage": {"input_tokens": 12, "output_tokens": 0},
}
submitted = self.client.post("/api/jobs", json=request).json()
claimed = self.worker.post("claim", {
"worker_id": self.worker.worker_id,
"wait_seconds": 0,
})["job"]
self.worker.post("result", {
"worker_id": self.worker.worker_id,
"id": claimed["id"],
"lease_token": claimed["lease_token"],
"result": {"kind": "http_runtime_v1", "response": strict},
})
completed = self.client.get("/api/jobs/" + submitted["id"]).json()
self.assertEqual(completed["status"], "succeeded")
self.assertEqual(completed["result"], dict(strict, model=CANONICAL))
class WorkerConfigurationTests(unittest.TestCase):
def test_claim_transport_accepts_a_large_valid_batch_body(self):
gateway = Gateway("https://example.test", TOKEN)
reply = json.dumps({"job": {"body": {"state": "x" * (1200 * 1024)}}}).encode()
class Opener:
def open(self, request, timeout):
self.request = request
return io.BytesIO(reply)
gateway.opener = Opener()
response = gateway.post("claim", {"worker_id": "1" * 32})
self.assertEqual(len(response["job"]["body"]["state"]), 1200 * 1024)
self.assertTrue(gateway.opener.request.full_url.endswith("/internal/worker/claim"))
def test_every_wire_id_binds_to_its_exact_canonical_model(self):
registry = [{
"id": model,
"label": model,
"version": "1.0",
"manifest_sha256": format(index + 1, "x") * 64,
"revision": format(index + 1, "x") * 40,
} for index, model in enumerate(MODEL_ORDER)]
for model in MODEL_ORDER:
with self.subTest(model=model), patch.dict(os.environ, {
"DECISION_MODEL_REGISTRY_V2": json.dumps(registry),
"DECISION_WORKER_MODEL": model,
"DECISION_RUNTIME_URL": ORIGIN,
}):
runtime = runtime_from_environment()
self.assertEqual(runtime.model, model)
self.assertEqual(runtime.canonical_model, PROFILES[model].get(
"runtime_model", PROFILES[model]["repo_id"]))
self.assertEqual(runtime.manifest, registry[MODEL_ORDER.index(model)]["manifest_sha256"])
with patch.dict(os.environ, {
"DECISION_MODEL_REGISTRY_V2": json.dumps(registry),
"DECISION_WORKER_MODEL": "decision-nano-preview",
"DECISION_RUNTIME_URL": ORIGIN,
}), self.assertRaises(ValueError):
runtime_from_environment()
class PullQueueTetrisReadinessTests(unittest.TestCase):
def test_offline_queue_is_unavailable_in_config_and_race_admission(self):
lux_manifest = "d" * 64
relays = {
MODEL: Relay(TOKEN, MANIFEST, model=MODEL),
"decision-lux": Relay(
TOKEN, lux_manifest, model="decision-lux",
complete_input_tokens=16384,
),
}
registry = [
{"id": MODEL, "label": "Kai", "version": "1.0",
"manifest_sha256": MANIFEST},
{"id": "decision-lux", "label": "Lux", "version": "1.0",
"manifest_sha256": lux_manifest},
]
with TestClient(create_app(
mode="pull_queue", relays=relays, registry=registry,
)) as client:
headers = {"Authorization": "Bearer " + TOKEN}
ready = client.post("/internal/worker/heartbeat", json={
"model": MODEL,
"manifest_sha256": MANIFEST,
"worker_id": "1" * 32,
"phase": "ready",
"capabilities": ["context_batch_v1"],
}, headers=headers)
self.assertEqual(ready.status_code, 200)
config = client.get("/api/tetris/config").json()
by_id = {item["id"]: item["ready"] for item in config["competitors"]}
self.assertTrue(by_id["kai"])
self.assertFalse(by_id["lux"])
self.assertNotIn("nox", by_id)
race = {"left": "kai", "right": "lux", "mode": "steps", "max_steps": 1}
self.assertEqual(client.post("/api/tetris/races", json=race).status_code, 503)
ready = client.post("/internal/worker/heartbeat", json={
"model": "decision-lux",
"manifest_sha256": lux_manifest,
"worker_id": "2" * 32,
"phase": "ready",
"capabilities": ["context_batch_v1"],
}, headers=headers)
self.assertEqual(ready.status_code, 200)
config = client.get("/api/tetris/config").json()
by_id = {item["id"]: item["ready"] for item in config["competitors"]}
self.assertTrue(by_id["lux"])
class PullQueueArenaTests(unittest.TestCase):
KAI2 = {"id": "decision2-kai", "label": "Kai", "version": "2.0", "manifest_sha256": "e" * 64}
WORKER = {"model": "decision2-kai", "manifest_sha256": "e" * 64, "worker_id": "3" * 32}
HEADERS = {"Authorization": "Bearer " + TOKEN}
UNUSED_URL = "https://unused.example.test/v1/systemone"
@staticmethod
def placement(criteria):
tail = 0.4 / (len(criteria) - 1)
probabilities = [0.6] + [tail] * (len(criteria) - 1)
entropy = -sum(p * math.log(p) for p in probabilities)
return {"type": "choice", "choice": criteria[0],
"probabilities": dict(zip(criteria, probabilities)),
"confidence": min(1.0, max(0.0, 1.0 - entropy / math.log(len(criteria))))}
def play(self, *, fail):
canonical = PROFILES["decision2-kai"]["repo_id"]
with patch.dict(os.environ, {
"DECISION_WORKER_TOKEN": TOKEN,
"DECISION_STUDIO_GENERATION": "2.0",
"TETRIS_LOCAL_API_URL": self.UNUSED_URL,
"TETRIS_LOCAL_KAI2_API_URL": self.UNUSED_URL,
"LOCAL_API_URL": self.UNUSED_URL,
}, clear=True), TestClient(create_app(mode="pull_queue", registry=[self.KAI2])) as client:
adapter = client.app.state.tetris.adapter
self.assertEqual(adapter._endpoints["kai2"].url, "")
heartbeat = client.post("/internal/worker/heartbeat", headers=self.HEADERS, json={
**self.WORKER, "phase": "ready", "capabilities": ["context_batch_v1"]})
self.assertEqual(heartbeat.status_code, 200, heartbeat.text)
created = client.post("/api/tetris/races", json={
"left": "kai2", "right": "kai2", "mode": "steps", "max_steps": 2, "seed": 7})
self.assertEqual(created.status_code, 201, created.text)
race_id, turns = created.json()["id"], []
deadline = time.monotonic() + 30
while time.monotonic() < deadline:
job = client.post("/internal/worker/claim", headers=self.HEADERS,
json={**self.WORKER, "wait_seconds": 0.5}).json()["job"]
if job is None:
if client.get(f"/api/tetris/races/{race_id}").json()["status"] not in {
"created", "running"}:
break
continue
body = job["body"]
self.assertEqual(set(body), {"model", "state", "questions"})
self.assertEqual(body["model"], "decision2-kai")
outcome = {"error_code": "inference_failed"} if fail else {"result": {
"kind": "http_runtime_v1",
"response": {"model": canonical,
"answers": {"placement": self.placement(
list(body["questions"]["placement"]["criteria"]))},
"usage": {"input_tokens": 64, "output_tokens": 0}}}}
done = client.post("/internal/worker/result", headers=self.HEADERS, json={
**self.WORKER, "id": job["id"], "lease_token": job["lease_token"], **outcome})
self.assertEqual(done.status_code, 200, done.text)
turns.append(body)
trace = client.get(f"/api/tetris/races/{race_id}/trace").json()
self.assertIsNone(adapter._client)
return canonical, turns, trace
def test_race_turns_use_the_queue_and_never_a_configured_url(self):
canonical, turns, trace = self.play(fail=False)
self.assertEqual(trace["status"], "finished")
steps = trace["traces"]["left"] + trace["traces"]["right"]
self.assertEqual(len(steps), 4)
self.assertTrue(all(step.get("error") is None for step in steps))
requested = [step for step in steps if step.get("provider_model")]
self.assertEqual(len(requested), len(turns))
self.assertTrue(turns)
self.assertEqual({step["provider_model"] for step in requested}, {canonical})
def test_a_failed_queue_turn_ends_that_side(self):
_, turns, trace = self.play(fail=True)
self.assertEqual(trace["status"], "incomplete")
self.assertEqual(len(turns), 2)
for side in ("left", "right"):
with self.subTest(side=side):
self.assertEqual(trace["traces"][side][-1]["error"],
"The selected model could not complete this turn.")
if __name__ == "__main__":
unittest.main()
|