perf(bus): measure the posture the relay ships, and prove it took
Some checks failed
helm chart / lint + unittest + render matrix (push) Has been cancelled
helm chart / install on kind (gated) (push) Has been cancelled
helm chart / publish chart to GHCR (push) Has been cancelled
Meridian Harness / Build (aarch64-unknown-linux-musl) (push) Has been cancelled
Meridian Harness / Build (x86_64-unknown-linux-musl) (push) Has been cancelled
Meridian Harness / Publish rolling release (push) Has been cancelled
Meridian Harness / Publish tagged release (push) Has been cancelled
control plane / chart (push) Has been cancelled
control plane / test (push) Has been cancelled
control plane / browser-e2e (push) Has been cancelled
control plane / Build control plane image (linux/amd64) (push) Has been cancelled
control plane / Build control plane image (linux/arm64) (push) Has been cancelled
control plane / Publish signed control plane image (push) Has been cancelled
CI / Detect Changed Paths (push) Has been cancelled
CI / Dead Token Reference Guard (push) Has been cancelled
Docker image / Build (linux/amd64) (push) Has been cancelled
Docker image / Build (linux/arm64) (push) Has been cancelled
Docker image / Build public push gateway (linux/amd64) (push) Has been cancelled
Docker image / Build public push gateway (linux/arm64) (push) Has been cancelled
CI / Rust Lint (push) Has been cancelled
CI / Unit Tests (push) Has been cancelled
CI / Isolated DB Gate (push) Has been cancelled
CI / Desktop Core (push) Has been cancelled
CI / Desktop Smoke E2E (1) (push) Has been cancelled
CI / Desktop Smoke E2E (2) (push) Has been cancelled
CI / Desktop Smoke E2E (3) (push) Has been cancelled
CI / Desktop Smoke E2E (4) (push) Has been cancelled
CI / Desktop (push) Has been cancelled
CI / Desktop E2E Relay (push) Has been cancelled
CI / Desktop E2E Integration (1/2) (push) Has been cancelled
CI / Desktop E2E Integration (2/2) (push) Has been cancelled
CI / Desktop E2E Integration (push) Has been cancelled
CI / Backend Integration (relay e2e) (push) Has been cancelled
CI / Relay E2E (push) Has been cancelled
CI / Web (push) Has been cancelled
CI / Admin Web (push) Has been cancelled
CI / Mobile (push) Has been cancelled
CI / Security (push) Has been cancelled
CI / Server Cross-Compile (push) Has been cancelled
CI / Server Cross-Compile-1 (push) Has been cancelled
CI / Windows Rust (x86_64-pc-windows-msvc) (push) Has been cancelled
CI / Desktop Build (macOS) (push) Has been cancelled
Docker image / Merge release multi-arch manifest (push) Has been cancelled
Docker image / Merge debug multi-arch manifest (push) Has been cancelled
Docker image / Publish public push gateway image (push) Has been cancelled
Some checks failed
helm chart / lint + unittest + render matrix (push) Has been cancelled
helm chart / install on kind (gated) (push) Has been cancelled
helm chart / publish chart to GHCR (push) Has been cancelled
Meridian Harness / Build (aarch64-unknown-linux-musl) (push) Has been cancelled
Meridian Harness / Build (x86_64-unknown-linux-musl) (push) Has been cancelled
Meridian Harness / Publish rolling release (push) Has been cancelled
Meridian Harness / Publish tagged release (push) Has been cancelled
control plane / chart (push) Has been cancelled
control plane / test (push) Has been cancelled
control plane / browser-e2e (push) Has been cancelled
control plane / Build control plane image (linux/amd64) (push) Has been cancelled
control plane / Build control plane image (linux/arm64) (push) Has been cancelled
control plane / Publish signed control plane image (push) Has been cancelled
CI / Detect Changed Paths (push) Has been cancelled
CI / Dead Token Reference Guard (push) Has been cancelled
Docker image / Build (linux/amd64) (push) Has been cancelled
Docker image / Build (linux/arm64) (push) Has been cancelled
Docker image / Build public push gateway (linux/amd64) (push) Has been cancelled
Docker image / Build public push gateway (linux/arm64) (push) Has been cancelled
CI / Rust Lint (push) Has been cancelled
CI / Unit Tests (push) Has been cancelled
CI / Isolated DB Gate (push) Has been cancelled
CI / Desktop Core (push) Has been cancelled
CI / Desktop Smoke E2E (1) (push) Has been cancelled
CI / Desktop Smoke E2E (2) (push) Has been cancelled
CI / Desktop Smoke E2E (3) (push) Has been cancelled
CI / Desktop Smoke E2E (4) (push) Has been cancelled
CI / Desktop (push) Has been cancelled
CI / Desktop E2E Relay (push) Has been cancelled
CI / Desktop E2E Integration (1/2) (push) Has been cancelled
CI / Desktop E2E Integration (2/2) (push) Has been cancelled
CI / Desktop E2E Integration (push) Has been cancelled
CI / Backend Integration (relay e2e) (push) Has been cancelled
CI / Relay E2E (push) Has been cancelled
CI / Web (push) Has been cancelled
CI / Admin Web (push) Has been cancelled
CI / Mobile (push) Has been cancelled
CI / Security (push) Has been cancelled
CI / Server Cross-Compile (push) Has been cancelled
CI / Server Cross-Compile-1 (push) Has been cancelled
CI / Windows Rust (x86_64-pc-windows-msvc) (push) Has been cancelled
CI / Desktop Build (macOS) (push) Has been cancelled
Docker image / Merge release multi-arch manifest (push) Has been cancelled
Docker image / Merge debug multi-arch manifest (push) Has been cancelled
Docker image / Publish public push gateway image (push) Has been cancelled
Every Zenoh figure this repo has taken was measured with
`scouting/multicast/enabled: false` and nothing else set. Three keys the relay
pins were left at their upstream `true`, and one of them moved numbers.
`transport.shared_memory.enabled` and
`transport.shared_memory.transport_optimization.enabled` BOTH default true, and
the `eclipse-zenoh` wheel is built with the `shared-memory` feature. Verified on
this binding: with both unset, `zenoh_transport::unicast::manager` logs
`shm: Some(TransportShmConfig { .. })` and `zenoh_shm` allocates the 16 MiB pool
on the first >=3 KB put; with both false it logs `shm: None`. So every peer-mode
payload at or above the 3,072 B `message_size_threshold` -- 4 KB, 16 KB, 64 KB,
256 KB, which is the entire upper half of the crossover curve -- travelled
through POSIX shared memory, on a path the relay build cannot take at all
because `shared-memory` is not in its `zenoh` feature list.
Gossip was the second: measured with one publisher and four subscribers, gossip
autoconnect opened 24 transports where the shipped posture opens 8.
The harness now applies the relay's embedded posture key-for-key from
`bus-peer.json5`, READS EVERY KEY BACK (zenoh 1.9/1.10 accept `routing.peer.mode`
and silently drop it, so a clean insert proves nothing), prints the whole posture
once per run, and stamps a digest beside every table. A unit test pins the table
against the fixture so the two cannot diverge again.
Arms: the deployed shape is `client` against a real zenohd, so that arm runs
first and `--zenoh-session-mode` defaults to `client`; peer-direct arms are
labelled a transport floor rather than a routing path; `peer` + `--zenoh-connect`
is refused by name as the meridian-2m45 trap. `--zenoh-posture legacy-0a2` and
matrix arms J/K reproduce the old posture so the delta is measured, not argued --
the same role arm H plays for the superseded synchronous Redis publisher.
The deployed router config left `transport_optimization` unset while the tested
fixture pinned it false. It is fixed, along with two more drifts the new check
found (`listen.exit_on_failure`, `timestamping.drop_future_timestamp`), and
`daemon-config-check.sh` now runs `deploy/compose/zenoh/zenohd.json5` through the
same daemon and diffs its EFFECTIVE config against the fixture's. Differences are
declared with reasons; anything undeclared fails, and a declared difference that
no longer exists fails too. Reverting the fix reproduces exit 1.
Phase 0A.2's figures are marked, not deleted: which survive the posture change
(256 B FAIL, the sub-2 KB rows, the Redis-only arms, the Zenoh-vs-Zenoh ratios)
and which must be re-taken (the whole >=4 KB crossover, the express table above
4 KB, and the Docker-boundary reading of the peer-vs-router divergence).
No measurement was run. Host load ~30.
Signed-off-by: Joshua Belke <joshua@innovationhub-act.org>
This commit is contained in:
parent
445b241202
commit
af4b92eba4
8 changed files with 945 additions and 60 deletions
|
|
@ -56,10 +56,33 @@ that `Initial conf` line, not the file it handed over. It also fails if the
|
|||
daemon ever stops force-enabling those two, so the workaround cannot outlive
|
||||
the quirk it works around.
|
||||
|
||||
## The deployed router is checked here too, and against this one
|
||||
|
||||
`deploy/compose/zenoh/zenohd.json5` — the file Compose mounts — is run through
|
||||
the same daemon and the same assertions as `bus-router.json5`, and then the two
|
||||
resolved configurations are **diffed**.
|
||||
|
||||
That diff is the point. A key omitted from one file resolves to the daemon's
|
||||
default, so it can differ from a key the other file pins without either file
|
||||
ever mentioning the same line — which is exactly how the deployed router came to
|
||||
leave `transport.shared_memory.transport_optimization.enabled` unset (daemon
|
||||
default `true`) while this directory pinned it `false`. Nothing failed, because
|
||||
nothing compared them; the only way to see it was to read two files side by
|
||||
side. `enabled: false` neutralised it in practice
|
||||
(`zenoh-transport-1.8.0/src/shm_context.rs:73` returns before reading the second
|
||||
switch), so the drift was latent rather than live — which is why it survived.
|
||||
|
||||
Every legitimate difference between the two is **declared** in
|
||||
`assert_declared_parity`, with the reason, and anything else fails. A declared
|
||||
difference that has been resolved also fails, so the list cannot become a place
|
||||
to park drift. It currently holds three entries, one of which is a genuine
|
||||
disagreement recorded rather than settled: the deployed router stamps an HLC
|
||||
(`timestamping.enabled.router: true`) and this fixture does not.
|
||||
|
||||
## Running the daemon side
|
||||
|
||||
```bash
|
||||
./daemon-config-check.sh # standalone; needs Docker + python3
|
||||
./daemon-config-check.sh # fixtures + deployed router + parity
|
||||
MERIDIAN_ZENOH_DAEMON_CHECK=1 \
|
||||
cargo test -p meridian-pubsub --features zenoh --test zenoh_config
|
||||
```
|
||||
|
|
|
|||
|
|
@ -16,7 +16,17 @@
|
|||
# overturned. This script therefore parses the "Initial conf" line zenohd logs at
|
||||
# start-up, which is the configuration it actually runs.
|
||||
#
|
||||
# ./daemon-config-check.sh # positives + negatives
|
||||
# It also runs `deploy/compose/zenoh/zenohd.json5` — the file Compose actually
|
||||
# mounts — through the same daemon and the same assertions, and then diffs the
|
||||
# two effective router configurations against a DECLARED difference list. Before
|
||||
# that existed, the deployed file left
|
||||
# `transport.shared_memory.transport_optimization.enabled` unset (daemon default
|
||||
# `true`) while `bus-router.json5` pinned it `false`, and nothing anywhere
|
||||
# noticed: the tested router and the deployed router disagreed on a switch the
|
||||
# fixture's own comment calls load-bearing, and the only way to find it was to
|
||||
# read two files side by side.
|
||||
#
|
||||
# ./daemon-config-check.sh # positives + negatives + deploy parity
|
||||
#
|
||||
# Requires: docker, python3. Exits non-zero on the first violation.
|
||||
set -euo pipefail
|
||||
|
|
@ -56,13 +66,21 @@ note "$IMAGE @ $ACTUAL_DIGEST"
|
|||
# only fails on a loaded machine — and `sleep` is unavailable in some sandboxes
|
||||
# this runs in, where it exits 127 and `set -e` aborts with no message at all,
|
||||
# which reads exactly like "the daemon rejected the config".
|
||||
run_daemon() { # run_daemon <fixture-path> <logfile>
|
||||
local fixture="$1" logfile="$2" cid attempts=0
|
||||
run_daemon() { # run_daemon <config-path> <logfile>
|
||||
# The path may be relative to this directory (a fixture) or absolute (the
|
||||
# deployed router config). Both are handed to the SAME daemon invocation, so
|
||||
# a difference in the result is a difference in the config and nothing else.
|
||||
local fixture="$1" logfile="$2" cid attempts=0 src
|
||||
case "$fixture" in
|
||||
/*) src="$fixture" ;;
|
||||
*) src="$HERE/$fixture" ;;
|
||||
esac
|
||||
[ -f "$src" ] || fail "$fixture: no such config at $src"
|
||||
cid="$(docker create --network none --entrypoint /zenohd "$IMAGE" -c /fixture.json5)" \
|
||||
|| fail "$fixture: could not create a container. That is a Docker problem, not
|
||||
a verdict on the config — re-run before believing it."
|
||||
docker cp "$HERE/$fixture" "$cid:/fixture.json5" >/dev/null \
|
||||
|| fail "$fixture: could not copy the fixture into $cid"
|
||||
docker cp "$src" "$cid:/fixture.json5" >/dev/null \
|
||||
|| fail "$fixture: could not copy the config into $cid"
|
||||
docker start "$cid" >/dev/null \
|
||||
|| fail "$fixture: could not start $cid"
|
||||
# Four polls per second, up to STARTUP_WAIT seconds. Exits early the moment the
|
||||
|
|
@ -79,12 +97,19 @@ run_daemon() { # run_daemon <fixture-path> <logfile>
|
|||
docker rm -f "$cid" >/dev/null 2>&1 || true
|
||||
}
|
||||
|
||||
assert_effective_posture() { # assert_effective_posture <label> <logfile> <mode>
|
||||
local label="$1" logfile="$2" mode="$3"
|
||||
python3 - "$label" "$logfile" "$mode" <<'PY'
|
||||
# assert_effective_posture <label> <logfile> <mode> <expect-json>
|
||||
#
|
||||
# `expect-json` carries only the values on which the router fixture, the peer
|
||||
# fixture and the deployed router config are ALLOWED to differ. Everything else
|
||||
# below is an invariant every one of them must satisfy, which is the point: a
|
||||
# posture key that becomes negotiable stops being enforced.
|
||||
assert_effective_posture() {
|
||||
local label="$1" logfile="$2" mode="$3" expect="$4"
|
||||
python3 - "$label" "$logfile" "$mode" "$expect" <<'PY'
|
||||
import json, re, sys
|
||||
|
||||
label, logfile, mode = sys.argv[1], sys.argv[2], sys.argv[3]
|
||||
expect = json.loads(sys.argv[4])
|
||||
text = open(logfile, encoding="utf-8", errors="replace").read()
|
||||
match = re.search(r"Initial conf: (\{.*\})", text)
|
||||
if not match:
|
||||
|
|
@ -113,8 +138,9 @@ want("transport/unicast/lowlatency", False,
|
|||
"incompatible with QoS; zenoh-transport bails at start-up")
|
||||
want("transport/link/tx/queue/batching/enabled", True,
|
||||
"adaptive batching is the small-payload throughput path")
|
||||
want("transport/link/protocols", ["tcp"],
|
||||
"only the link the relay-side build compiles")
|
||||
if expect.get("link_protocols") is not None:
|
||||
want("transport/link/protocols", expect["link_protocols"],
|
||||
"only the link the relay-side build compiles")
|
||||
# `zenohd/src/main.rs:203-214` reads an *unset* multicast key as consent and
|
||||
# force-enables it — `(true, false) => set_enabled(Some(true))`. An explicitly
|
||||
# false key takes the `(false, false)` arm and survives, which is what this
|
||||
|
|
@ -139,8 +165,9 @@ for path in ("scouting/multicast/autoconnect", "scouting/gossip/autoconnect",
|
|||
if any(entries for entries in lists if entries):
|
||||
sys.exit(f"FAIL: {label}: {path} is {value!r}, must name no node type")
|
||||
|
||||
want("timestamping/enabled", {"router": False, "peer": False, "client": False},
|
||||
"ordering and identity come from the event envelope, not a router HLC")
|
||||
want("timestamping/enabled", expect["timestamping"],
|
||||
"ordering and identity come from the event envelope; a router HLC is a "
|
||||
"deliberate exception and has to be declared, never inherited")
|
||||
|
||||
# The part zenohd overrides, and what actually holds instead.
|
||||
if at("plugins_loading/enabled") is not True:
|
||||
|
|
@ -169,18 +196,129 @@ print(f" {label}: effective posture OK (mode={mode})")
|
|||
PY
|
||||
}
|
||||
|
||||
echo "== positive fixtures =="
|
||||
for spec in "bus-router.json5:router" "bus-peer.json5:peer"; do
|
||||
fixture="${spec%%:*}"; mode="${spec##*:}"
|
||||
# assert_declared_parity <tested-router-log> <deployed-router-log>
|
||||
#
|
||||
# The tested router config and the deployed router config are two files, and a
|
||||
# reviewer comparing them reads two files. This compares what the DAEMON
|
||||
# resolved from each, so an omitted key -- which resolves to a default the other
|
||||
# file pinned -- shows up as a difference even though the two files never
|
||||
# mention the same line. Every legitimate difference is declared below with its
|
||||
# reason; anything else fails, and a declared difference that has stopped
|
||||
# existing fails too, so the list cannot rot into an excuse.
|
||||
assert_declared_parity() {
|
||||
python3 - "$1" "$2" <<'PY'
|
||||
import json, re, sys
|
||||
|
||||
def initial_conf(path, label):
|
||||
text = open(path, encoding="utf-8", errors="replace").read()
|
||||
match = re.search(r"Initial conf: (\{.*\})", text)
|
||||
if not match:
|
||||
sys.exit(f"FAIL: {label}: zenohd never logged an initial configuration")
|
||||
return json.loads(match.group(1))
|
||||
|
||||
tested = initial_conf(sys.argv[1], "bus-router.json5")
|
||||
deployed = initial_conf(sys.argv[2], "deploy/compose/zenoh/zenohd.json5")
|
||||
|
||||
DECLARED = {
|
||||
"metadata": (
|
||||
"the deployed router names itself `meridian-zenohd`; the fixture is anonymous"
|
||||
),
|
||||
"timestamping/enabled/router": (
|
||||
"DISAGREEMENT, declared rather than resolved: the deployed router stamps an "
|
||||
"HLC on data that arrives without one, the fixture stamps nothing. Whichever "
|
||||
"is right, both files should say it -- this entry is the standing record "
|
||||
"that they currently do not, so the next reader finds it here instead of by "
|
||||
"diffing two configs"
|
||||
),
|
||||
"transport/link/protocols": (
|
||||
"the fixture pins the one link the relay-side build compiles; the deployed "
|
||||
"router leaves the daemon's full list. Both listen only on tcp/0.0.0.0:7447"
|
||||
),
|
||||
}
|
||||
|
||||
def flatten(node, prefix=""):
|
||||
if isinstance(node, dict) and node:
|
||||
for key, value in node.items():
|
||||
yield from flatten(value, f"{prefix}/{key}" if prefix else key)
|
||||
else:
|
||||
yield prefix, node
|
||||
|
||||
left, right = dict(flatten(tested)), dict(flatten(deployed))
|
||||
ABSENT = "<absent>"
|
||||
differences = []
|
||||
for key in sorted(set(left) | set(right)):
|
||||
a, b = left.get(key, ABSENT), right.get(key, ABSENT)
|
||||
if a != b:
|
||||
differences.append((key, a, b))
|
||||
|
||||
def covered_by(key):
|
||||
for declared in DECLARED:
|
||||
if key == declared or key.startswith(declared + "/"):
|
||||
return declared
|
||||
return None
|
||||
|
||||
undeclared = [d for d in differences if covered_by(d[0]) is None]
|
||||
seen = {covered_by(key) for key, _a, _b in differences} - {None}
|
||||
stale = sorted(set(DECLARED) - seen)
|
||||
|
||||
if undeclared:
|
||||
lines = "\n".join(
|
||||
f" {key}: tested={json.dumps(a)} deployed={json.dumps(b)}"
|
||||
for key, a, b in undeclared
|
||||
)
|
||||
sys.exit(
|
||||
"FAIL: the tested router config and the deployed router config resolve to "
|
||||
"different effective postures on undeclared keys:\n" + lines + "\n"
|
||||
" Make the two agree, or add the key to DECLARED in daemon-config-check.sh "
|
||||
"with the reason it may differ. Silence is what let "
|
||||
"transport.shared_memory.transport_optimization drift."
|
||||
)
|
||||
if stale:
|
||||
sys.exit(
|
||||
"FAIL: these differences are declared in daemon-config-check.sh but no longer "
|
||||
"exist: " + ", ".join(stale) + ". Delete the entries -- a declared difference "
|
||||
"that has been resolved is a stale excuse waiting to cover the next one."
|
||||
)
|
||||
|
||||
for key, reason in sorted(DECLARED.items()):
|
||||
print(f" declared difference — {key}: {reason}")
|
||||
print(" tested and deployed router configs agree on every other effective key")
|
||||
PY
|
||||
}
|
||||
|
||||
# Every node MERIDIAN runs, checked against the daemon that runs it.
|
||||
# bus-router.json5 — the router posture the Rust library is tested against
|
||||
# bus-peer.json5 — the posture the relay binary embeds
|
||||
# zenohd.json5 — the file Compose actually mounts
|
||||
ALL_OFF='{"router": false, "peer": false, "client": false}'
|
||||
DEPLOYED_ROUTER_CONFIG="$(cd "$HERE/../../../../.." && pwd)/deploy/compose/zenoh/zenohd.json5"
|
||||
FIXTURE_ROUTER_LOG="$(mktemp)"
|
||||
DEPLOYED_ROUTER_LOG="$(mktemp)"
|
||||
|
||||
echo "== positive configs =="
|
||||
for spec in \
|
||||
"bus-router.json5|router|{\"timestamping\": $ALL_OFF, \"link_protocols\": [\"tcp\"]}" \
|
||||
"bus-peer.json5|peer|{\"timestamping\": $ALL_OFF, \"link_protocols\": [\"tcp\"]}" \
|
||||
"$DEPLOYED_ROUTER_CONFIG|router|{\"timestamping\": {\"router\": true, \"peer\": false, \"client\": false}, \"link_protocols\": null}" \
|
||||
; do
|
||||
fixture="${spec%%|*}"; rest="${spec#*|}"; mode="${rest%%|*}"; expect="${rest#*|}"
|
||||
log="$(mktemp)"
|
||||
run_daemon "$fixture" "$log"
|
||||
read -r running exit_code < "$log.state"
|
||||
[ "$running" = "true" ] || fail "$fixture: zenohd exited (code $exit_code); it must accept this config
|
||||
$(cat "$log")"
|
||||
assert_effective_posture "$fixture" "$log" "$mode"
|
||||
assert_effective_posture "$(basename "$fixture")" "$log" "$mode" "$expect"
|
||||
case "$fixture" in
|
||||
*bus-router.json5) cp "$log" "$FIXTURE_ROUTER_LOG" ;;
|
||||
"$DEPLOYED_ROUTER_CONFIG") cp "$log" "$DEPLOYED_ROUTER_LOG" ;;
|
||||
esac
|
||||
rm -f "$log" "$log.state"
|
||||
done
|
||||
|
||||
echo "== tested router vs deployed router =="
|
||||
assert_declared_parity "$FIXTURE_ROUTER_LOG" "$DEPLOYED_ROUTER_LOG"
|
||||
rm -f "$FIXTURE_ROUTER_LOG" "$DEPLOYED_ROUTER_LOG"
|
||||
|
||||
echo "== negative fixtures =="
|
||||
# `unsupported-endpoint-protocol.json5` is deliberately NOT in this list: the
|
||||
# daemon image is built with the TLS link compiled in, so `tls/` is valid there.
|
||||
|
|
|
|||
|
|
@ -45,15 +45,31 @@
|
|||
// Explicit endpoints only — both ambient discovery mechanisms are off below.
|
||||
// A router that discovers its own peers is a federation primitive; this
|
||||
// feature deploys routers INSIDE one Meridian deployment (§ Security).
|
||||
listen: { endpoints: ["tcp/0.0.0.0:7447"] },
|
||||
// `exit_on_failure` is pinned rather than defaulted for the same reason the
|
||||
// scouting keys are: an unset key is a bet on an upstream default staying put.
|
||||
// A router that cannot bind its endpoint must die at start-up rather than run
|
||||
// as an unreachable no-op that answers every liveness probe.
|
||||
listen: { endpoints: ["tcp/0.0.0.0:7447"], exit_on_failure: true },
|
||||
connect: { endpoints: [] },
|
||||
|
||||
scouting: {
|
||||
// Docker multicast between host and container does not work anyway, but
|
||||
// that is not why this is off: ambient discovery is refused on principle.
|
||||
multicast: { enabled: false },
|
||||
// The autoconnect/listen lists below are inert while `enabled` is false, and
|
||||
// that is exactly why they are here: they mean a single flag flip cannot
|
||||
// turn this router into a node that dials whatever it hears about.
|
||||
multicast: {
|
||||
enabled: false,
|
||||
autoconnect: { router: [], peer: [], client: [] },
|
||||
listen: { router: false, peer: false, client: false },
|
||||
},
|
||||
// No CLI equivalent. Removing this line silently re-enables gossip.
|
||||
gossip: { enabled: false },
|
||||
gossip: {
|
||||
enabled: false,
|
||||
multihop: false,
|
||||
target: { router: [], peer: [] },
|
||||
autoconnect: { router: [], peer: [], client: [] },
|
||||
},
|
||||
},
|
||||
|
||||
transport: {
|
||||
|
|
@ -78,12 +94,42 @@
|
|||
// gated on a measured crossover plus `ulimit -l unlimited`/CAP_IPC_LOCK,
|
||||
// so leaving this unset silently activates a phase that has not been
|
||||
// decided. Explicit false, and Phase 6 flips it behind its own flag.
|
||||
shared_memory: { enabled: false },
|
||||
//
|
||||
// `transport_optimization` is a SECOND, separately-defaulted switch (also
|
||||
// `true`) that implicitly routes every message at or above
|
||||
// `message_size_threshold` — 3,072 B — into shared memory. It used to be
|
||||
// absent here while `crates/meridian-pubsub/tests/fixtures/zenoh/bus-router.json5`
|
||||
// pinned it false, so the tested router and the deployed router disagreed on
|
||||
// a switch that file's own comment calls load-bearing. `enabled: false`
|
||||
// already neutralises it today (`zenoh-transport-1.8.0/src/shm_context.rs:73`
|
||||
// returns before reading it), which is why the drift was latent rather than
|
||||
// live — and precisely why it could sit here unnoticed until someone flipped
|
||||
// `enabled`. Both switches now say what they mean, and
|
||||
// `daemon-config-check.sh` compares this file's EFFECTIVE posture against
|
||||
// the fixture's, so the next divergence fails a gate instead of waiting for
|
||||
// a human to read two files side by side.
|
||||
shared_memory: {
|
||||
enabled: false,
|
||||
mode: "lazy",
|
||||
transport_optimization: { enabled: false },
|
||||
},
|
||||
},
|
||||
|
||||
// Routers stamp an HLC timestamp on data that arrives without one. Peers and
|
||||
// clients do not — the relay is the authority for event time.
|
||||
timestamping: { enabled: { router: true, peer: false, client: false } },
|
||||
//
|
||||
// NOTE, and it is recorded rather than resolved here: the tested router
|
||||
// fixture (`crates/meridian-pubsub/tests/fixtures/zenoh/bus-router.json5`)
|
||||
// sets `router: false` on the grounds that ordering and identity come from the
|
||||
// event envelope, not a router HLC. These two files therefore disagree about
|
||||
// whether the router stamps. `daemon-config-check.sh` carries that
|
||||
// disagreement in its DECLARED list so it is visible rather than latent;
|
||||
// whichever answer is right, both files should give the same one.
|
||||
timestamping: {
|
||||
enabled: { router: true, peer: false, client: false },
|
||||
// Retimestamp rather than drop, pinned rather than inherited.
|
||||
drop_future_timestamp: false,
|
||||
},
|
||||
|
||||
// Overridden to `true` by zenohd (see header). `--adminspace-permissions
|
||||
// none` in the compose `command:` is what closes it; these values keep the
|
||||
|
|
|
|||
|
|
@ -68,7 +68,70 @@ pip install eclipse-zenoh
|
|||
|
||||
If the package is missing, `--mode zenoh` exits 2 with that install line. Unit tests skip live zenoh rather than failing CI.
|
||||
|
||||
The arm uses **peer mode**, explicit `connect`/`listen` on `tcp/127.0.0.1:<ephemeral>`, and `scouting/multicast/enabled: false`. It does **not** start or require `zenohd` (that is Phase 5 / `meridian-0ps`) and it does not enable shared-memory transport.
|
||||
### The posture is the profile
|
||||
|
||||
Every Zenoh session opens under **the posture the relay binary embeds** —
|
||||
`RELAY_ZENOH_POSTURE` in `relay_bus_scaling.py`, copied key-for-key from
|
||||
`crates/meridian-pubsub/tests/fixtures/zenoh/bus-peer.json5` and pinned against it
|
||||
by a unit test. That is: both scouting mechanisms present-and-false, timestamping
|
||||
off, TCP only, QoS on, lowlatency off, adaptive batching on, **shared memory and
|
||||
its large-message `transport_optimization` both off**, no plugins, no admin space,
|
||||
explicit endpoints.
|
||||
|
||||
Each key is inserted **and read back** from the live `zenoh.Config`. A key that is
|
||||
refused, dropped, or coerced fails the run rather than becoming a footnote — zenoh
|
||||
1.9 and 1.10 accept `routing.peer.mode` and silently discard it, so "the insert did
|
||||
not raise" is not evidence that a posture took.
|
||||
|
||||
The harness prints the whole posture once per run and stamps a one-line digest
|
||||
beside every table:
|
||||
|
||||
```text
|
||||
zenoh posture=relay#19f85373a192 (multicast=off gossip=off shared_memory=off transport_optimization=off)
|
||||
```
|
||||
|
||||
Two of those keys are traps rather than preferences, and both were measured on this
|
||||
binding rather than read off a document:
|
||||
|
||||
- **`scouting/*/enabled` must be present.** An absent key is consent: the upstream
|
||||
default is `true` for *both* mechanisms. Disabling multicast alone leaves gossip
|
||||
on, and with one publisher and four subscribers that opened **24 transports
|
||||
instead of 8** — the peers gossip-autoconnect into a full mesh inside the
|
||||
measured process.
|
||||
- **Shared memory needs two keys, not one.** `transport.shared_memory.enabled` and
|
||||
`transport.shared_memory.transport_optimization.enabled` both default `true`, and
|
||||
the `eclipse-zenoh` wheel *is* built with the `shared-memory` feature. Left unset,
|
||||
every payload at or above the 3,072 B `message_size_threshold` travels through
|
||||
POSIX shared memory: `zenoh_transport::unicast::manager` logs
|
||||
`shm: Some(TransportShmConfig { .. })` and `zenoh_shm` allocates a 16 MiB pool on
|
||||
the first large put. With both false it logs `shm: None` and allocates nothing.
|
||||
The relay build cannot use SHM at all — `shared-memory` is not in its `zenoh`
|
||||
feature list — so an unset key measures a transport the product does not have, on
|
||||
exactly the payload sizes a crossover curve is read from.
|
||||
|
||||
`--zenoh-posture legacy-0a2` reproduces what this harness applied before
|
||||
2026-08-21 (multicast off, nothing else). It exists so the delta is measurable
|
||||
rather than assumed, it is never the default, and every table it prints says it is
|
||||
not the shipped posture.
|
||||
|
||||
### Topology: which arm is the deployed one
|
||||
|
||||
The relay's default session mode is `client` against the router tier
|
||||
(`ZenohMode::default()`), so:
|
||||
|
||||
- **`--zenoh-connect tcp/…` + `--zenoh-session-mode client`** — a real `zenohd`.
|
||||
This is the shipped shape and the only arm that measures the deployed routing
|
||||
path. `--zenoh-session-mode` now defaults to `client` for that reason.
|
||||
- **no `--zenoh-connect`** — a direct peer-to-peer link on
|
||||
`tcp/127.0.0.1:<ephemeral>`. A useful transport floor, and **not** a routing path
|
||||
the relay uses; the harness labels it that way in its own output.
|
||||
- **`--zenoh-session-mode peer` together with `--zenoh-connect` is refused.** It is
|
||||
the `meridian-2m45` trap: a router will not propagate a subscription declaration
|
||||
between two Peer faces while gossip is off, so both sessions open, both report
|
||||
healthy, and nothing is delivered. `--zenoh-allow-peer-through-router` reproduces
|
||||
it on purpose.
|
||||
|
||||
The harness does **not** start or require `zenohd` (that is Phase 5 / `meridian-0ps`).
|
||||
|
||||
Publish keys follow the feature spec (`/` separator, kind in the key):
|
||||
|
||||
|
|
@ -120,7 +183,8 @@ Zenoh operating points, one flag each, all off the same grid:
|
|||
| cold vs prewarmed | `--zenoh-prewarm N` (default 200) warms the session on an unmatched key and reports the first publication separately; `--zenoh-prewarm 0` leaves the cold cost inside the measured window |
|
||||
| settle | `--zenoh-settle` seconds after peer links, before publishing |
|
||||
| 1 vs N subscribers | `--pods 1,2,4` — one zenoh session with one subscriber per pod |
|
||||
| peer vs router | `--zenoh-connect tcp/127.0.0.1:7447 --zenoh-session-mode client` points every session at a running `zenohd` instead of the in-process peer listener |
|
||||
| peer vs router | `--zenoh-connect tcp/127.0.0.1:7447` points every session at a running `zenohd` instead of the in-process peer listener; `--zenoh-session-mode` already defaults to `client`, which is the deployed shape |
|
||||
| posture | `--zenoh-posture relay` (default, what the relay ships) or `legacy-0a2` (the pre-2026-08-21 artifact — multicast off, gossip and shared memory left ON) |
|
||||
|
||||
Prewarm traffic is published under `meridian-prewarm/v1/…`, which no measured
|
||||
interest matches, so it warms the session and link without touching any count.
|
||||
|
|
@ -146,14 +210,95 @@ python3 -m unittest discover -s perf -p 'test_*.py'
|
|||
|
||||
The unit tests pin the default 1/2/4-pod 64× contract and include a mutant row that represents scoped mode receiving irrelevant global-firehose traffic; that row must fail the assertion. They also pin the zero-delivery rejection, the "deliveries, not per second" column labelling, the zenoh key shape, the payload-sweep sizes, the express/prewarm operating points, and the 5× gate recorder. Live zenoh is skipped unless `eclipse-zenoh` is importable.
|
||||
|
||||
They also pin the posture, which is the part that had no test at all when it drifted:
|
||||
every key of `RELAY_ZENOH_POSTURE` is compared against
|
||||
`crates/meridian-pubsub/tests/fixtures/zenoh/bus-peer.json5` (so the harness cannot
|
||||
diverge from what the relay embeds), both scouting switches and both shared-memory
|
||||
switches are asserted present-and-false by name, a dropped or coerced key is proved
|
||||
to fail the run, `--zenoh-session-mode peer` with `--zenoh-connect` is proved to be
|
||||
refused with `meridian-2m45` in the message, and every printed table is proved to
|
||||
carry its posture stamp.
|
||||
|
||||
The Phase 0A.2 matrix result is recorded below, with its profile. Raw logs are in
|
||||
`perf/results/0a2-20260819/`; `perf/run_0a2_matrix.sh` re-runs the matrix and
|
||||
`perf/summarize_0a2.py` re-derives every table in this document from those logs.
|
||||
|
||||
The matrix now runs the deployed shape (arm **D**, `client` sessions against a real
|
||||
`zenohd`) first, because it is the only arm in it that measures the routing path the
|
||||
relay uses; arms A, B, C and E are direct peer-to-peer links, which is a transport
|
||||
floor. Arms **J** and **K** re-run A and D under `--zenoh-posture legacy-0a2` so the
|
||||
posture correction is a measured delta rather than an argument — the same role arm H
|
||||
plays for the superseded synchronous Redis publisher.
|
||||
|
||||
---
|
||||
|
||||
## Phase 0A.2 result — directional Redis-vs-Zenoh crossover (2026-08-19)
|
||||
|
||||
> ### ⚠ TAKEN UNDER THE WRONG POSTURE — read this before quoting any number below
|
||||
>
|
||||
> Every Zenoh figure in this section was measured with
|
||||
> `scouting/multicast/enabled: false` **and nothing else set**. Three keys the relay
|
||||
> pins were therefore left at their upstream `true`:
|
||||
>
|
||||
> | key | this run | what the relay ships | consequence |
|
||||
> | --- | --- | --- | --- |
|
||||
> | `scouting.gossip.enabled` | **on** | off | measured with 24 transports where the shipped posture opens 8 |
|
||||
> | `transport.shared_memory.enabled` | **on** | off | SHM negotiated on every link (`shm: Some(..)`) |
|
||||
> | `transport.shared_memory.transport_optimization.enabled` | **on** | off | every payload ≥ 3,072 B implicitly routed through a 16 MiB POSIX SHM pool |
|
||||
>
|
||||
> The relay build **cannot use shared memory at all** — `shared-memory` is not in its
|
||||
> `zenoh` feature list — so the ≥ 4 KB rows below were taken on a transport the
|
||||
> product does not have. Nothing here is deleted: these runs are reproducible
|
||||
> evidence, `perf/results/0a2-20260819/` is intact, and `--zenoh-posture legacy-0a2`
|
||||
> plus matrix arms **J** and **K** re-take them under exactly this configuration so
|
||||
> the delta is measured rather than argued. The correction landed 2026-08-21; the
|
||||
> re-run under the shipped posture is a quiet-machine act and has not happened yet.
|
||||
>
|
||||
> **What must be re-taken before it licenses anything:**
|
||||
>
|
||||
> - **The crossover curve (~2× at 2 KB, ~5× at 16 KB, 6.82×/8.21× at 64/256 KB).**
|
||||
> This is the figure a lane migration would be justified on, and it is the figure
|
||||
> most exposed: every point at and above 4,096 B sits above Zenoh's 3,072 B
|
||||
> `message_size_threshold`, so the Zenoh side of each of those ratios had a
|
||||
> shared-memory fast path that the relay cannot use, while the Redis side had no
|
||||
> equivalent. The direction of the error is *toward Zenoh*. Treat the whole ≥ 4 KB
|
||||
> half of the curve as unmeasured, and with it the "16,384 B is where 5× arrives"
|
||||
> claim and the "**Router mode … prices the Docker boundary**" section — SHM works
|
||||
> between two host processes and cannot cross into the Docker VM, so an unknown
|
||||
> share of the peer-vs-router divergence at ≥ 64 KB is the SHM path dropping out,
|
||||
> not the boundary appearing.
|
||||
> - **The `express` table at 4,096 B and above**, for the same threshold reason.
|
||||
>
|
||||
> **What the posture plausibly does not move — plausibly, because nobody has
|
||||
> checked, which is why arms J and K exist:**
|
||||
>
|
||||
> - **The 256 B verdict (1.15–1.51×, FAIL against the ≥5× gate).** 256 B is well
|
||||
> below the 3,072 B SHM threshold, so no SHM path was available to either arm at
|
||||
> that size, and the gate fails by a factor of 3.3 — far outside the ±20–24%
|
||||
> run-to-run spread reported below. Gossip's extra transports are the residual
|
||||
> risk, and they cost the measured process work rather than saving it, so the
|
||||
> correction is unlikely to *raise* this number.
|
||||
> - **The 64 B–2 KB rows generally**, all of which are below the threshold, and all
|
||||
> of which the "both arms are receiver-bound below ~4 KB" finding already says are
|
||||
> a comparison of two Python receive paths rather than of two transports.
|
||||
> - **The Redis-only arms (F, G) and the sync-publisher artifact (H)**, which contain
|
||||
> no Zenoh session at all. The 73.3×/59.2× publisher finding and the 26.87× → 1.51×
|
||||
> correction stand unchanged.
|
||||
> - **The Zenoh-versus-Zenoh ratios**, which hold one configuration on both sides:
|
||||
> cold start 15.5× at 256 B, and the `express`-versus-batched ratios below 4 KB.
|
||||
> The 4,096 B cold-start row (17.9×) is the exception — its first publication also
|
||||
> paid the 16 MiB SHM pool allocation, so that one number is contaminated by an
|
||||
> allocation the relay never performs.
|
||||
> - **The 1-vs-N subscriber column at 256 B.** The 65,536 B and 262,144 B rows of
|
||||
> that table are above the threshold and are not exempt.
|
||||
>
|
||||
> **Which arm was in the deployed shape:** only **D** (client sessions against a real
|
||||
> `zenohd`). Gossip does not apply to a `client` session at all — upstream states
|
||||
> that instances in client mode do not participate in gossip — and SHM cannot cross
|
||||
> the Docker VM boundary, so arm D is the arm the posture correction is *least*
|
||||
> likely to move. That is also the arm whose numbers matter most, and it was run at
|
||||
> `--pods 1` only.
|
||||
|
||||
**This is a Python-binding probe, not the Rust production adapter's capacity, and it
|
||||
licenses no end-to-end MERIDIAN throughput claim.** Every figure below compares two
|
||||
Python clients on one shared laptop. The native number is Phase 0B.1's job. Nothing
|
||||
|
|
@ -177,6 +322,7 @@ this block.
|
|||
| topology | 64 communities × 100 events = one 6,400-event burst per arm; 1 subscribed community; every pod interested in it; pods 1, 2, 4 |
|
||||
| traffic class | opaque bytes on a bus — **nothing is signed, verified, parsed or stored** |
|
||||
| warmup | `--zenoh-prewarm 200` puts on an unmatched key, `--zenoh-settle 0.3` s |
|
||||
| **zenoh posture** | **`legacy-0a2` — multicast scouting off, and nothing else. Gossip scouting, shared memory and its large-message optimization were all left ON. See the warning above.** |
|
||||
| measured window | one 6,400-event burst per arm, 0.04–1.6 s wall clock depending on payload |
|
||||
| sample count | 6,400 deliveries per arm per pod count |
|
||||
| repetitions | **5**; every table reports median with the full min–max across them |
|
||||
|
|
@ -413,8 +559,14 @@ reconnect loss.
|
|||
Python subscriber callback bound most of these measurements. Phase 0B.1 owns that.
|
||||
- **Nothing end-to-end.** Nothing here is signed, verified, parsed, or stored. The
|
||||
BIP-340 ceiling in the root `AGENTS.md` is untouched by any of it.
|
||||
- **Nothing about shared memory.** SHM was not enabled and would not apply across the
|
||||
Docker boundary anyway.
|
||||
- ~~**Nothing about shared memory.** SHM was not enabled and would not apply across
|
||||
the Docker boundary anyway.~~ **This was wrong, and it is the finding that dates
|
||||
this section.** SHM *was* enabled — by omission, because both of its switches
|
||||
default `true` and the `eclipse-zenoh` wheel is built with the feature. It applied
|
||||
to every peer-mode payload at or above 3,072 B. It did not apply across the Docker
|
||||
boundary, which is the half of the original sentence that was true, and which is
|
||||
why the peer and router columns must now be read as two different transports at
|
||||
large payloads rather than as one transport either side of a boundary.
|
||||
- **No verdict on Zenoh at 256 B.** The gate fails there; the crossover carries the
|
||||
story, and the traffic classes that live above 4 KB are where this transport pays.
|
||||
|
||||
|
|
|
|||
|
|
@ -15,12 +15,19 @@ Default scenario:
|
|||
Modes:
|
||||
* redis: measured Redis PUB/SUB delivery using only Python stdlib (default).
|
||||
* model: deterministic no-service arithmetic model.
|
||||
* zenoh: measured Eclipse Zenoh peer-mode delivery (optional eclipse-zenoh).
|
||||
* zenoh: measured Eclipse Zenoh delivery (optional eclipse-zenoh).
|
||||
|
||||
Every Zenoh session opens under the relay's shipped posture -- both scouting
|
||||
mechanisms off, shared memory and its large-message optimization off, TCP only,
|
||||
explicit endpoints, no plugins. That posture is applied AND read back, and it is
|
||||
stamped on every table, because a throughput number whose posture is not beside
|
||||
it says nothing (root AGENTS.md, "Published Throughput Ceilings").
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import resource
|
||||
|
|
@ -43,8 +50,8 @@ DEFAULT_GAIN_THRESHOLD = 5.0
|
|||
ZENOH_INSTALL_HINT = (
|
||||
"--mode zenoh requires the optional 'eclipse-zenoh' package.\n"
|
||||
"Install it with: pip install eclipse-zenoh\n"
|
||||
"This arm uses peer mode with explicit connect and multicast scouting off; "
|
||||
"it does not start zenohd (that is Phase 5 / meridian-0ps).\n"
|
||||
"Every session opens under the relay's shipped posture (see RELAY_ZENOH_POSTURE); "
|
||||
"`--zenoh-connect` points it at a running zenohd, which this harness never starts.\n"
|
||||
"SKIPPED, not zero: nothing below is a measurement when the binding is absent."
|
||||
)
|
||||
# Latency is stamped into the payload as an offset from this process epoch so a
|
||||
|
|
@ -62,6 +69,101 @@ DEFAULT_ZENOH_SETTLE_S = 0.3
|
|||
# against Zenoh's batched put, which returns on enqueue. See RELAY_BUS_SCALING.md.
|
||||
DEFAULT_REDIS_PIPELINE_DEPTH = 256
|
||||
|
||||
# --------------------------------------------------------------------------
|
||||
# Zenoh session posture
|
||||
# --------------------------------------------------------------------------
|
||||
# The relay does not run on Zenoh's defaults, and almost every default it
|
||||
# overrides is the opposite of what MERIDIAN needs. A harness measuring a
|
||||
# different posture is measuring a different transport, so the table below is
|
||||
# copied key-for-key from the config the relay binary embeds:
|
||||
#
|
||||
# crates/meridian-pubsub/tests/fixtures/zenoh/bus-peer.json5
|
||||
# -> meridian_pubsub::zenoh::config::PINNED_PEER_CONFIG
|
||||
#
|
||||
# Two entries are traps rather than preferences, and both were verified against
|
||||
# this binding rather than read off a document:
|
||||
#
|
||||
# * `scouting/*/enabled` must be PRESENT and false. An absent key is read as
|
||||
# consent: the upstream default is `true` for BOTH mechanisms, and `zenohd`
|
||||
# goes further and force-enables multicast when the key is missing. Disabling
|
||||
# multicast alone leaves gossip on -- which, measured here with one publisher
|
||||
# and four subscribers, opened 24 transports instead of 8, because the peers
|
||||
# gossip-autoconnect into a full mesh inside the measured process.
|
||||
# * `transport/shared_memory/enabled` and
|
||||
# `transport/shared_memory/transport_optimization/enabled` BOTH default
|
||||
# `true` (zenoh-config-1.8.0 DEFAULT_CONFIG.json5), and the `eclipse-zenoh`
|
||||
# wheel IS built with the `shared-memory` feature (`zenoh.shm` imports).
|
||||
# Left unset, every payload at or above `message_size_threshold` (3,072 B)
|
||||
# goes through POSIX shared memory: `zenoh_transport::unicast::manager` logs
|
||||
# `shm: Some(TransportShmConfig { .. })` and `zenoh_shm` allocates the 16 MiB
|
||||
# pool on the first large put. With both keys false it logs `shm: None` and
|
||||
# allocates nothing. The relay build cannot use SHM at all -- `shared-memory`
|
||||
# is not in its `zenoh` feature list -- so an unset key measures a transport
|
||||
# the product does not have, on exactly the payload sizes the crossover curve
|
||||
# is read from.
|
||||
#
|
||||
# Deliberately NOT copied from the fixture: `connect.exit_on_failure: false` and
|
||||
# `listen.exit_on_failure: true`. Those are relay *availability* policy -- the
|
||||
# relay must serve NIP-01 with the bus down -- and copying them would turn a
|
||||
# mistyped endpoint into a silent zero-delivery timeout instead of an error.
|
||||
RELAY_ZENOH_POSTURE: tuple[tuple[str, object], ...] = (
|
||||
# No implicit discovery of any kind. Both switches are ON upstream.
|
||||
("scouting/multicast/enabled", False),
|
||||
("scouting/multicast/autoconnect", {"router": [], "peer": [], "client": []}),
|
||||
("scouting/multicast/listen", {"router": False, "peer": False, "client": False}),
|
||||
("scouting/gossip/enabled", False),
|
||||
("scouting/gossip/multihop", False),
|
||||
("scouting/gossip/target", {"router": [], "peer": []}),
|
||||
("scouting/gossip/autoconnect", {"router": [], "peer": [], "client": []}),
|
||||
# Ordering and identity come from the event envelope, not a router HLC.
|
||||
("timestamping/enabled", {"router": False, "peer": False, "client": False}),
|
||||
("timestamping/drop_future_timestamp", False),
|
||||
# `open()` must not block on a peer that may not exist yet.
|
||||
("open/return_conditions/connect_scouted", False),
|
||||
("open/return_conditions/declares", False),
|
||||
# lowlatency is incompatible with qos, and the Q0-Q7 lane mapping needs qos.
|
||||
("transport/unicast/lowlatency", False),
|
||||
("transport/unicast/qos/enabled", True),
|
||||
# Exactly the link the relay build compiles.
|
||||
("transport/link/protocols", ["tcp"]),
|
||||
# Adaptive batching lives under `queue`, not directly under `tx`.
|
||||
("transport/link/tx/queue/batching/enabled", True),
|
||||
("transport/link/tx/queue/batching/time_limit", 1),
|
||||
# Charter Phase 6 is unscheduled; see the note above for what unset costs.
|
||||
("transport/shared_memory/enabled", False),
|
||||
("transport/shared_memory/mode", "lazy"),
|
||||
("transport/shared_memory/transport_optimization/enabled", False),
|
||||
# No plugin runtime, no storage manager, no REST, no admin space.
|
||||
("plugins_loading/enabled", False),
|
||||
("plugins_loading/search_dirs", []),
|
||||
("adminspace/enabled", False),
|
||||
("adminspace/permissions", {"read": False, "write": False}),
|
||||
)
|
||||
|
||||
# The posture the harness actually applied for the Phase 0A.2 run of 2026-08-19:
|
||||
# multicast off, nothing else. Kept selectable so the delta the correction
|
||||
# introduces is MEASURABLE rather than assumed -- the same reason
|
||||
# `--redis-publish sync` keeps the discarded round-trip-bound publisher
|
||||
# reproducible instead of merely described. Never the default, never a result.
|
||||
LEGACY_0A2_ZENOH_POSTURE: tuple[tuple[str, object], ...] = (
|
||||
("scouting/multicast/enabled", False),
|
||||
)
|
||||
|
||||
ZENOH_POSTURES: dict[str, tuple[tuple[str, object], ...]] = {
|
||||
"relay": RELAY_ZENOH_POSTURE,
|
||||
"legacy-0a2": LEGACY_0A2_ZENOH_POSTURE,
|
||||
}
|
||||
DEFAULT_ZENOH_POSTURE = "relay"
|
||||
ZENOH_POSTURE_SOURCE = "crates/meridian-pubsub/tests/fixtures/zenoh/bus-peer.json5"
|
||||
# The four keys that decide whether this harness is measuring the shipped
|
||||
# transport. Stamped beside every table; the rest print once per run.
|
||||
ZENOH_POSTURE_HEADLINE_KEYS = (
|
||||
"scouting/multicast/enabled",
|
||||
"scouting/gossip/enabled",
|
||||
"transport/shared_memory/enabled",
|
||||
"transport/shared_memory/transport_optimization/enabled",
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Scenario:
|
||||
|
|
@ -358,17 +460,106 @@ def pick_tcp_port() -> int:
|
|||
return int(sock.getsockname()[1])
|
||||
|
||||
|
||||
def zenoh_peer_config(
|
||||
class ZenohPostureError(RuntimeError):
|
||||
"""A posture key was refused, dropped, or coerced by the binding."""
|
||||
|
||||
|
||||
def zenoh_posture(name: str) -> tuple[tuple[str, object], ...]:
|
||||
try:
|
||||
return ZENOH_POSTURES[name]
|
||||
except KeyError:
|
||||
known = ", ".join(sorted(ZENOH_POSTURES))
|
||||
raise ValueError(f"unknown zenoh posture `{name}`; known postures: {known}") from None
|
||||
|
||||
|
||||
def apply_zenoh_posture(conf, posture: tuple[tuple[str, object], ...]) -> None:
|
||||
"""Insert every posture key, then read every one of them back.
|
||||
|
||||
The read-back is the point. `routing/peer/mode` is *accepted and silently
|
||||
dropped* by zenoh 1.9 and 1.10 -- the key never appears in the daemon's own
|
||||
resolved config and nothing errors -- so "the insert did not raise" is no
|
||||
evidence that a key took effect. A posture that cannot be read back is a
|
||||
failed run, not a footnote under a number.
|
||||
"""
|
||||
for key, value in posture:
|
||||
try:
|
||||
conf.insert_json5(key, json.dumps(value))
|
||||
except Exception as exc: # noqa: BLE001 - the binding raises bare exceptions
|
||||
raise ZenohPostureError(
|
||||
f"zenoh refused posture key `{key}` = {value!r}: {exc}. "
|
||||
"The posture is the measurement's profile; a run without it is void."
|
||||
) from exc
|
||||
for key, value in posture:
|
||||
try:
|
||||
observed = json.loads(conf.get_json(key))
|
||||
except Exception as exc: # noqa: BLE001
|
||||
raise ZenohPostureError(
|
||||
f"zenoh accepted posture key `{key}` and then could not read it back: {exc}. "
|
||||
"That is the shape of a silently dropped key -- do not trust this run."
|
||||
) from exc
|
||||
if observed != value:
|
||||
raise ZenohPostureError(
|
||||
f"zenoh posture key `{key}` reads back as {observed!r}, not {value!r}. "
|
||||
"The binding rewrote the posture; the session is not the one being reported."
|
||||
)
|
||||
|
||||
|
||||
def zenoh_posture_digest(posture: tuple[tuple[str, object], ...]) -> str:
|
||||
"""Short stable digest of the applied posture, printed beside every table."""
|
||||
payload = json.dumps([[key, value] for key, value in posture], separators=(",", ":"))
|
||||
return hashlib.sha256(payload.encode("utf-8")).hexdigest()[:12]
|
||||
|
||||
|
||||
def zenoh_posture_stamp(name: str) -> str:
|
||||
"""One line that identifies which transport a number was taken on."""
|
||||
posture = zenoh_posture(name)
|
||||
applied = dict(posture)
|
||||
headline = " ".join(
|
||||
f"{key.rsplit('/', 2)[-2] if key.endswith('/enabled') else key}="
|
||||
f"{'on' if applied.get(key, True) else 'off'}"
|
||||
for key in ZENOH_POSTURE_HEADLINE_KEYS
|
||||
)
|
||||
return f"posture={name}#{zenoh_posture_digest(posture)} ({headline})"
|
||||
|
||||
|
||||
def print_zenoh_posture_block(name: str) -> None:
|
||||
"""Every applied key, once per run, so a result carries its own profile."""
|
||||
posture = zenoh_posture(name)
|
||||
print(f"zenoh posture: {name} (digest {zenoh_posture_digest(posture)})")
|
||||
if name == DEFAULT_ZENOH_POSTURE:
|
||||
print(f" source: {ZENOH_POSTURE_SOURCE} — the config the relay binary embeds")
|
||||
else:
|
||||
print(
|
||||
" NOT the shipped posture. This is an artifact arm kept so the delta "
|
||||
"against the relay posture stays measurable; it licenses nothing."
|
||||
)
|
||||
print(" applied and read back from the live zenoh.Config:")
|
||||
for key, value in posture:
|
||||
print(f" {key} = {json.dumps(value)}")
|
||||
print(
|
||||
" not copied from the fixture: connect.exit_on_failure, listen.exit_on_failure "
|
||||
"(relay availability policy, not transport posture — the harness fails on its "
|
||||
"own timeout instead of running blind)"
|
||||
)
|
||||
|
||||
|
||||
def zenoh_session_config(
|
||||
zenoh,
|
||||
*,
|
||||
listen: list[str] | None = None,
|
||||
connect: list[str] | None = None,
|
||||
mode: str = "peer",
|
||||
posture: str = DEFAULT_ZENOH_POSTURE,
|
||||
):
|
||||
"""Explicit connect, multicast scouting off. `mode=client` targets a zenohd router."""
|
||||
"""A session under the relay's shipped posture, with explicit endpoints.
|
||||
|
||||
`mode=client` targets a running zenohd and is the deployed shape
|
||||
(`ZenohMode::default()`); `mode=peer` is a direct pod-to-pod link and is a
|
||||
transport floor, not a routing path the relay uses.
|
||||
"""
|
||||
conf = zenoh.Config()
|
||||
conf.insert_json5("mode", json.dumps(mode))
|
||||
conf.insert_json5("scouting/multicast/enabled", json.dumps(False))
|
||||
apply_zenoh_posture(conf, zenoh_posture(posture))
|
||||
if listen:
|
||||
conf.insert_json5("listen/endpoints", json.dumps(listen))
|
||||
if connect:
|
||||
|
|
@ -1018,12 +1209,13 @@ def open_zenoh_pods(
|
|||
relevant: set[str],
|
||||
timeout: float,
|
||||
mode: str = "peer",
|
||||
posture: str = DEFAULT_ZENOH_POSTURE,
|
||||
) -> list[ZenohPod]:
|
||||
opened: list[ZenohPod] = []
|
||||
try:
|
||||
for pod_index in range(pods):
|
||||
pod = ZenohPod(pod=pod_index, relevant_communities=relevant)
|
||||
conf = zenoh_peer_config(zenoh, connect=connect, mode=mode)
|
||||
conf = zenoh_session_config(zenoh, connect=connect, mode=mode, posture=posture)
|
||||
session = zenoh.open(conf)
|
||||
pod.session = session
|
||||
for key_expr in key_exprs:
|
||||
|
|
@ -1142,16 +1334,39 @@ def zenoh_endpoints(args: argparse.Namespace) -> list[str]:
|
|||
return [endpoint.strip() for endpoint in raw.split(",") if endpoint.strip()]
|
||||
|
||||
|
||||
def zenoh_pod_mode(args: argparse.Namespace, external: bool) -> str:
|
||||
"""The session mode every session in this run opens with.
|
||||
|
||||
With `--zenoh-connect` this is `--zenoh-session-mode`, whose default is
|
||||
`client` because that is `ZenohMode::default()` in the relay. Without it
|
||||
there is no router to be a client of, so both ends are peers on a direct
|
||||
link -- a transport floor, and explicitly not the deployed routing path.
|
||||
"""
|
||||
return args.zenoh_session_mode if external else "peer"
|
||||
|
||||
|
||||
def zenoh_topology_label(args: argparse.Namespace) -> str:
|
||||
endpoints = zenoh_endpoints(args)
|
||||
if not endpoints:
|
||||
return "peer-direct, in-process listener (transport floor, NOT the deployed routing path)"
|
||||
shape = "the deployed shape" if args.zenoh_session_mode == "client" else "NOT a deployed shape"
|
||||
return f"{args.zenoh_session_mode} mode -> {','.join(endpoints)} ({shape})"
|
||||
|
||||
|
||||
def open_zenoh_publisher(zenoh, args: argparse.Namespace) -> tuple[object, list[str], bool]:
|
||||
"""Return (publisher session, endpoints the pods connect to, external_router)."""
|
||||
external = zenoh_endpoints(args)
|
||||
posture = args.zenoh_posture
|
||||
if external:
|
||||
session = zenoh.open(
|
||||
zenoh_peer_config(zenoh, connect=external, mode=args.zenoh_session_mode)
|
||||
zenoh_session_config(
|
||||
zenoh, connect=external, mode=args.zenoh_session_mode, posture=posture
|
||||
)
|
||||
)
|
||||
return session, external, True
|
||||
listen = [f"tcp/127.0.0.1:{pick_tcp_port()}"]
|
||||
return zenoh.open(zenoh_peer_config(zenoh, listen=listen)), listen, False
|
||||
conf = zenoh_session_config(zenoh, listen=listen, mode="peer", posture=posture)
|
||||
return zenoh.open(conf), listen, False
|
||||
|
||||
|
||||
def measure_zenoh_for_pods(
|
||||
|
|
@ -1181,7 +1396,8 @@ def measure_zenoh_for_pods(
|
|||
key_exprs=[zenoh_firehose_key()],
|
||||
relevant=relevant,
|
||||
timeout=timeout,
|
||||
mode=args.zenoh_session_mode if external else "peer",
|
||||
mode=zenoh_pod_mode(args, external),
|
||||
posture=args.zenoh_posture,
|
||||
)
|
||||
old_expected = pods * args.communities * events_per_community
|
||||
old_stats = _zenoh_phase(
|
||||
|
|
@ -1210,7 +1426,8 @@ def measure_zenoh_for_pods(
|
|||
key_exprs=[zenoh_community_key(community) for community in communities[: args.subscribed_communities]],
|
||||
relevant=relevant,
|
||||
timeout=timeout,
|
||||
mode=args.zenoh_session_mode if external else "peer",
|
||||
mode=zenoh_pod_mode(args, external),
|
||||
posture=args.zenoh_posture,
|
||||
)
|
||||
new_expected = interested * args.subscribed_communities * events_per_community
|
||||
new_stats = _zenoh_phase(
|
||||
|
|
@ -1261,7 +1478,7 @@ def measure_zenoh_kind_scope(
|
|||
timeout = args.zenoh_timeout
|
||||
put_kwargs = zenoh_put_kwargs(zenoh, args)
|
||||
pub_session, connect, external = open_zenoh_publisher(zenoh, args)
|
||||
pod_mode = args.zenoh_session_mode if external else "peer"
|
||||
pod_mode = zenoh_pod_mode(args, external)
|
||||
peer_count = 0 if external else 1
|
||||
community_pods: list[ZenohPod] = []
|
||||
kind_pods: list[ZenohPod] = []
|
||||
|
|
@ -1275,6 +1492,7 @@ def measure_zenoh_kind_scope(
|
|||
relevant=relevant,
|
||||
timeout=timeout,
|
||||
mode=pod_mode,
|
||||
posture=args.zenoh_posture,
|
||||
)
|
||||
wait_for_zenoh_peers(pub_session, peer_count, timeout)
|
||||
time.sleep(args.zenoh_settle)
|
||||
|
|
@ -1437,17 +1655,15 @@ def print_rows(args: argparse.Namespace, rows: list[Measurement], payload_bytes:
|
|||
print(f"redis: {args.redis_url}, payload={payload_bytes} B")
|
||||
print(f"publisher: {redis_publish_label(args)}")
|
||||
if args.mode == "zenoh":
|
||||
endpoints = zenoh_endpoints(args)
|
||||
topology = (
|
||||
f"{args.zenoh_session_mode} mode -> {','.join(endpoints)}"
|
||||
if endpoints
|
||||
else "peer mode, in-process listener"
|
||||
)
|
||||
print(
|
||||
f"zenoh: {topology}, multicast scouting off, "
|
||||
f"zenoh: {zenoh_topology_label(args)}, "
|
||||
f"payload={payload_bytes} B, "
|
||||
f"publish key={zenoh_global_kind_key('{community}', args.kind)}"
|
||||
)
|
||||
# The two-ceilings rule: a number whose profile is not beside it is not
|
||||
# readable. This line travels with every table so a row lifted out of a
|
||||
# log still says which transport produced it.
|
||||
print(f"zenoh {zenoh_posture_stamp(args.zenoh_posture)}")
|
||||
print(
|
||||
"operating point: "
|
||||
f"{'express (batching bypassed)' if args.zenoh_express else 'batched (default)'}, "
|
||||
|
|
@ -1620,9 +1836,10 @@ def resolve_payload_sweep(args: argparse.Namespace) -> list[int]:
|
|||
return sizes
|
||||
|
||||
|
||||
def print_sweep_summary(mode: str, rows: list[SweepRow]) -> None:
|
||||
def print_sweep_summary(mode: str, rows: list[SweepRow], posture: str | None = None) -> None:
|
||||
print()
|
||||
print(f"Payload sweep summary — mode={mode} (firehose arm, elapsed measurements)")
|
||||
stamp = f" — {zenoh_posture_stamp(posture)}" if mode == "zenoh" and posture else ""
|
||||
print(f"Payload sweep summary — mode={mode}{stamp} (firehose arm, elapsed measurements)")
|
||||
print(
|
||||
"| payload B | delivered msg/s | delivered MiB/s | scoped msg/s | p50 µs | "
|
||||
"p95 µs | p99 µs | drops | cold 1st pub µs | gain vs redis | gate |"
|
||||
|
|
@ -1642,6 +1859,11 @@ def print_sweep_summary(mode: str, rows: list[SweepRow]) -> None:
|
|||
"- State the crossover, not one size. A single-payload gate can reject a "
|
||||
"transport that is correct for the traffic class it was chosen for."
|
||||
)
|
||||
if mode == "zenoh" and posture and posture != DEFAULT_ZENOH_POSTURE:
|
||||
print(
|
||||
f"- ARTIFACT ARM: posture `{posture}` is not what the relay ships. These rows "
|
||||
"exist to price the posture delta and license nothing on their own."
|
||||
)
|
||||
|
||||
|
||||
def run_one_payload(
|
||||
|
|
@ -1701,6 +1923,32 @@ def run_one_payload(
|
|||
)
|
||||
|
||||
|
||||
def validate_zenoh_topology(args: argparse.Namespace) -> None:
|
||||
"""Refuse the one arm that reports a healthy run and measures nothing.
|
||||
|
||||
`peer` sessions do not deliver through a `zenohd` router while gossip is
|
||||
off: the router will not propagate a subscription declaration between two
|
||||
Peer faces unless `failover_brokering` answers true, and that needs the
|
||||
gossip link-state net this posture deliberately does not build
|
||||
(`zenoh-1.8.0/src/net/routing/hat/router/pubsub.rs:200-222`,
|
||||
`hat/router/mod.rs:287-303,373`). That is meridian-2m45. The harness would
|
||||
otherwise sit out its full `--zenoh-timeout` and report a timeout with no
|
||||
hint that the topology, not the machine, was wrong.
|
||||
"""
|
||||
if not zenoh_endpoints(args) or args.zenoh_session_mode != "peer":
|
||||
return
|
||||
if getattr(args, "zenoh_allow_peer_through_router", False):
|
||||
return
|
||||
raise ValueError(
|
||||
"--zenoh-session-mode peer with --zenoh-connect is the meridian-2m45 trap: a "
|
||||
"zenohd router does not propagate subscription declarations between two Peer "
|
||||
"faces while gossip scouting is off, so both sessions open, both report "
|
||||
"healthy, and nothing is delivered. The relay ships `client` "
|
||||
"(ZenohMode::default()); use that. Pass "
|
||||
"--zenoh-allow-peer-through-router only to reproduce the trap on purpose."
|
||||
)
|
||||
|
||||
|
||||
def run(args: argparse.Namespace) -> int:
|
||||
if args.communities <= 0:
|
||||
raise ValueError("--communities must be positive")
|
||||
|
|
@ -1714,11 +1962,16 @@ def run(args: argparse.Namespace) -> int:
|
|||
sweep = resolve_payload_sweep(args)
|
||||
|
||||
if args.mode == "zenoh":
|
||||
# Both of these are cheap and both are fatal, so they run before the
|
||||
# first byte moves rather than after a full sweep has been spent.
|
||||
validate_zenoh_topology(args)
|
||||
zenoh_posture(args.zenoh_posture)
|
||||
try:
|
||||
require_zenoh()
|
||||
except RuntimeError as exc:
|
||||
print(str(exc), file=sys.stderr)
|
||||
return 2
|
||||
print_zenoh_posture_block(args.zenoh_posture)
|
||||
|
||||
summary: list[SweepRow] = []
|
||||
for index, payload_bytes in enumerate(sweep):
|
||||
|
|
@ -1729,7 +1982,7 @@ def run(args: argparse.Namespace) -> int:
|
|||
if row is not None:
|
||||
summary.append(row)
|
||||
if len(summary) > 1:
|
||||
print_sweep_summary(args.mode, summary)
|
||||
print_sweep_summary(args.mode, summary, getattr(args, "zenoh_posture", None))
|
||||
return 0
|
||||
|
||||
|
||||
|
|
@ -1832,8 +2085,36 @@ def build_parser() -> argparse.ArgumentParser:
|
|||
parser.add_argument(
|
||||
"--zenoh-session-mode",
|
||||
choices=["peer", "client"],
|
||||
default="peer",
|
||||
help="session mode used with --zenoh-connect (default peer)",
|
||||
default="client",
|
||||
help=(
|
||||
"session mode used with --zenoh-connect. Default `client`, because that is "
|
||||
"ZenohMode::default() in the relay and the only mode that delivers through a "
|
||||
"router while gossip is off (meridian-2m45). Ignored without --zenoh-connect: "
|
||||
"with no router there is nothing to be a client of, so both ends are peers "
|
||||
"on a direct link"
|
||||
),
|
||||
)
|
||||
parser.add_argument(
|
||||
"--zenoh-allow-peer-through-router",
|
||||
action="store_true",
|
||||
default=False,
|
||||
help=(
|
||||
"permit --zenoh-session-mode peer together with --zenoh-connect. Refused by "
|
||||
"default: it is the meridian-2m45 trap, where both sessions report healthy "
|
||||
"and the router delivers nothing"
|
||||
),
|
||||
)
|
||||
parser.add_argument(
|
||||
"--zenoh-posture",
|
||||
choices=sorted(ZENOH_POSTURES),
|
||||
default=DEFAULT_ZENOH_POSTURE,
|
||||
help=(
|
||||
"which Zenoh configuration every session opens under. `relay` (default) is "
|
||||
f"the posture the relay binary embeds ({ZENOH_POSTURE_SOURCE}). `legacy-0a2` "
|
||||
"reproduces what this harness applied before 2026-08-21 — multicast off and "
|
||||
"nothing else, so gossip scouting and shared memory were BOTH left at their "
|
||||
"upstream `true` — and is kept only so the delta stays measurable"
|
||||
),
|
||||
)
|
||||
parser.add_argument(
|
||||
"--kind",
|
||||
|
|
|
|||
|
|
@ -48,7 +48,19 @@ run() {
|
|||
for rep in $(seq 1 "$REPS"); do
|
||||
echo "=== repetition $rep/$REPS ==="
|
||||
|
||||
# D. The SHIPPED shape, and the only arm on this list that is one: `client`
|
||||
# sessions against a containerised zenohd is `ZenohMode::default()` plus the
|
||||
# deployed router config, and Zenoh then pays the same Docker VM boundary
|
||||
# Dragonfly does. Any figure that would license a lane migration has to come
|
||||
# from here. It runs first for that reason.
|
||||
run "r${rep}-D-zenoh-router-client-vs-dragonfly" \
|
||||
--mode zenoh --payload-sweep --pods 1 --no-measure-kind-scope \
|
||||
--zenoh-connect "$ZENOHD_ENDPOINT" --zenoh-session-mode client \
|
||||
--redis-url "$DRAGONFLY" --redis-timeout "$TIMEOUT"
|
||||
|
||||
# A. Zenoh peer, batched, against Dragonfly-in-Docker (the deployed backing store).
|
||||
# A DIRECT pod-to-pod link: a transport floor, and not a routing path the
|
||||
# relay uses -- the relay reaches every other pod through the router tier.
|
||||
run "r${rep}-A-zenoh-peer-batched-vs-dragonfly" \
|
||||
--mode zenoh --payload-sweep --pods 1,2,4 --no-measure-kind-scope \
|
||||
--redis-url "$DRAGONFLY" --redis-timeout "$TIMEOUT"
|
||||
|
|
@ -64,13 +76,6 @@ for rep in $(seq 1 "$REPS"); do
|
|||
--mode zenoh --payload-sweep --pods 1 --no-measure-kind-scope --zenoh-express \
|
||||
--redis-url "$DRAGONFLY" --redis-timeout "$TIMEOUT"
|
||||
|
||||
# D. Router arm: both peers are clients of a containerised zenohd, so Zenoh pays
|
||||
# the same Docker VM boundary Dragonfly does.
|
||||
run "r${rep}-D-zenoh-router-client-vs-dragonfly" \
|
||||
--mode zenoh --payload-sweep --pods 1 --no-measure-kind-scope \
|
||||
--zenoh-connect "$ZENOHD_ENDPOINT" --zenoh-session-mode client \
|
||||
--redis-url "$DRAGONFLY" --redis-timeout "$TIMEOUT"
|
||||
|
||||
# E. Cold start: the first publication after session init stays inside the window.
|
||||
run "r${rep}-E-zenoh-peer-coldstart" \
|
||||
--mode zenoh --payload-sweep 64,256,4096 --pods 1 --no-measure-kind-scope \
|
||||
|
|
@ -89,6 +94,24 @@ for rep in $(seq 1 "$REPS"); do
|
|||
run "r${rep}-H-zenoh-vs-dragonfly-SYNC-publisher-artifact" \
|
||||
--mode zenoh --payload-sweep 64,256,4096 --pods 1 --no-measure-kind-scope \
|
||||
--redis-publish sync --redis-url "$DRAGONFLY" --redis-timeout "$TIMEOUT"
|
||||
|
||||
# J, K. Posture artifacts. `legacy-0a2` reproduces what this harness applied
|
||||
# for the 2026-08-19 run: multicast scouting off and nothing else, so gossip
|
||||
# scouting AND shared memory were both left at their upstream `true`. They
|
||||
# exist for the same reason arm H does -- a discarded number stays
|
||||
# reproducible rather than merely asserted -- and they turn "the posture
|
||||
# plausibly does not move this" into a measured delta against A and D.
|
||||
# They license nothing, and say so in their own output.
|
||||
run "r${rep}-J-zenoh-peer-LEGACY-POSTURE-artifact" \
|
||||
--mode zenoh --payload-sweep --pods 1 --no-measure-kind-scope \
|
||||
--zenoh-posture legacy-0a2 \
|
||||
--redis-url "$DRAGONFLY" --redis-timeout "$TIMEOUT"
|
||||
|
||||
run "r${rep}-K-zenoh-router-LEGACY-POSTURE-artifact" \
|
||||
--mode zenoh --payload-sweep --pods 1 --no-measure-kind-scope \
|
||||
--zenoh-posture legacy-0a2 \
|
||||
--zenoh-connect "$ZENOHD_ENDPOINT" --zenoh-session-mode client \
|
||||
--redis-url "$DRAGONFLY" --redis-timeout "$TIMEOUT"
|
||||
done
|
||||
|
||||
echo "raw logs: $OUT"
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ from __future__ import annotations
|
|||
import argparse
|
||||
import contextlib
|
||||
import io
|
||||
import json
|
||||
import types
|
||||
import unittest
|
||||
from unittest.mock import patch
|
||||
|
|
@ -34,7 +35,9 @@ class RelayBusScalingTests(unittest.TestCase):
|
|||
zenoh_prewarm=harness.DEFAULT_ZENOH_PREWARM,
|
||||
zenoh_settle=harness.DEFAULT_ZENOH_SETTLE_S,
|
||||
zenoh_connect=None,
|
||||
zenoh_session_mode="peer",
|
||||
zenoh_session_mode="client",
|
||||
zenoh_allow_peer_through_router=False,
|
||||
zenoh_posture=harness.DEFAULT_ZENOH_POSTURE,
|
||||
kind=harness.DEFAULT_KIND,
|
||||
gain_threshold=harness.DEFAULT_GAIN_THRESHOLD,
|
||||
measure_kind_scope=False,
|
||||
|
|
@ -522,5 +525,205 @@ class WaitForCountsTests(unittest.TestCase):
|
|||
harness.wait_for_counts(subs, 100, timeout=0.05)
|
||||
|
||||
|
||||
class FakeZenohConfig:
|
||||
"""A `zenoh.Config` that records inserts and answers `get_json` from them.
|
||||
|
||||
Deliberately *not* a mock that always agrees: the drop/coerce behaviours the
|
||||
real binding can exhibit are what `apply_zenoh_posture` exists to catch, so
|
||||
the fake has to be able to exhibit them too.
|
||||
"""
|
||||
|
||||
def __init__(self, *, drop: set[str] | None = None, coerce: dict[str, object] | None = None):
|
||||
self.values: dict[str, object] = {}
|
||||
self.order: list[str] = []
|
||||
self._drop = drop or set()
|
||||
self._coerce = coerce or {}
|
||||
|
||||
def insert_json5(self, key: str, value: str) -> None:
|
||||
self.order.append(key)
|
||||
if key in self._drop:
|
||||
return
|
||||
self.values[key] = self._coerce.get(key, json.loads(value))
|
||||
|
||||
def get_json(self, key: str) -> str:
|
||||
if key not in self.values:
|
||||
raise KeyError(key)
|
||||
return json.dumps(self.values[key])
|
||||
|
||||
|
||||
class FakeZenoh:
|
||||
def __init__(self, **kwargs: object) -> None:
|
||||
self._kwargs = kwargs
|
||||
self.configs: list[FakeZenohConfig] = []
|
||||
|
||||
def Config(self) -> FakeZenohConfig: # noqa: N802 - mirrors the binding's name
|
||||
conf = FakeZenohConfig(**self._kwargs) # type: ignore[arg-type]
|
||||
self.configs.append(conf)
|
||||
return conf
|
||||
|
||||
|
||||
class ZenohPostureTests(unittest.TestCase):
|
||||
"""The posture is the profile. These pin what the harness claims to have run."""
|
||||
|
||||
def test_relay_posture_closes_both_scouting_mechanisms_present_and_false(self) -> None:
|
||||
applied = dict(harness.RELAY_ZENOH_POSTURE)
|
||||
# Present-and-false, never absent: an absent key is read as consent and
|
||||
# the upstream default (true) applies -- zenohd force-enables multicast.
|
||||
self.assertIs(applied["scouting/multicast/enabled"], False)
|
||||
self.assertIs(applied["scouting/gossip/enabled"], False)
|
||||
|
||||
def test_relay_posture_closes_both_shared_memory_switches(self) -> None:
|
||||
applied = dict(harness.RELAY_ZENOH_POSTURE)
|
||||
# Both default true in zenoh 1.8, and the eclipse-zenoh wheel is built
|
||||
# with the shared-memory feature. Unset means every payload >= 3,072 B
|
||||
# travels through POSIX SHM -- which the relay build cannot do at all.
|
||||
self.assertIs(applied["transport/shared_memory/enabled"], False)
|
||||
self.assertIs(applied["transport/shared_memory/transport_optimization/enabled"], False)
|
||||
|
||||
def test_relay_posture_matches_the_fixture_the_relay_binary_embeds(self) -> None:
|
||||
"""Drift between this table and `bus-peer.json5` is the whole defect."""
|
||||
import pathlib
|
||||
import re
|
||||
|
||||
fixture = (
|
||||
pathlib.Path(__file__).resolve().parents[1] / harness.ZENOH_POSTURE_SOURCE
|
||||
)
|
||||
if not fixture.exists(): # pragma: no cover - the fixture is tracked
|
||||
self.skipTest(f"{harness.ZENOH_POSTURE_SOURCE} not present")
|
||||
body = re.sub(r"(?m)//.*$", "", fixture.read_text(encoding="utf-8"))
|
||||
body = re.sub(r"([{,]\s*)([A-Za-z_][A-Za-z0-9_]*)(\s*:)", r'\1"\2"\3', body)
|
||||
body = re.sub(r"(?m)^(\s*)([A-Za-z_][A-Za-z0-9_]*)(\s*):", r'\1"\2"\3:', body)
|
||||
body = re.sub(r",(\s*[}\]])", r"\1", body)
|
||||
fixture_conf = json.loads(body)
|
||||
|
||||
def dig(path: str) -> object:
|
||||
node: object = fixture_conf
|
||||
for part in path.split("/"):
|
||||
self.assertIsInstance(node, dict, f"{path} is not reachable in the fixture")
|
||||
self.assertIn(part, node, f"`{path}` is absent from {harness.ZENOH_POSTURE_SOURCE}")
|
||||
node = node[part] # type: ignore[index]
|
||||
return node
|
||||
|
||||
for key, value in harness.RELAY_ZENOH_POSTURE:
|
||||
self.assertEqual(
|
||||
dig(key),
|
||||
value,
|
||||
f"`{key}` disagrees with {harness.ZENOH_POSTURE_SOURCE}; the harness would "
|
||||
"measure a posture the relay does not ship",
|
||||
)
|
||||
|
||||
def test_legacy_posture_is_only_what_the_0a2_run_actually_applied(self) -> None:
|
||||
self.assertEqual(
|
||||
harness.LEGACY_0A2_ZENOH_POSTURE,
|
||||
(("scouting/multicast/enabled", False),),
|
||||
)
|
||||
self.assertNotEqual(harness.DEFAULT_ZENOH_POSTURE, "legacy-0a2")
|
||||
|
||||
def test_session_config_applies_every_posture_key(self) -> None:
|
||||
zenoh = FakeZenoh()
|
||||
conf = harness.zenoh_session_config(zenoh, connect=["tcp/127.0.0.1:7447"], mode="client")
|
||||
self.assertEqual(json.loads(conf.get_json("mode")), "client")
|
||||
for key, value in harness.RELAY_ZENOH_POSTURE:
|
||||
self.assertEqual(json.loads(conf.get_json(key)), value)
|
||||
self.assertEqual(
|
||||
json.loads(conf.get_json("connect/endpoints")), ["tcp/127.0.0.1:7447"]
|
||||
)
|
||||
|
||||
def test_a_silently_dropped_posture_key_fails_the_run(self) -> None:
|
||||
"""zenoh 1.9/1.10 accept `routing/peer/mode` and drop it. Same shape."""
|
||||
zenoh = FakeZenoh(drop={"scouting/gossip/enabled"})
|
||||
with self.assertRaises(harness.ZenohPostureError) as caught:
|
||||
harness.zenoh_session_config(zenoh, listen=["tcp/127.0.0.1:0"])
|
||||
self.assertIn("scouting/gossip/enabled", str(caught.exception))
|
||||
|
||||
def test_a_coerced_posture_key_fails_the_run(self) -> None:
|
||||
zenoh = FakeZenoh(coerce={"transport/shared_memory/enabled": True})
|
||||
with self.assertRaises(harness.ZenohPostureError) as caught:
|
||||
harness.zenoh_session_config(zenoh, listen=["tcp/127.0.0.1:0"])
|
||||
self.assertIn("transport/shared_memory/enabled", str(caught.exception))
|
||||
|
||||
def test_posture_digest_changes_when_the_posture_changes(self) -> None:
|
||||
relay = harness.zenoh_posture_digest(harness.RELAY_ZENOH_POSTURE)
|
||||
legacy = harness.zenoh_posture_digest(harness.LEGACY_0A2_ZENOH_POSTURE)
|
||||
self.assertNotEqual(relay, legacy)
|
||||
self.assertEqual(len(relay), 12)
|
||||
|
||||
def test_every_zenoh_table_is_stamped_with_its_posture(self) -> None:
|
||||
args = RelayBusScalingTests().args(mode="zenoh", zenoh_connect="tcp/127.0.0.1:7447")
|
||||
rows = harness.model_measurements(RelayBusScalingTests().args(), [1])
|
||||
out = io.StringIO()
|
||||
with contextlib.redirect_stdout(out):
|
||||
harness.print_rows(args, rows, 256)
|
||||
text = out.getvalue()
|
||||
self.assertIn("posture=relay#", text)
|
||||
self.assertIn("multicast=off", text)
|
||||
self.assertIn("gossip=off", text)
|
||||
self.assertIn("shared_memory=off", text)
|
||||
self.assertIn("transport_optimization=off", text)
|
||||
|
||||
def test_posture_block_names_its_source_and_lists_every_key(self) -> None:
|
||||
out = io.StringIO()
|
||||
with contextlib.redirect_stdout(out):
|
||||
harness.print_zenoh_posture_block("relay")
|
||||
text = out.getvalue()
|
||||
self.assertIn(harness.ZENOH_POSTURE_SOURCE, text)
|
||||
for key, _value in harness.RELAY_ZENOH_POSTURE:
|
||||
self.assertIn(key, text)
|
||||
|
||||
def test_posture_block_marks_the_artifact_arm_as_licensing_nothing(self) -> None:
|
||||
out = io.StringIO()
|
||||
with contextlib.redirect_stdout(out):
|
||||
harness.print_zenoh_posture_block("legacy-0a2")
|
||||
self.assertIn("NOT the shipped posture", out.getvalue())
|
||||
|
||||
def test_unknown_posture_names_the_known_ones(self) -> None:
|
||||
with self.assertRaises(ValueError) as caught:
|
||||
harness.zenoh_posture("whatever")
|
||||
self.assertIn("relay", str(caught.exception))
|
||||
self.assertIn("legacy-0a2", str(caught.exception))
|
||||
|
||||
|
||||
class ZenohTopologyTests(unittest.TestCase):
|
||||
"""`client` is what the relay ships; `peer` through a router delivers nothing."""
|
||||
|
||||
def args(self, **overrides: object) -> argparse.Namespace:
|
||||
return RelayBusScalingTests().args(mode="zenoh", **overrides)
|
||||
|
||||
def test_default_session_mode_is_the_relays_default(self) -> None:
|
||||
parsed = harness.build_parser().parse_args([])
|
||||
self.assertEqual(parsed.zenoh_session_mode, "client")
|
||||
self.assertEqual(parsed.zenoh_posture, "relay")
|
||||
|
||||
def test_peer_through_a_router_is_refused_by_name(self) -> None:
|
||||
args = self.args(zenoh_connect="tcp/127.0.0.1:7447", zenoh_session_mode="peer")
|
||||
with self.assertRaises(ValueError) as caught:
|
||||
harness.validate_zenoh_topology(args)
|
||||
message = str(caught.exception)
|
||||
self.assertIn("meridian-2m45", message)
|
||||
self.assertIn("client", message)
|
||||
|
||||
def test_peer_through_a_router_is_reproducible_on_purpose(self) -> None:
|
||||
args = self.args(
|
||||
zenoh_connect="tcp/127.0.0.1:7447",
|
||||
zenoh_session_mode="peer",
|
||||
zenoh_allow_peer_through_router=True,
|
||||
)
|
||||
harness.validate_zenoh_topology(args)
|
||||
|
||||
def test_peer_direct_ignores_the_session_mode_flag(self) -> None:
|
||||
args = self.args(zenoh_session_mode="client")
|
||||
harness.validate_zenoh_topology(args)
|
||||
self.assertEqual(harness.zenoh_pod_mode(args, external=False), "peer")
|
||||
self.assertEqual(harness.zenoh_pod_mode(args, external=True), "client")
|
||||
|
||||
def test_topology_label_says_which_shape_is_deployed(self) -> None:
|
||||
floor = harness.zenoh_topology_label(self.args())
|
||||
self.assertIn("NOT the deployed routing path", floor)
|
||||
shipped = harness.zenoh_topology_label(
|
||||
self.args(zenoh_connect="tcp/127.0.0.1:7447", zenoh_session_mode="client")
|
||||
)
|
||||
self.assertIn("the deployed shape", shipped)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
|
|
|
|||
|
|
@ -13,6 +13,13 @@ of the same session -- so the number describes a real route, not an unmatched on
|
|||
Each trial opens a fresh session pair, so the spike is measured repeatedly rather
|
||||
than once.
|
||||
|
||||
Both sessions open under the relay's shipped posture, imported from
|
||||
``relay_bus_scaling`` so ``perf/`` holds exactly one copy of it. That matters here
|
||||
even though the headline figure is a Zenoh-to-Zenoh ratio: with the shared-memory
|
||||
keys left unset -- which is what this probe did before 2026-08-21 -- a 4,096 B
|
||||
payload exceeds Zenoh's 3,072 B ``message_size_threshold``, so the first put also
|
||||
pays the 16 MiB SHM pool allocation, on a path the relay build cannot use at all.
|
||||
|
||||
./perf/zenoh_coldstart_probe.py [trials] [payload_bytes]
|
||||
"""
|
||||
|
||||
|
|
@ -23,13 +30,21 @@ import statistics
|
|||
import sys
|
||||
import time
|
||||
|
||||
from relay_bus_scaling import (
|
||||
DEFAULT_ZENOH_POSTURE,
|
||||
apply_zenoh_posture,
|
||||
zenoh_posture,
|
||||
zenoh_posture_digest,
|
||||
)
|
||||
|
||||
MEASURED_KEY = "meridian/v1/coldstart/global/k/9"
|
||||
STEADY_PUTS = 200
|
||||
|
||||
|
||||
def open_session(zenoh, listen: str | None, connect: str | None):
|
||||
config = zenoh.Config()
|
||||
config.insert_json5("scouting/multicast/enabled", "false")
|
||||
config.insert_json5("mode", json.dumps("peer"))
|
||||
apply_zenoh_posture(config, zenoh_posture(DEFAULT_ZENOH_POSTURE))
|
||||
if listen is not None:
|
||||
config.insert_json5("listen/endpoints", json.dumps([listen]))
|
||||
if connect is not None:
|
||||
|
|
@ -88,7 +103,11 @@ def main(argv: list[str]) -> int:
|
|||
|
||||
ratios = [f / s for f, s in zip(firsts, steadies)]
|
||||
print()
|
||||
print(f"payload {payload_bytes} B, {trials} fresh session pairs, {STEADY_PUTS} steady puts each")
|
||||
print(
|
||||
f"payload {payload_bytes} B, {trials} fresh session pairs, {STEADY_PUTS} steady puts each, "
|
||||
f"posture={DEFAULT_ZENOH_POSTURE}#"
|
||||
f"{zenoh_posture_digest(zenoh_posture(DEFAULT_ZENOH_POSTURE))}"
|
||||
)
|
||||
print(
|
||||
f"first publication : median {statistics.median(firsts):,.1f} us "
|
||||
f"[{min(firsts):,.1f}-{max(firsts):,.1f}]"
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue