diff --git a/parakeet/src/loaders/feed.rs b/parakeet/src/loaders/feed.rs index 333a2ebc..5456a94f 100644 --- a/parakeet/src/loaders/feed.rs +++ b/parakeet/src/loaders/feed.rs @@ -40,6 +40,7 @@ impl std::ops::Deref for EnrichedFeedGen { /// Build SQL query for batch loading feed generators with computed fields /// /// This function is public for testing purposes. +/// Note: service_did resolved via IdCache to avoid actors JOIN pub fn build_feedgens_batch_query() -> &'static str { "SELECT f.actor_id, @@ -55,16 +56,17 @@ pub fn build_feedgens_batch_query() -> &'static str { f.avatar_cid, f.accepts_interactions, f.status, - f.like_count, - service.did as service_did + f.like_count FROM feedgens f - INNER JOIN actors service ON f.service_actor_id = service.id WHERE f.owner_actor_id = ANY($1) AND f.rkey::text = ANY($2) AND f.status = 'complete'" } -pub struct FeedGenLoader(pub(super) Pool); +pub struct FeedGenLoader( + pub(super) Pool, + pub(super) std::sync::Arc, +); impl BatchFn for FeedGenLoader { async fn load(&mut self, keys: &[FeedGenKey]) -> HashMap { let mut conn = self.0.get().await.unwrap(); @@ -105,8 +107,6 @@ impl BatchFn for FeedGenLoader { status: parakeet_db::types::FeedgenStatus, #[diesel(sql_type = diesel::sql_types::Integer)] like_count: i32, - #[diesel(sql_type = diesel::sql_types::Text)] - service_did: String, } let res: Vec = diesel_async::RunQueryDsl::load( @@ -121,7 +121,28 @@ impl BatchFn for FeedGenLoader { vec![] }); + // Collect unique service_actor_ids for DID resolution + let service_actor_ids: Vec = res + .iter() + .map(|row| row.service_actor_id) + .collect::>() + .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() + }; + HashMap::from_iter(res.into_iter().map(|row| { + // Get service DID from resolved data + let service_did = service_actor_data + .get(&row.service_actor_id) + .map(|d| d.did.clone()) + .unwrap_or_else(|| String::from("did:unknown")); + let enriched = EnrichedFeedGen { feedgen: models::FeedGen { actor_id: row.actor_id, @@ -141,7 +162,7 @@ impl BatchFn for FeedGenLoader { }, at_uri: String::new(), // Will be constructed by caller owner: String::new(), // Will be resolved by caller - service_did: row.service_did, + service_did, cid: row.cid, created_at: row.created_at, like_count: row.like_count, diff --git a/parakeet/src/loaders/misc.rs b/parakeet/src/loaders/misc.rs index d4192393..d2cd929a 100644 --- a/parakeet/src/loaders/misc.rs +++ b/parakeet/src/loaders/misc.rs @@ -96,22 +96,20 @@ fn build_starterpack_record( /// Build SQL query for batch loading starter packs with computed fields /// /// This function is public for testing purposes. +/// Note: created_at is computed in Rust from rkey, list_owner_did resolved via IdCache pub fn build_starterpacks_batch_query() -> &'static str { "SELECT sp.actor_id, sp.rkey, sp.cid, - tid_timestamp(sp.rkey) as created_at, sp.owner_actor_id, sp.list_actor_id, sp.list_rkey, sp.name, sp.description, sp.description_facets, - sp.status, - list_owner.did as list_owner_did + sp.status FROM starterpacks sp - LEFT JOIN actors list_owner ON sp.list_actor_id = list_owner.id WHERE sp.owner_actor_id = ANY($1) AND sp.rkey = ANY($2)" } @@ -131,7 +129,10 @@ pub fn build_starterpack_feeds_query() -> &'static str { ORDER BY sf.starterpack_id, sf.position" } -pub struct StarterPackLoader(pub(super) Pool); +pub struct StarterPackLoader( + pub(super) Pool, + pub(super) std::sync::Arc, +); pub type StarterPackLoaderRet = EnrichedStarterPack; impl BatchFn for StarterPackLoader { async fn load(&mut self, keys: &[StarterPackKey]) -> HashMap { @@ -151,8 +152,6 @@ impl BatchFn for StarterPackLoader { rkey: i64, #[diesel(sql_type = diesel::sql_types::Binary)] cid: Vec, - #[diesel(sql_type = diesel::sql_types::Timestamptz)] - created_at: chrono::DateTime, #[diesel(sql_type = diesel::sql_types::Integer)] owner_actor_id: i32, #[diesel(sql_type = diesel::sql_types::Nullable)] @@ -167,8 +166,6 @@ impl BatchFn for StarterPackLoader { description_facets: Option, #[diesel(sql_type = parakeet_db::schema::sql_types::RecordStatus)] status: parakeet_db::types::RecordStatus, - #[diesel(sql_type = diesel::sql_types::Nullable)] - list_owner_did: Option, } let starterpacks: Vec = diesel_async::RunQueryDsl::load( @@ -183,6 +180,21 @@ impl BatchFn for StarterPackLoader { vec![] }); + // Collect unique list_actor_ids for DID resolution via IdCache + let list_actor_ids: Vec = starterpacks + .iter() + .filter_map(|sp| sp.list_actor_id) + .collect::>() + .into_iter() + .collect(); + + // Batch resolve list owner DIDs via IdCache + let list_owner_data = if !list_actor_ids.is_empty() { + self.1.get_actor_data_many(&list_actor_ids).await + } else { + std::collections::HashMap::new() + }; + // Load feeds for all starterpacks (with URIs already resolved) let sp_rkeys: Vec = starterpacks.iter().map(|sp| sp.rkey).collect(); @@ -214,9 +226,17 @@ impl BatchFn for StarterPackLoader { }; HashMap::from_iter(starterpacks.into_iter().map(|row| { - // Build list URI if list exists (still needed for the record) - let list_uri = if let (Some(list_owner_did), Some(list_rkey)) = (&row.list_owner_did, &row.list_rkey) { - format!("at://{}/app.bsky.graph.list/{}", list_owner_did, list_rkey) + // Compute created_at from rkey TID in Rust + let created_at = parakeet_db::tid_util::tid_to_datetime(row.rkey); + + // Build list URI if list exists (using resolved DID from IdCache) + let list_uri = if let (Some(list_actor_id), Some(list_rkey)) = (row.list_actor_id, &row.list_rkey) { + if let Some(list_owner_data) = list_owner_data.get(&list_actor_id) { + format!("at://{}/app.bsky.graph.list/{}", list_owner_data.did, list_rkey) + } else { + // Fallback: could not resolve list owner DID + String::new() + } } else { String::new() }; @@ -230,7 +250,7 @@ impl BatchFn for StarterPackLoader { &row.description_facets, &list_uri, &feeds, - &row.created_at, + &created_at, ); let enriched = EnrichedStarterPack { @@ -251,7 +271,7 @@ impl BatchFn for StarterPackLoader { list: list_uri, feeds, cid: row.cid, - created_at: row.created_at, // Computed from tid_timestamp(rkey) in SQL + created_at, // Computed from rkey in Rust record, }; (StarterPackKey(row.owner_actor_id, row.rkey), enriched) @@ -262,23 +282,19 @@ impl BatchFn for StarterPackLoader { /// Build SQL query for batch loading verifications with computed fields /// /// This function is public for testing purposes. +/// Note: created_at computed in Rust from rkey, DIDs resolved via IdCache pub fn build_verifications_batch_query() -> &'static str { "SELECT v.id, v.actor_id, v.rkey, v.cid, - tid_timestamp(v.rkey) as created_at, v.verifier_actor_id, v.subject_actor_id, v.handle, - v.display_name, - subject_a.did as subject_did, - verifier_a.did as verifier_did + v.display_name FROM verification v - INNER JOIN actors subject_a ON v.subject_actor_id = subject_a.id - INNER JOIN actors verifier_a ON v.verifier_actor_id = verifier_a.id - WHERE subject_a.did = ANY($1)" + WHERE v.subject_actor_id = ANY($1)" } pub struct VerificationLoader(pub(super) Pool, pub(super) std::sync::Arc); @@ -286,8 +302,15 @@ impl BatchFn> for VerificationLoader { async fn load(&mut self, keys: &[String]) -> HashMap> { let mut conn = self.0.get().await.unwrap(); - // Build SQL query with URI reconstruction - let refs: Vec<&str> = keys.iter().map(|s| s.as_str()).collect(); + // Resolve DIDs to actor_ids using IdCache + let did_to_actor = self.1.get_actor_ids(keys).await; + + // Collect actor_ids for query + let subject_actor_ids: Vec = did_to_actor.values().map(|a| a.actor_id).collect(); + + if subject_actor_ids.is_empty() { + return HashMap::new(); + } let query = build_verifications_batch_query(); @@ -301,8 +324,6 @@ impl BatchFn> for VerificationLoader { rkey: i64, #[diesel(sql_type = diesel::sql_types::Binary)] cid: Vec, - #[diesel(sql_type = diesel::sql_types::Timestamptz)] - created_at: chrono::DateTime, #[diesel(sql_type = diesel::sql_types::Integer)] verifier_actor_id: i32, #[diesel(sql_type = diesel::sql_types::Integer)] @@ -311,15 +332,11 @@ impl BatchFn> for VerificationLoader { handle: String, #[diesel(sql_type = diesel::sql_types::Text)] display_name: String, - #[diesel(sql_type = diesel::sql_types::Text)] - subject_did: String, - #[diesel(sql_type = diesel::sql_types::Text)] - verifier_did: String, } let verifications: Vec = diesel_async::RunQueryDsl::load( diesel::sql_query(query) - .bind::, _>(&refs), + .bind::, _>(&subject_actor_ids), &mut conn, ) .await @@ -328,12 +345,35 @@ impl BatchFn> for VerificationLoader { vec![] }); + // Collect unique actor_ids for DID resolution + let mut all_actor_ids: Vec = Vec::new(); + for row in &verifications { + all_actor_ids.push(row.verifier_actor_id); + all_actor_ids.push(row.subject_actor_id); + } + 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; + // Group by subject DID let mut result: HashMap> = HashMap::new(); for row in verifications { + // Get DIDs from resolved data + let subject_did = actor_data.get(&row.subject_actor_id) + .map(|d| d.did.clone()) + .unwrap_or_else(|| String::from("did:unknown")); + let verifier_did = actor_data.get(&row.verifier_actor_id) + .map(|d| d.did.clone()) + .unwrap_or_else(|| String::from("did:unknown")); + + // Compute created_at from rkey TID in Rust + let created_at = parakeet_db::tid_util::tid_to_datetime(row.rkey); + // Encode TID using Rust utility function let encoded_rkey = parakeet_db::tid_util::encode_tid(row.rkey); - let at_uri = format!("at://{}/dev.bsky.verification.verification/{}", row.subject_did, encoded_rkey); + let at_uri = format!("at://{}/dev.bsky.verification.verification/{}", subject_did, encoded_rkey); let enriched = EnrichedVerification { verification: models::Verification { @@ -346,11 +386,11 @@ impl BatchFn> for VerificationLoader { handle: row.handle, display_name: row.display_name, }, - verifier: row.verifier_did, + verifier: verifier_did, at_uri, - created_at: row.created_at, // Computed from tid_timestamp(rkey) in SQL + created_at, // Computed from rkey in Rust }; - result.entry(row.subject_did).or_default().push(enriched); + result.entry(subject_did.clone()).or_default().push(enriched); } result diff --git a/parakeet/src/loaders/mod.rs b/parakeet/src/loaders/mod.rs index 0fa5c652..66e7b23d 100644 --- a/parakeet/src/loaders/mod.rs +++ b/parakeet/src/loaders/mod.rs @@ -97,10 +97,10 @@ impl Dataloaders { // 1 hour TTL, 50k capacity for profiles profile_by_id: new_plc_loader(ProfileByIdLoader(pool.clone(), id_cache.clone()), "profile_id:", 3600, 50_000), // 1 hour TTL, 10k capacity for feeds/lists/etc - feedgen: new_plc_loader(FeedGenLoader(pool.clone()), "feedgen:", 3600, 10_000), + feedgen: new_plc_loader(FeedGenLoader(pool.clone(), id_cache.clone()), "feedgen:", 3600, 10_000), labeler: new_plc_loader(LabelServiceLoader(pool.clone()), "labeler:", 3600, 10_000), list: new_plc_loader(ListLoader(pool.clone()), "list:", 3600, 10_000), - starterpacks: new_plc_loader(StarterPackLoader(pool.clone()), "starterpacks:", 3600, 10_000), + starterpacks: new_plc_loader(StarterPackLoader(pool.clone(), id_cache.clone()), "starterpacks:", 3600, 10_000), verification: new_plc_loader(VerificationLoader(pool.clone(), id_cache.clone()), "verification:", 3600, 10_000), verification_raw: VerificationLoader(pool.clone(), id_cache.clone()),