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
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:
parent
c5afeed1c1
commit
a9b30ea276
1 changed files with 497 additions and 1 deletions
|
|
@ -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).
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue