diff --git a/parakeet-db/src/actor_cache.rs b/parakeet-db/src/actor_cache.rs deleted file mode 100644 index 1b0ad212..00000000 --- a/parakeet-db/src/actor_cache.rs +++ /dev/null @@ -1,265 +0,0 @@ -//! Unified actor cache using DashMap -//! -//! This module provides a thread-safe cache for basic actor metadata -//! (DID, handle, status, sync_state) that can be shared across both -//! consumer and parakeet services. -//! -//! Uses DashMap for efficient concurrent access. DashMap is internally sharded -//! (64 shards by default) which provides excellent parallel performance without -//! the complexity of manual sharding and bloom filters. - -use dashmap::DashMap; -use std::sync::Arc; - -use crate::types::{ActorStatus, ActorSyncState}; - -/// Basic actor metadata stored in the cache -#[derive(Debug, Clone)] -pub struct ActorMetadata { - /// The actor's DID - pub did: String, - /// The actor's handle (if known) - pub handle: Option, - /// The actor's status (active, deleted, etc.) - pub status: ActorStatus, - /// The actor's sync state (synced, dirty, partial, etc.) - pub sync_state: ActorSyncState, - /// The actor's repository revision (if known) - pub repo_rev: Option, -} - -/// A thread-safe cached actor metadata store -/// -/// Uses DashMap for efficient concurrent access with automatic sharding. -#[derive(Clone, Debug)] -pub struct CachedActorStore { - /// Map of DID -> ActorMetadata (internally sharded by DashMap) - actors: Arc>, -} - -impl CachedActorStore { - /// Create a new empty cached actor store - pub fn new() -> Self { - Self { - actors: Arc::new(DashMap::new()), - } - } - - /// Update the cache with new actor metadata (replaces existing data) - pub fn update_cache(&self, actors: Vec) { - self.actors.clear(); - for actor in actors { - self.actors.insert(actor.did.clone(), actor); - } - } - - /// Check if an actor exists in the cache - pub fn contains_actor(&self, did: &str) -> bool { - self.actors.contains_key(did) - } - - /// Get actor metadata from the cache - pub fn get_actor(&self, did: &str) -> Option { - self.actors.get(did).map(|entry| entry.clone()) - } - - /// Get multiple actors at once - pub fn get_actors_batch(&self, dids: &[&str]) -> Vec> { - dids.iter().map(|did| self.get_actor(did)).collect() - } - - /// Add or update an actor in the cache - pub fn upsert_actor(&self, actor: ActorMetadata) { - self.actors.insert(actor.did.clone(), actor); - } - - /// Remove an actor from the cache - pub fn remove_actor(&self, did: &str) { - self.actors.remove(did); - } - - /// Get the number of entries in the cache - pub fn len(&self) -> usize { - self.actors.len() - } - - /// Check if the cache is empty - pub fn is_empty(&self) -> bool { - self.actors.is_empty() - } -} - -impl Default for CachedActorStore { - fn default() -> Self { - Self::new() - } -} - -#[cfg(test)] -mod tests { - use super::*; - - fn make_test_actor(did: &str) -> ActorMetadata { - ActorMetadata { - did: did.to_string(), - handle: Some(format!("{}.bsky.social", did)), - status: ActorStatus::Active, - sync_state: ActorSyncState::Synced, - repo_rev: Some("test_rev".to_string()), - } - } - - #[test] - fn test_contains_rejects_unknown_actors() { - let cache = CachedActorStore::new(); - - let actors = vec![ - make_test_actor("did:plc:user1"), - make_test_actor("did:plc:user2"), - ]; - - cache.update_cache(actors); - - // Unknown actor should not be found - assert!( - !cache.contains_actor("did:plc:unknown"), - "Should reject unknown actors" - ); - } - - #[test] - fn test_contains_finds_cached_actors() { - // Question: Does contains_actor correctly identify cached actors? - let cache = CachedActorStore::new(); - - let actors = vec![ - make_test_actor("did:plc:alice"), - make_test_actor("did:plc:bob"), - make_test_actor("did:plc:charlie"), - ]; - - cache.update_cache(actors); - - // All cached actors should be found - assert!(cache.contains_actor("did:plc:alice")); - assert!(cache.contains_actor("did:plc:bob")); - assert!(cache.contains_actor("did:plc:charlie")); - } - - #[test] - fn test_get_actor_returns_metadata() { - // Question: Does get_actor return the correct metadata? - let cache = CachedActorStore::new(); - - let actor = make_test_actor("did:plc:testuser"); - cache.update_cache(vec![actor.clone()]); - - let retrieved = cache.get_actor("did:plc:testuser"); - assert!(retrieved.is_some()); - - let retrieved = retrieved.unwrap(); - assert_eq!(retrieved.did, "did:plc:testuser"); - assert_eq!( - retrieved.handle, - Some("did:plc:testuser.bsky.social".to_string()) - ); - assert_eq!(retrieved.status, ActorStatus::Active); - } - - #[test] - fn test_upsert_adds_new_actor() { - // Question: Does upsert correctly add new actors? - let cache = CachedActorStore::new(); - - assert!(!cache.contains_actor("did:plc:newuser")); - - cache.upsert_actor(make_test_actor("did:plc:newuser")); - - assert!( - cache.contains_actor("did:plc:newuser"), - "Upserted actor should be found" - ); - } - - #[test] - fn test_upsert_updates_existing_actor() { - // Question: Does upsert correctly update existing actors? - let cache = CachedActorStore::new(); - - let mut actor = make_test_actor("did:plc:user1"); - cache.upsert_actor(actor.clone()); - - // Update the actor - actor.handle = Some("newhandle.bsky.social".to_string()); - cache.upsert_actor(actor); - - let retrieved = cache.get_actor("did:plc:user1").unwrap(); - assert_eq!(retrieved.handle, Some("newhandle.bsky.social".to_string())); - } - - #[test] - fn test_remove_actor_deletes_from_cache() { - // Question: Does remove_actor correctly remove actors from the cache? - let cache = CachedActorStore::new(); - - cache.upsert_actor(make_test_actor("did:plc:user1")); - cache.upsert_actor(make_test_actor("did:plc:user2")); - - assert!(cache.contains_actor("did:plc:user1")); - - cache.remove_actor("did:plc:user1"); - - assert!( - !cache.contains_actor("did:plc:user1"), - "Removed actor should not be found" - ); - assert!( - cache.contains_actor("did:plc:user2"), - "Other actors should still be present" - ); - } - - #[test] - fn test_batch_get() { - // Question: Does get_actors_batch correctly retrieve multiple actors? - let cache = CachedActorStore::new(); - - cache.update_cache(vec![ - make_test_actor("did:plc:alice"), - make_test_actor("did:plc:bob"), - ]); - - let dids = vec!["did:plc:alice", "did:plc:bob", "did:plc:unknown"]; - let results = cache.get_actors_batch(&dids); - - assert_eq!(results.len(), 3); - assert!(results[0].is_some()); - assert!(results[1].is_some()); - assert!(results[2].is_none()); - } - - #[test] - fn test_update_cache_replaces_old_data() { - // Question: Does update_cache correctly replace old data? - let cache = CachedActorStore::new(); - - cache.update_cache(vec![make_test_actor("did:plc:old")]); - assert!(cache.contains_actor("did:plc:old")); - - cache.update_cache(vec![make_test_actor("did:plc:new")]); - assert!(!cache.contains_actor("did:plc:old")); - assert!(cache.contains_actor("did:plc:new")); - } - - - #[test] - fn test_empty_cache() { - // Question: Does an empty cache correctly report being empty? - let cache = CachedActorStore::new(); - - assert_eq!(cache.len(), 0); - assert!(cache.is_empty()); - assert!(!cache.contains_actor("did:plc:anyone")); - } - -} diff --git a/parakeet-db/src/id_cache.rs b/parakeet-db/src/id_cache.rs deleted file mode 100644 index aa100ace..00000000 --- a/parakeet-db/src/id_cache.rs +++ /dev/null @@ -1,722 +0,0 @@ -//! ID cache using Moka for fast, lock-free lookups with TTL -//! -//! This module provides a thread-safe cache for ID mappings that are used -//! frequently throughout the system: -//! - DID → (actor_id, is_allowlisted) -//! - AT URI → post_id (i64) -//! -//! These are immutable mappings (IDs never change once created), making them -//! perfect for caching. The cache is populated on-demand as lookups occur. -//! -//! Uses Moka for efficient lock-free concurrent access with automatic TTL expiration. -//! Moka uses a lock-free concurrent hash table (not RwLocks like DashMap), providing -//! better performance under high contention. -//! -//! Memory usage with 50k capacity: -//! - Actors: 50k × ~60 bytes = ~3MB -//! - Posts: 50k × ~100 bytes = ~5MB -//! - Total: ~8MB max -//! -//! TTL: 5 minutes (balances freshness with database load reduction) -//! Performance benefit: PostgreSQL query (1-5ms) → Moka lookup (~10ns) = 100,000x faster - -use moka::future::Cache; -use std::collections::HashMap; -use std::time::Duration; - -/// Actor cache TTL: 5 minutes -/// Short enough to catch status changes, long enough to reduce DB load significantly -const ACTOR_CACHE_TTL_SECS: u64 = 300; - -/// Actor cache capacity: 50k entries -const ACTOR_CACHE_CAPACITY: u64 = 50_000; - -/// Post cache TTL: 24 hours -/// Posts are immutable once created, so longer TTL is safe -/// Increased from 10 minutes to 24 hours to reduce database load -const POST_CACHE_TTL_SECS: u64 = 86400; - -/// Post cache capacity: 50k entries -const POST_CACHE_CAPACITY: u64 = 50_000; - -/// Cached actor information (forward lookup: DID/handle → actor data) -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct CachedActor { - /// Internal actor ID - pub actor_id: i32, - /// Whether actor is allowlisted (sync_state IN ('synced', 'dirty', 'processing')) - pub is_allowlisted: bool, -} - -/// Cached actor data for reverse lookups (actor_id → actor data) -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct CachedActorData { - /// Actor DID - pub did: String, - /// Actor handle (optional, may not be set) - pub handle: Option, -} - -/// Cached post data for reverse lookups (post_id → post data) -/// -/// Use `CachedPostData::not_found()` to create a sentinel value for posts that don't exist. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct CachedPostData { - /// Actor ID who created the post (0 for "not found" sentinel) - pub actor_id: i32, - /// Actor DID who created the post (empty for "not found" sentinel) - pub did: String, - /// Post rkey (TID) (0 for "not found" sentinel) - pub rkey: i64, - /// Post CID bytes (empty for "not found" sentinel) - pub cid: Vec, -} - -impl CachedPostData { - /// Create a sentinel value representing "post not found" - /// - /// Use this to cache the fact that a post ID doesn't exist (deleted/missing). - /// This prevents repeated database queries for the same non-existent post. - pub fn not_found() -> Self { - Self { - actor_id: 0, - did: String::new(), - rkey: 0, - cid: Vec::new(), - } - } - - /// Check if this is a "not found" sentinel value - pub fn is_not_found(&self) -> bool { - self.actor_id == 0 && self.did.is_empty() - } -} - -/// A thread-safe cache for ID mappings -/// -/// Caches four types of mappings: -/// - Actor forward: DID/handle (String) → (actor_id, is_allowlisted) -/// - Actor reverse: actor_id (i32) → (did, handle) -/// - Post: AT URI (String) → post_id (i64) -/// - Post data: post_id (i64) → (actor_id, did, rkey, cid) for reverse lookups -/// -/// Uses Moka for lock-free concurrent access with automatic TTL expiration. -#[derive(Clone)] -pub struct IdCache { - /// DID/handle → (actor_id, is_allowlisted) - Forward lookup - actor_ids: Cache, - /// actor_id → (did, handle) - Reverse lookup - actor_data: Cache, - /// AT URI → internal post ID - post_ids: Cache, - /// post_id → (actor_id, did, rkey, cid) for reverse lookups - post_data: Cache, - /// (actor_id, rkey) → (did, cid) for natural key lookups - post_data_by_key: Cache<(i32, i64), CachedPostData>, -} - -impl IdCache { - /// Create a new ID cache with default TTL and capacity - /// - /// - Actor forward cache: 5min TTL, 50k capacity - /// - Actor reverse cache: 5min TTL, 50k capacity - /// - Post cache: 10min TTL, 50k capacity - /// - Post data cache: 10min TTL, 50k capacity - pub fn new() -> Self { - Self { - actor_ids: Cache::builder() - .max_capacity(ACTOR_CACHE_CAPACITY) - .time_to_live(Duration::from_secs(ACTOR_CACHE_TTL_SECS)) - .build(), - actor_data: Cache::builder() - .max_capacity(ACTOR_CACHE_CAPACITY) - .time_to_live(Duration::from_secs(ACTOR_CACHE_TTL_SECS)) - .build(), - post_ids: Cache::builder() - .max_capacity(POST_CACHE_CAPACITY) - .time_to_live(Duration::from_secs(POST_CACHE_TTL_SECS)) - .build(), - post_data: Cache::builder() - .max_capacity(POST_CACHE_CAPACITY) - .time_to_live(Duration::from_secs(POST_CACHE_TTL_SECS)) - .build(), - post_data_by_key: Cache::builder() - .max_capacity(POST_CACHE_CAPACITY) - .time_to_live(Duration::from_secs(POST_CACHE_TTL_SECS)) - .build(), - } - } - - // ======================================================================== - // Actor ID Cache - // ======================================================================== - - /// Get actor info from cache - pub async fn get_actor_id(&self, did: &str) -> Option { - self.actor_ids.get(did).await - } - - /// Get actor ID (just the ID, ignoring allowlist status) - /// - /// Compatibility method for code that only needs the ID - pub async fn get_actor_id_only(&self, did: &str) -> Option { - self.actor_ids.get(did).await.map(|cached| cached.actor_id) - } - - /// Get multiple actor infos from cache - /// - /// Returns a HashMap with only the DIDs that were found in cache. - /// Missing DIDs are omitted from the result. - pub async fn get_actor_ids(&self, dids: &[String]) -> HashMap { - let mut result = HashMap::new(); - for did in dids { - if let Some(actor) = self.actor_ids.get(did).await { - result.insert(did.clone(), actor); - } - } - result - } - - /// Set actor info in cache - pub async fn set_actor_id(&self, did: String, actor: CachedActor) { - self.actor_ids.insert(did, actor).await; - } - - /// Set actor ID with allowlist status - /// - /// Convenience method that constructs CachedActor - pub async fn set_actor_id_with_allowlist( - &self, - did: String, - actor_id: i32, - is_allowlisted: bool, - ) { - self.actor_ids - .insert( - did, - CachedActor { - actor_id, - is_allowlisted, - }, - ) - .await; - } - - /// Set multiple actor infos in cache - pub async fn set_actor_ids(&self, mappings: &HashMap) { - for (did, actor) in mappings { - self.actor_ids.insert(did.clone(), *actor).await; - } - } - - /// Invalidate a specific actor from cache - /// - /// Use when actor status changes (e.g., tombstoned, sync state updated) - pub async fn invalidate_actor(&self, did: &str) { - self.actor_ids.invalidate(did).await; - } - - // ======================================================================== - // Actor Data Cache (Reverse Lookups: actor_id → actor data) - // ======================================================================== - - /// Get actor data from cache by actor ID (reverse lookup) - pub async fn get_actor_data(&self, actor_id: i32) -> Option { - self.actor_data.get(&actor_id).await - } - - /// Get multiple actor data from cache by actor IDs - /// - /// Returns a HashMap with only the actor IDs that were found in cache. - /// Missing actor IDs are omitted from the result. - pub async fn get_actor_data_many(&self, actor_ids: &[i32]) -> HashMap { - let mut result = HashMap::new(); - for actor_id in actor_ids { - if let Some(data) = self.actor_data.get(actor_id).await { - result.insert(*actor_id, data); - } - } - result - } - - /// Set actor data in cache - pub async fn set_actor_data(&self, actor_id: i32, data: CachedActorData) { - self.actor_data.insert(actor_id, data).await; - } - - /// Set multiple actor data in cache - pub async fn set_actor_data_many(&self, mappings: &HashMap) { - for (actor_id, data) in mappings { - self.actor_data.insert(*actor_id, data.clone()).await; - } - } - - /// Invalidate a specific actor data from cache by actor ID - pub async fn invalidate_actor_data(&self, actor_id: i32) { - self.actor_data.invalidate(&actor_id).await; - } - - /// Get the number of cached actor data entries - pub fn actor_data_len(&self) -> u64 { - self.actor_data.entry_count() - } - - // ======================================================================== - // Post ID Cache - // ======================================================================== - - /// Get post ID from cache - pub async fn get_post_id(&self, uri: &str) -> Option { - self.post_ids.get(uri).await - } - - /// Get multiple post IDs from cache - /// - /// Returns a HashMap with only the URIs that were found in cache. - /// Missing URIs are omitted from the result. - pub async fn get_post_ids(&self, uris: &[String]) -> HashMap { - let mut result = HashMap::new(); - for uri in uris { - if let Some(post_id) = self.post_ids.get(uri).await { - result.insert(uri.clone(), post_id); - } - } - result - } - - /// Set post ID in cache - pub async fn set_post_id(&self, uri: String, post_id: i64) { - self.post_ids.insert(uri, post_id).await; - } - - /// Set multiple post IDs in cache - pub async fn set_post_ids(&self, mappings: &HashMap) { - for (uri, post_id) in mappings { - self.post_ids.insert(uri.clone(), *post_id).await; - } - } - - /// Invalidate a specific post from cache - /// - /// Use when post is deleted or status changes - pub async fn invalidate_post(&self, uri: &str) { - self.post_ids.invalidate(uri).await; - } - - // ======================================================================== - // Post Data Cache (Reverse Lookups) - // ======================================================================== - - /// Get post data from cache by post ID (reverse lookup) - pub async fn get_post_data(&self, post_id: i64) -> Option { - self.post_data.get(&post_id).await - } - - /// Get multiple post data from cache by post IDs - /// - /// Returns a HashMap with only the post IDs that were found in cache. - /// Missing post IDs are omitted from the result. - pub async fn get_post_data_many(&self, post_ids: &[i64]) -> HashMap { - let mut result = HashMap::new(); - for post_id in post_ids { - if let Some(data) = self.post_data.get(post_id).await { - result.insert(*post_id, data); - } - } - result - } - - /// Set post data in cache - pub async fn set_post_data(&self, post_id: i64, data: CachedPostData) { - self.post_data.insert(post_id, data).await; - } - - /// Set multiple post data in cache - pub async fn set_post_data_many(&self, mappings: &HashMap) { - for (post_id, data) in mappings { - self.post_data.insert(*post_id, data.clone()).await; - } - } - - /// Invalidate a specific post data from cache by post ID - pub async fn invalidate_post_data(&self, post_id: i64) { - self.post_data.invalidate(&post_id).await; - } - - // ======================================================================== - // Post Data By Natural Key Cache - // ======================================================================== - - /// Get post data from cache by natural key (actor_id, rkey) - pub async fn get_post_data_by_key(&self, actor_id: i32, rkey: i64) -> Option { - self.post_data_by_key.get(&(actor_id, rkey)).await - } - - /// Get multiple post data from cache by natural keys - /// - /// Returns a HashMap with only the keys that were found in cache. - /// Missing keys are omitted from the result. - pub async fn get_post_data_by_keys_many(&self, keys: &[(i32, i64)]) -> HashMap<(i32, i64), CachedPostData> { - let mut result = HashMap::new(); - for key in keys { - if let Some(data) = self.post_data_by_key.get(key).await { - result.insert(*key, data); - } - } - result - } - - /// Set post data in cache by natural key - pub async fn set_post_data_by_key(&self, actor_id: i32, rkey: i64, data: CachedPostData) { - self.post_data_by_key.insert((actor_id, rkey), data).await; - } - - /// Set multiple post data in cache by natural keys - pub async fn set_post_data_by_keys_many(&self, mappings: &HashMap<(i32, i64), CachedPostData>) { - for (key, data) in mappings { - self.post_data_by_key.insert(*key, data.clone()).await; - } - } - - /// Invalidate a specific post data from cache by natural key - pub async fn invalidate_post_data_by_key(&self, actor_id: i32, rkey: i64) { - self.post_data_by_key.invalidate(&(actor_id, rkey)).await; - } - - // ======================================================================== - // Statistics - // ======================================================================== - - /// Get the number of cached actor IDs - /// - /// Note: This includes entries that may have expired but not yet been evicted. - /// Moka evicts entries lazily. - pub fn actor_ids_len(&self) -> u64 { - self.actor_ids.entry_count() - } - - /// Get the number of cached post IDs - /// - /// Note: This includes entries that may have expired but not yet been evicted. - /// Moka evicts entries lazily. - pub fn post_ids_len(&self) -> u64 { - self.post_ids.entry_count() - } - - /// Get the number of cached post data entries - /// - /// Note: This includes entries that may have expired but not yet been evicted. - /// Moka evicts entries lazily. - pub fn post_data_len(&self) -> u64 { - self.post_data.entry_count() - } - - /// Estimate memory usage in bytes - /// - /// This is a rough estimate: - /// - Actor IDs: ~60 bytes per DID (string + CachedActor) - /// - Post IDs: ~100 bytes per URI (string + i64) - pub fn estimated_memory_bytes(&self) -> u64 { - // Rough estimates including Moka overhead - let actor_memory = self.actor_ids.entry_count() * 60; - let post_memory = self.post_ids.entry_count() * 100; - actor_memory + post_memory - } -} - -impl Default for IdCache { - fn default() -> Self { - Self::new() - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[tokio::test] - async fn test_actor_id_cache() { - let cache = IdCache::new(); - - // Empty cache - assert_eq!(cache.get_actor_id("did:plc:alice").await, None); - - // Set and get - cache - .set_actor_id( - "did:plc:alice".to_string(), - CachedActor { - actor_id: 123, - is_allowlisted: true, - }, - ) - .await; - let result = cache.get_actor_id("did:plc:alice").await; - assert!(result.is_some()); - let actor = result.unwrap(); - assert_eq!(actor.actor_id, 123); - assert!(actor.is_allowlisted); - - // Update - cache - .set_actor_id( - "did:plc:alice".to_string(), - CachedActor { - actor_id: 456, - is_allowlisted: false, - }, - ) - .await; - let result = cache.get_actor_id("did:plc:alice").await; - assert!(result.is_some()); - let actor = result.unwrap(); - assert_eq!(actor.actor_id, 456); - assert!(!actor.is_allowlisted); - } - - #[tokio::test] - async fn test_actor_id_only() { - let cache = IdCache::new(); - - cache - .set_actor_id_with_allowlist("did:plc:alice".to_string(), 123, true) - .await; - - assert_eq!( - cache.get_actor_id_only("did:plc:alice").await, - Some(123) - ); - } - - #[tokio::test] - async fn test_post_id_cache() { - let cache = IdCache::new(); - - let uri = "at://did:plc:alice/app.bsky.feed.post/3jzfcijpj2z2a"; - - // Empty cache - assert_eq!(cache.get_post_id(uri).await, None); - - // Set and get - cache.set_post_id(uri.to_string(), 999).await; - assert_eq!(cache.get_post_id(uri).await, Some(999)); - } - - #[tokio::test] - async fn test_batch_actor_ids() { - let cache = IdCache::new(); - - cache - .set_actor_id( - "did:plc:alice".to_string(), - CachedActor { - actor_id: 1, - is_allowlisted: true, - }, - ) - .await; - cache - .set_actor_id( - "did:plc:bob".to_string(), - CachedActor { - actor_id: 2, - is_allowlisted: false, - }, - ) - .await; - cache - .set_actor_id( - "did:plc:charlie".to_string(), - CachedActor { - actor_id: 3, - is_allowlisted: true, - }, - ) - .await; - - let dids = vec![ - "did:plc:alice".to_string(), - "did:plc:bob".to_string(), - "did:plc:unknown".to_string(), - ]; - - let result = cache.get_actor_ids(&dids).await; - - assert_eq!(result.len(), 2); - assert_eq!(result.get("did:plc:alice").unwrap().actor_id, 1); - assert!(result.get("did:plc:alice").unwrap().is_allowlisted); - assert_eq!(result.get("did:plc:bob").unwrap().actor_id, 2); - assert!(!result.get("did:plc:bob").unwrap().is_allowlisted); - assert_eq!(result.get("did:plc:unknown"), None); - } - - #[tokio::test] - async fn test_batch_post_ids() { - let cache = IdCache::new(); - - let uri1 = "at://did:plc:alice/app.bsky.feed.post/abc123".to_string(); - let uri2 = "at://did:plc:bob/app.bsky.feed.post/def456".to_string(); - let uri3 = "at://did:plc:unknown/app.bsky.feed.post/xyz789".to_string(); - - cache.set_post_id(uri1.clone(), 100).await; - cache.set_post_id(uri2.clone(), 200).await; - - let uris = vec![uri1.clone(), uri2.clone(), uri3.clone()]; - - let result = cache.get_post_ids(&uris).await; - - assert_eq!(result.len(), 2); - assert_eq!(result.get(&uri1), Some(&100)); - assert_eq!(result.get(&uri2), Some(&200)); - assert_eq!(result.get(&uri3), None); - } - - #[tokio::test] - async fn test_set_actor_ids_batch() { - let cache = IdCache::new(); - - let mut mappings = HashMap::new(); - mappings.insert( - "did:plc:alice".to_string(), - CachedActor { - actor_id: 1, - is_allowlisted: true, - }, - ); - mappings.insert( - "did:plc:bob".to_string(), - CachedActor { - actor_id: 2, - is_allowlisted: false, - }, - ); - - cache.set_actor_ids(&mappings).await; - - assert_eq!( - cache.get_actor_id("did:plc:alice").await.unwrap().actor_id, - 1 - ); - assert_eq!( - cache.get_actor_id("did:plc:bob").await.unwrap().actor_id, - 2 - ); - } - - #[tokio::test] - async fn test_set_post_ids_batch() { - let cache = IdCache::new(); - - let mut mappings = HashMap::new(); - let uri1 = "at://did:plc:alice/app.bsky.feed.post/abc123".to_string(); - let uri2 = "at://did:plc:bob/app.bsky.feed.post/def456".to_string(); - mappings.insert(uri1.clone(), 100); - mappings.insert(uri2.clone(), 200); - - cache.set_post_ids(&mappings).await; - - assert_eq!(cache.get_post_id(&uri1).await, Some(100)); - assert_eq!(cache.get_post_id(&uri2).await, Some(200)); - } - - #[tokio::test] - async fn test_statistics() { - let cache = IdCache::new(); - - assert_eq!(cache.actor_ids_len(), 0); - assert_eq!(cache.post_ids_len(), 0); - - cache - .set_actor_id( - "did:plc:alice".to_string(), - CachedActor { - actor_id: 1, - is_allowlisted: true, - }, - ) - .await; - cache - .set_post_id("at://did:plc:bob/app.bsky.feed.post/123".to_string(), 100) - .await; - - // Sync pending tasks (Moka processes inserts asynchronously) - cache.actor_ids.run_pending_tasks().await; - cache.post_ids.run_pending_tasks().await; - - assert_eq!(cache.actor_ids_len(), 1); - assert_eq!(cache.post_ids_len(), 1); - - let mem = cache.estimated_memory_bytes(); - assert!(mem > 0); - assert!(mem < 1000); // Should be small with just 2 entries - } - - #[tokio::test] - async fn test_invalidate_actor() { - let cache = IdCache::new(); - - cache - .set_actor_id( - "did:plc:alice".to_string(), - CachedActor { - actor_id: 123, - is_allowlisted: true, - }, - ) - .await; - assert!(cache.get_actor_id("did:plc:alice").await.is_some()); - - cache.invalidate_actor("did:plc:alice").await; - assert_eq!(cache.get_actor_id("did:plc:alice").await, None); - } - - #[tokio::test] - async fn test_invalidate_post() { - let cache = IdCache::new(); - - let uri = "at://did:plc:alice/app.bsky.feed.post/123"; - cache.set_post_id(uri.to_string(), 999).await; - assert_eq!(cache.get_post_id(uri).await, Some(999)); - - cache.invalidate_post(uri).await; - assert_eq!(cache.get_post_id(uri).await, None); - } - - #[tokio::test] - async fn test_concurrent_access() { - use std::sync::Arc; - - let cache = Arc::new(IdCache::new()); - let mut handles = vec![]; - - // Spawn 10 tasks that all write to the cache - for i in 0..10 { - let cache = cache.clone(); - let handle = tokio::spawn(async move { - for j in 0..100 { - let did = format!("did:plc:user{}_{}", i, j); - cache - .set_actor_id( - did, - CachedActor { - actor_id: i * 100 + j, - is_allowlisted: j % 2 == 0, - }, - ) - .await; - } - }); - handles.push(handle); - } - - for handle in handles { - handle.await.unwrap(); - } - - // Sync pending tasks (Moka processes inserts asynchronously) - cache.actor_ids.run_pending_tasks().await; - - // Should have 1000 entries (10 tasks × 100 entries each) - assert_eq!(cache.actor_ids_len(), 1000); - } -} diff --git a/parakeet-db/src/lib.rs b/parakeet-db/src/lib.rs index 8e7fffda..70cf9215 100644 --- a/parakeet-db/src/lib.rs +++ b/parakeet-db/src/lib.rs @@ -1,9 +1,7 @@ -pub mod actor_cache; pub mod at_uri_util; pub mod cid_util; pub mod composite_types; pub mod compression; -pub mod id_cache; pub mod jacquard_views; pub mod models; pub mod notifications;