mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
Merge origin/main into eva/serverless-elasticache-compat
Brings in #4124: relay-membership checks served from the read replica via bounded routed read. No conflicts. Co-authored-by: npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d <011987e296fd5006292d2f930b574be47c7801048d1983c46c425d3c95f0cffd@buzz.block.builderlab.xyz> Signed-off-by: npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d <011987e296fd5006292d2f930b574be47c7801048d1983c46c425d3c95f0cffd@buzz.block.builderlab.xyz>
This commit is contained in:
commit
fc200bf453
@@ -3999,8 +3999,33 @@ impl Db {
|
||||
}
|
||||
|
||||
/// Returns `true` if `pubkey` (64-char hex) is a member of `community`.
|
||||
///
|
||||
/// Replica-routed on the bounded arm — the one PERMISSION read routed by
|
||||
/// explicit product decision (bounded-stale membership beats the 10s
|
||||
/// cache it replaced). Admits and revokes may lag by at most the budget
|
||||
/// `B`; everything else fails closed to the writer, exactly like
|
||||
/// [`Db::query_events_routed_bounded`]. Not precedent for routing other
|
||||
/// permission reads.
|
||||
pub async fn is_relay_member(&self, community: CommunityId, pubkey: &str) -> Result<bool> {
|
||||
relay_members::is_relay_member(&self.pool, community, pubkey).await
|
||||
let path = "relay_membership";
|
||||
match self.route_read(path, RoutePredicate::Bounded).await {
|
||||
RouteDecision::Replica(mut tx, _entry, reason) => {
|
||||
match relay_members::is_relay_member_on(&mut tx, community, pubkey).await {
|
||||
Ok(is_member) => {
|
||||
Self::record_route(path, "replica", reason);
|
||||
Ok(is_member)
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(path, "replica read failed; re-running on writer: {e}");
|
||||
Self::record_route(path, "writer", "replica_error");
|
||||
relay_members::is_relay_member(&self.pool, community, pubkey).await
|
||||
}
|
||||
}
|
||||
}
|
||||
RouteDecision::Writer => {
|
||||
relay_members::is_relay_member(&self.pool, community, pubkey).await
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns the relay member record for `pubkey` in `community`, or `None` if not found.
|
||||
@@ -7414,6 +7439,78 @@ mod tests {
|
||||
drop_scratch_db(&admin, writer, &wname).await;
|
||||
}
|
||||
|
||||
/// Routed relay-membership check: budget unset ⇒ writer; budget set +
|
||||
/// fresh proved entry ⇒ replica (bounded arm); over-budget entry ⇒
|
||||
/// writer. Divergent membership rows prove which pool answered.
|
||||
#[tokio::test]
|
||||
#[ignore = "requires Postgres"]
|
||||
async fn is_relay_member_is_bounded_routed_and_fails_closed() {
|
||||
let admin = PgPool::connect(&admin_url().await)
|
||||
.await
|
||||
.expect("connect admin");
|
||||
let (writer, wname) = create_scratch_db(&admin, "mem_w").await;
|
||||
let (replica, rname) = create_scratch_db(&admin, "mem_r").await;
|
||||
|
||||
let community = Uuid::new_v4();
|
||||
for pool in [&writer, &replica] {
|
||||
sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)")
|
||||
.bind(community)
|
||||
.bind(format!("member-routing-{}.example", community.simple()))
|
||||
.execute(pool)
|
||||
.await
|
||||
.expect("insert community");
|
||||
}
|
||||
let cid = CommunityId::from_uuid(community);
|
||||
let writer_only = "aa".repeat(32);
|
||||
let replica_only = "bb".repeat(32);
|
||||
relay_members::add_relay_member(&writer, cid, &writer_only, "member", None)
|
||||
.await
|
||||
.expect("seed writer member");
|
||||
relay_members::add_relay_member(&replica, cid, &replica_only, "member", None)
|
||||
.await
|
||||
.expect("seed replica member");
|
||||
|
||||
let mut db = Db::from_pools(writer.clone(), replica.clone());
|
||||
db.fence().force_open_for_tests(chrono::Utc::now());
|
||||
|
||||
// Budget unset ⇒ bounded arm disabled ⇒ writer.
|
||||
assert!(
|
||||
db.is_relay_member(cid, &writer_only)
|
||||
.await
|
||||
.expect("gate off"),
|
||||
"budget unset must answer from the writer"
|
||||
);
|
||||
assert!(!db.is_relay_member(cid, &replica_only).await.unwrap());
|
||||
|
||||
// Budget set + fresh entry ⇒ replica.
|
||||
db.set_replica_read_max_age_for_tests(Some(std::time::Duration::from_secs(5)));
|
||||
assert!(
|
||||
db.is_relay_member(cid, &replica_only)
|
||||
.await
|
||||
.expect("gate on"),
|
||||
"budget set must answer from the replica"
|
||||
);
|
||||
assert!(!db.is_relay_member(cid, &writer_only).await.unwrap());
|
||||
|
||||
// Entry older than the budget ⇒ fail closed to the writer. Close
|
||||
// first so no prior fresh entry can be the one proved (matches the
|
||||
// count test; today `force_open_for_tests_at` also clears the ring).
|
||||
db.fence().close();
|
||||
db.fence().force_open_for_tests_at(
|
||||
chrono::Utc::now(),
|
||||
std::time::Instant::now() - std::time::Duration::from_secs(10),
|
||||
);
|
||||
assert!(
|
||||
db.is_relay_member(cid, &writer_only)
|
||||
.await
|
||||
.expect("entry too old"),
|
||||
"an over-budget entry must fail closed to the writer"
|
||||
);
|
||||
|
||||
drop_scratch_db(&admin, replica, &rname).await;
|
||||
drop_scratch_db(&admin, writer, &wname).await;
|
||||
}
|
||||
|
||||
/// Community separation across every routed seam, verified on
|
||||
/// REPLICA-SERVED reads.
|
||||
///
|
||||
|
||||
@@ -29,10 +29,22 @@ pub struct RelayMember {
|
||||
|
||||
/// Returns `true` if `pubkey` (64-char hex) is a member of `community`.
|
||||
pub async fn is_relay_member(pool: &PgPool, community: CommunityId, pubkey: &str) -> Result<bool> {
|
||||
let mut conn = pool.acquire().await?;
|
||||
is_relay_member_on(&mut conn, community, pubkey).await
|
||||
}
|
||||
|
||||
/// [`is_relay_member`] on a specific session — the replica-routing path runs
|
||||
/// the lookup on the exact reader connection whose heartbeat observation
|
||||
/// proved fence coverage.
|
||||
pub(crate) async fn is_relay_member_on(
|
||||
conn: &mut sqlx::PgConnection,
|
||||
community: CommunityId,
|
||||
pubkey: &str,
|
||||
) -> Result<bool> {
|
||||
let row = sqlx::query("SELECT 1 FROM relay_members WHERE community_id = $1 AND pubkey = $2")
|
||||
.bind(community.as_uuid())
|
||||
.bind(pubkey)
|
||||
.fetch_optional(pool)
|
||||
.fetch_optional(conn)
|
||||
.await?;
|
||||
Ok(row.is_some())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user