feat(bus): implement ZenohEventBus, its shadow adapter, and the trust ladder
Some checks failed
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
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
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

Phase 3 of the Zenoh transport charter. The second `EventBus` implementor, the
shadow adapter that composes it with Redis, and the MIP-XP profile type the
relay's cross-node re-verification branch will be written against.

The default build is unchanged and proved so: `cargo tree -p meridian-pubsub`
is byte-identical before and after (312 lines, same sha256). `postcard` is the
only new dependency edge and it was already in `Cargo.lock` via
`meridian-relay-mesh`, so no package is added to any graph.

What is here, and why each piece is shaped the way it is:

- **`zenoh/topic.rs`** — the key space, exactly as the charter states it. Zenoh's
  wildcards are `*` and `**`; `+` is a legal *literal* chunk, so an MQTT-shaped
  subscription declares fine and then receives silence. That is named as its own
  error, and asserted against Zenoh's own matcher rather than string equality.
  Parsing accepts only canonical spellings: `Uuid::parse_str` also takes braced,
  URN and unhyphenated forms and `"+9".parse::<u32>()` is `Ok(9)`, and every one
  of those renders a *different* key for the same logical topic. `v1` is a wire
  fence — `meridian/v1/**` provably does not intersect a `v2` key.

- **`zenoh/codec.rs`** — postcard, and the routing attachment. The attachment
  carries `community`, `channel`, `kind`, `event_id` in 71 bytes, so a receiver
  answers "is this my own echo?" without deserializing the payload. That is
  charter Gap #5, and it is enforced by a signature rather than by call order:
  `decide_from_attachment` takes no payload, so it cannot decode the event even
  if a later edit wanted it to. Three types needed postcard-friendly mirrors —
  `nostr::Event` (self-describing `Deserialize`), and `CacheInvalidation` /
  `ConnControl`, which are `#[serde(tag = "op")]` and cannot be decoded from a
  non-self-describing format at all. A test fails if that ever stops being true.

- **`zenoh/config.rs`** — one statement of the posture. `tests/fixtures/zenoh/
  bus-peer.json5` is embedded with `include_str!` and `check_posture` is the only
  implementation of the policy check; `tests/zenoh_config.rs` now delegates to it
  instead of keeping a copy. Fails closed on: no endpoint, unknown mode (`router`
  is not a mode a relay may take), a link this build did not compile, TLS
  material with no TLS link, a lost pin, and **any non-loopback endpoint** — the
  charter requires the re-verification branch before a non-local endpoint exists,
  that branch is the relay slice, so until it lands this adapter refuses to be
  reachable from off-host.

- **`zenoh/mod.rs`** — `ZenohEventBus`. One long-lived session; declared
  publishers and subscribers created once and cached, because the first
  publication of a session costs ~15.5× the steady rate. `retain_topic` /
  `release_topic` are now the declaration's lifetime: the retain count and the
  live `Subscriber` are one map entry, so "declared" and "desired" cannot drift,
  and there is no 500 ms debounce because undeclare is local. QoS comes from
  `meridian_core::qos::kind_to_lane`, translated into Zenoh's types here and
  nowhere else; `meridian-core` keeps no bearer dependency. `express` is taken
  only where a lane earned it — Q0/Q1 always, Q2 never without a measurement,
  because express costs about half the throughput below 16 KB.

- **`zenoh/health.rs`** — readiness is configured connects **plus** required
  declarations **plus** a fresh application loopback, and `readiness_gap` names
  the first unmet condition rather than answering a bare `false`.

- **`zenoh/shadow.rs`** — dual publish, Redis-authoritative delivery, per-lane
  parity counters. A Zenoh success never rescues a Redis failure: the caller sees
  the primary's error, because a shadow that could rescue the primary makes the
  primary's health unobservable.

- **`trust.rs`** — the MIP-XP ladder, `LaneFloors`, and `admit(link, claimed,
  lane, floors, proof) -> Admission`. Not feature-gated: the branch it exists for
  compiles in every build. Effective trust is the meet of claim and link, and an
  unknown wire value maps to least trust.

- **`bus.rs`** — `AnyEventBus`, the enum dispatcher. `EventBus` is RPITIT and not
  dyn-compatible; `async-trait` would put a heap allocation on the publish path
  of a bus built to remove per-event cost. `bridge` gets no variant, because
  nothing constructs one yet.

The subscriber→`broadcast::Sender` adaptation needs no runner:
`broadcast::Sender::send` is synchronous, so the Zenoh callback *is* the adapter.
`run_broadcast_consumer` remains the only place lag is interpreted.

Not in this commit, by scope: relay wiring (`config.rs`, `main.rs`,
`reconcile.rs`, `metrics.rs`), the `BusBackend` variants, `.env.example` and
Helm. `bridge` is unimplemented.

Gates: `cargo test -p meridian-pubsub` 0; with `--features zenoh` 0 (92 lib + 21
zenoh_bus + 11 zenoh_config + 1 doc); clippy both ways 0; `cargo fmt -p
meridian-pubsub -- --check` 0; `just check-unwrap-budget` 0 (no new
unwrap/expect in any production path); `just check-alert-runbooks` 0;
`just check-architecture-map` 0.

Signed-off-by: Joshua Belke <joshua@innovationhub-act.org>
This commit is contained in:
Josh Belke 2026-08-19 15:59:51 -04:00
commit 25349f15c3
14 changed files with 5346 additions and 166 deletions

1
Cargo.lock generated
View file

@ -4669,6 +4669,7 @@ dependencies = [
"meridian-core",
"metrics",
"nostr",
"postcard",
"redis",
"serde",
"serde_json",

View file

@ -27,6 +27,9 @@ metrics = { workspace = true }
# checkout's default build and dependency graph are unchanged. See the workspace
# `Cargo.toml` for the exact-version pin and the per-feature justification.
zenoh = { workspace = true, optional = true }
# Bus payload codec. Not workspace-inherited only because it is optional here;
# the pin matches the workspace's (`postcard = "1", use-std`).
postcard = { version = "1", default-features = false, features = ["use-std"], optional = true }
[features]
default = []
@ -35,7 +38,7 @@ default = []
# cost every contributor pays. Measured: `cargo tree -p meridian-pubsub` is
# byte-identical with the feature off (144 packages) and 295 with it on;
# `Cargo.lock` carries 73 more entries either way.
zenoh = ["dep:zenoh"]
zenoh = ["dep:zenoh", "dep:postcard"]
[dev-dependencies]
tokio = { workspace = true }
tokio = { workspace = true, features = ["macros", "rt-multi-thread", "time"] }

View file

@ -410,12 +410,160 @@ impl EventBus for RedisEventBus {
}
}
/// Compile-time assertion that the only backend still satisfies the seam,
/// including its `Send + Sync` supertraits — a non-`Send` field added to
/// [`RedisEventBus`] fails here rather than at a distant spawn site.
/// Every constructed [`EventBus`] backend, as one dispatcher.
///
/// # Why an enum and not `Box<dyn EventBus>`
///
/// [`EventBus`] is RPITIT (`-> impl Future + Send`) and therefore **not**
/// dyn-compatible. Making it dyn-compatible would mean `async-trait`, which
/// boxes a future on every call — a heap allocation per publish, on the publish
/// path of a bus whose entire purpose is to remove per-event cost. The enum
/// costs an exhaustive `match` instead, and the exhaustiveness is itself the
/// guard: a variant cannot be added without every method being taught what to
/// do with it.
///
/// **Reversal evidence:** a measured need for runtime backend selection that an
/// enum cannot express, weighed against the allocation it would add to every
/// publish. Reversing means adding `async-trait` to this crate's manifest,
/// which is a separate dependency decision.
///
/// # What is deliberately not here
///
/// `bridge` — dual publish *and* dual consume with bounded dedup — has no
/// variant, because nothing constructs one yet. The same rule the relay's
/// `MERIDIAN_BUS` parser follows: a mode gets a variant in the commit that
/// builds it, never before.
pub enum AnyEventBus {
/// Redis/Dragonfly PUB/SUB — the fresh-checkout default.
Redis(RedisEventBus),
/// Zenoh only. Serves from Zenoh; Redis keeps presence and rate limits.
#[cfg(feature = "zenoh")]
Zenoh(crate::zenoh::ZenohEventBus),
/// Dual publish, Redis-authoritative delivery, parity counted.
#[cfg(feature = "zenoh")]
Shadow(crate::zenoh::ShadowEventBus),
}
impl AnyEventBus {
/// Operator-facing label for logs, readiness output and metrics.
#[must_use]
pub const fn label(&self) -> &'static str {
match self {
AnyEventBus::Redis(_) => "redis",
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(_) => "zenoh",
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(_) => "shadow",
}
}
}
impl EventBus for AnyEventBus {
async fn publish_event(
&self,
ctx: &TenantContext,
topic: EventTopic,
event: &nostr::Event,
) -> Result<(), BusError> {
match self {
AnyEventBus::Redis(bus) => bus.publish_event(ctx, topic, event).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => bus.publish_event(ctx, topic, event).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => bus.publish_event(ctx, topic, event).await,
}
}
fn subscribe_local(&self) -> broadcast::Receiver<ChannelEvent> {
match self {
AnyEventBus::Redis(bus) => EventBus::subscribe_local(bus),
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => EventBus::subscribe_local(bus),
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => EventBus::subscribe_local(bus),
}
}
async fn retain_topic(&self, ctx: &TenantContext, topic: EventTopic) {
match self {
AnyEventBus::Redis(bus) => bus.retain_topic(ctx, topic).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => bus.retain_topic(ctx, topic).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => bus.retain_topic(ctx, topic).await,
}
}
async fn release_topic(&self, ctx: &TenantContext, topic: EventTopic) {
match self {
AnyEventBus::Redis(bus) => bus.release_topic(ctx, topic).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => bus.release_topic(ctx, topic).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => bus.release_topic(ctx, topic).await,
}
}
async fn publish_cache_invalidation(
&self,
ctx: &TenantContext,
invalidation: &CacheInvalidation,
) -> Result<(), BusError> {
match self {
AnyEventBus::Redis(bus) => bus.publish_cache_invalidation(ctx, invalidation).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => bus.publish_cache_invalidation(ctx, invalidation).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => bus.publish_cache_invalidation(ctx, invalidation).await,
}
}
async fn publish_conn_control(
&self,
ctx: &TenantContext,
command: &ConnControl,
) -> Result<(), BusError> {
match self {
AnyEventBus::Redis(bus) => bus.publish_conn_control(ctx, command).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => bus.publish_conn_control(ctx, command).await,
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => bus.publish_conn_control(ctx, command).await,
}
}
fn subscribe_cache_invalidations(&self) -> broadcast::Receiver<ScopedCacheInvalidation> {
match self {
AnyEventBus::Redis(bus) => EventBus::subscribe_cache_invalidations(bus),
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => EventBus::subscribe_cache_invalidations(bus),
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => EventBus::subscribe_cache_invalidations(bus),
}
}
fn subscribe_conn_control(&self) -> broadcast::Receiver<ScopedConnControl> {
match self {
AnyEventBus::Redis(bus) => EventBus::subscribe_conn_control(bus),
#[cfg(feature = "zenoh")]
AnyEventBus::Zenoh(bus) => EventBus::subscribe_conn_control(bus),
#[cfg(feature = "zenoh")]
AnyEventBus::Shadow(bus) => EventBus::subscribe_conn_control(bus),
}
}
}
/// Compile-time assertion that every constructed backend satisfies the seam,
/// including its `Send + Sync` supertraits — a non-`Send` field added to any of
/// them fails here rather than at a distant spawn site.
const _: fn() = || {
fn assert_event_bus<T: EventBus>() {}
assert_event_bus::<RedisEventBus>();
assert_event_bus::<AnyEventBus>();
#[cfg(feature = "zenoh")]
assert_event_bus::<crate::zenoh::ZenohEventBus>();
#[cfg(feature = "zenoh")]
assert_event_bus::<crate::zenoh::ShadowEventBus>();
};
#[cfg(test)]

View file

@ -26,6 +26,16 @@ pub enum PubSubError {
/// A Redis channel key could not be parsed as a valid channel ID.
#[error("Invalid channel key: {0}")]
InvalidChannelKey(String),
/// A Zenoh bus operation failed.
///
/// Feature-gated so the default build's error enum is byte-identical:
/// `BusError` is a public alias of this type, and a variant that only
/// exists when a backend exists keeps every `match` in the relay exhaustive
/// without a catch-all arm.
#[cfg(feature = "zenoh")]
#[error(transparent)]
Zenoh(#[from] crate::zenoh::ZenohBusError),
}
impl From<tokio::sync::broadcast::error::RecvError> for PubSubError {

View file

@ -56,7 +56,18 @@ pub mod rate_limiter;
pub mod subscriber;
/// Community-scoped Redis event topics.
pub mod topic;
pub use bus::{BusError, EventBus, RedisEventBus};
/// MIP-XP exchange profiles — the trust ladder the cross-node re-verification
/// branch is written against. Not feature-gated: that branch compiles in every
/// build, so the type it reads must too.
pub mod trust;
/// Zenoh bus transport. Behind the off-by-default `zenoh` cargo feature, so a
/// fresh checkout's default build and dependency graph are unchanged.
#[cfg(feature = "zenoh")]
pub mod zenoh;
pub use bus::{AnyEventBus, BusError, EventBus, RedisEventBus};
pub use error::PubSubError;
use std::collections::HashMap;

View file

@ -0,0 +1,367 @@
//! MIP-XP exchange profiles — the closed trust ladder that the cross-node
//! re-verification branch is written against.
//!
//! This module is deliberately **not** behind the `zenoh` feature and carries
//! no transport dependency. The branch it exists for lives in the relay's
//! `fan_out_pubsub_event`, which compiles in every build, so a profile type
//! that only existed under an optional feature could not be used there.
//!
//! # What this replaces
//!
//! The former sketch was a `trusted_bus: bool`. A boolean cannot express "P2 to
//! that peer, P0 to this one", which is the exact shape a router topology
//! produces the moment more than one link exists. The ladder below is closed
//! and ordered least-to-most trusted, so "at least this much" is a comparison
//! rather than a convention.
//!
//! # The two rules that make it safe
//!
//! - **Effective trust is the meet of the claim and the authenticated link.**
//! An envelope may *claim* a profile; it may never be believed above what the
//! link itself authenticated. [`TrustProfile::meet`] is that rule, and
//! [`admit`] applies it before anything else.
//! - **Unknown maps to least trust.** [`TrustProfile::from_wire`] answers
//! [`TrustProfile::P0None`] for every value it does not recognise, so a newer
//! peer inventing a profile number cannot promote itself.
use meridian_core::qos::Lane;
/// One rung of the MIP-XP profile ladder, ordered least-to-most trusted.
///
/// The derived [`Ord`] is the ladder: `P0None < P1Mac < P2Ed25519 < P3Schnorr`.
/// Variants are never reordered — the comparison is the enforcement.
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Default)]
pub enum TrustProfile {
/// Link identity only (mTLS / QUIC certificate). Legal for server-to-server
/// traffic inside one operator's trust domain. This is what the relay's
/// cross-pod fan-out already runs at today, undeclared.
#[default]
P0None,
/// A symmetric MAC derived from the handshake (HMAC-SHA256 or keyed
/// BLAKE3). Not third-party verifiable.
P1Mac,
/// Ed25519 over the link key. Third-party verifiable against the *link*
/// key, never against an author key, so it never appears on a client
/// connection.
P2Ed25519,
/// BIP-340 secp256k1 — the NIP-01 native profile. Mandatory for anything
/// stored and served over a standard Nostr `REQ`.
P3Schnorr,
}
impl TrustProfile {
/// Every rung, least-trusted first.
pub const ALL: [TrustProfile; 4] = [
TrustProfile::P0None,
TrustProfile::P1Mac,
TrustProfile::P2Ed25519,
TrustProfile::P3Schnorr,
];
/// The wire byte for this profile.
#[must_use]
pub const fn wire_value(self) -> u8 {
match self {
TrustProfile::P0None => 0,
TrustProfile::P1Mac => 1,
TrustProfile::P2Ed25519 => 2,
TrustProfile::P3Schnorr => 3,
}
}
/// Decode a wire byte. **Every unknown value maps to [`Self::P0None`]** —
/// least trust — rather than erroring, because the caller's only safe
/// response to an unrecognised claim is to disbelieve it, and an error path
/// invites a caller to skip the comparison entirely.
#[must_use]
pub const fn from_wire(value: u8) -> TrustProfile {
match value {
1 => TrustProfile::P1Mac,
2 => TrustProfile::P2Ed25519,
3 => TrustProfile::P3Schnorr,
_ => TrustProfile::P0None,
}
}
/// Metric-safe label.
#[must_use]
pub const fn label(self) -> &'static str {
match self {
TrustProfile::P0None => "p0-none",
TrustProfile::P1Mac => "p1-mac",
TrustProfile::P2Ed25519 => "p2-ed25519",
TrustProfile::P3Schnorr => "p3-schnorr",
}
}
/// The lattice meet — the *lower* of two profiles.
///
/// This is the "a claimed envelope profile is never trusted above the
/// authenticated link's profile" rule, as a function.
#[must_use]
pub fn meet(self, other: TrustProfile) -> TrustProfile {
if self <= other {
self
} else {
other
}
}
}
impl std::fmt::Display for TrustProfile {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.label())
}
}
/// The per-lane floor vector a link negotiated (`XP-ACCEPT`).
///
/// A link does not pick one profile; it picks a *vector* of floors indexed by
/// traffic class, and enforcement is per message against the floor for that
/// message's lane. One link can therefore carry P0 typing frames and P3 stored
/// chat at the same time.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LaneFloors([TrustProfile; 8]);
impl LaneFloors {
/// The floors for a link **inside one MERIDIAN deployment**.
///
/// Every lane is [`TrustProfile::P0None`]. That is not a relaxation: it
/// writes down the profile the cross-pod fan-out already runs at, where the
/// link is a private in-cluster service and no re-verification happens
/// today. Raising a floor is what a federation or external-router link
/// does, and it is a configured value, never a default.
pub const IN_DEPLOYMENT: LaneFloors = LaneFloors([TrustProfile::P0None; 8]);
/// A uniform floor across every lane.
#[must_use]
pub const fn uniform(profile: TrustProfile) -> LaneFloors {
LaneFloors([profile; 8])
}
/// The floor for one lane.
#[must_use]
pub const fn floor(&self, lane: Lane) -> TrustProfile {
self.0[lane.wire_value() as usize]
}
/// Raise (or lower) one lane's floor, returning the new vector.
#[must_use]
pub const fn with_floor(mut self, lane: Lane, profile: TrustProfile) -> LaneFloors {
self.0[lane.wire_value() as usize] = profile;
self
}
}
impl Default for LaneFloors {
fn default() -> Self {
LaneFloors::IN_DEPLOYMENT
}
}
/// What proof the payload carries about itself, independent of the link.
///
/// This is the input that decides whether a below-floor message can be rescued
/// by re-verification or must be refused outright.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PayloadProof {
/// A NIP-01 event carrying a BIP-340 signature over its own id. Below-floor
/// arrivals of this shape are re-verified rather than dropped.
SignedNostrEvent,
/// The payload proves nothing about itself — a cache-key drop, a
/// connection-control command. There is nothing to re-verify, so a
/// below-floor arrival is rejected.
None,
}
/// The decision [`admit`] returns for one arriving bus message.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Admission {
/// Transport attribution is accepted for this lane. Tenant and access
/// enforcement still run receiver-side — this is not a delivery decision.
Accept,
/// `verify_event` MUST run before the event is fanned out.
Reverify,
/// Refuse and count. Never "vaguely reverified".
Reject,
}
impl Admission {
/// Metric-safe label.
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Admission::Accept => "accept",
Admission::Reverify => "reverify",
Admission::Reject => "reject",
}
}
}
/// The cross-node admission decision, as a pure function.
///
/// This is the signature the re-verification branch in the relay's
/// `fan_out_pubsub_event` is written against. It makes no I/O, reads no
/// environment, and returns the same answer for the same inputs, so the branch
/// is testable without a bus.
///
/// ```
/// use meridian_core::qos::Lane;
/// use meridian_pubsub::trust::{admit, Admission, LaneFloors, PayloadProof, TrustProfile};
///
/// // In-deployment: every lane floors at P0, so an unauthenticated link is
/// // still accepted — this is what runs today, now written down.
/// assert_eq!(
/// admit(
/// TrustProfile::P0None,
/// TrustProfile::P0None,
/// Lane::Q4,
/// &LaneFloors::IN_DEPLOYMENT,
/// PayloadProof::SignedNostrEvent,
/// ),
/// Admission::Accept,
/// );
///
/// // A federation link floors chat at P3. A P0 link carrying a signed event
/// // is re-verified, not trusted and not dropped.
/// let federated = LaneFloors::IN_DEPLOYMENT.with_floor(Lane::Q4, TrustProfile::P3Schnorr);
/// assert_eq!(
/// admit(
/// TrustProfile::P0None,
/// TrustProfile::P3Schnorr, // the claim is ignored: it exceeds the link
/// Lane::Q4,
/// &federated,
/// PayloadProof::SignedNostrEvent,
/// ),
/// Admission::Reverify,
/// );
///
/// // The same link carrying a cache drop has nothing to re-verify.
/// let control = LaneFloors::IN_DEPLOYMENT.with_floor(Lane::Q0, TrustProfile::P2Ed25519);
/// assert_eq!(
/// admit(
/// TrustProfile::P0None,
/// TrustProfile::P2Ed25519,
/// Lane::Q0,
/// &control,
/// PayloadProof::None,
/// ),
/// Admission::Reject,
/// );
/// ```
#[must_use]
pub fn admit(
link: TrustProfile,
claimed: TrustProfile,
lane: Lane,
floors: &LaneFloors,
proof: PayloadProof,
) -> Admission {
let effective = link.meet(claimed);
if effective >= floors.floor(lane) {
return Admission::Accept;
}
match proof {
PayloadProof::SignedNostrEvent => Admission::Reverify,
PayloadProof::None => Admission::Reject,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_ladder_is_ordered_least_to_most_trusted() {
assert!(TrustProfile::P0None < TrustProfile::P1Mac);
assert!(TrustProfile::P1Mac < TrustProfile::P2Ed25519);
assert!(TrustProfile::P2Ed25519 < TrustProfile::P3Schnorr);
assert_eq!(TrustProfile::default(), TrustProfile::P0None);
}
#[test]
fn unknown_wire_values_map_to_least_trust() {
for value in 0u8..=255 {
let profile = TrustProfile::from_wire(value);
if !(1..=3).contains(&value) {
assert_eq!(
profile,
TrustProfile::P0None,
"wire value {value} must not promote itself"
);
}
}
for profile in TrustProfile::ALL {
assert_eq!(TrustProfile::from_wire(profile.wire_value()), profile);
}
}
#[test]
fn a_claim_can_never_exceed_the_authenticated_link() {
for link in TrustProfile::ALL {
for claimed in TrustProfile::ALL {
let effective = link.meet(claimed);
assert!(effective <= link, "claim {claimed} promoted link {link}");
assert!(effective <= claimed);
assert_eq!(effective, claimed.meet(link), "meet must be commutative");
}
}
}
#[test]
fn in_deployment_floors_every_lane_at_p0() {
for lane in Lane::ALL {
assert_eq!(
LaneFloors::IN_DEPLOYMENT.floor(lane),
TrustProfile::P0None,
"{lane} must record the profile the fan-out already runs at"
);
}
}
#[test]
fn a_raised_floor_moves_only_its_own_lane() {
let floors = LaneFloors::IN_DEPLOYMENT.with_floor(Lane::Q4, TrustProfile::P3Schnorr);
assert_eq!(floors.floor(Lane::Q4), TrustProfile::P3Schnorr);
for lane in Lane::ALL {
if lane != Lane::Q4 {
assert_eq!(floors.floor(lane), TrustProfile::P0None);
}
}
}
#[test]
fn below_floor_control_is_rejected_never_reverified() {
let floors = LaneFloors::uniform(TrustProfile::P2Ed25519);
assert_eq!(
admit(
TrustProfile::P1Mac,
TrustProfile::P3Schnorr,
Lane::Q0,
&floors,
PayloadProof::None,
),
Admission::Reject,
);
}
#[test]
fn at_or_above_floor_is_accepted_on_every_lane() {
let floors = LaneFloors::uniform(TrustProfile::P1Mac);
for lane in Lane::ALL {
for proof in [PayloadProof::SignedNostrEvent, PayloadProof::None] {
assert_eq!(
admit(
TrustProfile::P2Ed25519,
TrustProfile::P1Mac,
lane,
&floors,
proof,
),
Admission::Accept,
"{lane} at floor must accept"
);
}
}
}
}

View file

@ -0,0 +1,674 @@
//! Postcard wire encoding for the bus, and the routing attachment that lets a
//! receiver decide **before** decoding the payload.
//!
//! # The attachment is the point
//!
//! Every publication carries two independent byte strings: a small
//! [`RoutingAttachment`] in Zenoh's attachment slot, and the postcard payload.
//! The receiving pod decodes the attachment — ~70 bytes, no allocation beyond
//! the decoded struct — and can then answer *"is this my own echo?"* from
//! `event_id` alone. On a hit it returns without ever looking at the payload.
//!
//! That is charter Gap #5. Today the local-echo check happens after a full
//! JSON deserialize of an event this pod published seconds earlier, so a
//! two-pod deployment pays a parse for every event twice: once to publish and
//! once to throw away.
//!
//! # Why not just serde the existing types
//!
//! Postcard is **not self-describing**. Three things in this crate look
//! encodable and are not:
//!
//! - `nostr::Event`'s `Deserialize` is written for a self-describing format, so
//! [`WireEvent`] mirrors it as a fixed-shape struct with byte arrays instead
//! of hex strings — which is also 2× smaller on the wire.
//! - [`crate::cache_invalidation::CacheInvalidation`] and
//! [`crate::conn_control::ConnControl`] are `#[serde(tag = "op")]`.
//! Internally-tagged enums cannot be deserialized from a non-self-describing
//! format at all. [`WireCacheInvalidation`] and [`WireConnControl`] are their
//! externally-tagged mirrors, and the round-trip tests below are what hold
//! the two shapes together.
//!
//! Every mirror is total: adding a variant to the source enum fails to compile
//! here rather than silently losing a case, because the `From` impls match
//! exhaustively.
use meridian_core::CommunityId;
use serde::{Deserialize, Serialize};
use thiserror::Error;
use uuid::Uuid;
use crate::cache_invalidation::CacheInvalidation;
use crate::conn_control::ConnControl;
/// Wire schema of the bus envelope. Unknown values are rejected and counted,
/// never guessed.
pub const BUS_SCHEMA: u8 = 1;
/// Largest payload this codec will decode, in bytes.
///
/// Tracks the relay's own `DEFAULT_MAX_FRAME_BYTES` (512 KiB): an event too
/// large to enter through the client edge must not be admitted from the bus
/// either, or the bus becomes a way around the edge's limit.
pub const MAX_PAYLOAD_BYTES: usize = 512 * 1024;
/// Largest attachment this codec will decode, in bytes.
///
/// The attachment is a fixed shape whose only variable-length member is a
/// `u32` varint, so 128 is generous by an order of magnitude. It is bounded
/// because the attachment is read *before* anything is authenticated, which
/// makes it the cheapest thing on the path to attack.
pub const MAX_ATTACHMENT_BYTES: usize = 128;
/// Why a bus message could not be encoded or decoded.
#[derive(Debug, Error)]
pub enum CodecError {
/// Postcard failed to encode or decode.
#[error("postcard codec error: {0}")]
Postcard(#[from] postcard::Error),
/// The message decoded, but bytes were left over. A trailing-garbage
/// message is a different message, not a tolerable one.
#[error("{remaining} trailing byte(s) after a complete {what}")]
TrailingBytes {
/// What was being decoded.
what: &'static str,
/// How many bytes were left.
remaining: usize,
},
/// The envelope declares a schema this build does not implement.
#[error("bus schema {found} is not schema {BUS_SCHEMA}")]
UnknownSchema {
/// The schema byte that was found.
found: u8,
},
/// The message is larger than this codec will decode.
#[error("{what} is {found} bytes, over the {limit}-byte bound")]
TooLarge {
/// What was being decoded.
what: &'static str,
/// Observed size.
found: usize,
/// Configured bound.
limit: usize,
},
/// A fixed-width field did not reconstruct a valid Nostr value.
#[error("invalid {field} in a decoded event: {reason}")]
InvalidEventField {
/// Which field.
field: &'static str,
/// Why it was rejected.
reason: String,
},
}
/// Routing and dedup metadata, carried in the Zenoh attachment.
///
/// Field order is the wire order and must not be reordered — postcard encodes
/// structs positionally with no field names.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct RoutingAttachment {
/// Envelope schema; always [`BUS_SCHEMA`] on publish.
pub schema: u8,
/// Server-resolved community, as raw UUID bytes.
pub community: [u8; 16],
/// Exact channel when the topic is channel-scoped, `None` for global.
pub channel: Option<[u8; 16]>,
/// The event kind — the same value as the key expression's `/k/` chunk.
pub kind: u32,
/// The event id, which is what the local-echo check compares.
pub event_id: [u8; 32],
}
impl RoutingAttachment {
/// Build an attachment for a routed event.
#[must_use]
pub fn new(
community: CommunityId,
channel: Option<Uuid>,
kind: u32,
event_id: [u8; 32],
) -> Self {
Self {
schema: BUS_SCHEMA,
community: *community.as_uuid().as_bytes(),
channel: channel.map(|id| *id.as_bytes()),
kind,
event_id,
}
}
/// The community this attachment claims, as a typed id.
///
/// A *claim*, deliberately: the key expression the sample arrived on is the
/// routing fact, and `filter_fanout_by_access` is the enforcement point.
/// Nothing here authorizes anything.
#[must_use]
pub fn community_id(&self) -> CommunityId {
CommunityId::from_uuid(Uuid::from_bytes(self.community))
}
/// The channel this attachment claims, if it is channel-scoped.
#[must_use]
pub fn channel_id(&self) -> Option<Uuid> {
self.channel.map(Uuid::from_bytes)
}
/// Encode for the Zenoh attachment slot.
pub fn encode(&self) -> Result<Vec<u8>, CodecError> {
Ok(postcard::to_stdvec(self)?)
}
/// Decode an attachment, rejecting unknown schemas, oversize input and
/// trailing bytes.
///
/// This is the **only** decode on the local-echo path. It must stay cheap
/// and must never grow a dependency on the payload.
pub fn decode(bytes: &[u8]) -> Result<Self, CodecError> {
if bytes.len() > MAX_ATTACHMENT_BYTES {
return Err(CodecError::TooLarge {
what: "attachment",
found: bytes.len(),
limit: MAX_ATTACHMENT_BYTES,
});
}
let (attachment, rest) = postcard::take_from_bytes::<RoutingAttachment>(bytes)?;
if !rest.is_empty() {
return Err(CodecError::TrailingBytes {
what: "attachment",
remaining: rest.len(),
});
}
if attachment.schema != BUS_SCHEMA {
return Err(CodecError::UnknownSchema {
found: attachment.schema,
});
}
Ok(attachment)
}
}
/// Postcard-encodable mirror of `nostr::Event`.
///
/// Fixed-width byte arrays rather than hex strings: the id, pubkey and
/// signature are 128 bytes here against 256 hex characters in JSON, and the
/// conversion back rejects anything the Nostr types will not accept.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WireEvent {
/// Event id (32 bytes).
pub id: [u8; 32],
/// Author x-only public key (32 bytes).
pub pubkey: [u8; 32],
/// Unix seconds.
pub created_at: u64,
/// Event kind.
pub kind: u16,
/// Tag list, each tag a non-empty list of strings.
pub tags: Vec<Vec<String>>,
/// Event content.
pub content: String,
/// BIP-340 signature, as two 32-byte halves.
///
/// Not `[u8; 64]`: `serde` implements its array traits only up to length
/// 32, so a 64-byte array does not compile. `[[u8; 32]; 2]` keeps the field
/// fixed-width — postcard writes 64 bytes with no length prefix, exactly as
/// `[u8; 64]` would — where a `Vec<u8>` would add a varint and admit a
/// wrong-length signature into the decoder.
pub sig: [[u8; 32]; 2],
}
impl From<&nostr::Event> for WireEvent {
fn from(event: &nostr::Event) -> Self {
Self {
id: event.id.to_bytes(),
pubkey: event.pubkey.to_bytes(),
created_at: event.created_at.as_secs(),
kind: event.kind.as_u16(),
tags: event
.tags
.as_slice()
.iter()
.map(|tag| tag.as_slice().to_vec())
.collect(),
content: event.content.clone(),
sig: split_signature(event.sig.serialize()),
}
}
}
impl TryFrom<WireEvent> for nostr::Event {
type Error = CodecError;
fn try_from(wire: WireEvent) -> Result<Self, Self::Error> {
let pubkey = nostr::PublicKey::from_slice(&wire.pubkey).map_err(|e| {
CodecError::InvalidEventField {
field: "pubkey",
reason: e.to_string(),
}
})?;
let sig = nostr::secp256k1::schnorr::Signature::from_slice(&join_signature(wire.sig))
.map_err(|e| CodecError::InvalidEventField {
field: "sig",
reason: e.to_string(),
})?;
let mut tags = Vec::with_capacity(wire.tags.len());
for tag in wire.tags {
tags.push(
nostr::Tag::parse(tag).map_err(|e| CodecError::InvalidEventField {
field: "tags",
reason: e.to_string(),
})?,
);
}
Ok(nostr::Event::new(
nostr::EventId::from_byte_array(wire.id),
pubkey,
nostr::Timestamp::from_secs(wire.created_at),
nostr::Kind::from_u16(wire.kind),
tags,
wire.content,
sig,
))
}
}
fn split_signature(sig: [u8; 64]) -> [[u8; 32]; 2] {
let mut halves = [[0u8; 32]; 2];
halves[0].copy_from_slice(&sig[..32]);
halves[1].copy_from_slice(&sig[32..]);
halves
}
fn join_signature(halves: [[u8; 32]; 2]) -> [u8; 64] {
let mut sig = [0u8; 64];
sig[..32].copy_from_slice(&halves[0]);
sig[32..].copy_from_slice(&halves[1]);
sig
}
/// Encode an event for the bus payload slot.
pub fn encode_event(event: &nostr::Event) -> Result<Vec<u8>, CodecError> {
let bytes = postcard::to_stdvec(&WireEvent::from(event))?;
if bytes.len() > MAX_PAYLOAD_BYTES {
return Err(CodecError::TooLarge {
what: "event payload",
found: bytes.len(),
limit: MAX_PAYLOAD_BYTES,
});
}
Ok(bytes)
}
/// Decode an event from the bus payload slot.
///
/// Never called on the local-echo path — that is the whole point of
/// [`RoutingAttachment`].
pub fn decode_event(bytes: &[u8]) -> Result<nostr::Event, CodecError> {
let wire: WireEvent = decode_bounded(bytes, "event payload", MAX_PAYLOAD_BYTES)?;
nostr::Event::try_from(wire)
}
/// Externally-tagged mirror of [`CacheInvalidation`].
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum WireCacheInvalidation {
/// See [`CacheInvalidation::Membership`].
Membership {
/// Channel whose membership changed.
channel_id: [u8; 16],
/// Affected member's pubkey bytes.
pubkey: Vec<u8>,
},
/// See [`CacheInvalidation::AccessibleAll`].
AccessibleAll,
/// See [`CacheInvalidation::Visibility`].
Visibility {
/// Channel whose visibility changed.
channel_id: [u8; 16],
},
/// See [`CacheInvalidation::ChannelDeleted`].
ChannelDeleted,
}
impl From<&CacheInvalidation> for WireCacheInvalidation {
fn from(value: &CacheInvalidation) -> Self {
match value {
CacheInvalidation::Membership { channel_id, pubkey } => {
WireCacheInvalidation::Membership {
channel_id: *channel_id.as_bytes(),
pubkey: pubkey.clone(),
}
}
CacheInvalidation::AccessibleAll => WireCacheInvalidation::AccessibleAll,
CacheInvalidation::Visibility { channel_id } => WireCacheInvalidation::Visibility {
channel_id: *channel_id.as_bytes(),
},
CacheInvalidation::ChannelDeleted => WireCacheInvalidation::ChannelDeleted,
}
}
}
impl From<WireCacheInvalidation> for CacheInvalidation {
fn from(value: WireCacheInvalidation) -> Self {
match value {
WireCacheInvalidation::Membership { channel_id, pubkey } => {
CacheInvalidation::Membership {
channel_id: Uuid::from_bytes(channel_id),
pubkey,
}
}
WireCacheInvalidation::AccessibleAll => CacheInvalidation::AccessibleAll,
WireCacheInvalidation::Visibility { channel_id } => CacheInvalidation::Visibility {
channel_id: Uuid::from_bytes(channel_id),
},
WireCacheInvalidation::ChannelDeleted => CacheInvalidation::ChannelDeleted,
}
}
}
/// Externally-tagged mirror of [`ConnControl`].
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum WireConnControl {
/// See [`ConnControl::DisconnectCommunity`].
DisconnectCommunity,
/// See [`ConnControl::DisconnectPubkey`].
DisconnectPubkey {
/// Banned member's pubkey bytes.
pubkey: Vec<u8>,
/// Id echoed in the closing `OK` frame.
event_id: String,
/// Human-readable close reason.
reason: String,
},
/// See [`ConnControl::EvictChannelSubscriptions`].
EvictChannelSubscriptions {
/// Channel whose access was revoked.
channel_id: [u8; 16],
/// Removed member's pubkey bytes.
pubkey: Vec<u8>,
},
}
impl From<&ConnControl> for WireConnControl {
fn from(value: &ConnControl) -> Self {
match value {
ConnControl::DisconnectCommunity => WireConnControl::DisconnectCommunity,
ConnControl::DisconnectPubkey {
pubkey,
event_id,
reason,
} => WireConnControl::DisconnectPubkey {
pubkey: pubkey.clone(),
event_id: event_id.clone(),
reason: reason.clone(),
},
ConnControl::EvictChannelSubscriptions { channel_id, pubkey } => {
WireConnControl::EvictChannelSubscriptions {
channel_id: *channel_id.as_bytes(),
pubkey: pubkey.clone(),
}
}
}
}
}
impl From<WireConnControl> for ConnControl {
fn from(value: WireConnControl) -> Self {
match value {
WireConnControl::DisconnectCommunity => ConnControl::DisconnectCommunity,
WireConnControl::DisconnectPubkey {
pubkey,
event_id,
reason,
} => ConnControl::DisconnectPubkey {
pubkey,
event_id,
reason,
},
WireConnControl::EvictChannelSubscriptions { channel_id, pubkey } => {
ConnControl::EvictChannelSubscriptions {
channel_id: Uuid::from_bytes(channel_id),
pubkey,
}
}
}
}
}
/// Encode a cache-key drop for the bus payload slot.
pub fn encode_cache_invalidation(value: &CacheInvalidation) -> Result<Vec<u8>, CodecError> {
Ok(postcard::to_stdvec(&WireCacheInvalidation::from(value))?)
}
/// Decode a cache-key drop from the bus payload slot.
pub fn decode_cache_invalidation(bytes: &[u8]) -> Result<CacheInvalidation, CodecError> {
let wire: WireCacheInvalidation =
decode_bounded(bytes, "cache invalidation", MAX_PAYLOAD_BYTES)?;
Ok(wire.into())
}
/// Encode a connection-control command for the bus payload slot.
pub fn encode_conn_control(value: &ConnControl) -> Result<Vec<u8>, CodecError> {
Ok(postcard::to_stdvec(&WireConnControl::from(value))?)
}
/// Decode a connection-control command from the bus payload slot.
pub fn decode_conn_control(bytes: &[u8]) -> Result<ConnControl, CodecError> {
let wire: WireConnControl = decode_bounded(bytes, "connection control", MAX_PAYLOAD_BYTES)?;
Ok(wire.into())
}
fn decode_bounded<'a, T: Deserialize<'a>>(
bytes: &'a [u8],
what: &'static str,
limit: usize,
) -> Result<T, CodecError> {
if bytes.len() > limit {
return Err(CodecError::TooLarge {
what,
found: bytes.len(),
limit,
});
}
let (value, rest) = postcard::take_from_bytes::<T>(bytes)?;
if !rest.is_empty() {
return Err(CodecError::TrailingBytes {
what,
remaining: rest.len(),
});
}
Ok(value)
}
#[cfg(test)]
mod tests {
use super::*;
use nostr::{EventBuilder, Keys};
fn sample_event(kind: u16, content: &str) -> nostr::Event {
let keys = Keys::generate();
EventBuilder::new(nostr::Kind::from_u16(kind), content)
.tags([
nostr::Tag::parse(["h", "1a2b"]).expect("tag"),
nostr::Tag::parse(["e", "cafe", "", "root"]).expect("tag"),
])
.sign_with_keys(&keys)
.expect("signing a local event")
}
fn community() -> CommunityId {
CommunityId::from_uuid(Uuid::from_u128(0xaaaa))
}
#[test]
fn attachment_round_trips() {
let channel = Uuid::from_u128(0xbbbb);
for channel in [Some(channel), None] {
let attachment = RoutingAttachment::new(community(), channel, 24200, [7u8; 32]);
let bytes = attachment.encode().expect("encode");
assert_eq!(
RoutingAttachment::decode(&bytes).expect("decode"),
attachment
);
assert_eq!(
RoutingAttachment::decode(&bytes)
.expect("decode")
.channel_id(),
channel
);
}
}
/// The attachment is read on every arriving sample, before anything is
/// trusted. Its cost is therefore a property, not an accident.
#[test]
fn an_attachment_stays_small_enough_to_read_on_every_sample() {
let attachment = RoutingAttachment::new(
community(),
Some(Uuid::from_u128(u128::MAX)),
u32::MAX,
[0xffu8; 32],
);
let bytes = attachment.encode().expect("encode");
assert!(
bytes.len() <= MAX_ATTACHMENT_BYTES,
"worst-case attachment is {} bytes",
bytes.len()
);
// 1 schema + 16 community + 1 tag + 16 channel + 5 varint kind + 32 id.
assert_eq!(bytes.len(), 71, "attachment wire size changed");
}
#[test]
fn an_unknown_schema_is_rejected_not_guessed() {
let mut attachment = RoutingAttachment::new(community(), None, 9, [1u8; 32]);
attachment.schema = 2;
let bytes = postcard::to_stdvec(&attachment).expect("encode");
assert!(matches!(
RoutingAttachment::decode(&bytes),
Err(CodecError::UnknownSchema { found: 2 })
));
}
#[test]
fn trailing_bytes_are_rejected_on_every_decoder() {
let attachment = RoutingAttachment::new(community(), None, 9, [1u8; 32]);
let mut bytes = attachment.encode().expect("encode");
bytes.push(0);
assert!(matches!(
RoutingAttachment::decode(&bytes),
Err(CodecError::TrailingBytes { .. })
));
let mut payload = encode_event(&sample_event(9, "hi")).expect("encode");
payload.push(0);
assert!(matches!(
decode_event(&payload),
Err(CodecError::TrailingBytes { .. })
));
let mut control = encode_conn_control(&ConnControl::DisconnectCommunity).expect("encode");
control.push(0);
assert!(matches!(
decode_conn_control(&control),
Err(CodecError::TrailingBytes { .. })
));
}
#[test]
fn oversize_input_is_refused_before_it_is_parsed() {
assert!(matches!(
RoutingAttachment::decode(&[0u8; MAX_ATTACHMENT_BYTES + 1]),
Err(CodecError::TooLarge { .. })
));
assert!(matches!(
decode_event(&vec![0u8; MAX_PAYLOAD_BYTES + 1]),
Err(CodecError::TooLarge { .. })
));
}
#[test]
fn events_round_trip_through_postcard() {
for (kind, content) in [
(9u16, "hello"),
(24200, ""),
(39000, "a longer body with unicode — ✅ and \"quotes\""),
] {
let event = sample_event(kind, content);
let bytes = encode_event(&event).expect("encode");
let decoded = decode_event(&bytes).expect("decode");
assert_eq!(decoded, event, "kind {kind} did not round-trip");
assert_eq!(decoded.id, event.id);
assert_eq!(decoded.sig, event.sig);
assert_eq!(decoded.tags.as_slice(), event.tags.as_slice());
}
}
/// Postcard is smaller than the JSON the Redis path carries. Recorded as a
/// test rather than a claim, because "postcard, not JSON" is a charter
/// decision that should stop being true loudly if it ever stops being true.
#[test]
fn postcard_is_smaller_than_the_json_the_redis_path_publishes() {
use nostr::JsonUtil;
let event = sample_event(9, "a representative chat message body");
let postcard_len = encode_event(&event).expect("encode").len();
let json_len = event.as_json().len();
assert!(
postcard_len < json_len,
"postcard {postcard_len} B is not smaller than JSON {json_len} B"
);
}
#[test]
fn cache_invalidations_round_trip_through_their_mirror() {
let channel_id = Uuid::from_u128(0xcccc);
for value in [
CacheInvalidation::Membership {
channel_id,
pubkey: vec![1, 2, 3],
},
CacheInvalidation::AccessibleAll,
CacheInvalidation::Visibility { channel_id },
CacheInvalidation::ChannelDeleted,
] {
let bytes = encode_cache_invalidation(&value).expect("encode");
assert_eq!(decode_cache_invalidation(&bytes).expect("decode"), value);
}
}
#[test]
fn conn_control_commands_round_trip_through_their_mirror() {
let channel_id = Uuid::from_u128(0xdddd);
for value in [
ConnControl::DisconnectCommunity,
ConnControl::DisconnectPubkey {
pubkey: vec![9; 32],
event_id: "abc".to_string(),
reason: "banned: spam".to_string(),
},
ConnControl::EvictChannelSubscriptions {
channel_id,
pubkey: vec![4; 32],
},
] {
let bytes = encode_conn_control(&value).expect("encode");
assert_eq!(decode_conn_control(&bytes).expect("decode"), value);
}
}
/// The internally-tagged source enums are what makes the mirrors
/// necessary. If a future serde release makes them postcard-decodable, this
/// test fails and the mirrors can be deleted — which is the only way that
/// decision gets revisited.
#[test]
fn the_source_enums_still_cannot_be_decoded_by_postcard() {
let value = CacheInvalidation::AccessibleAll;
let encoded = postcard::to_stdvec(&value);
let decoded = encoded
.as_ref()
.ok()
.and_then(|bytes| postcard::from_bytes::<CacheInvalidation>(bytes).ok());
assert!(
decoded.is_none(),
"`#[serde(tag)]` became postcard-decodable; the wire mirrors are now redundant"
);
}
}

View file

@ -0,0 +1,688 @@
//! Typed bus configuration, and the one statement of the pinned Zenoh posture.
//!
//! # There is exactly one copy of the posture
//!
//! [`PINNED_PEER_CONFIG`] is `tests/fixtures/zenoh/bus-peer.json5`, embedded at
//! compile time. The session this crate opens in production is that file plus
//! the deployment's endpoints — not a second, hand-maintained set of the same
//! decisions. The direction matters: `src/` reads the fixture the Phase 1
//! contract tests and `daemon-config-check.sh` already defend, so the library,
//! the in-process test suite and the `zenohd` image cannot drift apart one
//! reviewed diff at a time.
//!
//! [`check_posture`] is the other half of that. It is the *only* implementation
//! of the policy check; `tests/zenoh_config.rs` calls this function rather than
//! keeping its own copy, so a value that stops being defended stops being
//! defended in one place.
//!
//! # Everything here fails closed
//!
//! Each refusal below exists because its permissive form is silent:
//!
//! - **No endpoint** — an empty `connect.endpoints` with scouting off is a
//! session that opens, reports healthy, and reaches nothing.
//! - **Unknown mode** — `router` parses fine and turns a relay pod into a
//! routing node in a fabric that did not plan for it.
//! - **A non-`tcp` endpoint** — this build is `default-features = false,
//! features = ["transport_tcp"]`, so `tls/` and `quic/` fail at `open` with a
//! protocol error that reads like a network fault.
//! - **TLS material with no TLS link compiled** — an operator who supplies
//! certificates has stated an intent this binary cannot honour, and honouring
//! it partially (TCP, unauthenticated) is the worst of the three outcomes.
//! - **A non-loopback endpoint** — see [`ZenohConfigError::NonLocalEndpoint`].
//! The charter's Phase 3 non-negotiable #3 requires the re-verification
//! branch in the relay's `fan_out_pubsub_event` *before* any non-local
//! endpoint can be configured. That branch is a separate slice, so until it
//! lands the only configuration this adapter accepts is one that cannot reach
//! another trust domain.
use std::time::Duration;
use thiserror::Error;
use crate::trust::{LaneFloors, TrustProfile};
/// The pinned relay-side Zenoh posture, embedded from the Phase 1 fixture.
///
/// Read the fixture's own `README.md` for what each value defends. Every one of
/// them has an upstream default that is the opposite of what MERIDIAN needs.
pub const PINNED_PEER_CONFIG: &str = include_str!("../../tests/fixtures/zenoh/bus-peer.json5");
/// Comma-separated Zenoh endpoints this pod connects to. Required.
pub const ENV_ENDPOINTS: &str = "MERIDIAN_ZENOH_ENDPOINTS";
/// Session mode: `peer` (default) or `client`.
pub const ENV_MODE: &str = "MERIDIAN_ZENOH_MODE";
/// Optional explicit listen endpoint. Defaults to the fixture's ephemeral port.
pub const ENV_LISTEN: &str = "MERIDIAN_ZENOH_LISTEN";
/// TLS material variables. Set any of these and this build fails closed,
/// because it compiled no TLS link.
pub const ENV_TLS_MATERIAL: [&str; 3] = [
"MERIDIAN_ZENOH_TLS_ROOT_CA",
"MERIDIAN_ZENOH_TLS_CERT",
"MERIDIAN_ZENOH_TLS_KEY",
];
/// 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.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ZenohMode {
/// Peer of the router tier — the production and benchmark shape.
#[default]
Peer,
/// Client of a router. Retained for a bounded single-node profile.
Client,
}
impl ZenohMode {
/// The Zenoh config value for this mode.
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
ZenohMode::Peer => "peer",
ZenohMode::Client => "client",
}
}
/// Parse a mode. Unknown values are an error, never a default.
pub fn parse(raw: &str) -> Result<ZenohMode, ZenohConfigError> {
match raw.trim() {
"peer" => Ok(ZenohMode::Peer),
"client" => Ok(ZenohMode::Client),
other => Err(ZenohConfigError::UnknownMode(other.to_string())),
}
}
}
/// Why a Zenoh bus configuration is refused.
#[derive(Debug, Clone, PartialEq, Eq, Error)]
pub enum ZenohConfigError {
/// No endpoint was configured. Topology is declared, never discovered, and
/// both scouting mechanisms are off — so an endpoint-less session reaches
/// nothing while reporting a healthy process.
#[error("{ENV_ENDPOINTS} is required: topology is declared, never discovered")]
NoEndpoints,
/// The session mode is not one this relay may take.
#[error("`{0}` is not a session mode; the relay is a `peer` or a `client`, never a router")]
UnknownMode(String),
/// An endpoint names a link protocol this build did not compile.
#[error(
"endpoint `{0}` needs a link this build did not compile \
(`zenoh` is default-features = false, features = [\"transport_tcp\"])"
)]
UnsupportedLinkProtocol(String),
/// An endpoint is not a Zenoh locator.
#[error("endpoint `{0}` is not a Zenoh locator (expected `tcp/host:port`)")]
MalformedEndpoint(String),
/// TLS material was supplied to a build with no TLS link.
#[error(
"{0} is set but this build compiled no TLS link; \
remove it or build with the TLS link and its reviewed pin"
)]
TlsMaterialWithoutTlsLink(String),
/// A non-loopback endpoint was configured before the trust-profile branch
/// exists.
///
/// The charter's dashed red edge: with no re-verification branch, anything
/// reaching a peered router is delivered to clients as an authentic event.
/// The branch is `fan_out_pubsub_event`'s and lands in the relay slice;
/// [`crate::trust::admit`] is the signature it is written against. Until
/// then this adapter refuses to be reachable from outside the host.
#[error(
"endpoint `{0}` is not loopback. A non-local endpoint needs the cross-node \
re-verification branch in `fan_out_pubsub_event` (see `meridian_pubsub::trust::admit`), \
which has not landed. Without it a peered router is an event-injection primitive."
)]
NonLocalEndpoint(String),
/// The built session config violates the pinned posture.
#[error("pinned Zenoh posture violated: {0}")]
Posture(String),
/// The embedded fixture or an inserted value did not parse.
#[error("Zenoh config error: {0}")]
Invalid(String),
}
/// Everything the Zenoh bus needs to open a session and report readiness.
#[derive(Debug, Clone)]
pub struct ZenohBusConfig {
/// Session mode.
pub mode: ZenohMode,
/// Explicitly configured connect endpoints. Never empty.
pub endpoints: Vec<String>,
/// Explicit listen endpoint, or `None` for the fixture's ephemeral port.
pub listen: Option<String>,
/// The profile this link authenticated to.
///
/// `P0None` while the only legal endpoints are loopback: link identity is
/// the host itself. It rises when a TLS/QUIC link with mutual
/// authentication is compiled and configured.
pub link_profile: TrustProfile,
/// Per-lane trust floors this link negotiated.
pub lane_floors: LaneFloors,
/// How long [`crate::zenoh::ZenohEventBus::wait_ready`] waits for connects
/// and declarations.
pub readiness_timeout: Duration,
/// How stale the loopback probe may be before readiness fails.
pub loopback_freshness: Duration,
/// Bound on the declared-publisher cache.
pub publisher_cache_capacity: usize,
/// Bound on the local-echo id set.
pub local_echo_capacity: usize,
}
impl ZenohBusConfig {
/// Default readiness budget.
pub const DEFAULT_READINESS_TIMEOUT: Duration = Duration::from_secs(10);
/// Default loopback freshness budget.
pub const DEFAULT_LOOPBACK_FRESHNESS: Duration = Duration::from_secs(60);
/// Default declared-publisher cache bound.
pub const DEFAULT_PUBLISHER_CACHE_CAPACITY: usize = 4096;
/// Default local-echo id-set bound.
pub const DEFAULT_LOCAL_ECHO_CAPACITY: usize = 16_384;
/// Build a configuration for explicit endpoints, validating every rule.
pub fn new<I, S>(mode: ZenohMode, endpoints: I) -> Result<Self, ZenohConfigError>
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
let endpoints: Vec<String> = endpoints
.into_iter()
.map(Into::into)
.map(|e| e.trim().to_string())
.filter(|e| !e.is_empty())
.collect();
let config = Self {
mode,
endpoints,
listen: None,
link_profile: TrustProfile::P0None,
lane_floors: LaneFloors::IN_DEPLOYMENT,
readiness_timeout: Self::DEFAULT_READINESS_TIMEOUT,
loopback_freshness: Self::DEFAULT_LOOPBACK_FRESHNESS,
publisher_cache_capacity: Self::DEFAULT_PUBLISHER_CACHE_CAPACITY,
local_echo_capacity: Self::DEFAULT_LOCAL_ECHO_CAPACITY,
};
config.validate()?;
Ok(config)
}
/// Read the configuration from the environment.
///
/// `MERIDIAN_BUS` selection itself belongs to the relay; this reads only
/// the Zenoh-specific variables, so the two can be tested apart.
pub fn from_env() -> Result<Self, ZenohConfigError> {
for name in ENV_TLS_MATERIAL {
if std::env::var_os(name).is_some_and(|v| !v.is_empty()) {
return Err(ZenohConfigError::TlsMaterialWithoutTlsLink(
name.to_string(),
));
}
}
let mode = match std::env::var(ENV_MODE) {
Ok(raw) if !raw.trim().is_empty() => ZenohMode::parse(&raw)?,
_ => ZenohMode::default(),
};
let raw = std::env::var(ENV_ENDPOINTS).unwrap_or_default();
let mut config = Self::new(mode, raw.split(','))?;
if let Ok(listen) = std::env::var(ENV_LISTEN) {
let listen = listen.trim().to_string();
if !listen.is_empty() {
validate_endpoint(&listen)?;
config.listen = Some(listen);
}
}
Ok(config)
}
/// Override the readiness budget.
#[must_use]
pub fn with_readiness_timeout(mut self, timeout: Duration) -> Self {
self.readiness_timeout = timeout;
self
}
/// Override the loopback freshness budget.
#[must_use]
pub fn with_loopback_freshness(mut self, freshness: Duration) -> Self {
self.loopback_freshness = freshness;
self
}
/// Override the explicit listen endpoint.
pub fn with_listen(mut self, listen: impl Into<String>) -> Result<Self, ZenohConfigError> {
let listen = listen.into();
validate_endpoint(&listen)?;
self.listen = Some(listen);
Ok(self)
}
/// Re-check every rule. Called by every constructor.
pub fn validate(&self) -> Result<(), ZenohConfigError> {
if self.endpoints.is_empty() {
return Err(ZenohConfigError::NoEndpoints);
}
for endpoint in &self.endpoints {
validate_endpoint(endpoint)?;
}
Ok(())
}
/// Build the Zenoh session config: the pinned posture plus this
/// deployment's endpoints, re-checked against [`check_posture`].
///
/// The posture check runs **after** the overrides, not before, because an
/// override is exactly how a pinned value gets lost.
pub fn session_config(&self) -> Result<::zenoh::Config, ZenohConfigError> {
self.validate()?;
let mut config = parse_pinned_posture()?;
config
.insert_json5("mode", &format!("\"{}\"", self.mode.as_str()))
.map_err(|e| ZenohConfigError::Invalid(e.to_string()))?;
config
.insert_json5("connect/endpoints", &endpoint_json(&self.endpoints))
.map_err(|e| ZenohConfigError::Invalid(e.to_string()))?;
if let Some(listen) = &self.listen {
config
.insert_json5(
"listen/endpoints",
&endpoint_json(std::slice::from_ref(listen)),
)
.map_err(|e| ZenohConfigError::Invalid(e.to_string()))?;
}
check_posture(&config).map_err(ZenohConfigError::Posture)?;
Ok(config)
}
}
/// Parse the embedded pinned posture into a Zenoh config.
pub fn parse_pinned_posture() -> Result<::zenoh::Config, ZenohConfigError> {
::zenoh::Config::from_json5(PINNED_PEER_CONFIG)
.map_err(|e| ZenohConfigError::Invalid(e.to_string()))
}
fn endpoint_json(endpoints: &[String]) -> String {
let quoted: Vec<String> = endpoints
.iter()
.map(|e| format!("\"{}\"", e.replace('\\', "\\\\").replace('"', "\\\"")))
.collect();
format!("[{}]", quoted.join(","))
}
/// Validate one endpoint against everything this build can honour.
fn validate_endpoint(endpoint: &str) -> Result<(), ZenohConfigError> {
let Some((protocol, address)) = endpoint.split_once('/') else {
return Err(ZenohConfigError::MalformedEndpoint(endpoint.to_string()));
};
if protocol != "tcp" {
return Err(ZenohConfigError::UnsupportedLinkProtocol(
endpoint.to_string(),
));
}
// Strip any locator metadata (`tcp/host:port?foo=bar#meta`).
let address = address
.split(['?', '#'])
.next()
.unwrap_or_default()
.to_string();
if address.is_empty() {
return Err(ZenohConfigError::MalformedEndpoint(endpoint.to_string()));
}
if !is_loopback(&address) {
return Err(ZenohConfigError::NonLocalEndpoint(endpoint.to_string()));
}
Ok(())
}
fn is_loopback(address: &str) -> bool {
let host = if let Some(rest) = address.strip_prefix('[') {
// `[::1]:7447`
match rest.split_once(']') {
Some((host, _)) => host.to_string(),
None => return false,
}
} else {
match address.rsplit_once(':') {
Some((host, _)) => host.to_string(),
None => address.to_string(),
}
};
if host.eq_ignore_ascii_case("localhost") {
return true;
}
host.parse::<std::net::IpAddr>()
.is_ok_and(|ip| ip.is_loopback() || ip.is_unspecified())
}
/// The pinned MERIDIAN posture, checked against a parsed Zenoh config.
///
/// Returns the **first** violation rather than panicking, so the same code path
/// serves the positive fixtures, the negative fixtures and the config this
/// crate builds at runtime. Every assertion goes through `get_json`, so a key
/// that Zenoh silently ignored reads back as the upstream default and fails
/// here.
///
/// Kept in one function on purpose: the peer session, the runtime-built config
/// and the `zenohd` router must not be allowed to drift into different
/// postures, and a shared check is the only thing that stops that happening one
/// reviewed diff at a time.
pub fn check_posture(config: &::zenoh::Config) -> Result<(), String> {
let json_at = |key: &str| -> Result<String, String> {
config
.get_json(key)
.map_err(|e| format!("config key `{key}` is not readable in zenoh 1.8.0: {e}"))
};
let expect = |key: &str, want: &str, why: &str| -> Result<(), String> {
let got = json_at(key)?;
if got == want {
Ok(())
} else {
Err(format!("{key} must be {want} but is {got} — {why}"))
}
};
// QoS on. Without it there is one transmission queue and the Q0-Q7 lane
// mapping is a no-op that no test downstream of here would notice.
expect(
"transport/unicast/qos/enabled",
"true",
"the eight MERIDIAN lanes need Zenoh's priority queues",
)?;
// Low-latency transport off. `qos && lowlatency` is rejected outright by
// zenoh-transport 1.8.0 at `zenoh::open`, not at parse.
expect(
"transport/unicast/lowlatency",
"false",
"incompatible with QoS in zenoh 1.8.0",
)?;
// Batching on. It is back-pressure-driven and is the throughput path for the
// small-payload traffic class; `express` opts a single publication out.
expect(
"transport/link/tx/queue/batching/enabled",
"true",
"adaptive batching is the small-payload throughput path",
)?;
// Shared memory off, both switches. Zenoh 1.x defaults BOTH to `true`
// (`zenoh-config-1.8.0/src/lib.rs:795-816`), so a config that merely omits
// them has switched on charter Phase 6.
expect(
"transport/shared_memory/enabled",
"false",
"SHM is charter Phase 6 and defaults ON — it must be refused explicitly",
)?;
expect(
"transport/shared_memory/transport_optimization/enabled",
"false",
"a second SHM switch that silently routes large messages through shared memory",
)?;
// Both scouting mechanisms off. Upstream defaults are `true` for BOTH.
// `false` and not "anything but true" on purpose: an omitted key reads back
// as `null` here and fails, which is the property that matters, because
// `zenohd/src/main.rs:203-214` reads an *unset* multicast key as consent.
expect(
"scouting/multicast/enabled",
"false",
"topology is declared, never discovered; an omitted key is read as consent by zenohd",
)?;
expect(
"scouting/gossip/enabled",
"false",
"gossip stays on when only multicast is disabled (Gap #17)",
)?;
// Belt and braces on the same gap: nothing may be autoconnected. Checked
// structurally rather than by string match — the value is mode-dependent.
for key in [
"scouting/multicast/autoconnect",
"scouting/gossip/autoconnect",
"scouting/gossip/target",
] {
let raw = json_at(key)?;
let parsed: serde_json::Value = serde_json::from_str(&raw)
.map_err(|e| format!("{key} did not read back as JSON ({raw}): {e}"))?;
let empty = match &parsed {
serde_json::Value::Null => true,
serde_json::Value::Array(items) => items.is_empty(),
serde_json::Value::Object(per_mode) => per_mode
.values()
.all(|v| v.as_array().is_some_and(|items| items.is_empty())),
_ => false,
};
if !empty {
return Err(format!(
"{key} must name no node type — implicit peering is refused, got {raw}"
));
}
}
// No plugin runtime. The pinned image ships a REST plugin and a
// storage-manager plugin in its root.
expect(
"plugins_loading/enabled",
"false",
"no plugin runtime, no REST surface, no second storage mechanism",
)?;
expect(
"plugins",
"{}",
"storage manager and REST are refused by name, not just left unloaded",
)?;
expect(
"plugins_loading/search_dirs",
"[]",
"zenohd re-enables plugin loading regardless of config; an empty search path is what survives",
)?;
// Admin space off, and its permissions refused — `zenohd` forces `enabled`
// back to `true`, so the permissions are the part that holds on the daemon.
expect(
"adminspace/enabled",
"false",
"the admin space is not part of the bus contract",
)?;
expect(
"adminspace/permissions/read",
"false",
"zenohd force-enables the admin space; read permission is what survives",
)?;
expect(
"adminspace/permissions/write",
"false",
"runtime config mutation over the admin space is refused outright",
)?;
// Only the link this build actually compiles.
expect(
"transport/link/protocols",
"[\"tcp\"]",
"this build is default-features=false + transport_tcp",
)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_embedded_posture_is_the_fixture_and_it_passes_its_own_check() {
let embedded = parse_pinned_posture().expect("the embedded fixture must parse");
check_posture(&embedded).expect("the embedded fixture must satisfy the pinned posture");
let from_file = ::zenoh::Config::from_file(
std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/zenoh/bus-peer.json5"),
)
.expect("the fixture on disk must parse");
assert_eq!(
serde_json::to_string(&embedded).expect("serialize"),
serde_json::to_string(&from_file).expect("serialize"),
"the embedded posture drifted from the fixture on disk"
);
}
#[test]
fn no_endpoint_is_refused() {
assert_eq!(
ZenohBusConfig::new(ZenohMode::Peer, Vec::<String>::new()).unwrap_err(),
ZenohConfigError::NoEndpoints
);
assert_eq!(
ZenohBusConfig::new(ZenohMode::Peer, ["", " "]).unwrap_err(),
ZenohConfigError::NoEndpoints
);
}
#[test]
fn an_unknown_mode_is_refused_and_router_is_not_a_mode() {
assert!(matches!(
ZenohMode::parse("router"),
Err(ZenohConfigError::UnknownMode(_))
));
assert!(matches!(
ZenohMode::parse("Peer"),
Err(ZenohConfigError::UnknownMode(_))
));
assert_eq!(ZenohMode::parse("peer").expect("peer"), ZenohMode::Peer);
assert_eq!(
ZenohMode::parse(" client ").expect("client"),
ZenohMode::Client
);
}
#[test]
fn a_link_this_build_did_not_compile_is_refused_before_open() {
for endpoint in [
"tls/127.0.0.1:7447",
"quic/127.0.0.1:7447",
"udp/127.0.0.1:7447",
] {
assert!(
matches!(
ZenohBusConfig::new(ZenohMode::Peer, [endpoint]),
Err(ZenohConfigError::UnsupportedLinkProtocol(_))
),
"{endpoint} must be refused at config time, not at open"
);
}
}
#[test]
fn a_malformed_endpoint_is_refused() {
for endpoint in ["127.0.0.1:7447", "tcp/"] {
assert!(
matches!(
ZenohBusConfig::new(ZenohMode::Peer, [endpoint]),
Err(ZenohConfigError::MalformedEndpoint(_))
),
"{endpoint} must be refused as a non-locator"
);
}
}
/// The charter's dashed red edge, as a config refusal.
#[test]
fn a_non_loopback_endpoint_is_refused_until_the_reverify_branch_lands() {
for endpoint in [
"tcp/10.0.0.5:7447",
"tcp/router.meridian.example:7447",
"tcp/[2001:db8::1]:7447",
] {
assert!(
matches!(
ZenohBusConfig::new(ZenohMode::Peer, [endpoint]),
Err(ZenohConfigError::NonLocalEndpoint(_))
),
"{endpoint} must be refused while `fan_out_pubsub_event` cannot re-verify"
);
}
for endpoint in [
"tcp/127.0.0.1:7447",
"tcp/localhost:7447",
"tcp/[::1]:7447",
"tcp/0.0.0.0:0",
] {
assert!(
ZenohBusConfig::new(ZenohMode::Peer, [endpoint]).is_ok(),
"{endpoint} is loopback and must be accepted"
);
}
}
#[test]
fn the_built_session_config_keeps_the_pinned_posture_and_takes_the_endpoints() {
let config = ZenohBusConfig::new(ZenohMode::Peer, ["tcp/127.0.0.1:7451"])
.expect("loopback endpoint")
.with_listen("tcp/127.0.0.1:0")
.expect("loopback listen");
let session = config.session_config().expect("session config");
check_posture(&session).expect("the built config must satisfy the pinned posture");
assert_eq!(
session.get_json("mode").expect("mode is readable"),
"\"peer\""
);
let connect = session
.get_json("connect/endpoints")
.expect("endpoints are readable");
assert!(
connect.contains("127.0.0.1:7451"),
"the configured endpoint did not survive: {connect}"
);
let listen = session
.get_json("listen/endpoints")
.expect("listen is readable");
assert!(
listen.contains("127.0.0.1"),
"listen override lost: {listen}"
);
}
#[test]
fn client_mode_survives_the_override() {
let session = ZenohBusConfig::new(ZenohMode::Client, ["tcp/127.0.0.1:7452"])
.expect("loopback endpoint")
.session_config()
.expect("session config");
assert_eq!(
session.get_json("mode").expect("mode is readable"),
"\"client\""
);
check_posture(&session).expect("mode must not disturb the posture");
}
/// The posture check is not decoration: a config that loses a pinned value
/// must be rejected by the same function the positive paths pass.
#[test]
fn the_posture_check_rejects_a_lost_pin() {
let mut config = parse_pinned_posture().expect("parse");
config
.insert_json5("scouting/gossip/enabled", "true")
.expect("insert");
let violation = check_posture(&config).expect_err("gossip on must be refused");
assert!(
violation.contains("scouting/gossip/enabled"),
"expected gossip to be the named violation, got: {violation}"
);
}
#[test]
fn defaults_are_bounded_rather_than_unlimited() {
let config =
ZenohBusConfig::new(ZenohMode::Peer, ["tcp/127.0.0.1:7447"]).expect("loopback");
assert!(config.publisher_cache_capacity > 0);
assert!(config.local_echo_capacity > 0);
assert_eq!(config.link_profile, TrustProfile::P0None);
assert_eq!(config.lane_floors, LaneFloors::IN_DEPLOYMENT);
}
}

View file

@ -0,0 +1,384 @@
//! Zenoh bus readiness, and the metrics that make it answerable.
//!
//! # A live process is not a ready bus
//!
//! The pinned posture opens a session even when the router tier is absent
//! (`connect.exit_on_failure: false`, both `open.return_conditions` off), on
//! purpose: the relay must serve NIP-01 whether or not the bus is reachable.
//! That is exactly why "the session opened" cannot be readiness. Three separate
//! things must be true, and each has been the whole failure on its own:
//!
//! 1. **Configured connects.** With scouting off, a session with no transport
//! reaches nothing and reports nothing wrong.
//! 2. **Required declarations.** The control lanes are declared subscribers. A
//! session with a live transport and no declaration receives no cache drops
//! and no ban enforcement, silently.
//! 3. **A fresh application loopback.** A declaration can exist while nothing
//! moves through it. The probe publishes on this session's own health key
//! and requires the sample back, so readiness is evidence of delivery rather
//! than evidence of configuration.
//!
//! [`ZenohHealth::readiness_gap`] names the *first* unmet condition rather than
//! answering a bare `false` — explicit denial, never a silent empty result.
use std::sync::RwLock;
use std::time::{Duration, Instant};
/// Point-in-time readiness of the Zenoh bus.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct ZenohHealthSnapshot {
/// Whether `zenoh::open` returned and the session has not been closed.
pub session_open: bool,
/// How many endpoints the configuration named.
pub endpoints_configured: usize,
/// How many transports the session currently has.
pub endpoints_connected: usize,
/// How many declarations must exist before this lane set can serve.
pub declarations_required: usize,
/// How many currently exist.
pub declarations_active: usize,
/// Wall clock of the last successful loopback, for operator diagnostics.
pub last_loopback_unix: Option<u64>,
/// Monotonic instant of the last successful loopback, for the freshness
/// budget. Wall clock is not used for the comparison because it can move.
pub last_loopback_at: Option<Instant>,
}
impl ZenohHealthSnapshot {
/// The first unmet readiness condition, or `None` when the bus is ready.
#[must_use]
pub fn readiness_gap(&self, freshness: Duration) -> Option<ReadinessGap> {
if !self.session_open {
return Some(ReadinessGap::SessionClosed);
}
if self.endpoints_connected < 1 && self.endpoints_configured > 0 {
return Some(ReadinessGap::NoTransport);
}
if self.declarations_active < self.declarations_required {
return Some(ReadinessGap::MissingDeclarations);
}
match self.last_loopback_at {
None => Some(ReadinessGap::NoLoopback),
Some(at) if at.elapsed() > freshness => Some(ReadinessGap::StaleLoopback),
Some(_) => None,
}
}
/// Whether every readiness condition holds.
#[must_use]
pub fn is_ready(&self, freshness: Duration) -> bool {
self.readiness_gap(freshness).is_none()
}
}
/// The first readiness condition that does not hold.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReadinessGap {
/// `zenoh::open` has not returned, or the session has been closed.
SessionClosed,
/// Endpoints were configured and none is connected. With both scouting
/// mechanisms off there is no other way to acquire one.
NoTransport,
/// A required declaration is missing, so a lane exists that nothing feeds.
MissingDeclarations,
/// No loopback sample has ever returned on this session.
NoLoopback,
/// The last loopback is older than the configured freshness budget.
StaleLoopback,
}
impl ReadinessGap {
/// Metric-safe label.
#[must_use]
pub const fn label(self) -> &'static str {
match self {
ReadinessGap::SessionClosed => "session_closed",
ReadinessGap::NoTransport => "no_transport",
ReadinessGap::MissingDeclarations => "missing_declarations",
ReadinessGap::NoLoopback => "no_loopback",
ReadinessGap::StaleLoopback => "stale_loopback",
}
}
/// Operator-facing explanation.
#[must_use]
pub const fn describe(self) -> &'static str {
match self {
ReadinessGap::SessionClosed => "the Zenoh session is not open",
ReadinessGap::NoTransport => {
"endpoints are configured but no transport is established; \
scouting is off, so there is no other way to acquire one"
}
ReadinessGap::MissingDeclarations => {
"a required declared subscriber is missing; a lane would silently receive nothing"
}
ReadinessGap::NoLoopback => "no loopback sample has returned on this session yet",
ReadinessGap::StaleLoopback => {
"the last loopback sample is older than the freshness budget"
}
}
}
}
impl std::fmt::Display for ReadinessGap {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.describe())
}
}
/// Shared readiness registry for one Zenoh bus.
///
/// Deliberately not [`crate::SubscriberHealth`]: that type's lane set is the
/// Redis subscriber topology (three reconnecting loops with subscription
/// generations), and a declared-subscriber transport has neither generations
/// nor a reconnect loop of its own. Reusing it would have forced fake
/// generations onto a readiness question that is actually about declarations
/// and loopback freshness.
#[derive(Debug)]
pub struct ZenohHealth {
inner: RwLock<ZenohHealthSnapshot>,
freshness: Duration,
}
impl ZenohHealth {
/// Build a registry for a session with `endpoints_configured` endpoints and
/// `declarations_required` mandatory declarations.
#[must_use]
pub fn new(
endpoints_configured: usize,
declarations_required: usize,
freshness: Duration,
) -> Self {
let health = Self {
inner: RwLock::new(ZenohHealthSnapshot {
endpoints_configured,
declarations_required,
..ZenohHealthSnapshot::default()
}),
freshness,
};
health.publish();
health
}
/// The freshness budget this registry compares the loopback against.
#[must_use]
pub const fn freshness(&self) -> Duration {
self.freshness
}
fn update(&self, update: impl FnOnce(&mut ZenohHealthSnapshot)) {
if let Ok(mut snapshot) = self.inner.write() {
update(&mut snapshot);
}
self.publish();
}
/// Record whether the session is open.
pub fn set_session_open(&self, open: bool) {
self.update(|snapshot| snapshot.session_open = open);
}
/// Record how many transports the session currently has.
pub fn set_endpoints_connected(&self, connected: usize) {
self.update(|snapshot| snapshot.endpoints_connected = connected);
}
/// Record how many declarations currently exist.
pub fn set_declarations_active(&self, active: usize) {
self.update(|snapshot| snapshot.declarations_active = active);
}
/// Record a loopback sample returning on this session.
pub fn mark_loopback(&self) {
let now = Instant::now();
let unix = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.ok()
.map(|d| d.as_secs());
self.update(|snapshot| {
snapshot.last_loopback_at = Some(now);
snapshot.last_loopback_unix = unix;
});
}
/// Current snapshot.
#[must_use]
pub fn snapshot(&self) -> ZenohHealthSnapshot {
self.inner
.read()
.map(|snapshot| *snapshot)
.unwrap_or_default()
}
/// The first unmet readiness condition, or `None`.
#[must_use]
pub fn readiness_gap(&self) -> Option<ReadinessGap> {
self.snapshot().readiness_gap(self.freshness)
}
/// Whether the bus is ready to serve.
#[must_use]
pub fn is_ready(&self) -> bool {
self.readiness_gap().is_none()
}
/// Publish the current snapshot to the metrics registry.
///
/// Every name below is a string literal at the `metrics::*!` call site, as
/// `crates/meridian-relay/AGENTS.md` requires: the alert→metric guard
/// resolves declarations by reading that literal, so a computed name drops
/// the metric out of the guard.
fn publish(&self) {
let snapshot = self.snapshot();
metrics::gauge!("meridian_zenoh_session_open").set(if snapshot.session_open {
1.0
} else {
0.0
});
metrics::gauge!("meridian_zenoh_endpoints_configured")
.set(snapshot.endpoints_configured as f64);
metrics::gauge!("meridian_zenoh_endpoints_connected")
.set(snapshot.endpoints_connected as f64);
metrics::gauge!("meridian_zenoh_declarations_required")
.set(snapshot.declarations_required as f64);
metrics::gauge!("meridian_zenoh_declarations_active")
.set(snapshot.declarations_active as f64);
if let Some(unix) = snapshot.last_loopback_unix {
metrics::gauge!("meridian_zenoh_last_loopback_unix").set(unix as f64);
}
let gap = snapshot.readiness_gap(self.freshness);
metrics::gauge!("meridian_zenoh_ready").set(if gap.is_none() { 1.0 } else { 0.0 });
// The gap is a *labelled* gauge rather than a family of booleans so an
// operator reads one series and gets the reason, not five series and a
// deduction.
for candidate in [
ReadinessGap::SessionClosed,
ReadinessGap::NoTransport,
ReadinessGap::MissingDeclarations,
ReadinessGap::NoLoopback,
ReadinessGap::StaleLoopback,
] {
metrics::gauge!("meridian_zenoh_readiness_gap", "gap" => candidate.label())
.set(if gap == Some(candidate) { 1.0 } else { 0.0 });
}
}
}
#[cfg(test)]
mod tests {
use super::*;
const FRESHNESS: Duration = Duration::from_secs(60);
fn ready_snapshot() -> ZenohHealthSnapshot {
ZenohHealthSnapshot {
session_open: true,
endpoints_configured: 1,
endpoints_connected: 1,
declarations_required: 3,
declarations_active: 3,
last_loopback_unix: Some(1),
last_loopback_at: Some(Instant::now()),
}
}
#[test]
fn a_live_session_alone_is_not_ready() {
let snapshot = ZenohHealthSnapshot {
session_open: true,
endpoints_configured: 1,
..ZenohHealthSnapshot::default()
};
assert_eq!(
snapshot.readiness_gap(FRESHNESS),
Some(ReadinessGap::NoTransport)
);
assert!(!snapshot.is_ready(FRESHNESS));
}
#[test]
fn a_connected_session_with_a_missing_declaration_is_not_ready() {
let snapshot = ZenohHealthSnapshot {
declarations_active: 2,
..ready_snapshot()
};
assert_eq!(
snapshot.readiness_gap(FRESHNESS),
Some(ReadinessGap::MissingDeclarations)
);
}
#[test]
fn declarations_without_a_loopback_are_not_ready() {
let snapshot = ZenohHealthSnapshot {
last_loopback_at: None,
last_loopback_unix: None,
..ready_snapshot()
};
assert_eq!(
snapshot.readiness_gap(FRESHNESS),
Some(ReadinessGap::NoLoopback)
);
}
#[test]
fn a_stale_loopback_fails_readiness() {
let snapshot = ready_snapshot();
assert!(snapshot.is_ready(FRESHNESS));
assert_eq!(
snapshot.readiness_gap(Duration::ZERO),
Some(ReadinessGap::StaleLoopback),
"a zero budget must make any past loopback stale"
);
}
#[test]
fn the_gap_order_reports_the_root_cause_first() {
let snapshot = ZenohHealthSnapshot::default();
assert_eq!(
snapshot.readiness_gap(FRESHNESS),
Some(ReadinessGap::SessionClosed),
"a closed session must not be reported as a missing declaration"
);
}
#[test]
fn the_registry_tracks_each_condition_independently() {
let health = ZenohHealth::new(1, 3, FRESHNESS);
assert_eq!(health.readiness_gap(), Some(ReadinessGap::SessionClosed));
health.set_session_open(true);
assert_eq!(health.readiness_gap(), Some(ReadinessGap::NoTransport));
health.set_endpoints_connected(1);
assert_eq!(
health.readiness_gap(),
Some(ReadinessGap::MissingDeclarations)
);
health.set_declarations_active(3);
assert_eq!(health.readiness_gap(), Some(ReadinessGap::NoLoopback));
health.mark_loopback();
assert_eq!(health.readiness_gap(), None);
assert!(health.is_ready());
health.set_session_open(false);
assert!(!health.is_ready(), "readiness must be revocable");
}
#[test]
fn every_gap_has_a_distinct_label_and_description() {
let gaps = [
ReadinessGap::SessionClosed,
ReadinessGap::NoTransport,
ReadinessGap::MissingDeclarations,
ReadinessGap::NoLoopback,
ReadinessGap::StaleLoopback,
];
let mut labels: Vec<&str> = gaps.iter().map(|g| g.label()).collect();
labels.sort_unstable();
labels.dedup();
assert_eq!(labels.len(), gaps.len(), "gap labels must be distinct");
for gap in gaps {
assert!(!gap.describe().is_empty());
}
}
}

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,501 @@
//! Shadow mode — publish to both, serve from Redis, count the difference.
//!
//! This is the reversible-rollout mechanism the whole transport plan rests on.
//! Every publish goes to Redis **and** Zenoh; every delivery comes from Redis.
//! A Zenoh failure is therefore never a request failure, and a Zenoh success is
//! never evidence of anything until the parity counters say it happened every
//! time for a full deploy cycle.
//!
//! # What "authoritative" means here, precisely
//!
//! - [`ShadowEventBus::publish_event`] returns **Redis's** result. If Redis
//! succeeded and Zenoh failed, the caller sees `Ok`. If Redis failed, the
//! caller sees the Redis error even when Zenoh succeeded — a shadow that
//! could rescue a primary failure would make the primary's health
//! unobservable, which is the one thing a soak must not do.
//! - The three `subscribe_*` methods return the **Redis** receivers. Nothing
//! Zenoh delivers is ever served.
//!
//! # What the parity counters do and do not prove
//!
//! [`ParitySnapshot`] is **publish** parity: for each lane, how often both
//! backends accepted, and how often exactly one did. That is what gates the
//! `Redis → Shadow → Bridge` transition against a *publisher-side* regression.
//!
//! It is **not** delivery parity. Proving that every event published on Redis
//! also arrived over Zenoh needs the correlation probe from the rollout
//! contract — one event, one cache invalidation and one connection-control
//! command per step, with zero missing ids on both transports — and that probe
//! is not in this crate. [`ShadowEventBus::run_shadow_receipt_observer`] is the
//! honest half of it: it drains the Zenoh lanes (which nothing else consumes,
//! so it changes no delivery semantics) and counts receipts, so a shadow lane
//! that receives *nothing* is visible immediately rather than at cutover.
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use meridian_core::TenantContext;
use tokio::sync::broadcast;
use crate::bus::{BusError, EventBus, RedisEventBus};
use crate::cache_invalidation::{CacheInvalidation, ScopedCacheInvalidation};
use crate::conn_control::{ConnControl, ScopedConnControl};
use crate::zenoh::ZenohEventBus;
use crate::{ChannelEvent, EventTopic};
/// Which lane a parity observation belongs to.
///
/// The same three lanes as [`crate::SubscriberLane`], and the labels match, so
/// a dashboard can join the two without a translation table.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ParityLane {
/// Cross-pod event fan-out.
Events,
/// Membership/visibility cache invalidation.
CacheInvalidation,
/// Live connection-control commands.
ConnectionControl,
}
impl ParityLane {
/// Every lane.
pub const ALL: [ParityLane; 3] = [
ParityLane::Events,
ParityLane::CacheInvalidation,
ParityLane::ConnectionControl,
];
const fn index(self) -> usize {
match self {
ParityLane::Events => 0,
ParityLane::CacheInvalidation => 1,
ParityLane::ConnectionControl => 2,
}
}
/// Metric-safe label, identical to [`crate::SubscriberLane::label`].
#[must_use]
pub const fn label(self) -> &'static str {
match self {
ParityLane::Events => "events",
ParityLane::CacheInvalidation => "cache_invalidation",
ParityLane::ConnectionControl => "connection_control",
}
}
}
/// Publish-parity counts for one lane.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct LaneParity {
/// Both backends accepted the publication.
pub both: u64,
/// Redis accepted and Zenoh did not. The soak's primary signal.
pub redis_only: u64,
/// Zenoh accepted and Redis did not. Serving is unaffected — the caller
/// still sees the Redis error — but it means the primary is degraded.
pub zenoh_only: u64,
/// Neither accepted.
pub neither: u64,
/// Messages observed arriving on the Zenoh side of this lane.
pub shadow_received: u64,
}
impl LaneParity {
/// Publications where exactly one backend accepted.
#[must_use]
pub const fn divergent(&self) -> u64 {
self.redis_only + self.zenoh_only
}
}
/// Point-in-time parity across all three lanes.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct ParitySnapshot {
/// Per-lane counts, indexed by [`ParityLane`].
lanes: [LaneParity; 3],
}
impl ParitySnapshot {
/// Counts for one lane.
#[must_use]
pub const fn lane(&self, lane: ParityLane) -> LaneParity {
self.lanes[lane.index()]
}
/// Whether every lane saw identical publish outcomes.
///
/// This is the gate for `Shadow → Bridge`, and it is deliberately strict:
/// one divergent publication anywhere is a miss.
#[must_use]
pub fn is_clean(&self) -> bool {
self.lanes.iter().all(|lane| lane.divergent() == 0)
}
}
#[derive(Debug, Default)]
struct LaneCounters {
both: AtomicU64,
redis_only: AtomicU64,
zenoh_only: AtomicU64,
neither: AtomicU64,
shadow_received: AtomicU64,
}
impl LaneCounters {
fn snapshot(&self) -> LaneParity {
LaneParity {
both: self.both.load(Ordering::Relaxed),
redis_only: self.redis_only.load(Ordering::Relaxed),
zenoh_only: self.zenoh_only.load(Ordering::Relaxed),
neither: self.neither.load(Ordering::Relaxed),
shadow_received: self.shadow_received.load(Ordering::Relaxed),
}
}
}
/// Shared parity registry. Cloneable so a probe or an observer can hold one.
#[derive(Debug, Default)]
pub struct ParityCounters {
lanes: [LaneCounters; 3],
}
impl ParityCounters {
/// Record one publish outcome pair.
pub fn record(&self, lane: ParityLane, redis_ok: bool, zenoh_ok: bool) {
let counters = &self.lanes[lane.index()];
let (counter, outcome) = match (redis_ok, zenoh_ok) {
(true, true) => (&counters.both, "both"),
(true, false) => (&counters.redis_only, "redis_only"),
(false, true) => (&counters.zenoh_only, "zenoh_only"),
(false, false) => (&counters.neither, "neither"),
};
counter.fetch_add(1, Ordering::Relaxed);
metrics::counter!(
"meridian_bus_shadow_publish_total",
"lane" => lane.label(),
"outcome" => outcome,
)
.increment(1);
if redis_ok != zenoh_ok {
metrics::counter!(
"meridian_bus_shadow_divergence_total",
"lane" => lane.label(),
"outcome" => outcome,
)
.increment(1);
}
}
/// Record one message observed on the Zenoh side of a lane.
pub fn record_shadow_receipt(&self, lane: ParityLane) {
self.lanes[lane.index()]
.shadow_received
.fetch_add(1, Ordering::Relaxed);
metrics::counter!("meridian_bus_shadow_received_total", "lane" => lane.label())
.increment(1);
}
/// Current counts.
#[must_use]
pub fn snapshot(&self) -> ParitySnapshot {
ParitySnapshot {
lanes: [
self.lanes[0].snapshot(),
self.lanes[1].snapshot(),
self.lanes[2].snapshot(),
],
}
}
}
/// Dual-publish, Redis-serve, count-the-difference composition of the two
/// backends.
///
/// Composed as a concrete struct rather than `Box<dyn EventBus>`: the seam is
/// RPITIT and not dyn-compatible, and a boxed future per call would put a heap
/// allocation on the publish path of a bus built to remove per-event cost.
pub struct ShadowEventBus {
primary: RedisEventBus,
shadow: ZenohEventBus,
parity: Arc<ParityCounters>,
}
impl ShadowEventBus {
/// Compose an already-constructed Redis and Zenoh bus.
#[must_use]
pub fn new(primary: RedisEventBus, shadow: ZenohEventBus) -> Self {
Self {
primary,
shadow,
parity: Arc::new(ParityCounters::default()),
}
}
/// The authoritative backend. Serving, presence and the three Redis
/// subscriber loops all belong to it.
#[must_use]
pub fn primary(&self) -> &RedisEventBus {
&self.primary
}
/// The observed backend. Never serves.
#[must_use]
pub fn shadow(&self) -> &ZenohEventBus {
&self.shadow
}
/// Shared parity registry.
#[must_use]
pub fn parity(&self) -> Arc<ParityCounters> {
Arc::clone(&self.parity)
}
/// Current parity counts.
#[must_use]
pub fn parity_snapshot(&self) -> ParitySnapshot {
self.parity.snapshot()
}
/// Drain the Zenoh lanes and count receipts. Runs forever — spawn it.
///
/// This is an **observer**, not a delivery runner. It forwards nothing and
/// makes no readiness transition, so it is not a second copy of
/// [`crate::broadcast_consumer::run_broadcast_consumer`], which remains the
/// only place a lane's lag is interpreted and compensated. It exists
/// because in shadow mode nothing else reads the Zenoh receivers: without
/// it those channels fill, every shadow lane lags permanently, and the
/// soak measures a backlog rather than the transport.
///
/// A lagged receiver here means *observation* was lost, never delivery, so
/// it is counted as `lagged` and the loop continues.
pub async fn run_shadow_receipt_observer(&self) {
let mut events = EventBus::subscribe_local(&self.shadow);
let mut cache = EventBus::subscribe_cache_invalidations(&self.shadow);
let mut control = EventBus::subscribe_conn_control(&self.shadow);
loop {
tokio::select! {
received = events.recv() => {
if !observe(received, &self.parity, ParityLane::Events) {
return;
}
}
received = cache.recv() => {
if !observe(received, &self.parity, ParityLane::CacheInvalidation) {
return;
}
}
received = control.recv() => {
if !observe(received, &self.parity, ParityLane::ConnectionControl) {
return;
}
}
}
}
}
}
/// Returns `false` when the lane's channel closed and the observer must stop.
fn observe<T>(
received: Result<T, broadcast::error::RecvError>,
parity: &ParityCounters,
lane: ParityLane,
) -> bool {
match received {
Ok(_) => {
parity.record_shadow_receipt(lane);
true
}
Err(broadcast::error::RecvError::Lagged(skipped)) => {
metrics::counter!("meridian_bus_shadow_observer_lagged_total", "lane" => lane.label())
.increment(skipped);
true
}
Err(broadcast::error::RecvError::Closed) => false,
}
}
impl EventBus for ShadowEventBus {
async fn publish_event(
&self,
ctx: &TenantContext,
topic: EventTopic,
event: &nostr::Event,
) -> Result<(), BusError> {
let primary = self.primary.publish_event(ctx, topic, event).await;
let shadow = self.shadow.publish_event(ctx, topic, event).await;
self.record(ParityLane::Events, &primary, shadow);
primary
}
fn subscribe_local(&self) -> broadcast::Receiver<ChannelEvent> {
EventBus::subscribe_local(&self.primary)
}
async fn retain_topic(&self, ctx: &TenantContext, topic: EventTopic) {
// Both sides retain: the shadow must be receiving the same topics, or
// its receipt counts measure the subscription rather than the
// transport.
self.primary.retain_topic(ctx, topic).await;
self.shadow.retain_topic(ctx, topic).await;
}
async fn release_topic(&self, ctx: &TenantContext, topic: EventTopic) {
self.primary.release_topic(ctx, topic).await;
self.shadow.release_topic(ctx, topic).await;
}
async fn publish_cache_invalidation(
&self,
ctx: &TenantContext,
invalidation: &CacheInvalidation,
) -> Result<(), BusError> {
let primary = self
.primary
.publish_cache_invalidation(ctx, invalidation)
.await;
let shadow = self
.shadow
.publish_cache_invalidation(ctx, invalidation)
.await;
self.record(ParityLane::CacheInvalidation, &primary, shadow);
primary
}
async fn publish_conn_control(
&self,
ctx: &TenantContext,
command: &ConnControl,
) -> Result<(), BusError> {
let primary = self.primary.publish_conn_control(ctx, command).await;
let shadow = self.shadow.publish_conn_control(ctx, command).await;
self.record(ParityLane::ConnectionControl, &primary, shadow);
primary
}
fn subscribe_cache_invalidations(&self) -> broadcast::Receiver<ScopedCacheInvalidation> {
EventBus::subscribe_cache_invalidations(&self.primary)
}
fn subscribe_conn_control(&self) -> broadcast::Receiver<ScopedConnControl> {
EventBus::subscribe_conn_control(&self.primary)
}
}
impl ShadowEventBus {
fn record(
&self,
lane: ParityLane,
primary: &Result<(), BusError>,
shadow: Result<(), BusError>,
) {
let redis_ok = primary.is_ok();
let zenoh_ok = shadow.is_ok();
if let Err(e) = shadow {
// Logged, never returned: the shadow is observed, not served.
tracing::warn!(
lane = lane.label(),
error = %e,
"shadow bus publication failed; Redis remains authoritative"
);
}
self.parity.record(lane, redis_ok, zenoh_ok);
}
}
/// Compile-time assertion that the composed backend still satisfies the seam.
const _: fn() = || {
fn assert_event_bus<T: EventBus>() {}
assert_event_bus::<ShadowEventBus>();
};
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parity_starts_clean_and_a_single_divergence_dirties_it() {
let counters = ParityCounters::default();
assert!(counters.snapshot().is_clean());
counters.record(ParityLane::Events, true, true);
assert!(counters.snapshot().is_clean());
assert_eq!(counters.snapshot().lane(ParityLane::Events).both, 1);
counters.record(ParityLane::Events, true, false);
let snapshot = counters.snapshot();
assert!(!snapshot.is_clean(), "one miss must fail the gate");
assert_eq!(snapshot.lane(ParityLane::Events).redis_only, 1);
assert_eq!(snapshot.lane(ParityLane::Events).divergent(), 1);
assert!(
snapshot.lane(ParityLane::CacheInvalidation).divergent() == 0,
"a miss on one lane must not be attributed to another"
);
}
#[test]
fn a_primary_failure_is_recorded_even_when_the_shadow_succeeded() {
let counters = ParityCounters::default();
counters.record(ParityLane::ConnectionControl, false, true);
let snapshot = counters.snapshot();
assert_eq!(snapshot.lane(ParityLane::ConnectionControl).zenoh_only, 1);
assert!(!snapshot.is_clean());
}
#[test]
fn both_failing_is_not_counted_as_divergence() {
let counters = ParityCounters::default();
counters.record(ParityLane::CacheInvalidation, false, false);
let snapshot = counters.snapshot();
assert_eq!(snapshot.lane(ParityLane::CacheInvalidation).neither, 1);
assert_eq!(snapshot.lane(ParityLane::CacheInvalidation).divergent(), 0);
assert!(
snapshot.is_clean(),
"a shared outage is not a parity miss; it is an outage"
);
}
#[test]
fn shadow_receipts_are_counted_per_lane() {
let counters = ParityCounters::default();
counters.record_shadow_receipt(ParityLane::Events);
counters.record_shadow_receipt(ParityLane::Events);
assert_eq!(
counters.snapshot().lane(ParityLane::Events).shadow_received,
2
);
assert_eq!(
counters.snapshot().lane(ParityLane::Events).divergent(),
0,
"a receipt is an observation, not a publish outcome"
);
}
#[test]
fn parity_lane_labels_match_the_subscriber_lane_labels() {
assert_eq!(
ParityLane::ALL.map(|lane| lane.label()),
crate::SubscriberLane::ALL.map(|lane| lane.label()),
"a dashboard joins these two by label; they must not drift"
);
}
#[test]
fn a_lagged_observation_is_not_a_closed_lane() {
let counters = ParityCounters::default();
assert!(observe::<()>(
Err(broadcast::error::RecvError::Lagged(3)),
&counters,
ParityLane::Events
));
assert!(!observe::<()>(
Err(broadcast::error::RecvError::Closed),
&counters,
ParityLane::Events
));
assert_eq!(
counters.snapshot().lane(ParityLane::Events).shadow_received,
0,
"lag must not be counted as a receipt"
);
}
}

View file

@ -0,0 +1,613 @@
//! The MERIDIAN bus key-expression space, and nothing else.
//!
//! ```text
//! meridian/v1/{community_uuid}/ch/{channel_uuid}/k/{kind}
//! meridian/v1/{community_uuid}/global/k/{kind}
//! meridian/v1/{community_uuid}/ctl/{conn,cache,fence}
//! ```
//!
//! # Three properties this module exists to hold
//!
//! **The separator is `/` and the wildcards are `*` and `**`.** Zenoh has no
//! MQTT-style `+`; a `+` chunk is a perfectly legal *literal* chunk that
//! matches only the string `"+"`, so a key expression written with MQTT habits
//! parses, declares, subscribes, and then silently receives nothing. A
//! reference project shipped exactly that. [`BusKey::parse`] therefore names
//! `+` and `#` as their own error rather than letting them fall through to
//! "not a UUID".
//!
//! **`v1` is a wire fence, not a comment.** A subscriber declared on
//! `meridian/v1/**` must never see a `v2` publication, which is what makes a
//! future incompatible payload encoding a coexisting rollout rather than a
//! flag day. [`super::topic`]'s tests assert the non-intersection against
//! Zenoh's own matcher, not against string equality.
//!
//! **Rendering is canonical, and parsing accepts only the canonical form.**
//! `Uuid::parse_str` accepts braced, URN and unhyphenated spellings, and
//! `u32::from_str` accepts a leading `+`. All of those produce a *different*
//! key expression for the same logical topic, so a publisher and a subscriber
//! that disagree about spelling silently stop matching. Every chunk parser here
//! re-renders and compares.
//!
//! Key expressions remain a **routing label, never an isolation boundary** —
//! the same law as `crate::topic`. `filter_fanout_by_access` is the enforcement
//! point and is unchanged by this module.
use std::fmt;
use meridian_core::CommunityId;
use thiserror::Error;
use uuid::Uuid;
use crate::EventTopic;
/// First chunk of every MERIDIAN key expression.
pub const ROOT: &str = "meridian";
/// The wire-fence version chunk.
///
/// A payload encoding change that a `v1` subscriber could misread takes `v2`,
/// and the two coexist for the length of a rollout because a `v1/**`
/// declaration does not intersect a `v2` key.
pub const VERSION: &str = "v1";
/// Chunk that scopes exact-channel events.
pub const CHANNEL_SCOPE: &str = "ch";
/// Chunk that scopes community-global events.
pub const GLOBAL_SCOPE: &str = "global";
/// Chunk that introduces the kind component.
pub const KIND_SCOPE: &str = "k";
/// Chunk that scopes the non-event control lanes.
pub const CONTROL_SCOPE: &str = "ctl";
/// Chunk that scopes a session's own liveness probe.
///
/// Deliberately starts with `_`, which no canonical UUID can, so a health key
/// can never be confused with a community and [`BusKey::parse`] rejects it
/// without a special case.
pub const HEALTH_SCOPE: &str = "_health";
/// Zenoh's single-chunk wildcard.
pub const CHUNK_WILDCARD: &str = "*";
/// Zenoh's any-chunks wildcard.
pub const TREE_WILDCARD: &str = "**";
/// A non-event control lane.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum ControlTopic {
/// Connection-control commands (`crate::conn_control::ConnControl`).
Conn,
/// Cache-key drops (`crate::cache_invalidation::CacheInvalidation`).
Cache,
/// Fence/generation changes. Declared here so the space is complete; the
/// relay slice wires the producer.
Fence,
}
impl ControlTopic {
/// Every control lane.
pub const ALL: [ControlTopic; 3] =
[ControlTopic::Conn, ControlTopic::Cache, ControlTopic::Fence];
/// The key-expression chunk for this lane.
#[must_use]
pub const fn chunk(self) -> &'static str {
match self {
ControlTopic::Conn => "conn",
ControlTopic::Cache => "cache",
ControlTopic::Fence => "fence",
}
}
/// Parse a control chunk. Unknown chunks are `None`, never a default.
#[must_use]
pub fn from_chunk(chunk: &str) -> Option<ControlTopic> {
match chunk {
"conn" => Some(ControlTopic::Conn),
"cache" => Some(ControlTopic::Cache),
"fence" => Some(ControlTopic::Fence),
_ => None,
}
}
}
impl fmt::Display for ControlTopic {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.chunk())
}
}
/// Why a string is not a MERIDIAN bus key expression.
///
/// Every variant names one confusion that produced a silently-matchless
/// subscription somewhere. `Malformed` is the catch-all and is deliberately the
/// least useful, so a recognised mistake gets its own name instead.
#[derive(Debug, Clone, PartialEq, Eq, Error)]
pub enum KeyExprError {
/// An MQTT wildcard (`+` or `#`) appeared in a Zenoh key expression. Both
/// are legal literal characters to Zenoh and match nothing.
#[error("`{0}` uses an MQTT wildcard; Zenoh's wildcards are `*` (one chunk) and `**` (any)")]
MqttWildcard(String),
/// The key does not start with `meridian/v1`.
#[error("`{0}` is not under `{ROOT}/{VERSION}`")]
WrongRoot(String),
/// The key is under `meridian` but names a different wire version.
#[error("`{found}` is wire version `{found_version}`, not `{VERSION}`")]
WrongVersion {
/// The full key that was rejected.
found: String,
/// The version chunk that was found.
found_version: String,
},
/// A UUID chunk was absent, unparseable, or not in canonical lowercase
/// hyphenated form.
#[error("`{0}` is not a canonical lowercase hyphenated UUID chunk")]
NonCanonicalUuid(String),
/// A kind chunk was absent, non-numeric, out of `u32` range, or carried a
/// sign or a leading zero.
#[error("`{0}` is not a canonical decimal kind chunk")]
NonCanonicalKind(String),
/// The scope chunk after the community is not one this grammar defines.
#[error("`{0}` is not one of `{CHANNEL_SCOPE}`, `{GLOBAL_SCOPE}`, `{CONTROL_SCOPE}`")]
UnknownScope(String),
/// A wildcard appeared where a concrete key was required.
#[error("`{0}` contains a wildcard; a concrete publication key was required")]
Wildcard(String),
/// The chunk count or shape does not match any arm of the grammar.
#[error("`{0}` is not a MERIDIAN bus key expression")]
Malformed(String),
}
/// A fully-parsed concrete MERIDIAN bus key expression.
///
/// "Concrete" is load-bearing: this type can only ever describe a single
/// publication key. Subscription patterns are produced by the free functions
/// below and are deliberately not representable here, so a wildcard cannot be
/// handed to a publisher by accident.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum BusKey {
/// An event on a community-scoped routing topic, at one exact kind.
Event {
/// Server-resolved community.
community: CommunityId,
/// Tenant-local routing scope.
topic: EventTopic,
/// The event kind, the last routing chunk.
kind: u32,
},
/// A non-event control message.
Control {
/// Server-resolved community.
community: CommunityId,
/// Which control lane.
control: ControlTopic,
},
}
impl BusKey {
/// The community this key is fenced to.
#[must_use]
pub const fn community(&self) -> CommunityId {
match self {
BusKey::Event { community, .. } | BusKey::Control { community, .. } => *community,
}
}
/// Render the canonical key expression.
#[must_use]
pub fn render(&self) -> String {
match self {
BusKey::Event {
community,
topic: EventTopic::Channel(channel),
kind,
} => format!(
"{ROOT}/{VERSION}/{community}/{CHANNEL_SCOPE}/{channel}/{KIND_SCOPE}/{kind}"
),
BusKey::Event {
community,
topic: EventTopic::Global,
kind,
} => format!("{ROOT}/{VERSION}/{community}/{GLOBAL_SCOPE}/{KIND_SCOPE}/{kind}"),
BusKey::Control { community, control } => {
format!("{ROOT}/{VERSION}/{community}/{CONTROL_SCOPE}/{control}")
}
}
}
/// Parse a concrete key expression received from the bus.
///
/// Rejects wildcards, MQTT punctuation, every non-canonical spelling of a
/// UUID or a kind, and any wire version other than [`VERSION`].
pub fn parse(raw: &str) -> Result<BusKey, KeyExprError> {
reject_mqtt_wildcards(raw)?;
if raw.contains('*') {
return Err(KeyExprError::Wildcard(raw.to_string()));
}
let mut chunks = raw.split('/');
match chunks.next() {
Some(ROOT) => {}
_ => return Err(KeyExprError::WrongRoot(raw.to_string())),
}
match chunks.next() {
Some(VERSION) => {}
Some(other) => {
return Err(KeyExprError::WrongVersion {
found: raw.to_string(),
found_version: other.to_string(),
})
}
None => return Err(KeyExprError::WrongRoot(raw.to_string())),
}
let community = CommunityId::from_uuid(parse_uuid_chunk(
chunks.next().ok_or_else(|| malformed(raw))?,
)?);
let key = match chunks.next().ok_or_else(|| malformed(raw))? {
CHANNEL_SCOPE => {
let channel = parse_uuid_chunk(chunks.next().ok_or_else(|| malformed(raw))?)?;
expect_chunk(chunks.next(), KIND_SCOPE, raw)?;
let kind = parse_kind_chunk(chunks.next().ok_or_else(|| malformed(raw))?)?;
BusKey::Event {
community,
topic: EventTopic::Channel(channel),
kind,
}
}
GLOBAL_SCOPE => {
expect_chunk(chunks.next(), KIND_SCOPE, raw)?;
let kind = parse_kind_chunk(chunks.next().ok_or_else(|| malformed(raw))?)?;
BusKey::Event {
community,
topic: EventTopic::Global,
kind,
}
}
CONTROL_SCOPE => {
let chunk = chunks.next().ok_or_else(|| malformed(raw))?;
let control = ControlTopic::from_chunk(chunk)
.ok_or_else(|| KeyExprError::UnknownScope(chunk.to_string()))?;
BusKey::Control { community, control }
}
other => return Err(KeyExprError::UnknownScope(other.to_string())),
};
if chunks.next().is_some() {
return Err(malformed(raw));
}
Ok(key)
}
}
impl fmt::Display for BusKey {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.render())
}
}
/// The declared-subscriber pattern for one tenant-local event topic.
///
/// This is the one-to-one replacement for the relay's `channel_wildcard_index`
/// / `global_wildcard_index`: `.../k/*` matches every kind on that topic, so
/// kind filtering that used to happen after the bytes crossed the wire can
/// happen before, one declaration per kind, whenever a caller wants it.
#[must_use]
pub fn event_subscription(community: CommunityId, topic: EventTopic) -> String {
match topic {
EventTopic::Channel(channel) => format!(
"{ROOT}/{VERSION}/{community}/{CHANNEL_SCOPE}/{channel}/{KIND_SCOPE}/{CHUNK_WILDCARD}"
),
EventTopic::Global => {
format!("{ROOT}/{VERSION}/{community}/{GLOBAL_SCOPE}/{KIND_SCOPE}/{CHUNK_WILDCARD}")
}
}
}
/// The declared-subscriber pattern for one exact kind on one topic.
#[must_use]
pub fn event_kind_subscription(community: CommunityId, topic: EventTopic, kind: u32) -> String {
BusKey::Event {
community,
topic,
kind,
}
.render()
}
/// The declared-subscriber pattern for one control lane.
#[must_use]
pub fn control_key(community: CommunityId, control: ControlTopic) -> String {
BusKey::Control { community, control }.render()
}
/// Everything a community publishes, under the current wire version.
///
/// The `v1` chunk is inside the pattern on purpose: a `v2` publication does not
/// intersect it.
#[must_use]
pub fn community_subscription(community: CommunityId) -> String {
format!("{ROOT}/{VERSION}/{community}/{TREE_WILDCARD}")
}
/// This session's own liveness-probe key.
///
/// Outside the community space by construction — `_health` cannot be a
/// canonical UUID — so a probe never lands on a tenant lane and
/// [`BusKey::parse`] rejects it rather than mis-scoping it.
#[must_use]
pub fn health_key(session_id: &str) -> String {
format!("{ROOT}/{VERSION}/{HEALTH_SCOPE}/{session_id}")
}
fn malformed(raw: &str) -> KeyExprError {
KeyExprError::Malformed(raw.to_string())
}
fn expect_chunk(found: Option<&str>, want: &str, raw: &str) -> Result<(), KeyExprError> {
match found {
Some(chunk) if chunk == want => Ok(()),
Some(chunk) => Err(KeyExprError::UnknownScope(chunk.to_string())),
None => Err(malformed(raw)),
}
}
fn reject_mqtt_wildcards(raw: &str) -> Result<(), KeyExprError> {
if raw.contains('+') || raw.contains('#') {
return Err(KeyExprError::MqttWildcard(raw.to_string()));
}
Ok(())
}
/// Parse a UUID chunk, accepting **only** the canonical lowercase hyphenated
/// spelling.
///
/// `Uuid::parse_str` also accepts `{braced}`, `urn:uuid:` and unhyphenated
/// forms. Each of those is a different key expression for the same logical
/// topic, so accepting them here would let a publisher and a subscriber agree
/// on the community and still never match.
fn parse_uuid_chunk(chunk: &str) -> Result<Uuid, KeyExprError> {
let parsed =
Uuid::parse_str(chunk).map_err(|_| KeyExprError::NonCanonicalUuid(chunk.to_string()))?;
if parsed.to_string() != chunk {
return Err(KeyExprError::NonCanonicalUuid(chunk.to_string()));
}
Ok(parsed)
}
/// Parse a kind chunk, accepting **only** canonical decimal digits.
///
/// `"+9".parse::<u32>()` is `Ok(9)` and `"09".parse::<u32>()` is `Ok(9)`, and
/// both render back as `9` — so a permissive parser would report a match for a
/// key expression that Zenoh routes nowhere near the one it renders.
fn parse_kind_chunk(chunk: &str) -> Result<u32, KeyExprError> {
if chunk.is_empty() || !chunk.bytes().all(|b| b.is_ascii_digit()) {
return Err(KeyExprError::NonCanonicalKind(chunk.to_string()));
}
if chunk.len() > 1 && chunk.starts_with('0') {
return Err(KeyExprError::NonCanonicalKind(chunk.to_string()));
}
chunk
.parse::<u32>()
.map_err(|_| KeyExprError::NonCanonicalKind(chunk.to_string()))
}
#[cfg(test)]
mod tests {
use super::*;
fn community(n: u128) -> CommunityId {
CommunityId::from_uuid(Uuid::from_u128(n))
}
fn sample_keys() -> Vec<String> {
let c = community(0xaaaa);
let ch = Uuid::from_u128(0xbbbb);
let mut keys = vec![
event_subscription(c, EventTopic::Channel(ch)),
event_subscription(c, EventTopic::Global),
community_subscription(c),
health_key("abcdef0123456789"),
];
for kind in [0u32, 1, 9, 24200, 39000, u32::MAX] {
keys.push(event_kind_subscription(c, EventTopic::Channel(ch), kind));
keys.push(event_kind_subscription(c, EventTopic::Global, kind));
}
for control in ControlTopic::ALL {
keys.push(control_key(c, control));
}
keys
}
/// The contract test this module exists for. A reference project shipped
/// MQTT punctuation into a Zenoh key space and lost every match.
#[test]
fn no_generated_key_ever_carries_mqtt_punctuation() {
for key in sample_keys() {
assert!(
!key.contains('+'),
"{key} carries an MQTT single-level wildcard"
);
assert!(
!key.contains('#'),
"{key} carries an MQTT multi-level wildcard"
);
assert!(!key.contains(':'), "{key} carries a Redis-style separator");
assert!(
key.starts_with(&format!("{ROOT}/{VERSION}/")),
"{key} escaped the wire fence"
);
}
}
#[test]
fn mqtt_wildcards_are_rejected_by_name_not_as_a_bad_uuid() {
let c = community(0xaaaa);
for raw in [
format!("{ROOT}/{VERSION}/+/{CHANNEL_SCOPE}/x/{KIND_SCOPE}/9"),
format!("{ROOT}/{VERSION}/{c}/{CHANNEL_SCOPE}/+/{KIND_SCOPE}/9"),
format!("{ROOT}/{VERSION}/{c}/{CHANNEL_SCOPE}/x/{KIND_SCOPE}/+"),
format!("{ROOT}/{VERSION}/{c}/#"),
] {
assert!(
matches!(BusKey::parse(&raw), Err(KeyExprError::MqttWildcard(_))),
"{raw} must be named as an MQTT wildcard, got {:?}",
BusKey::parse(&raw)
);
}
}
#[test]
fn subscription_patterns_use_zenoh_wildcards() {
let c = community(0xaaaa);
let ch = Uuid::from_u128(0xbbbb);
assert!(event_subscription(c, EventTopic::Channel(ch)).ends_with("/k/*"));
assert!(event_subscription(c, EventTopic::Global).ends_with("/k/*"));
assert!(community_subscription(c).ends_with("/**"));
}
#[test]
fn every_concrete_key_round_trips() {
let c = community(0xaaaa);
let ch = Uuid::from_u128(0xbbbb);
let mut keys = vec![];
for kind in [0u32, 1, 9, 24200, 39000, u32::MAX] {
keys.push(BusKey::Event {
community: c,
topic: EventTopic::Channel(ch),
kind,
});
keys.push(BusKey::Event {
community: c,
topic: EventTopic::Global,
kind,
});
}
for control in ControlTopic::ALL {
keys.push(BusKey::Control {
community: c,
control,
});
}
for key in keys {
let rendered = key.render();
assert_eq!(
BusKey::parse(&rendered).expect("canonical render must parse"),
key,
"{rendered} did not round-trip"
);
}
}
#[test]
fn v2_is_refused_by_name_so_a_v1_reader_never_guesses() {
let c = community(0xaaaa);
let raw = format!("{ROOT}/v2/{c}/{GLOBAL_SCOPE}/{KIND_SCOPE}/9");
match BusKey::parse(&raw) {
Err(KeyExprError::WrongVersion { found_version, .. }) => {
assert_eq!(found_version, "v2");
}
other => panic!("expected a version refusal, got {other:?}"),
}
}
#[test]
fn non_canonical_uuid_spellings_are_refused() {
let c = community(0xaaaa);
let ch = Uuid::from_u128(0xbbbb);
for spelling in [
ch.simple().to_string(),
format!("{{{ch}}}"),
format!("urn:uuid:{ch}"),
ch.to_string().to_uppercase(),
] {
let raw = format!("{ROOT}/{VERSION}/{c}/{CHANNEL_SCOPE}/{spelling}/{KIND_SCOPE}/9");
assert!(
BusKey::parse(&raw).is_err(),
"{spelling} must not be accepted as a channel chunk"
);
}
}
#[test]
fn non_canonical_kind_spellings_are_refused() {
let c = community(0xaaaa);
// `"+9".parse::<u32>()` is `Ok(9)` and `"09".parse::<u32>()` is `Ok(9)`.
// Both render back as `9`, so a permissive parser would report a match
// for a key Zenoh routes nowhere near the one it renders. `"+9"` is
// caught one step earlier, by the MQTT check, which is why the
// assertion below accepts either refusal — what matters is that none
// of these is accepted.
for kind in ["+9", "09", "-9", "9.0", "0x9", "", " 9", "4294967296"] {
let raw = format!("{ROOT}/{VERSION}/{c}/{GLOBAL_SCOPE}/{KIND_SCOPE}/{kind}");
assert!(
BusKey::parse(&raw).is_err(),
"kind chunk {kind:?} must be refused"
);
}
assert!(matches!(
BusKey::parse(&format!(
"{ROOT}/{VERSION}/{c}/{GLOBAL_SCOPE}/{KIND_SCOPE}/09"
)),
Err(KeyExprError::NonCanonicalKind(_))
));
// `0` is canonical and must survive the leading-zero rule.
let raw = format!("{ROOT}/{VERSION}/{c}/{GLOBAL_SCOPE}/{KIND_SCOPE}/0");
assert!(BusKey::parse(&raw).is_ok());
}
#[test]
fn the_legacy_redis_key_shape_is_refused() {
let c = community(0xaaaa);
let ch = Uuid::from_u128(0xbbbb);
for raw in [
format!("meridian:{c}:channel:{ch}"),
format!("meridian:{c}:global"),
format!("meridian:{c}:cache-invalidate"),
] {
assert!(
BusKey::parse(&raw).is_err(),
"the Redis key shape must not parse as a Zenoh bus key: {raw}"
);
}
}
#[test]
fn a_wildcard_is_not_a_publication_key() {
let c = community(0xaaaa);
for raw in [
event_subscription(c, EventTopic::Global),
community_subscription(c),
] {
assert!(
matches!(BusKey::parse(&raw), Err(KeyExprError::Wildcard(_))),
"{raw} must not parse as a concrete publication key"
);
}
}
#[test]
fn a_health_key_is_not_a_community_key() {
let key = health_key("abcdef");
assert!(BusKey::parse(&key).is_err());
assert!(key.starts_with(&format!("{ROOT}/{VERSION}/{HEALTH_SCOPE}/")));
}
#[test]
fn two_communities_never_share_a_key() {
let ch = Uuid::from_u128(0xcccc);
assert_ne!(
event_kind_subscription(community(0xaaaa), EventTopic::Channel(ch), 9),
event_kind_subscription(community(0xbbbb), EventTopic::Channel(ch), 9)
);
assert_ne!(
community_subscription(community(0xaaaa)),
community_subscription(community(0xbbbb))
);
}
}

View file

@ -0,0 +1,798 @@
//! The `ZenohEventBus` contract: key expressions against Zenoh's own matcher,
//! attachment-only local-echo rejection, two-node delivery on all three lanes,
//! declared-subscriber lifetime, fail-closed readiness, and shadow parity.
//!
//! Gated behind the off-by-default `zenoh` cargo feature:
//! `cargo test -p meridian-pubsub --features zenoh --test zenoh_bus`.
//!
//! Every Zenoh test is `#[tokio::test(flavor = "multi_thread")]`. This is not
//! taste: `zenoh-runtime-1.8.0/src/lib.rs:149` **panics** on Tokio's
//! current-thread scheduler, so a single-threaded test does not fail, it aborts.
#![cfg(feature = "zenoh")]
// Wrapped in a module so the test *names* carry `zenoh_bus`, not just the test
// binary — the same reason `zenoh_config.rs` does it. Unwrapped, filtering by
// name matches nothing and reports a green "0 passed; N filtered out".
mod zenoh_bus {
use std::time::Duration;
use meridian_core::qos::Lane;
use meridian_core::{CommunityId, TenantContext};
use meridian_pubsub::cache_invalidation::CacheInvalidation;
use meridian_pubsub::conn_control::ConnControl;
use meridian_pubsub::trust::{admit, Admission, LaneFloors, PayloadProof, TrustProfile};
use meridian_pubsub::zenoh::codec::{self, RoutingAttachment};
use meridian_pubsub::zenoh::config::{ZenohBusConfig, ZenohConfigError, ZenohMode};
use meridian_pubsub::zenoh::health::ReadinessGap;
use meridian_pubsub::zenoh::topic::{self, BusKey, ControlTopic};
use meridian_pubsub::zenoh::{
decide_from_attachment, qos_for_kind, AttachmentDecision, AttachmentRejection,
LocalEchoFilter, ShadowEventBus, ZenohBusError, ZenohEventBus,
};
use meridian_pubsub::{
AnyEventBus, ChannelEvent, EventBus, EventTopic, PubSubConfig, RedisEventBus,
};
use uuid::Uuid;
use zenoh::key_expr::KeyExpr;
const COMMUNITY_A: u128 = 0x0a11;
const COMMUNITY_B: u128 = 0x0b22;
fn community(n: u128) -> CommunityId {
CommunityId::from_uuid(Uuid::from_u128(n))
}
fn ctx(n: u128) -> TenantContext {
TenantContext::resolved(community(n), "bus.example")
}
fn key(raw: &str) -> KeyExpr<'static> {
KeyExpr::try_from(raw.to_string())
.unwrap_or_else(|e| panic!("`{raw}` is not a valid Zenoh key expression: {e}"))
}
fn sample_event(kind: u16, content: &str) -> nostr::Event {
let keys = nostr::Keys::generate();
nostr::EventBuilder::new(nostr::Kind::from_u16(kind), content)
.sign_with_keys(&keys)
.expect("signing a local event")
}
/// A free loopback TCP port. Bound and released, which is racy in theory and
/// has never been in practice for a test that binds it milliseconds later.
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()
}
// ---------------------------------------------------------------- key space
/// The contract test this whole key module exists for.
///
/// Zenoh's wildcards are `*` and `**`. `+` is a perfectly legal *literal*
/// chunk, so an MQTT-shaped subscription declares successfully, matches the
/// string `"+"` and nothing else, and receives silence. A reference project
/// shipped exactly that. Asserted here against Zenoh's own matcher rather
/// than against string equality, because string equality would pass on the
/// day the matcher changed.
#[test]
fn an_mqtt_wildcard_matches_nothing_meridian_publishes() {
let c = community(COMMUNITY_A);
let channel = Uuid::from_u128(0xc0ffee);
let mqtt = key(&format!("meridian/v1/{c}/ch/+/k/9"));
let mqtt_kind = key(&format!("meridian/v1/{c}/ch/{channel}/k/+"));
for kind in [9u32, 24200, 39000] {
let published = key(&topic::event_kind_subscription(
c,
EventTopic::Channel(channel),
kind,
));
assert!(
!mqtt.intersects(&published),
"the MQTT `+` chunk must not match {published}"
);
assert!(
!mqtt_kind.intersects(&published),
"the MQTT `+` chunk must not match {published}"
);
}
// And the Zenoh wildcard in the same position does match.
let zenoh_wildcard = key(&topic::event_subscription(c, EventTopic::Channel(channel)));
assert!(
zenoh_wildcard.intersects(&key(&topic::event_kind_subscription(
c,
EventTopic::Channel(channel),
9
)))
);
}
#[test]
fn a_kind_wildcard_declaration_matches_every_kind_on_its_own_topic_only() {
let c = community(COMMUNITY_A);
let channel = Uuid::from_u128(0xc0ffee);
let other_channel = Uuid::from_u128(0xdecaf);
let declaration = key(&topic::event_subscription(c, EventTopic::Channel(channel)));
for kind in [0u32, 9, 24200, 39000, u32::MAX] {
assert!(
declaration.intersects(&key(&topic::event_kind_subscription(
c,
EventTopic::Channel(channel),
kind
))),
"kind {kind} must match the topic's `/k/*` declaration"
);
}
assert!(
!declaration.intersects(&key(&topic::event_kind_subscription(
c,
EventTopic::Channel(other_channel),
9
))),
"another channel must not match"
);
assert!(
!declaration.intersects(&key(&topic::event_kind_subscription(
community(COMMUNITY_B),
EventTopic::Channel(channel),
9
))),
"another community must not match — the community chunk is the tenant fence"
);
assert!(
!declaration.intersects(&key(&topic::event_kind_subscription(
c,
EventTopic::Global,
9
))),
"the global scope must not match a channel declaration"
);
}
/// `v1` is a wire fence, not a comment. A payload encoding that a `v1`
/// reader would misread takes `v2`, and the two coexist for the length of a
/// rollout precisely because this assertion holds.
#[test]
fn a_v1_subscriber_never_sees_v2() {
let c = community(COMMUNITY_A);
let v1_everything = key(&topic::community_subscription(c));
let v2_event = key(&format!("meridian/v2/{c}/global/k/9"));
let v1_event = key(&topic::event_kind_subscription(c, EventTopic::Global, 9));
assert!(v1_everything.intersects(&v1_event));
assert!(
!v1_everything.intersects(&v2_event),
"the version chunk must fence the two encodings apart"
);
}
#[test]
fn the_control_lane_wildcard_covers_every_community_and_only_its_own_lane() {
let cache_declaration = key("meridian/v1/*/ctl/cache");
for id in [COMMUNITY_A, COMMUNITY_B] {
let c = community(id);
assert!(cache_declaration.intersects(&key(&topic::control_key(c, ControlTopic::Cache))));
assert!(
!cache_declaration.intersects(&key(&topic::control_key(c, ControlTopic::Conn))),
"the cache lane must not receive connection control"
);
assert!(
!cache_declaration.intersects(&key(&topic::control_key(c, ControlTopic::Fence))),
"the cache lane must not receive fence changes"
);
}
}
/// Every key this crate generates must be a key expression Zenoh accepts.
/// A grammar that produces something `KeyExpr::try_from` rejects fails at
/// declare time, in production, on the first tenant with that shape.
#[test]
fn every_generated_key_is_a_valid_zenoh_key_expression() {
let c = community(COMMUNITY_A);
let channel = Uuid::from_u128(0xc0ffee);
let mut keys = vec![
topic::event_subscription(c, EventTopic::Channel(channel)),
topic::event_subscription(c, EventTopic::Global),
topic::community_subscription(c),
topic::health_key("a1b2c3d4e5f6"),
"meridian/v1/*/ctl/cache".to_string(),
"meridian/v1/*/ctl/conn".to_string(),
];
for kind in [0u32, 9, 24200, u32::MAX] {
keys.push(topic::event_kind_subscription(
c,
EventTopic::Channel(channel),
kind,
));
keys.push(topic::event_kind_subscription(c, EventTopic::Global, kind));
}
for control in ControlTopic::ALL {
keys.push(topic::control_key(c, control));
}
for raw in keys {
let _ = key(&raw);
}
}
// ------------------------------------------------------- attachment / dedup
/// Charter Gap #5, proved by the signature rather than by the call order.
///
/// `decide_from_attachment` takes no payload. It therefore *cannot*
/// deserialize the event, so "the echo check runs before the payload is
/// decoded" is a property of the type. The payload below is deliberate
/// garbage: had the decision needed it, this test would fail rather than
/// pass by accident.
#[test]
fn a_local_echo_is_rejected_from_the_attachment_alone() {
let filter = LocalEchoFilter::new(64);
let event = sample_event(9, "echo me");
let event_id = event.id.to_bytes();
let attachment = RoutingAttachment::new(
community(COMMUNITY_A),
Some(Uuid::from_u128(0xc0ffee)),
9,
event_id,
)
.encode()
.expect("encode");
// The payload that would have to be decoded if the attachment were not
// enough. It is not valid postcard for any of our types.
let undecodable_payload = vec![0xffu8; 64];
assert!(
codec::decode_event(&undecodable_payload).is_err(),
"fixture drift: this payload only proves anything while it is undecodable"
);
assert_eq!(
decide_from_attachment(Some(&attachment), &filter),
AttachmentDecision::Deliver(RoutingAttachment::decode(&attachment).expect("decode")),
"an unseen id must be handed on for payload decode"
);
filter.record(event_id);
assert_eq!(
decide_from_attachment(Some(&attachment), &filter),
AttachmentDecision::LocalEcho,
"a recorded id must be rejected with no payload in scope at all"
);
}
#[test]
fn a_missing_unknown_or_corrupt_attachment_is_named_not_guessed() {
let filter = LocalEchoFilter::new(4);
assert_eq!(
decide_from_attachment(None, &filter),
AttachmentDecision::Reject(AttachmentRejection::Missing)
);
assert_eq!(
decide_from_attachment(Some(&[0xff, 0xff, 0xff]), &filter),
AttachmentDecision::Reject(AttachmentRejection::Undecodable)
);
let mut attachment = RoutingAttachment::new(community(COMMUNITY_A), None, 9, [3u8; 32]);
attachment.schema = 99;
let bytes = postcard::to_stdvec(&attachment).expect("encode");
assert_eq!(
decide_from_attachment(Some(&bytes), &filter),
AttachmentDecision::Reject(AttachmentRejection::UnknownSchema)
);
}
// ------------------------------------------------------------------- config
#[test]
fn the_qos_for_a_kind_comes_from_the_landed_lane_table() {
assert_eq!(
qos_for_kind(meridian_core::kind::KIND_STREAM_MESSAGE_V2)
.expect("chat has a lane")
.lane,
Lane::Q4
);
assert!(
qos_for_kind(65535).is_err(),
"an unclassified kind must fail closed, not default onto Data"
);
}
#[test]
fn a_non_loopback_endpoint_is_refused_before_a_session_can_exist() {
assert!(matches!(
ZenohBusConfig::new(ZenohMode::Peer, ["tcp/10.1.2.3:7447"]),
Err(ZenohConfigError::NonLocalEndpoint(_))
));
}
// ------------------------------------------------------------ live sessions
fn peer_config(listen_port: u16, connect_port: u16) -> ZenohBusConfig {
ZenohBusConfig::new(ZenohMode::Peer, [format!("tcp/127.0.0.1:{connect_port}")])
.expect("loopback endpoint")
.with_listen(format!("tcp/127.0.0.1:{listen_port}"))
.expect("loopback listen")
.with_readiness_timeout(Duration::from_secs(15))
.with_loopback_freshness(Duration::from_secs(30))
}
/// Two mutually-connected peers. Zenoh keeps one link; both sides end up
/// with one transport, which is what readiness requires.
async fn peer_pair() -> (ZenohEventBus, ZenohEventBus) {
let port_a = free_port();
let port_b = free_port();
let a = ZenohEventBus::open(peer_config(port_a, port_b))
.await
.expect("peer A must open");
let b = ZenohEventBus::open(peer_config(port_b, port_a))
.await
.expect("peer B must open");
(a, b)
}
async fn recv_within(
rx: &mut tokio::sync::broadcast::Receiver<ChannelEvent>,
budget: Duration,
) -> Option<ChannelEvent> {
tokio::time::timeout(budget, rx.recv()).await.ok()?.ok()
}
/// A session that opens with no reachable endpoint must NOT report ready.
///
/// The pinned posture opens anyway on purpose — the relay serves NIP-01
/// whether or not the bus tier is up — so this is the assertion that stops
/// "the process is alive" from being read as "the bus works".
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn readiness_fails_closed_when_no_transport_exists() {
let config =
ZenohBusConfig::new(ZenohMode::Peer, [format!("tcp/127.0.0.1:{}", free_port())])
.expect("loopback endpoint")
.with_listen(format!("tcp/127.0.0.1:{}", free_port()))
.expect("loopback listen")
.with_readiness_timeout(Duration::from_millis(400));
let bus = ZenohEventBus::open(config)
.await
.expect("the session must open even with the bus tier absent");
match bus.wait_ready().await {
Err(ZenohBusError::NotReady(ReadinessGap::NoTransport)) => {}
other => panic!("expected a named NoTransport refusal, got {other:?}"),
}
assert!(!bus.health().is_ready());
bus.close().await.expect("shutdown must be clean");
}
/// The loopback probe is an *application* check: the sample has to come
/// back through Zenoh, not merely be handed to a socket.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_connected_pair_reaches_readiness_including_the_loopback() {
let (a, b) = peer_pair().await;
a.wait_ready().await.expect("peer A must reach readiness");
b.wait_ready().await.expect("peer B must reach readiness");
let snapshot = a.health().snapshot();
assert!(snapshot.session_open);
assert!(snapshot.endpoints_connected >= 1);
assert_eq!(snapshot.declarations_active, snapshot.declarations_required);
assert!(
snapshot.last_loopback_unix.is_some(),
"readiness must be evidence of delivery, not of configuration"
);
a.close().await.expect("clean shutdown");
b.close().await.expect("clean shutdown");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn an_event_published_on_one_peer_arrives_on_the_other() {
let (a, b) = peer_pair().await;
a.wait_ready().await.expect("peer A ready");
b.wait_ready().await.expect("peer B ready");
let ctx = ctx(COMMUNITY_A);
let channel = Uuid::from_u128(0xc0ffee);
let topic = EventTopic::Channel(channel);
a.retain_topic(&ctx, topic).await;
let mut rx = EventBus::subscribe_local(&a);
// Give the remote declaration a moment to propagate.
tokio::time::sleep(Duration::from_millis(200)).await;
let event = sample_event(9, "hello from peer B");
b.publish_event(&ctx, topic, &event)
.await
.expect("publish must be accepted");
let received = recv_within(&mut rx, Duration::from_secs(5))
.await
.expect("peer A must receive peer B's event");
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");
}
/// The publishing pod must not deliver its own publication back to its own
/// consumers — and the drop must happen from the attachment.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_pods_own_publication_is_not_delivered_back_to_it() {
let (a, b) = peer_pair().await;
a.wait_ready().await.expect("peer A ready");
b.wait_ready().await.expect("peer B ready");
let ctx = ctx(COMMUNITY_A);
let topic = EventTopic::Global;
a.retain_topic(&ctx, topic).await;
let mut rx = EventBus::subscribe_local(&a);
tokio::time::sleep(Duration::from_millis(200)).await;
let mine = sample_event(9, "published by A");
a.publish_event(&ctx, topic, &mine)
.await
.expect("publish must be accepted");
assert!(
recv_within(&mut rx, Duration::from_millis(600))
.await
.is_none(),
"a pod must not receive its own publication back"
);
// A different pod's event on the same topic still arrives, so the drop
// above is the echo check and not a broken subscription.
let theirs = sample_event(9, "published by B");
b.publish_event(&ctx, topic, &theirs)
.await
.expect("publish must be accepted");
let received = recv_within(&mut rx, Duration::from_secs(5))
.await
.expect("the remote event must still arrive");
assert_eq!(received.event, theirs);
a.close().await.expect("clean shutdown");
b.close().await.expect("clean shutdown");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn the_control_lanes_carry_their_own_payloads_scoped_by_key() {
let (a, b) = peer_pair().await;
a.wait_ready().await.expect("peer A ready");
b.wait_ready().await.expect("peer B ready");
let ctx = ctx(COMMUNITY_B);
let mut cache_rx = EventBus::subscribe_cache_invalidations(&a);
let mut conn_rx = EventBus::subscribe_conn_control(&a);
tokio::time::sleep(Duration::from_millis(200)).await;
let invalidation = CacheInvalidation::Membership {
channel_id: Uuid::from_u128(0xfeed),
pubkey: vec![7; 32],
};
b.publish_cache_invalidation(&ctx, &invalidation)
.await
.expect("cache invalidation must be accepted");
let command = ConnControl::DisconnectPubkey {
pubkey: vec![9; 32],
event_id: "abc".to_string(),
reason: "banned: spam".to_string(),
};
b.publish_conn_control(&ctx, &command)
.await
.expect("connection control must be accepted");
let scoped_cache = tokio::time::timeout(Duration::from_secs(5), cache_rx.recv())
.await
.expect("cache invalidation must arrive")
.expect("lane open");
assert_eq!(scoped_cache.community_id, ctx.community());
assert_eq!(scoped_cache.invalidation, invalidation);
let scoped_conn = tokio::time::timeout(Duration::from_secs(5), conn_rx.recv())
.await
.expect("connection control must arrive")
.expect("lane open");
assert_eq!(scoped_conn.community_id, ctx.community());
assert_eq!(scoped_conn.command, command);
a.close().await.expect("clean shutdown");
b.close().await.expect("clean shutdown");
}
/// A pod that never retained a topic must not receive it. This is the
/// `/k/{kind}` and per-topic declaration win: interest is expressed in the
/// key space, so uninterested pods stop receiving the bytes entirely.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn an_unretained_topic_is_not_delivered() {
let (a, b) = peer_pair().await;
a.wait_ready().await.expect("peer A ready");
b.wait_ready().await.expect("peer B ready");
let ctx = ctx(COMMUNITY_A);
let retained = EventTopic::Channel(Uuid::from_u128(0x1111));
let ignored = EventTopic::Channel(Uuid::from_u128(0x2222));
a.retain_topic(&ctx, retained).await;
let mut rx = EventBus::subscribe_local(&a);
tokio::time::sleep(Duration::from_millis(200)).await;
b.publish_event(&ctx, ignored, &sample_event(9, "unwanted"))
.await
.expect("publish must be accepted");
assert!(
recv_within(&mut rx, Duration::from_millis(600))
.await
.is_none(),
"an unretained topic must not be delivered"
);
b.publish_event(&ctx, retained, &sample_event(9, "wanted"))
.await
.expect("publish must be accepted");
assert!(
recv_within(&mut rx, Duration::from_secs(5)).await.is_some(),
"the retained topic must still be delivered"
);
a.close().await.expect("clean shutdown");
b.close().await.expect("clean shutdown");
}
/// Interest is the declaration's lifetime. There is no `DesiredTopics` and
/// no 500 ms debounce: the retain count and the live `Subscriber` are one
/// map entry, so "declared" and "desired" cannot disagree.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn interest_is_the_declared_subscribers_lifetime() {
let (a, b) = peer_pair().await;
a.wait_ready().await.expect("peer A ready");
b.wait_ready().await.expect("peer B ready");
let ctx = ctx(COMMUNITY_A);
let topic = EventTopic::Channel(Uuid::from_u128(0x3333));
assert_eq!(a.topic_refcount(&ctx, topic).await, 0);
a.retain_topic(&ctx, topic).await;
a.retain_topic(&ctx, topic).await;
assert_eq!(a.topic_refcount(&ctx, topic).await, 2);
a.release_topic(&ctx, topic).await;
assert_eq!(
a.topic_refcount(&ctx, topic).await,
1,
"a balanced retain/release pair must be a no-op"
);
let mut rx = EventBus::subscribe_local(&a);
tokio::time::sleep(Duration::from_millis(200)).await;
b.publish_event(&ctx, topic, &sample_event(9, "still wanted"))
.await
.expect("publish must be accepted");
assert!(
recv_within(&mut rx, Duration::from_secs(5)).await.is_some(),
"one outstanding retain must keep the declaration alive"
);
a.release_topic(&ctx, topic).await;
assert_eq!(a.topic_refcount(&ctx, topic).await, 0);
tokio::time::sleep(Duration::from_millis(300)).await;
b.publish_event(&ctx, topic, &sample_event(9, "no longer wanted"))
.await
.expect("publish must be accepted");
assert!(
recv_within(&mut rx, Duration::from_millis(800))
.await
.is_none(),
"dropping the last retain must undeclare, with no debounce window"
);
a.close().await.expect("clean shutdown");
b.close().await.expect("clean shutdown");
}
/// A kind with no MIP-QC lane has no scheduling class, so there is nothing
/// to publish it on. Explicit denial rather than a convenient `Data`
/// default — the exact mechanism by which a durable kind would otherwise
/// acquire best-effort delivery.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn an_unclassified_kind_is_refused_at_publish() {
let port = free_port();
let bus = ZenohEventBus::open(peer_config(port, free_port()))
.await
.expect("session must open");
let unclassified = sample_event(65535, "no lane");
assert!(
qos_for_kind(65535).is_err(),
"fixture drift: 65535 must remain unclassified for this test to mean anything"
);
match bus
.publish_event(&ctx(COMMUNITY_A), EventTopic::Global, &unclassified)
.await
{
Err(e) => {
let text = e.to_string();
assert!(
text.contains("65535") && text.contains("lane"),
"the refusal must name the kind and the missing lane, got: {text}"
);
}
Ok(()) => panic!("an unclassified kind must not be published"),
}
bus.close().await.expect("clean shutdown");
}
// ------------------------------------------------------------------- shadow
fn offline_redis_bus_pool() -> deadpool_redis::Pool {
// Port 1 is reserved and nothing listens on it, so every Redis command
// fails immediately. That is the point: it makes "Redis is
// authoritative" a deterministic assertion instead of one that needs a
// container.
deadpool_redis::Config::from_url("redis://127.0.0.1:1")
.create_pool(Some(deadpool_redis::Runtime::Tokio1))
.expect("pool construction does not connect")
}
/// The whole point of shadow mode: Redis decides, Zenoh is observed.
///
/// Redis is unreachable here and Zenoh is live, so the shadow *succeeds*
/// where the primary fails. The caller must still see the primary's error —
/// a shadow that could rescue a primary failure makes the primary's health
/// unobservable, which is the one thing a soak must not do.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn shadow_mode_serves_redis_and_only_counts_zenoh() {
let zenoh = ZenohEventBus::open(peer_config(free_port(), free_port()))
.await
.expect("session must open");
let redis = RedisEventBus::with_config(
PubSubConfig::new("redis://127.0.0.1:1")
.with_unsubscribe_debounce(Duration::from_millis(1)),
offline_redis_bus_pool(),
)
.await
.expect("bus construction does not connect");
let shadow = ShadowEventBus::new(redis, zenoh);
assert!(shadow.parity_snapshot().is_clean());
let ctx = ctx(COMMUNITY_A);
let event = sample_event(9, "dual published");
let result = shadow.publish_event(&ctx, EventTopic::Global, &event).await;
assert!(
result.is_err(),
"Redis is authoritative: its failure must reach the caller even though Zenoh succeeded"
);
let parity = shadow.parity_snapshot();
let lane = parity.lane(meridian_pubsub::zenoh::shadow::ParityLane::Events);
assert_eq!(lane.zenoh_only, 1, "the divergence must be recorded");
assert_eq!(lane.both, 0);
assert!(!parity.is_clean(), "one miss must fail the cutover gate");
// Delivery still comes from Redis, so the shadow's own receivers are
// never what `subscribe_local` hands out.
let _redis_receiver = EventBus::subscribe_local(&shadow);
}
// ------------------------------------------------------------- dispatcher
/// The enum dispatcher, exercised through the trait for every variant that
/// is actually constructed. `Box<dyn EventBus>` is impossible here — the
/// seam is RPITIT — and this is what replaces it.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn the_enum_dispatcher_carries_every_constructed_backend() {
async fn lanes_are_independent<B: EventBus>(bus: &B) {
let _events = EventBus::subscribe_local(bus);
let _cache = EventBus::subscribe_cache_invalidations(bus);
let _control = EventBus::subscribe_conn_control(bus);
assert_eq!(EventBus::subscribe_local(bus).len(), 0);
}
let redis = AnyEventBus::Redis(
RedisEventBus::with_config(
PubSubConfig::new("redis://127.0.0.1:1"),
offline_redis_bus_pool(),
)
.await
.expect("bus construction does not connect"),
);
assert_eq!(redis.label(), "redis");
lanes_are_independent(&redis).await;
let zenoh = AnyEventBus::Zenoh(
ZenohEventBus::open(peer_config(free_port(), free_port()))
.await
.expect("session must open"),
);
assert_eq!(zenoh.label(), "zenoh");
lanes_are_independent(&zenoh).await;
let shadow = AnyEventBus::Shadow(ShadowEventBus::new(
RedisEventBus::with_config(
PubSubConfig::new("redis://127.0.0.1:1"),
offline_redis_bus_pool(),
)
.await
.expect("bus construction does not connect"),
ZenohEventBus::open(peer_config(free_port(), free_port()))
.await
.expect("session must open"),
));
assert_eq!(shadow.label(), "shadow");
lanes_are_independent(&shadow).await;
}
// ------------------------------------------------------------------- trust
/// The signature the relay's cross-node re-verification branch is written
/// against. Asserted here so the branch has a contract to land onto rather
/// than a shape to guess at.
#[test]
fn the_trust_ladder_decides_accept_reverify_reject() {
// In-deployment today: every lane floors at P0, so an unauthenticated
// link is accepted. That is not a relaxation, it is the profile the
// cross-pod fan-out already runs at, written down.
for lane in Lane::ALL {
assert_eq!(
admit(
TrustProfile::P0None,
TrustProfile::P0None,
lane,
&LaneFloors::IN_DEPLOYMENT,
PayloadProof::SignedNostrEvent,
),
Admission::Accept
);
}
// Raise a lane's floor and a below-floor signed event is re-verified.
let federated = LaneFloors::IN_DEPLOYMENT.with_floor(Lane::Q4, TrustProfile::P3Schnorr);
assert_eq!(
admit(
TrustProfile::P0None,
TrustProfile::P3Schnorr,
Lane::Q4,
&federated,
PayloadProof::SignedNostrEvent,
),
Admission::Reverify,
"a claim must never be believed above the authenticated link"
);
// A control envelope has nothing to re-verify, so it is refused.
let control = LaneFloors::IN_DEPLOYMENT.with_floor(Lane::Q0, TrustProfile::P2Ed25519);
assert_eq!(
admit(
TrustProfile::P0None,
TrustProfile::P2Ed25519,
Lane::Q0,
&control,
PayloadProof::None,
),
Admission::Reject
);
}
#[test]
fn a_bus_key_round_trips_through_the_grammar_the_relay_will_parse() {
let c = community(COMMUNITY_A);
let channel = Uuid::from_u128(0xc0ffee);
let key = BusKey::Event {
community: c,
topic: EventTopic::Channel(channel),
kind: 24200,
};
assert_eq!(BusKey::parse(&key.render()).expect("round trip"), key);
assert_eq!(key.community(), c);
}
}

View file

@ -47,169 +47,17 @@ mod zenoh_config {
/// The policy every positive fixture pins, checked against the parsed config.
///
/// Delegates to `meridian_pubsub::zenoh::config::check_posture`, which is
/// the **only** implementation. This used to be a copy living here, and a
/// copy is the failure the whole directory exists to prevent: the library
/// now builds its runtime session config from the same fixture and checks
/// it with the same function, so a value that stops being defended stops
/// being defended in one place rather than in one of two.
///
/// Returns the first violation rather than panicking, so the same code path
/// serves both the positive fixtures and `policy_rejects_a_plugin_enabled_config`.
/// Kept in one function on purpose: the peer session and the `zenohd` router
/// must not be allowed to drift into different postures, and a shared check is
/// the only thing that stops that happening one reviewed diff at a time.
fn check_meridian_policy(config: &Config) -> Result<(), String> {
let expect = |key: &str, want: &str, why: &str| -> Result<(), String> {
let got = json_at(config, key);
if got == want {
Ok(())
} else {
Err(format!("{key} must be {want} but is {got} — {why}"))
}
};
// QoS on. Without it there is one transmission queue and the Q0-Q7 lane
// mapping is a no-op that no test downstream of here would notice.
expect(
"transport/unicast/qos/enabled",
"true",
"the eight MERIDIAN lanes need Zenoh's priority queues",
)?;
// Low-latency transport off. `qos && lowlatency` is rejected outright by
// zenoh-transport 1.8.0; see `rejects_qos_plus_lowlatency`.
expect(
"transport/unicast/lowlatency",
"false",
"incompatible with QoS in zenoh 1.8.0",
)?;
// Batching on. It is back-pressure-driven and is the throughput path for the
// small-payload traffic class; `express` opts a single publication out.
expect(
"transport/link/tx/queue/batching/enabled",
"true",
"adaptive batching is the small-payload throughput path",
)?;
// Shared memory off, both switches. Zenoh 1.x defaults BOTH to `true`
// (`zenoh-config-1.8.0/src/lib.rs:795-816`), so a config that merely omits
// them has switched on charter Phase 6 — unscheduled, gated on a measured
// payload crossover, and needing `ulimit -l unlimited` / `CAP_IPC_LOCK`.
// The relay-side build cannot use SHM (`shared-memory` is not in the feature
// list) but the pinned `zenohd` image is built with it, so the daemon would
// announce support on a link the other end cannot honour.
expect(
"transport/shared_memory/enabled",
"false",
"SHM is charter Phase 6 and defaults ON — it must be refused explicitly",
)?;
expect(
"transport/shared_memory/transport_optimization/enabled",
"false",
"a second SHM switch that silently routes large messages through shared memory",
)?;
// Both scouting mechanisms off. Upstream defaults are `true` for BOTH
// (`zenoh-config-1.8.0/src/defaults.rs:75` and `:105`) — disabling multicast
// alone leaves gossip discovering and autoconnecting peers, which is the
// charter's Gap #17, and gossip has no CLI flag at all, so the config file
// is the only place it can be closed.
//
// `false` and not "anything but true" on purpose. Both keys are
// `Option<bool>`: an omitted key reads back as `null` here and fails, which
// is the property that matters, because `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. Deleting these two lines from
// a fixture would therefore turn discovery back on, and only this assertion
// notices.
expect(
"scouting/multicast/enabled",
"false",
"topology is declared, never discovered; an omitted key is read as consent by zenohd",
)?;
expect(
"scouting/gossip/enabled",
"false",
"gossip stays on when only multicast is disabled (Gap #17)",
)?;
// Belt and braces on the same gap: even if a peer is reached by some other
// path, nothing may be autoconnected. Checked structurally rather than by
// string match — the value is mode-dependent
// (`{"router":[],"peer":[],"client":[]}`) and every one of those lists must
// be empty.
for key in [
"scouting/multicast/autoconnect",
"scouting/gossip/autoconnect",
"scouting/gossip/target",
] {
let raw = json_at(config, key);
let parsed: serde_json::Value = serde_json::from_str(&raw)
.map_err(|e| format!("{key} did not read back as JSON ({raw}): {e}"))?;
let empty = match &parsed {
serde_json::Value::Null => true,
serde_json::Value::Array(items) => items.is_empty(),
serde_json::Value::Object(per_mode) => per_mode
.values()
.all(|v| v.as_array().is_some_and(|items| items.is_empty())),
_ => false,
};
if !empty {
return Err(format!(
"{key} must name no node type — implicit peering is refused, got {raw}"
));
}
}
// No plugin runtime. The pinned image ships a REST plugin and a
// storage-manager plugin in its root; this is the switch that keeps them
// unreachable.
expect(
"plugins_loading/enabled",
"false",
"no plugin runtime, no REST surface, no second storage mechanism",
)?;
// No plugin is configured either, and the search path is empty. Both matter:
// `zenohd/src/main.rs:136-137` forces `plugins_loading.enabled` back to
// `true` after reading the config file, so the switch above holds for the
// library only. These two are what keep the REST and storage-manager
// plugins shipped in the image root unreachable on the daemon.
expect(
"plugins",
"{}",
"storage manager and REST are refused by name, not just left unloaded",
)?;
expect(
"plugins_loading/search_dirs",
"[]",
"zenohd re-enables plugin loading regardless of config; an empty search path is what survives",
)?;
// Admin space off: it is a query surface over the node, not part of the bus
// contract, and it makes no authorization decision MERIDIAN would honour.
// `zenohd` forces `enabled` back to `true`, so the permissions are the part
// that actually holds on the daemon — pin both.
expect(
"adminspace/enabled",
"false",
"the admin space is not part of the bus contract",
)?;
expect(
"adminspace/permissions/read",
"false",
"zenohd force-enables the admin space; read permission is what survives",
)?;
expect(
"adminspace/permissions/write",
"false",
"runtime config mutation over the admin space is refused outright",
)?;
// Only the link this build actually compiles.
expect(
"transport/link/protocols",
"[\"tcp\"]",
"this build is default-features=false + transport_tcp",
)?;
Ok(())
meridian_pubsub::zenoh::config::check_posture(config)
}
fn assert_meridian_policy(config: &Config, label: &str) {