mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
+1








14fba21e57
Signed-off-by: tlongwell-block <109685178+tlongwell-block@users.noreply.github.com> Signed-off-by: npub1jh9wn95s0472h86ahapupaf7m6kx4v9sx2n0atj2hltcfer8k06s5n3pyf <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co> Signed-off-by: Tyler Longwell <tlongwell@block.xyz> Signed-off-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@sprout-oss.stage.blox.sqprod.co> Signed-off-by: npub17jjz49l9jjmhhk7cac63j8yt9z555n9cw8vk7v5jz4vzw4ppld5qgj57cc <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Eva <011987e296fd5006292d2f930b574be47c7801048d1983c46c425d3c95f0cffd@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Mari <95cae996907d7cab9f5dbf43c0f53edeac6ab0b032a6feae4abfd784e467b3f5@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Sami <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Max <d8473ee32b973aa31a21a65adddcc4b69cc2a8a4dee8121ecd51926e0cddbc02@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Quinn <96f056ad5f2305c8ddf637dc65d048aa4c12d7daeb8867690e34fca46b0ef64c@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Dawn <c6237ef84fa537c78dcee78efd2d4e59f728859c7f194da42ac51ededfa0be05@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Tyler Longwell <tlongwell@block.xyz> Co-authored-by: Sami <sami@sprout-oss.stage.blox.sqprod.co> Co-authored-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@sprout-oss.stage.blox.sqprod.co>
252 lines
8.6 KiB
Rust
252 lines
8.6 KiB
Rust
//! Cross-pod cache-key invalidation over Redis pub/sub.
|
|
//!
|
|
//! Each relay pod keeps in-memory (moka) membership / accessible-channels /
|
|
//! visibility caches. A membership or visibility change is applied to the local
|
|
//! caches only on the pod that processed the write; other pods would otherwise
|
|
//! rely on the 10s TTL to expire stale entries. This module carries the same
|
|
//! key drops to every pod immediately.
|
|
//!
|
|
//! The message is a pure cache-key drop — never an "evict these subscriptions"
|
|
//! payload. The per-event access gate (`filter_fanout_by_access`) is the
|
|
//! universal delivery-enforcement point, so dropping the stale key is
|
|
//! sufficient: the next read re-fetches authoritative state from the DB.
|
|
|
|
use buzz_core::{CommunityId, TenantContext};
|
|
use futures_util::StreamExt;
|
|
use serde::{Deserialize, Serialize};
|
|
use tokio::sync::broadcast;
|
|
use uuid::Uuid;
|
|
|
|
use crate::topic::BUZZ_PREFIX;
|
|
|
|
/// Tenant-local Redis pub/sub channel suffix for cache-invalidation messages.
|
|
pub const CACHE_INVALIDATION_SUFFIX: &str = "cache-invalidate";
|
|
|
|
/// Pattern used by the subscriber to receive cache invalidations for all
|
|
/// communities this pod may have cached locally.
|
|
pub const CACHE_INVALIDATION_PATTERN: &str = "buzz:*:cache-invalidate";
|
|
|
|
/// Redis pub/sub channel for cache-invalidation messages under `ctx`.
|
|
pub fn cache_invalidation_channel(ctx: &TenantContext) -> String {
|
|
format!(
|
|
"{BUZZ_PREFIX}:{}:{CACHE_INVALIDATION_SUFFIX}",
|
|
ctx.community()
|
|
)
|
|
}
|
|
|
|
/// Parse a cache-invalidation Redis channel into its scoped community id.
|
|
pub fn parse_cache_invalidation_channel(channel: &str) -> Option<CommunityId> {
|
|
let mut parts = channel.split(':');
|
|
if parts.next()? != BUZZ_PREFIX {
|
|
return None;
|
|
}
|
|
let community_id = Uuid::parse_str(parts.next()?).ok()?;
|
|
if parts.next()? != CACHE_INVALIDATION_SUFFIX {
|
|
return None;
|
|
}
|
|
if parts.next().is_some() {
|
|
return None;
|
|
}
|
|
Some(CommunityId::from_uuid(community_id))
|
|
}
|
|
|
|
/// A cache-key drop to apply on every pod. Each variant mirrors exactly one of
|
|
/// the relay's local `invalidate_*` operations. The community is carried by
|
|
/// [`ScopedCacheInvalidation`], not by the tenant-local operation.
|
|
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
|
#[serde(tag = "op")]
|
|
pub enum CacheInvalidation {
|
|
/// Drop the `(channel_id, pubkey)` membership entry and the user's
|
|
/// accessible-channels entry. Mirrors `invalidate_membership`.
|
|
Membership {
|
|
/// Channel whose membership changed.
|
|
channel_id: Uuid,
|
|
/// Affected member's pubkey bytes.
|
|
pubkey: Vec<u8>,
|
|
},
|
|
/// Drop every user's accessible-channels entry. Mirrors
|
|
/// `invalidate_all_accessible_channels` (e.g. a new open channel).
|
|
AccessibleAll,
|
|
/// Drop the cached visibility for a single channel. Mirrors
|
|
/// `invalidate_channel_visibility` (e.g. an open→private flip).
|
|
Visibility {
|
|
/// Channel whose visibility changed.
|
|
channel_id: Uuid,
|
|
},
|
|
/// Drop all membership / accessible / visibility caches. Mirrors
|
|
/// `invalidate_channel_deleted`.
|
|
ChannelDeleted,
|
|
}
|
|
|
|
/// A cache invalidation received from a community-scoped Redis channel.
|
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
|
pub struct ScopedCacheInvalidation {
|
|
/// Community whose local cache key should be dropped.
|
|
pub community_id: CommunityId,
|
|
/// Tenant-local cache invalidation operation.
|
|
pub invalidation: CacheInvalidation,
|
|
}
|
|
|
|
/// Initial reconnect backoff (1 second).
|
|
const BACKOFF_INITIAL_SECS: u64 = 1;
|
|
/// Maximum reconnect backoff (30 seconds).
|
|
const BACKOFF_MAX_SECS: u64 = 30;
|
|
|
|
/// Subscribes to `buzz:*:cache-invalidate` and forwards scoped drops to the broadcast.
|
|
///
|
|
/// Mirrors `subscriber::run_subscriber`: a reconnect loop with exponential
|
|
/// backoff (1s → 2s → 4s → … → 30s max). Never returns — runs for the lifetime
|
|
/// of the relay.
|
|
pub async fn run_cache_invalidation_subscriber(
|
|
redis_url: String,
|
|
broadcast_tx: broadcast::Sender<ScopedCacheInvalidation>,
|
|
) {
|
|
let mut backoff_secs = BACKOFF_INITIAL_SECS;
|
|
|
|
loop {
|
|
match connect_and_subscribe(&redis_url, &broadcast_tx).await {
|
|
Ok(()) => {
|
|
backoff_secs = BACKOFF_INITIAL_SECS;
|
|
tracing::warn!(
|
|
"Redis cache-invalidation stream ended (clean disconnect) — reconnecting in {backoff_secs}s"
|
|
);
|
|
}
|
|
Err(e) => {
|
|
tracing::error!(
|
|
"Redis cache-invalidation error: {e} — reconnecting in {backoff_secs}s"
|
|
);
|
|
}
|
|
}
|
|
|
|
tokio::time::sleep(tokio::time::Duration::from_secs(backoff_secs)).await;
|
|
backoff_secs = (backoff_secs * 2).min(BACKOFF_MAX_SECS);
|
|
|
|
tracing::info!("Attempting to reconnect to Redis cache-invalidation...");
|
|
}
|
|
}
|
|
|
|
async fn connect_and_subscribe(
|
|
redis_url: &str,
|
|
broadcast_tx: &broadcast::Sender<ScopedCacheInvalidation>,
|
|
) -> Result<(), redis::RedisError> {
|
|
let client = redis::Client::open(redis_url)?;
|
|
let mut conn = client.get_async_pubsub().await?;
|
|
|
|
conn.psubscribe(CACHE_INVALIDATION_PATTERN).await?;
|
|
|
|
tracing::info!(
|
|
"Redis cache-invalidation subscriber connected — listening on {CACHE_INVALIDATION_PATTERN}"
|
|
);
|
|
|
|
let mut stream = conn.on_message();
|
|
while let Some(msg) = stream.next().await {
|
|
let channel = msg.get_channel_name();
|
|
let Some(community_id) = parse_cache_invalidation_channel(channel) else {
|
|
tracing::warn!("Received cache-invalidation message on unexpected channel: {channel}");
|
|
continue;
|
|
};
|
|
|
|
let payload: String = match msg.get_payload() {
|
|
Ok(p) => p,
|
|
Err(e) => {
|
|
tracing::warn!("Failed to get cache-invalidation payload: {e}");
|
|
continue;
|
|
}
|
|
};
|
|
|
|
let invalidation: CacheInvalidation = match serde_json::from_str(&payload) {
|
|
Ok(v) => v,
|
|
Err(e) => {
|
|
tracing::warn!("Failed to deserialize cache-invalidation message: {e}");
|
|
continue;
|
|
}
|
|
};
|
|
|
|
let scoped = ScopedCacheInvalidation {
|
|
community_id,
|
|
invalidation,
|
|
};
|
|
|
|
if broadcast_tx.send(scoped).is_err() {
|
|
tracing::trace!("No cache-invalidation receivers — message dropped");
|
|
}
|
|
}
|
|
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
fn ctx(id: u128, host: &str) -> TenantContext {
|
|
TenantContext::resolved(CommunityId::from_uuid(Uuid::from_u128(id)), host)
|
|
}
|
|
|
|
#[test]
|
|
fn cache_invalidation_channel_is_community_scoped() {
|
|
let community_a = ctx(0xaaaa, "a.example");
|
|
let community_b = ctx(0xbbbb, "b.example");
|
|
|
|
assert_eq!(
|
|
cache_invalidation_channel(&community_a),
|
|
format!("buzz:{}:cache-invalidate", community_a.community())
|
|
);
|
|
assert_ne!(
|
|
cache_invalidation_channel(&community_a),
|
|
cache_invalidation_channel(&community_b)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn parses_cache_invalidation_channel() {
|
|
let community_id = CommunityId::from_uuid(Uuid::from_u128(0xaaaa));
|
|
let raw = format!("buzz:{community_id}:cache-invalidate");
|
|
|
|
assert_eq!(parse_cache_invalidation_channel(&raw), Some(community_id));
|
|
}
|
|
|
|
#[test]
|
|
fn rejects_bad_cache_invalidation_channels() {
|
|
for raw in [
|
|
"buzz:cache-invalidate",
|
|
"buzz:not-a-uuid:cache-invalidate",
|
|
"not-buzz:00000000-0000-0000-0000-00000000aaaa:cache-invalidate",
|
|
"buzz:00000000-0000-0000-0000-00000000aaaa:cache-invalidate:extra",
|
|
"buzz:00000000-0000-0000-0000-00000000aaaa:channel:00000000-0000-0000-0000-00000000bbbb",
|
|
] {
|
|
assert_eq!(parse_cache_invalidation_channel(raw), None);
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn membership_roundtrips_through_json() {
|
|
let msg = CacheInvalidation::Membership {
|
|
channel_id: Uuid::from_u128(0x1234),
|
|
pubkey: vec![1, 2, 3, 4],
|
|
};
|
|
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 [
|
|
CacheInvalidation::AccessibleAll,
|
|
CacheInvalidation::ChannelDeleted,
|
|
CacheInvalidation::Visibility {
|
|
channel_id: Uuid::from_u128(0xabcd),
|
|
},
|
|
] {
|
|
let json = serde_json::to_string(&msg).unwrap();
|
|
assert_eq!(
|
|
serde_json::from_str::<CacheInvalidation>(&json).unwrap(),
|
|
msg
|
|
);
|
|
}
|
|
}
|
|
}
|