diff --git a/consumer/src/db/sql/bookmarks_upsert.sql b/consumer/src/db/sql/bookmarks_upsert.sql index 3172ef03..5bce229f 100644 --- a/consumer/src/db/sql/bookmarks_upsert.sql +++ b/consumer/src/db/sql/bookmarks_upsert.sql @@ -5,6 +5,7 @@ WITH post_parts AS ( SELECT SPLIT_PART(SUBSTRING($3 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($3 FROM 6), '/', 2) as collection, SPLIT_PART(SUBSTRING($3 FROM 6), '/', 3) as rkey ), post_lookup AS ( @@ -13,6 +14,7 @@ post_lookup AS ( tid_to_i64(pp.rkey) as post_rkey FROM post_parts pp INNER JOIN actors a ON a.did = pp.did + WHERE pp.collection = 'app.bsky.feed.post' -- Only parse post URIs ) INSERT INTO bookmarks (actor_id, rkey, post_actor_id, post_rkey) SELECT diff --git a/consumer/src/db/sql/list_block_upsert.sql b/consumer/src/db/sql/list_block_upsert.sql index 80171fe2..b036c113 100644 --- a/consumer/src/db/sql/list_block_upsert.sql +++ b/consumer/src/db/sql/list_block_upsert.sql @@ -5,7 +5,9 @@ WITH list_parts AS ( SELECT SPLIT_PART(SUBSTRING($2 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($2 FROM 6), '/', 2) as collection, tid_to_i64(SPLIT_PART(SUBSTRING($2 FROM 6), '/', 3)) as rkey + WHERE SPLIT_PART(SUBSTRING($2 FROM 6), '/', 2) = 'app.bsky.graph.list' -- Only parse list URIs ), list_lookup AS ( SELECT l.id diff --git a/consumer/src/db/sql/postgate_upsert.sql b/consumer/src/db/sql/postgate_upsert.sql index 0294eb39..5a566971 100644 --- a/consumer/src/db/sql/postgate_upsert.sql +++ b/consumer/src/db/sql/postgate_upsert.sql @@ -38,9 +38,11 @@ detached_inserts AS ( CROSS JOIN LATERAL ( SELECT SPLIT_PART(SUBSTRING(uri FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING(uri FROM 6), '/', 2) as collection, tid_to_i64(SPLIT_PART(SUBSTRING(uri FROM 6), '/', 3)) as rkey ) parts INNER JOIN actors a ON a.did = parts.did + WHERE parts.collection = 'app.bsky.feed.post' -- Only parse post URIs, skip invalid ON CONFLICT DO NOTHING ) -- Return the operation result (1 row affected) diff --git a/consumer/src/db/sql/profile_upsert.sql b/consumer/src/db/sql/profile_upsert.sql index 566abeb9..6f3864a2 100644 --- a/consumer/src/db/sql/profile_upsert.sql +++ b/consumer/src/db/sql/profile_upsert.sql @@ -9,6 +9,7 @@ WITH pinned_lookup AS ( FROM actors a WHERE $7::text IS NOT NULL AND a.did = SPLIT_PART(SUBSTRING($7::text FROM 6), '/', 1) + AND SPLIT_PART(SUBSTRING($7::text FROM 6), '/', 2) = 'app.bsky.feed.post' -- Only parse post URIs ), joined_sp_lookup AS ( SELECT sp.id @@ -16,6 +17,7 @@ joined_sp_lookup AS ( INNER JOIN starterpacks sp ON sp.actor_id = a.id WHERE $8::text IS NOT NULL AND a.did = SPLIT_PART(SUBSTRING($8::text FROM 6), '/', 1) + AND SPLIT_PART(SUBSTRING($8::text FROM 6), '/', 2) = 'app.bsky.graph.starterpack' -- Only parse starterpack URIs AND sp.rkey = tid_to_i64(SPLIT_PART(SUBSTRING($8::text FROM 6), '/', 3)) ) UPDATE actors SET diff --git a/consumer/src/db/sql/status_upsert.sql b/consumer/src/db/sql/status_upsert.sql index 6531945a..6c694149 100644 --- a/consumer/src/db/sql/status_upsert.sql +++ b/consumer/src/db/sql/status_upsert.sql @@ -3,11 +3,13 @@ -- Parameters: $1=actor_id, $2=cid, $3=status, $4=duration, $5=embed_uri, $6=thumb_mime_type, $7=thumb_cid, $8=created_at -- NOTE: actor_id is provided by dispatcher after ensuring actor exists -- Note: No JOIN to posts table - we already have the natural key from URI parsing +-- Only parse post embeds - skip feed generators, profiles, etc. (non-post collections) WITH embed_lookup AS ( SELECT a.id AS actor_id, tid_to_i64(SPLIT_PART(SUBSTRING($5::text FROM 6), '/', 3)) as rkey FROM actors a WHERE $5::text IS NOT NULL AND a.did = SPLIT_PART(SUBSTRING($5::text FROM 6), '/', 1) + AND SPLIT_PART(SUBSTRING($5::text FROM 6), '/', 2) = 'app.bsky.feed.post' -- Only parse post URIs ) UPDATE actors SET status_cid = $2, diff --git a/consumer/src/db/sql/threadgate_upsert.sql b/consumer/src/db/sql/threadgate_upsert.sql index 3374ce89..0fb8ae4d 100644 --- a/consumer/src/db/sql/threadgate_upsert.sql +++ b/consumer/src/db/sql/threadgate_upsert.sql @@ -46,9 +46,11 @@ hidden_inserts AS ( CROSS JOIN LATERAL ( SELECT SPLIT_PART(SUBSTRING(uri FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING(uri FROM 6), '/', 2) as collection, tid_to_i64(SPLIT_PART(SUBSTRING(uri FROM 6), '/', 3)) as rkey ) parts INNER JOIN actors a ON a.did = parts.did + WHERE parts.collection = 'app.bsky.feed.post' -- Only parse post URIs, skip invalid ON CONFLICT DO NOTHING ), -- Insert allowed lists by resolving URIs to list IDs @@ -62,11 +64,13 @@ list_inserts AS ( CROSS JOIN LATERAL ( SELECT SPLIT_PART(SUBSTRING(uri FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING(uri FROM 6), '/', 2) as collection, tid_to_i64(SPLIT_PART(SUBSTRING(uri FROM 6), '/', 3)) as rkey ) parts INNER JOIN actors a ON a.did = parts.did INNER JOIN lists l ON l.actor_id = a.id AND l.rkey = parts.rkey WHERE $7 IS NOT NULL + AND parts.collection = 'app.bsky.graph.list' -- Only parse list URIs, skip invalid ON CONFLICT DO NOTHING ) -- Return the operation result (0 for insert, 1 for update) diff --git a/consumer/src/workers/backfill/resolve_bulk.rs b/consumer/src/workers/backfill/resolve_bulk.rs index e25c01d8..691aa75c 100644 --- a/consumer/src/workers/backfill/resolve_bulk.rs +++ b/consumer/src/workers/backfill/resolve_bulk.rs @@ -122,20 +122,27 @@ pub async fn resolve_all_references_bulk( } // Quoted post (from embed.record or embed.recordWithMedia) + // Note: Only add if it's actually a post - record embeds can be feedgens, lists, etc. if let Some(crate::types::records::EmbedOuter::Bsky(embed)) = &post.embed { use crate::types::records::AppBskyEmbed; match embed { AppBskyEmbed::Record(record_embed) => { - post_uris_with_cids.push(( - record_embed.record.uri.as_str(), - record_embed.record.cid.to_string(), - )); + // Only resolve as post if it's actually a post URI (not feedgen, list, etc.) + if record_embed.record.uri.contains("/app.bsky.feed.post/") { + post_uris_with_cids.push(( + record_embed.record.uri.as_str(), + record_embed.record.cid.to_string(), + )); + } } AppBskyEmbed::RecordWithMedia(rwm) => { - post_uris_with_cids.push(( - rwm.record.record.uri.as_str(), - rwm.record.record.cid.to_string(), - )); + // Only resolve as post if it's actually a post URI + if rwm.record.record.uri.contains("/app.bsky.feed.post/") { + post_uris_with_cids.push(( + rwm.record.record.uri.as_str(), + rwm.record.record.cid.to_string(), + )); + } } _ => {} }