diff --git a/parakeet/src/db.rs b/parakeet/src/db.rs index 870ce418..ff10e73a 100644 --- a/parakeet/src/db.rs +++ b/parakeet/src/db.rs @@ -32,7 +32,7 @@ pub mod threads; pub mod uri_reconstruction; // Re-export commonly used functions -pub use actors::{get_actor_ids_by_dids, get_actor_status, resolve_actor, ResolvedActor}; +pub use actors::{get_actor_data_by_ids, get_actor_ids_by_dids, get_actor_status, resolve_actor, ResolvedActor}; pub use bookmarks::get_user_bookmarks; pub use feeds::{get_author_feed, get_list_feed, get_quotes, get_reposted_by, get_timeline_posts, get_timeline_reposts, AuthorFeedFilter, AuthorFeedItem}; pub use feedgens::{get_actor_feedgens, get_all_feedgen_uris, get_all_feedgens_by_likes, get_feedgen_service_did}; diff --git a/parakeet/src/db/actors.rs b/parakeet/src/db/actors.rs index 1ddedde6..37ac65ea 100644 --- a/parakeet/src/db/actors.rs +++ b/parakeet/src/db/actors.rs @@ -248,3 +248,68 @@ pub async fn get_actor_ids_by_dids( Ok(result) } + +/// Batch resolve actor IDs to actor data (DID and handle) with cache support +/// +/// This function handles cache misses by querying the database and populating the cache. +pub async fn get_actor_data_by_ids( + conn: &mut AsyncPgConnection, + actor_ids: &[i32], + cache: ¶keet_db::id_cache::IdCache, +) -> QueryResult> { + use diesel::sql_types::{Array, Integer, Nullable, Text}; + + if actor_ids.is_empty() { + return Ok(std::collections::HashMap::new()); + } + + // Get cached data + let mut result = cache.get_actor_data_many(actor_ids).await; + + // If all found in cache, return early + if result.len() == actor_ids.len() { + return Ok(result); + } + + // Find missing IDs + let missing_ids: Vec = actor_ids + .iter() + .filter(|id| !result.contains_key(*id)) + .copied() + .collect(); + + if missing_ids.is_empty() { + return Ok(result); + } + + // Query database for missing actors + #[derive(QueryableByName)] + struct ActorDataRow { + #[diesel(sql_type = Integer)] + id: i32, + #[diesel(sql_type = Text)] + did: String, + #[diesel(sql_type = Nullable)] + handle: Option, + } + + let missing_actors: Vec = diesel::sql_query( + "SELECT id, did, handle FROM actors WHERE id = ANY($1)" + ) + .bind::, _>(&missing_ids) + .load(conn) + .await?; + + // Add to result and populate cache + for actor in missing_actors { + let data = parakeet_db::id_cache::CachedActorData { + did: actor.did.clone(), + handle: actor.handle.clone(), + }; + result.insert(actor.id, data.clone()); + // Populate cache for future requests + cache.set_actor_data(actor.id, data).await; + } + + Ok(result) +} diff --git a/parakeet/src/hydration/posts/feed.rs b/parakeet/src/hydration/posts/feed.rs index 32a84c30..1dc511fe 100644 --- a/parakeet/src/hydration/posts/feed.rs +++ b/parakeet/src/hydration/posts/feed.rs @@ -71,10 +71,19 @@ impl StatefulHydrator<'_> { } } - // Single batched async call to get all actor DIDs + // Single batched async call to get all actor DIDs (with database fallback for cache misses) let actor_ids: Vec = actor_ids_needed.into_iter().collect(); let id_cache = self.loaders.post_state.id_cache(); - let actor_data = id_cache.get_actor_data_many(&actor_ids).await; + let actor_data = if let Ok(mut conn) = self.loaders.post_state.get_conn().await { + crate::db::get_actor_data_by_ids(&mut conn, &actor_ids, id_cache) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve actor data in feed hydration: {e}"); + std::collections::HashMap::new() + }) + } else { + id_cache.get_actor_data_many(&actor_ids).await + }; // Build actor_id -> DID mapping for fast lookups let actor_id_to_did: HashMap = actor_data diff --git a/parakeet/src/hydration/posts/mod.rs b/parakeet/src/hydration/posts/mod.rs index 92f289a5..9e562f56 100644 --- a/parakeet/src/hydration/posts/mod.rs +++ b/parakeet/src/hydration/posts/mod.rs @@ -25,7 +25,19 @@ impl StatefulHydrator<'_> { } let labeler_ids: Vec = labels.iter().map(|l| l.labeler_actor_id).collect(); - let labeler_data = self.loaders.post_state.id_cache().get_actor_data_many(&labeler_ids).await; + + // Get connection and resolve with cache miss handling + let labeler_data = if let Ok(mut conn) = self.loaders.post_state.get_conn().await { + crate::db::get_actor_data_by_ids(&mut conn, &labeler_ids, self.loaders.post_state.id_cache()) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve label actor data: {e}"); + std::collections::HashMap::new() + }) + } else { + // Fallback to cache-only if can't get connection + self.loaders.post_state.id_cache().get_actor_data_many(&labeler_ids).await + }; labels.iter().filter_map(|record| { let labeler_did = labeler_data.get(&record.labeler_actor_id)?.did.clone(); @@ -163,7 +175,20 @@ impl StatefulHydrator<'_> { let reply_actor_ids = Self::collect_reply_actor_ids(&posts_with_stats); let reply_actor_cache = if !reply_actor_ids.is_empty() { let cache_start = std::time::Instant::now(); - let actor_data = self.loaders.post_state.id_cache().get_actor_data_many(&reply_actor_ids).await; + + // Get connection and resolve with cache miss handling + let actor_data = if let Ok(mut conn) = self.loaders.post_state.get_conn().await { + crate::db::get_actor_data_by_ids(&mut conn, &reply_actor_ids, self.loaders.post_state.id_cache()) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve reply actor data: {e}"); + std::collections::HashMap::new() + }) + } else { + // Fallback to cache-only if can't get connection + self.loaders.post_state.id_cache().get_actor_data_many(&reply_actor_ids).await + }; + let cache_time = cache_start.elapsed().as_secs_f64() * 1000.0; if cache_time > 5.0 { tracing::info!(" → Pre-fetch reply actor IDs: {:.1} ms ({} actors)", cache_time, reply_actor_ids.len()); diff --git a/parakeet/src/hydration/profile/mod.rs b/parakeet/src/hydration/profile/mod.rs index a721e8c9..4bb7a8a1 100644 --- a/parakeet/src/hydration/profile/mod.rs +++ b/parakeet/src/hydration/profile/mod.rs @@ -26,8 +26,18 @@ impl super::StatefulHydrator<'_> { // Collect unique labeler actor IDs let labeler_ids: Vec = labels.iter().map(|l| l.labeler_actor_id).collect(); - // Batch resolve labeler_actor_id → DID using IdCache - let labeler_data = self.loaders.profile_state.id_cache().get_actor_data_many(&labeler_ids).await; + // Batch resolve labeler_actor_id → DID with cache miss handling + let labeler_data = if let Ok(mut conn) = self.loaders.profile_state.get_conn().await { + crate::db::get_actor_data_by_ids(&mut conn, &labeler_ids, self.loaders.profile_state.id_cache()) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve profile label actor data: {e}"); + std::collections::HashMap::new() + }) + } else { + // Fallback to cache-only if can't get connection + self.loaders.profile_state.id_cache().get_actor_data_many(&labeler_ids).await + }; // Convert each label - use DID and profile URI as targets let mut result = Vec::new(); @@ -158,9 +168,18 @@ impl super::StatefulHydrator<'_> { }) .collect(); - // Batch resolve list owner actor_ids → DIDs via IdCache + // Batch resolve list owner actor_ids → DIDs via IdCache (with database fallback for cache misses) let list_owner_ids_vec: Vec = list_owner_ids.into_iter().collect(); - let id_to_actor_data = self.loaders.profile_state.id_cache().get_actor_data_many(&list_owner_ids_vec).await; + let id_to_actor_data = if let Ok(mut conn) = self.loaders.profile_state.get_conn().await { + crate::db::get_actor_data_by_ids(&mut conn, &list_owner_ids_vec, self.loaders.profile_state.id_cache()) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve list owner actor data: {e}"); + std::collections::HashMap::new() + }) + } else { + self.loaders.profile_state.id_cache().get_actor_data_many(&list_owner_ids_vec).await + }; // Build list URIs with resolved DIDs let mut list_uris = Vec::new(); diff --git a/parakeet/src/loaders/feed.rs b/parakeet/src/loaders/feed.rs index 5456a94f..ace6fe94 100644 --- a/parakeet/src/loaders/feed.rs +++ b/parakeet/src/loaders/feed.rs @@ -129,12 +129,13 @@ impl BatchFn for FeedGenLoader { .into_iter() .collect(); - // Batch resolve service DIDs via IdCache - let service_actor_data = if !service_actor_ids.is_empty() { - self.1.get_actor_data_many(&service_actor_ids).await - } else { - std::collections::HashMap::new() - }; + // Batch resolve service DIDs with automatic cache miss handling + let service_actor_data = crate::db::get_actor_data_by_ids(&mut conn, &service_actor_ids, &self.1) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve service actor data: {e}"); + std::collections::HashMap::new() + }); HashMap::from_iter(res.into_iter().map(|row| { // Get service DID from resolved data diff --git a/parakeet/src/loaders/labeler.rs b/parakeet/src/loaders/labeler.rs index 653fc824..40ccef22 100644 --- a/parakeet/src/loaders/labeler.rs +++ b/parakeet/src/loaders/labeler.rs @@ -170,10 +170,13 @@ impl LabelLoader { pub async fn load(&self, uri: &str, services: &[LabelConfigItem]) -> Vec { let mut conn = self.0.get().await.unwrap(); - // Parse URI to extract DID (at://did:plc:xxx/...) + // Parse URI to extract DID - handle both AT URIs (at://did:plc:xxx/...) and plain DIDs let subject_did = if uri.starts_with("at://") { uri.strip_prefix("at://") .and_then(|s| s.split('/').next()) + } else if uri.starts_with("did:") { + // Plain DID as the subject + Some(uri) } else { None }; @@ -181,7 +184,7 @@ impl LabelLoader { let subject_did = match subject_did { Some(did) => did, None => { - tracing::warn!("Invalid AT URI format for labels: {}", uri); + tracing::debug!("Invalid URI format for labels (expected at:// or did:): {}", uri); return Vec::new(); } }; @@ -257,7 +260,12 @@ impl LabelLoader { .into_iter() .collect(); - let labeler_data = self.1.get_actor_data_many(&unique_labeler_ids).await; + let labeler_data = crate::db::get_actor_data_by_ids(&mut conn, &unique_labeler_ids, &self.1) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve labeler actor data: {e}"); + std::collections::HashMap::new() + }); labels .into_iter() @@ -298,11 +306,21 @@ impl LabelLoader { let mut valid_uris = Vec::new(); let mut subject_dids = Vec::new(); for uri in uris { - if let Some(did) = uri.strip_prefix("at://").and_then(|s| s.split('/').next()) { + let did = if uri.starts_with("at://") { + // AT URI format: at://did:plc:xxx/... + uri.strip_prefix("at://").and_then(|s| s.split('/').next()) + } else if uri.starts_with("did:") { + // Plain DID as the subject + Some(uri.as_str()) + } else { + None + }; + + if let Some(did) = did { valid_uris.push(uri.as_str()); subject_dids.push(did.to_string()); } else { - tracing::warn!("Invalid AT URI format for labels: {}", uri); + tracing::debug!("Invalid URI format for labels (expected at:// or did:): {}", uri); } } @@ -384,7 +402,12 @@ impl LabelLoader { .into_iter() .collect(); - let labeler_data = self.1.get_actor_data_many(&unique_labeler_ids).await; + let labeler_data = crate::db::get_actor_data_by_ids(&mut conn, &unique_labeler_ids, &self.1) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve labeler actor data: {e}"); + std::collections::HashMap::new() + }); labels .into_iter() diff --git a/parakeet/src/loaders/misc.rs b/parakeet/src/loaders/misc.rs index d2cd929a..62fe9eb9 100644 --- a/parakeet/src/loaders/misc.rs +++ b/parakeet/src/loaders/misc.rs @@ -188,9 +188,14 @@ impl BatchFn for StarterPackLoader { .into_iter() .collect(); - // Batch resolve list owner DIDs via IdCache + // Batch resolve list owner DIDs with cache miss handling let list_owner_data = if !list_actor_ids.is_empty() { - self.1.get_actor_data_many(&list_actor_ids).await + crate::db::get_actor_data_by_ids(&mut conn, &list_actor_ids, &self.1) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve list owner actor data: {e}"); + std::collections::HashMap::new() + }) } else { std::collections::HashMap::new() }; @@ -354,8 +359,13 @@ impl BatchFn> for VerificationLoader { all_actor_ids.sort_unstable(); all_actor_ids.dedup(); - // Batch resolve actor_ids to DIDs via IdCache - let actor_data = self.1.get_actor_data_many(&all_actor_ids).await; + // Batch resolve actor_ids to DIDs with cache miss handling + let actor_data = crate::db::get_actor_data_by_ids(&mut conn, &all_actor_ids, &self.1) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve verification actor data: {e}"); + std::collections::HashMap::new() + }); // Group by subject DID let mut result: HashMap> = HashMap::new(); diff --git a/parakeet/src/loaders/post.rs b/parakeet/src/loaders/post.rs index c47ea0fe..57059916 100644 --- a/parakeet/src/loaders/post.rs +++ b/parakeet/src/loaders/post.rs @@ -978,10 +978,14 @@ impl BatchFn for PostLoader { mention_actor_ids.sort_unstable(); mention_actor_ids.dedup(); - // Batch lookup actor DIDs for mentions using id_cache + // Batch lookup actor DIDs for mentions with cache miss handling let mention_actor_lookup: HashMap = if !mention_actor_ids.is_empty() { - self.1.get_actor_data_many(&mention_actor_ids) + crate::db::get_actor_data_by_ids(&mut conn, &mention_actor_ids, &self.1) .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve mention actor data: {e}"); + std::collections::HashMap::new() + }) .into_iter() .map(|(actor_id, data)| (actor_id, data.did)) .collect() @@ -1004,10 +1008,14 @@ impl BatchFn for PostLoader { hidden_reply_actor_ids.sort_unstable(); hidden_reply_actor_ids.dedup(); - // Batch lookup actor DIDs for hidden replies using id_cache + // Batch lookup actor DIDs for hidden replies with cache miss handling let _hidden_reply_actor_lookup: HashMap = if !hidden_reply_actor_ids.is_empty() { - self.1.get_actor_data_many(&hidden_reply_actor_ids) + crate::db::get_actor_data_by_ids(&mut conn, &hidden_reply_actor_ids, &self.1) .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve hidden reply actor data: {e}"); + std::collections::HashMap::new() + }) .into_iter() .map(|(actor_id, data)| (actor_id, data.did)) .collect() @@ -1246,6 +1254,11 @@ impl PostStateLoader { &self.1 } + /// Get a database connection from the pool + pub async fn get_conn(&self) -> Result, diesel_async::pooled_connection::deadpool::PoolError> { + self.0.get().await + } + /// Load viewer cache by actor_id - single query to get bookmarks and pinned_post /// /// This is the ONLY place we should query the actors table for viewer data. diff --git a/parakeet/src/loaders/profile.rs b/parakeet/src/loaders/profile.rs index 10ff0969..3654e219 100644 --- a/parakeet/src/loaders/profile.rs +++ b/parakeet/src/loaders/profile.rs @@ -330,9 +330,14 @@ impl BatchFn for ProfileByIdLoader { embed_post_keys.sort_unstable(); embed_post_keys.dedup(); - // Batch resolve embed post actor DIDs + // Batch resolve embed post actor DIDs with cache miss handling let embed_actor_ids: Vec = embed_post_keys.iter().map(|(actor_id, _)| *actor_id).collect(); - let embed_actor_data = self.1.get_actor_data_many(&embed_actor_ids).await; + let embed_actor_data = crate::db::get_actor_data_by_ids(&mut conn, &embed_actor_ids, &self.1) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve embed actor data: {e}"); + std::collections::HashMap::new() + }); // Build map of (actor_id, rkey) -> embed_uri let embed_uri_map: std::collections::HashMap<(i32, i64), String> = embed_post_keys @@ -648,6 +653,11 @@ impl ProfileStateLoader { &self.0 } + /// Get a database connection from the pool + pub async fn get_conn(&self) -> Result, diesel_async::pooled_connection::deadpool::PoolError> { + self.0.get().await + } + /// Get single profile state (DID-based interface, internally uses actor_ids) pub async fn get(&self, viewer_did: &str, subject_did: &str) -> Option { let results = self.get_many(viewer_did, &vec![subject_did.to_string()]).await; diff --git a/parakeet/src/xrpc/app_bsky/feed/feedgen.rs b/parakeet/src/xrpc/app_bsky/feed/feedgen.rs index 1cc5d66c..5bd6b75d 100644 --- a/parakeet/src/xrpc/app_bsky/feed/feedgen.rs +++ b/parakeet/src/xrpc/app_bsky/feed/feedgen.rs @@ -266,9 +266,17 @@ pub async fn get_suggested_feeds( None }; - // Construct URIs from natural keys by resolving actor_ids to DIDs + // Construct URIs from natural keys by resolving actor_ids to DIDs (with database fallback for cache misses) let actor_ids: Vec = page_feedgens.iter().map(|k| k.0).collect(); - let actor_data = state.id_cache.get_actor_data_many(&actor_ids).await; + let actor_data = { + let mut conn = state.pool.get().await?; + crate::db::get_actor_data_by_ids(&mut conn, &actor_ids, &state.id_cache) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve actor data for feed generators: {e}"); + std::collections::HashMap::new() + }) + }; let page_uris: Vec = page_feedgens .iter() diff --git a/parakeet/src/xrpc/app_bsky/unspecced/mod.rs b/parakeet/src/xrpc/app_bsky/unspecced/mod.rs index 05c9ba36..c68fc365 100644 --- a/parakeet/src/xrpc/app_bsky/unspecced/mod.rs +++ b/parakeet/src/xrpc/app_bsky/unspecced/mod.rs @@ -174,9 +174,17 @@ pub async fn get_suggested_feeds( }) .collect(); - // Construct URIs from natural keys by resolving actor_ids to DIDs + // Construct URIs from natural keys by resolving actor_ids to DIDs (with database fallback for cache misses) let actor_ids: Vec = top_feedgens.iter().map(|k| k.0).collect(); - let actor_data = state.id_cache.get_actor_data_many(&actor_ids).await; + let actor_data = { + let mut conn = state.pool.get().await?; + crate::db::get_actor_data_by_ids(&mut conn, &actor_ids, &state.id_cache) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve actor data for trending feeds: {e}"); + std::collections::HashMap::new() + }) + }; let page_uris: Vec = top_feedgens .iter() @@ -429,9 +437,17 @@ pub async fn get_popular_feed_generators( None }; - // Construct URIs from natural keys by resolving actor_ids to DIDs + // Construct URIs from natural keys by resolving actor_ids to DIDs (with database fallback for cache misses) let actor_ids: Vec = page_feedgens.iter().map(|k| k.0).collect(); - let actor_data = state.id_cache.get_actor_data_many(&actor_ids).await; + let actor_data = { + let mut conn = state.pool.get().await?; + crate::db::get_actor_data_by_ids(&mut conn, &actor_ids, &state.id_cache) + .await + .unwrap_or_else(|e| { + tracing::warn!("Failed to resolve actor data for suggested feeds: {e}"); + std::collections::HashMap::new() + }) + }; let page_uris: Vec = page_feedgens .iter()