diff --git a/crates/buzz-media/src/storage.rs b/crates/buzz-media/src/storage.rs index 725b7ed7c..1e053917a 100644 --- a/crates/buzz-media/src/storage.rs +++ b/crates/buzz-media/src/storage.rs @@ -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); diff --git a/crates/buzz-pubsub/src/lib.rs b/crates/buzz-pubsub/src/lib.rs index f0d46a1a5..323b7b949 100644 --- a/crates/buzz-pubsub/src/lib.rs +++ b/crates/buzz-pubsub/src/lib.rs @@ -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, - presence: DashMap<(buzz_core::CommunityId, String), String>, + presence: DashMap<(buzz_core::CommunityId, String), PresenceLease>, + presence_ttl: Duration, broadcast_tx: broadcast::Sender, cache_invalidation_tx: broadcast::Sender, conn_control_tx: broadcast::Sender, } +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, 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 { + 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");