diff --git a/AGENTS.md b/AGENTS.md index d80fc0c5..3f8acb5d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -28,7 +28,7 @@ State / architecture / decisions / "why": Obsidian Vault, `AINode` (cluster ops: - **The legacy 110-item decision path (`--backend ainode|chat|jev`) scores typed decisions against labels, and its confidence numbers are the product.** Accuracy is the weakest number in the block: a wrong answer at 0.95 is the failure mode, so every block carries the Brier score on the labeled option, an expected calibration error with the five-bin reliability table behind it, and the count of wrong answers surviving a 0.8 and a 0.9 gate. **A probability nobody reported is absent, never assumed** (such a row is in the accuracy, out of the calibration, and counted in `no_confidence`), **an item that failed is one row with an `error`** and never a wrong answer, and **cost is a posted vendor rate over reported tokens or `0`** for a local backend, never an estimate. `bench/decide/items.json` is repo data versioned next to its results: loading is strict, a malformed item is a load error rather than a skipped item, and a set's items must share one kind, question and option set because a set is one measurement. Backends are split `request()` / `parse()` as pure functions so `tests/test_bench_decide.py` pins every request shape and response shape with canned payloads and no network. **The TypeSafe key is never printed, never written into a record and never put in a note**: a run reports only which of `--api-key`, `$TYPESAFE_API_KEY` or `~/.jev_api_key` it came from. Sets, metrics and flags: `bench/decide/README.md`. - **The speech bench (`ainode/bench/speech/`, `scripts/ainode-bench.py speech`) scores against committed audio, and both halves of that are load-bearing.** A word error rate is only comparable over the same bytes, so the ten clips live in `bench/speech/clips/` as repo data (1.4 MB) rather than being synthesised per run, and **`clips.CLIPS_VERSION` is bumped on any edit to a text, a voice or a WAV**; `--generate-clips` rebuilds the set with macOS `say` plus `afconvert` and is a maintenance step a run never takes. The reference is the exact string handed to `say`, fixed before the run, and **nothing adjusts a reference after a transcript is seen**: a reference edited to match what a model said makes the rate a statement about the editor. The normaliser is part of the measurement, so it is versioned (`metrics.NORMALIZER_VERSION`) and the record carries BOTH rates, `wer` with number words folded to digits and `wer_orthographic` with case and punctuation only, because a transcript that heard every word and wrote "9" for "nine" is not a hearing error and one number alone hides which kind it was. Nothing is folded that changes a word: no stopword list, no stemming, no synonym map, no per-clip exception. `wer` is pooled over words, never a mean of per-clip rates. **A clip that failed is one row with an `error` and nulls for every number**, counted out of every rate, percentile and factor, never folded in as a 100 percent error rate: a transport failure inside a figure a reader takes as the model's is the one mistake this section can make. It is the one bench whose request body is not JSON (`client.py` assembles the multipart itself), so a run through a node's `:3000/v1` exercises the fleet's own audio path; the block is `bench/SCHEMA.md`. - **One proxy handler serves every forwarded inference path**: `proxy_to_vllm` is registered for `POST /v1/chat/completions`, `/v1/completions`, `/v1/messages`, `/v1/messages/count_tokens`, `/v1/responses`, `/v1/rerank`, `/v1/score`, `/v1/audio/transcriptions`, `/v1/audio/translations`, `/tokenize` and `/detokenize` (`GET /v1/models` is the federated union and does not forward; `POST /v1/embeddings` keeps its own handler because it validates the body first, and `/v1/decide` plus `/v1/systemone` compose their own completions). Add a path by registering it on that handler, never by writing a second proxy: routing on the body's `model`, transport failover, the multimodal ordering below, SSE passthrough and header passthrough are all protocol-agnostic and already there. Two of the paths are NOT under `/v1` because vLLM does not serve them there (`/tokenize`, `/detokenize`), so the table is the authority on the path and not a prefix rule. **The two audio paths are the ones whose body is NOT JSON**: OpenAI's speech-to-text API is a `multipart/form-data` upload with the model id as a form field, so the handler reads it with `api/multipart.py::form_fields` over the body it already buffered (never `request.multipart()`, which consumes the stream the proxy still has to forward) and forwards the bytes UNCHANGED under the caller's own `Content-Type`: a multipart body is only parseable against the boundary in its own header, so re-encoding the parts hands the engine a body the forwarded header no longer describes. A multipart body that names no `model` is a 400 naming the field, never a fallback to this node's own model, which would send someone's audio to a chat engine. They are also the first pair whose existence is MODEL-CONDITIONAL: vLLM attaches its speech-to-text router only when the served model reports the `transcription` task, so no chat or pooling engine's `/openapi.json` lists them and the check below cannot be run on one. For a path like that, the evidence is the router's own declaration in the engine image the recipe pins (`entrypoints/openai/speech_to_text/api_router.py`), and the `openapi.json` check still applies the first time such a model actually serves. Forward `request.path_qs`, not `request.path`: Claude Code posts to `/v1/messages?beta=true`. This is a route table and not a catch-all: an unregistered path stays a 404, which is why a new path is added only after `curl http://:/openapi.json` on a real engine says the engine answers it. A body with no `model` falls back to this node's own model ONLY when it serves a primary (`api/server.py::own_model`: `app["engine"]` set and `config.model` named); a node that serves nothing, a routing-only master above all, answers 400 `missing_model_field` and forwards nothing (`tests/test_proxy_no_model.py`). -- **`/v1/decide`'s response shape is a contract, not an implementation detail** (`api/decide.py`). Its bench is written against the exact shape (`model`, `node`, `latency_ms`, `decisions[key] = {answer, confidence, distribution, latency_ms}`, `usage = {prompt_tokens, completion_tokens, calls}`), so a `200` always carries every question asked: a bad request is a `400` and an engine that cannot answer is a `503`, never a partial `decisions` block. Probabilities are keyed by the caller's OPTION strings, never by the letters used to constrain the engine. It is the one `/v1` path deliberately NOT on `proxy_to_vllm`, because it composes N grammar-constrained chat completions of its own from one request and has no caller body to forward; it still routes through the proxy's own `_routing_candidates` and the shared `app["client_session"]`, so never give it its own routing rule or HTTP stack. The constraint field is vLLM 0.27.1's `structured_outputs: {"choice": [...]}`: the legacy `guided_choice` is accepted by that image and then silently ignored, so sending it instead would produce free prose with no error. The engine-facing half of a decision request is `decide.py::run_questions`, shared with `/v1/systemone`: one path to the engines and one way the probabilities are read, so a route decides the status and the shape and nothing else. **A decision adapter's own temperatures are applied there and nowhere else** (#276): when the served model's directory in THIS node's store (`decide.py::model_store_dir`, which goes through `models/registry.py::snapshot_dir_for`, never a hardcoded path) carries `temperatures.json`, each question's label logprobs are divided by its kind's temperature (`choice`, `noul`, `score`; decide's `boolean` is `noul`, its `score` is `score`, an `options` list is `choice`) before the softmax, `"calibration": "raw"` opts out, and both routes answer a top-level `calibration: {applied, temperatures}` block so a caller can refit. A node with no copy of the model (the master, usually) takes the table from the node it routes to, `GET /api/decide/calibration` over the fleet key (`decide.py::fetch_peer_temperatures`, cached `PEER_TEMPERATURES_TTL_S`, a peer that did not answer remembered `PEER_TEMPERATURES_RETRY_S`), so the answer is the same wherever the request entered; a failure there never fails a decision, it answers raw with `temperatures: null`. **A decision model warms on bind** (#277): `models/api_routes.py::_wait_for_bind` hands every bound engine to `decide.py::schedule_decision_warmup`, which sends one constrained question per kind through `ask_one` in the background when the directory carries `prompt_contract.json` or `temperatures.json`, and `/api/status` reports `warm` per instance (null for a model with nothing to warm). An engine adopted at restart never binds, so adoption warms too: the boot primary in `_await_primary_bind` and every adopted stacked instance in `_warm_adopted_stacked`. Readiness is not held for it. A connected engine call that times out is reported as the grammar compiling, never as an unreachable node. +- **`/v1/decide`'s response shape is a contract, not an implementation detail** (`api/decide.py`). Its bench is written against the exact shape (`model`, `node`, `latency_ms`, `decisions[key] = {answer, confidence, distribution, latency_ms}`, `usage = {prompt_tokens, completion_tokens, calls}`), so a `200` always carries every question asked: a bad request is a `400` and an engine that cannot answer is a `503`, never a partial `decisions` block. Probabilities are keyed by the caller's OPTION strings, never by the letters used to constrain the engine. It is the one `/v1` path deliberately NOT on `proxy_to_vllm`, because it composes N grammar-constrained chat completions of its own from one request and has no caller body to forward; it still routes through the proxy's own `_routing_candidates` and the shared `app["client_session"]`, so never give it its own routing rule or HTTP stack. The constraint field is vLLM 0.27.1's `structured_outputs: {"choice": [...]}`: the legacy `guided_choice` is accepted by that image and then silently ignored, so sending it instead would produce free prose with no error. The engine-facing half of a decision request is `decide.py::run_questions`, shared with `/v1/systemone`: one path to the engines and one way the probabilities are read, so a route decides the status and the shape and nothing else. **A decision adapter's own temperatures are applied there and nowhere else** (#276): when the served model's directory in THIS node's store (`decide.py::model_store_dir`, which goes through `models/registry.py::snapshot_dir_for`, never a hardcoded path) carries `temperatures.json`, each question's label logprobs are divided by its kind's temperature (`choice`, `noul`, `score`; decide's `boolean` is `noul`, its `score` is `score`, an `options` list is `choice`) before the softmax, `"calibration": "raw"` opts out, and both routes answer a top-level `calibration: {applied, temperatures}` block so a caller can refit. When no owner answers a forwarded request (below) and this node calls the engines itself, a node with no copy of the model takes the table from the node it routes to, `GET /api/decide/calibration` over the fleet key (`decide.py::fetch_peer_temperatures`, cached `PEER_TEMPERATURES_TTL_S`, a peer that did not answer remembered `PEER_TEMPERATURES_RETRY_S`), so the answer is the same wherever the request entered; a failure there never fails a decision, it answers raw with `temperatures: null`. **A decision model warms on bind** (#277): `models/api_routes.py::_wait_for_bind` hands every bound engine to `decide.py::schedule_decision_warmup`, which sends one constrained question per kind through `ask_one` in the background when the directory carries `prompt_contract.json` or `temperatures.json`, and `/api/status` reports `warm` per instance (null for a model with nothing to warm). An engine adopted at restart never binds, so adoption warms too: the boot primary in `_await_primary_bind` and every adopted stacked instance in `_warm_adopted_stacked`. Readiness is not held for it. A connected engine call that times out is reported as the grammar compiling, never as an unreachable node. **The owner answers** (`decide.py::forward_to_owner`): a request for a model this node does not serve is posted whole to the serving node's AINode (same path, fleet key, `X-AINode-Forwarded-By`), replicas tried in routing order past a transport error, a 5xx or a 401/403/404; any other status is the owner's own answer and is returned as it came. A forwarded request is never forwarded again and is not counted by the owner's rate limiter (`ratelimit/middleware.py::is_forwarded_by_fleet`, fleet key only), because the node the caller reached already counted it. Only when no owner answers does this node call the engines itself, with the fetched table when it can get one. A model served locally never forwards. **A decision model warms on bind** (#277): `models/api_routes.py::_wait_for_bind` hands every bound engine to `decide.py::schedule_decision_warmup`, which sends one constrained question per kind through `ask_one` in the background when the directory carries `prompt_contract.json` or `temperatures.json`, and `/api/status` reports `warm` per instance (null for a model with nothing to warm). Readiness is not held for it. A connected engine call that times out is reported as the grammar compiling, never as an unreachable node. - **`/v1/systemone` is TypeSafe's wire format and nothing else** (`api/systemone.py`). It exists so a client written for the hosted System One endpoint (Titanium's JDE `jevJudge({endpoint, model})`, browser-use's jev-ultrafast, the TypeSafe SDK, the playground) answers off a model on this fleet with its endpoint changed and nothing else, which makes the request and response shapes THEIRS and not ours to tidy: `{state, model, questions}` in, `{model, answers, usage: {input_tokens, output_tokens}, latency_ms}` out (plus AINode's `calibration` block, #276), a `choice` answering with the caller's own criteria key, a `noul` answering P(true), a `score` answering the expected level with a `legend` from position to level name, a `422` naming the field on a malformed request and a `503` when no node serves the model. JDE's `parseAnswers` discards the WHOLE answer set over one answer it cannot read, so every probability is finite and inside [0, 1] (`systemone.py::probability`) and a question the engine left unanswered is a `503`, never a 200 carrying an invented option. **`confidence` is NOT the picked option's probability**: the hosted service reports it chance corrected, `(n * p_max - 1) / (n - 1)`, which is INFERRED from every example in TypeSafe's published docs and SDK types (`systemone.py::normalized_confidence`) because no document states it, and a caller's bands are tuned against those numbers; the raw distribution goes out untouched beside it, so the formula is revisable against evidence and the numbers it came from are never lost. **A question may carry at most `TOP_LOGPROBS` criteria**, which is the engine request's ceiling and not the format's 255: only the top 20 labels come back with a probability, so a wider option set is a `422` naming the cap rather than a distribution missing its tail, and raising it means a second pass over the remaining labels, which is a measurement and not a constant. Translate in and out around `run_questions`: a gap here is never answered with a second engine path, a second routing rule or a second reading of the logprobs. **Calibration is the model's**: the one correction on this path is the served adapter's own `temperatures.json`, applied in `run_questions` as the `/v1/decide` bullet says, and nothing here adds a correction of its own. The way to find out what a local model's confidence is worth stays `scripts/ainode-bench.py decide`. `tests/test_systemone.py` carries JDE's own reader rule for rule and replays a real case from its blind set, read where it lives (`JDE_COMPLETION_CASES`, default `/Users/sem/code/jde/cases/`) and never copied in: it is someone else's measurement data, and the test skips where it is absent. - **A multimodal chat request consults the capability cache before routing** (`api/server.py::proxy_to_vllm`). A body carrying an `image_url` / `input_audio` / `video_url` / `file` part, or the Anthropic Messages spelling (an `image` / `document` block, including one nested in a `tool_result`), is never routed on the model id alone: order candidates accepting-first (cached `vision: true`), then never-probed, and drop instances cached `vision: false`; vLLM's `may be provided in one prompt` 400 is a routing miss, so record `vision: false` and fail over, while every other 4xx goes back to the caller untouched. A request with no media keeps the plain order (local hop first, then peers). There is ONE capability cache: `app["chat_caps_cache"]`, filled by `/api/models/caps` in `api/chat_routes.py`, which probes remote instances directly on their engine port. Never add a second. - **Every node-to-node request AINode makes carries the FLEET KEY, and there is one helper that puts it there** (`ainode/auth/fleet.py`: `fleet_headers(app)` off a running app, `fleet_key_headers(secret)` off a config). The key is `HMAC-SHA256(cluster_secret, "ainode-fleet-key-v1")`, so every node holding the secret computes the same one, the join flow already distributes it, rotation follows the secret, and NOTHING new is written to disk. The middleware accepts it as the caller id `fleet` (`auth/middleware.py::identify_caller`), reading the secret LIVE per request, so a node whose secret differs refuses its would-be peers exactly as it drops their datagrams. Auth was enabled-able and unusable on a cluster before this: a node with `auth.enabled` answered 401 to its own fan-outs, so the fleet ran open. Add a peer call and you add the header: an ENGINE port is not a peer call in this sense (the inference proxy, the capability probes and the embeddings route talk to a vLLM container, which never sees this middleware), and `POST /api/cluster/join` is keyless by construction. `tests/test_fleet_auth.py` WALKS THE SOURCE for peer URLs and fails on one whose function does not name the helper, with an exempt list that has to state a reason. diff --git a/README.md b/README.md index 4c5756bb..9118dedb 100644 --- a/README.md +++ b/README.md @@ -1294,12 +1294,15 @@ response says what was applied, and `POST /v1/systemone` carries the same block Send `"calibration": "raw"` to get the engine's own spread instead; the block then reads `"applied": false` and still lists the temperatures you opted out of, -so you can refit against the raw numbers. The file lives on the node that serves -the model. A node that routes the request there without a copy of its own (the -master, usually) asks that node for its table over the fleet key -(`GET /api/decide/calibration?model=`) and caches it for five minutes, so the -answer is tempered the same wherever the request enters. A peer on an older -release has no such route, and then the answer is raw with `"temperatures": null`. +so you can refit against the raw numbers. The node that serves the model is the +one that answers: a request that reaches a node without that model (the master, +usually) is handed whole to the serving node's AINode over the fleet key, so the +answer carries that node's temperatures and warm grammar wherever the request +entered. If the first node serving it is down the next replica answers; if no +node's AINode answers, the receiving node calls the engine itself and tempers +with the table it fetches from that node (`GET /api/decide/calibration?model=`, +cached five minutes); only when that fails too is the answer raw with +`"temperatures": null`. **A decision model warms up when it loads.** The first constrained request per question shape makes the engine compile the answer grammar, 60 to 90 s on a diff --git a/ainode/api/decide.py b/ainode/api/decide.py index 0bdbc5bd..7747b2f3 100644 --- a/ainode/api/decide.py +++ b/ainode/api/decide.py @@ -682,6 +682,103 @@ def candidates_for(request: web.Request, model: str) -> list: return candidates +# ------------------------------------------------------- the owner answers +# +# A decision model's temperatures.json and its warm state live on the node that +# serves it. When a request enters on a node that does not serve the model (the +# master, above all, which is the endpoint clients use), that node used to call +# the remote ENGINE itself, so the answer came back raw with +# ``calibration: {applied: false, temperatures: null}`` while the same request +# sent to the owner was tempered (2026-09-26, jebadiah-9b-v2 through Spark-1). +# Now the whole request goes to the owner's AINode, on the same route and under +# the fleet key, and the owner answers it exactly as if the caller had come +# straight to it. A request served locally takes the local path, unchanged. + +#: How long one owner has to answer a forwarded decision: its own engine limit, +#: plus room for the hop. +FORWARD_TIMEOUT_S = CALL_TIMEOUT_S + 30.0 + + +def owner_web_ports(app, candidates: list) -> list[tuple[str, int]]: + """``(host, web_port)`` of each remote node in *candidates*, in routing order. + + A candidate is an engine ``(fabric_ip, api_port)``, and a node with two + replicas stacked on it appears once. The web port is the one the node + announces, which is not always 3000 (Atlas serves on 3100). + """ + cluster = app.get("cluster_state") + by_host: dict = {} + for member in (cluster.members() if cluster is not None else []): + host = getattr(member, "fabric_ip", "") or "" + if host and host not in by_host: + by_host[host] = member + owners: list[tuple[str, int]] = [] + for host, _port in candidates: + if host == "localhost": # _routing_candidates' name for this node + continue + member = by_host.get(host) + if member is None: + continue + owner = (host, int(getattr(member, "web_port", 3000) or 3000)) + if owner not in owners: + owners.append(owner) + return owners + + +async def forward_to_owner(request: web.Request, body: dict, + candidates: list) -> Optional[web.Response]: + """Hand the request to the node that serves the model, and return its answer. + + None means "answer here": the model is served on this node (the local path + is unchanged), the request was itself forwarded (it is never forwarded + twice), no owner can be named, or no owner answered, in which case the old + path calls the engines directly and the answer is raw rather than missing. + + Owners are tried in routing order. A transport failure, a timeout, a 5xx + (including an owner whose engine is down), or a 401/403/404 (a key it does + not share, or a release without the route) moves on to the next replica. + Any other status is the owner's own answer, a 4xx naming a field included, + and is returned as it came. + """ + from ainode.auth.fleet import FORWARDED_BY_HEADER, fleet_headers + if request.headers.get(FORWARDED_BY_HEADER): + return None + if any(host == "localhost" for host, _ in candidates): + return None + app = request.app + owners = owner_web_ports(app, candidates) + session = app.get("client_session") + if not owners or session is None: + return None + config = app.get("config") + headers = fleet_headers(app, { + "Content-Type": "application/json", + FORWARDED_BY_HEADER: str(getattr(config, "node_id", "") or "peer"), + }) + timeout = aiohttp.ClientTimeout(total=FORWARD_TIMEOUT_S, + sock_connect=CONNECT_TIMEOUT_S) + payload = json.dumps(body).encode() + tried: list[str] = [] + for host, web_port in owners: + url = f"http://{host}:{web_port}{request.path}" + try: + async with session.post(url, data=payload, headers=headers, + timeout=timeout) as resp: + answer = await resp.read() + if resp.status >= 500 or resp.status in (401, 403, 404): + tried.append(f"{host}:{web_port} {resp.status}") + continue + return web.Response(status=resp.status, body=answer, + content_type="application/json") + except asyncio.CancelledError: + raise + except Exception as exc: + tried.append(f"{host}:{web_port} {type(exc).__name__}") + logger.warning("no owner of %s answered %s (%s); calling the engines directly", + body.get("model"), request.path, "; ".join(tried)) + return None + + async def ask_one(session: aiohttp.ClientSession, candidates: list, body: dict, timeout_s: float = CALL_TIMEOUT_S) -> tuple: """One question, with the proxy's failover. Returns (payload, cand, latency_ms). @@ -1001,6 +1098,11 @@ async def handle_decide(request: web.Request) -> web.Response: if not candidates: return unavailable(f"no node is serving '{model}'") + # A model another node serves is answered by that node, calibration and all. + forwarded = await forward_to_owner(request, dict(body, model=model), candidates) + if forwarded is not None: + return forwarded + run = await run_questions(request, model, questions, state, instructions, candidates, calibration) collector = request.app.get("metrics_collector") diff --git a/ainode/api/systemone.py b/ainode/api/systemone.py index 1589de6f..381b3e31 100644 --- a/ainode/api/systemone.py +++ b/ainode/api/systemone.py @@ -72,6 +72,7 @@ DecideError, calibration_mode, candidates_for, + forward_to_owner, merge_usage, normalize_questions, resolve_model, @@ -485,6 +486,11 @@ async def handle_systemone(request: web.Request) -> web.Response: if not candidates: return unavailable(f"no node is serving '{model}'") + # A model another node serves is answered by that node, calibration and all. + forwarded = await forward_to_owner(request, dict(body, model=model), candidates) + if forwarded is not None: + return forwarded + # No shared instructions block: in this format a question's own instructions # are the whole prompt for it, and the questions of one ask still never see # each other's answers. diff --git a/ainode/auth/fleet.py b/ainode/auth/fleet.py index a18c06e1..db1779b7 100644 --- a/ainode/auth/fleet.py +++ b/ainode/auth/fleet.py @@ -59,6 +59,12 @@ #: dashboard: revoking the fleet's access to a node means changing that node's #: ``cluster_secret``, which is the same act as removing it from the cluster. FLEET_KEY_ID = "fleet" +#: Set by a node that hands a decision request (/v1/decide, /v1/systemone) to the +#: node that owns the model. The owner answers it itself and never forwards it +#: again, and its rate limiter does not count it: the node the caller reached has +#: already counted the caller, and every forwarded request arrives under the one +#: fleet key, so counting them here would make the whole fleet one client. +FORWARDED_BY_HEADER = "X-AINode-Forwarded-By" def fleet_key(secret: Optional[str]) -> str: diff --git a/ainode/ratelimit/middleware.py b/ainode/ratelimit/middleware.py index 04e965d9..aa58da87 100644 --- a/ainode/ratelimit/middleware.py +++ b/ainode/ratelimit/middleware.py @@ -278,6 +278,19 @@ def client_key(request) -> str: return f"ip:{remote or 'unknown'}" +def is_forwarded_by_fleet(request) -> bool: + """A decision request another node forwarded here under the fleet key. + + Only the fleet key stamps ``api_key_id == "fleet"``, so a client cannot claim + the exemption by sending the header. See ``auth/fleet.py::FORWARDED_BY_HEADER``. + """ + from ainode.auth.fleet import FLEET_KEY_ID, FORWARDED_BY_HEADER + getter = getattr(request, "get", None) + key_id = getter("api_key_id", "") if callable(getter) else "" + headers = getattr(request, "headers", None) or {} + return key_id == FLEET_KEY_ID and bool(headers.get(FORWARDED_BY_HEADER)) + + def too_many_requests(decision: Decision, config: RateLimitConfig) -> web.Response: """The 429: a Retry-After header and a body that names the limit.""" if decision.limit == "max_inflight": @@ -319,6 +332,8 @@ async def rate_limit_middleware(request: web.Request, handler): limiter: Optional[RateLimiter] = request.app.get("rate_limiter") if limiter is None or not limiter.enabled or not is_limited_path(request.path): return await handler(request) + if is_forwarded_by_fleet(request): + return await handler(request) key = client_key(request) decision = limiter.admit(key) if not decision.allowed: diff --git a/tests/test_decide_forward.py b/tests/test_decide_forward.py new file mode 100644 index 00000000..7add1ce2 --- /dev/null +++ b/tests/test_decide_forward.py @@ -0,0 +1,249 @@ +"""A decision request is answered by the node that owns the model. + +The owner holds the model's ``temperatures.json`` and its warm state. A node that +does not serve the model (the master, which is the endpoint clients use) used to +call the owner's ENGINE itself and answered raw with ``temperatures: null``, +while the same request sent to the owner was tempered (2026-09-26, +jebadiah-9b-v2 through Spark-1). The receiving node now hands the whole request +to the owner's AINode under the fleet key. +""" + +from types import SimpleNamespace + +import pytest +import pytest_asyncio +from aiohttp import web +from aiohttp.test_utils import TestClient, TestServer + +from ainode.api import decide +from ainode.auth.fleet import FORWARDED_BY_HEADER +from ainode.ratelimit.middleware import is_forwarded_by_fleet +from tests.test_decide import MODEL, QUESTIONS, TICKET, FakeEngine, _app, _free_port +from tests.test_decide_calibration import TEMPS, _store +from tests.test_systemone import QUESTIONS as JEV_QUESTIONS + +SECRET = "fleet-secret-for-tests" + + +def _decide_body(**over): + body = {"model": MODEL, "state": TICKET, "questions": QUESTIONS} + body.update(over) + return body + + +def _jev_body(**over): + body = {"model": MODEL, "state": TICKET, "questions": JEV_QUESTIONS} + body.update(over) + return body + + +@pytest.fixture(autouse=True) +def _fresh_peer_cache(): + """The fallback path reads PR 284's peer-table cache; no test may inherit one.""" + decide._peer_temperatures.clear() + yield + decide._peer_temperatures.clear() + + +@pytest_asyncio.fixture +async def engine(): + fake = FakeEngine() + server = TestServer(fake.app()) + await server.start_server() + try: + yield fake, server.port + finally: + await server.close() + + +def _node(engine_port, models_dir, web_port=None): + """An AINode whose cluster view says a peer at 127.0.0.1 serves MODEL.""" + app = _app(engine_port) + app["config"].models_dir = str(models_dir) + app["config"].cluster_secret = SECRET + if web_port is not None: + for member in app["cluster_state"].members(): + member.web_port = web_port + return app + + +@pytest_asyncio.fixture +async def owner(engine, tmp_path): + """The node that holds MODEL's temperatures, as a real AINode app.""" + _, engine_port = engine + store = tmp_path / "owner-store" + store.mkdir() + _store(store) + server = TestServer(_node(engine_port, store)) + await server.start_server() + try: + yield server + finally: + await server.close() + + +async def _entry(engine_port, owner_port, tmp_path): + """The node the caller reached: no copy of the model, routes to the owner.""" + store = tmp_path / "entry-store" + store.mkdir(exist_ok=True) + client = TestClient(TestServer(_node(engine_port, store, web_port=owner_port))) + await client.start_server() + return client + + +@pytest.mark.asyncio +async def test_the_master_answers_with_the_owner_s_calibration(engine, owner, tmp_path): + fake, engine_port = engine + client = await _entry(engine_port, owner.port, tmp_path) + try: + resp = await client.post("/v1/decide", json=_decide_body()) + assert resp.status == 200 + assert (await resp.json())["calibration"] == {"applied": True, + "temperatures": TEMPS} + + resp = await client.post("/v1/systemone", json=_jev_body()) + assert resp.status == 200 + assert (await resp.json())["calibration"] == {"applied": True, + "temperatures": TEMPS} + + resp = await client.post("/v1/decide", json=_decide_body(calibration="raw")) + assert (await resp.json())["calibration"] == {"applied": False, + "temperatures": TEMPS} + finally: + await client.close() + assert fake.seen, "the owner reached its engine" + + +@pytest.mark.asyncio +async def test_the_owner_s_own_refusal_comes_back_as_it_is(engine, owner, tmp_path): + _, engine_port = engine + client = await _entry(engine_port, owner.port, tmp_path) + try: + resp = await client.post("/v1/systemone", json=_jev_body(calibration=1)) + assert resp.status == 422 + assert "'calibration'" in (await resp.json())["error"]["message"] + finally: + await client.close() + + +class Owner: + """An owner's web port that records what it was sent and answers as told.""" + + def __init__(self, status=200): + self.status = status + self.seen: list = [] + + def app(self): + app = web.Application() + app.router.add_post("/v1/decide", self.handle) + app.router.add_post("/v1/systemone", self.handle) + return app + + async def handle(self, request): + self.seen.append((request.path, dict(request.headers), await request.json())) + if self.status != 200: + return web.json_response({"error": {"message": "engine down"}}, + status=self.status) + return web.json_response({"model": MODEL, "node": "owner", "decisions": {}, + "calibration": {"applied": True, + "temperatures": TEMPS}}) + + +@pytest.mark.asyncio +async def test_a_down_owner_fails_over_to_the_next_replica(engine, tmp_path, + monkeypatch): + _, engine_port = engine + broken, healthy = Owner(status=503), Owner() + servers = [TestServer(broken.app()), TestServer(healthy.app())] + for s in servers: + await s.start_server() + dead_port = _free_port() # nothing listens: a node that is gone + monkeypatch.setattr(decide, "owner_web_ports", lambda _app, _c: [ + ("127.0.0.1", dead_port), ("127.0.0.1", servers[0].port), + ("127.0.0.1", servers[1].port)]) + client = await _entry(engine_port, servers[1].port, tmp_path) + try: + resp = await client.post("/v1/decide", json=_decide_body()) + assert resp.status == 200 + data = await resp.json() + assert data["node"] == "owner" + finally: + await client.close() + for s in servers: + await s.close() + assert len(broken.seen) == 1 and len(healthy.seen) == 1 + path, headers, body = healthy.seen[0] + assert path == "/v1/decide" + assert headers.get("Authorization", "").startswith("Bearer "), "the fleet key" + assert headers.get(FORWARDED_BY_HEADER) == "local-node" + assert body["model"] == MODEL + + +@pytest.mark.asyncio +async def test_no_owner_answering_falls_back_to_the_engines(engine, tmp_path): + """A peer on a release without the route, or every owner down: answer raw + off the engine rather than not at all.""" + fake, engine_port = engine + old = Owner(status=404) + server = TestServer(old.app()) + await server.start_server() + client = await _entry(engine_port, server.port, tmp_path) + try: + resp = await client.post("/v1/decide", json=_decide_body()) + assert resp.status == 200 + assert (await resp.json())["calibration"] == {"applied": False, + "temperatures": None} + finally: + await client.close() + await server.close() + assert len(old.seen) == 1 + assert fake.seen, "the fallback called the engine directly" + + +@pytest.mark.asyncio +async def test_a_forwarded_request_is_answered_where_it_lands(engine, tmp_path): + """Never forwarded twice: the owner answers even if its own view names + somebody else.""" + fake, engine_port = engine + never = Owner() + server = TestServer(never.app()) + await server.start_server() + client = await _entry(engine_port, server.port, tmp_path) + try: + from ainode.auth.fleet import fleet_key + resp = await client.post("/v1/decide", json=_decide_body(), headers={ + FORWARDED_BY_HEADER: "spark-1", + "Authorization": f"Bearer {fleet_key(SECRET)}"}) + assert resp.status == 200 + finally: + await client.close() + await server.close() + assert never.seen == [] + assert fake.seen + + +@pytest.mark.asyncio +async def test_a_model_served_here_takes_the_local_path(): + request = SimpleNamespace(headers={}, app={}, path="/v1/decide") + assert await decide.forward_to_owner( + request, {"model": MODEL}, [("localhost", 8002), ("10.0.0.5", 8002)]) is None + + +def test_owner_ports_come_from_what_each_node_announces(): + members = [SimpleNamespace(fabric_ip="10.100.0.17", web_port=3000), + SimpleNamespace(fabric_ip="10.100.0.20", web_port=3100)] + app = {"cluster_state": SimpleNamespace(members=lambda: members)} + cands = [("localhost", 8000), ("10.100.0.17", 8002), ("10.100.0.17", 8003), + ("10.100.0.20", 8100), ("10.9.9.9", 8000)] + assert decide.owner_web_ports(app, cands) == [("10.100.0.17", 3000), + ("10.100.0.20", 3100)] + + +def test_only_the_fleet_key_skips_the_owner_s_rate_limit(): + forwarded = {FORWARDED_BY_HEADER: "spark-1"} + fleet = SimpleNamespace(headers=forwarded, get=lambda k, d="": "fleet") + client_claiming = SimpleNamespace(headers=forwarded, get=lambda k, d="": "k-123") + plain_fleet = SimpleNamespace(headers={}, get=lambda k, d="": "fleet") + assert is_forwarded_by_fleet(fleet) + assert not is_forwarded_by_fleet(client_claiming) + assert not is_forwarded_by_fleet(plain_fleet)