mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(relay): add local pubsub and media backends
Signed-off-by: npub1z3hmzc9ryehxzedl5wzlvpyvja0d483peaja5zt6pd0209f9x2jspe2dxh <146fb160a3266e6165bfa385f6048c975eda9e21cf65da097a0b5ea7952532a5@buzz.block.builderlab.xyz>
This commit is contained in:
parent
984d696018
commit
5a55f23f1d
Generated
+1
@@ -1132,6 +1132,7 @@ dependencies = [
|
||||
"buzz-auth",
|
||||
"buzz-core",
|
||||
"chrono",
|
||||
"dashmap",
|
||||
"deadpool-redis",
|
||||
"futures-util",
|
||||
"nostr",
|
||||
|
||||
+329
-111
@@ -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<Box<dyn futures_core::Stream<Item = Result<Bytes, MediaError>> + 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<PathBuf>) -> 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<PathBuf, MediaError> {
|
||||
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<Vec<u8>, 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<Vec<u8>, 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<Bytes, MediaError>` 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<ByteStream, MediaError> {
|
||||
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<bool, MediaError> {
|
||||
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<Option<BlobHeadMeta>, 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<BlobMeta, MediaError> {
|
||||
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<String>,
|
||||
max_keys: usize,
|
||||
) -> Result<crate::bucket_index::Page, MediaError> {
|
||||
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);
|
||||
|
||||
@@ -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 }
|
||||
|
||||
@@ -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<RedisPubSubManager>),
|
||||
InProcess(Arc<InProcessPubSubManager>),
|
||||
}
|
||||
|
||||
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<Self>) {
|
||||
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<ChannelEvent> {
|
||||
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<ScopedCacheInvalidation> {
|
||||
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<ScopedConnControl> {
|
||||
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<i64, PubSubError> {
|
||||
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<i64, PubSubError> {
|
||||
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<Option<String>, 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<HashMap<String, String>, 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<EventTopicKey, usize>,
|
||||
presence: DashMap<(buzz_core::CommunityId, String), String>,
|
||||
broadcast_tx: broadcast::Sender<ChannelEvent>,
|
||||
cache_invalidation_tx: broadcast::Sender<ScopedCacheInvalidation>,
|
||||
conn_control_tx: broadcast::Sender<ScopedConnControl>,
|
||||
}
|
||||
|
||||
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<ChannelEvent> {
|
||||
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<ScopedCacheInvalidation> {
|
||||
self.cache_invalidation_tx.subscribe()
|
||||
}
|
||||
fn subscribe_conn_control(&self) -> broadcast::Receiver<ScopedConnControl> {
|
||||
self.conn_control_tx.subscribe()
|
||||
}
|
||||
async fn publish_cache_invalidation(
|
||||
&self,
|
||||
ctx: &TenantContext,
|
||||
invalidation: &CacheInvalidation,
|
||||
) -> Result<i64, PubSubError> {
|
||||
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<i64, PubSubError> {
|
||||
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<i64, PubSubError> {
|
||||
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<Option<String>, 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<HashMap<String, String>, 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");
|
||||
|
||||
@@ -46,9 +46,12 @@ pub fn build_router(state: Arc<AppState>) -> 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<AppState>) -> 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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user