mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(relay): expire local presence leases
Signed-off-by: npub1z3hmzc9ryehxzedl5wzlvpyvja0d483peaja5zt6pd0209f9x2jspe2dxh <146fb160a3266e6165bfa385f6048c975eda9e21cf65da097a0b5ea7952532a5@buzz.block.builderlab.xyz>
This commit is contained in:
parent
d2a4a31c11
commit
94e60f7b9b
@@ -442,14 +442,21 @@ impl MediaStorage {
|
||||
.await
|
||||
.map_err(|e| MediaError::StorageError(e.to_string()))?
|
||||
.map_err(Self::fs_error)?;
|
||||
let start = continuation_token
|
||||
.and_then(|token| {
|
||||
files
|
||||
.iter()
|
||||
.position(|(key, _)| key == &token)
|
||||
.map(|n| n + 1)
|
||||
})
|
||||
.unwrap_or(0);
|
||||
if max_keys == 0 {
|
||||
return Err(MediaError::StorageError(
|
||||
"list page size must be nonzero".to_owned(),
|
||||
));
|
||||
}
|
||||
let start = match continuation_token {
|
||||
Some(token) => files
|
||||
.iter()
|
||||
.position(|(key, _)| key == &token)
|
||||
.map(|n| n + 1)
|
||||
.ok_or_else(|| {
|
||||
MediaError::StorageError("unknown list continuation token".to_owned())
|
||||
})?,
|
||||
None => 0,
|
||||
};
|
||||
let remaining = files.split_off(start);
|
||||
let is_truncated = remaining.len() > max_keys;
|
||||
let objects: Vec<_> = remaining.into_iter().take(max_keys).collect();
|
||||
@@ -595,6 +602,29 @@ mod tests {
|
||||
assert!(!storage.head(&key).await.unwrap());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn filesystem_listing_rejects_zero_page_size_and_unknown_token() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let storage = MediaStorage::filesystem(dir.path());
|
||||
let key = format!("{}.bin", "b".repeat(64));
|
||||
storage
|
||||
.put(&key, b"x", "application/octet-stream")
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(storage
|
||||
.list_page(None, 0)
|
||||
.await
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("nonzero"));
|
||||
assert!(storage
|
||||
.list_page(Some("missing".to_owned()), 1)
|
||||
.await
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("unknown list continuation token"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sidecar_keys_are_community_scoped() {
|
||||
let a = tenant(1);
|
||||
|
||||
@@ -46,7 +46,7 @@ pub use error::PubSubError;
|
||||
use std::collections::HashMap;
|
||||
use std::future::pending;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use buzz_core::TenantContext;
|
||||
use dashmap::DashMap;
|
||||
@@ -131,8 +131,20 @@ impl PubSubManager {
|
||||
|
||||
/// Creates an in-process backend for a single relay process.
|
||||
pub fn in_process() -> Self {
|
||||
Self::in_process_with_presence_ttl(Duration::from_secs(presence::PRESENCE_TTL_SECS))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn in_process_with_presence_ttl(presence_ttl: Duration) -> Self {
|
||||
Self {
|
||||
backend: PubSubBackend::InProcess(Arc::new(InProcessPubSubManager::new())),
|
||||
backend: PubSubBackend::InProcess(Arc::new(InProcessPubSubManager::new(presence_ttl))),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(not(test))]
|
||||
fn in_process_with_presence_ttl(presence_ttl: Duration) -> Self {
|
||||
Self {
|
||||
backend: PubSubBackend::InProcess(Arc::new(InProcessPubSubManager::new(presence_ttl))),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -309,20 +321,27 @@ impl PubSubManager {
|
||||
/// Process-local pub/sub, presence, and control-plane fan-out for single-node relays.
|
||||
struct InProcessPubSubManager {
|
||||
topics: DashMap<EventTopicKey, usize>,
|
||||
presence: DashMap<(buzz_core::CommunityId, String), String>,
|
||||
presence: DashMap<(buzz_core::CommunityId, String), PresenceLease>,
|
||||
presence_ttl: Duration,
|
||||
broadcast_tx: broadcast::Sender<ChannelEvent>,
|
||||
cache_invalidation_tx: broadcast::Sender<ScopedCacheInvalidation>,
|
||||
conn_control_tx: broadcast::Sender<ScopedConnControl>,
|
||||
}
|
||||
|
||||
struct PresenceLease {
|
||||
status: String,
|
||||
expires_at: Instant,
|
||||
}
|
||||
|
||||
impl InProcessPubSubManager {
|
||||
fn new() -> Self {
|
||||
fn new(presence_ttl: Duration) -> Self {
|
||||
let (broadcast_tx, _) = broadcast::channel(4096);
|
||||
let (cache_invalidation_tx, _) = broadcast::channel(4096);
|
||||
let (conn_control_tx, _) = broadcast::channel(4096);
|
||||
Self {
|
||||
topics: DashMap::new(),
|
||||
presence: DashMap::new(),
|
||||
presence_ttl,
|
||||
broadcast_tx,
|
||||
cache_invalidation_tx,
|
||||
conn_control_tx,
|
||||
@@ -406,8 +425,13 @@ impl InProcessPubSubManager {
|
||||
pubkey: &PublicKey,
|
||||
status: &str,
|
||||
) -> Result<(), PubSubError> {
|
||||
self.presence
|
||||
.insert((ctx.community(), pubkey.to_hex()), status.to_owned());
|
||||
self.presence.insert(
|
||||
(ctx.community(), pubkey.to_hex()),
|
||||
PresenceLease {
|
||||
status: status.to_owned(),
|
||||
expires_at: Instant::now() + self.presence_ttl,
|
||||
},
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
async fn clear_presence(
|
||||
@@ -423,10 +447,8 @@ impl InProcessPubSubManager {
|
||||
ctx: &TenantContext,
|
||||
pubkey: &PublicKey,
|
||||
) -> Result<Option<String>, PubSubError> {
|
||||
Ok(self
|
||||
.presence
|
||||
.get(&(ctx.community(), pubkey.to_hex()))
|
||||
.map(|v| v.clone()))
|
||||
let key = (ctx.community(), pubkey.to_hex());
|
||||
Ok(self.active_presence(&key))
|
||||
}
|
||||
async fn get_presence_bulk(
|
||||
&self,
|
||||
@@ -437,12 +459,21 @@ impl InProcessPubSubManager {
|
||||
.iter()
|
||||
.filter_map(|p| {
|
||||
let hex = p.to_hex();
|
||||
self.presence
|
||||
.get(&(ctx.community(), hex.clone()))
|
||||
.map(|v| (hex, v.clone()))
|
||||
self.active_presence(&(ctx.community(), hex.clone()))
|
||||
.map(|status| (hex, status))
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
|
||||
fn active_presence(&self, key: &(buzz_core::CommunityId, String)) -> Option<String> {
|
||||
let lease = self.presence.get(key)?;
|
||||
if lease.expires_at > Instant::now() {
|
||||
return Some(lease.status.clone());
|
||||
}
|
||||
drop(lease);
|
||||
self.presence.remove(key);
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
/// Redis implementation behind [`PubSubManager`].
|
||||
@@ -1005,6 +1036,32 @@ mod tests {
|
||||
assert_eq!(manager.get_presence(&a, &key).await.unwrap(), None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn in_process_presence_lease_expires_without_disconnect() {
|
||||
let manager = PubSubManager::in_process_with_presence_ttl(Duration::from_millis(5));
|
||||
let tenant = ctx(1, "lease.example");
|
||||
let pubkey = Keys::generate().public_key();
|
||||
manager
|
||||
.set_presence(&tenant, &pubkey, "online")
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
manager
|
||||
.get_presence(&tenant, &pubkey)
|
||||
.await
|
||||
.unwrap()
|
||||
.as_deref(),
|
||||
Some("online")
|
||||
);
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
assert_eq!(manager.get_presence(&tenant, &pubkey).await.unwrap(), None);
|
||||
assert!(manager
|
||||
.get_presence_bulk(&tenant, &[pubkey])
|
||||
.await
|
||||
.unwrap()
|
||||
.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn config_defaults_debounce_but_allows_override() {
|
||||
let config = PubSubConfig::new("redis://example");
|
||||
|
||||
Reference in New Issue
Block a user