diff --git a/ainode/api/chat_routes.py b/ainode/api/chat_routes.py index 7a2062b7..e99634f9 100644 --- a/ainode/api/chat_routes.py +++ b/ainode/api/chat_routes.py @@ -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, diff --git a/ainode/api/decide.py b/ainode/api/decide.py index 7747b2f3..33013425 100644 --- a/ainode/api/decide.py +++ b/ainode/api/decide.py @@ -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 @@ -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 diff --git a/ainode/api/server.py b/ainode/api/server.py index e99678ab..532e1350 100644 --- a/ainode/api/server.py +++ b/ainode/api/server.py @@ -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) @@ -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"} @@ -1776,6 +1777,53 @@ 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). @@ -1783,34 +1831,38 @@ def _routing_candidates(cluster, model: str, local_node_id: str, local_port: int 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: @@ -1818,17 +1870,20 @@ def _routing_table(cluster, local_node_id: str, local_port: int) -> dict: 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)) diff --git a/ainode/bench/fleet.py b/ainode/bench/fleet.py index 2c38c0e9..dead26c9 100644 --- a/ainode/bench/fleet.py +++ b/ainode/bench/fleet.py @@ -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 []): diff --git a/ainode/models/api_routes.py b/ainode/models/api_routes.py index daac3a0c..e39c82e0 100644 --- a/ainode/models/api_routes.py +++ b/ainode/models/api_routes.py @@ -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" diff --git a/tests/test_federation.py b/tests/test_federation.py index bffc5604..730173f9 100644 --- a/tests/test_federation.py +++ b/tests/test_federation.py @@ -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)]