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
10 changes: 7 additions & 3 deletions ainode/api/chat_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -116,13 +116,17 @@ def add(entry: dict) -> None:
out[i] = entry
return

from ainode.api.server import _local_fabric_ip, member_hosts

local_fabric = _local_fabric_ip(cluster, local_id)
for node in (cluster.members() if cluster is not None else []):
if _status_of(node) not in _SERVING_STATES:
continue
is_local = node.node_id == local_id
host = "localhost" if is_local else (getattr(node, "fabric_ip", "") or "")
if not host:
continue # remote node with no fabric IP is unroutable, so unreportable
hosts = ["localhost"] if is_local else member_hosts(node, local_fabric)
if not hosts:
continue # remote node with no address is unroutable, so unreportable
host = hosts[0]
node_port = local_port if is_local else node.api_port
base = {
"node_id": node.node_id,
Expand Down
14 changes: 8 additions & 6 deletions ainode/api/decide.py
Original file line number Diff line number Diff line change
Expand Up @@ -566,10 +566,11 @@ def _peer_web_ports(app, candidates: list) -> list[tuple[str, int]]:
"""
cluster = app.get("cluster_state")
by_host: dict = {}
from ainode.api.server import member_addresses

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
for host in member_addresses(member):
by_host.setdefault(host, member)
peers: list[tuple[str, int]] = []
for host, _port in candidates:
if host == "localhost": # _routing_candidates names this node so
Expand Down Expand Up @@ -708,10 +709,11 @@ def owner_web_ports(app, candidates: list) -> list[tuple[str, int]]:
"""
cluster = app.get("cluster_state")
by_host: dict = {}
from ainode.api.server import member_addresses

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
for host in member_addresses(member):
by_host.setdefault(host, member)
owners: list[tuple[str, int]] = []
for host, _port in candidates:
if host == "localhost": # _routing_candidates' name for this node
Expand Down
93 changes: 74 additions & 19 deletions ainode/api/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -1001,7 +1001,7 @@ def pinned_candidate(cluster, config, node_id: Optional[str], port: Optional[int
"""
if node_id and node_id != config.node_id:
node = cluster.get_node(node_id) if cluster is not None else None
host = (getattr(node, "fabric_ip", "") or "") if node is not None else ""
host = best_host(cluster, config.node_id, node) if node is not None else ""
if not host:
return None
return (host, port or getattr(node, "api_port", None) or config.api_port)
Expand Down Expand Up @@ -1742,10 +1742,11 @@ async def json(self):
return await handler(_Shim(request, body))

node = cluster.get_node(node_id) if cluster is not None else None
host = (getattr(node, "fabric_ip", "") or "") if node else ""
local_id = getattr(request.app.get("config"), "node_id", None)
host = best_host(cluster, local_id, node) if node else ""
if not host:
return web.json_response(
{"error": f"node '{node_id}' not found or has no fabric IP"}, status=404)
{"error": f"node '{node_id}' not found or has no address"}, status=404)
url = f"http://{host}:{node.web_port}{path}"
session: aiohttp.ClientSession = request.app["client_session"]
fwd = {k: v for k, v in body.items() if k != "node_id"}
Expand Down Expand Up @@ -1776,59 +1777,113 @@ async def handle_cluster_unload(request: web.Request) -> web.Response:
return await _cluster_dispatch(request, "/api/models/unload")


def _same_subnet(a: str, b: str) -> bool:
"""True when two IPv4 addresses share a /24. Anything else answers False."""
if not a or not b or a.count(".") != 3 or b.count(".") != 3:
return False
return a.rsplit(".", 1)[0] == b.rsplit(".", 1)[0]


def _local_fabric_ip(cluster, local_node_id: str) -> str:
"""This node's own announced fabric IP, from cluster state, or ""."""
for n in (cluster.members() if cluster is not None else []):
if n.node_id == local_node_id:
return getattr(n, "fabric_ip", "") or ""
return ""


def member_hosts(node, local_fabric: str) -> list:
"""Addresses to reach a REMOTE member on, best first, or [] (#295).

A member announces its fabric IP, and the UDP source of that announcement
is its LAN address (``peer_ip``). The fabric is the right path between nodes
that share it. A router that is not on it (a master on the R750, whose own
fabric IP is its LAN address) has no route there, so every request would
time out. When this node's own fabric IP and the member's are on different
/24s, the LAN address goes first. The other one stays in the list as a
failover target, so a wrong guess costs a retry rather than the request.
With no fabric IP of our own known, the fabric stays first, as before.
"""
fabric = getattr(node, "fabric_ip", "") or ""
peer = getattr(node, "peer_ip", "") or ""
hosts = [h for h in (fabric, peer) if h]
if fabric and peer and local_fabric and not _same_subnet(local_fabric, fabric):
hosts = [peer, fabric]
return list(dict.fromkeys(hosts))


def best_host(cluster, local_node_id: str, node) -> str:
"""The one address to reach remote *node* on (``member_hosts``), or ""."""
hosts = member_hosts(node, _local_fabric_ip(cluster, local_node_id)) if node else []
return hosts[0] if hosts else ""


def member_addresses(node) -> tuple:
"""Every address *node* may appear under in a routing candidate."""
return tuple(h for h in (getattr(node, "fabric_ip", "") or "",
getattr(node, "peer_ip", "") or "") if h)


def _routing_candidates(cluster, model: str, local_node_id: str, local_port: int) -> list:
"""All (host, port) currently serving `model` (routing-truth).

Returns a LIST so proxy_to_vllm can fail over when the first target is a
stale/ghost claim (a node that crashed but still advertises the model). A
crashed node is indistinguishable from a live one in cluster state, so
failover — not ordering — is what makes routing robust. Local node first
(cheapest hop), then remote peers.
(cheapest hop), then remote peers on their best address (``member_hosts``),
then the same peers on their other address, so every node's best guess is
tried before any second guess.
"""
local, remote = [], []
local, remote, fallback = [], [], []
local_fabric = _local_fabric_ip(cluster, local_node_id)
for n in (cluster.members() if cluster is not None else []):
status = n.status.value if hasattr(n.status, "value") else str(n.status)
if status not in ("online", "serving", "member-ready"):
continue
is_local = n.node_id == local_node_id
host = "localhost" if is_local else (getattr(n, "fabric_ip", "") or "")
if not host:
hosts = ["localhost"] if is_local else member_hosts(n, local_fabric)
if not hosts:
continue
node_port = local_port if is_local else n.api_port
bucket = local if is_local else remote
seen = set()
ports = []
# The node's primary/solo model is served on its main api_port.
if getattr(n, "model", "") == model and node_port not in seen:
bucket.append((host, node_port))
seen.add(node_port)
if getattr(n, "model", "") == model:
ports.append(node_port)
# Each stacked instance is served on its OWN port — a co-resident 2nd
# model on this node lives at :8001, not the node's main :8000.
for inst in (getattr(n, "instances", []) or []):
if inst.get("model") != model:
continue
iport = inst.get("api_port") or node_port
if iport not in seen:
bucket.append((host, iport))
seen.add(iport)
return local + remote
if iport not in ports:
ports.append(iport)
primary = local if is_local else remote
for port in ports:
primary.append((hosts[0], port))
fallback.extend((host, port) for host in hosts[1:])
return local + remote + fallback


def _routing_table(cluster, local_node_id: str, local_port: int) -> dict:
"""model name → (host, port) for every model served across the fleet (F1).

Built from cluster broadcast state: each node advertises its solo model and
any instances it heads. The local node routes to localhost; remote nodes to
their fabric IP (reachable from the master over the cluster fabric).
their best address (``member_hosts``): the fabric IP when this node shares
the fabric, else the LAN address the announcement came from.
"""
table: dict = {}
local_fabric = _local_fabric_ip(cluster, local_node_id)
for n in cluster.members():
status = n.status.value if hasattr(n.status, "value") else str(n.status)
if status not in ("online", "serving", "member-ready"):
continue
is_local = n.node_id == local_node_id
host = "localhost" if is_local else (getattr(n, "fabric_ip", "") or "")
if not host:
hosts = ["localhost"] if is_local else member_hosts(n, local_fabric)
if not hosts:
continue
host = hosts[0]
port = local_port if is_local else n.api_port
if getattr(n, "model", ""):
table.setdefault(n.model, (host, port))
Expand Down
6 changes: 4 additions & 2 deletions ainode/bench/fleet.py
Original file line number Diff line number Diff line change
Expand Up @@ -340,10 +340,12 @@ def _owning_node(cluster, local_node_id, host, port):
the node. Matched the same way the candidate was built: the local node is
addressed as "localhost", peers by fabric IP.
"""
from ainode.api.server import member_addresses

for n in (cluster.members() if cluster is not None else []):
is_local = n.node_id == local_node_id
nhost = "localhost" if is_local else (getattr(n, "fabric_ip", "") or "")
if nhost != host:
nhosts = ("localhost",) if is_local else member_addresses(n)
if host not in nhosts:
continue
ports = {getattr(n, "api_port", 0)}
for inst in (getattr(n, "instances", []) or []):
Expand Down
4 changes: 3 additions & 1 deletion ainode/models/api_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -2275,10 +2275,12 @@ async def handle_model_unload(request: web.Request) -> web.Response:
remote_stopped = False
peers_reached = 0
if cluster is not None and session is not None and model:
from ainode.api.server import best_host

for node in cluster.members():
if node.node_id == config.node_id:
continue
host = node.fabric_ip or node.node_name
host = best_host(cluster, config.node_id, node) or node.node_name
if not host:
continue
url = f"http://{host}:{node.web_port}/api/models/unload?fanout=0"
Expand Down
64 changes: 64 additions & 0 deletions tests/test_federation.py
Original file line number Diff line number Diff line change
Expand Up @@ -187,3 +187,67 @@ def test_the_master_routes_to_an_announced_head_and_its_stack():
"spark1", 3000) == [("10.100.0.12", 8000)]
assert server._routing_candidates(cluster, "Qwen/Qwen3-Embedding-0.6B",
"spark1", 3000) == [("10.100.0.12", 8001)]


# #295: a router off the fabric (a master on the R750) reaches members on the
# LAN address their announcement came from, not their fabric IP.

def _peer_node(nid, model="", fabric="", peer=None, instances=None):
n = _node(nid, model=model, fabric=fabric, instances=instances)
n.peer_ip = peer
return n


def test_off_fabric_router_prefers_the_lan_address():
from ainode.api.server import _routing_candidates
c = _cluster([
_peer_node("atlas", fabric="192.168.0.98"),
_peer_node("spark1", model="Q", fabric="10.100.0.11", peer="192.168.0.10"),
])
assert _routing_table(c, "atlas", 8100)["Q"] == ("192.168.0.10", 8000)
# LAN first, the fabric kept as a last-resort failover target.
assert _routing_candidates(c, "Q", "atlas", 8100) == [
("192.168.0.10", 8000), ("10.100.0.11", 8000)]


def test_on_fabric_router_keeps_the_fabric_first():
from ainode.api.server import _routing_candidates
c = _cluster([
_peer_node("spark1", fabric="10.100.0.11", peer="192.168.0.10"),
_peer_node("spark2", model="Q", fabric="10.100.0.13", peer="192.168.0.11"),
])
assert _routing_table(c, "spark1", 8000)["Q"] == ("10.100.0.13", 8000)
assert _routing_candidates(c, "Q", "spark1", 8000) == [
("10.100.0.13", 8000), ("192.168.0.11", 8000)]


def test_every_best_guess_is_tried_before_any_fallback():
from ainode.api.server import _routing_candidates
inst = {"model": "Q", "api_port": 8001}
c = _cluster([
_peer_node("atlas", fabric="192.168.0.98"),
_peer_node("spark1", model="Q", fabric="10.100.0.11", peer="192.168.0.10"),
_peer_node("spark4", fabric="10.100.0.17", peer="192.168.0.199", instances=[inst]),
])
got = _routing_candidates(c, "Q", "atlas", 8100)
assert got[:2] == [("192.168.0.10", 8000), ("192.168.0.199", 8001)]
assert set(got[2:]) == {("10.100.0.11", 8000), ("10.100.0.17", 8001)}


def test_unknown_own_fabric_keeps_the_old_order():
from ainode.api.server import member_hosts
n = _peer_node("spark2", fabric="10.100.0.13", peer="192.168.0.11")
assert member_hosts(n, "") == ["10.100.0.13", "192.168.0.11"]


def test_lan_only_member_is_now_routable():
c = _cluster([_peer_node("box", model="Q", fabric="", peer="192.168.0.50")])
assert _routing_table(c, "spark1", 8000)["Q"] == ("192.168.0.50", 8000)


def test_decide_owner_lookup_matches_the_lan_address():
from ainode.api.decide import owner_web_ports
n = _peer_node("spark1", model="Q", fabric="10.100.0.11", peer="192.168.0.10")
n.web_port = 3000
app = {"cluster_state": _cluster([n])}
assert owner_web_ports(app, [("192.168.0.10", 8000)]) == [("192.168.0.10", 3000)]
Loading