Skip to content
Draft
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: 10 additions & 0 deletions runtime/glm53-spark-mtp3-mesh/MANAGED_MESH.md
Original file line number Diff line number Diff line change
Expand Up @@ -362,6 +362,7 @@ The bounded `--run-seconds` mode is for isolated diagnostics only.
|---|---|
| Child-process and peer checks | 1-second loop; peer HTTP timeout 2 seconds |
| Unavailable peer connection | 300-second grace after the first transport failure (covers a management-switch reboot); degraded health blocks model startup |
| Management-address probe | Startup: fail fast. Runtime: same grace bound as peer transport when the address matches the startup-validated identity; fabric checks continue and stay immediate; degraded health blocks model startup |
| Docker container status | One background query at a time, 3-second timeout; unknown status blocks model startup |
| MAC/IP, Ethernet MTU, sysfs GID/netdev, routes, qdiscs, TC state | 5-second periodic check |
| Full RDMA active-MTU probe | Startup and approximately every 60 seconds |
Expand All @@ -376,6 +377,15 @@ An authentication failure, explicit negative readiness, or changed process
generation does not receive transport-error grace: it triggers failure when
observed. Local marker exits also trigger failure without that grace.

A temporary loss of this rank's own management address shares the peer-transport
grace bound and does not itself stop serving. Fabric, marker, authentication,
generation and readiness checks continue during that interval and keep their
immediate-failure semantics; a management identity that differs from the
startup-validated one never receives grace. The model process survives the
outage, but API clients routed over the management network can still lose
connectivity for its duration, and a peer fault visible only through the
management path takes up to the grace interval to detect.

Docker status queries run outside the fabric-monitor loop. A slow or failed
query reports `docker_status_degraded: true`; it does not declare fabric
failure or interrupt existing serving. Marker, network, and authenticated
Expand Down
61 changes: 54 additions & 7 deletions runtime/glm53-spark-mtp3-mesh/managed_network.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,19 @@
fabric = profile.fabric


class ManagementAddressLoss(RuntimeError):
"""This rank's validated management address is absent from its netdev.

Raised only after every fabric check has passed, so callers can grant a
bounded grace interval without delaying any RoCE fault. Grace applies
only after startup validation: up() records the rank/address/netdev
tuple it verified, and check() raises this type only while the recorded
tuple matches the site's current assignment. Before the first successful
up(), an absent address raises ValueError instead and never receives
grace.
"""


def command(argv):
result = subprocess.run(argv, capture_output=True, text=True, timeout=15)
if result.returncode:
Expand Down Expand Up @@ -64,6 +77,7 @@ def __init__(self, site_path, rank, state_dir, *, runner=command,
for rule in self.plan.tc_rules:
if rule.intermediate_rank == rank:
self.objects["rule:" + rule.name] = ("rule", rule)
self.validated_management_identity = None

def _json(self, argv):
value = json.loads(self.runner(argv))
Expand Down Expand Up @@ -135,12 +149,31 @@ def _save(self):
finally:
if os.path.exists(temporary):
os.unlink(temporary)
def _management_identity(self):
"""The (address, netdev) tuple this site assigns to this rank."""
return (self.site["management_addresses"][self.rank], self.local.management_netdev)

def _links(self, *, verify_rdma_mtu=True):
def _management_address_present(self):
"""True when this rank's site-assigned management address is on its netdev."""
management = self._json(["ip", "-j", "-4", "addr", "show", "dev", self.local.management_netdev])
if self.site["management_addresses"][self.rank] not in {
x.get("local") for link in management for x in link.get("addr_info", [])}:
raise ValueError("Management address does not identify this rank")
return self.site["management_addresses"][self.rank] in {
x.get("local") for link in management for x in link.get("addr_info", [])}

def _links(self, *, verify_rdma_mtu=True, management_loss=None):
"""Verify the management address and the RoCE port fabric.

When management_loss is a list and this rank's startup-validated
management identity (the exact address/netdev tuple recorded by the
first successful up()) is absent from its netdev, the loss is
appended there and the fabric port checks still run. Any fabric
mismatch raises immediately: a RoCE fault must never wait behind a
management outage. Callers without a management_loss list keep
startup semantics and fail fast on the first absent address.
"""
if not self._management_address_present():
if management_loss is None or self.validated_management_identity != self._management_identity():
raise ValueError("Management address does not identify this rank")
management_loss.append(True)
for port in self.local.ports:
links = self._json(["ip", "-j", "addr", "show", "dev", port.netdev])
if (len(links) != 1 or links[0].get("mtu") != self.plan.expected_ethernet_mtu
Expand Down Expand Up @@ -209,23 +242,33 @@ def _argv(self, kind, obj, add):
return fabric.tc_rule_command(obj, add=add)
return ["tc", "qdisc", "add" if add else "del", "dev", obj, "clsact"]

def check(self, *, verify_rdma_mtu=True):
def check(self, *, verify_rdma_mtu=True, management_loss=None):
"""Verify local state; periodic callers may defer only the verbs MTU probe.

Startup and periodic full checks must retain the default. Disabling
the probe still verifies Ethernet MTU, GID/netdev identity and TC rules.

When management_loss is a list and this rank's startup-validated
management identity (the exact address/netdev tuple recorded by the
first successful up()) is absent from its netdev, every fabric check
still runs; if none raises, the typed ManagementAddressLoss is raised
for the caller's grace decision. An absent address without a matching
validated identity raises ValueError instead and never receives
grace. Fabric failures raise immediately regardless.
"""
self._links(verify_rdma_mtu=verify_rdma_mtu)
management_loss = management_loss if management_loss is not None else []
self._links(verify_rdma_mtu=verify_rdma_mtu, management_loss=management_loss)
missing = [key for key, (kind, obj) in self.objects.items() if not self._present(kind, obj)]
if missing:
raise ValueError(f"Missing mesh network objects: {missing}")
if management_loss:
raise ManagementAddressLoss("Management address does not identify this rank")
return {"ready": True, "rank": self.rank, "plan_sha256": self.plan.sha256,
"objects": len(self.objects)}

def up(self):
with self._lock():
self._links()
# Detect conflicting objects before making any network changes.
for kind, obj in self.objects.values():
self._present(kind, obj)
for key, (kind, obj) in self.objects.items():
Expand All @@ -241,6 +284,10 @@ def up(self):
self._save()
if not self._present(kind, obj):
raise ValueError(f"Created network object is absent: {key}")
# Record the validated identity only after every startup check
# above succeeded, so the grace bound covers exactly the
# configuration the supervisor verified at startup.
self.validated_management_identity = self._management_identity()
return {**self.check(), "ownership": dict(self.journal["objects"])}

def down(self):
Expand Down
33 changes: 29 additions & 4 deletions runtime/glm53-spark-mtp3-mesh/managed_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -127,16 +127,23 @@ def validate_group(rows):
view = digest(generations)
if any(row.get('phase') != 'armed' or row.get('view_digest') != view
or row.get('peer_health_degraded', False)
or row.get('management_degraded', False)
or row.get('docker_status_degraded', False) for row in rows):
raise RuntimeError('Mesh ranks have not armed the same process generation set')
return view


class PeerWatch:
"""Tolerate short connection loss, never authenticated negative readiness or a new generation."""
"""Tolerate short connection loss, never authenticated negative readiness or a new generation.

Peer transport and local management-address outages carry independent
latches: a successful peer poll clears only the peer latch, and a clean
fabric/management check clears only the management latch.
"""
def __init__(self):
self.generations = None
self.outage_started = None
self.mgmt_outage_started = None

def observe(self, rows):
observed = {str(row['rank']): row['generation'] for row in rows}
Expand All @@ -152,6 +159,16 @@ def transport_error(self, now):
if now - self.outage_started >= PEER_OUTAGE_GRACE:
raise RuntimeError('Authenticated peer transport remained unavailable beyond its grace interval')

def management_error(self, now):
"""Latch a local management-address outage; clears only on its own recovery."""
if self.mgmt_outage_started is None:
self.mgmt_outage_started = now
if now - self.mgmt_outage_started >= PEER_OUTAGE_GRACE:
raise RuntimeError('Local management address remained unavailable beyond its grace interval')

def management_recovered(self):
self.mgmt_outage_started = None


def notify(message):
address = os.environ.get('NOTIFY_SOCKET')
Expand Down Expand Up @@ -303,7 +320,8 @@ def __init__(self, config_path):
self.model = self.config['container_id']
self.state = {'protocol': PROTOCOL, 'rank': self.rank, 'epoch': self.config['epoch'],
'identity': self.identity, 'generation': self.generation,
'local_ready': False, 'phase': 'starting', 'view_digest': None}
'local_ready': False, 'phase': 'starting', 'view_digest': None,
'management_degraded': False}
self.lock = threading.Lock()
self.stop = threading.Event()
self.children = []
Expand Down Expand Up @@ -443,7 +461,7 @@ def run(self):
if docker_running(self.model):
raise RuntimeError('Stop the dependent model before starting mesh ownership')
self.owns_guard = True
from managed_network import NetworkManager
from managed_network import NetworkManager, ManagementAddressLoss
self.network = NetworkManager(Path(self.config['site_path']), self.rank, self.state_dir / 'network')
self.network.up()
self.start_markers()
Expand All @@ -460,7 +478,14 @@ def run(self):
raise RuntimeError('A managed source marker exited')
if time.monotonic() - last_network >= NETWORK_POLL_SECONDS:
full = time.monotonic() - last_full_network >= 60
self.network.check(verify_rdma_mtu=full)
try:
self.network.check(verify_rdma_mtu=full)
except ManagementAddressLoss:
peer_watch.management_error(time.monotonic())
self.publish(management_degraded=True)
else:
peer_watch.management_recovered()
self.publish(management_degraded=False)
if full:
last_full_network = time.monotonic()
last_network = time.monotonic()
Expand Down
97 changes: 97 additions & 0 deletions runtime/glm53-spark-mtp3-mesh/test_managed_network.py
Original file line number Diff line number Diff line change
Expand Up @@ -380,3 +380,100 @@ def test_root_required_by_default(rig, monkeypatch):
monkeypatch.setattr(network.os, "geteuid", lambda: 1000, raising=False)
with pytest.raises(PermissionError, match="root"):
manager.up()


def test_management_loss_raises_after_fabric_checks_pass(rig):
"""The real NetworkManager: an absent management address defers into the
typed loss only after every port, GID/netdev, object and (on request)
RDMA MTU check passed."""
manager, host = rig
manager.up() # records the validated identity
assert manager.validated_management_identity == manager._management_identity()
# Remove the address from the fake host's management netdev answer.
original_call = host.__call__
def without_address(argv):
if "addr" in argv and argv[-1] == manager.local.management_netdev:
return json.dumps([{"addr_info": []}])
return original_call(argv)
manager.runner = without_address
loss = []
outage_commands_before = len(host.commands)
with pytest.raises(network.ManagementAddressLoss, match="Management address"):
manager.check(management_loss=loss)
assert loss == [True]
# Every fabric check ran during THIS outage despite the loss: the
# per-port answers were read after startup activity was excluded.
outage_commands = host.commands[outage_commands_before:]
fabric_ports = [argv[-1] for argv in outage_commands
if "addr" in argv and argv[-1] != manager.local.management_netdev]
assert set(fabric_ports) == {p.netdev for p in manager.local.ports}
assert outage_commands # the object-presence queries also ran
def test_management_fail_fast_at_startup(rig):
"""Startup semantics: up() with the address already absent fails on the
first _links() call, before any validation is recorded."""
manager, host = rig
original_call = host.__call__
def without_address(argv):
if "addr" in argv and argv[-1] == manager.local.management_netdev:
return json.dumps([{"addr_info": []}])
return original_call(argv)
manager.runner = without_address
with pytest.raises(ValueError, match="Management address"):
manager.up()
assert manager.validated_management_identity is None


def test_management_loss_with_fabric_fault_raises_fabric_error(rig):
"""A fabric mismatch on the same round wins: no typed loss is raised."""
manager, host = rig
manager.up()
original_call = host.__call__
def without_address_and_wrong_mtu(argv):
if "addr" in argv and argv[-1] == manager.local.management_netdev:
return json.dumps([{"addr_info": []}])
if argv[0] == "ibv_devinfo":
return "\tactive_mtu: 1500 (1)\n"
return original_call(argv)
manager.runner = without_address_and_wrong_mtu
loss = []
with pytest.raises(ValueError, match="RoCE MTU differs"):
manager.check(management_loss=loss, verify_rdma_mtu=True)
# The loss was recorded during _links but the fabric fault took
# precedence: the raised error is the fabric one, not the typed loss.


def test_management_loss_with_missing_object_raises_object_error(rig):
"""A missing TC rule on the same round fails immediately, not as a loss."""
manager, host = rig
manager.up()
key = next(k for k in manager.objects if k.startswith("rule:"))
del host.inventory[key]
original_call = host.__call__
def without_address(argv):
if "addr" in argv and argv[-1] == manager.local.management_netdev:
return json.dumps([{"addr_info": []}])
return original_call(argv)
manager.runner = without_address
loss = []
with pytest.raises(ValueError, match="Missing mesh network objects"):
manager.check(management_loss=loss)
assert loss == [True] # recorded during _links, but objects error wins


def test_unvalidated_management_identity_never_receives_grace(rig):
"""Without a successful up(), an absent address raises ValueError even
through the loss channel; a site change after validation also fails."""
manager, host = rig
original_call = host.__call__
def without_address(argv):
if "addr" in argv and argv[-1] == manager.local.management_netdev:
return json.dumps([{"addr_info": []}])
return original_call(argv)
manager.runner = without_address
assert manager.validated_management_identity is None
loss = []
with pytest.raises(ValueError, match="Management address"):
manager.check(management_loss=loss)
assert loss == []


Loading
Loading