From 5a55f23f1d4032e3cedcf8ec39d0aeaf583f66bf Mon Sep 17 00:00:00 2001 From: npub1z3hmzc9ryehxzedl5wzlvpyvja0d483peaja5zt6pd0209f9x2jspe2dxh <146fb160a3266e6165bfa385f6048c975eda9e21cf65da097a0b5ea7952532a5@buzz.block.builderlab.xyz> Date: Sat, 1 Aug 2026 10:19:14 -0400 Subject: [PATCH] feat(relay): add local pubsub and media backends Signed-off-by: npub1z3hmzc9ryehxzedl5wzlvpyvja0d483peaja5zt6pd0209f9x2jspe2dxh <146fb160a3266e6165bfa385f6048c975eda9e21cf65da097a0b5ea7952532a5@buzz.block.builderlab.xyz> --- Cargo.lock | 1 + crates/buzz-media/src/storage.rs | 440 +++++++++++++++++++++++-------- crates/buzz-pubsub/Cargo.toml | 1 + crates/buzz-pubsub/src/lib.rs | 213 +++++++++++++++ crates/buzz-relay/src/router.rs | 20 +- 5 files changed, 557 insertions(+), 118 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index da86c89b8..739ba9de6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1132,6 +1132,7 @@ dependencies = [ "buzz-auth", "buzz-core", "chrono", + "dashmap", "deadpool-redis", "futures-util", "nostr", diff --git a/crates/buzz-media/src/storage.rs b/crates/buzz-media/src/storage.rs index ffed5ef07..725b7ed7c 100644 --- a/crates/buzz-media/src/storage.rs +++ b/crates/buzz-media/src/storage.rs @@ -1,6 +1,6 @@ //! S3/MinIO storage client. -use std::path::Path; +use std::path::{Path, PathBuf}; use std::pin::Pin; use buzz_core::tenant::{CommunityId, TenantContext}; @@ -8,9 +8,11 @@ use buzz_core::tenant::{CommunityId, TenantContext}; use crate::config::{MediaConfig, S3AddressingStyle}; use crate::error::MediaError; use bytes::Bytes; +use futures_util::StreamExt; use s3::creds::Credentials; use s3::{Bucket, Region}; use serde::{Deserialize, Serialize}; +use tokio_util::io::ReaderStream; /// A stream of byte chunks from S3, usable with `axum::body::Body::from_stream()`. pub type ByteStream = Pin> + Send>>; @@ -22,6 +24,11 @@ pub struct MediaStorage { enum MediaBackend { S3(S3MediaStorage), + Filesystem(FilesystemMediaStorage), +} + +struct FilesystemMediaStorage { + root: PathBuf, } struct S3MediaStorage { @@ -79,126 +86,232 @@ impl MediaStorage { }) } + /// Creates a filesystem-backed content-addressed media store rooted at `root`. + pub fn filesystem(root: impl Into) -> Self { + Self { + backend: MediaBackend::Filesystem(FilesystemMediaStorage { root: root.into() }), + } + } + + #[cfg(test)] fn s3(&self) -> &S3MediaStorage { match &self.backend { MediaBackend::S3(storage) => storage, + MediaBackend::Filesystem(_) => panic!("filesystem media storage has no S3 client"), + } + } + + fn filesystem_path(storage: &FilesystemMediaStorage, key: &str) -> Result { + if key.is_empty() + || key.starts_with('/') + || key + .split('/') + .any(|part| part.is_empty() || part == "." || part == "..") + { + return Err(MediaError::StorageError( + "invalid media storage key".to_owned(), + )); + } + let filename = key.rsplit('/').next().expect("validated non-empty key"); + let sha = filename.split('.').next().unwrap_or_default(); + if sha.len() >= 4 && sha.bytes().all(|byte| byte.is_ascii_hexdigit()) { + Ok(storage.root.join(&sha[..2]).join(&sha[2..4]).join(key)) + } else { + Ok(storage.root.join("_aux").join(key)) + } + } + + fn fs_error(error: std::io::Error) -> MediaError { + if error.kind() == std::io::ErrorKind::NotFound { + MediaError::NotFound + } else { + MediaError::Io(error.to_string()) } } /// Store an object from a byte slice. - /// - /// Used for images, sidecars, and thumbnails. For large video files use - /// [`put_file`] to avoid loading the entire blob into RAM. pub async fn put(&self, key: &str, bytes: &[u8], content_type: &str) -> Result<(), MediaError> { - self.s3() - .bucket - .put_object_with_content_type(key, bytes, content_type) - .await?; - Ok(()) + match &self.backend { + MediaBackend::S3(storage) => { + storage + .bucket + .put_object_with_content_type(key, bytes, content_type) + .await?; + Ok(()) + } + MediaBackend::Filesystem(storage) => { + let path = Self::filesystem_path(storage, key)?; + if let Some(parent) = path.parent() { + tokio::fs::create_dir_all(parent) + .await + .map_err(Self::fs_error)?; + } + tokio::fs::write(path, bytes).await.map_err(Self::fs_error) + } + } } - /// Stream a file from disk into S3 without loading it into RAM. - /// - /// Uses rust-s3's `put_object_stream_with_content_type` which reads from - /// the file incrementally via an 8 MiB `BufReader`. The full file is never - /// held in memory simultaneously. Intended for video blobs (up to 500 MB). + /// Stream a file from disk into storage without loading it into RAM. pub async fn put_file( &self, key: &str, path: &Path, content_type: &str, ) -> Result<(), MediaError> { - const BUF: usize = 8 * 1024 * 1024; // 8 MiB read buffer - - let file = tokio::fs::File::open(path) - .await - .map_err(|e| MediaError::Io(e.to_string()))?; - let mut reader = tokio::io::BufReader::with_capacity(BUF, file); - - self.s3() - .bucket - .put_object_stream_with_content_type(&mut reader, key, content_type) - .await?; - Ok(()) + match &self.backend { + MediaBackend::S3(storage) => { + const BUF: usize = 8 * 1024 * 1024; + let file = tokio::fs::File::open(path) + .await + .map_err(|e| MediaError::Io(e.to_string()))?; + let mut reader = tokio::io::BufReader::with_capacity(BUF, file); + storage + .bucket + .put_object_stream_with_content_type(&mut reader, key, content_type) + .await?; + Ok(()) + } + MediaBackend::Filesystem(storage) => { + let target = Self::filesystem_path(storage, key)?; + if let Some(parent) = target.parent() { + tokio::fs::create_dir_all(parent) + .await + .map_err(Self::fs_error)?; + } + tokio::fs::copy(path, target) + .await + .map(|_| ()) + .map_err(Self::fs_error) + } + } } /// Retrieve an object's bytes. pub async fn get(&self, key: &str) -> Result, MediaError> { - match self.s3().bucket.get_object(key).await { - Ok(response) => Ok(response.to_vec()), - Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Err(MediaError::NotFound), - Err(e) => Err(MediaError::StorageError(e.to_string())), + match &self.backend { + MediaBackend::S3(storage) => match storage.bucket.get_object(key).await { + Ok(response) => Ok(response.to_vec()), + Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Err(MediaError::NotFound), + Err(e) => Err(MediaError::StorageError(e.to_string())), + }, + MediaBackend::Filesystem(storage) => { + tokio::fs::read(Self::filesystem_path(storage, key)?) + .await + .map_err(Self::fs_error) + } } } - /// Retrieve a byte range from an object via S3-native `Range` GET. - /// - /// `start` and `end` are inclusive byte offsets. Only the requested slice - /// is transferred from S3 — the full object is never loaded into RAM. - /// Intended for HTTP 206 range responses on large video blobs. + /// Retrieve an inclusive byte range from an object. pub async fn get_range(&self, key: &str, start: u64, end: u64) -> Result, MediaError> { - match self - .s3() - .bucket - .get_object_range(key, start, Some(end)) - .await - { - Ok(response) => Ok(response.to_vec()), - Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Err(MediaError::NotFound), - Err(e) => Err(MediaError::StorageError(e.to_string())), + match &self.backend { + MediaBackend::S3(storage) => { + match storage.bucket.get_object_range(key, start, Some(end)).await { + Ok(response) => Ok(response.to_vec()), + Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Err(MediaError::NotFound), + Err(e) => Err(MediaError::StorageError(e.to_string())), + } + } + MediaBackend::Filesystem(storage) => { + let bytes = tokio::fs::read(Self::filesystem_path(storage, key)?) + .await + .map_err(Self::fs_error)?; + let start = usize::try_from(start) + .map_err(|_| MediaError::StorageError("range start overflow".to_owned()))?; + let end = usize::try_from(end) + .map_err(|_| MediaError::StorageError("range end overflow".to_owned()))?; + if start > end || end >= bytes.len() { + return Err(MediaError::StorageError("invalid media range".to_owned())); + } + Ok(bytes[start..=end].to_vec()) + } } } - /// Stream an object's bytes from S3 without loading into RAM. - /// - /// Returns a pinned stream of `Result` chunks. - /// The full object is never buffered — intended for streaming large - /// blobs (video) directly into HTTP responses via `Body::from_stream()`. + /// Stream an object's bytes without buffering its full content. pub async fn get_stream(&self, key: &str) -> Result { - let response = self - .s3() - .bucket - .get_object_stream(key) - .await - .map_err(|e| MediaError::StorageError(e.to_string()))?; - - if response.status_code == 404 { - return Err(MediaError::NotFound); + match &self.backend { + MediaBackend::S3(storage) => { + let response = storage + .bucket + .get_object_stream(key) + .await + .map_err(|e| MediaError::StorageError(e.to_string()))?; + if response.status_code == 404 { + return Err(MediaError::NotFound); + } + Ok(Box::pin(futures_util::StreamExt::map( + response.bytes, + |chunk| chunk.map_err(|e| MediaError::StorageError(e.to_string())), + ))) + } + MediaBackend::Filesystem(storage) => { + let file = tokio::fs::File::open(Self::filesystem_path(storage, key)?) + .await + .map_err(Self::fs_error)?; + Ok(Box::pin( + ReaderStream::new(file).map(|chunk| chunk.map_err(Self::fs_error)), + )) + } } - - let stream = futures_util::StreamExt::map(response.bytes, |chunk| { - chunk.map_err(|e| MediaError::StorageError(e.to_string())) - }); - Ok(Box::pin(stream)) } - /// Check if an object exists. Returns false on 404. + /// Check if an object exists. Returns false on absence. pub async fn head(&self, key: &str) -> Result { - match self.s3().bucket.head_object(key).await { - Ok(_) => Ok(true), - Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Ok(false), - Err(e) => Err(MediaError::StorageError(e.to_string())), + match &self.backend { + MediaBackend::S3(storage) => match storage.bucket.head_object(key).await { + Ok(_) => Ok(true), + Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Ok(false), + Err(e) => Err(MediaError::StorageError(e.to_string())), + }, + MediaBackend::Filesystem(storage) => { + match tokio::fs::metadata(Self::filesystem_path(storage, key)?).await { + Ok(_) => Ok(true), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false), + Err(error) => Err(Self::fs_error(error)), + } + } } } - /// Delete an object. Returns an error on failure — callers decide whether to propagate. + /// Delete an object. pub async fn delete(&self, key: &str) -> Result<(), MediaError> { - self.s3() - .bucket - .delete_object(key) - .await - .map_err(|e| MediaError::StorageError(e.to_string()))?; - Ok(()) + match &self.backend { + MediaBackend::S3(storage) => { + storage + .bucket + .delete_object(key) + .await + .map_err(|e| MediaError::StorageError(e.to_string()))?; + Ok(()) + } + MediaBackend::Filesystem(storage) => { + tokio::fs::remove_file(Self::filesystem_path(storage, key)?) + .await + .map_err(Self::fs_error) + } + } } /// HEAD with metadata — returns Content-Length (size). pub async fn head_with_metadata(&self, key: &str) -> Result, MediaError> { - match self.s3().bucket.head_object(key).await { - Ok((result, _)) => Ok(Some(BlobHeadMeta { - size: result.content_length.unwrap_or(0) as u64, - })), - Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Ok(None), - Err(e) => Err(MediaError::StorageError(e.to_string())), + match &self.backend { + MediaBackend::S3(storage) => match storage.bucket.head_object(key).await { + Ok((result, _)) => Ok(Some(BlobHeadMeta { + size: result.content_length.unwrap_or(0) as u64, + })), + Err(s3::error::S3Error::HttpFailWithBody(404, _)) => Ok(None), + Err(e) => Err(MediaError::StorageError(e.to_string())), + }, + MediaBackend::Filesystem(storage) => { + match tokio::fs::metadata(Self::filesystem_path(storage, key)?).await { + Ok(metadata) => Ok(Some(BlobHeadMeta { + size: metadata.len(), + })), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None), + Err(error) => Err(Self::fs_error(error)), + } + } } } @@ -224,8 +337,8 @@ impl MediaStorage { sha256: &str, ) -> Result { let key = Self::ctx_sidecar_key(ctx, sha256); - let resp = self.s3().bucket.get_object(&key).await?; - let meta: BlobMeta = serde_json::from_slice(&resp.to_vec())?; + let bytes = self.get(&key).await?; + let meta: BlobMeta = serde_json::from_slice(&bytes)?; Ok(meta) } @@ -259,39 +372,96 @@ impl MediaStorage { .map(|m| m.mime_type) } - /// One page of a full-bucket listing, for the storage sweep. Wraps - /// rust-s3's manual `list_page` (NOT the auto-paginating `list`, which - /// has no cap) and converts the result into the storage-agnostic - /// [`crate::bucket_index::Page`] shape the pure fold consumes. - /// - /// `max_keys` bounds one HTTP response, not the sweep's total object - /// cap — the caller (`fold_bucket_listing`) enforces the cumulative cap - /// across pages. + /// One page of a full-bucket listing. pub async fn list_page( &self, continuation_token: Option, max_keys: usize, ) -> Result { - let (result, _status) = self - .s3() - .bucket - .list_page( - String::new(), - None, - continuation_token, - None, - Some(max_keys), - ) - .await?; - Ok(crate::bucket_index::Page { - objects: result - .contents - .into_iter() - .map(|obj| (obj.key, obj.size)) - .collect(), - next_continuation_token: result.next_continuation_token, - is_truncated: result.is_truncated, - }) + match &self.backend { + MediaBackend::S3(storage) => { + let (result, _) = storage + .bucket + .list_page( + String::new(), + None, + continuation_token, + None, + Some(max_keys), + ) + .await?; + Ok(crate::bucket_index::Page { + objects: result + .contents + .into_iter() + .map(|obj| (obj.key, obj.size)) + .collect(), + next_continuation_token: result.next_continuation_token, + is_truncated: result.is_truncated, + }) + } + MediaBackend::Filesystem(storage) => { + let root = storage.root.clone(); + let mut files = tokio::task::spawn_blocking(move || { + fn walk( + root: &Path, + base: &Path, + out: &mut Vec<(String, u64)>, + ) -> std::io::Result<()> { + for entry in std::fs::read_dir(root)? { + let entry = entry?; + let path = entry.path(); + if path.is_dir() { + walk(&path, base, out)?; + } else { + let key = path + .strip_prefix(base) + .unwrap() + .to_string_lossy() + .replace('\\', "/"); + if let Some(key) = key.strip_prefix("_aux/") { + out.push((key.to_owned(), entry.metadata()?.len())); + } else if let Some((_, _, key)) = + key.split_once('/').and_then(|(a, rest)| { + rest.split_once('/').map(|(b, c)| (a, b, c)) + }) + { + out.push((key.to_owned(), entry.metadata()?.len())); + } + } + } + Ok(()) + } + let mut out = Vec::new(); + if root.exists() { + walk(&root, &root, &mut out)?; + } + out.sort_by(|a, b| a.0.cmp(&b.0)); + Ok::<_, std::io::Error>(out) + }) + .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); + let remaining = files.split_off(start); + let is_truncated = remaining.len() > max_keys; + let objects: Vec<_> = remaining.into_iter().take(max_keys).collect(); + let next_continuation_token = is_truncated + .then(|| objects.last().expect("nonempty truncated page").0.clone()); + Ok(crate::bucket_index::Page { + objects, + next_continuation_token, + is_truncated, + }) + } + } } } @@ -377,6 +547,54 @@ mod tests { ); } + #[tokio::test] + async fn filesystem_backend_roundtrips_cas_sidecars_ranges_and_pages() { + let dir = tempfile::tempdir().unwrap(); + let storage = MediaStorage::filesystem(dir.path()); + let sha = "a".repeat(64); + let key = format!("{sha}.bin"); + storage + .put(&key, b"hello", "application/octet-stream") + .await + .unwrap(); + assert_eq!(storage.get(&key).await.unwrap(), b"hello"); + assert_eq!(storage.get_range(&key, 1, 3).await.unwrap(), b"ell"); + assert_eq!( + storage + .head_with_metadata(&key) + .await + .unwrap() + .unwrap() + .size, + 5 + ); + assert!(dir.path().join("aa").join("aa").join(&key).exists()); + + let a = tenant(1); + let b = tenant(2); + storage + .put_sidecar( + &a, + &sha, + &BlobMeta { + mime_type: "text/plain".to_owned(), + ..Default::default() + }, + ) + .await + .unwrap(); + assert_eq!( + storage.read_sidecar_mime(&a, &key).await.as_deref(), + Some("text/plain") + ); + assert_eq!(storage.read_sidecar_mime(&b, &key).await, None); + let page = storage.list_page(None, 1).await.unwrap(); + assert_eq!(page.objects.len(), 1); + assert!(page.is_truncated); + storage.delete(&key).await.unwrap(); + assert!(!storage.head(&key).await.unwrap()); + } + #[test] fn sidecar_keys_are_community_scoped() { let a = tenant(1); diff --git a/crates/buzz-pubsub/Cargo.toml b/crates/buzz-pubsub/Cargo.toml index 2ee2d6b47..3b1a18b3b 100644 --- a/crates/buzz-pubsub/Cargo.toml +++ b/crates/buzz-pubsub/Cargo.toml @@ -17,6 +17,7 @@ serde = { workspace = true } serde_json = { workspace = true } uuid = { workspace = true } chrono = { workspace = true } +dashmap = { workspace = true } tracing = { workspace = true } thiserror = { workspace = true } nostr = { workspace = true } diff --git a/crates/buzz-pubsub/src/lib.rs b/crates/buzz-pubsub/src/lib.rs index b0dd9510f..f0d46a1a5 100644 --- a/crates/buzz-pubsub/src/lib.rs +++ b/crates/buzz-pubsub/src/lib.rs @@ -44,10 +44,12 @@ pub mod topic; pub use error::PubSubError; use std::collections::HashMap; +use std::future::pending; use std::sync::Arc; use std::time::Duration; use buzz_core::TenantContext; +use dashmap::DashMap; use nostr::PublicKey; use tokio::sync::{broadcast, mpsc, Mutex}; @@ -106,6 +108,7 @@ pub struct PubSubManager { enum PubSubBackend { Redis(Arc), + InProcess(Arc), } impl PubSubManager { @@ -126,11 +129,19 @@ impl PubSubManager { }) } + /// Creates an in-process backend for a single relay process. + pub fn in_process() -> Self { + Self { + backend: PubSubBackend::InProcess(Arc::new(InProcessPubSubManager::new())), + } + } + /// Runs the backend event subscriber. Process-local backends may implement /// this as a pending task because publication already reaches local receivers. pub async fn run_subscriber(self: Arc) { match &self.backend { PubSubBackend::Redis(backend) => Arc::clone(backend).run_subscriber().await, + PubSubBackend::InProcess(_) => pending::<()>().await, } } @@ -142,6 +153,7 @@ impl PubSubManager { .run_cache_invalidation_subscriber() .await } + PubSubBackend::InProcess(_) => pending::<()>().await, } } @@ -151,6 +163,7 @@ impl PubSubManager { PubSubBackend::Redis(backend) => { Arc::clone(backend).run_conn_control_subscriber().await } + PubSubBackend::InProcess(_) => pending::<()>().await, } } @@ -158,6 +171,7 @@ impl PubSubManager { pub fn subscribe_local(&self) -> broadcast::Receiver { match &self.backend { PubSubBackend::Redis(backend) => backend.subscribe_local(), + PubSubBackend::InProcess(backend) => backend.subscribe_local(), } } @@ -165,6 +179,7 @@ impl PubSubManager { pub async fn retain_topic(&self, ctx: &TenantContext, topic: EventTopic) { match &self.backend { PubSubBackend::Redis(backend) => backend.retain_topic(ctx, topic).await, + PubSubBackend::InProcess(backend) => backend.retain_topic(ctx, topic).await, } } @@ -172,6 +187,7 @@ impl PubSubManager { pub async fn release_topic(&self, ctx: &TenantContext, topic: EventTopic) { match &self.backend { PubSubBackend::Redis(backend) => backend.release_topic(ctx, topic).await, + PubSubBackend::InProcess(backend) => backend.release_topic(ctx, topic).await, } } @@ -179,6 +195,7 @@ impl PubSubManager { pub async fn topic_refcount(&self, ctx: &TenantContext, topic: EventTopic) -> usize { match &self.backend { PubSubBackend::Redis(backend) => backend.topic_refcount(ctx, topic).await, + PubSubBackend::InProcess(backend) => backend.topic_refcount(ctx, topic).await, } } @@ -186,6 +203,7 @@ impl PubSubManager { pub fn subscribe_cache_invalidations(&self) -> broadcast::Receiver { match &self.backend { PubSubBackend::Redis(backend) => backend.subscribe_cache_invalidations(), + PubSubBackend::InProcess(backend) => backend.subscribe_cache_invalidations(), } } @@ -193,6 +211,7 @@ impl PubSubManager { pub fn subscribe_conn_control(&self) -> broadcast::Receiver { match &self.backend { PubSubBackend::Redis(backend) => backend.subscribe_conn_control(), + PubSubBackend::InProcess(backend) => backend.subscribe_conn_control(), } } @@ -206,6 +225,9 @@ impl PubSubManager { PubSubBackend::Redis(backend) => { backend.publish_cache_invalidation(ctx, invalidation).await } + PubSubBackend::InProcess(backend) => { + backend.publish_cache_invalidation(ctx, invalidation).await + } } } @@ -217,6 +239,7 @@ impl PubSubManager { ) -> Result { match &self.backend { PubSubBackend::Redis(backend) => backend.publish_conn_control(ctx, command).await, + PubSubBackend::InProcess(backend) => backend.publish_conn_control(ctx, command).await, } } @@ -229,6 +252,7 @@ impl PubSubManager { ) -> Result { match &self.backend { PubSubBackend::Redis(backend) => backend.publish_event(ctx, topic, event).await, + PubSubBackend::InProcess(backend) => backend.publish_event(ctx, topic, event).await, } } @@ -241,6 +265,7 @@ impl PubSubManager { ) -> Result<(), PubSubError> { match &self.backend { PubSubBackend::Redis(backend) => backend.set_presence(ctx, pubkey, status).await, + PubSubBackend::InProcess(backend) => backend.set_presence(ctx, pubkey, status).await, } } @@ -252,6 +277,7 @@ impl PubSubManager { ) -> Result<(), PubSubError> { match &self.backend { PubSubBackend::Redis(backend) => backend.clear_presence(ctx, pubkey).await, + PubSubBackend::InProcess(backend) => backend.clear_presence(ctx, pubkey).await, } } @@ -263,6 +289,7 @@ impl PubSubManager { ) -> Result, PubSubError> { match &self.backend { PubSubBackend::Redis(backend) => backend.get_presence(ctx, pubkey).await, + PubSubBackend::InProcess(backend) => backend.get_presence(ctx, pubkey).await, } } @@ -274,10 +301,150 @@ impl PubSubManager { ) -> Result, PubSubError> { match &self.backend { PubSubBackend::Redis(backend) => backend.get_presence_bulk(ctx, pubkeys).await, + PubSubBackend::InProcess(backend) => backend.get_presence_bulk(ctx, pubkeys).await, } } } +/// 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>, + broadcast_tx: broadcast::Sender, + cache_invalidation_tx: broadcast::Sender, + conn_control_tx: broadcast::Sender, +} + +impl InProcessPubSubManager { + fn new() -> 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(), + broadcast_tx, + cache_invalidation_tx, + conn_control_tx, + } + } + fn subscribe_local(&self) -> broadcast::Receiver { + self.broadcast_tx.subscribe() + } + async fn retain_topic(&self, ctx: &TenantContext, topic: EventTopic) { + self.topics + .entry(EventTopicKey::from_context(ctx, topic)) + .and_modify(|n| *n += 1) + .or_insert(1); + } + async fn release_topic(&self, ctx: &TenantContext, topic: EventTopic) { + let key = EventTopicKey::from_context(ctx, topic); + if let Some(mut entry) = self.topics.get_mut(&key) { + *entry -= 1; + if *entry == 0 { + drop(entry); + self.topics.remove(&key); + } + } + } + async fn topic_refcount(&self, ctx: &TenantContext, topic: EventTopic) -> usize { + self.topics + .get(&EventTopicKey::from_context(ctx, topic)) + .map(|n| *n) + .unwrap_or(0) + } + fn subscribe_cache_invalidations(&self) -> broadcast::Receiver { + self.cache_invalidation_tx.subscribe() + } + fn subscribe_conn_control(&self) -> broadcast::Receiver { + self.conn_control_tx.subscribe() + } + async fn publish_cache_invalidation( + &self, + ctx: &TenantContext, + invalidation: &CacheInvalidation, + ) -> Result { + Ok(self + .cache_invalidation_tx + .send(ScopedCacheInvalidation { + community_id: ctx.community(), + invalidation: invalidation.clone(), + }) + .unwrap_or(0) as i64) + } + async fn publish_conn_control( + &self, + ctx: &TenantContext, + command: &ConnControl, + ) -> Result { + Ok(self + .conn_control_tx + .send(ScopedConnControl { + community_id: ctx.community(), + command: command.clone(), + }) + .unwrap_or(0) as i64) + } + async fn publish_event( + &self, + ctx: &TenantContext, + topic: EventTopic, + event: &nostr::Event, + ) -> Result { + Ok(self + .broadcast_tx + .send(ChannelEvent { + community_id: ctx.community(), + topic, + event: event.clone(), + }) + .unwrap_or(0) as i64) + } + async fn set_presence( + &self, + ctx: &TenantContext, + pubkey: &PublicKey, + status: &str, + ) -> Result<(), PubSubError> { + self.presence + .insert((ctx.community(), pubkey.to_hex()), status.to_owned()); + Ok(()) + } + async fn clear_presence( + &self, + ctx: &TenantContext, + pubkey: &PublicKey, + ) -> Result<(), PubSubError> { + self.presence.remove(&(ctx.community(), pubkey.to_hex())); + Ok(()) + } + async fn get_presence( + &self, + ctx: &TenantContext, + pubkey: &PublicKey, + ) -> Result, PubSubError> { + Ok(self + .presence + .get(&(ctx.community(), pubkey.to_hex())) + .map(|v| v.clone())) + } + async fn get_presence_bulk( + &self, + ctx: &TenantContext, + pubkeys: &[PublicKey], + ) -> Result, PubSubError> { + Ok(pubkeys + .iter() + .filter_map(|p| { + let hex = p.to_hex(); + self.presence + .get(&(ctx.community(), hex.clone())) + .map(|v| (hex, v.clone())) + }) + .collect()) + } +} + /// Redis implementation behind [`PubSubManager`]. struct RedisPubSubManager { pool: deadpool_redis::Pool, @@ -792,6 +959,52 @@ mod tests { assert_eq!(manager.topic_refcount(&ctx, topic).await, 0); } + #[tokio::test] + async fn in_process_fanout_presence_and_control_are_community_scoped() { + let manager = PubSubManager::in_process(); + let a = ctx(1, "a.example"); + let b = ctx(2, "b.example"); + let topic = EventTopic::Global; + manager.retain_topic(&a, topic).await; + manager.retain_topic(&a, topic).await; + assert_eq!(manager.topic_refcount(&a, topic).await, 2); + manager.release_topic(&a, topic).await; + assert_eq!(manager.topic_refcount(&a, topic).await, 1); + + let mut events = manager.subscribe_local(); + let mut invalidations = manager.subscribe_cache_invalidations(); + let mut controls = manager.subscribe_conn_control(); + let event = EventBuilder::new(Kind::TextNote, "local") + .tags([]) + .sign_with_keys(&Keys::generate()) + .unwrap(); + assert_eq!(manager.publish_event(&a, topic, &event).await.unwrap(), 1); + assert_eq!(events.recv().await.unwrap().community_id, a.community()); + manager + .publish_cache_invalidation(&a, &CacheInvalidation::AccessibleAll) + .await + .unwrap(); + assert_eq!( + invalidations.recv().await.unwrap().community_id, + a.community() + ); + manager + .publish_conn_control(&b, &ConnControl::DisconnectCommunity) + .await + .unwrap(); + assert_eq!(controls.recv().await.unwrap().community_id, b.community()); + + let key = Keys::generate().public_key(); + manager.set_presence(&a, &key, "online").await.unwrap(); + assert_eq!( + manager.get_presence(&a, &key).await.unwrap().as_deref(), + Some("online") + ); + assert_eq!(manager.get_presence(&b, &key).await.unwrap(), None); + manager.clear_presence(&a, &key).await.unwrap(); + assert_eq!(manager.get_presence(&a, &key).await.unwrap(), None); + } + #[test] fn config_defaults_debounce_but_allows_override() { let config = PubSubConfig::new("redis://example"); diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index 379f0ead1..f56796efa 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -46,9 +46,12 @@ pub fn build_router(state: Arc) -> Router { .layer(RequestBodyLimitLayer::new(media_body_limit)) .with_state(state.clone()); - let git_router = api::git::git_router(state.clone()); - - let git_policy_router = api::git::git_policy_router(state.clone()); + // Git is deliberately unavailable in the single-node profile: its only + // durable backend is S3. Do not build its routes (or internal hook policy + // endpoint) until a local GitStore exists. + let git_enabled = !state.config.profile.is_single_node(); + let git_router = git_enabled.then(|| api::git::git_router(state.clone())); + let git_policy_router = git_enabled.then(|| api::git::git_policy_router(state.clone())); let admin_enabled = state.config.admin.is_some(); let admin_web_dir = state @@ -133,10 +136,13 @@ pub fn build_router(state: Arc) -> Router { // Merge — each sub-router carries its own body limit. // Metrics → Trace → CORS applied once over the combined router. - let mut merged = api_router - .merge(media_router) - .merge(git_router) - .merge(git_policy_router); + let mut merged = api_router.merge(media_router); + if let Some(git_router) = git_router { + merged = merged.merge(git_router); + } + if let Some(git_policy_router) = git_policy_router { + merged = merged.merge(git_policy_router); + } if let Some(admin_router) = admin_router { merged = merged.merge(admin_router); }