"""Depth-aware 503 semantics (plan 6.2.1) and never-proxy-half-awake.""" from __future__ import annotations import asyncio import httpx from conftest import raw_asgi async def _sleep_at_level(stack, level: int) -> None: result = await stack.manager.sleep_service("text", level, reason="test") assert result["ok"], result async def test_503_from_level1_sleep(pub, backend, stack): await _sleep_at_level(stack, 1) stack.cfg.hold_sleep_s = 0.05 backend.services["vllm-text"]["wake_delay"] = 0.4 r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 503 assert r.headers["retry-after"] == "10" err = r.json()["error"] assert err["code"] == "model_waking" assert err["type"] == "model_waking" assert err["sleep_depth"] == "sleeping" assert isinstance(err["estimated_wake_seconds"], int) assert err["estimated_wake_seconds"] == 6 assert "Qwen3.6-35B-A3B-FP8" in err["message"] async def test_503_from_level2_offload(pub, backend, stack): await _sleep_at_level(stack, 2) stack.cfg.hold_offload_s = 0.05 backend.services["vllm-text"]["wake_delay"] = 0.4 r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 503 assert r.headers["retry-after"] == "60" err = r.json()["error"] assert err["sleep_depth"] == "offloaded" assert err["estimated_wake_seconds"] == 60 async def test_503_when_container_restarting(pub, backend, stack): backend.services["vllm-text"]["reachable"] = False stack.cfg.hold_restart_s = 0.05 r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 503 assert r.headers["retry-after"] == "600" err = r.json()["error"] assert err["sleep_depth"] == "restarting" assert err["estimated_wake_seconds"] == 600 async def test_unknown_depth_is_conservative_offloaded(pub, backend, stack): """Router restarted / slept behind our back: depth unknown -> offloaded.""" backend.set_sleeping("ocr", True, level=1) # actually only level 1 asleep stack.cfg.hold_offload_s = 0.05 backend.services["vllm-ocr"]["wake_delay"] = 0.3 r = await pub.post("/v1/chat/completions", json={"model": "ocr"}) assert r.status_code == 503 err = r.json()["error"] # conservative: worst-case depth, longest client wait assert err["sleep_depth"] == "offloaded" assert r.headers["retry-after"] == "60" async def test_wake_sequence_retried_once_then_503(pub, backend, stack): backend.set_sleeping("text", True) backend.services["vllm-text"]["wake_fails"] = 2 r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 503 assert r.headers["retry-after"] == "60" err = r.json()["error"] assert err["code"] == "model_waking" assert err["sleep_depth"] == "offloaded" # exactly one retry (plan 6.2.1) assert backend.count("POST vllm-text/wake_up") == 2 # never proxied to a half-awake backend assert backend.count("POST vllm-text/v1/") == 0 async def test_never_proxy_before_health_ok(pub, backend): """The /v1 call must happen after the wake sequence, never before it. Requests admitted between wake_up and reload_weights return 200 + garbage, so the whole sequence has to finish first. """ backend.set_sleeping("text", True) backend.services["vllm-text"]["wake_delay"] = 0.1 r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 200 svc = backend.services["vllm-text"] assert svc["wake_up_calls"] == 1 assert svc["reload_calls"] == 1 assert svc["reset_calls"] == 1 calls = backend.calls idx = {needle: backend.last(needle) for needle in ( "POST vllm-text/wake_up", "POST vllm-text/collective_rpc", "POST vllm-text/reset_prefix_cache", "POST vllm-text/v1/chat/completions", )} assert idx["POST vllm-text/v1/chat/completions"] > idx["POST vllm-text/reset_prefix_cache"] assert idx["POST vllm-text/reset_prefix_cache"] > idx["POST vllm-text/collective_rpc"] assert idx["POST vllm-text/collective_rpc"] > idx["POST vllm-text/wake_up"] assert calls[-1] == "POST vllm-text/v1/chat/completions" async def test_level1_wake_uses_the_fast_path(pub, backend, stack): """From level-1 sleep, /wake_up ALONE is enough (calibration 2026-08-17): bit-identical output at temp 0, and ~20s cheaper than the reload sequence (23.4s -> 2.5-3.8s on the text model).""" assert (await stack.manager.sleep_service("text", 1))["ok"] is True assert stack.manager.services["text"].depth == "sleeping" backend.calls.clear() svc = backend.services["vllm-text"] for field in ("wake_up_calls", "reload_calls", "reset_calls"): svc[field] = 0 r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 200 assert svc["wake_up_calls"] == 1 # wake_up only ... assert svc["reload_calls"] == 0 # ... no reload_weights ... assert svc["reset_calls"] == 0 # ... and no prefix-cache reset assert backend.count("POST vllm-text/v1/chat/completions") == 1 async def test_unknown_depth_uses_the_full_sequence(pub, backend, stack): """Router restarted / slept out of band: unknown depth -> conservative level-2 treatment, full sequence.""" backend.set_sleeping("text", True) # router depth stays unknown r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 200 svc = backend.services["vllm-text"] assert (svc["wake_up_calls"], svc["reload_calls"], svc["reset_calls"]) == (1, 1, 1) async def test_readiness_is_gated_on_is_sleeping_not_health(pub, backend, stack): """/health answers 200 on a sleeping backend, so /is_sleeping is the gate; nothing is proxied while the backend still reports is_sleeping=true (a request sent to a sleeping backend hangs instead of erroring).""" assert (await stack.manager.sleep_service("text", 2))["ok"] is True backend.calls.clear() backend.services["vllm-text"]["hold_sleeping_polls"] = 3 # wake "in flight" r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 200 calls = backend.calls probes = [i for i, c in enumerate(calls) if c == "GET vllm-text/is_sleeping"] assert len(probes) >= 4 # kept polling until it flipped assert calls[-1] == "POST vllm-text/v1/chat/completions" # proxied last assert backend.services["vllm-text"]["api_calls"] == 1 async def test_backend_dies_between_health_and_proxy(pub, backend, stack): """Transport error mid-proxy -> service marked restarting -> depth 503.""" text = backend.services["vllm-text"] original = backend.handle_async_request async def flaky(request: httpx.Request) -> httpx.Response: if request.url.path.startswith("/v1/"): raise httpx.ConnectError("backend gone", request=request) return await original(request) backend.handle_async_request = flaky stack.cfg.hold_restart_s = 0.05 r = await pub.post("/v1/chat/completions", json={"model": "text"}) assert r.status_code == 503 assert r.headers["retry-after"] == "600" assert r.json()["error"]["sleep_depth"] == "restarting" assert text["wake_up_calls"] == 0 async def test_wake_after_level2_runs_full_sequence(pub, backend): backend.set_sleeping("embed", True, level=2) r = await pub.post("/v1/embeddings", json={"input": "hi"}) assert r.status_code == 200 svc = backend.services["vllm-embed"] assert svc["wake_up_calls"] == 1 assert svc["reload_calls"] == 1 # reload_weights is mandatory after L2 assert svc["reset_calls"] == 1 # prefix cache reset too assert svc["api_calls"] == 1 async def test_error_body_is_openai_shaped(pub, backend, stack): await _sleep_at_level(stack, 1) stack.cfg.hold_sleep_s = 0.01 backend.services["vllm-text"]["wake_delay"] = 0.2 r = await pub.post("/v1/chat/completions", json={"model": "text"}) payload = r.json() assert set(payload) == {"error"} assert set(payload["error"]) == { "type", "code", "message", "sleep_depth", "estimated_wake_seconds" } assert r.headers["content-type"].startswith("application/json") async def test_depth_survives_raw_traversal_requests(stack, backend): """Traversal requests are rejected before any backend contact.""" backend.set_sleeping("text", True) for raw in ("/v1/../sleep", "/v1%2f..%2fsleep", "//sleep", "/v1/../../wake_up"): status, _ = await raw_asgi(stack.public, "POST", raw) assert status == 404, raw assert backend.calls == []