mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
perf(relay): cache relay-membership checks off the writer pool
check_relay_membership runs a writer-pool SELECT on every authenticated HTTP request and WS AUTH — after the replica read routing shipped, this is the largest remaining read volume on the bb-public writer. Add a 10s-TTL moka cache (mirroring the existing channel-membership cache) in front of is_relay_member for the two hot callsites, with cross-pod Redis invalidation via a new CacheInvalidation::RelayMembership variant. All four membership-mutation flows drop the affected key: - v2 invite claim (Joined outcome) - v1 invite claim (was_inserted) - admin add/remove (kind 9030/9031) - NIP-43 self-leave Negative results are cached too, and invalidated on join, so a fresh member is admitted immediately rather than after TTL expiry. The two invite-flow verification reads stay on the direct DB path (pre-mutation, staleness unacceptable). A missed cross-pod publish degrades to the 10s TTL, same contract as the existing caches. Verified: cargo test -p buzz-relay (803 passed) and -p buzz-pubsub (25 passed); clippy -D warnings clean. Co-authored-by: tlongwell-block <109685178+tlongwell-block@users.noreply.github.com> Signed-off-by: tlongwell-block <109685178+tlongwell-block@users.noreply.github.com>
This commit is contained in:
co-authored by
tlongwell-block
parent
468647a51f
commit
223cfe5903
@@ -76,6 +76,13 @@ pub enum CacheInvalidation {
|
||||
/// Drop all membership / accessible / visibility caches. Mirrors
|
||||
/// `invalidate_channel_deleted`.
|
||||
ChannelDeleted,
|
||||
/// Drop one pubkey's relay-membership entry. Mirrors
|
||||
/// `invalidate_relay_membership` (invite claim, admin add/remove,
|
||||
/// self-leave).
|
||||
RelayMembership {
|
||||
/// Affected member's pubkey bytes.
|
||||
pubkey: Vec<u8>,
|
||||
},
|
||||
}
|
||||
|
||||
/// A cache invalidation received from a community-scoped Redis channel.
|
||||
@@ -232,6 +239,18 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn relay_membership_roundtrips_through_json() {
|
||||
let msg = CacheInvalidation::RelayMembership {
|
||||
pubkey: vec![5, 6, 7, 8],
|
||||
};
|
||||
let json = serde_json::to_string(&msg).unwrap();
|
||||
assert_eq!(
|
||||
serde_json::from_str::<CacheInvalidation>(&json).unwrap(),
|
||||
msg
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unit_variants_roundtrip_through_json() {
|
||||
for msg in [
|
||||
|
||||
@@ -406,6 +406,9 @@ pub async fn claim_invite(
|
||||
member = %claimer_hex,
|
||||
"relay member added via v2 invite"
|
||||
);
|
||||
// Drop any cached negative membership so the new member is
|
||||
// admitted immediately, not after cache TTL.
|
||||
state.invalidate_relay_membership(&tenant, &claimer_hex);
|
||||
// NIP-43 side effects only on Joined, never on other outcomes.
|
||||
if let Err(e) = publish_nip43_member_added(&tenant, &state, &claimer_hex).await {
|
||||
tracing::warn!(
|
||||
@@ -484,6 +487,9 @@ pub async fn claim_invite(
|
||||
member = %claimer_hex,
|
||||
"relay member added via invite"
|
||||
);
|
||||
// Drop any cached negative membership so the new member is admitted
|
||||
// immediately, not after cache TTL.
|
||||
state.invalidate_relay_membership(&tenant, &claimer_hex);
|
||||
if let Err(e) = publish_nip43_member_added(&tenant, &state, &claimer_hex).await {
|
||||
tracing::warn!("failed to publish NIP-43 member-added delta after claim: {e}");
|
||||
}
|
||||
|
||||
@@ -70,8 +70,7 @@ pub mod relay_members {
|
||||
|
||||
let pubkey_hex = hex::encode(pubkey_bytes);
|
||||
let is_member = state
|
||||
.db
|
||||
.is_relay_member(community, &pubkey_hex)
|
||||
.is_relay_member_cached(community, &pubkey_hex)
|
||||
.await
|
||||
.map_err(|e| format!("relay membership check failed: {e}"))?;
|
||||
if is_member {
|
||||
@@ -87,8 +86,7 @@ pub mod relay_members {
|
||||
Ok(owner_pubkey) => {
|
||||
let owner_hex = owner_pubkey.to_hex();
|
||||
let owner_is_member = state
|
||||
.db
|
||||
.is_relay_member(community, &owner_hex)
|
||||
.is_relay_member_cached(community, &owner_hex)
|
||||
.await
|
||||
.map_err(|e| format!("relay membership check (owner) failed: {e}"))?;
|
||||
if owner_is_member {
|
||||
|
||||
@@ -1981,6 +1981,9 @@ async fn ingest_event_inner(
|
||||
}
|
||||
|
||||
// Publish NIP-43 announcements — fire-and-forget.
|
||||
// Drop the cached positive membership so the leave takes effect
|
||||
// fleet-wide now, not after cache TTL.
|
||||
state.invalidate_relay_membership(tenant, &sender_hex);
|
||||
if let Err(e) =
|
||||
crate::handlers::side_effects::publish_nip43_member_removed(tenant, state, &sender_hex)
|
||||
.await
|
||||
|
||||
@@ -297,6 +297,9 @@ async fn execute_relay_admin_command(
|
||||
// Only publish NIP-43 announcements when the row was actually inserted —
|
||||
// skip on no-op re-adds to avoid spurious kind:8000 events.
|
||||
if was_inserted {
|
||||
// Drop any cached negative membership so the new member is
|
||||
// admitted immediately, not after cache TTL.
|
||||
state.invalidate_relay_membership(tenant, &target_hex);
|
||||
if let Err(e) = publish_nip43_member_added(tenant, state, &target_hex).await {
|
||||
warn!(error = %e, "failed to publish NIP-43 member added event");
|
||||
}
|
||||
@@ -357,6 +360,10 @@ async fn execute_relay_admin_command(
|
||||
"relay member removed"
|
||||
);
|
||||
|
||||
// Drop the cached positive membership so removal takes effect
|
||||
// fleet-wide now, not after cache TTL.
|
||||
state.invalidate_relay_membership(tenant, &target_hex);
|
||||
|
||||
if let Err(e) = publish_nip43_member_removed(tenant, state, &target_hex).await {
|
||||
warn!(error = %e, "failed to publish NIP-43 member removed event");
|
||||
}
|
||||
|
||||
@@ -543,6 +543,11 @@ pub struct AppState {
|
||||
/// Short TTL (10s) — membership changes are rare but must propagate.
|
||||
#[allow(clippy::type_complexity)]
|
||||
pub membership_cache: Arc<moka::sync::Cache<(CommunityId, Uuid, Vec<u8>), bool>>,
|
||||
/// Relay membership cache: (community_id, pubkey_bytes) → is_relay_member.
|
||||
/// Short TTL (10s) — checked on every authenticated HTTP request and WS
|
||||
/// AUTH, so this keeps the hot auth path off the writer pool. Invalidated
|
||||
/// on invite claim, admin add/remove, and self-leave.
|
||||
pub relay_membership_cache: Arc<moka::sync::Cache<(CommunityId, Vec<u8>), bool>>,
|
||||
/// Accessible channel IDs cache: (community_id, pubkey_bytes) → channel UUIDs.
|
||||
/// Short TTL (10s) — invalidated on membership or channel visibility changes.
|
||||
#[allow(clippy::type_complexity)]
|
||||
@@ -745,6 +750,13 @@ impl AppState {
|
||||
.support_invalidation_closures()
|
||||
.build(),
|
||||
),
|
||||
relay_membership_cache: Arc::new(
|
||||
moka::sync::Cache::builder()
|
||||
.max_capacity(100_000)
|
||||
.time_to_live(std::time::Duration::from_secs(10))
|
||||
.support_invalidation_closures()
|
||||
.build(),
|
||||
),
|
||||
accessible_channels_cache: Arc::new(
|
||||
moka::sync::Cache::builder()
|
||||
.max_capacity(10_000)
|
||||
@@ -842,6 +854,57 @@ impl AppState {
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
/// Check relay membership with a 10-second cache. Falls back to DB on miss.
|
||||
///
|
||||
/// Runs on every authenticated HTTP request and WS AUTH (see
|
||||
/// `check_relay_membership`), so a hit removes a writer-pool SELECT from
|
||||
/// the hottest auth path in the relay. Negative results are cached too:
|
||||
/// they are dropped by `invalidate_relay_membership` when the pubkey joins
|
||||
/// (invite claim, admin add), so a fresh member is admitted immediately
|
||||
/// rather than after TTL expiry.
|
||||
pub async fn is_relay_member_cached(
|
||||
&self,
|
||||
community_id: CommunityId,
|
||||
pubkey_hex: &str,
|
||||
) -> Result<bool, buzz_db::DbError> {
|
||||
let key = (community_id, pubkey_hex.as_bytes().to_vec());
|
||||
if let Some(cached) = self.relay_membership_cache.get(&key) {
|
||||
metrics::counter!("buzz_relay_membership_cache_hits_total").increment(1);
|
||||
return Ok(cached);
|
||||
}
|
||||
metrics::counter!("buzz_relay_membership_cache_misses_total").increment(1);
|
||||
let result = self.db.is_relay_member(community_id, pubkey_hex).await?;
|
||||
self.relay_membership_cache.insert(key, result);
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
/// Invalidate one pubkey's relay-membership entry after a membership
|
||||
/// change (invite claim, admin add/remove, self-leave).
|
||||
///
|
||||
/// Local drop plus fire-and-forget cross-pod publish, same contract as
|
||||
/// [`invalidate_membership`]: a dropped publish degrades to the 10s TTL,
|
||||
/// never a permanent leak.
|
||||
pub fn invalidate_relay_membership(&self, tenant: &TenantContext, pubkey_hex: &str) {
|
||||
self.invalidate_relay_membership_local(tenant.community(), pubkey_hex.as_bytes());
|
||||
self.spawn_cache_invalidation(
|
||||
tenant,
|
||||
CacheInvalidation::RelayMembership {
|
||||
pubkey: pubkey_hex.as_bytes().to_vec(),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
/// Local-only relay-membership drop. The cross-pod consumer calls this
|
||||
/// directly so applying a received drop never re-publishes it.
|
||||
pub(crate) fn invalidate_relay_membership_local(
|
||||
&self,
|
||||
community_id: CommunityId,
|
||||
pubkey: &[u8],
|
||||
) {
|
||||
self.relay_membership_cache
|
||||
.invalidate(&(community_id, pubkey.to_vec()));
|
||||
}
|
||||
|
||||
/// Invalidate caches after a membership change (add/remove member).
|
||||
///
|
||||
/// Drops the local moka entries AND fire-and-forget publishes the same drop
|
||||
@@ -996,6 +1059,9 @@ impl AppState {
|
||||
CacheInvalidation::ChannelDeleted => {
|
||||
self.invalidate_channel_deleted_local(community_id);
|
||||
}
|
||||
CacheInvalidation::RelayMembership { pubkey } => {
|
||||
self.invalidate_relay_membership_local(community_id, &pubkey);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1477,6 +1543,81 @@ mod tests {
|
||||
assert_eq!(mgr.pubkey_for_conn(Uuid::new_v4()), None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn relay_membership_cache_hit_and_scoped_invalidation() {
|
||||
let state = test_state().await;
|
||||
let community_a = CommunityId::from_uuid(Uuid::from_u128(0xAAAA));
|
||||
let community_b = CommunityId::from_uuid(Uuid::from_u128(0xBBBB));
|
||||
let member_hex = "aa".repeat(32);
|
||||
let other_hex = "bb".repeat(32);
|
||||
|
||||
// Seed as the cached method would after DB reads.
|
||||
state
|
||||
.relay_membership_cache
|
||||
.insert((community_a, member_hex.as_bytes().to_vec()), true);
|
||||
state
|
||||
.relay_membership_cache
|
||||
.insert((community_a, other_hex.as_bytes().to_vec()), false);
|
||||
state
|
||||
.relay_membership_cache
|
||||
.insert((community_b, member_hex.as_bytes().to_vec()), true);
|
||||
|
||||
// A cached entry is served without touching the DB: the lazy test
|
||||
// pool points at nothing, so a DB fallback would error, not return.
|
||||
assert!(state
|
||||
.is_relay_member_cached(community_a, &member_hex)
|
||||
.await
|
||||
.expect("cache hit must not touch the DB"),);
|
||||
// Negative entries are cached and served the same way.
|
||||
assert!(!state
|
||||
.is_relay_member_cached(community_a, &other_hex)
|
||||
.await
|
||||
.expect("negative cache hit must not touch the DB"),);
|
||||
|
||||
// Invalidation drops exactly the (community, pubkey) entry: same
|
||||
// pubkey in another community and other pubkeys are untouched.
|
||||
state.invalidate_relay_membership_local(community_a, member_hex.as_bytes());
|
||||
assert_eq!(
|
||||
state
|
||||
.relay_membership_cache
|
||||
.get(&(community_a, member_hex.as_bytes().to_vec())),
|
||||
None,
|
||||
"invalidated entry must be dropped"
|
||||
);
|
||||
assert_eq!(
|
||||
state
|
||||
.relay_membership_cache
|
||||
.get(&(community_b, member_hex.as_bytes().to_vec())),
|
||||
Some(true),
|
||||
"same pubkey in another community must survive"
|
||||
);
|
||||
assert_eq!(
|
||||
state
|
||||
.relay_membership_cache
|
||||
.get(&(community_a, other_hex.as_bytes().to_vec())),
|
||||
Some(false),
|
||||
"other pubkeys in the same community must survive"
|
||||
);
|
||||
|
||||
// The cross-pod path applies the same local drop.
|
||||
state
|
||||
.relay_membership_cache
|
||||
.insert((community_a, member_hex.as_bytes().to_vec()), true);
|
||||
state.apply_cache_invalidation(
|
||||
community_a,
|
||||
buzz_pubsub::cache_invalidation::CacheInvalidation::RelayMembership {
|
||||
pubkey: member_hex.as_bytes().to_vec(),
|
||||
},
|
||||
);
|
||||
assert_eq!(
|
||||
state
|
||||
.relay_membership_cache
|
||||
.get(&(community_a, member_hex.as_bytes().to_vec())),
|
||||
None,
|
||||
"cross-pod drop must clear the entry"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn accessible_channel_invalidation_is_scoped_to_community() {
|
||||
let state = test_state().await;
|
||||
|
||||
Reference in New Issue
Block a user