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
2 changes: 1 addition & 1 deletion AGENTS.md

Large diffs are not rendered by default.

15 changes: 9 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -1294,12 +1294,15 @@ response says what was applied, and `POST /v1/systemone` carries the same block

Send `"calibration": "raw"` to get the engine's own spread instead; the block
then reads `"applied": false` and still lists the temperatures you opted out of,
so you can refit against the raw numbers. The file lives on the node that serves
the model. A node that routes the request there without a copy of its own (the
master, usually) asks that node for its table over the fleet key
(`GET /api/decide/calibration?model=<id>`) and caches it for five minutes, so the
answer is tempered the same wherever the request enters. A peer on an older
release has no such route, and then the answer is raw with `"temperatures": null`.
so you can refit against the raw numbers. The node that serves the model is the
one that answers: a request that reaches a node without that model (the master,
usually) is handed whole to the serving node's AINode over the fleet key, so the
answer carries that node's temperatures and warm grammar wherever the request
entered. If the first node serving it is down the next replica answers; if no
node's AINode answers, the receiving node calls the engine itself and tempers
with the table it fetches from that node (`GET /api/decide/calibration?model=<id>`,
cached five minutes); only when that fails too is the answer raw with
`"temperatures": null`.

**A decision model warms up when it loads.** The first constrained request per
question shape makes the engine compile the answer grammar, 60 to 90 s on a
Expand Down
102 changes: 102 additions & 0 deletions ainode/api/decide.py
Original file line number Diff line number Diff line change
Expand Up @@ -682,6 +682,103 @@ def candidates_for(request: web.Request, model: str) -> list:
return candidates


# ------------------------------------------------------- the owner answers
#
# A decision model's temperatures.json and its warm state live on the node that
# serves it. When a request enters on a node that does not serve the model (the
# master, above all, which is the endpoint clients use), that node used to call
# the remote ENGINE itself, so the answer came back raw with
# ``calibration: {applied: false, temperatures: null}`` while the same request
# sent to the owner was tempered (2026-09-26, jebadiah-9b-v2 through Spark-1).
# Now the whole request goes to the owner's AINode, on the same route and under
# the fleet key, and the owner answers it exactly as if the caller had come
# straight to it. A request served locally takes the local path, unchanged.

#: How long one owner has to answer a forwarded decision: its own engine limit,
#: plus room for the hop.
FORWARD_TIMEOUT_S = CALL_TIMEOUT_S + 30.0


def owner_web_ports(app, candidates: list) -> list[tuple[str, int]]:
"""``(host, web_port)`` of each remote node in *candidates*, in routing order.

A candidate is an engine ``(fabric_ip, api_port)``, and a node with two
replicas stacked on it appears once. The web port is the one the node
announces, which is not always 3000 (Atlas serves on 3100).
"""
cluster = app.get("cluster_state")
by_host: dict = {}
for member in (cluster.members() if cluster is not None else []):
host = getattr(member, "fabric_ip", "") or ""
if host and host not in by_host:
by_host[host] = member
owners: list[tuple[str, int]] = []
for host, _port in candidates:
if host == "localhost": # _routing_candidates' name for this node
continue
member = by_host.get(host)
if member is None:
continue
owner = (host, int(getattr(member, "web_port", 3000) or 3000))
if owner not in owners:
owners.append(owner)
return owners


async def forward_to_owner(request: web.Request, body: dict,
candidates: list) -> Optional[web.Response]:
"""Hand the request to the node that serves the model, and return its answer.

None means "answer here": the model is served on this node (the local path
is unchanged), the request was itself forwarded (it is never forwarded
twice), no owner can be named, or no owner answered, in which case the old
path calls the engines directly and the answer is raw rather than missing.

Owners are tried in routing order. A transport failure, a timeout, a 5xx
(including an owner whose engine is down), or a 401/403/404 (a key it does
not share, or a release without the route) moves on to the next replica.
Any other status is the owner's own answer, a 4xx naming a field included,
and is returned as it came.
"""
from ainode.auth.fleet import FORWARDED_BY_HEADER, fleet_headers
if request.headers.get(FORWARDED_BY_HEADER):
return None
if any(host == "localhost" for host, _ in candidates):
return None
app = request.app
owners = owner_web_ports(app, candidates)
session = app.get("client_session")
if not owners or session is None:
return None
config = app.get("config")
headers = fleet_headers(app, {
"Content-Type": "application/json",
FORWARDED_BY_HEADER: str(getattr(config, "node_id", "") or "peer"),
})
timeout = aiohttp.ClientTimeout(total=FORWARD_TIMEOUT_S,
sock_connect=CONNECT_TIMEOUT_S)
payload = json.dumps(body).encode()
tried: list[str] = []
for host, web_port in owners:
url = f"http://{host}:{web_port}{request.path}"
try:
async with session.post(url, data=payload, headers=headers,
timeout=timeout) as resp:
answer = await resp.read()
if resp.status >= 500 or resp.status in (401, 403, 404):
tried.append(f"{host}:{web_port} {resp.status}")
continue
return web.Response(status=resp.status, body=answer,
content_type="application/json")
except asyncio.CancelledError:
raise
except Exception as exc:
tried.append(f"{host}:{web_port} {type(exc).__name__}")
logger.warning("no owner of %s answered %s (%s); calling the engines directly",
body.get("model"), request.path, "; ".join(tried))
return None


async def ask_one(session: aiohttp.ClientSession, candidates: list, body: dict,
timeout_s: float = CALL_TIMEOUT_S) -> tuple:
"""One question, with the proxy's failover. Returns (payload, cand, latency_ms).
Expand Down Expand Up @@ -1001,6 +1098,11 @@ async def handle_decide(request: web.Request) -> web.Response:
if not candidates:
return unavailable(f"no node is serving '{model}'")

# A model another node serves is answered by that node, calibration and all.
forwarded = await forward_to_owner(request, dict(body, model=model), candidates)
if forwarded is not None:
return forwarded

run = await run_questions(request, model, questions, state, instructions,
candidates, calibration)
collector = request.app.get("metrics_collector")
Expand Down
6 changes: 6 additions & 0 deletions ainode/api/systemone.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@
DecideError,
calibration_mode,
candidates_for,
forward_to_owner,
merge_usage,
normalize_questions,
resolve_model,
Expand Down Expand Up @@ -485,6 +486,11 @@ async def handle_systemone(request: web.Request) -> web.Response:
if not candidates:
return unavailable(f"no node is serving '{model}'")

# A model another node serves is answered by that node, calibration and all.
forwarded = await forward_to_owner(request, dict(body, model=model), candidates)
if forwarded is not None:
return forwarded

# No shared instructions block: in this format a question's own instructions
# are the whole prompt for it, and the questions of one ask still never see
# each other's answers.
Expand Down
6 changes: 6 additions & 0 deletions ainode/auth/fleet.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,12 @@
#: dashboard: revoking the fleet's access to a node means changing that node's
#: ``cluster_secret``, which is the same act as removing it from the cluster.
FLEET_KEY_ID = "fleet"
#: Set by a node that hands a decision request (/v1/decide, /v1/systemone) to the
#: node that owns the model. The owner answers it itself and never forwards it
#: again, and its rate limiter does not count it: the node the caller reached has
#: already counted the caller, and every forwarded request arrives under the one
#: fleet key, so counting them here would make the whole fleet one client.
FORWARDED_BY_HEADER = "X-AINode-Forwarded-By"


def fleet_key(secret: Optional[str]) -> str:
Expand Down
15 changes: 15 additions & 0 deletions ainode/ratelimit/middleware.py
Original file line number Diff line number Diff line change
Expand Up @@ -278,6 +278,19 @@ def client_key(request) -> str:
return f"ip:{remote or 'unknown'}"


def is_forwarded_by_fleet(request) -> bool:
"""A decision request another node forwarded here under the fleet key.

Only the fleet key stamps ``api_key_id == "fleet"``, so a client cannot claim
the exemption by sending the header. See ``auth/fleet.py::FORWARDED_BY_HEADER``.
"""
from ainode.auth.fleet import FLEET_KEY_ID, FORWARDED_BY_HEADER
getter = getattr(request, "get", None)
key_id = getter("api_key_id", "") if callable(getter) else ""
headers = getattr(request, "headers", None) or {}
return key_id == FLEET_KEY_ID and bool(headers.get(FORWARDED_BY_HEADER))


def too_many_requests(decision: Decision, config: RateLimitConfig) -> web.Response:
"""The 429: a Retry-After header and a body that names the limit."""
if decision.limit == "max_inflight":
Expand Down Expand Up @@ -319,6 +332,8 @@ async def rate_limit_middleware(request: web.Request, handler):
limiter: Optional[RateLimiter] = request.app.get("rate_limiter")
if limiter is None or not limiter.enabled or not is_limited_path(request.path):
return await handler(request)
if is_forwarded_by_fleet(request):
return await handler(request)
key = client_key(request)
decision = limiter.admit(key)
if not decision.allowed:
Expand Down
Loading
Loading