mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 19:18:25 +00:00
Assert datagram delivery over Ethernet in the ethernet-mesh chaos scenario
No chaos scenario sent a datagram across an Ethernet link: both Ethernet scenarios asserted only that the tree formed, so framing, AEAD over AF_PACKET and delivery through the Ethernet transport were untested. A new delivery assertion pings node pairs at each configured payload size at teardown, after flapped links and stopped nodes are restored, one packet at a time until min_replies replies arrive or deadline_secs passes. The verdict is whether that happened, not a loss ratio, so a busy host slows the probe without failing it. With require_transport and transport_node it also reads that node's peers before and after the probe and fails if the node has no peers or any peer on another transport, so traversal is checked rather than assumed. A probe that never ran fails as a harness failure. ethernet-mesh asserts delivery to and from n04, whose only edges are Ethernet, at payloads 0 and 1200: n04 to n06 and back cross Ethernet then UDP, and n04 to n05 is one Ethernet hop. Break-checks: dropping echo requests on n06's TUN fails n04->n06 with zero replies while baseline passes; adding a UDP edge to n04 fails the traversal check. The payload sizes cannot exercise the receiver's padding trim on veth: an empty echo is a 226-byte frame and a 1200-byte echo 1342 bytes, and the veth does not pad shorter frames (54-byte control frames and 48-byte beacons arrived unpadded; no received data frame carried bytes past its length field). The trim is covered by the unit tests on the receive loop's parse.
This commit is contained in:
+21
-1
@@ -105,7 +105,10 @@ Explicit topologies exercising non-UDP transports.
|
||||
- **ethernet-only**: 4-node ring on raw Ethernet (AF_PACKET). Peers discovered
|
||||
via beacons, not static config. Minimal netem (1-5ms delay).
|
||||
- **ethernet-mesh**: Mirrors `tcp-mesh` topology but with Ethernet instead of
|
||||
TCP. UDP edges use static config; Ethernet edges use beacon discovery.
|
||||
TCP. UDP edges use static config; Ethernet edges use beacon discovery. The
|
||||
only scenario that asserts datagram delivery over Ethernet: its `delivery`
|
||||
assertion pings to and from n04, whose only edges are Ethernet, at two
|
||||
payload sizes, and checks that n04 has no peer on another transport.
|
||||
- **tcp-mesh**: 6-node mesh with 4 UDP and 3 TCP edges. Both transports use
|
||||
static peer config. Netem mutation (30% fraction, every 20-40s) and link
|
||||
flaps (1 link max, 10-20s down).
|
||||
@@ -258,6 +261,23 @@ The assertion thresholds in the shipped file are calibrated against
|
||||
recorded runs at the invocation CI uses, and the file's own comments say
|
||||
what they were derived from. Read those before retuning them.
|
||||
|
||||
A `delivery` assertion checks the data plane rather than the tree. It runs
|
||||
at teardown, after flapped links and stopped nodes are restored, and pings
|
||||
each pair at each payload size until `min_replies` replies arrive or
|
||||
`deadline_secs` passes:
|
||||
|
||||
```yaml
|
||||
assertions:
|
||||
delivery:
|
||||
pairs: [[n04, n06], [n06, n04]] # [src, dst] node ids
|
||||
payload_bytes: [0, 1200] # ICMPv6 echo payload sizes, 0-1400
|
||||
min_replies: 3 # default 3
|
||||
deadline_secs: 60 # default 60, per pair and size
|
||||
require_transport: ethernet # optional, with transport_node:
|
||||
transport_node: n04 # fail if n04 has no peers, or any
|
||||
# peer on another transport
|
||||
```
|
||||
|
||||
## Topology Algorithms
|
||||
|
||||
| Algorithm | Parameters | Description |
|
||||
|
||||
@@ -3,7 +3,8 @@
|
||||
# Exercises both transports in a single mesh. UDP edges use static
|
||||
# peer config; Ethernet edges use beacon discovery. Tests that the
|
||||
# spanning tree converges across heterogeneous transports with netem
|
||||
# and link flaps active.
|
||||
# and link flaps active, and that datagrams are delivered over the
|
||||
# Ethernet links (the delivery assertion below).
|
||||
#
|
||||
# Topology:
|
||||
#
|
||||
@@ -62,20 +63,46 @@ link_flaps:
|
||||
traffic:
|
||||
enabled: false
|
||||
|
||||
# Baseline: the mesh came up, agreed on a root, and took parents. This
|
||||
# asserts nothing about Ethernet link behaviour under flaps; it exists
|
||||
# so that a run in which the mesh never formed cannot report success,
|
||||
# which until now it could, because this scenario carried no assertions
|
||||
# at all.
|
||||
# Baseline: the mesh came up, agreed on a root, and took parents. It
|
||||
# says nothing about the data plane; it exists so that a run in which
|
||||
# the mesh never formed cannot report success. Six nodes, one root, five
|
||||
# parented in all six provably-completed archived runs.
|
||||
#
|
||||
# Six nodes, one root, five parented in all six provably-completed
|
||||
# archived runs. The other assertion covering Ethernet transport in
|
||||
# CI; see ethernet-only for why that matters.
|
||||
# Delivery: the one assertion anywhere in CI that a datagram crossed an
|
||||
# Ethernet link. n04's only edges are Ethernet (n01-n04, n04-n05), so a
|
||||
# probe to or from n04 must cross one, provided n04 has no peer on
|
||||
# another transport; its container also sits on the docker network, so
|
||||
# that is checked, not assumed: n04's peers are read before and after
|
||||
# the probe, and the assertion fails if it has none or any is not
|
||||
# Ethernet. n04 -> n06 and n06 -> n04 cross Ethernet and then UDP;
|
||||
# n04 -> n05 is a direct Ethernet hop.
|
||||
#
|
||||
# Load-robust by shape: each pair and size is pinged one packet at a
|
||||
# time until 3 replies or 60 s, and the verdict is whether that
|
||||
# happened, not a loss ratio, so a busy host slows it without failing
|
||||
# it. It runs at teardown, after flapped links are restored.
|
||||
#
|
||||
# Sizes: payload 0 is the smallest ICMPv6 echo and 1200 is near the
|
||||
# 1280-byte TUN MTU. Neither can reach the receiver's padding trim on a
|
||||
# veth link. Measured on a live run of this scenario (2026-09-19): an
|
||||
# empty echo is a 226-byte Ethernet frame and a 1200-byte one 1342
|
||||
# bytes, far above the 60-byte minimum; and the veth does not pad the
|
||||
# frames that are shorter (54-byte control frames and 48-byte beacons
|
||||
# arrived as sent; none of 250 received data frames carried bytes past
|
||||
# its length field). That trim is covered instead by the unit tests on the receive loop's
|
||||
# data-frame parse (data_payload in src/transport/ethernet/mod.rs).
|
||||
assertions:
|
||||
baseline:
|
||||
min_nodes_reporting: 6
|
||||
max_roots: 1
|
||||
min_nodes_parented: 5
|
||||
delivery:
|
||||
pairs: [[n04, n06], [n06, n04], [n04, n05]]
|
||||
payload_bytes: [0, 1200]
|
||||
min_replies: 3
|
||||
deadline_secs: 60
|
||||
require_transport: ethernet
|
||||
transport_node: n04
|
||||
|
||||
logging:
|
||||
rust_log: "info"
|
||||
|
||||
@@ -50,11 +50,12 @@ traffic:
|
||||
# A 4-node mesh forms a spanning tree: one root and three nodes
|
||||
# with a parent. All six provably-completed archived runs show exactly
|
||||
# that, so these are the shape of a converged mesh rather than a
|
||||
# tolerance fitted to observations. This is the only assertion covering
|
||||
# Ethernet transport anywhere in CI, and it covers the control plane
|
||||
# only: traffic is disabled above, so no datagram crosses an Ethernet
|
||||
# link in any test. Framing, the length field that trims NIC minimum-
|
||||
# frame padding, and AEAD over Ethernet are all unexercised as a result.
|
||||
# tolerance fitted to observations. This scenario covers the Ethernet
|
||||
# control plane only (beacon discovery, peering, the tree): traffic is
|
||||
# disabled above and no datagram is probed here. Data-plane delivery
|
||||
# over Ethernet (framing, AEAD over AF_PACKET) is asserted in
|
||||
# ethernet-mesh's delivery assertion, and the length field that trims
|
||||
# minimum-frame padding by unit tests on the receive loop's parse.
|
||||
assertions:
|
||||
baseline:
|
||||
min_nodes_reporting: 4
|
||||
|
||||
@@ -10,6 +10,15 @@ Currently supported assertions:
|
||||
- ``bloom_send_rate``: per-node trailing-window ceiling on
|
||||
``stats.bloom.sent`` delta. Calibrated for the bloom-storm
|
||||
regression scenario but generally usable.
|
||||
- ``min_parent_switches`` / ``max_parent_switches``: bounds on the
|
||||
parent switches counted from the node logs.
|
||||
- ``max_errors``: ceiling on ERROR-level log lines (applied by default).
|
||||
- ``baseline``: the mesh formed, agreed on a root and took parents.
|
||||
- ``tree_parents``: named nodes ended with the named parents.
|
||||
- ``congestion_signals``: floors on how many nodes saw each congestion
|
||||
counter.
|
||||
- ``delivery``: datagrams of each configured size were delivered
|
||||
between node pairs, optionally proven to cross one transport.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -22,6 +31,7 @@ from .scenario import (
|
||||
BaselineAssertion,
|
||||
BloomSendRateAssertion,
|
||||
CongestionSignalsAssertion,
|
||||
DeliveryAssertion,
|
||||
MaxErrorsAssertion,
|
||||
MaxParentSwitchesAssertion,
|
||||
MinParentSwitchesAssertion,
|
||||
@@ -370,6 +380,92 @@ def evaluate_congestion_signals(
|
||||
)
|
||||
|
||||
|
||||
_TRANSPORT_LABEL = {"ethernet": "Ethernet", "udp": "UDP", "tcp": "TCP"}
|
||||
|
||||
|
||||
def _transport_failure(transport: dict) -> str | None:
|
||||
"""Return why the transport proof failed, or None when it holds."""
|
||||
node, want = transport["node"], transport["want"]
|
||||
for read in transport["reads"]:
|
||||
when, peers = read["when"], read["peers"]
|
||||
if peers is None:
|
||||
return (
|
||||
f"traversal unproven: show_peers on {node} failed {when} the "
|
||||
f"probe"
|
||||
)
|
||||
if not peers:
|
||||
return (
|
||||
f"traversal unproven: {node} reported zero peers {when} the "
|
||||
f"probe, after waiting {read['waited_s']:.0f}s"
|
||||
)
|
||||
other = [
|
||||
f"{str(p.get('npub', '?'))[:16]} via "
|
||||
f"{p.get('transport_type', 'no transport_type')}"
|
||||
for p in peers if p.get("transport_type") != want
|
||||
]
|
||||
if other:
|
||||
return (
|
||||
f"{node} has a non-{_TRANSPORT_LABEL.get(want, want)} peer "
|
||||
f"{when} the probe ({'; '.join(other)}), so a probe to or "
|
||||
f"from it may not have crossed {want}"
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
def evaluate_delivery(
|
||||
cfg: DeliveryAssertion,
|
||||
probe: dict | None,
|
||||
) -> AssertionOutcome:
|
||||
"""Every configured pair and size got ``min_replies`` replies in time.
|
||||
|
||||
``probe`` is the runner's delivery probe result. None means the probe
|
||||
never ran, which fails as a harness failure: no probe and no delivery
|
||||
would otherwise look alike.
|
||||
"""
|
||||
if probe is None:
|
||||
return AssertionOutcome(
|
||||
name="delivery",
|
||||
passed=False,
|
||||
detail=(
|
||||
"FAIL delivery: the delivery probe never ran, so nothing was "
|
||||
"observed. This is a harness failure, not a statement about "
|
||||
"delivery."
|
||||
),
|
||||
)
|
||||
|
||||
failures = []
|
||||
transport = probe.get("transport")
|
||||
if transport is not None:
|
||||
why = _transport_failure(transport)
|
||||
if why:
|
||||
failures.append(why)
|
||||
|
||||
parts = []
|
||||
for r in probe["results"]:
|
||||
label = f"{r['src']}->{r['dst']} {r['size']}B"
|
||||
if r["replies"] >= cfg.min_replies:
|
||||
parts.append(f"{label} {r['replies']}/{r['attempts']} in {r['elapsed_s']:.1f}s")
|
||||
else:
|
||||
failures.append(
|
||||
f"{label} got {r['replies']} of {cfg.min_replies} replies in "
|
||||
f"{r['elapsed_s']:.0f}s ({r['attempts']} attempts); last ping: "
|
||||
f"{r['last_output']!r}"
|
||||
)
|
||||
|
||||
if failures:
|
||||
return AssertionOutcome(
|
||||
name="delivery",
|
||||
passed=False,
|
||||
detail=f"FAIL delivery: {'; '.join(failures)}",
|
||||
)
|
||||
via = f" via {transport['want']} ({transport['node']})" if transport else ""
|
||||
return AssertionOutcome(
|
||||
name="delivery",
|
||||
passed=True,
|
||||
detail=f"PASS delivery{via}: {'; '.join(parts)}",
|
||||
)
|
||||
|
||||
|
||||
def evaluate_max_errors(
|
||||
cfg: MaxErrorsAssertion,
|
||||
errors: list[tuple[str, str]],
|
||||
|
||||
@@ -17,6 +17,7 @@ from .assertions import (
|
||||
BloomSendRateMonitor,
|
||||
evaluate_baseline,
|
||||
evaluate_congestion_signals,
|
||||
evaluate_delivery,
|
||||
evaluate_max_errors,
|
||||
evaluate_max_parent_switches,
|
||||
evaluate_min_parent_switches,
|
||||
@@ -25,6 +26,7 @@ from .assertions import (
|
||||
from .compose import generate_compose
|
||||
from .config_gen import write_configs
|
||||
from .control import (
|
||||
query_peers,
|
||||
query_status,
|
||||
snapshot_all_congestion,
|
||||
snapshot_all_mmp,
|
||||
@@ -96,6 +98,9 @@ class SimRunner:
|
||||
# than as an absence of congestion.
|
||||
self.final_congestion: dict | None = None
|
||||
self.final_tree: dict | None = None
|
||||
# Set by _probe_delivery. None means the probe never ran, which the
|
||||
# delivery assertion reports as a harness failure.
|
||||
self.delivery_probe: dict | None = None
|
||||
|
||||
def _evaluate_max_parent_switches(
|
||||
self, cfg, parent_switches: list[tuple[str, str]]
|
||||
@@ -712,6 +717,17 @@ class SimRunner:
|
||||
# it as absent. Wait for every node to answer first.
|
||||
self._wait_control_sockets(CONTROL_WAIT_SECS)
|
||||
|
||||
# Every flapped link and stopped node is back, so the delivery
|
||||
# probe measures the restored mesh.
|
||||
if self.scenario.assertions.delivery is not None:
|
||||
try:
|
||||
self._probe_delivery()
|
||||
except Exception:
|
||||
# Leaves delivery_probe None, which the assertion
|
||||
# reports as a harness failure; the rest of teardown
|
||||
# still has to run.
|
||||
log.exception("Delivery probe failed")
|
||||
|
||||
# Collect iperf3 throughput results before containers stop
|
||||
if self.traffic_mgr:
|
||||
iperf_results = self.traffic_mgr.collect_results()
|
||||
@@ -797,6 +813,15 @@ class SimRunner:
|
||||
else:
|
||||
log.error("%s", outcome.detail)
|
||||
|
||||
dl_cfg = self.scenario.assertions.delivery
|
||||
if dl_cfg is not None:
|
||||
outcome = evaluate_delivery(dl_cfg, self.delivery_probe)
|
||||
self.assertion_outcomes.append(outcome)
|
||||
if outcome.passed:
|
||||
log.info("%s", outcome.detail)
|
||||
else:
|
||||
log.error("%s", outcome.detail)
|
||||
|
||||
cong_cfg = self.scenario.assertions.congestion_signals
|
||||
if cong_cfg is not None:
|
||||
outcome = evaluate_congestion_signals(
|
||||
@@ -878,6 +903,76 @@ class SimRunner:
|
||||
else:
|
||||
log.info("Control socket wait: all nodes answered after %.1fs", elapsed)
|
||||
|
||||
def _read_peers(self, node_id: str, when: str, wait_secs: float) -> dict:
|
||||
"""Read a node's peers, waiting up to wait_secs for at least one.
|
||||
|
||||
A node restored at teardown may not have re-peered yet, which is
|
||||
not the failure the transport check is for, so an empty list is
|
||||
re-read until the wait runs out. A failed read is retried the same
|
||||
way and reported as None if it never succeeds.
|
||||
"""
|
||||
container = self.topology.container_name(node_id)
|
||||
start = time.monotonic()
|
||||
while True:
|
||||
data = query_peers(container)
|
||||
peers = None if data is None else data.get("peers", [])
|
||||
waited = time.monotonic() - start
|
||||
if peers or waited >= wait_secs:
|
||||
return {"when": when, "peers": peers, "waited_s": waited}
|
||||
time.sleep(2)
|
||||
|
||||
def _ping_until(self, src: str, dst: str, size: int, cfg) -> dict:
|
||||
"""Ping dst from src until cfg.min_replies replies or the deadline."""
|
||||
container = self.topology.container_name(src)
|
||||
target = f"{self.topology.nodes[dst].npub}.fips"
|
||||
cmd = ["docker", "exec", container, "ping6", "-c", "1", "-W", "2",
|
||||
"-s", str(size), target]
|
||||
start = time.monotonic()
|
||||
replies = attempts = 0
|
||||
last = ""
|
||||
while replies < cfg.min_replies and time.monotonic() - start < cfg.deadline_secs:
|
||||
attempts += 1
|
||||
try:
|
||||
res = subprocess.run(cmd, capture_output=True, text=True, timeout=15)
|
||||
lines = (res.stdout + res.stderr).strip().splitlines()
|
||||
last = lines[-1] if lines else ""
|
||||
if res.returncode == 0:
|
||||
replies += 1
|
||||
continue
|
||||
except subprocess.TimeoutExpired:
|
||||
last = "docker exec ping6 timed out"
|
||||
time.sleep(1)
|
||||
return {
|
||||
"src": src, "dst": dst, "size": size, "replies": replies,
|
||||
"attempts": attempts, "elapsed_s": time.monotonic() - start,
|
||||
"last_output": last,
|
||||
}
|
||||
|
||||
def _probe_delivery(self):
|
||||
"""Run the delivery assertion's probes and keep the result."""
|
||||
cfg = self.scenario.assertions.delivery
|
||||
log.info("Probing delivery: %d pair(s) x %d size(s)",
|
||||
len(cfg.pairs), len(cfg.payload_bytes))
|
||||
transport = None
|
||||
if cfg.require_transport:
|
||||
transport = {
|
||||
"node": cfg.transport_node, "want": cfg.require_transport,
|
||||
"reads": [self._read_peers(cfg.transport_node, "before",
|
||||
cfg.deadline_secs)],
|
||||
}
|
||||
results = [
|
||||
self._ping_until(src, dst, size, cfg)
|
||||
for src, dst in cfg.pairs
|
||||
for size in cfg.payload_bytes
|
||||
]
|
||||
if transport is not None:
|
||||
transport["reads"].append(
|
||||
self._read_peers(cfg.transport_node, "after", 0)
|
||||
)
|
||||
self.delivery_probe = {"transport": transport, "results": results}
|
||||
with open(os.path.join(self.output_dir, "delivery-probe.json"), "w") as f:
|
||||
json.dump(self.delivery_probe, f, indent=2)
|
||||
|
||||
def _take_snapshot(self, label: str):
|
||||
"""Query all nodes via control socket and save tree/MMP/congestion snapshots."""
|
||||
if not self.topology:
|
||||
|
||||
@@ -301,6 +301,32 @@ class CongestionSignalsAssertion:
|
||||
min_nodes_ce_received: int | None = None
|
||||
|
||||
|
||||
@dataclass
|
||||
class DeliveryAssertion:
|
||||
"""Datagrams of each size must be delivered between each node pair.
|
||||
|
||||
The probe runs at teardown, after every flapped link and stopped node
|
||||
has been restored. For each pair and payload size it repeats a
|
||||
one-packet ping6 over the mesh until ``min_replies`` replies have come
|
||||
back or ``deadline_secs`` has passed. The verdict is binary per pair
|
||||
and size, not a loss ratio, so a busy host slows the probe without
|
||||
failing it.
|
||||
|
||||
``require_transport`` with ``transport_node`` proves the probes
|
||||
crossed that transport: the node's peers are read before and after
|
||||
the probe, and the assertion fails if the node has no peers or any
|
||||
peer on another transport. Choose a node whose only edges use that
|
||||
transport and put it in every pair.
|
||||
"""
|
||||
|
||||
pairs: list[tuple[str, str]] = field(default_factory=list)
|
||||
payload_bytes: list[int] = field(default_factory=list)
|
||||
min_replies: int = 3
|
||||
deadline_secs: int = 60
|
||||
require_transport: str | None = None
|
||||
transport_node: str | None = None
|
||||
|
||||
|
||||
@dataclass
|
||||
class MaxErrorsAssertion:
|
||||
"""Ceiling on ERROR-level lines across every node's log.
|
||||
@@ -334,6 +360,7 @@ class AssertionsConfig:
|
||||
congestion_signals: CongestionSignalsAssertion | None = None
|
||||
tree_parents: TreeParentsAssertion | None = None
|
||||
baseline: BaselineAssertion | None = None
|
||||
delivery: DeliveryAssertion | None = None
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -414,6 +441,7 @@ _SECTION_KEYS = {
|
||||
"assertions": {
|
||||
"bloom_send_rate", "min_parent_switches", "max_parent_switches",
|
||||
"max_errors", "congestion_signals", "tree_parents", "baseline",
|
||||
"delivery",
|
||||
},
|
||||
"logging": {"rust_log", "output_dir"},
|
||||
}
|
||||
@@ -428,7 +456,14 @@ _ASSERTION_KEYS = {
|
||||
"baseline": {
|
||||
"min_nodes_reporting", "max_roots", "min_nodes_parented", "min_sessions",
|
||||
},
|
||||
"delivery": {
|
||||
"pairs", "payload_bytes", "min_replies", "deadline_secs",
|
||||
"require_transport", "transport_node",
|
||||
},
|
||||
}
|
||||
# Largest delivery probe payload: an ICMPv6 echo of this size fits the
|
||||
# 1280-byte TUN MTU with FIPS overhead to spare.
|
||||
_DELIVERY_MAX_PAYLOAD = 1400
|
||||
_NETEM_POLICY_KEYS = {
|
||||
"delay_ms", "jitter_ms", "loss_pct", "duplicate_pct", "reorder_pct",
|
||||
"corrupt_pct",
|
||||
@@ -776,6 +811,10 @@ def load_scenario(path: str) -> Scenario:
|
||||
"min_nodes_ce_received); a block with none asserts nothing"
|
||||
)
|
||||
s.assertions.congestion_signals = CongestionSignalsAssertion(**floors)
|
||||
if "delivery" in asrt:
|
||||
s.assertions.delivery = _parse_delivery(
|
||||
asrt["delivery"], s.topology.num_nodes
|
||||
)
|
||||
if "tree_parents" in asrt:
|
||||
tp = asrt["tree_parents"]
|
||||
if not isinstance(tp, dict) or not tp:
|
||||
@@ -861,6 +900,106 @@ def load_scenario(path: str) -> Scenario:
|
||||
return s
|
||||
|
||||
|
||||
def _delivery_node(val, num_nodes: int, where: str) -> str:
|
||||
"""Validate one node id named by the delivery assertion."""
|
||||
if not isinstance(val, str) or not _NODE_ID_RE.fullmatch(val):
|
||||
raise ValueError(
|
||||
f"assertions.delivery.{where}: {val!r} is not a node id of the "
|
||||
f"form 'n04'"
|
||||
)
|
||||
idx = int(val[1:])
|
||||
if idx < 1 or idx > num_nodes:
|
||||
raise ValueError(
|
||||
f"assertions.delivery.{where}: '{val}' is outside this "
|
||||
f"scenario's {num_nodes} nodes"
|
||||
)
|
||||
return val
|
||||
|
||||
|
||||
def _delivery_int(dl: dict, key: str, default: int) -> int:
|
||||
"""Read a positive integer setting of the delivery assertion."""
|
||||
val = dl.get(key, default)
|
||||
if isinstance(val, bool) or not isinstance(val, int) or val < 1:
|
||||
raise ValueError(
|
||||
f"assertions.delivery.{key}: must be a positive integer, got {val!r}"
|
||||
)
|
||||
return val
|
||||
|
||||
|
||||
def _parse_delivery(dl, num_nodes: int) -> DeliveryAssertion:
|
||||
"""Parse and validate the delivery assertion block."""
|
||||
if not isinstance(dl, dict):
|
||||
raise ValueError("assertions.delivery: must be a mapping")
|
||||
_reject_unknown(dl, _ASSERTION_KEYS["delivery"], "assertions.delivery")
|
||||
|
||||
raw_pairs = dl.get("pairs")
|
||||
if not isinstance(raw_pairs, list) or not raw_pairs:
|
||||
raise ValueError(
|
||||
"assertions.delivery.pairs: give at least one [src, dst] pair; "
|
||||
"an empty list asserts nothing"
|
||||
)
|
||||
pairs = []
|
||||
for pair in raw_pairs:
|
||||
if not isinstance(pair, list) or len(pair) != 2:
|
||||
raise ValueError(
|
||||
f"assertions.delivery.pairs: each entry must be [src, dst], "
|
||||
f"got {pair!r}"
|
||||
)
|
||||
src = _delivery_node(pair[0], num_nodes, "pairs")
|
||||
dst = _delivery_node(pair[1], num_nodes, "pairs")
|
||||
if src == dst:
|
||||
raise ValueError(
|
||||
f"assertions.delivery.pairs: [{src}, {dst}] pings a node from "
|
||||
f"itself, which crosses no link"
|
||||
)
|
||||
pairs.append((src, dst))
|
||||
|
||||
sizes = dl.get("payload_bytes")
|
||||
if not isinstance(sizes, list) or not sizes:
|
||||
raise ValueError(
|
||||
"assertions.delivery.payload_bytes: give at least one size; an "
|
||||
"empty list asserts nothing"
|
||||
)
|
||||
for size in sizes:
|
||||
if (isinstance(size, bool) or not isinstance(size, int)
|
||||
or not 0 <= size <= _DELIVERY_MAX_PAYLOAD):
|
||||
raise ValueError(
|
||||
f"assertions.delivery.payload_bytes: each size must be an "
|
||||
f"integer from 0 to {_DELIVERY_MAX_PAYLOAD}, got {size!r}"
|
||||
)
|
||||
|
||||
transport = dl.get("require_transport")
|
||||
node = dl.get("transport_node")
|
||||
if (transport is None) != (node is None):
|
||||
raise ValueError(
|
||||
"assertions.delivery: require_transport and transport_node go "
|
||||
"together; one without the other proves nothing"
|
||||
)
|
||||
if transport is not None:
|
||||
if transport not in VALID_TRANSPORTS:
|
||||
raise ValueError(
|
||||
f"assertions.delivery.require_transport: {transport!r} is not "
|
||||
f"one of {', '.join(VALID_TRANSPORTS)}"
|
||||
)
|
||||
node = _delivery_node(node, num_nodes, "transport_node")
|
||||
missing = [p for p in pairs if node not in p]
|
||||
if missing:
|
||||
raise ValueError(
|
||||
f"assertions.delivery: pairs {missing} do not include "
|
||||
f"transport_node '{node}', so the transport check says "
|
||||
f"nothing about them"
|
||||
)
|
||||
|
||||
return DeliveryAssertion(
|
||||
pairs=pairs,
|
||||
payload_bytes=list(sizes),
|
||||
min_replies=_delivery_int(dl, "min_replies", 3),
|
||||
deadline_secs=_delivery_int(dl, "deadline_secs", 60),
|
||||
require_transport=transport,
|
||||
transport_node=node,
|
||||
)
|
||||
|
||||
|
||||
_SUPPRESSING_LOG_LEVELS = ("off", "error", "warn")
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user