diff --git a/lexica/src/app_bsky/actor.rs b/lexica/src/app_bsky/actor.rs index 9b63526f..0fffa6c5 100644 --- a/lexica/src/app_bsky/actor.rs +++ b/lexica/src/app_bsky/actor.rs @@ -197,7 +197,7 @@ pub struct ProfileView { pub indexed_at: DateTime, } -#[derive(Debug, Serialize, Deserialize)] +#[derive(Clone, Debug, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct ProfileViewDetailed { pub did: String, diff --git a/parakeet/src/cache_listener.rs b/parakeet/src/cache_listener.rs index 1522f7b3..d1e5d73e 100644 --- a/parakeet/src/cache_listener.rs +++ b/parakeet/src/cache_listener.rs @@ -136,45 +136,65 @@ async fn handle_cache_invalidation(state: &GlobalState, cache_key: &str) { } // Profile invalidation: "profile:{actor_id}" - if cache_key.starts_with("profile:") { - // No-op for now - we don't cache profiles yet - // When we add profile cache, parse actor_id and invalidate - return; + if let Some(actor_id_str) = cache_key.strip_prefix("profile:") { + if let Ok(actor_id) = actor_id_str.parse::() { + state.profile_cache.invalidate(actor_id).await; + info!(actor_id = actor_id, "Invalidated profile cache"); + return; + } } // Post invalidation: "post:{actor_id}:{rkey}" - if cache_key.starts_with("post:") { - // No-op for now - we don't cache individual posts yet - // When we add post cache, parse actor_id:rkey and invalidate - return; + if let Some(rest) = cache_key.strip_prefix("post:") { + if let Some((actor_id_str, rkey_str)) = rest.split_once(':') { + if let (Ok(actor_id), Ok(rkey)) = (actor_id_str.parse::(), rkey_str.parse::()) { + state.post_cache.invalidate(actor_id, rkey).await; + info!(actor_id = actor_id, rkey = rkey, "Invalidated post cache"); + return; + } + } } // Feedgen invalidation: "feedgen:{actor_id}:{rkey}" - if cache_key.starts_with("feedgen:") { - // No-op for now - we don't cache feedgens yet - // When we add feedgen cache, parse actor_id:rkey and invalidate - return; + if let Some(rest) = cache_key.strip_prefix("feedgen:") { + if let Some((actor_id_str, rkey)) = rest.split_once(':') { + if let Ok(actor_id) = actor_id_str.parse::() { + state.feedgen_cache.invalidate(actor_id, rkey).await; + info!(actor_id = actor_id, rkey = rkey, "Invalidated feedgen cache"); + return; + } + } } // List invalidation: "list:{actor_id}:{rkey}" - if cache_key.starts_with("list:") { - // No-op for now - we don't cache lists yet - // When we add list cache, parse actor_id:rkey and invalidate - return; + if let Some(rest) = cache_key.strip_prefix("list:") { + if let Some((actor_id_str, rkey)) = rest.split_once(':') { + if let Ok(actor_id) = actor_id_str.parse::() { + state.list_cache.invalidate(actor_id, rkey).await; + info!(actor_id = actor_id, rkey = rkey, "Invalidated list cache"); + return; + } + } } // Starterpack invalidation: "starterpack:{actor_id}:{rkey}" - if cache_key.starts_with("starterpack:") { - // No-op for now - we don't cache starterpacks yet - // When we add starterpack cache, parse actor_id:rkey and invalidate - return; + if let Some(rest) = cache_key.strip_prefix("starterpack:") { + if let Some((actor_id_str, rkey)) = rest.split_once(':') { + if let Ok(actor_id) = actor_id_str.parse::() { + state.starterpack_cache.invalidate(actor_id, rkey).await; + info!(actor_id = actor_id, rkey = rkey, "Invalidated starterpack cache"); + return; + } + } } // Labeler invalidation: "labeler:{actor_id}" - if cache_key.starts_with("labeler:") { - // No-op for now - we don't cache labelers yet - // When we add labeler cache, parse actor_id and invalidate - return; + if let Some(actor_id_str) = cache_key.strip_prefix("labeler:") { + if let Ok(actor_id) = actor_id_str.parse::() { + state.labeler_cache.invalidate(actor_id).await; + info!(actor_id = actor_id, "Invalidated labeler cache"); + return; + } } warn!(cache_key = cache_key, "Unknown cache invalidation pattern"); diff --git a/parakeet/src/entity_cache.rs b/parakeet/src/entity_cache.rs new file mode 100644 index 00000000..459c4456 --- /dev/null +++ b/parakeet/src/entity_cache.rs @@ -0,0 +1,805 @@ +//! Entity caching for various AT Protocol objects +//! +//! This is the ONLY approved way to get hydrated data in the application. +//! Direct hydration through StatefulHydrator should be avoided to ensure +//! all data access goes through the cache. +//! +//! ## Usage +//! +//! Instead of: +//! ``` +//! let hydrator = StatefulHydrator::new(...); +//! let profile = hydrator.hydrate_profile_detailed(did).await; // ❌ Bypasses cache! +//! ``` +//! +//! Always use: +//! ``` +//! let profile = profile_cache.get_or_hydrate(actor_id, did, &hydrator).await; // ✅ Uses cache! +//! ``` +//! +//! ## Cache Strategy +//! +//! We cache the expensive "Detailed" variants: +//! - `ProfileViewDetailed` - Full profile with counts, used in getProfile +//! - `PostView` - Full post with stats, embeds, and viewer state +//! - `GeneratorView`, `ListView`, etc. - Full entities with all metadata +//! +//! Cache invalidation is handled by database triggers via pg_notify. + +use moka::future::Cache; +use std::time::Duration; + +// Import the actual types used in the XRPC endpoints +use lexica::app_bsky::actor::ProfileViewDetailed; +use lexica::app_bsky::feed::PostView; +use lexica::app_bsky::feed::GeneratorView; +use lexica::app_bsky::graph::ListView; +use lexica::app_bsky::graph::StarterPackViewBasic; + +use crate::hydration::StatefulHydrator; +use crate::id_cache_helpers; +use diesel_async::pooled_connection::deadpool::Pool; +use diesel_async::AsyncPgConnection; +use parakeet_db::id_cache::IdCache; +use std::sync::Arc; + +/// Parsed components of an AT-URI +pub struct ParsedAtUri { + pub did: String, + pub collection: String, + pub rkey: String, +} + +/// Parse an AT-URI (at://did/collection/rkey) into its components +pub fn parse_at_uri(uri: &str) -> Option { + // Expected format: at://did/collection/rkey + let uri = uri.strip_prefix("at://")?; + let parts: Vec<&str> = uri.split('/').collect(); + + if parts.len() != 3 { + return None; + } + + Some(ParsedAtUri { + did: parts[0].to_string(), + collection: parts[1].to_string(), + rkey: parts[2].to_string(), + }) +} + +/// Profile cache that handles hydration automatically +/// +/// Uses actor_id as the cache key for consistency with database triggers +/// +/// Usage: +/// ``` +/// let profile = profile_cache.get_or_hydrate(actor_id, did, &hyd).await?; +/// ``` +#[derive(Clone)] +pub struct ProfileCache { + cache: Cache, // Key is actor_id +} + +impl ProfileCache { + pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { + let cache = Cache::builder() + .max_capacity(max_capacity) + .time_to_live(Duration::from_secs(ttl_secs)) + .support_invalidation_closures() + .build(); + + Self { cache } + } + + /// Get a profile from cache or hydrate it if not cached + /// + /// This is the primary API - it handles all caching logic internally + /// Takes both actor_id (for caching) and DID (for hydration) + pub async fn get_or_hydrate( + &self, + actor_id: i32, + did: String, + hydrator: &StatefulHydrator<'_>, + ) -> Option { + // Check cache first using actor_id + if let Some(profile) = self.cache.get(&actor_id).await { + tracing::debug!(actor_id, "Profile cache hit"); + return Some(profile); + } + + // Cache miss - hydrate the profile using DID + tracing::debug!(actor_id, did = %did, "Profile cache miss, hydrating"); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let profile = hydrator.hydrate_profile_detailed(did).await?; + + // Store in cache using actor_id as key + self.cache.insert(actor_id, profile.clone()).await; + tracing::debug!(actor_id, "Profile cached"); + + Some(profile) + } + + /// Get multiple profiles with caching + /// + /// This is the recommended way to get multiple profiles, ensuring each one + /// uses the cache individually + pub async fn get_or_hydrate_batch( + &self, + dids: Vec, + pool: &Pool, + id_cache: &Arc, + hydrator: &StatefulHydrator<'_>, + ) -> Vec { + let mut profiles = Vec::with_capacity(dids.len()); + + for did in dids { + // Get actor_id for each DID + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch(pool, id_cache, &did).await { + if let Some(profile) = self.get_or_hydrate(actor_id, did, hydrator).await { + profiles.push(profile); + } + } + } + + profiles + } + + pub async fn invalidate(&self, actor_id: i32) { + self.cache.invalidate(&actor_id).await; + tracing::debug!(actor_id, "Profile cache invalidated"); + } + + /// Invalidate all entries (for testing/admin) + pub async fn invalidate_all(&self) { + self.cache.invalidate_all(); + tracing::info!("All profile cache entries invalidated"); + } +} + +/// Post cache that handles hydration automatically +/// +/// Uses (actor_id, rkey) as the cache key for consistency with database triggers +/// +/// Usage: +/// ``` +/// let post = post_cache.get_or_hydrate(actor_id, rkey, uri, &hyd).await; +/// ``` +#[derive(Clone)] +pub struct PostCache { + cache: Cache<(i32, i64), PostView>, // Key is (actor_id, rkey) +} + +impl PostCache { + pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { + let cache = Cache::builder() + .max_capacity(max_capacity) + .time_to_live(Duration::from_secs(ttl_secs)) + .support_invalidation_closures() + .build(); + + Self { cache } + } + + /// Get a single post from cache or hydrate it + /// + /// Takes both (actor_id, rkey) for caching and URI for hydration + pub async fn get_or_hydrate_single( + &self, + actor_id: i32, + rkey: i64, + uri: String, + hydrator: &StatefulHydrator<'_>, + ) -> Option { + // Check cache first using (actor_id, rkey) + if let Some(post) = self.cache.get(&(actor_id, rkey)).await { + tracing::debug!(actor_id, rkey, "Post cache hit"); + return Some(post); + } + + // Cache miss - hydrate the post using URI + tracing::debug!(actor_id, rkey, uri = %uri, "Post cache miss, hydrating"); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let posts = hydrator.hydrate_posts(vec![uri]).await; + let post = posts.into_values().next()?; + + // Store in cache using (actor_id, rkey) as key + self.cache.insert((actor_id, rkey), post.clone()).await; + tracing::debug!(actor_id, rkey, "Post cached"); + + Some(post) + } + + /// Get multiple posts with caching by parsing URIs + /// + /// Parses AT-URIs (at://did/collection/rkey) to extract actor_id and rkey, + /// enabling individual post caching + /// + /// Returns a HashMap for efficient lookups while preserving the ability to iterate + pub async fn get_or_hydrate_from_uris( + &self, + uris: Vec, + pool: &Pool, + id_cache: &Arc, + hydrator: &StatefulHydrator<'_>, + ) -> std::collections::HashMap { + let mut results = Vec::with_capacity(uris.len()); + let mut missing_uris = Vec::new(); + let mut missing_indices = Vec::new(); + + // Try to get each post from cache + for (idx, uri) in uris.iter().enumerate() { + // Parse AT-URI: at://did/collection/rkey + if let Some(parsed) = parse_at_uri(uri) { + // Only process posts (app.bsky.feed.post collection) + if parsed.collection == "app.bsky.feed.post" { + // Get actor_id from DID + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + // Try to parse rkey as i64 (TID) + if let Ok(rkey) = parsed.rkey.parse::() { + // Check cache + if let Some(post) = self.cache.get(&(actor_id, rkey)).await { + tracing::debug!(actor_id, rkey, "Post cache hit"); + results.push(Some(post)); + continue; + } + } + } + } + } + + // Cache miss or couldn't parse - need to hydrate + missing_uris.push(uri.clone()); + missing_indices.push(idx); + results.push(None); + } + + // Hydrate missing posts if any + if !missing_uris.is_empty() { + tracing::debug!("Hydrating {} missing posts", missing_uris.len()); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let hydrated = hydrator.hydrate_posts(missing_uris).await; + + // Store hydrated posts in cache and results + for (idx, uri) in missing_indices.into_iter().zip(hydrated.keys()) { + if let Some(post) = hydrated.get(uri) { + // Try to cache it if we can parse the URI + if let Some(parsed) = parse_at_uri(uri) { + if parsed.collection == "app.bsky.feed.post" { + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + if let Ok(rkey) = parsed.rkey.parse::() { + self.cache.insert((actor_id, rkey), post.clone()).await; + tracing::debug!(actor_id, rkey, "Post cached"); + } + } + } + } + + results[idx] = Some(post.clone()); + } + } + } + + // Build HashMap from successfully loaded posts + let mut posts_map = std::collections::HashMap::new(); + for (uri, post_opt) in uris.into_iter().zip(results.into_iter()) { + if let Some(post) = post_opt { + posts_map.insert(uri, post); + } + } + posts_map + } + + pub async fn invalidate(&self, actor_id: i32, rkey: i64) { + self.cache.invalidate(&(actor_id, rkey)).await; + tracing::debug!(actor_id, rkey, "Post cache invalidated"); + } + + /// Invalidate all entries (for testing/admin) + pub async fn invalidate_all(&self) { + self.cache.invalidate_all(); + tracing::info!("All post cache entries invalidated"); + } +} + +/// Feedgen cache using (actor_id, rkey) as key +#[derive(Clone)] +pub struct FeedgenCache { + cache: Cache<(i32, String), GeneratorView>, +} + +impl FeedgenCache { + pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { + let cache = Cache::builder() + .max_capacity(max_capacity) + .time_to_live(Duration::from_secs(ttl_secs)) + .support_invalidation_closures() + .build(); + + Self { cache } + } + + /// Get a single feedgen from cache or hydrate it + pub async fn get_or_hydrate_single( + &self, + actor_id: i32, + rkey: String, + uri: String, + hydrator: &StatefulHydrator<'_>, + ) -> Option { + // Check cache first using (actor_id, rkey) + if let Some(feedgen) = self.cache.get(&(actor_id, rkey.clone())).await { + tracing::debug!(actor_id, rkey, "Feedgen cache hit"); + return Some(feedgen); + } + + // Cache miss - hydrate the feedgen using URI + tracing::debug!(actor_id, rkey, uri = %uri, "Feedgen cache miss, hydrating"); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let feedgen = hydrator.hydrate_feedgen(uri).await?; + + // Store in cache using (actor_id, rkey) as key + self.cache.insert((actor_id, rkey.clone()), feedgen.clone()).await; + tracing::debug!(actor_id, rkey, "Feedgen cached"); + + Some(feedgen) + } + + /// Get multiple feedgens with caching by parsing URIs + /// Returns a HashMap for efficient lookups + pub async fn get_or_hydrate_from_uris( + &self, + uris: Vec, + pool: &Pool, + id_cache: &Arc, + hydrator: &StatefulHydrator<'_>, + ) -> std::collections::HashMap { + let mut results = Vec::with_capacity(uris.len()); + let mut missing_uris = Vec::new(); + let mut missing_indices = Vec::new(); + + // Try to get each feedgen from cache + for (idx, uri) in uris.iter().enumerate() { + // Parse AT-URI: at://did/collection/rkey + if let Some(parsed) = parse_at_uri(uri) { + // Only process feedgens (app.bsky.feed.generator collection) + if parsed.collection == "app.bsky.feed.generator" { + // Get actor_id from DID + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + // Check cache + if let Some(feedgen) = self.cache.get(&(actor_id, parsed.rkey.clone())).await { + tracing::debug!(actor_id, rkey = parsed.rkey, "Feedgen cache hit"); + results.push(Some(feedgen)); + continue; + } + } + } + } + + // Cache miss or couldn't parse - need to hydrate + missing_uris.push(uri.clone()); + missing_indices.push(idx); + results.push(None); + } + + // Hydrate missing feedgens if any + if !missing_uris.is_empty() { + tracing::debug!("Hydrating {} missing feedgens", missing_uris.len()); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let hydrated = hydrator.hydrate_feedgens(missing_uris).await; + + // Store hydrated feedgens in cache and results + for (idx, uri) in missing_indices.into_iter().zip(hydrated.keys()) { + if let Some(feedgen) = hydrated.get(uri) { + // Try to cache it if we can parse the URI + if let Some(parsed) = parse_at_uri(uri) { + if parsed.collection == "app.bsky.feed.generator" { + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + self.cache.insert((actor_id, parsed.rkey.clone()), feedgen.clone()).await; + tracing::debug!(actor_id, rkey = parsed.rkey, "Feedgen cached"); + } + } + } + + results[idx] = Some(feedgen.clone()); + } + } + } + + // Build HashMap from successfully loaded feedgens + let mut feedgens_map = std::collections::HashMap::new(); + for (uri, feedgen_opt) in uris.into_iter().zip(results.into_iter()) { + if let Some(feedgen) = feedgen_opt { + feedgens_map.insert(uri, feedgen); + } + } + feedgens_map + } + + pub async fn get(&self, actor_id: i32, rkey: &str) -> Option { + let feedgen = self.cache.get(&(actor_id, rkey.to_string())).await?; + tracing::debug!(actor_id, rkey, "Feedgen cache hit"); + Some(feedgen) + } + + pub async fn set(&self, actor_id: i32, rkey: &str, feedgen: GeneratorView) { + self.cache.insert((actor_id, rkey.to_string()), feedgen).await; + tracing::debug!(actor_id, rkey, "Feedgen cached"); + } + + pub async fn invalidate(&self, actor_id: i32, rkey: &str) { + self.cache.invalidate(&(actor_id, rkey.to_string())).await; + tracing::debug!(actor_id, rkey, "Feedgen cache invalidated"); + } + + /// Invalidate all entries (for testing/admin) + pub async fn invalidate_all(&self) { + self.cache.invalidate_all(); + tracing::info!("All feedgen cache entries invalidated"); + } +} + +/// List cache using (actor_id, rkey) as key +#[derive(Clone)] +pub struct ListCache { + cache: Cache<(i32, String), ListView>, +} + +impl ListCache { + pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { + let cache = Cache::builder() + .max_capacity(max_capacity) + .time_to_live(Duration::from_secs(ttl_secs)) + .support_invalidation_closures() + .build(); + + Self { cache } + } + + /// Get a single list from cache or hydrate it + pub async fn get_or_hydrate_single( + &self, + actor_id: i32, + rkey: String, + uri: String, + hydrator: &StatefulHydrator<'_>, + ) -> Option { + // Check cache first using (actor_id, rkey) + if let Some(list) = self.cache.get(&(actor_id, rkey.clone())).await { + tracing::debug!(actor_id, rkey, "List cache hit"); + return Some(list); + } + + // Cache miss - hydrate the list using URI + tracing::debug!(actor_id, rkey, uri = %uri, "List cache miss, hydrating"); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let list = hydrator.hydrate_list(uri).await?; + + // Store in cache using (actor_id, rkey) as key + self.cache.insert((actor_id, rkey.clone()), list.clone()).await; + tracing::debug!(actor_id, rkey, "List cached"); + + Some(list) + } + + /// Get multiple lists with caching by parsing URIs + /// Returns a HashMap for efficient lookups + pub async fn get_or_hydrate_from_uris( + &self, + uris: Vec, + pool: &Pool, + id_cache: &Arc, + hydrator: &StatefulHydrator<'_>, + ) -> std::collections::HashMap { + let mut results = Vec::with_capacity(uris.len()); + let mut missing_uris = Vec::new(); + let mut missing_indices = Vec::new(); + + // Try to get each list from cache + for (idx, uri) in uris.iter().enumerate() { + // Parse AT-URI: at://did/collection/rkey + if let Some(parsed) = parse_at_uri(uri) { + // Only process lists (app.bsky.graph.list collection) + if parsed.collection == "app.bsky.graph.list" { + // Get actor_id from DID + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + // Check cache + if let Some(list) = self.cache.get(&(actor_id, parsed.rkey.clone())).await { + tracing::debug!(actor_id, rkey = parsed.rkey, "List cache hit"); + results.push(Some(list)); + continue; + } + } + } + } + + // Cache miss or couldn't parse - need to hydrate + missing_uris.push(uri.clone()); + missing_indices.push(idx); + results.push(None); + } + + // Hydrate missing lists if any + if !missing_uris.is_empty() { + tracing::debug!("Hydrating {} missing lists", missing_uris.len()); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let hydrated = hydrator.hydrate_lists(missing_uris).await; + + // Store hydrated lists in cache and results + for (idx, uri) in missing_indices.into_iter().zip(hydrated.keys()) { + if let Some(list) = hydrated.get(uri) { + // Try to cache it if we can parse the URI + if let Some(parsed) = parse_at_uri(uri) { + if parsed.collection == "app.bsky.graph.list" { + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + self.cache.insert((actor_id, parsed.rkey.clone()), list.clone()).await; + tracing::debug!(actor_id, rkey = parsed.rkey, "List cached"); + } + } + } + + results[idx] = Some(list.clone()); + } + } + } + + // Build HashMap from successfully loaded lists + let mut lists_map = std::collections::HashMap::new(); + for (uri, list_opt) in uris.into_iter().zip(results.into_iter()) { + if let Some(list) = list_opt { + lists_map.insert(uri, list); + } + } + lists_map + } + + pub async fn get(&self, actor_id: i32, rkey: &str) -> Option { + let list = self.cache.get(&(actor_id, rkey.to_string())).await?; + tracing::debug!(actor_id, rkey, "List cache hit"); + Some(list) + } + + pub async fn set(&self, actor_id: i32, rkey: &str, list: ListView) { + self.cache.insert((actor_id, rkey.to_string()), list).await; + tracing::debug!(actor_id, rkey, "List cached"); + } + + pub async fn invalidate(&self, actor_id: i32, rkey: &str) { + self.cache.invalidate(&(actor_id, rkey.to_string())).await; + tracing::debug!(actor_id, rkey, "List cache invalidated"); + } + + /// Invalidate all entries (for testing/admin) + pub async fn invalidate_all(&self) { + self.cache.invalidate_all(); + tracing::info!("All list cache entries invalidated"); + } +} + +/// Starterpack cache using (actor_id, rkey) as key +#[derive(Clone)] +pub struct StarterpackCache { + cache: Cache<(i32, String), StarterPackViewBasic>, +} + +impl StarterpackCache { + pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { + let cache = Cache::builder() + .max_capacity(max_capacity) + .time_to_live(Duration::from_secs(ttl_secs)) + .support_invalidation_closures() + .build(); + + Self { cache } + } + + /// Get a single starterpack from cache or hydrate it + pub async fn get_or_hydrate_single( + &self, + actor_id: i32, + rkey: String, + uri: String, + hydrator: &StatefulHydrator<'_>, + ) -> Option { + // Check cache first using (actor_id, rkey) + if let Some(pack) = self.cache.get(&(actor_id, rkey.clone())).await { + tracing::debug!(actor_id, rkey, "Starterpack cache hit"); + return Some(pack); + } + + // Cache miss - hydrate the starterpack using URI + tracing::debug!(actor_id, rkey, uri = %uri, "Starterpack cache miss, hydrating"); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let packs = hydrator.hydrate_starterpacks_basic(vec![uri]).await; + let pack = packs.into_values().next()?; + + // Store in cache using (actor_id, rkey) as key + self.cache.insert((actor_id, rkey.clone()), pack.clone()).await; + tracing::debug!(actor_id, rkey, "Starterpack cached"); + + Some(pack) + } + + /// Get multiple starterpacks with caching by parsing URIs + /// Returns a HashMap for efficient lookups + pub async fn get_or_hydrate_from_uris( + &self, + uris: Vec, + pool: &Pool, + id_cache: &Arc, + hydrator: &StatefulHydrator<'_>, + ) -> std::collections::HashMap { + let mut results = Vec::with_capacity(uris.len()); + let mut missing_uris = Vec::new(); + let mut missing_indices = Vec::new(); + + // Try to get each starterpack from cache + for (idx, uri) in uris.iter().enumerate() { + // Parse AT-URI: at://did/collection/rkey + if let Some(parsed) = parse_at_uri(uri) { + // Only process starterpacks (app.bsky.graph.starterpack collection) + if parsed.collection == "app.bsky.graph.starterpack" { + // Get actor_id from DID + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + // Check cache + if let Some(pack) = self.cache.get(&(actor_id, parsed.rkey.clone())).await { + tracing::debug!(actor_id, rkey = parsed.rkey, "Starterpack cache hit"); + results.push(Some(pack)); + continue; + } + } + } + } + + // Cache miss or couldn't parse - need to hydrate + missing_uris.push(uri.clone()); + missing_indices.push(idx); + results.push(None); + } + + // Hydrate missing starterpacks if any + if !missing_uris.is_empty() { + tracing::debug!("Hydrating {} missing starterpacks", missing_uris.len()); + + // We're allowed to call the deprecated method here since we're the cache + #[allow(deprecated)] + let hydrated = hydrator.hydrate_starterpacks_basic(missing_uris).await; + + // Store hydrated starterpacks in cache and results + for (idx, uri) in missing_indices.into_iter().zip(hydrated.keys()) { + if let Some(pack) = hydrated.get(uri) { + // Try to cache it if we can parse the URI + if let Some(parsed) = parse_at_uri(uri) { + if parsed.collection == "app.bsky.graph.starterpack" { + if let Ok(actor_id) = id_cache_helpers::get_actor_id_or_fetch( + pool, + id_cache, + &parsed.did + ).await { + self.cache.insert((actor_id, parsed.rkey.clone()), pack.clone()).await; + tracing::debug!(actor_id, rkey = parsed.rkey, "Starterpack cached"); + } + } + } + + results[idx] = Some(pack.clone()); + } + } + } + + // Build HashMap from successfully loaded starterpacks + let mut packs_map = std::collections::HashMap::new(); + for (uri, pack_opt) in uris.into_iter().zip(results.into_iter()) { + if let Some(pack) = pack_opt { + packs_map.insert(uri, pack); + } + } + packs_map + } + + pub async fn get(&self, actor_id: i32, rkey: &str) -> Option { + let pack = self.cache.get(&(actor_id, rkey.to_string())).await?; + tracing::debug!(actor_id, rkey, "Starterpack cache hit"); + Some(pack) + } + + pub async fn set(&self, actor_id: i32, rkey: &str, pack: StarterPackViewBasic) { + self.cache.insert((actor_id, rkey.to_string()), pack).await; + tracing::debug!(actor_id, rkey, "Starterpack cached"); + } + + pub async fn invalidate(&self, actor_id: i32, rkey: &str) { + self.cache.invalidate(&(actor_id, rkey.to_string())).await; + tracing::debug!(actor_id, rkey, "Starterpack cache invalidated"); + } + + /// Invalidate all entries (for testing/admin) + pub async fn invalidate_all(&self) { + self.cache.invalidate_all(); + tracing::info!("All starterpack cache entries invalidated"); + } +} + +/// Labeler cache using actor_id as key +/// Note: Labelers always use "self" as rkey +#[derive(Clone)] +pub struct LabelerCache { + cache: Cache, // Using Value since we don't have HydratedLabeler type +} + +impl LabelerCache { + pub fn new(ttl_secs: u64, max_capacity: u64) -> Self { + let cache = Cache::builder() + .max_capacity(max_capacity) + .time_to_live(Duration::from_secs(ttl_secs)) + .support_invalidation_closures() + .build(); + + Self { cache } + } + + pub async fn get(&self, actor_id: i32) -> Option { + let labeler = self.cache.get(&actor_id).await?; + tracing::debug!(actor_id, "Labeler cache hit"); + Some(labeler) + } + + pub async fn set(&self, actor_id: i32, labeler: serde_json::Value) { + self.cache.insert(actor_id, labeler).await; + tracing::debug!(actor_id, "Labeler cached"); + } + + pub async fn invalidate(&self, actor_id: i32) { + self.cache.invalidate(&actor_id).await; + tracing::debug!(actor_id, "Labeler cache invalidated"); + } + + /// Invalidate all entries (for testing/admin) + pub async fn invalidate_all(&self) { + self.cache.invalidate_all(); + tracing::info!("All labeler cache entries invalidated"); + } +} \ No newline at end of file diff --git a/parakeet/src/hydration/feedgen.rs b/parakeet/src/hydration/feedgen.rs index fed7f469..0fdbf3c1 100644 --- a/parakeet/src/hydration/feedgen.rs +++ b/parakeet/src/hydration/feedgen.rs @@ -53,6 +53,10 @@ fn build_feedgen( } impl super::StatefulHydrator<'_> { + #[deprecated( + since = "0.1.0", + note = "Use FeedgenCache::get_or_hydrate_single() to ensure caching. Direct hydration bypasses the cache." + )] pub async fn hydrate_feedgen(&self, feedgen: String) -> Option { let labels = self.get_label(&feedgen).await; let viewer = self.get_feedgen_viewer_state(&feedgen).await; @@ -85,6 +89,10 @@ impl super::StatefulHydrator<'_> { )) } + #[deprecated( + since = "0.1.0", + note = "Use FeedgenCache::get_or_hydrate_from_uris() to ensure caching. Direct hydration bypasses the cache." + )] pub async fn hydrate_feedgens(&self, feedgens: Vec) -> HashMap { let labels = self.get_label_many(&feedgens).await; let viewers = self.get_feedgen_viewer_states(&feedgens).await; diff --git a/parakeet/src/hydration/list.rs b/parakeet/src/hydration/list.rs index 6cb485cd..3c370372 100644 --- a/parakeet/src/hydration/list.rs +++ b/parakeet/src/hydration/list.rs @@ -160,6 +160,10 @@ impl StatefulHydrator<'_> { .collect() } + #[deprecated( + since = "0.1.0", + note = "Use ListCache::get_or_hydrate_single() to ensure caching. Direct hydration bypasses the cache." + )] pub async fn hydrate_list(&self, list: String) -> Option { let labels = self.get_label(&list).await; let viewer = self.get_list_viewer_state(&list).await; @@ -189,6 +193,10 @@ impl StatefulHydrator<'_> { build_listview(enriched, count, profile, labels, viewer, &self.cdn) } + #[deprecated( + since = "0.1.0", + note = "Use ListCache::get_or_hydrate_from_uris() to ensure caching. Direct hydration bypasses the cache." + )] pub async fn hydrate_lists(&self, lists: Vec) -> HashMap { if lists.is_empty() { return HashMap::new(); diff --git a/parakeet/src/hydration/posts/mod.rs b/parakeet/src/hydration/posts/mod.rs index 4fe822e0..6e588fa5 100644 --- a/parakeet/src/hydration/posts/mod.rs +++ b/parakeet/src/hydration/posts/mod.rs @@ -308,6 +308,14 @@ impl StatefulHydrator<'_> { (results, reply_actor_cache) } + /// Hydrate multiple posts + /// + /// **DEPRECATED**: This method bypasses caching. Use PostCache::get_or_hydrate_from_uris() instead. + /// Direct hydration should only be called from within the cache system. + #[deprecated( + since = "0.1.0", + note = "Use PostCache::get_or_hydrate_from_uris() to ensure caching. Direct hydration bypasses the cache." + )] pub async fn hydrate_posts(&self, posts: Vec) -> HashMap { let (posts_data, actor_cache) = self.hydrate_posts_inner(posts).await; diff --git a/parakeet/src/hydration/profile/mod.rs b/parakeet/src/hydration/profile/mod.rs index f098e612..5fbaa6fd 100644 --- a/parakeet/src/hydration/profile/mod.rs +++ b/parakeet/src/hydration/profile/mod.rs @@ -280,6 +280,14 @@ impl super::StatefulHydrator<'_> { .collect() } + /// Get detailed profile data + /// + /// **DEPRECATED**: This method bypasses caching. Use ProfileCache::get_or_hydrate() instead. + /// Direct hydration should only be called from within the cache system. + #[deprecated( + since = "0.1.0", + note = "Use ProfileCache::get_or_hydrate() to ensure caching. Direct hydration bypasses the cache." + )] pub async fn hydrate_profile_detailed(&self, did: String) -> Option { let labels = self.get_profile_label(&did).await; let viewer = self.get_profile_viewer_state(&did).await; diff --git a/parakeet/src/hydration/starter_packs.rs b/parakeet/src/hydration/starter_packs.rs index 91b3419c..2bbd3fc5 100644 --- a/parakeet/src/hydration/starter_packs.rs +++ b/parakeet/src/hydration/starter_packs.rs @@ -99,6 +99,10 @@ impl StatefulHydrator<'_> { Some(build_basic(enriched, creator, labels, list_item_count)) } + #[deprecated( + since = "0.1.0", + note = "Use StarterpackCache::get_or_hydrate_from_uris() to ensure caching. Direct hydration bypasses the cache." + )] pub async fn hydrate_starterpacks_basic( &self, packs: Vec, @@ -205,6 +209,10 @@ impl StatefulHydrator<'_> { .collect() } + #[deprecated( + since = "0.1.0", + note = "Use starterpack-specific cache or hydrate_starterpack_basic. StarterPackView (full) and StarterPackViewBasic are different types." + )] pub async fn hydrate_starterpack(&self, pack: String) -> Option { let labels = self.get_label(&pack).await; diff --git a/parakeet/src/lib.rs b/parakeet/src/lib.rs index b39994c2..6f619bcf 100644 --- a/parakeet/src/lib.rs +++ b/parakeet/src/lib.rs @@ -11,6 +11,7 @@ pub mod cache; pub mod cache_listener; pub mod config; pub mod db; +pub mod entity_cache; pub mod hydration; pub mod id_cache_helpers; pub mod loaders; @@ -33,5 +34,11 @@ pub struct GlobalState { pub rate_limit_config: config::ConfigRateLimit, pub timeline_cache: Arc, pub author_feed_cache: Arc, + pub profile_cache: Arc, + pub post_cache: Arc, + pub feedgen_cache: Arc, + pub list_cache: Arc, + pub starterpack_cache: Arc, + pub labeler_cache: Arc, pub http_client: reqwest::Client, } diff --git a/parakeet/src/main.rs b/parakeet/src/main.rs index aa8835d4..3d8418be 100644 --- a/parakeet/src/main.rs +++ b/parakeet/src/main.rs @@ -110,6 +110,14 @@ async fn main() -> eyre::Result<()> { // Initialize author feed cache (60 second TTL, 10k max items) let author_feed_cache = Arc::new(timeline_cache::AuthorFeedCache::new(60, 10_000)); + // Initialize entity caches (60 second TTL, varying capacities) + let profile_cache = Arc::new(entity_cache::ProfileCache::new(60, 5_000)); + let post_cache = Arc::new(entity_cache::PostCache::new(60, 10_000)); + let feedgen_cache = Arc::new(entity_cache::FeedgenCache::new(60, 1_000)); + let list_cache = Arc::new(entity_cache::ListCache::new(60, 2_000)); + let starterpack_cache = Arc::new(entity_cache::StarterpackCache::new(60, 1_000)); + let labeler_cache = Arc::new(entity_cache::LabelerCache::new(60, 500)); + GlobalState { pool, dataloaders, @@ -122,6 +130,12 @@ async fn main() -> eyre::Result<()> { rate_limit_config: conf.rate_limit.clone(), timeline_cache, author_feed_cache, + profile_cache, + post_cache, + feedgen_cache, + list_cache, + starterpack_cache, + labeler_cache, http_client, } }; diff --git a/parakeet/src/unified_cache.rs b/parakeet/src/unified_cache.rs new file mode 100644 index 00000000..d0370776 --- /dev/null +++ b/parakeet/src/unified_cache.rs @@ -0,0 +1,198 @@ +//! Unified caching system that owns hydration +//! +//! This module provides the ONLY way to get hydrated data in the application. +//! All hydration must go through the cache to ensure proper caching behavior. + +use moka::future::Cache; +use std::sync::Arc; +use std::time::Duration; + +use diesel_async::pooled_connection::deadpool::Pool; +use diesel_async::AsyncPgConnection; + +use lexica::app_bsky::actor::ProfileViewDetailed; +use lexica::app_bsky::feed::PostView; + +use crate::hydration::StatefulHydrator; +use crate::loaders::Dataloaders; +use crate::xrpc::cdn::BskyCdn; +use crate::id_cache_helpers; +use parakeet_db::id_cache::IdCache; + +/// The unified cache system that owns all hydration +/// +/// This is the ONLY way to get hydrated data in the application. +/// Direct access to hydration is not allowed to ensure caching is always used. +#[derive(Clone)] +pub struct UnifiedCache { + profile_cache: Cache, + post_cache: Cache<(i32, i64), PostView>, + + // Dependencies needed for hydration + dataloaders: Arc, + cdn: Arc, + pool: Pool, + id_cache: Arc, +} + +impl UnifiedCache { + pub fn new( + dataloaders: Arc, + cdn: Arc, + pool: Pool, + id_cache: Arc, + profile_ttl: u64, + post_ttl: u64, + ) -> Self { + let profile_cache = Cache::builder() + .max_capacity(5_000) + .time_to_live(Duration::from_secs(profile_ttl)) + .support_invalidation_closures() + .build(); + + let post_cache = Cache::builder() + .max_capacity(10_000) + .time_to_live(Duration::from_secs(post_ttl)) + .support_invalidation_closures() + .build(); + + Self { + profile_cache, + post_cache, + dataloaders, + cdn, + pool, + id_cache, + } + } + + /// Get a profile - the ONLY way to get profile data + /// + /// This method handles caching and hydration internally. + /// No direct access to hydration is allowed. + pub async fn get_profile( + &self, + did: String, + labelers: &[String], + viewer_did: Option, + ) -> Option { + // First, get the actor_id for caching + let actor_id = id_cache_helpers::get_actor_id_or_fetch(&self.pool, &self.id_cache, &did).await.ok()?; + + // Check cache + if let Some(profile) = self.profile_cache.get(&actor_id).await { + tracing::debug!(actor_id, "Profile cache hit"); + return Some(profile); + } + + // Cache miss - hydrate using internal hydrator + tracing::debug!(actor_id, "Profile cache miss, hydrating"); + + // Get viewer actor_id if provided + let viewer_actor_id = if let Some(ref viewer) = viewer_did { + id_cache_helpers::get_actor_id_or_fetch(&self.pool, &self.id_cache, viewer).await.ok() + } else { + None + }; + + // Create hydrator for this request + let hydrator = StatefulHydrator::new( + &self.dataloaders, + &self.cdn, + labelers, + viewer_did, + viewer_actor_id, + ).await; + + // Hydrate the profile + let profile = hydrator.hydrate_profile_detailed(did).await?; + + // Store in cache + self.profile_cache.insert(actor_id, profile.clone()).await; + tracing::debug!(actor_id, "Profile cached"); + + Some(profile) + } + + /// Get multiple profiles + /// + /// For now, this doesn't use individual caching but could be optimized + pub async fn get_profiles( + &self, + dids: Vec, + labelers: &[String], + viewer_did: Option, + ) -> Vec { + // Get viewer actor_id if provided + let viewer_actor_id = if let Some(ref viewer) = viewer_did { + id_cache_helpers::get_actor_id_or_fetch(&self.pool, &self.id_cache, viewer).await.ok() + } else { + None + }; + + // Create hydrator for this request + let hydrator = StatefulHydrator::new( + &self.dataloaders, + &self.cdn, + labelers, + viewer_did, + viewer_actor_id, + ).await; + + // TODO: Use individual caching per profile + hydrator.hydrate_profiles_detailed(dids) + .await + .into_values() + .collect() + } + + /// Get posts from AT-URIs + pub async fn get_posts_from_uris( + &self, + uris: Vec, + labelers: &[String], + viewer_did: Option, + ) -> Vec { + // This would use the same pattern as profiles + // For now, keeping it simple + + let viewer_actor_id = if let Some(ref viewer) = viewer_did { + id_cache_helpers::get_actor_id_or_fetch(&self.pool, &self.id_cache, viewer).await.ok() + } else { + None + }; + + let hydrator = StatefulHydrator::new( + &self.dataloaders, + &self.cdn, + labelers, + viewer_did, + viewer_actor_id, + ).await; + + // TODO: Parse URIs and use individual post caching + hydrator.hydrate_posts(uris) + .await + .into_values() + .collect() + } + + /// Invalidate a profile by actor_id + pub async fn invalidate_profile(&self, actor_id: i32) { + self.profile_cache.invalidate(&actor_id).await; + tracing::debug!(actor_id, "Profile cache invalidated"); + } + + /// Invalidate a post by actor_id and rkey + pub async fn invalidate_post(&self, actor_id: i32, rkey: i64) { + self.post_cache.invalidate(&(actor_id, rkey)).await; + tracing::debug!(actor_id, rkey, "Post cache invalidated"); + } + + /// Clear all caches (for admin/testing) + pub async fn clear_all(&self) { + self.profile_cache.invalidate_all(); + self.post_cache.invalidate_all(); + tracing::info!("All caches cleared"); + } +} \ No newline at end of file diff --git a/parakeet/src/xrpc/app_bsky/actor.rs b/parakeet/src/xrpc/app_bsky/actor.rs index e92ee09b..771875eb 100644 --- a/parakeet/src/xrpc/app_bsky/actor.rs +++ b/parakeet/src/xrpc/app_bsky/actor.rs @@ -46,9 +46,12 @@ pub async fn get_profile( let mut conn = state.pool.get().await?; check_actor_status(&state.pool, &state.id_cache, &did).await?; - // Hydrate the profile from our data - let profile = hyd - .hydrate_profile_detailed(did) + // Get actor_id for cache key + let actor_id = crate::id_cache_helpers::get_actor_id_or_fetch(&state.pool, &state.id_cache, &did).await?; + + // Use the ergonomic cache API - it handles both caching and hydration + let profile = state.profile_cache + .get_or_hydrate(actor_id, did, &hyd) .await .ok_or_else(Error::not_found)?; @@ -85,11 +88,10 @@ pub async fn get_profiles( let dids = get_actor_dids(&state.dataloaders, query.actors).await; - let profiles = hyd - .hydrate_profiles_detailed(dids) - .await - .into_values() - .collect(); + // Use the cache's batch method instead of direct hydration + let profiles = state.profile_cache + .get_or_hydrate_batch(dids, &state.pool, &state.id_cache, &hyd) + .await; Ok(Json(GetProfilesRes { profiles }).into_response()) } diff --git a/parakeet/src/xrpc/app_bsky/bookmark.rs b/parakeet/src/xrpc/app_bsky/bookmark.rs index a45367ac..04ed9027 100644 --- a/parakeet/src/xrpc/app_bsky/bookmark.rs +++ b/parakeet/src/xrpc/app_bsky/bookmark.rs @@ -214,7 +214,13 @@ pub async fn get_bookmarks( }) .collect(); - let mut posts = hyd.hydrate_posts(uris).await; + // Use cache for post hydration (returns HashMap directly) + let mut posts = state.post_cache.get_or_hydrate_from_uris( + uris, + &state.pool, + &state.id_cache, + &hyd, + ).await; let bookmarks = results .into_iter() diff --git a/parakeet/src/xrpc/app_bsky/feed/feedgen.rs b/parakeet/src/xrpc/app_bsky/feed/feedgen.rs index a0afa231..1cc5d66c 100644 --- a/parakeet/src/xrpc/app_bsky/feed/feedgen.rs +++ b/parakeet/src/xrpc/app_bsky/feed/feedgen.rs @@ -75,14 +75,20 @@ pub async fn get_actor_feeds( }) .collect(); - let mut feeds = hyd.hydrate_feedgens(at_uris).await; + // Use cache for feedgen hydration (returns HashMap directly) + let mut feeds_map = state.feedgen_cache.get_or_hydrate_from_uris( + at_uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let feeds = results .into_iter() .filter_map(|r| { let did = actor_id_to_did.get(&r.1)?; let at_uri = format!("at://{}/app.bsky.feed.generator/{}", did, r.2); - feeds.remove(&at_uri) + feeds_map.remove(&at_uri) }) .collect(); @@ -117,7 +123,33 @@ pub async fn get_feed_generator( }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let Some(view) = hyd.hydrate_feedgen(query.feed).await else { + // Parse the feed URI to extract actor_id and rkey for caching + let view = if let Some(parsed) = crate::entity_cache::parse_at_uri(&query.feed) { + if parsed.collection == "app.bsky.feed.generator" { + // Get actor_id from DID + if let Ok(actor_id) = crate::id_cache_helpers::get_actor_id_or_fetch( + &state.pool, + &state.id_cache, + &parsed.did + ).await { + // Use cache for single feedgen + state.feedgen_cache.get_or_hydrate_single( + actor_id, + parsed.rkey, + query.feed.clone(), + &hyd, + ).await + } else { + None + } + } else { + None + } + } else { + None + }; + + let Some(view) = view else { return Err(Error::not_found()); }; @@ -154,11 +186,16 @@ pub async fn get_feed_generators( }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let feeds = hyd - .hydrate_feedgens(query.feeds) - .await - .into_values() - .collect(); + // Use cache for batch feedgen hydration (returns HashMap) + let feeds_map = state.feedgen_cache.get_or_hydrate_from_uris( + query.feeds, + &state.pool, + &state.id_cache, + &hyd, + ).await; + + // Convert to Vec for API response + let feeds = feeds_map.into_values().collect(); Ok(Json(GetFeedGeneratorsRes { feeds })) } @@ -251,7 +288,14 @@ pub async fn get_suggested_feeds( (None, None) }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let mut feeds_map = hyd.hydrate_feedgens(page_uris.clone()).await; + + // Use cache for batch feedgen hydration (returns HashMap directly) + let mut feeds_map = state.feedgen_cache.get_or_hydrate_from_uris( + page_uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let feeds: Vec = page_uris .into_iter() diff --git a/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs b/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs index 7d9dae9e..3805105e 100644 --- a/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs +++ b/parakeet/src/xrpc/app_bsky/feed/get_timeline.rs @@ -5,7 +5,7 @@ use crate::xrpc::extract::{AtpAcceptLabelers, AtpAuth}; use crate::GlobalState; use axum::extract::{Query, State}; use axum::Json; -use lexica::app_bsky::feed::{FeedReasonRepost, FeedViewPost, FeedViewPostReason}; +use lexica::app_bsky::feed::{FeedReasonRepost, FeedViewPost, FeedViewPostReason, PostView}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; @@ -189,10 +189,15 @@ pub async fn get_timeline( // OPTIMIZATION: Use the original (actor_id, rkey) results for repost query instead of re-parsing URIs let post_keys: Vec<(i32, i64)> = results.iter().map(|(_, actor_id, rkey)| (*actor_id, *rkey)).collect(); - // Parallelize: hydrate posts and fetch repost data concurrently + // Parallelize: hydrate posts (using cache) and fetch repost data concurrently step_timer = std::time::Instant::now(); let (mut post_views, reposts_results) = tokio::join!( - hyd.hydrate_posts(at_uris.clone()), + state.post_cache.get_or_hydrate_from_uris( + at_uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ), async { crate::db::get_timeline_reposts(&mut conn, &followed_actor_ids, &post_keys) .await @@ -321,7 +326,12 @@ async fn hydrate_timeline_feed( tracing::error!("Failed to get DB connection for repost hydration: {}", e); // Return posts without repost info by hydrating them let uri_count = at_uris.len(); - let mut post_views = hyd.hydrate_posts(at_uris.clone()).await; + let mut post_views = state.post_cache.get_or_hydrate_from_uris( + at_uris.clone(), + &state.pool, + &state.id_cache, + hyd, + ).await; let feed: Vec = at_uris .into_iter() .filter_map(|uri| { @@ -346,7 +356,12 @@ async fn hydrate_timeline_feed( // Parallelize: hydrate posts and get followed actor_ids concurrently let (mut post_views, follows) = tokio::join!( - hyd.hydrate_posts(at_uris.clone()), + state.post_cache.get_or_hydrate_from_uris( + at_uris.clone(), + &state.pool, + &state.id_cache, + hyd, + ), async { crate::db::get_followed_dids(&mut conn, user_actor_id) .await diff --git a/parakeet/src/xrpc/app_bsky/feed/posts/queries.rs b/parakeet/src/xrpc/app_bsky/feed/posts/queries.rs index 394c64ec..79846e9b 100644 --- a/parakeet/src/xrpc/app_bsky/feed/posts/queries.rs +++ b/parakeet/src/xrpc/app_bsky/feed/posts/queries.rs @@ -35,11 +35,19 @@ pub async fn get_posts( (None, None) }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let posts = hyd.hydrate_posts(query.uris).await; - Ok(Json(PostsRes { - posts: posts.into_values().collect(), - })) + // Use the new method that parses URIs and caches individual posts + let posts_map = state.post_cache.get_or_hydrate_from_uris( + query.uris, + &state.pool, + &state.id_cache, + &hyd, + ).await; + + // Convert to Vec for API response + let posts = posts_map.into_values().collect(); + + Ok(Json(PostsRes { posts })) } #[derive(Debug, Deserialize)] @@ -166,7 +174,13 @@ pub async fn get_quotes( let cursor = uris.last().cloned(); - let mut posts_map = hyd.hydrate_posts(uris.clone()).await; + // Use cache for post hydration (returns HashMap directly) + let mut posts_map = state.post_cache.get_or_hydrate_from_uris( + uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let posts = uris .into_iter() diff --git a/parakeet/src/xrpc/app_bsky/feed/posts/threads.rs b/parakeet/src/xrpc/app_bsky/feed/posts/threads.rs index 71225d4a..86de9996 100644 --- a/parakeet/src/xrpc/app_bsky/feed/posts/threads.rs +++ b/parakeet/src/xrpc/app_bsky/feed/posts/threads.rs @@ -1,6 +1,6 @@ use axum::extract::{Query, State}; use axum::Json; -use lexica::app_bsky::feed::{BlockedAuthor, ThreadViewPost, ThreadViewPostType, ThreadgateView}; +use lexica::app_bsky::feed::{BlockedAuthor, PostView, ThreadViewPost, ThreadViewPostType, ThreadgateView}; use serde::{Deserialize, Serialize}; use std::collections::HashMap; @@ -142,10 +142,20 @@ pub async fn get_post_thread( .map(|item: &ThreadItem| item.at_uri.clone()) .collect(); - // Parallelize: hydrate replies and parents concurrently + // Parallelize: hydrate replies and parents concurrently using cache (returns HashMap directly) let (mut replies_hydrated, mut parents_hydrated) = tokio::join!( - hyd.hydrate_posts(reply_uris), - hyd.hydrate_posts(parent_uris) + state.post_cache.get_or_hydrate_from_uris( + reply_uris, + &state.pool, + &state.id_cache, + &hyd, + ), + state.post_cache.get_or_hydrate_from_uris( + parent_uris, + &state.pool, + &state.id_cache, + &hyd, + ) ); let mut tmpbuf: HashMap<_, Vec<_>> = HashMap::new(); diff --git a/parakeet/src/xrpc/app_bsky/feed/search.rs b/parakeet/src/xrpc/app_bsky/feed/search.rs index 0a109e80..9160f0d2 100644 --- a/parakeet/src/xrpc/app_bsky/feed/search.rs +++ b/parakeet/src/xrpc/app_bsky/feed/search.rs @@ -192,12 +192,19 @@ pub async fn search_posts( (None, None) }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let mut posts_map = hyd.hydrate_posts(uris.clone()).await; + + // Use cache for post hydration (returns HashMap directly) + let posts_map = state.post_cache.get_or_hydrate_from_uris( + uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; // Maintain search result order let posts: Vec = uris .into_iter() - .filter_map(|uri| posts_map.remove(&uri)) + .filter_map(|uri| posts_map.get(&uri).cloned()) .collect(); // Calculate cursor (rank of last result) diff --git a/parakeet/src/xrpc/app_bsky/graph/lists.rs b/parakeet/src/xrpc/app_bsky/graph/lists.rs index 858fae91..76bfc978 100644 --- a/parakeet/src/xrpc/app_bsky/graph/lists.rs +++ b/parakeet/src/xrpc/app_bsky/graph/lists.rs @@ -74,13 +74,19 @@ pub async fn get_lists( .map(|r| format!("at://{}/app.bsky.graph.list/{}", did, r.1)) .collect(); - let mut lists = hyd.hydrate_lists(at_uris).await; + // Use cache for list hydration (returns HashMap directly) + let mut lists_map = state.list_cache.get_or_hydrate_from_uris( + at_uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let lists = results .into_iter() .filter_map(|r| { let at_uri = format!("at://{}/app.bsky.graph.list/{}", did, r.1); - lists.remove(&at_uri) + lists_map.remove(&at_uri) }) .collect(); @@ -111,7 +117,33 @@ pub async fn get_list( }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let Some(list) = hyd.hydrate_list(query.list).await else { + // Parse the list URI to extract actor_id and rkey for caching + let list = if let Some(parsed) = crate::entity_cache::parse_at_uri(&query.list) { + if parsed.collection == "app.bsky.graph.list" { + // Get actor_id from DID + if let Ok(actor_id) = crate::id_cache_helpers::get_actor_id_or_fetch( + &state.pool, + &state.id_cache, + &parsed.did + ).await { + // Use cache for single list + state.list_cache.get_or_hydrate_single( + actor_id, + parsed.rkey.clone(), + query.list.clone(), + &hyd, + ).await + } else { + None + } + } else { + None + } + } else { + None + }; + + let Some(list) = list else { return Err(Error::not_found()); }; @@ -187,10 +219,18 @@ pub async fn get_list_mutes( .last() .map(|last| last.0.timestamp_millis().to_string()); - let uris = results.iter().map(|r| r.1.clone()).collect(); + let uris: Vec = results.iter().map(|r| r.1.clone()).collect(); - let lists = hyd.hydrate_lists(uris).await; - let lists = lists.into_values().collect::>(); + // Use cache for list hydration (returns HashMap) + let lists_map = state.list_cache.get_or_hydrate_from_uris( + uris, + &state.pool, + &state.id_cache, + &hyd, + ).await; + + // Convert to Vec for API response + let lists = lists_map.into_values().collect(); Ok(Json(GetListsRes { cursor, lists })) } @@ -216,10 +256,18 @@ pub async fn get_list_blocks( .last() .map(|last| last.0.timestamp_millis().to_string()); - let uris = results.iter().map(|r| r.1.clone()).collect(); + let uris: Vec = results.iter().map(|r| r.1.clone()).collect(); + + // Use cache for list hydration (returns HashMap) + let lists_map = state.list_cache.get_or_hydrate_from_uris( + uris, + &state.pool, + &state.id_cache, + &hyd, + ).await; - let lists = hyd.hydrate_lists(uris).await; - let lists = lists.into_values().collect::>(); + // Convert to Vec for API response + let lists = lists_map.into_values().collect(); Ok(Json(GetListsRes { cursor, lists })) } diff --git a/parakeet/src/xrpc/app_bsky/graph/search.rs b/parakeet/src/xrpc/app_bsky/graph/search.rs index 9f514921..7e6bbecf 100644 --- a/parakeet/src/xrpc/app_bsky/graph/search.rs +++ b/parakeet/src/xrpc/app_bsky/graph/search.rs @@ -108,7 +108,14 @@ pub async fn search_starter_packs( (None, None) }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let mut packs_map = hyd.hydrate_starterpacks_basic(uris.clone()).await; + + // Use cache for starterpack hydration (returns HashMap directly) + let mut packs_map = state.starterpack_cache.get_or_hydrate_from_uris( + uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; // Maintain search result order let starter_packs: Vec = uris diff --git a/parakeet/src/xrpc/app_bsky/graph/starter_packs.rs b/parakeet/src/xrpc/app_bsky/graph/starter_packs.rs index fdd3ec51..75ff66cd 100644 --- a/parakeet/src/xrpc/app_bsky/graph/starter_packs.rs +++ b/parakeet/src/xrpc/app_bsky/graph/starter_packs.rs @@ -72,11 +72,17 @@ pub async fn get_actor_starter_packs( }) .collect(); - let mut starter_packs = hyd.hydrate_starterpacks_basic(uris.clone()).await; + // Use cache for starterpack hydration (returns HashMap directly) + let mut starter_packs_map = state.starterpack_cache.get_or_hydrate_from_uris( + uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let starter_packs = uris .into_iter() - .filter_map(|uri| starter_packs.remove(&uri)) + .filter_map(|uri| starter_packs_map.remove(&uri)) .collect(); Ok(Json(StarterPacksRes { @@ -112,7 +118,31 @@ pub async fn get_starter_pack( }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let Some(starter_pack) = hyd.hydrate_starterpack(query.starter_pack).await else { + // Parse the starterpack URI to extract actor_id and rkey for caching + let starter_pack = if let Some(parsed) = crate::entity_cache::parse_at_uri(&query.starter_pack) { + if parsed.collection == "app.bsky.graph.starterpack" { + // Get actor_id from DID + if let Ok(actor_id) = crate::id_cache_helpers::get_actor_id_or_fetch( + &state.pool, + &state.id_cache, + &parsed.did + ).await { + // Use cache for single starterpack, but it returns StarterPackViewBasic + // We need StarterPackView, so we still use hydrate_starterpack directly + // This is a limitation - the cache stores Basic views but this needs full View + #[allow(deprecated)] + hyd.hydrate_starterpack(query.starter_pack.clone()).await + } else { + None + } + } else { + None + } + } else { + None + }; + + let Some(starter_pack) = starter_pack else { return Err(Error::not_found()); }; @@ -139,11 +169,16 @@ pub async fn get_starter_packs( }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let starter_packs = hyd - .hydrate_starterpacks_basic(query.uris) - .await - .into_values() - .collect(); + // Use cache for batch starterpack hydration (returns HashMap) + let packs_map = state.starterpack_cache.get_or_hydrate_from_uris( + query.uris, + &state.pool, + &state.id_cache, + &hyd, + ).await; + + // Convert to Vec for API response + let starter_packs = packs_map.into_values().collect(); Ok(Json(StarterPacksRes { starter_packs, diff --git a/parakeet/src/xrpc/app_bsky/notification/mod.rs b/parakeet/src/xrpc/app_bsky/notification/mod.rs index 789b7734..4a89f09a 100644 --- a/parakeet/src/xrpc/app_bsky/notification/mod.rs +++ b/parakeet/src/xrpc/app_bsky/notification/mod.rs @@ -226,7 +226,19 @@ pub async fn list_notifications( Some(did), Some(actor_id), ).await; - let profiles_map = hyd.hydrate_profiles_detailed(author_dids.clone()).await; + // Use cache for profile hydration + let profiles_vec = state.profile_cache.get_or_hydrate_batch( + author_dids.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; + + // Convert Vec to HashMap for compatibility with existing code + let profiles_map: std::collections::HashMap = + profiles_vec.into_iter() + .map(|profile| (profile.did.clone(), profile)) + .collect(); tracing::info!(" → Hydrate profiles: {:.1}ms ({} profiles)", profile_start.elapsed().as_secs_f64() * 1000.0, profiles_map.len()); // Build notification responses diff --git a/parakeet/src/xrpc/app_bsky/unspecced/mod.rs b/parakeet/src/xrpc/app_bsky/unspecced/mod.rs index 21b02d0e..519dbd22 100644 --- a/parakeet/src/xrpc/app_bsky/unspecced/mod.rs +++ b/parakeet/src/xrpc/app_bsky/unspecced/mod.rs @@ -196,7 +196,14 @@ pub async fn get_suggested_feeds( (None, None) }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let mut feeds_map = hyd.hydrate_feedgens(page_uris.clone()).await; + + // Use cache for feedgen hydration (returns HashMap directly) + let mut feeds_map = state.feedgen_cache.get_or_hydrate_from_uris( + page_uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let feeds: Vec = page_uris .into_iter() @@ -423,7 +430,14 @@ pub async fn get_popular_feed_generators( (None, None) }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let mut feeds_map = hyd.hydrate_feedgens(page_uris.clone()).await; + + // Use cache for feedgen hydration (returns HashMap directly) + let mut feeds_map = state.feedgen_cache.get_or_hydrate_from_uris( + page_uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let feeds: Vec = page_uris .into_iter() @@ -523,7 +537,14 @@ pub async fn get_suggested_starter_packs( (None, None) }; let hyd = StatefulHydrator::new(&state.dataloaders, &state.cdn, &labelers, maybe_did, maybe_actor_id).await; - let mut packs_map = hyd.hydrate_starterpacks_basic(page_uris.clone()).await; + + // Use cache for starterpack hydration (returns HashMap directly) + let mut packs_map = state.starterpack_cache.get_or_hydrate_from_uris( + page_uris.clone(), + &state.pool, + &state.id_cache, + &hyd, + ).await; let starter_packs: Vec = page_uris .into_iter() diff --git a/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/other_replies.rs b/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/other_replies.rs index 062cb97c..d475a1ca 100644 --- a/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/other_replies.rs +++ b/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/other_replies.rs @@ -96,7 +96,12 @@ pub async fn get_post_thread_other_v2( &state.id_cache ).await?; let reply_uris: Vec = replies.iter().map(|item| item.at_uri.clone()).collect(); - let replies_hydrated = hyd.hydrate_posts(reply_uris).await; + let replies_hydrated = state.post_cache.get_or_hydrate_from_uris( + reply_uris, + &state.pool, + &state.id_cache, + &hyd, + ).await; // Build a map of parent_uri -> [child_posts] let mut replies_by_parent: HashMap> = HashMap::new(); diff --git a/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/post_thread.rs b/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/post_thread.rs index a58e9a44..9cb8ffa0 100644 --- a/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/post_thread.rs +++ b/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/post_thread.rs @@ -57,6 +57,7 @@ pub async fn get_post_thread_v2( // Create a thread builder let builder = ThreadBuilder { hydrater: &hyd, + post_cache: &state.post_cache, anchor_uri: uri, anchor_post, threadgate: threadgate.clone(), @@ -67,6 +68,7 @@ pub async fn get_post_thread_v2( prioritize_followed_users: query.prioritize_followed_users, is_authenticated, id_cache: &state.id_cache, + pool: &state.pool, }; // Build the thread diff --git a/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/thread_builder.rs b/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/thread_builder.rs index 7e170d21..d5280790 100644 --- a/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/thread_builder.rs +++ b/parakeet/src/xrpc/app_bsky/unspecced/thread_v2/thread_builder.rs @@ -1,6 +1,11 @@ +use crate::entity_cache::PostCache; use crate::hydration::StatefulHydrator; +use diesel_async::pooled_connection::deadpool::Pool; +use diesel_async::AsyncPgConnection; use lexica::app_bsky::feed::{PostView, ThreadgateView}; +use parakeet_db::id_cache::IdCache; use std::collections::HashMap; +use std::sync::Arc; use super::models::{PostThreadSort, ThreadItemPost, ThreadV2Item, ThreadV2ItemType}; use super::sorting::sort_replies; @@ -9,6 +14,7 @@ use super::sorting::sort_replies; #[expect(dead_code, reason = "threadgate and prioritize_followed_users are infrastructure for future thread filtering features")] pub struct ThreadBuilder<'a> { pub hydrater: &'a StatefulHydrator<'a>, + pub post_cache: &'a Arc, pub anchor_uri: String, pub anchor_post: PostView, pub threadgate: Option, @@ -18,7 +24,8 @@ pub struct ThreadBuilder<'a> { pub sort: PostThreadSort, pub prioritize_followed_users: bool, pub is_authenticated: bool, - pub id_cache: &'a parakeet_db::id_cache::IdCache, + pub id_cache: &'a Arc, + pub pool: &'a Pool, } impl ThreadBuilder<'_> { @@ -175,7 +182,12 @@ impl ThreadBuilder<'_> { let hydrate_start = std::time::Instant::now(); let parent_uris: Vec = parents.iter().map(|item| item.at_uri.clone()).collect(); - let parents_hydrated = self.hydrater.hydrate_posts(parent_uris).await; + let parents_hydrated = self.post_cache.get_or_hydrate_from_uris( + parent_uris, + self.pool, + self.id_cache, + self.hydrater, + ).await; tracing::info!(" → Hydrate parent posts: {:.1} ms ({} posts)", hydrate_start.elapsed().as_secs_f64() * 1000.0, parents_hydrated.len()); // Check if anchor post has a parent that wasn't returned by get_thread_parents @@ -268,7 +280,12 @@ impl ThreadBuilder<'_> { // Hydrate all the posts let hydrate_start = std::time::Instant::now(); let reply_uris: Vec = replies.iter().map(|item| item.at_uri.clone()).collect(); - let replies_hydrated = self.hydrater.hydrate_posts(reply_uris).await; + let replies_hydrated = self.post_cache.get_or_hydrate_from_uris( + reply_uris, + self.pool, + self.id_cache, + self.hydrater, + ).await; tracing::info!(" → Hydrate posts: {:.1} ms ({} posts)", hydrate_start.elapsed().as_secs_f64() * 1000.0, replies_hydrated.len()); // Build a map: parent_uri -> [(child_uri, sql_depth)]