test(relay): prove the bus tenancy binding ACCEPTS, not just that it refuses
Some checks failed
helm chart / lint + unittest + render matrix (push) Has been cancelled
helm chart / install on kind (gated) (push) Has been cancelled
helm chart / publish chart to GHCR (push) Has been cancelled
Meridian Harness / Build (aarch64-unknown-linux-musl) (push) Has been cancelled
Meridian Harness / Build (x86_64-unknown-linux-musl) (push) Has been cancelled
Meridian Harness / Publish rolling release (push) Has been cancelled
Meridian Harness / Publish tagged release (push) Has been cancelled
CI / Detect Changed Paths (push) Has been cancelled
CI / Dead Token Reference Guard (push) Has been cancelled
control plane / chart (push) Has been cancelled
control plane / test (push) Has been cancelled
control plane / browser-e2e (push) Has been cancelled
control plane / Build control plane image (linux/amd64) (push) Has been cancelled
control plane / Build control plane image (linux/arm64) (push) Has been cancelled
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
control plane / Publish signed control plane image (push) Has been cancelled
Docker image / Merge release multi-arch manifest (push) Has been cancelled
Docker image / Merge debug multi-arch manifest (push) Has been cancelled
Docker image / Publish public push gateway image (push) Has been cancelled

`meridian-f43l` closed a cross-tenant injection path by making a below-floor
bus event's tenancy come from its row in the system of record rather than from
the publisher-chosen key expression. Every refusal on that path was covered
without infrastructure; `CommunityBinding::Bound` was covered by code
inspection only, because the row it asks for cannot exist without Postgres.

That is the wrong way round for a security control: a binding that refused
everything would have passed the entire existing suite. Three tests, in the
crate's `--lib` module because `fan_out_pubsub_event_with` — the only entry
point that reaches the `Reverify` arm without a Zenoh session — is
`pub(crate)`, and gated `#[ignore = "requires Postgres"]` so
`just test-relay-db-run` (CI's "Relay DB/Redis gate") runs them and the
infra-free suite stays infra-free:

- `reverify_delivers_the_community_whose_row_binds_the_event_and_no_other`.
  One event, one signature, one channel UUID, two real communities. Under A,
  where the event's row exists carrying the routed channel scope, it is
  re-verified and delivered. Under B — a real community holding a real *open*
  channel with the same UUID and a subscriber on it, which is the shape an
  attacker gets for free from caller-chosen channel ids and a
  `(community_id, id)` primary key — the identical bytes on the identical
  topic are refused. Without the delivered half, "refused" would be equally
  consistent with a binding that refuses everything.
- `a_relay_derived_channel_kind_is_bound_by_its_row_scope_and_nothing_else`.
  `bind_scope` answers `Silent` for reactions, deletions, gift wraps, 9007 and
  44100/44101, so for those kinds the routed channel is checked in exactly one
  place: `stored.channel_id == routed_channel`. The stream-message pair cannot
  reach that arm — its refusal is decided before any lookup and its acceptance
  has tag and row agreeing — so this test is the one that fails when the
  comparison is deleted. Verified: with it deleted the pair still passes and
  this test fails.
- `unknown_tenancy_refuses_and_withdraws_the_replay_id_that_unbound_spends`.
  A database error refuses rather than accepts, and is told apart from a
  decided `Unbound` by the replay set: `Unbound` spends the id, `Unknown`
  withdraws it so one Postgres blip is not a ten-minute hole. Needs no live
  database — `Unbound` is reached by an ephemeral kind and `Unknown` by a pool
  pinned at an address nothing listens on, so it runs in the default suite.
  The seen-set assertion is also what proves the refusal came from the tenancy
  check rather than from the visibility gate failing closed one step later.

Mutation-verified in both directions: forcing `Bound` fails the paired refusal
and the `Unknown` test; forcing `Unbound` fails the acceptance.

The `--lib` pool is untouched. `test_state` still resolves through
`test_config` and connects lazily; `postgres_state` is a sibling builder that
resolves `meridian_test_db::disposable_url()` and connects eagerly, failing
loudly rather than skipping, the same repair `api::operator` got in
meridian-wxqw. Redis stays unconnectable for all three — the cross-node
*receive* path publishes nothing, and a live one would pull these into
meridian-h0tg's non-exiting test binary. The seeded rows are removed on the
way out, because that gate's database is shared and nothing resets it.

Refs: meridian-f43l
Signed-off-by: Joshua Belke <joshua@innovationhub-act.org>
This commit is contained in:
Josh Belke 2026-08-21 04:47:12 -04:00
commit a9b30ea276

View file

@ -2695,7 +2695,7 @@ mod tests {
// build is a Redis one, and Redis answers `IN_DEPLOYMENT`/`P0None`,
// which is `Accept` for every lane by construction.
use meridian_core::kind::KIND_STREAM_MESSAGE;
use meridian_core::kind::{KIND_REACTION, KIND_STREAM_MESSAGE};
use meridian_pubsub::trust::{LaneFloors, TrustProfile};
use crate::handlers::event::fan_out_pubsub_event_with;
@ -3035,6 +3035,435 @@ mod tests {
);
}
// --- The acceptance path, against a real system of record. ---------
//
// Everything above is proved by refusal, and a binding that refused
// everything would pass all of it. `CommunityBinding::Bound` is the one
// verdict that cannot be reached without Postgres, because the row it
// asks for is the only tenancy evidence a NIP-01 event has.
/// Register a subscriber under a **named** community.
///
/// `register_stream_sub` pins the nil community through
/// `SubscriptionRegistry::register`'s test-only wrapper, and the nil
/// UUID is refused by `chk_communities_id_not_nil`, so no row can ever
/// bind an event to it. A test that asks the system of record a real
/// question therefore needs a real community on the connection as well
/// as on the event: `filter_fanout_by_access` drops, at the send
/// chokepoint, every recipient whose connection is not labelled with
/// the event's community.
fn register_sub_in(
state: &AppState,
community: meridian_core::tenant::CommunityId,
sub_id: &str,
kind: u32,
channel_id: Uuid,
) -> (Uuid, mpsc::Receiver<Message>) {
let conn_id = Uuid::new_v4();
let (tx, rx) = mpsc::channel(10);
let (ctrl_tx, _ctrl_rx) = mpsc::channel(10);
state.conn_manager.register(
conn_id,
tx,
ctrl_tx,
CancellationToken::new(),
community,
Arc::new(AtomicU8::new(0)),
Default::default(),
3,
);
state.sub_registry.register_scoped(
community,
conn_id,
sub_id.to_string(),
vec![Filter::new().kind(Kind::Custom(kind as u16))],
Some(channel_id),
);
(conn_id, rx)
}
/// Seed a community and one **open** channel in it under `channel_id`.
///
/// The channel row is real rather than a seeded
/// `channel_visibility_cache` entry. The refusal tests above seed the
/// cache to keep an unreachable pool from being the thing that stops
/// delivery, but for the acceptance test the visibility gate is part of
/// what has to be shown working: a cache entry is precisely the input
/// that would let it pass with the gate removed.
async fn seed_community_with_open_channel(
pool: &sqlx::PgPool,
label: &str,
channel_id: Uuid,
) -> meridian_core::tenant::CommunityId {
let community_uuid = Uuid::new_v4();
let host = format!("bus-bound-{label}-{}.example", community_uuid.simple());
sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)")
.bind(community_uuid)
.bind(&host)
.execute(pool)
.await
.expect("seed community");
sqlx::query(
"INSERT INTO channels (id, community_id, name, visibility, created_by) \
VALUES ($1, $2, $3, 'open', $4)",
)
.bind(channel_id)
.bind(community_uuid)
.bind(format!("bus-bound-{label}"))
.bind(vec![7u8; 32])
.execute(pool)
.await
.expect("seed open channel");
meridian_core::tenant::CommunityId::from_uuid(community_uuid)
}
/// `just test-relay-db-run` targets a shared disposable database that
/// nothing resets between runs, so a test that writes rows removes its
/// own. Ordered by the foreign keys: events and channels both reference
/// `communities(id)`.
async fn drop_seeded_community(
pool: &sqlx::PgPool,
community: meridian_core::tenant::CommunityId,
) {
for statement in [
"DELETE FROM events WHERE community_id = $1",
"DELETE FROM channels WHERE community_id = $1",
"DELETE FROM communities WHERE id = $1",
] {
sqlx::query(statement)
.bind(community.as_uuid())
.execute(pool)
.await
.expect("clean up seeded rows");
}
}
/// **The acceptance path — the one verdict the rest of this module
/// cannot reach, paired with the refusal that makes it mean anything.**
///
/// One event, one signature, one channel UUID, two real communities:
///
/// * under A — where the event's own row exists, carrying exactly the
/// routed channel scope — it is re-verified and **delivered**;
/// * under B — a real community holding a real open channel with the
/// **same** UUID and a subscriber attached to it — the identical
/// bytes on the identical topic are **refused**.
///
/// B is built that way deliberately. `kind:9007` reaches
/// `create_channel_with_id`, so channel UUIDs are caller-chosen, and
/// the channels primary key is `(community_id, id)` — an attacker who
/// can create a channel in B mints A's UUID there and every check
/// downstream of this one then passes. Only the event row tells the two
/// apart, which is why `community_binding` asks for it and not for the
/// channel, and this pair is the evidence that the distinction is real
/// rather than argued.
///
/// Without the delivered half, "refused" would be equally consistent
/// with a binding that refuses everything, which would pass every other
/// test on this path.
#[tokio::test]
#[ignore = "requires Postgres"]
async fn reverify_delivers_the_community_whose_row_binds_the_event_and_no_other() {
let (state, pool) = super::fanout_access::postgres_state().await;
let (link, floors) = below_floor();
let channel = Uuid::new_v4();
let community_a = seed_community_with_open_channel(&pool, "a", channel).await;
let community_b = seed_community_with_open_channel(&pool, "b", channel).await;
let event = stream_message(channel);
let event_id = event.id;
// The row that is the only tenancy evidence available: stored under
// A, carrying the channel the signature names.
let (_stored, inserted) = state
.db
.insert_event(community_a, &event, Some(channel))
.await
.expect("store the event under community A");
assert!(
inserted,
"the binding row must be freshly written, not an ON CONFLICT no-op"
);
let (_a_conn, mut a_rx) =
register_sub_in(&state, community_a, "bound-a", KIND_STREAM_MESSAGE, channel);
let (_b_conn, mut b_rx) =
register_sub_in(&state, community_b, "bound-b", KIND_STREAM_MESSAGE, channel);
// (1) Bound: delivered.
fan_out_pubsub_event_with(
&state,
ChannelEvent {
community_id: community_a,
topic: EventTopic::Channel(channel),
event: event.clone(),
},
link,
floors,
)
.await;
let delivered = event_from_ws_message(a_rx.try_recv().expect(
"an event whose stored row places it in this community, routed to the channel \
its signature names, must be delivered below the floor",
));
assert_eq!(delivered.id, event_id);
// (2) The negative that makes (1) mean something: same bytes, same
// topic, same channel UUID, different asserted tenancy.
fan_out_pubsub_event_with(
&state,
ChannelEvent {
community_id: community_b,
topic: EventTopic::Channel(channel),
event: event.clone(),
},
link,
floors,
)
.await;
assert!(
b_rx.try_recv().is_err(),
"the same signed event asserted under a community no row places it in must not \
fan out, even though that community holds a real open channel with the same \
UUID and a subscriber on it"
);
assert!(
a_rx.try_recv().is_err(),
"B's arrival must not reach A's subscriber either"
);
// `Unbound` is a decided answer, so the arrival keeps its replay
// entry and a re-send of the same bytes costs this pod neither
// another BIP-340 verify nor another Postgres round trip.
// `Unknown` withdraws it instead — see the test below.
assert!(
state
.bus_reverify_seen
.get(&(community_b, event_id.to_bytes()))
.is_some(),
"a decided Unbound refusal spends the replay id for the window"
);
drop_seeded_community(&pool, community_a).await;
drop_seeded_community(&pool, community_b).await;
}
/// The `Silent` half of the same claim: for a kind whose channel the
/// relay resolves from a stored row rather than from the signed `h`
/// tag, **the row is the entire binding** — both ways round.
///
/// `bind_scope` answers `Silent` for reactions, deletions, gift wraps,
/// `kind:9007` and the relay-authored 44100/44101, so the pure
/// comparison refuses none of them and every topic is shape-legal. That
/// is deliberate — guessing from the tag there would have cut
/// legitimate cross-pod traffic — and it means the routed *channel*
/// scope of a `Silent` kind is checked in exactly one place: the
/// `stored.channel_id == routed_channel` arm of `community_binding`.
///
/// The stream-message pair above cannot reach that arm. Its refusal is
/// decided by `bind_scope` before any lookup, and its acceptance has
/// the signed tag and the row agreeing, so the comparison could be
/// deleted and both halves would still pass. This test is the one that
/// fails if it is: one reaction, one community, two real open channels,
/// delivered under the channel its row names and refused under the
/// other.
#[tokio::test]
#[ignore = "requires Postgres"]
async fn a_relay_derived_channel_kind_is_bound_by_its_row_scope_and_nothing_else() {
let (state, pool) = super::fanout_access::postgres_state().await;
let (link, floors) = below_floor();
let stored_channel = Uuid::new_v4();
let other_channel = Uuid::new_v4();
let community = seed_community_with_open_channel(&pool, "silent", stored_channel).await;
sqlx::query(
"INSERT INTO channels (id, community_id, name, visibility, created_by) \
VALUES ($1, $2, 'bus-bound-silent-other', 'open', $3)",
)
.bind(other_channel)
.bind(community.as_uuid())
.bind(vec![7u8; 32])
.execute(&pool)
.await
.expect("seed the second open channel");
// A reaction carrying the *other* channel in its `h` tag, to make
// the point that the tag is not what binds it: the row is.
let reaction = EventBuilder::new(Kind::Custom(KIND_REACTION as u16), "+")
.tags([nostr::Tag::parse(["h", &other_channel.to_string()]).expect("h tag")])
.sign_with_keys(&Keys::generate())
.expect("sign reaction");
let reaction_id = reaction.id;
assert_eq!(
crate::handlers::bus_scope::bind_scope(
&reaction,
EventTopic::Channel(other_channel)
),
crate::handlers::bus_scope::ScopeBinding::Silent,
"test precondition: the signature must decide nothing for this kind"
);
let (_stored, inserted) = state
.db
.insert_event(community, &reaction, Some(stored_channel))
.await
.expect("store the reaction under the channel its row names");
assert!(inserted, "the binding row must be freshly written");
let (_stored_conn, mut stored_rx) = register_sub_in(
&state,
community,
"silent-stored",
KIND_REACTION,
stored_channel,
);
let (_other_conn, mut other_rx) = register_sub_in(
&state,
community,
"silent-other",
KIND_REACTION,
other_channel,
);
// Refused first: the topic the `h` tag names, which is *not* the
// scope the row carries. Order matters — the replay set is keyed on
// `(community, event_id)`, so the accepted arrival below would
// otherwise make this one a replay drop and prove nothing about
// tenancy.
fan_out_pubsub_event_with(
&state,
ChannelEvent {
community_id: community,
topic: EventTopic::Channel(other_channel),
event: reaction.clone(),
},
link,
floors,
)
.await;
assert!(
other_rx.try_recv().is_err(),
"a relay-derived-channel kind routed to a real open channel its stored row does \
not name must not fan out"
);
// The replay id is spent by that refusal, exactly as `Unbound`
// does everywhere else. Withdraw it so the accepted arrival below
// is decided by the binding rather than by the replay guard.
state
.bus_reverify_seen
.invalidate(&(community, reaction_id.to_bytes()));
fan_out_pubsub_event_with(
&state,
ChannelEvent {
community_id: community,
topic: EventTopic::Channel(stored_channel),
event: reaction,
},
link,
floors,
)
.await;
assert_eq!(
event_from_ws_message(stored_rx.try_recv().expect(
"the same reaction routed to the channel its stored row names must be \
delivered",
))
.id,
reaction_id
);
drop_seeded_community(&pool, community).await;
}
/// `Unknown` fails closed, and the only thing that tells it apart from
/// `Unbound` is the replay set.
///
/// Both refuse and both record the same `reason` label,
/// `COMMUNITY_UNBOUND` — and `MeridianBusCommunityUnbound` in
/// `deploy/charts/meridian/templates/prometheusrule.yaml` selects on
/// exactly that value, so a Postgres outage reaches an operator wearing
/// a tenancy alert's name. The distinguishing signal today is the
/// `warn!` inside `community_binding`, not the metric.
///
/// What the two arms genuinely do differently is load-bearing rather
/// than cosmetic: a decided `Unbound` spends the replay id for the
/// window, while `Unknown` withdraws it, so one Postgres blip is not a
/// ten-minute hole through which no redelivery of that event can pass.
/// This test pins that difference, which is also what proves the
/// refusal came from the tenancy check and not from the visibility gate
/// failing closed one step later.
///
/// Neither half needs a live database. `Unbound` is reached without one
/// by an ephemeral kind, which `community_binding` answers from the
/// kind registry with no round trip; `Unknown` is reached by pinning
/// the pool at an address nothing listens on — a property this state
/// asserts rather than inherits, because `test_config` resolves a
/// *live* `MERIDIAN_TEST_DATABASE_URL` whenever the gate sets one.
#[tokio::test]
async fn unknown_tenancy_refuses_and_withdraws_the_replay_id_that_unbound_spends() {
let state = super::fanout_access::unreachable_db_state().await;
let (link, floors) = below_floor();
let community = meridian_core::tenant::CommunityId::from_uuid(Uuid::nil());
// Unknown: a stored kind, so the lookup runs — and cannot answer.
let channel = Uuid::new_v4();
let (_conn, mut rx) = register_stream_sub(&state, "unknown", Some(channel));
let event = stream_message(channel);
let event_id = event.id;
fan_out_pubsub_event_with(
&state,
ChannelEvent {
community_id: community,
topic: EventTopic::Channel(channel),
event,
},
link,
floors,
)
.await;
assert!(
rx.try_recv().is_err(),
"a tenancy question the system of record cannot answer must refuse, not accept"
);
assert!(
state
.bus_reverify_seen
.get(&(community, event_id.to_bytes()))
.is_none(),
"an Unknown refusal must withdraw the replay id, or one Postgres blip becomes a \
ten-minute hole for that event"
);
// Unbound, for contrast, on the same state: ephemeral, so the kind
// registry answers and no round trip is spent.
let (_presence_conn, mut presence_rx) = register_presence_sub(&state, "unbound");
let presence = presence_event("online");
let presence_id = presence.id;
fan_out_pubsub_event_with(
&state,
ChannelEvent {
community_id: community,
topic: EventTopic::Global,
event: presence,
},
link,
floors,
)
.await;
assert!(
presence_rx.try_recv().is_err(),
"an event no row can ever place here must not fan out below the floor"
);
assert!(
state
.bus_reverify_seen
.get(&(community, presence_id.to_bytes()))
.is_some(),
"a decided Unbound refusal keeps its replay id"
);
}
/// Meridian's own Redis, or a loud failure. Never a silent skip, and
/// never `127.0.0.1:6379`.
///
@ -3436,6 +3865,17 @@ mod tests {
let mut config = test_config();
config.redis_url = redis_url.to_string();
let pool = sqlx::PgPool::connect_lazy(&config.database_url).expect("lazy pg pool");
state_over(config, pool).await
}
/// Assemble an [`AppState`] over an already-resolved pool.
///
/// Split out of [`test_state_with_redis_url`] so the Postgres-backed
/// builder below shares one construction with the infra-free one. The
/// two must differ **only** in how the pool was obtained: a test that
/// reaches the database is worth nothing if it is exercising a
/// differently-assembled state from the module it belongs to.
async fn state_over(config: crate::config::Config, pool: sqlx::PgPool) -> Arc<AppState> {
let db = meridian_db::Db::from_pool(pool.clone());
let redis_pool = deadpool_redis::Config::from_url(&config.redis_url)
.create_pool(Some(deadpool_redis::Runtime::Tokio1))
@ -3473,6 +3913,62 @@ mod tests {
test_state_with_redis_url("redis://127.0.0.1:1").await
}
/// State over a **live** disposable Postgres, plus the pool that seeded
/// it, for the tests whose subject is a question put to the system of
/// record.
///
/// Fails loudly rather than skipping, and that is the whole point.
/// [`test_config`] resolves `MERIDIAN_TEST_DATABASE_URL` when it is set
/// and falls back to an address nothing listens on when it is not,
/// which is the right posture for the fan-out tests that never query —
/// but a test whose subject *is* the query must not report PASS having
/// asserted nothing, which is exactly how `api::operator`'s twelve
/// tests passed against no database at all (meridian-wxqw). Callers are
/// `#[ignore = "requires Postgres"]` and run under
/// `just test-relay-db-run`, where a missing database is an error.
///
/// Redis stays deliberately unconnectable: the cross-node *receive*
/// path publishes nothing, so a live Redis would buy no coverage while
/// pulling these tests into `meridian-h0tg`'s non-exiting test binary.
pub(super) async fn postgres_state() -> (Arc<AppState>, sqlx::PgPool) {
let mut config = test_config();
config.redis_url = "redis://127.0.0.1:1".to_string();
config.database_url =
meridian_test_db::disposable_url().unwrap_or_else(|e| panic!("{e}"));
let pool = sqlx::PgPool::connect(&config.database_url)
.await
.unwrap_or_else(|e| {
panic!(
"connecting to the disposable test database named by \
MERIDIAN_TEST_DATABASE_URL failed: {e}"
)
});
(state_over(config, pool.clone()).await, pool)
}
/// State whose pool can never connect, whatever the environment says.
///
/// [`test_state`]'s pool is unreachable on a developer's box and
/// *reachable* inside `just test-relay-db-run`, because `test_config`
/// honours `MERIDIAN_TEST_DATABASE_URL`. A test whose subject is the
/// database-**error** verdict needs the same answer in both, so it pins
/// the address rather than inheriting one.
///
/// The acquire timeout is pinned too. `sqlx` retries a failed connect
/// until the acquire deadline, so on the default 30 s a test that only
/// wants one `Err` back spends half a minute re-proving that nothing is
/// listening on port 1.
pub(super) async fn unreachable_db_state() -> Arc<AppState> {
let mut config = test_config();
config.redis_url = "redis://127.0.0.1:1".to_string();
config.database_url = "postgres://127.0.0.1:1/meridian_test".to_string();
let pool = sqlx::postgres::PgPoolOptions::new()
.acquire_timeout(std::time::Duration::from_secs(2))
.connect_lazy(&config.database_url)
.expect("lazy pg pool");
state_over(config, pool).await
}
/// Real-PG, real-Redis state that hands back the audit shutdown handle so
/// a test can drain queued audit entries before asserting on `audit_log`.
/// `None` when Postgres or Redis is unavailable (test skips).