From 18aced01bf6585b4848d0da913f9c81c892fd3c8 Mon Sep 17 00:00:00 2001 From: webdevtodayjason Date: Fri, 25 Sep 2026 20:09:30 -0500 Subject: [PATCH] Make a master that serves no model safe (Atlas control plane, phase 1) Four fixes so AINode can run as a routing-only master on Atlas: - A request with no model field no longer falls back to config.model on a node that serves nothing; it answers 400 missing_model_field. A node that serves a primary keeps its default. - Unloading or ejecting the primary never promotes a survivor. The slot and its per-load overrides are cleared and stacked instances stay stacked (Whisper was promoted into Spark-4's config.model with another model's image and flags). - ainode doctor passes on a node that loads no model: "no model loaded" is INFO, and docker and engine backend findings are INFO there. - An empty account list is never replicated: the master's push, the worker's pull, the CLI push and /api/auth/users/sync all refuse it and log. --- AGENTS.md | 8 +- ainode/api/server.py | 49 ++++++++++- ainode/api/server_routes.py | 16 ++-- ainode/auth/replication.py | 46 ++++++++++ ainode/auth/session_routes.py | 11 +++ ainode/cli/doctor.py | 70 +++++++++++++-- ainode/models/api_routes.py | 42 ++++++--- tests/test_auth_replication.py | 72 ++++++++++++++++ tests/test_doctor.py | 8 +- tests/test_doctor_no_model.py | 136 ++++++++++++++++++++++++++++++ tests/test_primary_release.py | 150 +++++++++++++++++++++++++++++++++ tests/test_proxy_no_model.py | 114 +++++++++++++++++++++++++ tests/test_session_routes.py | 23 +++++ 13 files changed, 709 insertions(+), 36 deletions(-) create mode 100644 tests/test_doctor_no_model.py create mode 100644 tests/test_primary_release.py create mode 100644 tests/test_proxy_no_model.py diff --git a/AGENTS.md b/AGENTS.md index 34fb7dcf..a571cddd 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -27,7 +27,7 @@ State / architecture / decisions / "why": Obsidian Vault, `AINode` (cluster ops: - **The decision bench runs TWO measurements under one subcommand, and a record says which one produced it.** The **Jevals recipe** (`ainode/bench/decide/{jevals,sets,suite}.py`, `--suite`/`--questions` plus `--transport decide|systemone`) scores the public question sets the independent boards use, with their formulas, so an AINode-served model can be read next to Jev and its clones; the **legacy 110-item path** (`--backend`) is the bullet below. Mixing the two flag sets is an error rather than a guess. **`bench/decide/JEVALS.md` is authoritative for what a recipe number means**: it records every formula against the page and the date it was read (jevals.com/methodology and jevals.com/policy, JevBench's README, the typed-decisions card, all 2026-09-21) and names every deviation, and a formula change here is a change there. Rules with no exemption: **the model is never shown the answer** (labels live in a question file's separate `labels` map, a question carrying an answer-key-shaped field is a load error, and `suite.wire_leaks` re-checks the assembled body and REFUSES to send rather than producing a flattering score, skipping only the caller's own `state`); **every state is verified against its `state_sha256` on load**, which is what makes a number comparable to a board number and is how the three Jevals state constructions were recovered in the first place (`{"message": ...}` for Banking77, not `text`; `{"question": ..., "context": [...]}` for PubMedQA, not `context.contexts`); **item text is never committed**, only the manifests under `bench/decide/sets/`, because upstream does not republish it either and the upstream licences are not ours to relicense (`decide download` writes the gitignored `bench/decide/cache/`); **every metrics block names its recipe** (`recipe`, `recipe_of_record`, `probability_source` of `logprob`/`native`/`verbalized`) because two boards use one word for different arithmetic; **a set holding more than one ANSWER SPACE is broken down by answer space** (one `(type, options)` pair) with a per-primitive roll-up, and its Decision Score is the mean over the spaces, because the label prior is the base rates of the labels in ONE option list and pooling two of them builds the baseline that defines 0 over an answer space neither question has (`typed-decisions` asks `action` with four options in one workflow and five in another, so this is a wrong number and not a presentation choice); **a set in a listed system's published training data carries that finding with its primary source** (`sets.CONTAMINATION`); and **`decide.overall` on such a record carries no `brier`, `ece`, `bins` or `thresholds`**, because those names mean the legacy definitions and the table renders them "not measured" rather than borrowing a number computed another way. `tests/test_bench_decide_jevals.py` pins the three anchors of the scale against hand arithmetic (a calibrated system that is exactly the prior scores 0, a confidently wrong one scores negative, a perfect one scores 100), both transports against a loopback server, and that no wire body ever carried an answer key. - **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. +- **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 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. @@ -35,7 +35,7 @@ State / architecture / decisions / "why": Obsidian Vault, `AINode` (cluster ops: - **A fresh install requires a key, and an update never changes an installed node's access control** (`scripts/install.sh`). The installer mints one key on a node with no `config.json`, stores the SHA-256 the way `AuthConfig` does, prints the plaintext ONCE in a box, and the summary reads "API protected, one key"; `AINODE_AUTH=off` keeps the old open behaviour and prints what that choice means. Both gates matter: minting on an existing home would lock out every client the operator already pointed at the node, with a key they never saw. Because of this, `ainode update` verifies the running release on `/api/health` and not `/api/status` (health is the one keyless route, which is why it carries `version`): reading a keyed route there made every update on a protected node pull, pin, restart and then report that it had not applied. - **`auth.json` and `config.json` are 0600, written temp-then-replace, and re-read when they change.** Both carry credentials (`config.json` holds `cluster_secret` and may hold `hf_token`) and both were written under the default umask as root. `AuthConfig.reload_if_changed()` stats the store per request the way `ClusterSecret` stats the config per datagram, which is what makes `ainode auth enable|disable|key create|key revoke` LIVE on a running node (it used to need a restart nobody was told to do, so the CLI said "Auth enabled" about a node that went on answering everybody). A file that cannot be stat'ed or parsed keeps the state in memory: the tolerant direction is never "let everybody in". A credential-shaped config key goes in `api/server.py::SCRUBBED_CONFIG_KEYS` and is masked by `ainode config`, or it leaves the node in a `GET /api/config`. - **A person logs in, a machine presents a key, and the two stores sit side by side** (`ainode/auth/accounts.py`, `auth/session_routes.py`, #261). `users.json` is 0600, written temp-then-replace and re-read per request exactly as `auth.json` is, because `ainode auth user add`, a password change and a revoked session all have to land on a running node. The ONLY crypto is the standard library's: `hashlib.scrypt` (n=2**14, r=8, p=1) over a per-password salt, and a session token stored as its SHA-256 alone, so the file a backup or a support bundle picks up logs nobody in and a login page is not the reason the image grows an argon2 wheel. **A session does not expire** (in until you log out), which makes REVOCATION the whole control and it has to be complete: a password change, a disable and a removal each take that account's open sessions with them, and `session_for_token` refuses one whose account has since been removed or disabled. **A cookie-authenticated request whose method is not GET, HEAD or OPTIONS must carry `X-AINode-Client: dashboard`**, and that header is the whole CSRF defence: a browser attaches a cookie to a cross-site request on its own, and a custom header cannot be set cross-origin without a preflight this node never approves, so dropping the rule would make every node an operator is logged into writable by any page on the internet. A Bearer token is tried FIRST and the cookie only when no key matched, so automation is untouched and a stale key in a browser does not cost a good session its request. Administration (`/api/auth/users*`, `/api/auth/keys*`, `/api/auth/enable|disable`) goes through `accounts.require_admin`: an admin session, an operator key, or the fleet key, never a member; the two replication routes are the FLEET key alone because they move password hashes; and the last ENABLED admin cannot be removed or disabled, enforced in the store (`UsersStore.is_last_admin`) so the CLI cannot get around it. `/api/auth/login` and `/api/auth/me` are the newest entries in `SKIP_PATHS`, and adding one there means naming it in the middleware docstring, `api/cluster_join.py`'s docstring, the README's keyless paragraph and the dashboard's API access panel, because four tests assert each of those. -- **Dashboard ACCOUNTS have one authority, the master, and SESSIONS are never replicated** (`ainode/auth/replication.py`). A cluster cannot hold one account list per node: an operator who adds a login on the master and then opens the dashboard of whichever node a bookmark points at would be told their password is wrong. So the master pushes its `export_users()` to `/api/auth/users/sync` on every peer after every change (`app["users_changed"]`, with `fleet_headers` like every other node-to-node call, because the export carries password hashes), retrying a peer that refused on the next change and on a 60 second tick while it is behind; a worker pulls `/api/auth/users/export` at startup and every 5 minutes, which is the safety net a push cannot be (a node that was down, or one that has just joined and has no accounts at all). "Behind" is a comparison and not a queue: the master remembers the stamp each peer accepted. A session is one browser's credential against ONE node, so nothing copies them and `ainode auth session ...` is a per-node command; a node with no `cluster_secret` or nobody to talk to does nothing at all, because it could not authenticate to a peer anyway. The FIRST account is made on the box (`ainode auth user add --admin`), never over HTTP: a route that mints the first admin is a route anybody who can reach the port calls before the operator does. `--password-stdin` is not optional garnish, it is the only path with no TTY, and the installer's wrapper runs `docker exec -it`. A worker records the master's own stamp in `users-sync.json` so `ainode doctor`'s `login.sync` compares two values of the same kind rather than two hashes computed by different code. Pinned by `tests/test_auth_replication.py` and `tests/test_cli_auth_users.py`. +- **Dashboard ACCOUNTS have one authority, the master, and SESSIONS are never replicated** (`ainode/auth/replication.py`). A cluster cannot hold one account list per node: an operator who adds a login on the master and then opens the dashboard of whichever node a bookmark points at would be told their password is wrong. So the master pushes its `export_users()` to `/api/auth/users/sync` on every peer after every change (`app["users_changed"]`, with `fleet_headers` like every other node-to-node call, because the export carries password hashes), retrying a peer that refused on the next change and on a 60 second tick while it is behind; a worker pulls `/api/auth/users/export` at startup and every 5 minutes, which is the safety net a push cannot be (a node that was down, or one that has just joined and has no accounts at all). "Behind" is a comparison and not a queue: the master remembers the stamp each peer accepted. A session is one browser's credential against ONE node, so nothing copies them and `ainode auth session ...` is a per-node command; a node with no `cluster_secret` or nobody to talk to does nothing at all, because it could not authenticate to a peer anyway. The FIRST account is made on the box (`ainode auth user add --admin`), never over HTTP: a route that mints the first admin is a route anybody who can reach the port calls before the operator does. `--password-stdin` is not optional garnish, it is the only path with no TTY, and the installer's wrapper runs `docker exec -it`. A worker records the master's own stamp in `users-sync.json` so `ainode doctor`'s `login.sync` compares two values of the same kind rather than two hashes computed by different code. **An EMPTY account list never crosses the wire**: `import_users` replaces the whole list, so the master's push, the worker's pull, the CLI push and `/api/auth/users/sync` (409 `empty_account_list`) each refuse one and log the fix, and removing the last account is a per-node act. Pinned by `tests/test_auth_replication.py` and `tests/test_cli_auth_users.py`. - **A discovery announcement is SIGNED when `cluster_secret` is set, and a node that has one drops what it cannot verify** (`discovery/signing.py`). The signature is one extra top-level key beside the payload's own fields, never a wrapper around them, so a peer that has never heard of it drops the unknown key in `from_json` and stays in the cluster view. Both directions read the secret per datagram through `ClusterSecret`, so rotating it is a `config.json` edit and not a fleet-wide restart: never capture the value at startup. A node with NO secret behaves exactly as it did and says so once per process, because nothing on the fleet sets one and a mandatory-signing release would partition every existing cluster on upgrade. - **The discovery port has one home: `core/config.py::DEFAULT_DISCOVERY_PORT` (5679).** `NodeConfig`, both discovery classes and the installer read it from there. They each carried their own literal, and the dataclass said 5678 while the installer, the fleet and the docs said 5679, so a source install listened where nobody spoke and vanished from every cluster view with nothing logged (#181). The port, the cluster id and whether the wire is signed are logged together at startup. - **`ainode_version` travels on the announcement and into every view that lists nodes** (`/api/nodes`, `/api/cluster/resources`, `/api/version/check`, `/api/cluster/update-status`). A peer that announces none reports `""` and is counted as unknown: filling it in with the local version is how a split fleet goes on looking like a healthy one (#171). @@ -50,9 +50,9 @@ State / architecture / decisions / "why": Obsidian Vault, `AINode` (cluster ops: - **The Metrics view draws a gap as a gap** (`web/static/js/metrics-data.js` for the shaping, `metrics.js` for the canvases). It reads `/api/metrics/history` and nothing else: the ring buffer that app.js appended live polls to is gone, because a buffer plus a store is two sources for one line and the range picker asks for windows a 20 minute buffer cannot answer. A null is a HOLE in the line, never a point at the axis, never the previous sample carried forward, and a series with nothing measured in the whole window is said in words rather than drawn (a GB10 exposes no utilisation counter, so that panel says so, which is the same rule as the two bullets above). The counters (`requests.total`, `requests.errors`, `requests.tokens_generated`) are drawn as rates, and the interval across a restart, where the counter went backwards, is null: never a negative rate and never a zero. The shaping is DOM-free so it runs under node, and the picker's steps are pinned against the store's own budget, because the store splits `MAX_HISTORY_POINTS` across every series asked for and would answer a too-large range at a coarser step than the picker claims. `/api/metrics/history?node=` answers for ANY node of the cluster by fetching THAT node's own history and passing it through unchanged: never merged, averaged or filled in, and a peer that cannot be reached is an error naming it rather than an empty grid, which would draw as a node that measured nothing. Pinned by `tests/test_metrics_chart.py`. - **Live load progress is assembled in one place**: `api/server.py::load_progress_fields` (with `load_started_at` / `expected_ready_minutes`) answers `load_started_at`, `load_elapsed_seconds` and `expected_ready_minutes` for `/api/status`, for each `/api/nodes` row and for the announcement, and the browser does the same math once in `app.js::loadProgress`. Elapsed is always the SERVER's arithmetic plus the time since that payload arrived: the browser, the head and a peer do not share a clock, so a remote start stamp subtracted from a local `now()` reports skew as progress. The expected time comes from the launch-time ledger first, then the catalog's `typical_ready_minutes`, then nothing at all: never a figure derived from the weight size. - **A download that cannot fit is refused before it starts** (`ainode/models/fit.py`, `api_routes::_refuse_if_it_will_not_fit`, 507 with `needed`, `free` and the PATH). A pull with nowhere to go dies part way and reads as an engine that crashed (#184 point 4). The size comes from the Hub's file metadata, then the safetensors dtype breakdown, and is remembered in `/repo-sizes.json`; a size nobody could learn is `None` and the download proceeds, because 0 would wave every pull through as fitting a full disk. Only a KNOWN size is cached, so a Hub outage is not frozen in for a month. The launch path asks with `learn=False` (cache and catalog only): a launch must not wait on an HTTP round trip to huggingface.co. `force` always gets past it. The suite never reaches the Hub or the real cache: `tests/conftest.py::no_hub_size_lookup` redirects both. -- **`ainode doctor` runs inside the container, so a check only the HOST can answer is INFO, never WARN** (`cli/doctor.py::running_in_container`, `host_service_state`). The documented deployment is the installer's wrapper running `docker exec ainode ainode ...`, so a WARN about there being no systemd in a container was a permanent yellow line on every node in the fleet, and a WARN nobody can clear teaches operators to ignore WARNs (#225). The wrapper's `doctor` case reads the unit state on the host and passes it in as `AINODE_HOST_SERVICE_STATE`, which makes `service.unit` a real answer (OK when active, WARN when a dead unit sits under a live container). Same rule for the docker, image-pin and engine-image checks when there is no docker socket in the container: INFO naming the host command. `--fix` is REFUSED with `--peer` rather than silently dropped. +- **`ainode doctor` runs inside the container, so a check only the HOST can answer is INFO, never WARN** (`cli/doctor.py::running_in_container`, `host_service_state`). The documented deployment is the installer's wrapper running `docker exec ainode ainode ...`, so a WARN about there being no systemd in a container was a permanent yellow line on every node in the fleet, and a WARN nobody can clear teaches operators to ignore WARNs (#225). The wrapper's `doctor` case reads the unit state on the host and passes it in as `AINODE_HOST_SERVICE_STATE`, which makes `service.unit` a real answer (OK when active, WARN when a dead unit sits under a live container). Same rule for the docker, image-pin and engine-image checks when there is no docker socket in the container: INFO naming the host command. **A node that loads no model under AINode passes** (`engine_expected`, `as_no_model_info`): "no model loaded" is INFO, and so are the engine-launch checks (docker, engine backend, engine image) when nothing is pinned, no stacked instance is in `instances.json` and no `distributed.json` is recorded, because a routing-only master launches nothing. `--fix` is REFUSED with `--peer` rather than silently dropped. - **`ainode update` verifies before it claims anything, and only then reclaims.** The host wrapper (`scripts/install.sh`) waits for this node's own `/api/status` to report the version it just installed and exits non-zero otherwise: it printed "Update complete" for an update that never applied. Only after that does it remove the images that release replaced, via `ainode prune-images` in the container that is now up. Prune decisions live in `core/image_prune.py` as a pure function of a `docker images` listing (`--images-from` replays another node's listing and touches nothing), removals go by `repo:tag` and never by image id (one image, three mirrored tags), only AINode's own app repos are ever considered, and `keep` generations of releases are kept so `ainode update ` has something to roll back to (#184). An update that cannot resolve a version from GHCR now exits NON-ZERO after the restart, because it pulled a floating tag and verified nothing, and "exit 0" is what a roll script reads as "this node is updated". `AINODE_IMAGE_REPOS` carries BOTH Docker Hub spellings on purpose: `argentaios/ainode` is the repository that exists there (stale at 0.4.7, no CI step pushes to it) and is the spelling every document uses, while `argentos/ainode` is the old mirror step's misspelling that is tagged locally on fleet nodes and would otherwise hold the bytes of every image the GHCR line thinks it removed. -- **Stacked instances replay only after the primary BINDS, and a dead engine's last lines are read before its container goes** (`models/api_routes.py::_await_primary_bind`, `_log_failed_engine_output`, `EngineBackend.log_tail`). vLLM sizes its KV cache from what is free when it profiles, so an engine that starts while the primary is still reserving is the one that dies; "started" is not the signal and neither is a clock that ran out (#235, #96). A primary that never binds does not cancel the replay, it is logged as what it is. Every bind window is multiplied by `bind_window_scale(loading)`, because concurrent loads share one node's bandwidth and memory. The tail is read BEFORE the relaunch, which stop/rm's the corpse by name: engines run without `--rm` precisely so there is one to read. +- **Stacked instances replay only after the primary BINDS, and a dead engine's last lines are read before its container goes** (`models/api_routes.py::_await_primary_bind`, `_log_failed_engine_output`, `EngineBackend.log_tail`). vLLM sizes its KV cache from what is free when it profiles, so an engine that starts while the primary is still reserving is the one that dies; "started" is not the signal and neither is a clock that ran out (#235, #96). A primary that never binds does not cancel the replay, it is logged as what it is. Every bind window is multiplied by `bind_window_scale(loading)`, because concurrent loads share one node's bandwidth and memory. The tail is read BEFORE the relaunch, which stop/rm's the corpse by name: engines run without `--rm` precisely so there is one to read. **Losing the primary never promotes a survivor** (`models/api_routes.py::release_primary`, used by unload and eject): `app["engine"]` and `config.model` go to None and the per-load overrides reset to defaults, and a stacked instance stays stacked on its own port. Promoting one made Whisper Spark-4's `config.model` with another model's image and flags (`tests/test_primary_release.py`). - **Cluster update job state is on disk** (`/cluster-updates.json`), because the master self-stops as the last step of the job it is reporting on (#182). Every mutation persists. Whether a node came BACK on the target version is answered at poll time from the version on the wire, never from the job record. - **Joining a cluster is a token, not a hand-edited `config.json`** (`cluster/join.py`, `api/cluster_join.py`). `ainode cluster token` mints 32 random bytes on the MASTER, stores only the SHA-256 hash in `/join-tokens.json` with an expiry and a single use, and prints the `ainode join` line to paste. `POST /api/cluster/join` is the only route besides `/api/health` and `/api/auth/status` exempt from the API-key rule, because a node that has not joined cannot hold this cluster's key: the token IS the credential, so a wrong, expired and spent token all get ONE identical 403 body (telling them apart makes the route an oracle) and the handler rate limits per source IP since there is no key in front of it. Never add a route that mints a token: minting stays on the box. - **A join writes its own keys into `config.json` and touches nothing else** (`cluster/join.py::merge_config_keys`). It merges into the raw JSON; a `NodeConfig.load(); setattr; save()` round trip rewrites the file from the dataclass, so it ADDS every key this release happens to default differently and DROPS every key the dataclass does not declare. A joiner is `cluster_role: "worker"` AND `distributed_mode: "member"`: those are the two spellings of one decision, and `member` is not a value `cluster_role` accepts (the `PATCH /api/config` validator rejects it). diff --git a/ainode/api/server.py b/ainode/api/server.py index d90f26a9..0c3a9a64 100644 --- a/ainode/api/server.py +++ b/ainode/api/server.py @@ -1887,6 +1887,40 @@ def _no_multimodal_instance(model: str, tried: list) -> web.Response: ) +def own_model(app, config) -> str: + """The model a request with no ``model`` field falls back to, or "". + + Only a node that SERVES a primary has one: ``config.model`` names what boots + on the node's own api_port, and ``app["engine"]`` is that engine. A config can + name a model the node does not serve (the dataclass default, a config.json + copied from another node, a primary whose launch failed), so the name alone is + not enough, and a routing-only master serves nothing at all. + """ + if app.get("engine") is None: + return "" + return (getattr(config, "model", "") or "").strip() + + +def _missing_json_model() -> web.Response: + """400 for a JSON request that names no model, on a node with none of its own. + + The old fallback routed it to ``config.model``, which on a node that serves + nothing is a name nobody is serving, and the caller saw a 404 or a 502 about a + model they never asked for. Names the field instead, the way the multipart + paths do. + """ + return web.json_response( + {"error": { + "message": ("this request names no model and this node serves none of " + "its own to default to: add a \"model\" field to the JSON " + "body (GET /v1/models lists what the fleet is serving)"), + "type": "invalid_request_error", + "param": "model", + "code": "missing_model_field"}}, + status=400, + ) + + def _missing_form_model() -> web.Response: """400 for a multipart request that never says which model to route to. @@ -1912,8 +1946,11 @@ async def proxy_to_vllm(request: web.Request) -> web.StreamResponse: config: NodeConfig = request.app["config"] session: aiohttp.ClientSession = request.app["client_session"] collector: MetricsCollector = request.app["metrics_collector"] - # Extract the model name first — it drives BOTH routing and metrics. - model = config.model or "unknown" + # Extract the model name first: it drives BOTH routing and metrics. A request + # that names none falls back to this node's OWN model, and only when it serves + # one (see own_model). + default_model = own_model(request.app, config) + model = default_model or "unknown" body_bytes = None body_obj: dict = {} if request.method == "POST": @@ -1938,9 +1975,15 @@ async def proxy_to_vllm(request: web.Request) -> web.StreamResponse: parsed = _json.loads(body_bytes) if isinstance(parsed, dict): body_obj = parsed - model = parsed.get("model", model) + named = parsed.get("model") + if isinstance(named, str) and named.strip(): + model = named except Exception: pass + if not isinstance(body_obj.get("model"), str) or not body_obj["model"].strip(): + if not default_model: + collector.record_request(model, 0.0, error=True) + return _missing_json_model() # Tag the request so the server-view log middleware can capture the model try: request["_log_model"] = model diff --git a/ainode/api/server_routes.py b/ainode/api/server_routes.py index d41692fd..8202d290 100644 --- a/ainode/api/server_routes.py +++ b/ainode/api/server_routes.py @@ -584,16 +584,18 @@ async def handle_server_eject(request: web.Request) -> web.Response: pass manager.remove(inst.record.instance_id) if request.app.get("engine") is inst.backend: - request.app["engine"] = None # the primary went away - # routing-truth: the node must stop claiming a model it no longer - # serves, or the master keeps advertising a ghost. + # The primary went away. routing-truth: the node must stop + # claiming a model it no longer serves, or the master keeps + # advertising a ghost, and its launch parameters go with it. + request.app["engine"] = None config = request.app.get("config") - if config is not None and getattr(config, "model", None) == model_id: - config.model = None + if config is not None: + from ainode.models.api_routes import release_primary try: - config.save() + release_primary(request.app, config) except Exception: - pass + logger.warning("eject: failed to persist the cleared primary", + exc_info=True) # Persist the shrunken instance set. Without this the eject was # memory-only: startup replay reads the manifest, so the ejected model # came BACK on the next reboot (spark-4, 2026-08-13 — an ejected 0.5B diff --git a/ainode/auth/replication.py b/ainode/auth/replication.py index eb083bfc..7df59b89 100644 --- a/ainode/auth/replication.py +++ b/ainode/auth/replication.py @@ -24,6 +24,13 @@ credential it never issued, and revoking a session would have to be a fan-out to be true. ``ainode auth session ...`` is therefore a per-node command, and that is the honest shape rather than a limitation. +* **An EMPTY account list never crosses the wire, in either direction.** The + receiving end REPLACES its list with what it is sent (``import_users``), so a + master whose ``users.json`` is missing or empty (a new master that was promoted + before its store was copied over) would log every dashboard user out of the + whole fleet on its first push. The master refuses to push it, a worker refuses + to pull it, and ``/api/auth/users/sync`` refuses to take it, each with a log + line naming the fix. Removing the last account is therefore a per-node act. * **A node with no ``cluster_secret``, or with nobody to talk to, does nothing.** The fleet key is derived from that secret (``ainode/auth/fleet.py``), so a node without one cannot authenticate to a peer and must not pretend to: it keeps the @@ -76,6 +83,12 @@ #: No hashes and no passwords go in here, so it is not a credential file. SYNC_STATE_NAME = "users-sync.json" +#: Why an empty list is not replicated. One sentence, shared by the master's push, +#: the worker's pull and the CLI, so the three log lines read the same. +EMPTY_LIST_REASON = "empty account store" +EMPTY_LIST_FIX = ("copy users.json from a node that has the accounts (or add one " + "with `ainode auth user add`) before this node replicates") + #: The fallback minimum password length, used only when the account store does not #: publish one of its own. The store is the authority and refuses a short password #: itself; the CLI asks for the number so it can say no at the prompt instead of @@ -440,6 +453,7 @@ def __init__(self, app, pull_interval: float = PULL_INTERVAL_SECONDS, self._wake = asyncio.Event() self._task: Optional[asyncio.Task] = None self._master_unreachable_logged = False + self._empty_list_logged = False self._next_pull = 0.0 # -- state ------------------------------------------------------------ @@ -483,6 +497,15 @@ async def broadcast(self) -> dict: logger.exception("could not export this node's accounts") return {"pushed": [], "failed": [], "reason": "export failed"} + if not users: + # import_users on the far side REPLACES its list, so this push would + # delete every account on every worker. Refuse, and say so once: the + # retry timer wakes this every 60 seconds. + self._refuse_empty("this master holds no accounts; not pushing an " + "empty list to the workers") + return {"pushed": [], "failed": [], "reason": EMPTY_LIST_REASON} + self._empty_list_logged = False + stamp = users_stamp(users) targets = peer_targets(self.app) # Forget peers that are no longer in the cluster view, so a node that left @@ -568,6 +591,13 @@ async def pull(self) -> dict: return self._master_unreachable(base, "answered a body with no user list") users = body["users"] + if not users: + # A master with an empty store (a new one promoted before users.json + # was copied to it) must not log everybody out of this node too. + self._refuse_empty(f"the master at {base} answered an empty account " + f"list; keeping this node's accounts") + return {"imported": False, "reason": EMPTY_LIST_REASON, "master": base} + self._empty_list_logged = False try: changed = bool(store.import_users(users)) except Exception: @@ -583,6 +613,14 @@ async def pull(self) -> dict: return {"imported": True, "changed": changed, "count": len(users), "stamp": stamp, "master": base} + def _refuse_empty(self, what: str) -> None: + """Log a refused empty list: a WARNING the first time, debug after.""" + if not self._empty_list_logged: + self._empty_list_logged = True + logger.warning("account replication refused: %s. %s", what, EMPTY_LIST_FIX) + else: + logger.debug("account replication still refused: %s", what) + def _master_unreachable(self, base: str, reason: str) -> dict: if not self._master_unreachable_logged: self._master_unreachable_logged = True @@ -837,6 +875,14 @@ def replicate_from_cli(config, users: list, info: Optional[dict] = None) -> dict tick, which is the mechanism that actually guarantees delivery. """ targets = peer_targets_from_info(info, config) or peer_targets_from_config(config) + if not users: + # The same refusal as the server's own push: the peers would REPLACE + # their lists with nothing. + logger.warning("account replication refused: no accounts to push. %s", + EMPTY_LIST_FIX) + return {"pushed": [], "failed": [], "peers": len(targets), + "reason": (f"{EMPTY_LIST_REASON}; the peers keep their accounts " + f"(remove the last account on each node by hand)")} headers = fleet_key_headers(getattr(config, "cluster_secret", "")) if not headers: return {"pushed": [], "failed": [], "peers": len(targets), diff --git a/ainode/auth/session_routes.py b/ainode/auth/session_routes.py index 3ca8c53d..a7d6be63 100644 --- a/ainode/auth/session_routes.py +++ b/ainode/auth/session_routes.py @@ -706,6 +706,11 @@ async def handle_sync_users(request: web.Request) -> web.Response: untouched, because a half-adopted list is a node whose logins quietly differ from the rest of the cluster. ``changed`` is false for a list this node already has, which is what makes a replication pass cheap to repeat. + + An EMPTY list is refused with a 409 and this node keeps its accounts. The + import replaces the whole list, so an empty push is a fleet-wide logout, and + the sender that makes one is a master whose own store is empty (promoted + before users.json reached it), or a release from before it refused to. """ refused = fleet_only_refusal(request) if refused is not None: @@ -715,6 +720,12 @@ async def handle_sync_users(request: web.Request) -> web.Response: body = await _body(request) if "users" not in body: return _error("Send a 'users' list.", "invalid_request", 400) + if isinstance(body.get("users"), list) and not body["users"]: + logger.warning("refused an empty account list from a peer; this node keeps " + "its %d account(s)", len(store.users)) + return _error("Refusing an empty account list: it would delete every account " + "on this node. The sender's account store is empty; copy " + "users.json to it first.", "empty_account_list", 409) try: changed = store.import_users(body.get("users")) except ValueError as exc: diff --git a/ainode/cli/doctor.py b/ainode/cli/doctor.py index e8e2e87f..9b2f18e7 100644 --- a/ainode/cli/doctor.py +++ b/ainode/cli/doctor.py @@ -19,6 +19,11 @@ :func:`udp_listeners`, :func:`http_json`, :func:`latest_image_tag`, :func:`probe_gpus`, :func:`running_in_container`, :func:`host_service_state`). Tests monkeypatch the seam, not the check. +* A node that loads no model under AINode (a routing-only master, a bare-metal + box whose GPU belongs to something else) passes. "No model loaded" is an INFO, + and the checks that only matter to an engine launch (docker, the engine + backend and its image) are INFO there too, with the reason on the line, because + nothing on that node will ever launch one until a model is loaded. * A check that needs something only the HOST can see says so as an INFO naming the host command, rather than a WARN. The documented deployment runs this inside the container, so a WARN about the absence of systemd in there was a @@ -574,13 +579,60 @@ def check_cluster_id(config: NodeConfig, peers_seen: int = 0) -> list[Check]: return [Check("config.cluster_id", OK, cluster_id, data=data)] +#: The launch state a node keeps under AINODE_HOME: the stacked instances replayed +#: at boot (``models/api_routes.py::_manifest_path``) and the distributed shape +#: (``engine/reconcile.py::RECORD_FILENAME``). +INSTANCE_MANIFEST_NAME = "instances.json" +DISTRIBUTED_RECORD_NAME = "distributed.json" + +#: Appended to an engine-only check on a node that loads no model. +NO_MODEL_NOTE = ("no model is loaded under AINode on this node, so this only " + "matters once one is") + + +def engine_expected(config: NodeConfig, home) -> bool: + """Will anything on this node launch an engine container? + + True when a primary is pinned (``config.model``), when stacked instances are + recorded to replay at boot, or when a distributed shape is recorded. False is + a node that serves nothing under AINode: a routing-only master, or a + bare-metal box whose GPU is used by something else. + """ + if (getattr(config, "model", "") or "").strip(): + return True + home = Path(home) + try: + manifest = json.loads((home / INSTANCE_MANIFEST_NAME).read_text()) + if isinstance(manifest, dict) and manifest.get("instances"): + return True + except (OSError, ValueError): + pass + return _exists(home / DISTRIBUTED_RECORD_NAME) + + +def as_no_model_info(checks: list[Check]) -> list[Check]: + """Engine-only checks on a node that loads no model: a finding becomes INFO. + + The fact is still printed, with the fix, so an operator about to load a model + knows what to sort out first; it just is not a reason for this node to fail. + """ + for check in checks: + if check.status in (WARN, FAIL): + check.data = dict(check.data, status_with_a_model=check.status, + engine_expected=False) + check.status = INFO + check.detail = f"{check.detail}; {NO_MODEL_NOTE}" + return checks + + def check_model(config: NodeConfig) -> list[Check]: model = (getattr(config, "model", "") or "").strip() models_dir = getattr(config, "models_dir", "") or "" data = {"model": model or None, "models_dir": models_dir} if not model: - return [Check("config.model", OK, - "no model pinned, so nothing loads at boot", data=data)] + return [Check("config.model", INFO, + "no model loaded under AINode: none is pinned, so nothing " + "loads at boot", data=data)] found = model_dir_on_disk(models_dir, model) if found is not None: return [Check("config.model", OK, f"{model} is on disk at {found}", @@ -928,7 +980,7 @@ def check_ports(config: NodeConfig, udp_bound: Optional[set[int]] = None) -> lis fix="launch the model from the dashboard, or watch ainode logs -f", data={"port": api_port, "listening": False, "model": model})) else: - checks.append(Check("port.engine", OK, + checks.append(Check("port.engine", INFO, f"{api_port} is free and no model is pinned, which is the " f"idle shape", data={"port": api_port, "listening": False})) @@ -1609,9 +1661,15 @@ def run_checks(home=None, config_path=None) -> list[Check]: docker_checks = check_docker(_engine_image_for(config)) docker_ok = bool(docker_checks[0].data.get("reachable")) image_present = docker_checks[0].data.get("image_present") - - checks += check_engine_backend(config, config_path, docker_ok, image_present, - in_container=inside) + backend_checks = check_engine_backend(config, config_path, docker_ok, image_present, + in_container=inside) + # Docker, the backend and its image are how an engine is LAUNCHED. A node + # that loads no model launches nothing, so they must not fail it. + if not engine_expected(config, home): + as_no_model_info(docker_checks) + as_no_model_info(backend_checks) + + checks += backend_checks checks += check_gpu_memory_utilization(config) checks += check_discovery_port(config) checks += check_model(config) diff --git a/ainode/models/api_routes.py b/ainode/models/api_routes.py index 2acb5e45..8c520c73 100644 --- a/ainode/models/api_routes.py +++ b/ainode/models/api_routes.py @@ -793,6 +793,31 @@ def _persist_primary_overrides(config, gmu, overrides) -> None: setattr(config, k, v) +def release_primary(app, config) -> None: + """The node's primary is gone: clear the slot and its launch parameters. + + The primary is the instance on the node's OWN api_port, the one config.json + boots. When it goes, nothing is promoted into its place. Promoting a survivor + was the old behaviour, and on Spark-4 it made Whisper, a stacked instance on + :8001, the node's ``config.model`` while config.json still carried the previous + primary's gpu_memory_utilization, image and vLLM flags: the node advertised + Whisper on :8000, where nothing served it, and the next boot would have + launched Whisper on :8000 with another model's engine parameters beside the + manifest's own copy on :8001. A survivor stays exactly what it was, a stacked + instance on its own port, replayed from the manifest. + + So ``app["engine"]`` and ``config.model`` go to None, every per-load override + goes back to its NodeConfig default (the same reset a bare load applies), and + a head flips back to solo. Saves the config; raises what ``save`` raises. + """ + app["engine"] = None + config.model = None + _persist_primary_overrides(config, None, None) + if getattr(config, "distributed_mode", "") == "head": + config.distributed_mode = "solo" + config.save() + + def _will_own_primary_port(manager, config, existing) -> bool: """Will this load end up on the node's OWN api_port, i.e. be the primary? @@ -2149,21 +2174,12 @@ async def handle_model_unload(request: web.Request) -> web.Response: except Exception as exc: errors.append(f"instance.stop(): {exc}") manager.remove(inst.record.instance_id) - # If the primary went away, repoint app["engine"]/config to a survivor - # so the status/proxy back-compat path doesn't dangle on a dead backend. + # If the primary went away, the node HAS no primary until a load takes + # the primary port again: see release_primary for why no survivor is + # promoted into its place. if request.app.get("engine") is inst.backend: - survivors = manager.instances() - if survivors: - keep = survivors[0] - request.app["engine"] = keep.backend - config.model = keep.record.model - else: - request.app["engine"] = None - config.model = None - if getattr(config, "distributed_mode", "") == "head": - config.distributed_mode = "solo" try: - config.save() + release_primary(request.app, config) except Exception as exc: errors.append(f"config.save: {exc}") # Persist the reduced set so a restart doesn't resurrect the unloaded one. diff --git a/tests/test_auth_replication.py b/tests/test_auth_replication.py index 4091f136..c70665a3 100644 --- a/tests/test_auth_replication.py +++ b/tests/test_auth_replication.py @@ -612,6 +612,78 @@ def test_a_cli_push_with_no_cluster_secret_says_so_and_sends_nothing(monkeypatch assert result["pushed"] == [] and "cluster_secret" in result["reason"] +# ============================================================================= +# 5b. An empty list never crosses the wire (the Atlas cutover guard) +# ============================================================================= +# +# import_users REPLACES the whole list, so a master promoted before its users.json +# was copied would log every dashboard user out of the fleet on its first push, +# and every worker would do the same to itself on its next pull. + +@pytest.mark.asyncio +async def test_a_master_with_an_empty_store_pushes_nothing_and_warns_once( + peer, session, caplog): + app = _master_app(session, peer.port, StubStore([])) + replicator = rep.AccountReplicator(app) + + with caplog.at_level(logging.WARNING, logger=rep.logger.name): + first = await replicator.tick() + second = await replicator.tick() + + assert peer.posts == [], "an empty list reached a worker" + assert peer.users == USERS + assert first["reason"] == rep.EMPTY_LIST_REASON + assert first["pushed"] == [] and second["pushed"] == [] + warnings = [r for r in caplog.records if r.levelno == logging.WARNING] + assert len(warnings) == 1, "the 60 second retry must not warn on every tick" + assert "no accounts" in warnings[0].getMessage() + assert "users.json" in warnings[0].getMessage() + + +@pytest.mark.asyncio +async def test_a_master_pushes_again_once_its_store_has_accounts(peer, session): + store = StubStore([]) + replicator = rep.AccountReplicator(_master_app(session, peer.port, store)) + await replicator.broadcast() + assert peer.posts == [] + + store.users = [dict(u) for u in USERS] + result = await replicator.broadcast() + + assert result["pushed"] == ["spark2"] + assert peer.posts[0]["body"] == {"users": USERS} + + +@pytest.mark.asyncio +async def test_a_worker_refuses_an_empty_list_from_its_master(peer, session, home, + caplog): + peer.users = [] + store = StubStore(USERS) + + with caplog.at_level(logging.WARNING, logger=rep.logger.name): + result = await rep.AccountReplicator( + _worker_app(session, peer.port, store)).pull() + + assert result["imported"] is False + assert result["reason"] == rep.EMPTY_LIST_REASON + assert store.imports == [], "the worker wiped its accounts to match an empty master" + assert store.users == USERS + assert rep.read_sync_state(home) == {}, "a refused pull is not a successful sync" + assert any("empty account list" in r.getMessage() for r in caplog.records) + + +def test_a_cli_push_of_an_empty_list_is_refused_and_sends_nothing(monkeypatch): + monkeypatch.setattr(rep, "http_json", + lambda *a, **kw: pytest.fail("an empty list went out")) + config = NodeConfig(node_id="m1", cluster_secret=SECRET, cluster_role="master", + peer_ips=["10.0.0.2"]) + + result = rep.replicate_from_cli(config, []) + + assert result["pushed"] == [] and result["failed"] == [] + assert rep.EMPTY_LIST_REASON in result["reason"] + + # ============================================================================= # 6. The stamp, and the store seam # ============================================================================= diff --git a/tests/test_doctor.py b/tests/test_doctor.py index 54ca9c80..513cd169 100644 --- a/tests/test_doctor.py +++ b/tests/test_doctor.py @@ -187,9 +187,11 @@ def test_an_empty_cluster_id_fails(): # ---------------------------------------------------------------------- model -def test_no_model_pinned_is_ok(): +def test_no_model_pinned_is_info(): + """A node that loads no model is a fact, not a finding (Atlas serves nothing).""" check = _one(doc.check_model(NodeConfig(model=None))) - assert check.status == OK + assert check.status == doc.INFO + assert "no model loaded" in check.detail assert "nothing loads at boot" in check.detail @@ -443,7 +445,7 @@ def test_the_expected_port_shape_on_an_idle_node(monkeypatch): monkeypatch.setattr(doc, "tcp_listening", lambda port, **kw: port == 3000) checks = doc.check_ports(NodeConfig(model=None, discovery_port=5679), udp_bound={5679}) - assert [c.status for c in checks] == [OK, OK, OK] + assert [c.status for c in checks] == [OK, doc.INFO, OK] assert "idle shape" in _by_name(checks, "port.engine").detail diff --git a/tests/test_doctor_no_model.py b/tests/test_doctor_no_model.py new file mode 100644 index 00000000..af5e50a6 --- /dev/null +++ b/tests/test_doctor_no_model.py @@ -0,0 +1,136 @@ +"""``ainode doctor`` passes on a node that loads no model under AINode. + +Atlas is the shape: a bare-metal source install with a GPU (an A40 that its own +llama.cpp services hold), running AINode as a routing-only master. Nothing on it +will launch an engine container, so neither docker nor the engine backend nor its +image is a reason for it to fail, and "no model loaded" is an INFO line. A node +that DOES load a model keeps every one of those findings exactly as before. +""" + +from __future__ import annotations + +import json + +import pytest + +from ainode.cli import doctor as doc + +A40 = {"index": 0, "name": "NVIDIA A40", "memory_total_mb": 46068, + "memory_free_mb": 3400, "unified_memory": False, "persistence_mode": False} + + +@pytest.fixture +def atlas(tmp_path, monkeypatch): + """Every seam faked as Atlas answers it: GPU present, no docker for this user, + no systemctl answer, dashboard up, engine port free, nobody else seen yet.""" + monkeypatch.setenv("AINODE_HOME", str(tmp_path)) + monkeypatch.setattr(doc, "running_in_container", lambda *a, **kw: False) + monkeypatch.setattr(doc, "run_command", lambda argv, timeout=10.0: (127, "absent")) + monkeypatch.setattr(doc, "disk_usage", lambda path: (1000, 900)) + monkeypatch.setattr(doc, "tcp_listening", lambda port, **kw: port == 3000) + monkeypatch.setattr(doc, "udp_listeners", lambda: {5679}) + monkeypatch.setattr(doc, "http_json", + lambda url, timeout=3.0, headers=None: {"nodes": []}) + monkeypatch.setattr(doc, "latest_image_tag", lambda: None) + monkeypatch.setattr(doc, "probe_gpus", lambda: [dict(A40)]) + (tmp_path / "models").mkdir() + + def write(**config): + base = {"node_id": "atlas", "node_name": "Atlas", "model": None, + "cluster_role": "master", "discovery_port": 5679, + "models_dir": str(tmp_path / "models")} + base.update(config) + (tmp_path / "config.json").write_text(json.dumps(base)) + return {c.name: c for c in doc.run_checks(tmp_path, tmp_path / "config.json")} + + return write, tmp_path + + +def test_a_node_that_loads_no_model_passes(atlas): + write, _ = atlas + checks = write() + + failed = [f"{c.name}: {c.detail}" for c in checks.values() if c.status == doc.FAIL] + assert failed == [] + assert doc.exit_code(list(checks.values())) == 0 + assert checks["gpu.devices"].status == doc.OK, "the GPU is still reported" + + +def test_no_model_loaded_is_info(atlas): + write, _ = atlas + checks = write() + + assert checks["config.model"].status == doc.INFO + assert "no model loaded" in checks["config.model"].detail + assert checks["port.engine"].status == doc.INFO + + +def test_the_engine_launch_checks_are_info_with_the_reason(atlas): + write, _ = atlas + checks = write() + + for name in ("docker.daemon", "config.engine_backend"): + check = checks[name] + assert check.status == doc.INFO, f"{name} is {check.status}" + assert doc.NO_MODEL_NOTE in check.detail + assert check.data["engine_expected"] is False + assert check.fix, "the fix is still printed for when a model is loaded" + # What it WOULD be on a node that loads one, kept for --json readers. + assert checks["docker.daemon"].data["status_with_a_model"] == doc.FAIL + + +def test_an_eugr_backend_with_no_vllm_does_not_fail_a_node_with_no_model(atlas, + monkeypatch): + write, _ = atlas + monkeypatch.setattr(doc.shutil, "which", lambda name: None) + checks = write(engine_backend="eugr") + + assert checks["config.engine_backend"].status == doc.INFO + assert doc.exit_code(list(checks.values())) == 0 + + +def test_a_node_that_pins_a_model_still_fails_without_docker(atlas): + """The regression guard: the relaxation is for nodes that serve nothing.""" + write, _ = atlas + checks = write(model="Qwen/Qwen3.8-27B") + + assert checks["docker.daemon"].status == doc.FAIL + assert doc.NO_MODEL_NOTE not in checks["docker.daemon"].detail + assert doc.exit_code(list(checks.values())) == 1 + + +def test_stacked_instances_on_disk_count_as_a_model(atlas): + write, home = atlas + (home / doc.INSTANCE_MANIFEST_NAME).write_text(json.dumps( + {"instances": [{"model": "openai/whisper-large-v3-turbo", + "gpu_memory_utilization": 0.2}]})) + checks = write() + + assert checks["docker.daemon"].status == doc.FAIL + + +def test_engine_expected_reads_the_pin_the_manifest_and_the_distributed_record(tmp_path): + from ainode.core.config import NodeConfig + + idle = NodeConfig(model=None) + assert doc.engine_expected(idle, tmp_path) is False + assert doc.engine_expected(NodeConfig(model="Qwen/Q"), tmp_path) is True + + (tmp_path / doc.INSTANCE_MANIFEST_NAME).write_text('{"instances": []}') + assert doc.engine_expected(idle, tmp_path) is False + (tmp_path / doc.INSTANCE_MANIFEST_NAME).write_text("not json") + assert doc.engine_expected(idle, tmp_path) is False + + (tmp_path / doc.DISTRIBUTED_RECORD_NAME).write_text("{}") + assert doc.engine_expected(idle, tmp_path) is True + + +def test_the_file_names_match_the_ones_the_node_writes(): + """One home per name: the doctor reads what the launch path writes.""" + from pathlib import Path + + import ainode.models.api_routes as mr + from ainode.engine import reconcile + + assert Path(mr._manifest_path()).name == doc.INSTANCE_MANIFEST_NAME + assert reconcile.RECORD_FILENAME == doc.DISTRIBUTED_RECORD_NAME diff --git a/tests/test_primary_release.py b/tests/test_primary_release.py new file mode 100644 index 00000000..224756ac --- /dev/null +++ b/tests/test_primary_release.py @@ -0,0 +1,150 @@ +"""Losing the primary never promotes a survivor (the Spark-4 Whisper leak). + +The primary is the instance on the node's OWN api_port, the one config.json boots. +Unloading it used to repoint ``app["engine"]`` and ``config.model`` at +``survivors[0]``. On Spark-4 that made Whisper, a stacked instance on :8001, the +node's model while config.json still held the old primary's image, vLLM flags and +memory share: the node advertised Whisper on :8000, where nothing served it, and +the next boot would have launched Whisper with another model's engine parameters. + +Pinned here, for both ways the primary can be taken down (``/api/models/unload`` +and the server view's eject): + +* the slot is EMPTY afterwards: ``app["engine"]`` and ``config.model`` are None; +* the old primary's per-load overrides are gone from the config, reset to the + NodeConfig defaults; +* the survivor is untouched: still running, still on its own port, still in the + manifest with its OWN parameters. +""" + +from __future__ import annotations + +import asyncio +import json + +import pytest + +import ainode.models.api_routes as mr +from ainode.core.config import NodeConfig +from ainode.discovery.cluster import ClusterState + +from tests.test_cluster_load import _patch_backend, _Req + +PRIMARY = "Qwen/Qwen3.8-27B" +WHISPER = "openai/whisper-large-v3-turbo" +PRIMARY_LOAD = { + "model": PRIMARY, + "gpu_memory_utilization": 0.55, + "max_model_len": 65536, + "engine_image": "vllm/vllm-openai:v0.27.1", + "extra_vllm_args": ["--reasoning-parser", "qwen3"], + "served_model_name": ["qwen"], +} + + +@pytest.fixture +def node(monkeypatch, tmp_path): + """A Spark-4 shaped node: a pinned primary on :8000, Whisper stacked on :8001.""" + import ainode.engine.reconcile as reconcile + + monkeypatch.setattr(mr, "_manifest_path", lambda: tmp_path / "instances.json") + monkeypatch.setattr(reconcile, "record_path", lambda: tmp_path / "distributed.json") + _patch_backend(monkeypatch) + cfg = NodeConfig(node_id="spark4", api_port=8000) + saves = [] + cfg.save = lambda: saves.append(cfg.model) + app = {"engine": None, "config": cfg, "cluster_state": ClusterState(), + "ray_autostart_state": None} + asyncio.run(mr.handle_model_load(_Req(app, dict(PRIMARY_LOAD)))) + asyncio.run(mr.handle_model_load(_Req(app, { + "model": WHISPER, "gpu_memory_utilization": 0.2}))) + manager = app["instances"] + assert manager.by_model(PRIMARY).record.api_port == 8000 + assert manager.by_model(WHISPER).record.api_port == 8001 + assert cfg.model == PRIMARY and cfg.engine_image == PRIMARY_LOAD["engine_image"] + return app, cfg, manager, saves, tmp_path + + +def _assert_slot_released(app, cfg, manager, tmp_path): + defaults = NodeConfig() + assert app["engine"] is None, "a survivor was promoted into the primary slot" + assert cfg.model is None, f"config.model leaked to {cfg.model!r}" + # The old primary's launch parameters do not outlive it. + assert cfg.engine_image == defaults.engine_image + assert cfg.extra_vllm_args == defaults.extra_vllm_args + assert cfg.max_model_len == defaults.max_model_len + assert cfg.served_model_name == defaults.served_model_name + assert cfg.gpu_memory_utilization == defaults.gpu_memory_utilization + + # Whisper stays what it was: stacked, on its own port, running. + whisper = manager.by_model(WHISPER) + assert whisper is not None and whisper.backend.stopped is False + assert whisper.record.api_port == 8001 + assert whisper.backend.config.gpu_memory_utilization == 0.2 + assert manager.by_model(PRIMARY) is None + + # The restart replays Whisper from the manifest with ITS parameters, and + # nothing boots on the primary port. + entries = json.loads((tmp_path / "instances.json").read_text())["instances"] + assert [e["model"] for e in entries] == [WHISPER] + assert entries[0]["gpu_memory_utilization"] == 0.2 + assert entries[0].get("engine_image") != PRIMARY_LOAD["engine_image"] + assert entries[0].get("extra_vllm_args") != PRIMARY_LOAD["extra_vllm_args"] + + +def test_unloading_the_primary_does_not_promote_whisper(node): + app, cfg, manager, saves, tmp_path = node + + resp = asyncio.run(mr.handle_model_unload(_Req(app, {"model": PRIMARY}))) + + body = json.loads(resp.body) + assert body["stopped"] is True and body["remaining"] == 1 + _assert_slot_released(app, cfg, manager, tmp_path) + assert saves and saves[-1] is None, "the cleared primary was not persisted" + + +def test_unloading_the_primary_by_port_does_not_promote_whisper(node): + app, cfg, manager, _saves, tmp_path = node + + asyncio.run(mr.handle_model_unload(_Req(app, {"api_port": 8000}))) + + _assert_slot_released(app, cfg, manager, tmp_path) + + +def test_ejecting_the_primary_does_not_promote_whisper(node): + from ainode.api.server_routes import handle_server_eject + + app, cfg, manager, _saves, tmp_path = node + + class _EjectReq: + def __init__(self, app, model_id): + self.app = app + self.match_info = {"model_id": model_id} + + resp = asyncio.run(handle_server_eject(_EjectReq(app, PRIMARY))) + + assert json.loads(resp.body)["ok"] is True + _assert_slot_released(app, cfg, manager, tmp_path) + + +def test_unloading_the_stacked_whisper_leaves_the_primary_alone(node): + app, cfg, manager, _saves, _tmp = node + primary_backend = manager.by_model(PRIMARY).backend + + asyncio.run(mr.handle_model_unload(_Req(app, {"model": WHISPER}))) + + assert app["engine"] is primary_backend + assert cfg.model == PRIMARY + assert cfg.engine_image == PRIMARY_LOAD["engine_image"] + assert cfg.extra_vllm_args == PRIMARY_LOAD["extra_vllm_args"] + + +def test_unloading_the_last_instance_still_empties_the_node(node): + app, cfg, manager, _saves, _tmp = node + + asyncio.run(mr.handle_model_unload(_Req(app, {"model": WHISPER}))) + asyncio.run(mr.handle_model_unload(_Req(app, {"model": PRIMARY}))) + + assert manager.is_empty() + assert app["engine"] is None and cfg.model is None + assert cfg.engine_image == "" diff --git a/tests/test_proxy_no_model.py b/tests/test_proxy_no_model.py new file mode 100644 index 00000000..609ede24 --- /dev/null +++ b/tests/test_proxy_no_model.py @@ -0,0 +1,114 @@ +"""A request with no ``model`` never defaults to a model this node does not serve. + +``proxy_to_vllm`` used to route a JSON body with no ``model`` to ``config.model``. +On a routing-only master (Atlas, which serves nothing) that name is the dataclass +default or whatever config.json was copied from, so the caller got a 404 or a 502 +about a model they never asked for. The fallback now needs a node that SERVES a +primary; without one the answer is a 400 naming the missing field, and nothing is +forwarded. A node that does serve one keeps today's behaviour exactly. +""" + +from __future__ import annotations + +import pytest + +from ainode.core.config import NodeConfig +from tests.test_routing_caps import _Collector, _Session, _cluster, _node, _proxy + +MODEL = "Qwen/Qwen3.8-27B" +CHAT = {"messages": [{"role": "user", "content": "hi"}]} + + +def _master(config_model=None, engine=None, peers=None): + """A node with the given config.model and primary engine, plus serving peers.""" + kwargs = {} if config_model is None else {"model": config_model} + config = NodeConfig(node_id="atlas", node_name="Atlas", api_port=8000, + web_port=3000, **kwargs) + nodes = [_node("atlas", model="")] + list(peers or []) + session = _Session() + return { + "config": config, + "engine": engine, + "cluster_state": _cluster(nodes), + "client_session": session, + "metrics_collector": _Collector(), + }, session + + +def _assert_missing_field(status, payload, session): + assert status == 400 + err = payload["error"] + assert err["code"] == "missing_model_field" + assert err["param"] == "model" + assert '"model"' in err["message"] + assert session.tried == [], "a request with no model was forwarded somewhere" + + +@pytest.mark.parametrize("config_model", [ + None, # the NodeConfig default, a name nothing here serves + "", # an explicitly cleared config + MODEL, # a config.json copied from a node that served it +]) +def test_a_node_serving_nothing_answers_400_naming_the_field(config_model): + peer = _node("spark1", model=MODEL, fabric="10.100.0.11") + app, session = _master(config_model=config_model, peers=[peer]) + + status, payload = _proxy(app, CHAT) + + _assert_missing_field(status, payload, session) + assert app["metrics_collector"].calls[-1][1] is True, "not counted as an error" + + +@pytest.mark.parametrize("body", [ + dict(CHAT, model=""), + dict(CHAT, model=" "), + dict(CHAT, model=None), +]) +def test_an_empty_model_is_as_missing_as_an_absent_one(body): + app, session = _master() + + status, payload = _proxy(app, body) + + _assert_missing_field(status, payload, session) + + +def test_a_body_that_is_not_a_json_object_is_refused_the_same_way(): + app, session = _master() + + status, payload = _proxy(app, ["not", "an", "object"]) + + _assert_missing_field(status, payload, session) + + +def test_a_named_model_still_routes_through_a_node_serving_nothing(): + peer = _node("spark1", model=MODEL, fabric="10.100.0.11") + app, session = _master(peers=[peer]) + + status, _ = _proxy(app, dict(CHAT, model=MODEL)) + + assert status == 200 + assert session.tried == ["http://10.100.0.11:8000/v1/chat/completions"] + + +def test_a_node_that_serves_a_primary_keeps_its_default(): + """Today's behaviour, unchanged: no model means this node's own model.""" + local = _node("atlas", model=MODEL) + app, session = _master(config_model=MODEL, engine=object()) + app["cluster_state"] = _cluster([local]) + + status, _ = _proxy(app, CHAT) + + assert status == 200 + assert session.tried == ["http://localhost:8000/v1/chat/completions"] + assert app["metrics_collector"].calls[-1] == (MODEL, False) + + +def test_an_empty_model_on_a_serving_node_also_takes_the_default(): + local = _node("atlas", model=MODEL) + app, session = _master(config_model=MODEL, engine=object()) + app["cluster_state"] = _cluster([local]) + + status, _ = _proxy(app, dict(CHAT, model="")) + + assert status == 200 + assert session.tried == ["http://localhost:8000/v1/chat/completions"] diff --git a/tests/test_session_routes.py b/tests/test_session_routes.py index d63b5b1d..664364d5 100644 --- a/tests/test_session_routes.py +++ b/tests/test_session_routes.py @@ -952,6 +952,29 @@ async def test_a_malformed_sync_leaves_the_node_alone(self, protected): assert resp.status == 400 assert client.app["users_store"].export_users() == before + @pytest.mark.asyncio + async def test_an_empty_sync_is_refused_and_the_node_keeps_its_accounts( + self, protected): + """import_users REPLACES the list, so an empty push is a fleet-wide logout. + + The sender is a master whose own store is empty (promoted before users.json + reached it), or an older release that still pushes one. + """ + client, _ = protected + headers = {"Authorization": f"Bearer {fleet_key(SECRET)}"} + before = client.app["users_store"].export_users() + + resp = await client.post("/api/auth/users/sync", json={"users": []}, + headers=headers) + + assert resp.status == 409 + body = await resp.json() + assert body["error"]["type"] == "empty_account_list" + assert "users.json" in body["error"]["message"] + assert client.app["users_store"].export_users() == before + assert client.app["users_store"].verify_password( + "jason", ADMIN_PASSWORD) is True + @pytest.mark.asyncio async def test_export_is_not_swallowed_by_the_name_route(self, protected): """``/users/export`` must not resolve as ``/users/{name}``."""