From 8045bcb5a07887b486e15ca319aa0e2ce7548554 Mon Sep 17 00:00:00 2001 From: phil Date: Sun, 25 May 2025 11:30:32 -0400 Subject: [PATCH] get all collections will add paging eventually --- ufos/src/server.rs | 18 ++++++++++++++++++ ufos/src/storage.rs | 2 ++ ufos/src/storage_fjall.rs | 25 +++++++++++++++++++++++++ ufos/src/storage_mem.rs | 3 +++ 4 files changed, 48 insertions(+) diff --git a/ufos/src/server.rs b/ufos/src/server.rs index 3ad04c6..3e72c1d 100644 --- a/ufos/src/server.rs +++ b/ufos/src/server.rs @@ -213,6 +213,23 @@ async fn get_records_total_seen( ok_cors(seen_by_collection) } +/// Get all collections +/// +/// TODO: paginate +#[endpoint { + method = GET, + path = "/collections/all" +}] +async fn get_all_collections(ctx: RequestContext) -> OkCorsResponse> { + let Context { storage, .. } = ctx.context(); + let collections = storage + .get_all_collections(QueryPeriod::all_time()) + .await + .map_err(|e| HttpError::for_internal_error(format!("oh shoot: {e:?}")))?; + + ok_cors(collections) +} + /// Get top collections by record count #[endpoint { method = GET, @@ -274,6 +291,7 @@ pub async fn serve(storage: impl StoreReader + 'static) -> Result<(), String> { api.register(get_meta_info).unwrap(); api.register(get_records_by_collections).unwrap(); api.register(get_records_total_seen).unwrap(); + api.register(get_all_collections).unwrap(); api.register(get_top_collections_by_count).unwrap(); api.register(get_top_collections_by_dids).unwrap(); api.register(get_top_collections).unwrap(); diff --git a/ufos/src/storage.rs b/ufos/src/storage.rs index 275b3f8..aa14c2f 100644 --- a/ufos/src/storage.rs +++ b/ufos/src/storage.rs @@ -72,6 +72,8 @@ pub trait StoreReader: Send + Sync { async fn get_consumer_info(&self) -> StorageResult; + async fn get_all_collections(&self, period: QueryPeriod) -> StorageResult>; + async fn get_top_collections_by_count( &self, limit: usize, diff --git a/ufos/src/storage_fjall.rs b/ufos/src/storage_fjall.rs index 07c5f61..98714a7 100644 --- a/ufos/src/storage_fjall.rs +++ b/ufos/src/storage_fjall.rs @@ -372,6 +372,27 @@ impl FjallReader { }) } + fn get_all_collections(&self, period: QueryPeriod) -> StorageResult> { + Ok(if period.is_all_time() { + let snapshot = self.rollups.snapshot(); + let mut out = Vec::new(); + let prefix = AllTimeRollupKey::from_prefix_to_db_bytes(&Default::default())?; + for kv in snapshot.prefix(prefix) { + let (key_bytes, val_bytes) = kv?; + let key = db_complete::(&key_bytes)?; + let db_counts = db_complete::(&val_bytes)?; + out.push(Count { + thing: key.collection().to_string(), + records: db_counts.records(), + dids_estimate: db_counts.dids().estimate() as u64, + }); + } + out + } else { + todo!() + }) + } + fn get_top_collections_by_count( &self, limit: usize, @@ -573,6 +594,10 @@ impl StoreReader for FjallReader { let s = self.clone(); tokio::task::spawn_blocking(move || FjallReader::get_consumer_info(&s)).await? } + async fn get_all_collections(&self, period: QueryPeriod) -> StorageResult> { + let s = self.clone(); + tokio::task::spawn_blocking(move || FjallReader::get_all_collections(&s, period)).await? + } async fn get_top_collections_by_count( &self, limit: usize, diff --git a/ufos/src/storage_mem.rs b/ufos/src/storage_mem.rs index eb2cb69..83e896f 100644 --- a/ufos/src/storage_mem.rs +++ b/ufos/src/storage_mem.rs @@ -594,6 +594,9 @@ impl StoreReader for MemReader { let s = self.clone(); tokio::task::spawn_blocking(move || MemReader::get_top_collections(&s)).await? } + async fn get_all_collections(&self, _: QueryPeriod) -> StorageResult> { + todo!() + } async fn get_top_collections_by_count( &self, _: usize, -- 2.51.2