feat(mesh): cross-node routing + presence breaker (Phase 3)
- NodeConfig.address: routable host peers proxy to (wired registry + wire) - compute_routes emits 'remote' routes for consumed peer services (local requires ref satisfied by an online peer -> <address>:<port>) - _host_remote_block: fail-fast reverse_proxy (2s dial, passive health); presence expiry removes the route entirely = the primary circuit-breaker - mesh_gateway.py: API re-renders + reloads the Caddyfile on mesh change, iff content changed (no-op until a cross-node service is consumed) - tests: route emit / address fallback / breaker-absent / unconsumed Logic verified hermetically + against primer's real registry (castle-api -> primer:9020); live integration proven a no-op on civil.
This commit is contained in:
55
castle-api/src/castle_api/mesh_gateway.py
Normal file
55
castle-api/src/castle_api/mesh_gateway.py
Normal file
@@ -0,0 +1,55 @@
|
|||||||
|
"""Regenerate the gateway with cross-node (remote) routes from the live mesh.
|
||||||
|
|
||||||
|
Local routes come from `castle apply` (static). Remote routes are dynamic — they
|
||||||
|
appear/vanish as peers join/leave — so the API owns them: on a mesh change it
|
||||||
|
re-renders the Caddyfile (same generator `apply` uses, plus remote routes for
|
||||||
|
online peers) and reloads the gateway iff the content changed.
|
||||||
|
|
||||||
|
Safety: with no cross-node `requires`, the output equals the local-only Caddyfile,
|
||||||
|
so this is a verified no-op until this node actually consumes a peer service.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
import subprocess
|
||||||
|
|
||||||
|
from castle_core.config import SPECS_DIR
|
||||||
|
from castle_core.generators.caddyfile import generate_caddyfile_from_registry
|
||||||
|
|
||||||
|
from castle_api.config import get_registry
|
||||||
|
from castle_api.mesh import mesh_state
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
_GATEWAY_UNIT = "castle-castle-gateway.service"
|
||||||
|
|
||||||
|
|
||||||
|
def _regenerate(reload: bool) -> bool:
|
||||||
|
try:
|
||||||
|
reg = get_registry()
|
||||||
|
except Exception:
|
||||||
|
return False
|
||||||
|
remotes = {h: n.registry for h, n in mesh_state.all_nodes().items()}
|
||||||
|
try:
|
||||||
|
content = generate_caddyfile_from_registry(reg, remotes)
|
||||||
|
except Exception:
|
||||||
|
logger.exception("mesh gateway: Caddyfile generation failed")
|
||||||
|
return False
|
||||||
|
path = SPECS_DIR / "Caddyfile"
|
||||||
|
old = path.read_text() if path.exists() else ""
|
||||||
|
if content == old:
|
||||||
|
return False
|
||||||
|
path.write_text(content)
|
||||||
|
logger.info("mesh gateway: Caddyfile updated with cross-node routes")
|
||||||
|
if reload:
|
||||||
|
subprocess.run(
|
||||||
|
["systemctl", "--user", "reload", _GATEWAY_UNIT], check=False
|
||||||
|
)
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
|
async def refresh_remote_routes(reload: bool = True) -> bool:
|
||||||
|
"""Async wrapper — runs the blocking regen off the event loop."""
|
||||||
|
return await asyncio.to_thread(_regenerate, reload)
|
||||||
@@ -28,6 +28,8 @@ def registry_to_json(registry: NodeRegistry) -> str:
|
|||||||
"gateway_domain": registry.node.gateway_domain,
|
"gateway_domain": registry.node.gateway_domain,
|
||||||
# fleet role — so peers know which node is the config/secret authority.
|
# fleet role — so peers know which node is the config/secret authority.
|
||||||
"role": registry.node.role,
|
"role": registry.node.role,
|
||||||
|
# routable host peers proxy to for this node's services.
|
||||||
|
"address": registry.node.address,
|
||||||
},
|
},
|
||||||
"deployed": {},
|
"deployed": {},
|
||||||
}
|
}
|
||||||
@@ -76,6 +78,7 @@ def json_to_registry(payload: str) -> NodeRegistry:
|
|||||||
gateway_port=node_data.get("gateway_port", 9000),
|
gateway_port=node_data.get("gateway_port", 9000),
|
||||||
gateway_domain=node_data.get("gateway_domain"),
|
gateway_domain=node_data.get("gateway_domain"),
|
||||||
role=node_data.get("role", "follower"),
|
role=node_data.get("role", "follower"),
|
||||||
|
address=node_data.get("address"),
|
||||||
)
|
)
|
||||||
deployed: dict[str, Deployment] = {}
|
deployed: dict[str, Deployment] = {}
|
||||||
for key, comp_data in data.get("deployed", {}).items():
|
for key, comp_data in data.get("deployed", {}).items():
|
||||||
|
|||||||
@@ -28,6 +28,7 @@ from nats.js.api import KeyValueConfig
|
|||||||
from castle_core.registry import NodeRegistry
|
from castle_core.registry import NodeRegistry
|
||||||
|
|
||||||
from castle_api.mesh import mesh_state
|
from castle_api.mesh import mesh_state
|
||||||
|
from castle_api.mesh_gateway import refresh_remote_routes
|
||||||
from castle_api.mesh_wire import json_to_registry, registry_to_json
|
from castle_api.mesh_wire import json_to_registry, registry_to_json
|
||||||
from castle_api.stream import broadcast
|
from castle_api.stream import broadcast
|
||||||
|
|
||||||
@@ -93,6 +94,7 @@ class CastleNATSClient:
|
|||||||
await self.publish_registry(self._local_registry)
|
await self.publish_registry(self._local_registry)
|
||||||
await self._presence_kv.put(self._local_hostname, b"online")
|
await self._presence_kv.put(self._local_hostname, b"online")
|
||||||
await self._seed_existing()
|
await self._seed_existing()
|
||||||
|
await refresh_remote_routes() # establish any cross-node routes on startup
|
||||||
|
|
||||||
self._tasks = [
|
self._tasks = [
|
||||||
asyncio.create_task(self._watch_loop()),
|
asyncio.create_task(self._watch_loop()),
|
||||||
@@ -182,6 +184,7 @@ class CastleNATSClient:
|
|||||||
await broadcast(
|
await broadcast(
|
||||||
"mesh", {"event": "node_updated", "hostname": key}
|
"mesh", {"event": "node_updated", "hostname": key}
|
||||||
)
|
)
|
||||||
|
asyncio.create_task(refresh_remote_routes())
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("Error handling mesh entry for %s", key)
|
logger.exception("Error handling mesh entry for %s", key)
|
||||||
|
|
||||||
@@ -199,6 +202,7 @@ class CastleNATSClient:
|
|||||||
self._online.discard(hostname)
|
self._online.discard(hostname)
|
||||||
self._last_json.pop(hostname, None)
|
self._last_json.pop(hostname, None)
|
||||||
await broadcast("mesh", {"event": "node_offline", "hostname": hostname})
|
await broadcast("mesh", {"event": "node_offline", "hostname": hostname})
|
||||||
|
asyncio.create_task(refresh_remote_routes())
|
||||||
|
|
||||||
async def _heartbeat_loop(self) -> None:
|
async def _heartbeat_loop(self) -> None:
|
||||||
while True:
|
while True:
|
||||||
|
|||||||
@@ -121,8 +121,10 @@ def compute_routes(
|
|||||||
"""Build the ordered list of gateway routes. Every route is a host route whose
|
"""Build the ordered list of gateway routes. Every route is a host route whose
|
||||||
address is the service/frontend **name** (published at ``<name>.<domain>``);
|
address is the service/frontend **name** (published at ``<name>.<domain>``);
|
||||||
``proxy`` routes reverse-proxy a local port, ``static`` routes file-serve a
|
``proxy`` routes reverse-proxy a local port, ``static`` routes file-serve a
|
||||||
frontend's dist. Path routes no longer exist. ``remote_registries`` is accepted
|
frontend's dist. Path routes no longer exist. When ``remote_registries`` is
|
||||||
for signature compatibility but cross-node routing is out of scope here."""
|
given (online peers, keyed by hostname), ``remote`` routes are added for
|
||||||
|
services this node **consumes** (a local ``requires`` ref satisfied by a peer)
|
||||||
|
— so a consumed cross-node service is reachable at ``<ref>.<domain>``."""
|
||||||
if config is None:
|
if config is None:
|
||||||
try:
|
try:
|
||||||
from castle_core.config import load_config
|
from castle_core.config import load_config
|
||||||
@@ -141,9 +143,50 @@ def compute_routes(
|
|||||||
for name, kind, target in _local_routes(config, registry):
|
for name, kind, target in _local_routes(config, registry):
|
||||||
routes.append(GatewayRoute(name, kind, target, name, node))
|
routes.append(GatewayRoute(name, kind, target, name, node))
|
||||||
|
|
||||||
|
if remote_registries:
|
||||||
|
routes.extend(_remote_routes(config, registry, remote_registries))
|
||||||
|
|
||||||
return routes
|
return routes
|
||||||
|
|
||||||
|
|
||||||
|
def _remote_routes(
|
||||||
|
config: CastleConfig | None,
|
||||||
|
registry: NodeRegistry,
|
||||||
|
remote_registries: dict[str, NodeRegistry],
|
||||||
|
) -> list[GatewayRoute]:
|
||||||
|
"""Routes to services this node consumes from online peers.
|
||||||
|
|
||||||
|
A route is emitted for each local ``requires`` ref that (a) isn't satisfied
|
||||||
|
locally and (b) is provided by an exposed service on some peer. The route is
|
||||||
|
only present while the peer is (presence expiry removes the peer from
|
||||||
|
``remote_registries``, which *is* the circuit-breaker: gone → no route)."""
|
||||||
|
# Refs this node consumes.
|
||||||
|
consumed: set[str] = set()
|
||||||
|
local_names: set[str] = set()
|
||||||
|
if config is not None:
|
||||||
|
for _kind, name, dep in config.all_deployments():
|
||||||
|
local_names.add(name)
|
||||||
|
for req in getattr(dep, "requires", []) or []:
|
||||||
|
ref = getattr(req, "ref", None)
|
||||||
|
if ref and getattr(req, "kind", "deployment") == "deployment":
|
||||||
|
consumed.add(ref)
|
||||||
|
# Drop refs already satisfied locally.
|
||||||
|
consumed -= local_names
|
||||||
|
|
||||||
|
out: list[GatewayRoute] = []
|
||||||
|
for host, remote in sorted(remote_registries.items()):
|
||||||
|
addr = remote.node.address or host
|
||||||
|
for _kind, name, dep in remote.all():
|
||||||
|
if name not in consumed:
|
||||||
|
continue
|
||||||
|
if dep.subdomain and dep.port:
|
||||||
|
out.append(
|
||||||
|
GatewayRoute(name, "remote", f"{addr}:{dep.port}", name, host)
|
||||||
|
)
|
||||||
|
consumed.discard(name) # first online provider wins
|
||||||
|
return out
|
||||||
|
|
||||||
|
|
||||||
def _host_matcher_block(label: str, host: str, target: str) -> list[str]:
|
def _host_matcher_block(label: str, host: str, target: str) -> list[str]:
|
||||||
"""A `@host_X host <host> / handle @host_X { reverse_proxy <target> }` block.
|
"""A `@host_X host <host> / handle @host_X { reverse_proxy <target> }` block.
|
||||||
|
|
||||||
@@ -159,6 +202,27 @@ def _host_matcher_block(label: str, host: str, target: str) -> list[str]:
|
|||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def _host_remote_block(label: str, host: str, target: str) -> list[str]:
|
||||||
|
"""A remote (cross-node) host route with a fail-fast breaker: a short dial
|
||||||
|
timeout + passive health, so an unreachable peer 502s in ~2s instead of
|
||||||
|
hanging. (Presence removal drops the route entirely — this guards the
|
||||||
|
there-but-wedged case.)"""
|
||||||
|
matcher = f"@host_{label.replace('-', '_').replace('.', '_')}"
|
||||||
|
return [
|
||||||
|
f" {matcher} host {host}",
|
||||||
|
f" handle {matcher} {{",
|
||||||
|
f" reverse_proxy {target} {{",
|
||||||
|
" lb_try_duration 1s",
|
||||||
|
" fail_duration 30s",
|
||||||
|
" transport http {",
|
||||||
|
" dial_timeout 2s",
|
||||||
|
" }",
|
||||||
|
" }",
|
||||||
|
" }",
|
||||||
|
"",
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
def _host_static_block(label: str, host: str, serve_dir: str) -> list[str]:
|
def _host_static_block(label: str, host: str, serve_dir: str) -> list[str]:
|
||||||
"""A host matcher that file-serves a frontend's dist (with SPA fallback)."""
|
"""A host matcher that file-serves a frontend's dist (with SPA fallback)."""
|
||||||
matcher = f"@host_{label.replace('-', '_').replace('.', '_')}"
|
matcher = f"@host_{label.replace('-', '_').replace('.', '_')}"
|
||||||
@@ -237,6 +301,8 @@ def generate_caddyfile_from_registry(
|
|||||||
host = f"{r.address}.{domain}"
|
host = f"{r.address}.{domain}"
|
||||||
if r.kind == "static":
|
if r.kind == "static":
|
||||||
lines += _host_static_block(r.name or r.address, host, r.target)
|
lines += _host_static_block(r.name or r.address, host, r.target)
|
||||||
|
elif r.kind == "remote":
|
||||||
|
lines += _host_remote_block(r.name or r.address, host, r.target)
|
||||||
else:
|
else:
|
||||||
lines += _host_matcher_block(r.name or r.address, host, r.target)
|
lines += _host_matcher_block(r.name or r.address, host, r.target)
|
||||||
lines.append("}")
|
lines.append("}")
|
||||||
|
|||||||
@@ -32,6 +32,10 @@ class NodeConfig:
|
|||||||
tunnel_id: str | None = None
|
tunnel_id: str | None = None
|
||||||
# Emit the cert_obtained → `castle tls reconcile` hook (needs events-exec plugin).
|
# Emit the cert_obtained → `castle tls reconcile` hook (needs events-exec plugin).
|
||||||
cert_hook: bool = False
|
cert_hook: bool = False
|
||||||
|
# Routable host peers use to reach this node's services (LAN IP/hostname).
|
||||||
|
# Defaults to the hostname; set explicitly when the hostname isn't resolvable
|
||||||
|
# cross-node. Used to build `remote` gateway routes to this node.
|
||||||
|
address: str | None = None
|
||||||
# Fleet role: "authority" may write shared config/secrets to the mesh;
|
# Fleet role: "authority" may write shared config/secrets to the mesh;
|
||||||
# "follower" reconciles from it. Static (no election) — the authority is
|
# "follower" reconciles from it. Static (no election) — the authority is
|
||||||
# pinned in castle.yaml. When the authority is down, shared state is
|
# pinned in castle.yaml. When the authority is down, shared state is
|
||||||
@@ -158,6 +162,7 @@ def load_registry(path: Path | None = None) -> NodeRegistry:
|
|||||||
tunnel_id=node_data.get("tunnel_id"),
|
tunnel_id=node_data.get("tunnel_id"),
|
||||||
cert_hook=node_data.get("cert_hook", False),
|
cert_hook=node_data.get("cert_hook", False),
|
||||||
role=node_data.get("role", "follower"),
|
role=node_data.get("role", "follower"),
|
||||||
|
address=node_data.get("address"),
|
||||||
)
|
)
|
||||||
|
|
||||||
deployed: dict[str, Deployment] = {}
|
deployed: dict[str, Deployment] = {}
|
||||||
@@ -249,6 +254,8 @@ def save_registry(registry: NodeRegistry, path: Path | None = None) -> None:
|
|||||||
data["node"]["cert_hook"] = registry.node.cert_hook
|
data["node"]["cert_hook"] = registry.node.cert_hook
|
||||||
if registry.node.role and registry.node.role != "follower":
|
if registry.node.role and registry.node.role != "follower":
|
||||||
data["node"]["role"] = registry.node.role
|
data["node"]["role"] = registry.node.role
|
||||||
|
if registry.node.address:
|
||||||
|
data["node"]["address"] = registry.node.address
|
||||||
|
|
||||||
for key, comp in registry.deployed.items():
|
for key, comp in registry.deployed.items():
|
||||||
entry: dict = {
|
entry: dict = {
|
||||||
|
|||||||
87
core/tests/test_caddyfile_remote.py
Normal file
87
core/tests/test_caddyfile_remote.py
Normal file
@@ -0,0 +1,87 @@
|
|||||||
|
"""Cross-node (remote) gateway routes + presence breaker."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import yaml
|
||||||
|
from castle_core.config import load_config
|
||||||
|
from castle_core.generators.caddyfile import compute_routes
|
||||||
|
from castle_core.registry import Deployment, NodeConfig, NodeRegistry
|
||||||
|
|
||||||
|
|
||||||
|
def _config_requiring_widget(root: Path) -> None:
|
||||||
|
(root / "castle.yaml").write_text(yaml.safe_dump({"gateway": {"port": 18000}}))
|
||||||
|
svc_dir = root / "services"
|
||||||
|
svc_dir.mkdir()
|
||||||
|
(svc_dir / "consumer.yaml").write_text(
|
||||||
|
yaml.safe_dump(
|
||||||
|
{
|
||||||
|
"description": "consumes a remote widget",
|
||||||
|
"manager": "systemd",
|
||||||
|
"run": {"launcher": "command", "argv": ["consumer"]},
|
||||||
|
"requires": [{"kind": "deployment", "ref": "widget"}],
|
||||||
|
"manage": {"systemd": {}},
|
||||||
|
}
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _peer_with_widget(address: str | None) -> NodeRegistry:
|
||||||
|
widget = Deployment(
|
||||||
|
manager="systemd",
|
||||||
|
launcher="python",
|
||||||
|
run_cmd=[],
|
||||||
|
name="widget",
|
||||||
|
kind="service",
|
||||||
|
port=9099,
|
||||||
|
subdomain="widget",
|
||||||
|
managed=True,
|
||||||
|
)
|
||||||
|
return NodeRegistry(
|
||||||
|
node=NodeConfig(hostname="tower", address=address),
|
||||||
|
deployed={NodeRegistry.key("service", "widget"): widget},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _local() -> NodeRegistry:
|
||||||
|
return NodeRegistry(node=NodeConfig(hostname="civil"), deployed={})
|
||||||
|
|
||||||
|
|
||||||
|
def test_remote_route_emitted_for_consumed_peer_service(tmp_path: Path) -> None:
|
||||||
|
_config_requiring_widget(tmp_path)
|
||||||
|
config = load_config(tmp_path)
|
||||||
|
routes = compute_routes(
|
||||||
|
_local(), config, {"tower": _peer_with_widget("10.0.0.5")}
|
||||||
|
)
|
||||||
|
remote = [r for r in routes if r.kind == "remote"]
|
||||||
|
assert len(remote) == 1
|
||||||
|
assert remote[0].address == "widget"
|
||||||
|
assert remote[0].target == "10.0.0.5:9099"
|
||||||
|
assert remote[0].node == "tower"
|
||||||
|
|
||||||
|
|
||||||
|
def test_address_falls_back_to_hostname(tmp_path: Path) -> None:
|
||||||
|
_config_requiring_widget(tmp_path)
|
||||||
|
config = load_config(tmp_path)
|
||||||
|
routes = compute_routes(_local(), config, {"tower": _peer_with_widget(None)})
|
||||||
|
remote = [r for r in routes if r.kind == "remote"]
|
||||||
|
assert remote[0].target == "tower:9099"
|
||||||
|
|
||||||
|
|
||||||
|
def test_breaker_no_route_when_peer_absent(tmp_path: Path) -> None:
|
||||||
|
"""Presence expiry removes the peer from remote_registries -> no route."""
|
||||||
|
_config_requiring_widget(tmp_path)
|
||||||
|
config = load_config(tmp_path)
|
||||||
|
routes = compute_routes(_local(), config, {}) # no online peers
|
||||||
|
assert not [r for r in routes if r.kind == "remote"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_no_remote_route_for_unconsumed_service(tmp_path: Path) -> None:
|
||||||
|
"""A peer service nobody requires is not routed."""
|
||||||
|
(tmp_path / "castle.yaml").write_text(yaml.safe_dump({"gateway": {"port": 18000}}))
|
||||||
|
config = load_config(tmp_path) # no requires anywhere
|
||||||
|
routes = compute_routes(
|
||||||
|
_local(), config, {"tower": _peer_with_widget("10.0.0.5")}
|
||||||
|
)
|
||||||
|
assert not [r for r in routes if r.kind == "remote"]
|
||||||
@@ -247,6 +247,37 @@ uses the backend until `CASTLE_SECRET_BACKEND=openbao`):**
|
|||||||
3. **TLS hardening** — NATS + OpenBao are plaintext localhost. Required before
|
3. **TLS hardening** — NATS + OpenBao are plaintext localhost. Required before
|
||||||
cross-network or moving real secrets: NATS mTLS + auth, OpenBao TLS listener.
|
cross-network or moving real secrets: NATS mTLS + auth, OpenBao TLS listener.
|
||||||
|
|
||||||
|
### Phase 3 — cross-node routing + breaker: logic DONE + verified against a real peer (2026-07-07)
|
||||||
|
|
||||||
|
**Real second node:** `primer` (192.168.8.129) migrated onto the NATS mesh
|
||||||
|
(branch checked out, castle-api pointed at `nats://civil:4222` via a systemd
|
||||||
|
drop-in). civil ↔ primer mesh confirmed over the LAN — **Phase 1 two-node parity
|
||||||
|
done for real**, not simulated. (This also *restored* the civil↔primer mesh my
|
||||||
|
Phase 1 cutover had split — primer was still on MQTT→civil.)
|
||||||
|
|
||||||
|
- `NodeConfig.address` — routable host peers proxy to (wired through registry +
|
||||||
|
wire; falls back to hostname, which civil resolves for primer).
|
||||||
|
- `compute_routes` now emits **`remote`** routes for services this node
|
||||||
|
**consumes** (a local `requires` ref satisfied by an online peer), targeting
|
||||||
|
`<peer-address>:<port>`. `_host_remote_block` renders a **fail-fast breaker**
|
||||||
|
(2s dial timeout + passive health) for the there-but-wedged case; **presence
|
||||||
|
expiry removes the peer from the route set entirely** (gone → no route) — the
|
||||||
|
primary breaker.
|
||||||
|
- `castle_api/mesh_gateway.py` — on peer join/leave/change (+ startup) the API
|
||||||
|
re-renders the Caddyfile (same generator `apply` uses + remote routes for
|
||||||
|
online peers) and reloads the gateway **iff content changed**.
|
||||||
|
- **Verified:** 4 hermetic tests (route emitted / address fallback / breaker:
|
||||||
|
peer-absent → no route / unconsumed → no route); **resolution against primer's
|
||||||
|
real registry** pulled live from the mesh → `castle-api → primer:9020`; and the
|
||||||
|
live integration proven a **no-op** on civil (Caddyfile hash unchanged, gateway
|
||||||
|
healthy) since civil consumes nothing cross-node yet. Suites: **210 core + 93 api**.
|
||||||
|
|
||||||
|
**Remaining:** the full live curl+kill E2E (civil routing to a peer service, then
|
||||||
|
failing fast on kill) needs a peer-unique service civil consumes — every
|
||||||
|
underlying piece is verified, but the end-to-end demo needs that provisioning.
|
||||||
|
primer is now a permanent mesh member on the branch (revert: remove the drop-in +
|
||||||
|
`git checkout main` on primer).
|
||||||
|
|
||||||
## Decisions (resolved)
|
## Decisions (resolved)
|
||||||
|
|
||||||
1. **Discovery — both.** Keep mDNS for zero-config LAN peer discovery *and*
|
1. **Discovery — both.** Keep mDNS for zero-config LAN peer discovery *and*
|
||||||
|
|||||||
Reference in New Issue
Block a user