fix(bus): the relay is a client of the router tier, not a peer of it
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
meridian-2m45. Two relays in the documented default `peer` mode, both
connected to a `zenohd` router, exchanged NOTHING while both reported
`session_open`, `endpoints_connected 1` and `zenoh_ready 1`. Not a drop
either: no drop counter on the receiver was ever created, because the
sample never left the publisher.
The mechanism, in the pinned tree rather than in a blog post:
zenoh-1.8.0/src/net/routing/hat/router/pubsub.rs:200-222 — the router
refuses to propagate a subscription declaration from one Peer face to
another Peer face unless `failover_brokering(src, dst)`. A Client face
is exempt; the `src_face.whatami != WhatAmI::Peer` arm short-circuits.
hat/router/mod.rs:287-303 — `failover_brokering` needs
`linkstatepeers_net`, and hat/router/mod.rs:373 builds it only
`if peer_full_linkstate | gossip`. With `routing.peer.mode` at its
`peer_to_peer` default and `scouting.gossip.enabled: false` it is
`None`, so the answer is always false and the publisher never learns a
remote subscriber exists.
Upstream states it plainly at DEFAULT_CONFIG.json5:225-226 — "The
failover brokering only works if gossip discovery is enabled". So the
charter's gossip refusal and session mode `peer` are mutually
incompatible through a router. Confirmed both ways on a live router:
gossip on (autoconnect still empty) makes `peer` deliver; gossip off with
`routing.peer.mode: "linkstate"` on BOTH ends also makes it deliver.
`linkstate` was measured working and rejected anyway, unanimously
(Architecture & Scale council, parallel seats, minutes in the bead):
- it converts a per-pod value into a fleet-wide invariant that fails
SILENTLY in both directions when the halves disagree — measured, both
directions — and every rolling restart passes through that state;
- eclipse/zenoh 1.9.0 and 1.10.0 ACCEPT `routing.peer.mode` and
silently drop it (absent from the daemon's own `Initial conf`, router
starts healthy), so the fix evaporates on a minor upgrade and this P0
returns wearing a different hat;
- `endpoints` has no cardinality bound, so a second router endpoint
added for availability — posture-compliant, no refusal fires — would
make every relay a multihop forwarder between communities.
`client` needs no cross-node agreement: verified delivering against a
`peer_to_peer` router AND a `linkstate` one. `hat/client/` holds no
routing Network at all, so transit is impossible by type rather than by
topology. It is also the documented target shape — one pod, one regional
router (feature-zenoh-transport.md § Topology) — and the only shape this
repo has ever actually measured through a router
(perf/run_0a2_matrix.sh:71 already passes `--zenoh-session-mode client`).
`peer` stays a supported mode for a DIRECT pod-to-pod link, which is what
`zenoh_bus.rs::peer_pair` tests and what the benchmark arms use. Nothing
in that suite changes; all 21 still pass, because all 21 cross-connect
two peers directly and none of them ever touched a router. That is why
they were green through the whole defect.
The reversal condition, and the whole trade, sit on `ZenohMode` where the
next person to change it will be looking, not only in the tracker.
New gate `just zenoh-router-check` is the live half: it starts a real
`zenohd` from `deploy/compose/zenoh/zenohd.json5` with the deployed
command line, routes two `ZenohEventBus` sessions through it, and pins
BOTH the mode that delivers and the mode that silently does not. Proven
to fail against the pre-fix default. Opt-in behind
MERIDIAN_ZENOH_DAEMON_CHECK=1 and out of `just check`, same as
`zenoh-config-check`, because it needs Docker and the pinned image.
Signed-off-by: Joshua Belke <joshua@innovationhub-act.org>
This commit is contained in:
parent
128e5014a0
commit
b81fefb5da
4 changed files with 461 additions and 10 deletions
14
.env.example
14
.env.example
|
|
@ -143,9 +143,17 @@ REDIS_URL=redis://localhost:6379
|
|||
# them on Redis until a mutually-authenticated link profile exists.
|
||||
# MERIDIAN_ZENOH_ENDPOINTS=tcp/127.0.0.1:7447
|
||||
|
||||
# Session mode: `peer` (default) or `client`. `router` is deliberately not a
|
||||
# value — the relay is a peer of the router tier, never a router itself.
|
||||
# MERIDIAN_ZENOH_MODE=peer
|
||||
# Session mode: `client` (default) or `peer`. `router` is deliberately not a
|
||||
# value — the relay attaches to the router tier, never becomes one.
|
||||
#
|
||||
# `peer` is NOT a drop-in alternative here. Two relays in `peer` mode connected
|
||||
# to a `zenohd` router exchange nothing while both report healthy: with gossip
|
||||
# scouting off (this deployment refuses it) the router will not forward one
|
||||
# peer's subscription declaration to another peer. `peer` is for a DIRECT link
|
||||
# between two pods — the benchmark and single-node development shape. See
|
||||
# `ZenohMode` in crates/meridian-pubsub/src/zenoh/config.rs for the mechanism
|
||||
# and what would move this default back.
|
||||
# MERIDIAN_ZENOH_MODE=client
|
||||
|
||||
# Explicit listen endpoint. Defaults to the pinned fixture's ephemeral loopback
|
||||
# port (`tcp/127.0.0.1:0`), because a relay pod is dialled by the router rather
|
||||
|
|
|
|||
16
Justfile
16
Justfile
|
|
@ -313,6 +313,22 @@ perf-check:
|
|||
zenoh-config-check:
|
||||
MERIDIAN_ZENOH_DAEMON_CHECK=1 cargo test -p meridian-pubsub --features zenoh zenoh_config
|
||||
|
||||
# Prove the bus DELIVERS in the topology it is deployed in: two relay sessions,
|
||||
# one real zenohd, an event across. Nothing else does.
|
||||
#
|
||||
# `tests/zenoh_bus.rs` has 21 green tests and every one of them cross-connects
|
||||
# two peers DIRECTLY -- a shape that exercises none of the router's routing
|
||||
# tables. All 21 stayed green through meridian-2m45, where the documented
|
||||
# default mode delivered NOTHING through a router while both sessions reported
|
||||
# `zenoh_ready`. A test suite that never stands up the deployed topology cannot
|
||||
# see a defect that only exists in it, however many assertions it carries.
|
||||
#
|
||||
# Same gate as `zenoh-config-check` above and for the same reason: it needs
|
||||
# Docker and the pinned eclipse/zenoh:1.8.0 image, so it is opt-in behind
|
||||
# MERIDIAN_ZENOH_DAEMON_CHECK=1 and NOT in `just check`.
|
||||
zenoh-router-check:
|
||||
MERIDIAN_ZENOH_DAEMON_CHECK=1 cargo test -p meridian-pubsub --features zenoh --test zenoh_router -- --test-threads=1
|
||||
|
||||
# Compile and run the Zenoh backend. Nothing else does.
|
||||
#
|
||||
# `just clippy` is `--workspace --all-targets` with no `--all-features`, and
|
||||
|
|
|
|||
|
|
@ -56,7 +56,8 @@ pub const PINNED_PEER_CONFIG: &str = include_str!("../../tests/fixtures/zenoh/bu
|
|||
/// Comma-separated Zenoh endpoints this pod connects to. Required.
|
||||
pub const ENV_ENDPOINTS: &str = "MERIDIAN_ZENOH_ENDPOINTS";
|
||||
|
||||
/// Session mode: `peer` (default) or `client`.
|
||||
/// Session mode: `client` (default) or `peer`. See [`ZenohMode`] for why the
|
||||
/// default is `client` and what would move it back.
|
||||
pub const ENV_MODE: &str = "MERIDIAN_ZENOH_MODE";
|
||||
|
||||
/// Optional explicit listen endpoint. Defaults to the fixture's ephemeral port.
|
||||
|
|
@ -72,16 +73,80 @@ pub const ENV_TLS_MATERIAL: [&str; 3] = [
|
|||
|
||||
/// How the relay joins the fabric.
|
||||
///
|
||||
/// `router` is deliberately absent: the relay is a peer of the router tier,
|
||||
/// never a router itself, and a config that could express it would eventually
|
||||
/// be handed one.
|
||||
/// `router` is deliberately absent: the relay attaches to the router tier,
|
||||
/// never becomes one, and a config that could express it would eventually be
|
||||
/// handed one.
|
||||
///
|
||||
/// # The default is `Client`, and that is a recorded call — meridian-2m45
|
||||
///
|
||||
/// **`Peer` does not deliver through a `zenohd` router under this posture.**
|
||||
/// Two relays in `peer` mode, both connected to a router, both open and both
|
||||
/// reporting `zenoh_ready`, exchange **nothing** — and nothing counts it. The
|
||||
/// mechanism, in the pinned tree:
|
||||
///
|
||||
/// - `zenoh-1.8.0/src/net/routing/hat/router/pubsub.rs:200-222` — the router
|
||||
/// refuses to propagate a subscription declaration from one `Peer` face to
|
||||
/// another `Peer` face unless `failover_brokering(src, dst)` answers true. A
|
||||
/// `Client` face is exempt: the `src_face.whatami != WhatAmI::Peer` arm
|
||||
/// short-circuits, so a client's declaration reaches every non-router face.
|
||||
/// - `hat/router/mod.rs:287-303` — `failover_brokering` needs
|
||||
/// `linkstatepeers_net`, and `hat/router/mod.rs:373` builds it only
|
||||
/// `if peer_full_linkstate | gossip`.
|
||||
/// - So with `routing.peer.mode` at its `peer_to_peer` default and
|
||||
/// `scouting.gossip.enabled: false`, that net is `None` and the answer is
|
||||
/// always false. Upstream says so itself in `DEFAULT_CONFIG.json5:225-226`:
|
||||
/// "The failover brokering only works if gossip discovery is enabled and
|
||||
/// peers are configured with gossip target 'router'."
|
||||
///
|
||||
/// Gossip is refused by the charter, so the only way to keep `peer` was
|
||||
/// `routing.peer.mode: "linkstate"` on the relay **and** the router. Measured,
|
||||
/// that fixes delivery — and it was rejected anyway, unanimously, because:
|
||||
///
|
||||
/// - it turns a per-pod value into a fleet-wide invariant that fails **silently
|
||||
/// in both directions** when the halves disagree, which every rolling restart
|
||||
/// passes through (`DEFAULT_CONFIG.json5:243`: "This option needs to be set to
|
||||
/// the same value in all peers and routers of the subsystem");
|
||||
/// - `eclipse/zenoh` 1.9.0 and 1.10.0 **accept `routing.peer.mode` and silently
|
||||
/// drop it** — the key is absent from the daemon's own `Initial conf` and the
|
||||
/// router starts healthy — so the fix evaporates on a minor upgrade and this
|
||||
/// P0 returns wearing a different hat;
|
||||
/// - `endpoints` carries no cardinality bound, so adding a second router
|
||||
/// endpoint for availability — fully posture-compliant, no refusal fires —
|
||||
/// would make every relay a legitimate multihop forwarder between
|
||||
/// communities, which is the shape [`DIAGRAM.md`] Law 2 refuses.
|
||||
///
|
||||
/// `Client` needs no cross-node agreement (verified: it delivers against a
|
||||
/// `peer_to_peer` router *and* a `linkstate` router), and
|
||||
/// `zenoh-1.8.0/src/net/routing/hat/client/` holds no routing `Network` at all,
|
||||
/// so transit is impossible by type rather than by topology.
|
||||
///
|
||||
/// **What would reopen this:** a topology that requires one relay to hold
|
||||
/// concurrent transports to two or more routers. A client keeps exactly one
|
||||
/// link and reconnects across the endpoint list on session close
|
||||
/// (`zenoh-1.8.0/src/net/runtime/orchestrator.rs:1225-1227`, "Currently Client
|
||||
/// can have only one link"); a peer holds all of them at once. Today's target
|
||||
/// is one pod to one regional router
|
||||
/// (`.settings/features/feature-zenoh-transport.md` § Topology), so that
|
||||
/// requirement does not exist — if it appears, this call is back on the table
|
||||
/// and the router tier's own mesh has to be declared with it.
|
||||
///
|
||||
/// `crates/meridian-pubsub/tests/zenoh_router.rs` is the live gate on all of
|
||||
/// the above.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
|
||||
pub enum ZenohMode {
|
||||
/// Peer of the router tier — the production and benchmark shape.
|
||||
/// Client of the router tier — the deployed shape and the default.
|
||||
///
|
||||
/// One link at a time, reconnected across the endpoint list on close, and
|
||||
/// no routing table of its own.
|
||||
#[default]
|
||||
Peer,
|
||||
/// Client of a router. Retained for a bounded single-node profile.
|
||||
Client,
|
||||
/// Peer. Delivers over a **direct** link to another peer, which is the
|
||||
/// benchmark and single-node development shape
|
||||
/// (`tests/zenoh_bus.rs::peer_pair` cross-connects two of them).
|
||||
///
|
||||
/// It does **not** deliver through a `zenohd` router while gossip is off —
|
||||
/// see the type-level note above before selecting it.
|
||||
Peer,
|
||||
}
|
||||
|
||||
impl ZenohMode {
|
||||
|
|
|
|||
362
crates/meridian-pubsub/tests/zenoh_router.rs
Normal file
362
crates/meridian-pubsub/tests/zenoh_router.rs
Normal file
|
|
@ -0,0 +1,362 @@
|
|||
//! Two relay sessions routing through a **real `zenohd`** — the topology every
|
||||
//! other Zenoh test in this crate skips.
|
||||
//!
|
||||
//! `tests/zenoh_bus.rs` proves delivery between two peers that dial *each
|
||||
//! other* (`peer_pair()` cross-connects `listen`/`connect`). That is a direct
|
||||
//! link, and a direct link exercises none of the router's routing tables. The
|
||||
//! deployed topology is the opposite shape: every relay dials one router and no
|
||||
//! relay dials another relay, because both scouting mechanisms are off and
|
||||
//! `connect.endpoints` names only the router tier.
|
||||
//!
|
||||
//! Those two shapes disagree, and the disagreement is silent. With
|
||||
//! `mode: "peer"` plus `scouting.gossip.enabled: false`, `zenohd` will **not**
|
||||
//! forward a subscription declaration from one peer to another
|
||||
//! (`zenoh-1.8.0/src/net/routing/hat/router/pubsub.rs:200-222`), so a publisher
|
||||
//! never learns a remote subscriber exists and drops its own sample before it
|
||||
//! reaches the wire. Nothing on either side counts it: the publisher's
|
||||
//! `published_total` still increments and the receiver's drop counters are
|
||||
//! never even created. That was meridian-2m45: 21 green tests, zero delivery.
|
||||
//!
|
||||
//! The call was to make `client` the default rather than to keep `peer` working
|
||||
//! — the reasoning, the rejected alternative and the reversal condition live on
|
||||
//! `ZenohMode` in `src/zenoh/config.rs`, where the next person to change it will
|
||||
//! be looking. This file is the live half: it routes real sessions through a
|
||||
//! real router, against the two configs that actually ship — the embedded relay
|
||||
//! posture and `deploy/compose/zenoh/zenohd.json5`, the file Compose mounts —
|
||||
//! and it pins BOTH the mode that delivers and the mode that silently does not.
|
||||
//!
|
||||
//! Opt-in because it needs Docker and the pinned image, exactly like
|
||||
//! `zenoh_config::daemon_accepts_the_same_fixtures`:
|
||||
//! `MERIDIAN_ZENOH_DAEMON_CHECK=1 cargo test -p meridian-pubsub --features zenoh --test zenoh_router`
|
||||
//! `just zenoh-config-check` sets it for you.
|
||||
|
||||
#![cfg(feature = "zenoh")]
|
||||
|
||||
// Wrapped in a module so the test *names* carry `zenoh_router`, for the same
|
||||
// reason `zenoh_bus.rs` and `zenoh_config.rs` do it: filtering by an unwrapped
|
||||
// name matches nothing and reports a green "0 passed; N filtered out".
|
||||
mod zenoh_router {
|
||||
use std::process::Command;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use meridian_core::{CommunityId, TenantContext};
|
||||
use meridian_pubsub::zenoh::config::{ZenohBusConfig, ZenohMode};
|
||||
use meridian_pubsub::zenoh::ZenohEventBus;
|
||||
use meridian_pubsub::{ChannelEvent, EventBus, EventTopic};
|
||||
use uuid::Uuid;
|
||||
|
||||
const IMAGE: &str = "eclipse/zenoh:1.8.0";
|
||||
|
||||
fn repo_root() -> std::path::PathBuf {
|
||||
// crates/meridian-pubsub -> crates -> repo root
|
||||
std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
|
||||
.join("../..")
|
||||
.canonicalize()
|
||||
.expect("the crate must live two levels under the repo root")
|
||||
}
|
||||
|
||||
/// The router config Compose mounts — not a test-only copy of it. A fixture
|
||||
/// that agrees with the relay while the deployed file does not is exactly
|
||||
/// the drift this test exists to catch.
|
||||
fn deployed_router_config() -> std::path::PathBuf {
|
||||
repo_root().join("deploy/compose/zenoh/zenohd.json5")
|
||||
}
|
||||
|
||||
fn free_port() -> u16 {
|
||||
let listener =
|
||||
std::net::TcpListener::bind("127.0.0.1:0").expect("a loopback port must be available");
|
||||
listener
|
||||
.local_addr()
|
||||
.expect("a bound listener has an address")
|
||||
.port()
|
||||
}
|
||||
|
||||
/// A live `zenohd` on loopback, started from the deployed config with the
|
||||
/// deployed command line, and removed on drop.
|
||||
struct Router {
|
||||
name: String,
|
||||
port: u16,
|
||||
}
|
||||
|
||||
impl Router {
|
||||
fn start(label: &str) -> Router {
|
||||
let port = free_port();
|
||||
let name = format!("meridian-zenoh-router-test-{label}-{port}");
|
||||
// Idempotent: a previous panic may have leaked the name.
|
||||
let _ = Command::new("docker").args(["rm", "-f", &name]).output();
|
||||
|
||||
let mount = format!(
|
||||
"{}:/etc/zenoh/zenohd.json5:ro",
|
||||
deployed_router_config().display()
|
||||
);
|
||||
let publish = format!("127.0.0.1:{port}:7447");
|
||||
let output = Command::new("docker")
|
||||
.args([
|
||||
"run",
|
||||
"-d",
|
||||
"--name",
|
||||
&name,
|
||||
"--init",
|
||||
"-p",
|
||||
&publish,
|
||||
"-v",
|
||||
&mount,
|
||||
IMAGE,
|
||||
// Byte-for-byte the compose service's `command:`.
|
||||
"--config",
|
||||
"/etc/zenoh/zenohd.json5",
|
||||
"--no-multicast-scouting",
|
||||
"--adminspace-permissions",
|
||||
"none",
|
||||
])
|
||||
.output()
|
||||
.unwrap_or_else(|e| panic!("could not run docker: {e}"));
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"could not start {IMAGE}. That is a Docker problem, not a verdict on the \
|
||||
routing posture — re-run before believing it.\n{}",
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
|
||||
let router = Router { name, port };
|
||||
router.wait_listening();
|
||||
router
|
||||
}
|
||||
|
||||
fn wait_listening(&self) {
|
||||
let deadline = Instant::now() + Duration::from_secs(20);
|
||||
while Instant::now() < deadline {
|
||||
if std::net::TcpStream::connect_timeout(
|
||||
&format!("127.0.0.1:{}", self.port)
|
||||
.parse()
|
||||
.expect("loopback socket address"),
|
||||
Duration::from_millis(250),
|
||||
)
|
||||
.is_ok()
|
||||
{
|
||||
return;
|
||||
}
|
||||
std::thread::sleep(Duration::from_millis(100));
|
||||
}
|
||||
panic!(
|
||||
"zenohd never accepted a connection on 127.0.0.1:{}\n{}",
|
||||
self.port,
|
||||
self.logs()
|
||||
);
|
||||
}
|
||||
|
||||
fn logs(&self) -> String {
|
||||
Command::new("docker")
|
||||
.args(["logs", &self.name])
|
||||
.output()
|
||||
.map(|o| {
|
||||
format!(
|
||||
"--- {} stdout ---\n{}\n--- stderr ---\n{}",
|
||||
self.name,
|
||||
String::from_utf8_lossy(&o.stdout),
|
||||
String::from_utf8_lossy(&o.stderr)
|
||||
)
|
||||
})
|
||||
.unwrap_or_else(|e| format!("(could not read {} logs: {e})", self.name))
|
||||
}
|
||||
|
||||
fn endpoint(&self) -> String {
|
||||
format!("tcp/127.0.0.1:{}", self.port)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for Router {
|
||||
fn drop(&mut self) {
|
||||
let _ = Command::new("docker")
|
||||
.args(["rm", "-f", &self.name])
|
||||
.output();
|
||||
}
|
||||
}
|
||||
|
||||
/// A relay-shaped session: the embedded pinned posture plus this router's
|
||||
/// endpoint. No listen override — the fixture's ephemeral port is what a
|
||||
/// pod uses, and a pod is dialled by nobody.
|
||||
fn relay_config(router: &Router, mode: ZenohMode) -> ZenohBusConfig {
|
||||
ZenohBusConfig::new(mode, [router.endpoint()])
|
||||
.expect("a loopback router endpoint is accepted")
|
||||
.with_readiness_timeout(Duration::from_secs(20))
|
||||
.with_loopback_freshness(Duration::from_secs(60))
|
||||
}
|
||||
|
||||
fn ctx() -> TenantContext {
|
||||
TenantContext::resolved(
|
||||
CommunityId::from_uuid(Uuid::from_u128(0x0a11)),
|
||||
"bus.example",
|
||||
)
|
||||
}
|
||||
|
||||
async fn recv_within(
|
||||
rx: &mut tokio::sync::broadcast::Receiver<ChannelEvent>,
|
||||
budget: Duration,
|
||||
) -> Option<ChannelEvent> {
|
||||
tokio::time::timeout(budget, rx.recv()).await.ok()?.ok()
|
||||
}
|
||||
|
||||
fn sample_event(content: &str) -> nostr::Event {
|
||||
let keys = nostr::Keys::generate();
|
||||
nostr::EventBuilder::new(nostr::Kind::from_u16(9), content)
|
||||
.sign_with_keys(&keys)
|
||||
.expect("signing a local event")
|
||||
}
|
||||
|
||||
/// Two sessions, one router, one event. The whole finding in one assertion.
|
||||
async fn an_event_crosses_two_relays_through(router: &Router, mode: ZenohMode) {
|
||||
let a = ZenohEventBus::open(relay_config(router, mode))
|
||||
.await
|
||||
.expect("relay A must open");
|
||||
let b = ZenohEventBus::open(relay_config(router, mode))
|
||||
.await
|
||||
.expect("relay B must open");
|
||||
|
||||
a.wait_ready()
|
||||
.await
|
||||
.unwrap_or_else(|e| panic!("relay A must reach readiness through the router: {e:?}"));
|
||||
b.wait_ready()
|
||||
.await
|
||||
.unwrap_or_else(|e| panic!("relay B must reach readiness through the router: {e:?}"));
|
||||
|
||||
// Readiness is not reachability. Both sides report a transport to the
|
||||
// router whether or not the router will carry a declaration between
|
||||
// them, which is why this test does not stop here.
|
||||
let snapshot = a.health().snapshot();
|
||||
assert!(snapshot.session_open);
|
||||
assert!(snapshot.endpoints_connected >= 1);
|
||||
|
||||
let ctx = ctx();
|
||||
let topic = EventTopic::Channel(Uuid::from_u128(0xc0ffee));
|
||||
|
||||
a.retain_topic(&ctx, topic).await;
|
||||
let mut rx = EventBus::subscribe_local(&a);
|
||||
// A's declaration has to travel A -> router -> B before B's publication
|
||||
// can be routed back. That hop is the thing under test.
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
|
||||
let event = sample_event("across the router tier");
|
||||
b.publish_event(&ctx, topic, &event)
|
||||
.await
|
||||
.expect("publish must be accepted");
|
||||
|
||||
let received = recv_within(&mut rx, Duration::from_secs(10)).await;
|
||||
let received = received.unwrap_or_else(|| {
|
||||
panic!(
|
||||
"relay B's event never reached relay A through `{}` in `{}` mode.\n\
|
||||
Both sessions are open and connected, so this is a ROUTING refusal, not a \
|
||||
transport fault. The router decides which faces a subscription declaration \
|
||||
reaches in `zenoh-1.8.0/src/net/routing/hat/router/pubsub.rs:200-222`; a \
|
||||
Client face is exempt from the peer-to-peer gate there, so a `client` \
|
||||
session failing means something upstream of that gate changed — the \
|
||||
posture, the pinned image, or the deployed router config.\n{}",
|
||||
router.name,
|
||||
mode.as_str(),
|
||||
router.logs()
|
||||
)
|
||||
});
|
||||
assert_eq!(received.event, event);
|
||||
assert_eq!(received.community_id, ctx.community());
|
||||
assert_eq!(received.topic, topic);
|
||||
|
||||
a.close().await.expect("clean shutdown");
|
||||
b.close().await.expect("clean shutdown");
|
||||
}
|
||||
|
||||
fn skip_unless_opted_in() -> bool {
|
||||
if std::env::var("MERIDIAN_ZENOH_DAEMON_CHECK").as_deref() == Ok("1") {
|
||||
return false;
|
||||
}
|
||||
eprintln!(
|
||||
"skipped: set MERIDIAN_ZENOH_DAEMON_CHECK=1 (needs Docker and {IMAGE}) to route \
|
||||
two relay sessions through a real router"
|
||||
);
|
||||
true
|
||||
}
|
||||
|
||||
/// The default a deployment gets by doing nothing: `MERIDIAN_ZENOH_MODE` is
|
||||
/// commented out in `.env.example`, so `from_env` takes
|
||||
/// `ZenohMode::default()`.
|
||||
///
|
||||
/// It delivers because a Client face is exempt from the peer-to-peer gate —
|
||||
/// `propagate_simple_subscription_to`'s `src_face.whatami != WhatAmI::Peer`
|
||||
/// arm short-circuits, so the router forwards a client's declaration to
|
||||
/// every non-router face. Verified to hold against a `peer_to_peer` router
|
||||
/// **and** a `linkstate` one, which is the property that made it the call:
|
||||
/// it needs no agreement with the router's own config.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
async fn the_default_mode_routes_between_two_relays_through_a_real_router() {
|
||||
if skip_unless_opted_in() {
|
||||
return;
|
||||
}
|
||||
assert_eq!(
|
||||
ZenohMode::default(),
|
||||
ZenohMode::Client,
|
||||
"this test is the DEFAULT mode's delivery contract. If the default moves, \
|
||||
move this call with it — a default that no live test routes through a real \
|
||||
router is exactly how meridian-2m45 shipped with 21 green tests."
|
||||
);
|
||||
let router = Router::start("default");
|
||||
an_event_crosses_two_relays_through(&router, ZenohMode::default()).await;
|
||||
}
|
||||
|
||||
/// **This is the regression test for meridian-2m45.** It asserts the
|
||||
/// *absence* of delivery, because the defect is that `peer` mode looks
|
||||
/// perfectly healthy while carrying nothing: both sessions open, both
|
||||
/// connect, both reach readiness, the publisher counts a publication, and
|
||||
/// the receiver never even creates a drop counter.
|
||||
///
|
||||
/// Asserting the failure rather than deleting the case is deliberate. `peer`
|
||||
/// remains a supported mode for a *direct* pod-to-pod link
|
||||
/// (`zenoh_bus.rs::peer_pair`), so it cannot simply be removed — which means
|
||||
/// the trap stays reachable by anyone who sets `MERIDIAN_ZENOH_MODE=peer`
|
||||
/// against a router. This pins the trap's exact shape, so that the day
|
||||
/// upstream fixes it, or someone enables gossip or `routing.peer.mode`, the
|
||||
/// change announces itself here instead of being discovered in production.
|
||||
///
|
||||
/// If this test starts FAILING because the event *did* arrive, that is good
|
||||
/// news and the note on `ZenohMode` needs rewriting — do not simply invert
|
||||
/// the assertion without finding out which of the three levers moved.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
async fn peer_mode_silently_delivers_nothing_through_a_router_and_still_reports_ready() {
|
||||
if skip_unless_opted_in() {
|
||||
return;
|
||||
}
|
||||
let router = Router::start("peer");
|
||||
let a = ZenohEventBus::open(relay_config(&router, ZenohMode::Peer))
|
||||
.await
|
||||
.expect("relay A must open");
|
||||
let b = ZenohEventBus::open(relay_config(&router, ZenohMode::Peer))
|
||||
.await
|
||||
.expect("relay B must open");
|
||||
|
||||
// The half that lies. Both of these succeed in the broken topology,
|
||||
// which is why no probe, metric or dashboard caught it.
|
||||
a.wait_ready().await.expect("peer A reports ready");
|
||||
b.wait_ready().await.expect("peer B reports ready");
|
||||
assert!(a.health().snapshot().session_open);
|
||||
assert!(a.health().snapshot().endpoints_connected >= 1);
|
||||
|
||||
let ctx = ctx();
|
||||
let topic = EventTopic::Channel(Uuid::from_u128(0xc0ffee));
|
||||
a.retain_topic(&ctx, topic).await;
|
||||
let mut rx = EventBus::subscribe_local(&a);
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
|
||||
b.publish_event(&ctx, topic, &sample_event("into the void"))
|
||||
.await
|
||||
.expect("the publish is ACCEPTED — that is the trap");
|
||||
|
||||
assert!(
|
||||
recv_within(&mut rx, Duration::from_secs(3)).await.is_none(),
|
||||
"peer mode delivered through the router. Either upstream changed \
|
||||
hat/router/pubsub.rs:200-222, or gossip / routing.peer.mode is no longer at \
|
||||
the posture this repo pins. Find out which before touching this assertion — \
|
||||
and if peer-through-router now works, the ZenohMode note is stale."
|
||||
);
|
||||
|
||||
a.close().await.expect("clean shutdown");
|
||||
b.close().await.expect("clean shutdown");
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue