From b980236b4a2e4a48f2406b66d493eb3d69540f88 Mon Sep 17 00:00:00 2001 From: Timothy Quilling Date: Mon, 8 Dec 2025 16:34:51 -0500 Subject: [PATCH] feat: denormalization strategy; more tdd --- consumer/src/db/bulk_resolve/mod.rs | 64 ++++++++++ consumer/src/db/gates/queries.rs | 11 +- consumer/src/db/operations/feed/threadgate.rs | 52 +++++++- consumer/src/db/sql/bookmarks_upsert.sql | 7 -- consumer/src/db/sql/postgate_upsert.sql | 41 ------- consumer/src/db/sql/threadgate_upsert.sql | 59 ---------- .../db/sql/threadgate_upsert_denormalized.sql | 5 +- .../down.sql | 5 + .../up.sql | 10 ++ parakeet-db/src/schema.rs | 2 + parakeet/src/db/graph.rs | 54 +++++---- parakeet/src/db/notification_records.rs | 14 ++- parakeet/src/db/suggestions.rs | 2 +- parakeet/src/loaders/labeler.rs | 91 ++++++++------ parakeet/src/loaders/profile.rs | 6 +- parakeet/src/sql/list_states.sql | 59 +++++++--- parakeet/src/sql/post_state.sql | 17 +-- parakeet/src/sql/profile_state.sql | 111 ++++++++++++------ parakeet/src/sql/profile_state_by_ids.sql | 104 ++++++++++------ parakeet/src/xrpc/app_bsky/graph/lists.rs | 7 +- parakeet/tests/db_graph_test.rs | 4 +- parakeet/tests/db_uri_reconstruction_test.rs | 76 ------------ parakeet/tests/sql/batch_loading_test.rs | 15 ++- parakeet/tests/sql/db_module_test.rs | 6 +- 24 files changed, 448 insertions(+), 374 deletions(-) delete mode 100644 consumer/src/db/sql/bookmarks_upsert.sql delete mode 100644 consumer/src/db/sql/postgate_upsert.sql delete mode 100644 consumer/src/db/sql/threadgate_upsert.sql create mode 100644 migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/down.sql create mode 100644 migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/up.sql delete mode 100644 parakeet/tests/db_uri_reconstruction_test.rs diff --git a/consumer/src/db/bulk_resolve/mod.rs b/consumer/src/db/bulk_resolve/mod.rs index bdbe75c6..425360bc 100644 --- a/consumer/src/db/bulk_resolve/mod.rs +++ b/consumer/src/db/bulk_resolve/mod.rs @@ -501,3 +501,67 @@ pub async fn resolve_uri_ids_bulk( Ok(result) } + +/// Resolve multiple list AT-URIs to natural keys (actor_id, rkey) +/// +/// Returns a HashMap of AT-URI → (actor_id, rkey) for all URIs. +/// Lists use TEXT rkey (not i64 like posts), so rkey is returned as String. +/// +/// This function ensures actors exist (creates stubs if needed), then +/// returns natural key pairs based on parsing the URIs. +/// +/// # Performance +/// +/// Only needs to resolve actor DIDs, not query the lists table. +pub async fn resolve_list_uris_bulk( + conn: &C, + at_uris: &[&str], +) -> Result> { + if at_uris.is_empty() { + return Ok(HashMap::new()); + } + + // Parse URIs to extract (did, rkey) pairs + let mut uri_to_did_rkey: HashMap = HashMap::new(); + let mut dids_set: std::collections::HashSet = std::collections::HashSet::new(); + + for uri in at_uris { + let did = parakeet_db::at_uri_util::extract_did(uri) + .ok_or_else(|| eyre::eyre!("Invalid AT URI: missing DID in {}", uri))?; + let rkey = parakeet_db::at_uri_util::extract_rkey(uri) + .ok_or_else(|| eyre::eyre!("Invalid AT URI: missing rkey in {}", uri))?; + + uri_to_did_rkey.insert(uri.to_string(), (did.to_string(), rkey.to_string())); + dids_set.insert(did.to_string()); + } + + // Resolve all actor DIDs to actor_ids (creates stubs for missing actors) + let dids_vec: Vec<&str> = dids_set.iter().map(|s| s.as_str()).collect(); + let did_to_actor_id = resolve_actor_dids_bulk(conn, &dids_vec).await?; + + // Create actor stubs for any missing DIDs + let missing_dids: Vec<&str> = dids_vec + .iter() + .filter(|did| !did_to_actor_id.contains_key(**did)) + .copied() + .collect(); + + let mut did_to_actor_id = did_to_actor_id; + if !missing_dids.is_empty() { + let created = create_actor_stubs_bulk(conn, &missing_dids).await?; + did_to_actor_id.extend(created); + } + + // Build result HashMap mapping URI → (actor_id, rkey) + let mut result = HashMap::new(); + for (uri, (did, rkey)) in uri_to_did_rkey { + if let Some(&actor_id) = did_to_actor_id.get(&did) { + result.insert(uri, (actor_id, rkey)); + } else { + // This shouldn't happen since we create stubs for missing actors + eyre::bail!("Failed to resolve actor ID for DID {} in list URI {}", did, uri); + } + } + + Ok(result) +} diff --git a/consumer/src/db/gates/queries.rs b/consumer/src/db/gates/queries.rs index 7c23a5f2..ef002a9f 100644 --- a/consumer/src/db/gates/queries.rs +++ b/consumer/src/db/gates/queries.rs @@ -77,7 +77,8 @@ pub async fn get_post_mentions( /// Check if post author is in any of the allowed lists /// -/// Returns count of list memberships found +/// Returns count of list memberships found. +/// Uses natural keys (actor_id, rkey) to check list membership. pub async fn check_list_membership( conn: &C, allow_lists: &[String], @@ -98,14 +99,16 @@ pub async fn check_list_membership( SELECT a.id as list_owner_actor_id, lp.rkey as list_rkey FROM list_parts lp INNER JOIN actors a ON a.did = lp.did - INNER JOIN lists l ON l.actor_id = a.id AND l.rkey = lp.rkey + -- No need to join lists table - we're using natural keys + ), + post_author_id AS ( + SELECT id FROM actors WHERE did = $2 ) SELECT count(*) FROM list_items li - INNER JOIN actors a ON a.did = $2 INNER JOIN list_keys lk ON li.list_owner_actor_id = lk.list_owner_actor_id AND li.list_rkey = lk.list_rkey - WHERE li.subject_actor_id = a.id", + WHERE li.subject_actor_id = (SELECT id FROM post_author_id)", &[&allow_lists, &post_author], ) .await?; diff --git a/consumer/src/db/operations/feed/threadgate.rs b/consumer/src/db/operations/feed/threadgate.rs index f1c5525e..71f4ac6c 100644 --- a/consumer/src/db/operations/feed/threadgate.rs +++ b/consumer/src/db/operations/feed/threadgate.rs @@ -21,7 +21,8 @@ use ipld_core::cid::Cid; /// /// Returns (created_at, allow_rules, allowed_list_uris) if threadgate exists. /// The allow_rules are string representations like "app.bsky.feed.threadgate#mentionRule". -/// The allowed_list_uris are full AT URIs for list-based reply permissions. +/// The allowed_list_uris are full AT URIs for list-based reply permissions, reconstructed +/// from the natural keys (actor_id, rkey) stored in the posts table. /// /// DENORMALIZED: Reads threadgate_allow from posts table instead of separate threadgates table. pub async fn threadgate_get( @@ -39,7 +40,19 @@ pub async fn threadgate_get( SELECT tid_timestamp(p.rkey) as created_at, p.threadgate_allow::text[], - ARRAY[]::text[] as allowed_lists -- TODO: allowed_lists not yet implemented in denormalized schema + -- Reconstruct list URIs from natural keys + COALESCE( + ( + SELECT array_agg('at://' || a.did || '/app.bsky.graph.list/' || list_rkey) + FROM unnest( + COALESCE(p.threadgate_allowed_list_actor_ids, ARRAY[]::integer[]), + COALESCE(p.threadgate_allowed_list_rkeys, ARRAY[]::text[]) + ) AS t(list_actor_id, list_rkey) + LEFT JOIN actors a ON a.id = list_actor_id + WHERE list_actor_id IS NOT NULL + ), + ARRAY[]::text[] + ) as allowed_lists FROM post_parts pp INNER JOIN actors a ON a.did = pp.did INNER JOIN posts p ON p.actor_id = a.id AND p.rkey = tid_to_i64(pp.rkey) @@ -118,6 +131,35 @@ pub async fn threadgate_upsert( (vec![], vec![]) }; + // Extract allowed list URIs from rules and resolve to natural keys + let (allowed_list_actor_ids, allowed_list_rkeys): (Vec, Vec) = if let Some(ref allow_rules) = rec.allow { + let mut list_uris = Vec::new(); + for rule in allow_rules { + if let ThreadgateRule::List { list } = rule { + list_uris.push(list.as_str()); + } + } + + if !list_uris.is_empty() { + // Resolve list URIs to (actor_id, rkey) pairs using bulk_resolve + let resolved = crate::db::bulk_resolve::resolve_list_uris_bulk(conn, &list_uris).await?; + + let mut actor_ids = Vec::new(); + let mut rkeys = Vec::new(); + for uri in &list_uris { + if let Some(&(aid, ref rk)) = resolved.get(*uri) { + actor_ids.push(aid); + rkeys.push(rk.clone()); + } + } + (actor_ids, rkeys) + } else { + (vec![], vec![]) + } + } else { + (vec![], vec![]) + }; + // Update posts table with denormalized threadgate data conn.execute( include_str!("../../sql/threadgate_upsert_denormalized.sql"), @@ -127,6 +169,8 @@ pub async fn threadgate_upsert( &allow, &hidden_actor_ids, &hidden_rkeys, + &allowed_list_actor_ids, + &allowed_list_rkeys, ], ) .await @@ -146,7 +190,9 @@ pub async fn threadgate_delete( "UPDATE posts SET threadgate_allow = NULL, threadgate_hidden_actor_ids = NULL, - threadgate_hidden_rkeys = NULL + threadgate_hidden_rkeys = NULL, + threadgate_allowed_list_actor_ids = NULL, + threadgate_allowed_list_rkeys = NULL WHERE actor_id = $1 AND rkey = $2", &[&actor_id, &rkey], diff --git a/consumer/src/db/sql/bookmarks_upsert.sql b/consumer/src/db/sql/bookmarks_upsert.sql deleted file mode 100644 index f4ee6f98..00000000 --- a/consumer/src/db/sql/bookmarks_upsert.sql +++ /dev/null @@ -1,7 +0,0 @@ --- Insert bookmark for a post using natural keys --- Parameters: $1=actor_id, $2=rkey, $3=post_actor_id, $4=post_rkey --- NOTE: URI parsing and ID resolution now done in Rust -INSERT INTO bookmarks (actor_id, rkey, post_actor_id, post_rkey) -VALUES ($1, $2, $3, $4) --- Bookmarks are immutable - have unique (actor_id, rkey) key -ON CONFLICT (actor_id, rkey) DO NOTHING \ No newline at end of file diff --git a/consumer/src/db/sql/postgate_upsert.sql b/consumer/src/db/sql/postgate_upsert.sql deleted file mode 100644 index a9a943ee..00000000 --- a/consumer/src/db/sql/postgate_upsert.sql +++ /dev/null @@ -1,41 +0,0 @@ --- Optimized insert/update postgate --- Postgates control embedding permissions for posts --- Parameters: $1=actor_id, $2=rkey, $3=cid, $4=post_rkey, $5=rules --- $6=detached_actor_ids[], $7=detached_rkeys[] --- Note: created_at is derived from TID rkey --- OPTIMIZATION: URI parsing and actor resolution done in Rust, not in SQL -WITH upserted_postgate AS ( - INSERT INTO postgates (actor_id, rkey, cid, post_actor_id, post_rkey, rules) - VALUES ( - $1, -- actor_id - $2, -- rkey (postgate's own rkey) - $3::bytea, -- cid (embedded) - $1, -- post_actor_id (same actor - postgates reference own posts) - $4, -- post_rkey (parsed from URI in Rust) - $5::text[]::postgate_rule[] -- rules (cast array to enum array) - ) - ON CONFLICT (actor_id, rkey) DO UPDATE SET - cid=EXCLUDED.cid, - post_actor_id=EXCLUDED.post_actor_id, - post_rkey=EXCLUDED.post_rkey, - rules=EXCLUDED.rules - RETURNING post_actor_id, post_rkey -), --- Delete old detached posts before inserting new ones -deleted_detached AS ( - DELETE FROM postgate_detached - WHERE (post_actor_id, post_rkey) = (SELECT post_actor_id, post_rkey FROM upserted_postgate) -), --- Insert detached posts using pre-resolved actor IDs and rkeys (no JOINs!) -detached_inserts AS ( - INSERT INTO postgate_detached (post_actor_id, post_rkey, detached_post_actor_id, detached_post_rkey) - SELECT - (SELECT post_actor_id FROM upserted_postgate), - (SELECT post_rkey FROM upserted_postgate), - actor_id, - rkey - FROM UNNEST($6::int[], $7::bigint[]) AS t(actor_id, rkey) - WHERE $6 IS NOT NULL - ON CONFLICT DO NOTHING -) -SELECT 1 as result diff --git a/consumer/src/db/sql/threadgate_upsert.sql b/consumer/src/db/sql/threadgate_upsert.sql deleted file mode 100644 index 171292c2..00000000 --- a/consumer/src/db/sql/threadgate_upsert.sql +++ /dev/null @@ -1,59 +0,0 @@ --- Optimized insert/update threadgate --- Threadgates control reply permissions for posts --- Parameters: $1=actor_id, $2=rkey, $3=cid, $4=post_rkey, $5=allow, $6=record(unused) --- $7=hidden_actor_ids[], $8=hidden_rkeys[], $9=list_ids[] --- Note: created_at is derived from TID rkey --- OPTIMIZATION: URI parsing and actor resolution done in Rust, not in SQL -WITH unused_params AS ( - SELECT $6::jsonb as record -- Type hint for unused parameter -), -upserted_threadgate AS ( - INSERT INTO threadgates (actor_id, rkey, cid, post_actor_id, post_rkey, allow) - VALUES ( - $1, -- actor_id - $2, -- rkey (threadgate's own rkey) - $3, -- cid (bytea - stored as digest) - $1, -- post_actor_id (same actor - threadgates reference own posts) - $4, -- post_rkey (parsed from URI in Rust) - $5::text[]::threadgate_rule[] -- allow (cast to enum array) - ) - ON CONFLICT (actor_id, rkey) DO UPDATE SET - cid=EXCLUDED.cid, - post_actor_id=EXCLUDED.post_actor_id, - post_rkey=EXCLUDED.post_rkey, - allow=EXCLUDED.allow - RETURNING post_actor_id, post_rkey -), --- Delete old hidden replies and allowed lists before inserting new ones -deleted_hidden AS ( - DELETE FROM threadgate_hidden_replies - WHERE (post_actor_id, post_rkey) = (SELECT post_actor_id, post_rkey FROM upserted_threadgate) -), -deleted_lists AS ( - DELETE FROM threadgate_allowed_lists - WHERE (post_actor_id, post_rkey) = (SELECT post_actor_id, post_rkey FROM upserted_threadgate) -), --- Insert hidden replies using pre-resolved actor IDs and rkeys (no JOINs!) -hidden_inserts AS ( - INSERT INTO threadgate_hidden_replies (post_actor_id, post_rkey, hidden_post_actor_id, hidden_post_rkey) - SELECT - (SELECT post_actor_id FROM upserted_threadgate), - (SELECT post_rkey FROM upserted_threadgate), - actor_id, - rkey - FROM UNNEST($7::int[], $8::bigint[]) AS t(actor_id, rkey) - WHERE $7 IS NOT NULL - ON CONFLICT DO NOTHING -), --- Insert allowed lists using pre-resolved list IDs (no JOINs!) -list_inserts AS ( - INSERT INTO threadgate_allowed_lists (post_actor_id, post_rkey, list_id) - SELECT - (SELECT post_actor_id FROM upserted_threadgate), - (SELECT post_rkey FROM upserted_threadgate), - list_id - FROM UNNEST($9::bigint[]) AS list_id - WHERE $9 IS NOT NULL - ON CONFLICT DO NOTHING -) -SELECT 0 as result diff --git a/consumer/src/db/sql/threadgate_upsert_denormalized.sql b/consumer/src/db/sql/threadgate_upsert_denormalized.sql index 44cc9e11..bd0a805d 100644 --- a/consumer/src/db/sql/threadgate_upsert_denormalized.sql +++ b/consumer/src/db/sql/threadgate_upsert_denormalized.sql @@ -1,6 +1,7 @@ -- Denormalized threadgate upsert: Update posts table directly -- Threadgates control reply permissions - now stored in posts table -- Parameters: $1=actor_id, $2=post_rkey, $3=allow[], $4=hidden_actor_ids[], $5=hidden_rkeys[] +-- $6=allowed_list_actor_ids[], $7=allowed_list_rkeys[] -- Note: We update the post record ($1, $2) with threadgate data -- The threadgate's own (actor_id, rkey, cid) are no longer stored separately @@ -8,6 +9,8 @@ UPDATE posts SET threadgate_allow = $3::text[]::threadgate_rule[], threadgate_hidden_actor_ids = $4::integer[], - threadgate_hidden_rkeys = $5::bigint[] + threadgate_hidden_rkeys = $5::bigint[], + threadgate_allowed_list_actor_ids = $6::integer[], + threadgate_allowed_list_rkeys = $7::text[] WHERE actor_id = $1 AND rkey = $2 diff --git a/migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/down.sql b/migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/down.sql new file mode 100644 index 00000000..13a1e8d4 --- /dev/null +++ b/migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/down.sql @@ -0,0 +1,5 @@ +-- Revert threadgate allowed lists arrays + +ALTER TABLE posts + DROP COLUMN threadgate_allowed_list_actor_ids, + DROP COLUMN threadgate_allowed_list_rkeys; diff --git a/migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/up.sql b/migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/up.sql new file mode 100644 index 00000000..f53e86dc --- /dev/null +++ b/migrations/2025-12-08-205247_add_threadgate_allowed_lists_arrays/up.sql @@ -0,0 +1,10 @@ +-- Add threadgate allowed lists arrays to posts table +-- These store the natural keys (actor_id, rkey) of lists that are allowed to reply +-- when a threadgate has list-based reply rules + +ALTER TABLE posts + ADD COLUMN threadgate_allowed_list_actor_ids integer[], + ADD COLUMN threadgate_allowed_list_rkeys text[]; + +COMMENT ON COLUMN posts.threadgate_allowed_list_actor_ids IS 'Actor IDs of lists allowed to reply (parallel to threadgate_allowed_list_rkeys)'; +COMMENT ON COLUMN posts.threadgate_allowed_list_rkeys IS 'Rkeys of lists allowed to reply (parallel to threadgate_allowed_list_actor_ids)'; diff --git a/parakeet-db/src/schema.rs b/parakeet-db/src/schema.rs index c7529581..49bc88d0 100644 --- a/parakeet-db/src/schema.rs +++ b/parakeet-db/src/schema.rs @@ -495,6 +495,8 @@ diesel::table! { repost_count -> Int4, reply_count -> Int4, quote_count -> Int4, + threadgate_allowed_list_actor_ids -> Nullable>>, + threadgate_allowed_list_rkeys -> Nullable>>, } } diff --git a/parakeet/src/db/graph.rs b/parakeet/src/db/graph.rs index dd611a26..8c4da3a3 100644 --- a/parakeet/src/db/graph.rs +++ b/parakeet/src/db/graph.rs @@ -130,7 +130,7 @@ pub async fn get_actor_followers( } diesel::sql_query( - "SELECT (f).rkey, (f).actor_id as follower_actor_id + "SELECT (f).rkey, (f).subject_actor_id as follower_actor_id FROM actors, unnest(followers) AS f WHERE id = $1 AND ($2::bigint IS NULL OR (f).rkey < $2) @@ -244,9 +244,9 @@ pub async fn get_followed_by_batch( } diesel::sql_query( - "SELECT (f).actor_id as follower_actor_id, (f).rkey::text as rkey + "SELECT (f).subject_actor_id as follower_actor_id, (f).rkey::text as rkey FROM actors, unnest(followers) AS f - WHERE id = $1 AND (f).actor_id = ANY($2)" + WHERE id = $1 AND (f).subject_actor_id = ANY($2)" ) .bind::(actor_id) .bind::, _>(other_actor_ids) @@ -282,10 +282,10 @@ pub async fn get_mutual_followers( // Get target's followers that are in viewer's following list diesel::sql_query( - "SELECT tid_timestamp((f1).rkey) as created_at, (f1).actor_id as follower_actor_id + "SELECT tid_timestamp((f1).rkey) as created_at, (f1).subject_actor_id as follower_actor_id FROM actors target, unnest(target.followers) AS f1 WHERE target.id = $1 - AND (f1).actor_id IN ( + AND (f1).subject_actor_id IN ( SELECT (f2).subject_actor_id FROM actors viewer, unnest(viewer.following) AS f2 WHERE viewer.id = $2 @@ -356,10 +356,12 @@ pub async fn get_actor_lists( /// Get list items with cursor pagination /// +/// Uses natural keys (list_owner_actor_id, list_rkey) instead of list_id. /// Returns list of (created_at, at_uri, subject_did) tuples pub async fn get_list_items( conn: &mut AsyncPgConnection, - list_id: i64, + list_owner_actor_id: i32, + list_rkey: &str, cursor_timestamp: Option<&chrono::DateTime>, limit: u8, ) -> QueryResult, String, String)>> { @@ -381,12 +383,14 @@ pub async fn get_list_items( FROM list_items li INNER JOIN actors a ON li.actor_id = a.id INNER JOIN actors subject ON li.subject_actor_id = subject.id - WHERE li.list_id = $1 - AND ($2::timestamptz IS NULL OR tid_timestamp(li.rkey) < $2) + WHERE li.list_owner_actor_id = $1 + AND li.list_rkey = $2 + AND ($3::timestamptz IS NULL OR tid_timestamp(li.rkey) < $3) ORDER BY li.rkey DESC - LIMIT $3" + LIMIT $4" ) - .bind::(list_id) + .bind::(list_owner_actor_id) + .bind::(list_rkey) .bind::, _>(cursor_timestamp) .bind::(i64::from(limit)) .load::(conn) @@ -400,6 +404,7 @@ pub async fn get_list_items( /// Get muted lists for a user with cursor pagination /// +/// DENORMALIZED: Reads from actors.list_mutes[] array /// Returns list of (created_at, list_uri) tuples pub async fn get_user_list_mutes( conn: &mut AsyncPgConnection, @@ -417,13 +422,13 @@ pub async fn get_user_list_mutes( // Use .bind() for cursor parameter to prevent SQL injection diesel::sql_query( - "SELECT lm.created_at, 'at://' || a.did || '/app.bsky.graph.list/' || l.rkey::text as list_uri - FROM list_mutes lm - INNER JOIN lists l ON lm.list_id = l.id - INNER JOIN actors a ON l.actor_id = a.id - WHERE lm.actor_id = $1 - AND ($2::timestamptz IS NULL OR lm.created_at < $2) - ORDER BY lm.created_at DESC + "SELECT (lm).created_at as created_at, + 'at://' || a.did || '/app.bsky.graph.list/' || (lm).list_rkey as list_uri + FROM actors user_actor, unnest(user_actor.list_mutes) AS lm + INNER JOIN actors a ON (lm).list_actor_id = a.id + WHERE user_actor.id = $1 + AND ($2::timestamptz IS NULL OR (lm).created_at < $2) + ORDER BY (lm).created_at DESC LIMIT $3" ) .bind::(actor_id) @@ -440,6 +445,7 @@ pub async fn get_user_list_mutes( /// Get blocked lists for a user with cursor pagination /// +/// DENORMALIZED: Reads from actors.list_blocks[] array /// Returns list of (created_at, list_uri) tuples pub async fn get_user_list_blocks( conn: &mut AsyncPgConnection, @@ -457,13 +463,13 @@ pub async fn get_user_list_blocks( // Use .bind() for cursor parameter to prevent SQL injection diesel::sql_query( - "SELECT tid_timestamp(lb.rkey) as created_at, 'at://' || list_a.did || '/app.bsky.graph.list/' || l.rkey::text as list_uri - FROM list_blocks lb - INNER JOIN lists l ON lb.list_id = l.id - INNER JOIN actors list_a ON l.actor_id = list_a.id - WHERE lb.actor_id = $1 - AND ($2::timestamptz IS NULL OR tid_timestamp(lb.rkey) < $2) - ORDER BY lb.rkey DESC + "SELECT tid_timestamp((lb).rkey) as created_at, + 'at://' || list_a.did || '/app.bsky.graph.list/' || (lb).list_rkey as list_uri + FROM actors user_actor, unnest(user_actor.list_blocks) AS lb + INNER JOIN actors list_a ON (lb).list_actor_id = list_a.id + WHERE user_actor.id = $1 + AND ($2::timestamptz IS NULL OR tid_timestamp((lb).rkey) < $2) + ORDER BY (lb).rkey DESC LIMIT $3" ) .bind::(actor_id) diff --git a/parakeet/src/db/notification_records.rs b/parakeet/src/db/notification_records.rs index 31dbf378..26c362e8 100644 --- a/parakeet/src/db/notification_records.rs +++ b/parakeet/src/db/notification_records.rs @@ -228,8 +228,10 @@ pub async fn get_follow_records_batch( SELECT k.actor_id, k.rkey, a2.did as subject, tid_timestamp(k.rkey) as created_at FROM keys k - INNER JOIN follows f ON k.actor_id = f.actor_id AND k.rkey = f.rkey - INNER JOIN actors a2 ON f.subject_actor_id = a2.id", + INNER JOIN actors a ON k.actor_id = a.id + CROSS JOIN unnest(a.following) AS f + INNER JOIN actors a2 ON (f).subject_actor_id = a2.id + WHERE (f).rkey = k.rkey", ) .bind::, _>(actor_ids) .bind::, _>(rkeys) @@ -262,10 +264,10 @@ pub async fn get_follow_record( diesel::sql_query( "SELECT a2.did as subject, - tid_timestamp(f.rkey) as created_at - FROM follows f - INNER JOIN actors a2 ON f.subject_actor_id = a2.id - WHERE f.actor_id = $1 AND f.rkey = $2", + tid_timestamp((f).rkey) as created_at + FROM actors a, unnest(a.following) AS f + INNER JOIN actors a2 ON (f).subject_actor_id = a2.id + WHERE a.id = $1 AND (f).rkey = $2", ) .bind::(actor_id) .bind::(rkey_bigint) diff --git a/parakeet/src/db/suggestions.rs b/parakeet/src/db/suggestions.rs index 15ee3959..fc054627 100644 --- a/parakeet/src/db/suggestions.rs +++ b/parakeet/src/db/suggestions.rs @@ -194,7 +194,7 @@ pub async fn get_collaborative_filter_suggestions( diesel::sql_query( "WITH actor_followers AS ( - SELECT DISTINCT (f).actor_id as follower_id + SELECT DISTINCT (f).subject_actor_id as follower_id FROM actors, unnest(followers) AS f WHERE id = $1 LIMIT 1000 diff --git a/parakeet/src/loaders/labeler.rs b/parakeet/src/loaders/labeler.rs index 3b6ff2f3..a0584ffc 100644 --- a/parakeet/src/loaders/labeler.rs +++ b/parakeet/src/loaders/labeler.rs @@ -19,50 +19,69 @@ pub fn build_labeler_records_query(actor_ids_str: &str) -> String { ) } -/// Build SQL query for loading labels by URI +/// Build SQL query for loading labels by URI (DENORMALIZED) +/// +/// Labels are now stored as actor_label[] arrays on the actors table. +/// Each actor has a labels array containing labels they've applied. /// /// This function is public for testing purposes. pub fn build_labels_query() -> &'static str { - "SELECT - l.labeler_actor_id, - l.label, - l.uri, - l.self_label, - l.cid, - l.negated, - l.expires, - l.sig, - l.created_at, - a.did as labeler - FROM labels l - INNER JOIN actors a ON l.labeler_actor_id = a.id - WHERE l.uri = $1 - AND l.negated = false - AND (l.self_label = true OR a.did = ANY($2)) - ORDER BY l.created_at" + "WITH target_actor AS ( + SELECT id FROM actors + WHERE did = SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) + ) + SELECT + (lbl).labeler_actor_id, + (lbl).label as label, + $1::text as uri, + false as self_label, + NULL::bytea as cid, + (lbl).negated, + (lbl).expires, + NULL::bytea as sig, + (lbl).created_at, + labeler.did as labeler + FROM target_actor ta + CROSS JOIN actors labeler_subjects + CROSS JOIN unnest(labeler_subjects.labels) AS lbl + INNER JOIN actors labeler ON (lbl).labeler_actor_id = labeler.id + WHERE labeler_subjects.id = ta.id + AND (lbl).negated = false + AND labeler.did = ANY($2) + ORDER BY (lbl).created_at" } -/// Build SQL query for batch loading labels by multiple URIs +/// Build SQL query for batch loading labels by multiple URIs (DENORMALIZED) +/// +/// Labels are now stored as actor_label[] arrays on the actors table. +/// This query unnests labels from multiple target actors. /// /// This function is public for testing purposes. pub fn build_labels_many_query() -> &'static str { - "SELECT - l.labeler_actor_id, - l.label, - l.uri, - l.self_label, - l.cid, - l.negated, - l.expires, - l.sig, - l.created_at, - a.did as labeler - FROM labels l - INNER JOIN actors a ON l.labeler_actor_id = a.id - WHERE l.uri = ANY($1) - AND l.negated = false - AND (l.self_label = true OR a.did = ANY($2)) - ORDER BY l.created_at" + "WITH target_actors AS ( + SELECT unnest($1::text[]) as uri, + a.id as actor_id + FROM unnest($1::text[]) uri_val + INNER JOIN actors a ON a.did = SPLIT_PART(SUBSTRING(uri_val FROM 6), '/', 1) + ) + SELECT + (lbl).labeler_actor_id, + (lbl).label as label, + ta.uri as uri, + false as self_label, + NULL::bytea as cid, + (lbl).negated, + (lbl).expires, + NULL::bytea as sig, + (lbl).created_at, + labeler.did as labeler + FROM target_actors ta + INNER JOIN actors labeler_subjects ON labeler_subjects.id = ta.actor_id + CROSS JOIN unnest(labeler_subjects.labels) AS lbl + INNER JOIN actors labeler ON (lbl).labeler_actor_id = labeler.id + WHERE (lbl).negated = false + AND labeler.did = ANY($2) + ORDER BY (lbl).created_at" } // Enriched Labeler with reconstructed fields from Actor diff --git a/parakeet/src/loaders/profile.rs b/parakeet/src/loaders/profile.rs index 757aeddc..a971bd3f 100644 --- a/parakeet/src/loaders/profile.rs +++ b/parakeet/src/loaders/profile.rs @@ -60,10 +60,9 @@ pub fn build_profiles_batch_query() -> &'static str { a.profile_pronouns as pronouns, a.profile_website as website, a.chat_allow_incoming as allow_incoming, - l.actor_id as labeler_actor_id, + CASE WHEN a.labeler_cid IS NOT NULL THEN a.id ELSE NULL END as labeler_actor_id, a.notif_decl_allow_subscriptions as allow_subscriptions FROM actors a - LEFT JOIN labelers l ON l.actor_id = a.id WHERE a.did = ANY($1) AND a.status = 'active'::actor_status" } @@ -420,10 +419,9 @@ pub fn build_profiles_by_id_batch_query() -> &'static str { a.profile_pronouns as pronouns, a.profile_website as website, a.chat_allow_incoming as allow_incoming, - l.actor_id as labeler_actor_id, + CASE WHEN a.labeler_cid IS NOT NULL THEN a.id ELSE NULL END as labeler_actor_id, a.notif_decl_allow_subscriptions as allow_subscriptions FROM actors a - LEFT JOIN labelers l ON l.actor_id = a.id WHERE a.id = ANY($1)" } diff --git a/parakeet/src/sql/list_states.sql b/parakeet/src/sql/list_states.sql index 712a612a..1c704ecc 100644 --- a/parakeet/src/sql/list_states.sql +++ b/parakeet/src/sql/list_states.sql @@ -1,20 +1,47 @@ --- Get list states (blocks, mutes) for a viewer +-- Get list states (blocks, mutes) for a viewer (DENORMALIZED) -- $1 = viewer DID -- $2 = array of list dids -- $3 = array of list rkeys (text - lists use text rkeys) +with viewer as ( + select id from actors where did = $1 +), +lookup_lists as ( + select + a.did as list_did, + a.id as list_actor_id, + l.rkey as list_rkey + from lists l + inner join actors a on l.actor_id = a.id + inner join unnest($2::text[], $3::text[]) AS lookup(lookup_did, lookup_rkey) + ON a.did = lookup.lookup_did AND l.rkey = lookup.lookup_rkey +), +list_blocks_unnested as ( + select + ll.list_did, + ll.list_rkey, + va.did as block_actor_did, + (lb).rkey as block_rkey + from viewer v + inner join actors va on va.id = v.id + cross join unnest(va.list_blocks) as lb + inner join lookup_lists ll on (lb).list_actor_id = ll.list_actor_id and (lb).list_rkey = ll.list_rkey +), +list_mutes_unnested as ( + select + ll.list_did, + ll.list_rkey + from viewer v + inner join actors va on va.id = v.id + cross join unnest(va.list_mutes) as lm + inner join lookup_lists ll on (lm).list_actor_id = ll.list_actor_id and (lm).list_rkey = ll.list_rkey +) select - a.did as list_did, - l.rkey as list_rkey, - block_actor.did as block_actor_did, - lb.rkey as block_rkey, - lm.actor_id is not null as muted -from lists l -inner join actors a on l.actor_id = a.id -inner join unnest($2::text[], $3::text[]) AS lookup(lookup_did, lookup_rkey) - ON a.did = lookup.lookup_did AND l.rkey = lookup.lookup_rkey -left join list_blocks lb on lb.list_id = l.id - and lb.actor_id = (select id from actors where did = $1) -left join actors block_actor on lb.actor_id = block_actor.id -left join list_mutes lm on lm.list_id = l.id - and lm.actor_id = (select id from actors where did = $1) -where (lm.actor_id is not null or lb.id is not null) \ No newline at end of file + ll.list_did, + ll.list_rkey, + lb.block_actor_did, + lb.block_rkey, + lm.list_did is not null as muted +from lookup_lists ll +left join list_blocks_unnested lb on ll.list_did = lb.list_did and ll.list_rkey = lb.list_rkey +left join list_mutes_unnested lm on ll.list_did = lm.list_did and ll.list_rkey = lm.list_rkey +where (lb.block_actor_did is not null or lm.list_did is not null) \ No newline at end of file diff --git a/parakeet/src/sql/post_state.sql b/parakeet/src/sql/post_state.sql index bf2e9947..d538bdbf 100644 --- a/parakeet/src/sql/post_state.sql +++ b/parakeet/src/sql/post_state.sql @@ -67,14 +67,17 @@ viewer_reposts AS ( AND r.post_rkey <= $5 ), -- Query viewer's bookmarks for these specific posts only --- Uses rkey range to enable chunk exclusion on bookmarks +-- DENORMALIZED: Read from actors.bookmarks[] array viewer_bookmarks AS ( - SELECT b.post_actor_id, b.post_rkey - FROM bookmarks b - INNER JOIN target_posts tp ON b.post_actor_id = tp.actor_id AND b.post_rkey = tp.rkey - WHERE b.actor_id = $1 - AND b.post_rkey >= $4 - AND b.post_rkey <= $5 + SELECT + (b).post_actor_id as post_actor_id, + (b).post_rkey as post_rkey + FROM actors a + CROSS JOIN LATERAL unnest(COALESCE(a.bookmarks, ARRAY[]::bookmark_record[])) AS b + INNER JOIN target_posts tp ON (b).post_actor_id = tp.actor_id AND (b).post_rkey = tp.rkey + WHERE a.id = $1 + AND (b).post_rkey >= $4 + AND (b).post_rkey <= $5 ) SELECT tp.actor_id, diff --git a/parakeet/src/sql/profile_state.sql b/parakeet/src/sql/profile_state.sql index c4d74995..372142c1 100644 --- a/parakeet/src/sql/profile_state.sql +++ b/parakeet/src/sql/profile_state.sql @@ -1,5 +1,6 @@ --- Profile state query using normalized schema (no profile_states table) --- Reconstructs profile state from follows, blocks, mutes, and list_blocks/list_mutes +-- Profile state query using denormalized schema +-- Reconstructs profile state from actors.following, actors.blocks, actors.mutes arrays +-- and actors.list_blocks/list_mutes arrays -- Parameters: $1 = viewer DID, $2 = array of subject DIDs with viewer as ( @@ -9,68 +10,100 @@ subjects as ( select id, did from actors where did = any($2) ), -- Following relationship (return rkey bigint for encoding in Rust) +-- Read from viewer's following[] array following_cte as ( - select s.did as subject, f.rkey as following - from follows f - inner join subjects s on f.subject_actor_id = s.id - where f.actor_id = (select id from viewer) + select s.did as subject, (f).rkey as following + from viewer v + cross join lateral unnest(COALESCE( + (select following from actors where id = v.id), + ARRAY[]::follow_record[] + )) as f + inner join subjects s on (f).subject_actor_id = s.id ), -- Followed relationship (return rkey bigint for encoding in Rust) +-- Read from subject's following[] array (subject follows viewer) followed_cte as ( - select s.did as subject, f.rkey as followed - from follows f - inner join viewer v on f.subject_actor_id = v.id - inner join subjects s on f.actor_id = s.id + select s.did as subject, (f).rkey as followed + from subjects s + cross join lateral unnest(COALESCE( + (select following from actors where id = s.id), + ARRAY[]::follow_record[] + )) as f + inner join viewer v on (f).subject_actor_id = v.id ), -- Blocking relationship (return rkey bigint for encoding in Rust) +-- Read from viewer's blocks[] array blocking_cte as ( - select s.did as subject, b.rkey as blocking - from blocks b - inner join subjects s on b.subject_actor_id = s.id - where b.actor_id = (select id from viewer) + select s.did as subject, (b).rkey as blocking + from viewer v + cross join lateral unnest(COALESCE( + (select blocks from actors where id = v.id), + ARRAY[]::block_record[] + )) as b + inner join subjects s on (b).subject_actor_id = s.id ), --- Blocked relationship +-- Blocked relationship (subject has blocked viewer) +-- Read from subject's blocks[] array blocked_cte as ( select s.did as subject, true as blocked - from blocks b - inner join viewer v on b.subject_actor_id = v.id - inner join subjects s on b.actor_id = s.id + from subjects s + cross join lateral unnest(COALESCE( + (select blocks from actors where id = s.id), + ARRAY[]::block_record[] + )) as b + inner join viewer v on (b).subject_actor_id = v.id ), -- Muting relationship +-- Read from viewer's mutes[] array muting_cte as ( select s.did as subject, true as muting - from mutes m - inner join subjects s on m.subject_actor_id = s.id - where m.actor_id = (select id from viewer) + from viewer v + cross join lateral unnest(COALESCE( + (select mutes from actors where id = v.id), + ARRAY[]::mute_record[] + )) as m + inner join subjects s on (m).subject_actor_id = s.id ), -- List blocks (viewer has blocked subject via list) +-- Viewer has blocked lists → check if subject is in those lists vlb as ( - select s.did as subject, list_owner.did as list_owner_did, l.rkey as list_rkey - from list_blocks lb - inner join list_mutes lm on lm.list_id = lb.list_id - inner join lists l on l.id = lb.list_id - inner join actors list_owner on l.actor_id = list_owner.id - inner join subjects s on s.id = lb.actor_id - where lm.actor_id = (select id from viewer) + select s.did as subject, list_owner.did as list_owner_did, (lb).list_rkey as list_rkey + from viewer v + cross join lateral unnest(COALESCE( + (select list_blocks from actors where id = v.id), + ARRAY[]::list_block_record[] + )) as lb + inner join actors list_owner on list_owner.id = (lb).list_actor_id + inner join list_items li on li.list_owner_actor_id = (lb).list_actor_id + and li.list_rkey = (lb).list_rkey + inner join subjects s on s.id = li.subject_actor_id ), -- List blocks reverse (subject has blocked viewer via list) +-- Subject has blocked lists → check if viewer is in those lists vlb2 as ( select s.did as subject, true as blocked - from list_blocks lb - inner join list_mutes lm on lm.list_id = lb.list_id - inner join viewer v on v.id = lb.actor_id - inner join subjects s on s.id = lm.actor_id + from subjects s + cross join lateral unnest(COALESCE( + (select list_blocks from actors where id = s.id), + ARRAY[]::list_block_record[] + )) as lb + inner join list_items li on li.list_owner_actor_id = (lb).list_actor_id + and li.list_rkey = (lb).list_rkey + inner join viewer v on v.id = li.subject_actor_id ), -- List mutes (viewer has muted subject via list) +-- Viewer has muted lists → check if subject is in those lists vlm as ( - select a.did as subject, list_owner.did as list_owner_did, l.rkey as list_rkey - from list_items li - inner join list_mutes lm on lm.list_id = li.list_id - inner join lists l on l.id = li.list_id - inner join actors list_owner on l.actor_id = list_owner.id + select s.did as subject, list_owner.did as list_owner_did, (lm).list_rkey as list_rkey + from viewer v + cross join lateral unnest(COALESCE( + (select list_mutes from actors where id = v.id), + ARRAY[]::list_mute_record[] + )) as lm + inner join actors list_owner on list_owner.id = (lm).list_actor_id + inner join list_items li on li.list_owner_actor_id = (lm).list_actor_id + and li.list_rkey = (lm).list_rkey inner join subjects s on s.id = li.subject_actor_id - inner join actors a on a.id = s.id - where lm.actor_id = (select id from viewer) ) select $1 as did, diff --git a/parakeet/src/sql/profile_state_by_ids.sql b/parakeet/src/sql/profile_state_by_ids.sql index cfc1ccfc..95b642c4 100644 --- a/parakeet/src/sql/profile_state_by_ids.sql +++ b/parakeet/src/sql/profile_state_by_ids.sql @@ -1,5 +1,6 @@ -- Profile state query using actor_ids instead of DIDs (optimized for IdCache usage) --- Reconstructs profile state from follows, blocks, mutes, and list_blocks/list_mutes +-- Reconstructs profile state from actors.following, actors.blocks, actors.mutes arrays +-- and actors.list_blocks/list_mutes arrays -- Parameters: $1 = viewer actor_id (integer), $2 = array of subject actor_ids (integer[]) -- -- This is an optimized version that avoids decompressing the actors table by accepting @@ -12,65 +13,98 @@ subjects as ( select unnest($2::integer[]) as id ), -- Following relationship (return rkey bigint for encoding in Rust) +-- Read from viewer's following[] array following_cte as ( - select s.id as subject_id, f.rkey as following - from follows f - inner join subjects s on f.subject_actor_id = s.id - where f.actor_id = (select id from viewer) + select s.id as subject_id, (f).rkey as following + from viewer v + cross join lateral unnest(COALESCE( + (select following from actors where id = v.id), + ARRAY[]::follow_record[] + )) as f + inner join subjects s on (f).subject_actor_id = s.id ), -- Followed relationship (return rkey bigint for encoding in Rust) +-- Read from subject's following[] array (subject follows viewer) followed_cte as ( - select s.id as subject_id, f.rkey as followed - from follows f - inner join viewer v on f.subject_actor_id = v.id - inner join subjects s on f.actor_id = s.id + select s.id as subject_id, (f).rkey as followed + from subjects s + cross join lateral unnest(COALESCE( + (select following from actors where id = s.id), + ARRAY[]::follow_record[] + )) as f + inner join viewer v on (f).subject_actor_id = v.id ), -- Blocking relationship (return rkey bigint for encoding in Rust) +-- Read from viewer's blocks[] array blocking_cte as ( - select s.id as subject_id, b.rkey as blocking - from blocks b - inner join subjects s on b.subject_actor_id = s.id - where b.actor_id = (select id from viewer) + select s.id as subject_id, (b).rkey as blocking + from viewer v + cross join lateral unnest(COALESCE( + (select blocks from actors where id = v.id), + ARRAY[]::block_record[] + )) as b + inner join subjects s on (b).subject_actor_id = s.id ), --- Blocked relationship +-- Blocked relationship (subject has blocked viewer) +-- Read from subject's blocks[] array blocked_cte as ( select s.id as subject_id, true as blocked - from blocks b - inner join viewer v on b.subject_actor_id = v.id - inner join subjects s on b.actor_id = s.id + from subjects s + cross join lateral unnest(COALESCE( + (select blocks from actors where id = s.id), + ARRAY[]::block_record[] + )) as b + inner join viewer v on (b).subject_actor_id = v.id ), -- Muting relationship +-- Read from viewer's mutes[] array muting_cte as ( select s.id as subject_id, true as muting - from mutes m - inner join subjects s on m.subject_actor_id = s.id - where m.actor_id = (select id from viewer) + from viewer v + cross join lateral unnest(COALESCE( + (select mutes from actors where id = v.id), + ARRAY[]::mute_record[] + )) as m + inner join subjects s on (m).subject_actor_id = s.id ), -- List blocks (viewer has blocked subject via list) +-- Viewer has blocked lists → check if subject is in those lists vlb as ( - select s.id as subject_id, l.actor_id as list_owner_actor_id, l.rkey as list_rkey - from list_blocks lb - inner join list_mutes lm on lm.list_id = lb.list_id - inner join lists l on l.id = lb.list_id - inner join subjects s on s.id = lb.actor_id - where lm.actor_id = (select id from viewer) + select s.id as subject_id, (lb).list_actor_id as list_owner_actor_id, (lb).list_rkey as list_rkey + from viewer v + cross join lateral unnest(COALESCE( + (select list_blocks from actors where id = v.id), + ARRAY[]::list_block_record[] + )) as lb + inner join list_items li on li.list_owner_actor_id = (lb).list_actor_id + and li.list_rkey = (lb).list_rkey + inner join subjects s on s.id = li.subject_actor_id ), -- List blocks reverse (subject has blocked viewer via list) +-- Subject has blocked lists → check if viewer is in those lists vlb2 as ( select s.id as subject_id, true as blocked - from list_blocks lb - inner join list_mutes lm on lm.list_id = lb.list_id - inner join viewer v on v.id = lb.actor_id - inner join subjects s on s.id = lm.actor_id + from subjects s + cross join lateral unnest(COALESCE( + (select list_blocks from actors where id = s.id), + ARRAY[]::list_block_record[] + )) as lb + inner join list_items li on li.list_owner_actor_id = (lb).list_actor_id + and li.list_rkey = (lb).list_rkey + inner join viewer v on v.id = li.subject_actor_id ), -- List mutes (viewer has muted subject via list) +-- Viewer has muted lists → check if subject is in those lists vlm as ( - select s.id as subject_id, l.actor_id as list_owner_actor_id, l.rkey as list_rkey - from list_items li - inner join list_mutes lm on lm.list_id = li.list_id - inner join lists l on l.id = li.list_id + select s.id as subject_id, (lm).list_actor_id as list_owner_actor_id, (lm).list_rkey as list_rkey + from viewer v + cross join lateral unnest(COALESCE( + (select list_mutes from actors where id = v.id), + ARRAY[]::list_mute_record[] + )) as lm + inner join list_items li on li.list_owner_actor_id = (lm).list_actor_id + and li.list_rkey = (lm).list_rkey inner join subjects s on s.id = li.subject_actor_id - where lm.actor_id = (select id from viewer) ) select s.id as subject_id, diff --git a/parakeet/src/xrpc/app_bsky/graph/lists.rs b/parakeet/src/xrpc/app_bsky/graph/lists.rs index 9166f9ec..083a0ebc 100644 --- a/parakeet/src/xrpc/app_bsky/graph/lists.rs +++ b/parakeet/src/xrpc/app_bsky/graph/lists.rs @@ -123,12 +123,9 @@ pub async fn get_list( list_did, ).await?; - // Resolve list_id - let list_id = crate::db::get_list_id_by_uri(&mut conn, list_owner_actor_id, rkey_str).await?; - - // Query list items + // Query list items using natural keys (actor_id, rkey) let cursor_value = datetime_cursor(query.cursor.as_ref()); - let results = crate::db::get_list_items(&mut conn, list_id, cursor_value.as_ref(), limit).await?; + let results = crate::db::get_list_items(&mut conn, list_owner_actor_id, rkey_str, cursor_value.as_ref(), limit).await?; let cursor = results .last() diff --git a/parakeet/tests/db_graph_test.rs b/parakeet/tests/db_graph_test.rs index 2d6eb984..2a67298c 100644 --- a/parakeet/tests/db_graph_test.rs +++ b/parakeet/tests/db_graph_test.rs @@ -305,7 +305,7 @@ async fn test_get_list_items() -> eyre::Result<()> { let pool = common::test_diesel_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; - let _result = parakeet::db::get_list_items(&mut conn, 99999, None, 50) + let _result = parakeet::db::get_list_items(&mut conn, 99999, "test_rkey", None, 50) .await .wrap_err("SQL syntax error in get_list_items")?; @@ -319,7 +319,7 @@ async fn test_get_list_items_with_cursor() -> eyre::Result<()> { let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let cursor = chrono::Utc::now(); - let _result = parakeet::db::get_list_items(&mut conn, 99999, Some(&cursor), 50) + let _result = parakeet::db::get_list_items(&mut conn, 99999, "test_rkey", Some(&cursor), 50) .await .wrap_err("SQL syntax error in get_list_items with cursor")?; diff --git a/parakeet/tests/db_uri_reconstruction_test.rs b/parakeet/tests/db_uri_reconstruction_test.rs deleted file mode 100644 index d9210e50..00000000 --- a/parakeet/tests/db_uri_reconstruction_test.rs +++ /dev/null @@ -1,76 +0,0 @@ -//! Integration tests for parakeet/src/db/uri_reconstruction.rs SQL operations -//! -//! These tests verify SQL validity for URI reconstruction queries. -//! The goal is to ensure queries execute without syntax errors, NOT to verify -//! business logic or result correctness. - -mod common; - -use eyre::WrapErr; - -// ============================================================================ -// URI RECONSTRUCTION QUERY TESTS -// ============================================================================ - -// NOTE: get_post_uris_by_ids tests removed - function no longer exists after natural key migration -// Posts now use natural keys (actor_id, rkey) instead of synthetic IDs - -#[tokio::test] -async fn test_get_feedgen_uris_by_ids() -> eyre::Result<()> { - common::ensure_test_db_ready().await; - let pool = common::test_diesel_pool(); - let mut conn = pool.get().await.wrap_err("Failed to get connection")?; - - let feedgen_ids = vec![1, 2, 3]; - - let _result = parakeet::db::get_feedgen_uris_by_ids(&mut conn, &feedgen_ids) - .await - .wrap_err("SQL syntax error in get_feedgen_uris_by_ids")?; - - Ok(()) -} - -#[tokio::test] -async fn test_get_list_uris_by_ids() -> eyre::Result<()> { - common::ensure_test_db_ready().await; - let pool = common::test_diesel_pool(); - let mut conn = pool.get().await.wrap_err("Failed to get connection")?; - - let list_ids = vec![1, 2, 3]; - - let _result = parakeet::db::get_list_uris_by_ids(&mut conn, &list_ids) - .await - .wrap_err("SQL syntax error in get_list_uris_by_ids")?; - - Ok(()) -} - -#[tokio::test] -async fn test_get_labeler_uris_by_ids() -> eyre::Result<()> { - common::ensure_test_db_ready().await; - let pool = common::test_diesel_pool(); - let mut conn = pool.get().await.wrap_err("Failed to get connection")?; - - let labeler_ids = vec![1, 2, 3]; - - let _result = parakeet::db::get_labeler_uris_by_ids(&mut conn, &labeler_ids) - .await - .wrap_err("SQL syntax error in get_labeler_uris_by_ids")?; - - Ok(()) -} - -#[tokio::test] -async fn test_get_starterpack_uris_by_ids() -> eyre::Result<()> { - common::ensure_test_db_ready().await; - let pool = common::test_diesel_pool(); - let mut conn = pool.get().await.wrap_err("Failed to get connection")?; - - let starterpack_ids = vec![1, 2, 3]; - - let _result = parakeet::db::get_starterpack_uris_by_ids(&mut conn, &starterpack_ids) - .await - .wrap_err("SQL syntax error in get_starterpack_uris_by_ids")?; - - Ok(()) -} diff --git a/parakeet/tests/sql/batch_loading_test.rs b/parakeet/tests/sql/batch_loading_test.rs index b2938814..d0568309 100644 --- a/parakeet/tests/sql/batch_loading_test.rs +++ b/parakeet/tests/sql/batch_loading_test.rs @@ -225,24 +225,26 @@ async fn test_batch_starterpacks_query() -> eyre::Result<()> { Ok(()) } -/// Test feed URI resolution query (from loaders/misc.rs) +/// Test feedgen batch loading query includes URI construction (from loaders/feed.rs) #[tokio::test] -async fn test_feed_uris_query() -> eyre::Result<()> { +async fn test_feedgen_batch_query() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_diesel_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Use the ACTUAL query builder from the source code - let query = parakeet::loaders::build_feed_uris_query(); - let feed_ids = vec![1_i64, 2_i64, 3_i64]; + // This tests that feedgens query properly constructs URIs + // Query expects URIs (not natural keys), so construct AT-URI format + let query = parakeet::loaders::build_feedgens_batch_query(); + let uris = vec!["at://did:plc:test/app.bsky.feed.generator/test".to_string()]; diesel_async::RunQueryDsl::execute( diesel::sql_query(query) - .bind::, _>(&feed_ids), + .bind::, _>(&uris), &mut conn, ) .await - .wrap_err("Feed URI resolution query failed")?; + .wrap_err("Feedgen batch loading query failed")?; Ok(()) } @@ -271,6 +273,7 @@ async fn test_labeler_records_query() -> eyre::Result<()> { } /// Test single label loading query (from loaders/labeler.rs) +/// Labels are now stored as actor_label[] arrays on actors table #[tokio::test] async fn test_labels_single_query() -> eyre::Result<()> { common::ensure_test_db_ready().await; diff --git a/parakeet/tests/sql/db_module_test.rs b/parakeet/tests/sql/db_module_test.rs index 251db766..bb2668ce 100644 --- a/parakeet/tests/sql/db_module_test.rs +++ b/parakeet/tests/sql/db_module_test.rs @@ -1553,7 +1553,8 @@ async fn test_get_list_items_empty() -> eyre::Result<()> { // Query items for non-existent list let result = db::get_list_items( &mut conn, - 99999, // non-existent list_id + 99999, // non-existent actor_id + "test_rkey", // list rkey None, 50 ).await; @@ -1576,7 +1577,8 @@ async fn test_get_list_items_with_cursor() -> eyre::Result<()> { // Query with cursor timestamp let result = db::get_list_items( &mut conn, - 1, + 1, // actor_id + "test_rkey", // list rkey Some(&cursor_time), 50 ).await; -- 2.51.2