Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions AGENTS.md

Large diffs are not rendered by default.

49 changes: 46 additions & 3 deletions ainode/api/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand All @@ -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":
Expand All @@ -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
Expand Down
16 changes: 9 additions & 7 deletions ainode/api/server_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
46 changes: 46 additions & 0 deletions ainode/auth/replication.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 ------------------------------------------------------------
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand All @@ -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
Expand Down Expand Up @@ -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),
Expand Down
11 changes: 11 additions & 0 deletions ainode/auth/session_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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:
Expand Down
70 changes: 64 additions & 6 deletions ainode/cli/doctor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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}",
Expand Down Expand Up @@ -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}))
Expand Down Expand Up @@ -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)
Expand Down
42 changes: 29 additions & 13 deletions ainode/models/api_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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?

Expand Down Expand Up @@ -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.
Expand Down
Loading
Loading