diff --git a/parakeet/src/hydration/embed.rs b/parakeet/src/hydration/embed.rs index 4be621b4..65c28266 100644 --- a/parakeet/src/hydration/embed.rs +++ b/parakeet/src/hydration/embed.rs @@ -1,5 +1,5 @@ use crate::hydration::StatefulHydrator; -use crate::loaders::{EmbedLoaderRet, EnrichedPostEmbedRecord}; +use crate::loaders::{EmbedLoaderRet, EnrichedPostEmbedRecord, PostEmbedImage}; use crate::xrpc::cdn::BskyCdn; use itertools::Itertools as _; use lexica::app_bsky::embed::{ @@ -14,6 +14,66 @@ fn build_aspect_ratio(height: Option, width: Option) -> Option Vec { + [ + post.image_1.as_ref(), + post.image_2.as_ref(), + post.image_3.as_ref(), + post.image_4.as_ref(), + ] + .into_iter() + .enumerate() + .filter_map(|(idx, img_opt)| { + img_opt.map(|img| PostEmbedImage { + post_actor_id: post.post.actor_id, + post_rkey: post.post.rkey, + seq: idx as i16, + mime_type: img.mime_type, + cid: img.cid.clone(), + alt: img.alt.clone(), + width: img.width, + height: img.height, + }) + }) + .collect() +} + +/// Helper function to extract video embed from a HydratedPost +fn extract_post_video(post: &crate::loaders::HydratedPost) -> Option { + use crate::loaders::PostEmbedVideo; + + post.video_embed.as_ref().map(|ve| { + PostEmbedVideo { + post_actor_id: post.post.actor_id, + post_rkey: post.post.rkey, + mime_type: ve.mime_type, + cid: ve.cid.clone(), + alt: ve.alt.clone(), + width: ve.width, + height: ve.height, + } + }) +} + +/// Helper function to extract external embed from a HydratedPost +fn extract_post_external(post: &crate::loaders::HydratedPost) -> Option { + use crate::loaders::PostEmbedExt; + + post.ext_embed.as_ref().map(|ee| { + PostEmbedExt { + post_actor_id: post.post.actor_id, + post_rkey: post.post.rkey, + uri: ee.uri.clone(), + title: ee.title.clone().unwrap_or_default(), + description: ee.description.clone().unwrap_or_default(), + thumb_mime_type: ee.thumb_mime_type, + thumb_cid: ee.thumb_cid.clone(), + } + }) +} + fn build_record_view(post: PostView) -> RecordView { RecordView { uri: post.uri, @@ -84,34 +144,12 @@ impl StatefulHydrator<'_> { /// Extract embed from a HydratedPost (no database query needed) /// This replaces the slow EmbedLoader path that was querying the posts table again pub async fn hydrate_embed_from_post(&self, post: &crate::loaders::HydratedPost) -> Option { - use crate::loaders::{PostEmbedImage, PostEmbedVideo, PostEmbedExt}; let embed_type = post.post.embed_type.as_ref()?; let embed_ret = match embed_type { parakeet_db::types::EmbedType::Images => { - let images: Vec = [ - post.image_1.as_ref(), - post.image_2.as_ref(), - post.image_3.as_ref(), - post.image_4.as_ref(), - ] - .into_iter() - .enumerate() - .filter_map(|(idx, img_opt)| { - img_opt.map(|img| PostEmbedImage { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - seq: idx as i16, - mime_type: img.mime_type, - cid: img.cid.clone(), - alt: img.alt.clone(), - width: img.width, - height: img.height, - }) - }) - .collect(); - + let images = extract_post_images(post); if !images.is_empty() { Some(EmbedLoaderRet::Images(images)) } else { @@ -119,30 +157,10 @@ impl StatefulHydrator<'_> { } } parakeet_db::types::EmbedType::Video => { - post.video_embed.as_ref().map(|ve| { - EmbedLoaderRet::Video(PostEmbedVideo { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - mime_type: ve.mime_type, - cid: ve.cid.clone(), - alt: ve.alt.clone(), - width: ve.width, - height: ve.height, - }) - }) + extract_post_video(post).map(EmbedLoaderRet::Video) } parakeet_db::types::EmbedType::External => { - post.ext_embed.as_ref().map(|ee| { - EmbedLoaderRet::External(PostEmbedExt { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - uri: ee.uri.clone(), - title: ee.title.clone().unwrap_or_default(), - description: ee.description.clone().unwrap_or_default(), - thumb_mime_type: ee.thumb_mime_type, - thumb_cid: ee.thumb_cid.clone(), - }) - }) + extract_post_external(post).map(EmbedLoaderRet::External) } parakeet_db::types::EmbedType::Record => { post.embedded_rkey.and_then(|embedded_rkey_bigint| { @@ -179,28 +197,7 @@ impl StatefulHydrator<'_> { let media = match post.post.embed_subtype.as_ref()? { parakeet_db::types::EmbedType::Images => { - let images: Vec = [ - post.image_1.as_ref(), - post.image_2.as_ref(), - post.image_3.as_ref(), - post.image_4.as_ref(), - ] - .into_iter() - .enumerate() - .filter_map(|(idx, img_opt)| { - img_opt.map(|img| PostEmbedImage { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - seq: idx as i16, - mime_type: img.mime_type, - cid: img.cid.clone(), - alt: img.alt.clone(), - width: img.width, - height: img.height, - }) - }) - .collect(); - + let images = extract_post_images(post); if !images.is_empty() { Some(EmbedLoaderRet::Images(images)) } else { @@ -208,30 +205,10 @@ impl StatefulHydrator<'_> { } } parakeet_db::types::EmbedType::Video => { - post.video_embed.as_ref().map(|ve| { - EmbedLoaderRet::Video(PostEmbedVideo { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - mime_type: ve.mime_type, - cid: ve.cid.clone(), - alt: ve.alt.clone(), - width: ve.width, - height: ve.height, - }) - }) + extract_post_video(post).map(EmbedLoaderRet::Video) } parakeet_db::types::EmbedType::External => { - post.ext_embed.as_ref().map(|ee| { - EmbedLoaderRet::External(PostEmbedExt { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - uri: ee.uri.clone(), - title: ee.title.clone().unwrap_or_default(), - description: ee.description.clone().unwrap_or_default(), - thumb_mime_type: ee.thumb_mime_type, - thumb_cid: ee.thumb_cid.clone(), - }) - }) + extract_post_external(post).map(EmbedLoaderRet::External) } _ => None, }?; @@ -430,7 +407,6 @@ impl StatefulHydrator<'_> { &self, posts: &HashMap, Option)>, ) -> HashMap { - use crate::loaders::{PostEmbedImage, PostEmbedVideo, PostEmbedExt}; let conversion_start = std::time::Instant::now(); @@ -442,28 +418,7 @@ impl StatefulHydrator<'_> { let embed_ret = match embed_type { parakeet_db::types::EmbedType::Images => { - let images: Vec = [ - post.image_1.as_ref(), - post.image_2.as_ref(), - post.image_3.as_ref(), - post.image_4.as_ref(), - ] - .into_iter() - .enumerate() - .filter_map(|(idx, img_opt)| { - img_opt.map(|img| PostEmbedImage { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - seq: idx as i16, - mime_type: img.mime_type, - cid: img.cid.clone(), - alt: img.alt.clone(), - width: img.width, - height: img.height, - }) - }) - .collect(); - + let images = extract_post_images(post); if !images.is_empty() { Some(EmbedLoaderRet::Images(images)) } else { @@ -471,30 +426,10 @@ impl StatefulHydrator<'_> { } } parakeet_db::types::EmbedType::Video => { - post.video_embed.as_ref().map(|ve| { - EmbedLoaderRet::Video(PostEmbedVideo { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - mime_type: ve.mime_type, - cid: ve.cid.clone(), - alt: ve.alt.clone(), - width: ve.width, - height: ve.height, - }) - }) + extract_post_video(post).map(EmbedLoaderRet::Video) } parakeet_db::types::EmbedType::External => { - post.ext_embed.as_ref().map(|ee| { - EmbedLoaderRet::External(PostEmbedExt { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - uri: ee.uri.clone(), - title: ee.title.clone().unwrap_or_default(), - description: ee.description.clone().unwrap_or_default(), - thumb_mime_type: ee.thumb_mime_type, - thumb_cid: ee.thumb_cid.clone(), - }) - }) + extract_post_external(post).map(EmbedLoaderRet::External) } parakeet_db::types::EmbedType::Record => { post.embedded_rkey.and_then(|embedded_rkey_bigint| { @@ -531,55 +466,14 @@ impl StatefulHydrator<'_> { let media = match post.post.embed_subtype.as_ref()? { parakeet_db::types::EmbedType::Images => { - let images: Vec = [ - post.image_1.as_ref(), - post.image_2.as_ref(), - post.image_3.as_ref(), - post.image_4.as_ref(), - ] - .into_iter() - .enumerate() - .filter_map(|(idx, img_opt)| { - img_opt.map(|img| PostEmbedImage { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - seq: idx as i16, - mime_type: img.mime_type, - cid: img.cid.clone(), - alt: img.alt.clone(), - width: img.width, - height: img.height, - }) - }) - .collect(); - + let images = extract_post_images(post); (!images.is_empty()).then_some(EmbedLoaderRet::Images(images)) } parakeet_db::types::EmbedType::Video => { - post.video_embed.as_ref().map(|ve| { - EmbedLoaderRet::Video(PostEmbedVideo { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - mime_type: ve.mime_type, - cid: ve.cid.clone(), - alt: ve.alt.clone(), - width: ve.width, - height: ve.height, - }) - }) + extract_post_video(post).map(EmbedLoaderRet::Video) } parakeet_db::types::EmbedType::External => { - post.ext_embed.as_ref().map(|ee| { - EmbedLoaderRet::External(PostEmbedExt { - post_actor_id: post.post.actor_id, - post_rkey: post.post.rkey, - uri: ee.uri.clone(), - title: ee.title.clone().unwrap_or_default(), - description: ee.description.clone().unwrap_or_default(), - thumb_mime_type: ee.thumb_mime_type, - thumb_cid: ee.thumb_cid.clone(), - }) - }) + extract_post_external(post).map(EmbedLoaderRet::External) } _ => None, }?; diff --git a/parakeet/src/hydration/posts/builders.rs b/parakeet/src/hydration/posts/builders.rs index 9212c22a..94842814 100644 --- a/parakeet/src/hydration/posts/builders.rs +++ b/parakeet/src/hydration/posts/builders.rs @@ -7,6 +7,7 @@ use lexica::app_bsky::graph::ListViewBasic; use lexica::app_bsky::RecordStats; use parakeet_db::models; use parakeet_db::models::PostStats; +use std::collections::HashMap; pub(super) type HydratePostsRet = ( crate::loaders::HydratedPost, @@ -21,6 +22,16 @@ pub(super) type HydratePostsRet = ( pub(super) fn build_postview( (post, author, labels, embed, threadgate, viewer, stats): HydratePostsRet, id_cache: ¶keet_db::id_cache::IdCache, +) -> PostView { + build_postview_with_cache((post, author, labels, embed, threadgate, viewer, stats), id_cache, &HashMap::new()) +} + +/// Version of build_postview that accepts a pre-fetched actor ID to DID map +/// to avoid blocking async calls +pub(super) fn build_postview_with_cache( + (post, author, labels, embed, threadgate, viewer, stats): HydratePostsRet, + id_cache: ¶keet_db::id_cache::IdCache, + actor_id_to_did_cache: &HashMap, ) -> PostView { let stats = stats .map(|stats| RecordStats { @@ -40,9 +51,16 @@ pub(super) fn build_postview( // Convert _parent_key to parent URI if let Some(parent_key) = record.get("_parent_key").and_then(|v| v.as_array()) { if let (Some(actor_id), Some(rkey)) = (parent_key.get(0).and_then(|v| v.as_i64()), parent_key.get(1).and_then(|v| v.as_i64())) { - if let Some(actor_data) = futures::executor::block_on(id_cache.get_actor_data(actor_id as i32)) { + let actor_id = actor_id as i32; + // Try cache first, then fall back to blocking call + let did_opt = actor_id_to_did_cache.get(&actor_id) + .cloned() + .or_else(|| futures::executor::block_on(id_cache.get_actor_data(actor_id)) + .map(|data| data.did)); + + if let Some(did) = did_opt { let encoded_rkey = parakeet_db::tid_util::encode_tid(rkey); - let parent_uri = format!("at://{}/app.bsky.feed.post/{}", actor_data.did, encoded_rkey); + let parent_uri = format!("at://{}/app.bsky.feed.post/{}", did, encoded_rkey); let parent_cid = record.get("_parent_cid").and_then(|v| v.as_str()).unwrap_or("bafyrei_invalid_cid"); reply["parent"] = serde_json::json!({"uri": parent_uri, "cid": parent_cid}); } @@ -52,9 +70,16 @@ pub(super) fn build_postview( // Convert _root_key to root URI if let Some(root_key) = record.get("_root_key").and_then(|v| v.as_array()) { if let (Some(actor_id), Some(rkey)) = (root_key.get(0).and_then(|v| v.as_i64()), root_key.get(1).and_then(|v| v.as_i64())) { - if let Some(actor_data) = futures::executor::block_on(id_cache.get_actor_data(actor_id as i32)) { + let actor_id = actor_id as i32; + // Try cache first, then fall back to blocking call + let did_opt = actor_id_to_did_cache.get(&actor_id) + .cloned() + .or_else(|| futures::executor::block_on(id_cache.get_actor_data(actor_id)) + .map(|data| data.did)); + + if let Some(did) = did_opt { let encoded_rkey = parakeet_db::tid_util::encode_tid(rkey); - let root_uri = format!("at://{}/app.bsky.feed.post/{}", actor_data.did, encoded_rkey); + let root_uri = format!("at://{}/app.bsky.feed.post/{}", did, encoded_rkey); let root_cid = record.get("_root_cid").and_then(|v| v.as_str()).unwrap_or("bafyrei_invalid_cid"); reply["root"] = serde_json::json!({"uri": root_uri, "cid": root_cid}); } @@ -94,6 +119,18 @@ pub(super) fn build_threadgate_view( post_did: &str, lists: Vec, id_cache: ¶keet_db::id_cache::IdCache, +) -> ThreadgateView { + build_threadgate_view_with_cache(threadgate, post_did, lists, id_cache, &HashMap::new()) +} + +/// Version of build_threadgate_view that accepts a pre-fetched actor ID to DID map +/// to avoid blocking async calls +pub(super) fn build_threadgate_view_with_cache( + threadgate: EnrichedThreadgate, + post_did: &str, + lists: Vec, + id_cache: ¶keet_db::id_cache::IdCache, + actor_id_to_did_cache: &HashMap, ) -> ThreadgateView { // Construct threadgate URI at the edge let encoded_rkey = parakeet_db::tid_util::encode_tid(threadgate.rkey); @@ -115,14 +152,17 @@ pub(super) fn build_threadgate_view( None } }) { - // Look up DID from IdCache at hydration time (the edge!) - // Note: This is a synchronous blocking call - ideally we'd batch these - // but for now this is acceptable since threadgates with hidden replies are rare - if let Some(actor_data) = futures::executor::block_on( - id_cache.get_actor_data(actor_id as i32) - ) { + // Look up DID from cache first, then fall back to blocking call + let actor_id_i32 = actor_id as i32; + let did_opt = actor_id_to_did_cache.get(&actor_id_i32) + .cloned() + .or_else(|| futures::executor::block_on( + id_cache.get_actor_data(actor_id_i32) + ).map(|data| data.did)); + + if let Some(did) = did_opt { let reply_rkey_encoded = parakeet_db::tid_util::encode_tid(rkey); - hidden_uris.push(format!("at://{}/app.bsky.feed.post/{}", actor_data.did, reply_rkey_encoded)); + hidden_uris.push(format!("at://{}/app.bsky.feed.post/{}", did, reply_rkey_encoded)); } } } diff --git a/parakeet/src/hydration/posts/mod.rs b/parakeet/src/hydration/posts/mod.rs index 7b5e88d8..f9517f87 100644 --- a/parakeet/src/hydration/posts/mod.rs +++ b/parakeet/src/hydration/posts/mod.rs @@ -112,6 +112,40 @@ impl StatefulHydrator<'_> { self.hydrate_post(uri).await } + /// Collect actor IDs from reply fields that need DID lookups + fn collect_reply_actor_ids(posts: &HashMap, Option)>) -> Vec { + let mut actor_ids = std::collections::HashSet::new(); + + for (_, (post, threadgate_opt, _)) in posts { + // Check for parent and root actor IDs in the record + if let Some(parent_key) = post.record.get("_parent_key").and_then(|v| v.as_array()) { + if let Some(actor_id) = parent_key.get(0).and_then(|v| v.as_i64()) { + actor_ids.insert(actor_id as i32); + } + } + if let Some(root_key) = post.record.get("_root_key").and_then(|v| v.as_array()) { + if let Some(actor_id) = root_key.get(0).and_then(|v| v.as_i64()) { + actor_ids.insert(actor_id as i32); + } + } + + // Check for hidden reply actor IDs in threadgates + if let Some(threadgate) = threadgate_opt { + if let Some(hidden_keys) = threadgate.record.get("_hidden_reply_keys").and_then(|v| v.as_array()) { + for key in hidden_keys { + if let Some(arr) = key.as_array() { + if let Some(actor_id) = arr.get(0).and_then(|v| v.as_i64()) { + actor_ids.insert(actor_id as i32); + } + } + } + } + } + } + + actor_ids.into_iter().collect() + } + pub(super) async fn hydrate_posts_inner( &self, posts: Vec, @@ -124,6 +158,22 @@ impl StatefulHydrator<'_> { tracing::warn!(" → Slow load stats and posts: {:.1} ms", load_time); } + // Pre-fetch actor IDs needed for reply and threadgate construction + 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; + 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()); + } + actor_data.into_iter() + .map(|(id, data)| (id, data.did)) + .collect::>() + } else { + HashMap::new() + }; + let stats: HashMap = posts_with_stats .iter() .filter_map(|(uri, (_, _, stats_opt))| { @@ -254,6 +304,9 @@ impl StatefulHydrator<'_> { } pub async fn hydrate_posts(&self, posts: Vec) -> HashMap { + // For now, continue using the original build_postview which will fall back to blocking calls + // The optimization is in place but we'd need to refactor hydrate_posts_inner to return + // the actor cache to fully utilize it. This is a future optimization. self.hydrate_posts_inner(posts) .await .into_iter()