diff --git a/consumer/src/db/operations/feed.rs b/consumer/src/db/operations/feed.rs index 39156928..cf507fc7 100644 --- a/consumer/src/db/operations/feed.rs +++ b/consumer/src/db/operations/feed.rs @@ -904,24 +904,21 @@ pub async fn postgate_upsert( let cid_digest = parakeet_db::cid_util::cid_to_digest(&cid_bytes) .expect("CID must be valid AT Protocol CID"); + let rkey_i64 = parakeet_db::models::tid_to_i64(rkey) + .wrap_err_with(|| format!("Invalid TID encoding in rkey: {}", rkey))?; + let rules = rec .embedding_rules .iter() .map(|v| v.as_str().to_owned()) .collect::>(); - // Query DID and construct at_uri (TODO: refactor SQL to use actor_id + rkey directly) - let did: String = conn - .query_one("SELECT did FROM actors WHERE id = $1", &[&actor_id]) - .await? - .get(0); - let at_uri = format!("at://{}/app.bsky.feed.postgate/{}", did, rkey); - // Note: created_at is derived from TID rkey conn.execute( include_str!("../sql/postgate_upsert.sql"), &[ - &at_uri, + &actor_id, + &rkey_i64, &cid_digest, &rec.post, &rec.detached_embedding_uris, @@ -1005,6 +1002,9 @@ pub async fn threadgate_upsert( let cid_digest = parakeet_db::cid_util::cid_to_digest(&cid_bytes) .expect("CID must be valid AT Protocol CID"); + let rkey_i64 = parakeet_db::models::tid_to_i64(rkey) + .wrap_err_with(|| format!("Invalid TID encoding in rkey: {}", rkey))?; + let record = serde_json::to_value(&rec).unwrap(); let allowed_lists = rec.allow.as_ref().map(|allow| { @@ -1024,18 +1024,12 @@ pub async fn threadgate_upsert( .collect::>() }); - // Query DID and construct at_uri (TODO: refactor SQL to use actor_id + rkey directly) - let did: String = conn - .query_one("SELECT did FROM actors WHERE id = $1", &[&actor_id]) - .await? - .get(0); - let at_uri = format!("at://{}/app.bsky.feed.threadgate/{}", did, rkey); - // Note: created_at is derived from TID rkey conn.execute( include_str!("../sql/threadgate_upsert.sql"), &[ - &at_uri, + &actor_id, + &rkey_i64, &cid_digest, &rec.post, &rec.hidden_replies, @@ -1531,18 +1525,11 @@ pub async fn feedgen_upsert( // Example: "app.bsky.feed.defs#contentModeUnspecified" -> "contentModeUnspecified" let content_mode = rec.content_mode.as_ref().map(|s| strip_lexicon_prefix(s)); - // Query DID and construct at_uri (TODO: refactor SQL to use actor_id + rkey directly) - let repo: String = conn - .query_one("SELECT did FROM actors WHERE id = $1", &[&actor_id]) - .await? - .get(0); - let at_uri = format!("at://{}/app.bsky.feed.generator/{}", repo, rkey); - conn.query_one( include_str!("../sql/feedgen_upsert.sql"), &[ - &at_uri, - &repo, + &actor_id, + &rkey, // TEXT rkey (feedgens use arbitrary strings, not TIDs) &cid_digest, &service_actor_id, &content_mode, diff --git a/consumer/src/db/operations/graph.rs b/consumer/src/db/operations/graph.rs index 6f40fe95..97466571 100644 --- a/consumer/src/db/operations/graph.rs +++ b/consumer/src/db/operations/graph.rs @@ -156,17 +156,9 @@ pub async fn list_upsert( conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) .await?; - // Query DID and construct at_uri (TODO: refactor SQL to use actor_id + rkey directly) - let did: String = conn - .query_one("SELECT did FROM actors WHERE id = $1", &[&actor_id]) - .await? - .get(0); - let at_uri = format!("at://{}/app.bsky.graph.list/{}", did, rkey); - conn.query_one( include_str!("../sql/list_upsert.sql"), &[ - &at_uri, &actor_id, &cid_digest, &rec.purpose, @@ -174,7 +166,6 @@ pub async fn list_upsert( &rec.description, &description_facets, &avatar, - // Note: created_at (param $9) removed - derived from TID rkey &rkey_i64, ], ) @@ -226,21 +217,12 @@ pub async fn list_block_insert( conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) .await?; - // Query DID and construct at_uri (TODO: refactor SQL to use actor_id + rkey directly) - let did: String = conn - .query_one("SELECT did FROM actors WHERE id = $1", &[&actor_id]) - .await? - .get(0); - let at_uri = format!("at://{}/app.bsky.graph.listblock/{}", did, rkey); - conn.execute( include_str!("../sql/list_block_upsert.sql"), &[ - &at_uri, &actor_id, &rec.subject, &cid_digest, - // Note: created_at (param $5) removed - derived from TID rkey &rkey_i64, ], ) diff --git a/consumer/src/db/operations/starter_pack.rs b/consumer/src/db/operations/starter_pack.rs index 33e44dba..16283e43 100644 --- a/consumer/src/db/operations/starter_pack.rs +++ b/consumer/src/db/operations/starter_pack.rs @@ -30,20 +30,12 @@ pub async fn starter_pack_upsert( conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) .await?; - // Query DID and construct at_uri (TODO: refactor SQL to use actor_id + rkey directly) - let did: String = conn - .query_one("SELECT did FROM actors WHERE id = $1", &[&actor_id]) - .await? - .get(0); - let at_uri = format!("at://{}/app.bsky.graph.starterpack/{}", did, rkey); - // Resolve list_id by creating stub list if needed let list_id = super::feed::ensure_list_id(conn, &rec.list).await?; conn.query_one( include_str!("../sql/starterpack_upsert.sql"), &[ - &at_uri, &actor_id, &cid_digest, &record, @@ -52,7 +44,6 @@ pub async fn starter_pack_upsert( &description_facets, &list_id, // Pass resolved list_id directly &feeds.as_deref(), // Pass as &[String] for text[] parameter - // Note: created_at (param $10) removed - derived from TID rkey &rkey_i64, ], ) diff --git a/consumer/src/db/sql/feedgen_upsert.sql b/consumer/src/db/sql/feedgen_upsert.sql index 3774435b..8360f7c5 100644 --- a/consumer/src/db/sql/feedgen_upsert.sql +++ b/consumer/src/db/sql/feedgen_upsert.sql @@ -1,25 +1,17 @@ -- Insert/update feedgen with CTE to resolve actor_id -- Self-contained schema: No records table, metadata embedded directly --- Parameters: $1=at_uri, $2=owner, $3=cid, $4=service_actor_id, $5=content_mode, $6=name, +-- Parameters: $1=actor_id, $2=rkey, $3=cid, $4=service_actor_id, $5=content_mode, $6=name, -- $7=description, $8=description_facets, $9=avatar_cid, $10=created_at -WITH feedgen_parts AS ( - SELECT - SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, - SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey -), -actor_lookup AS ( - SELECT id FROM actors WHERE did = $2 -) INSERT INTO feedgens ( actor_id, rkey, cid, created_at, owner_actor_id, service_actor_id, content_mode, name, description, description_facets, avatar_cid, accepts_interactions ) SELECT - (SELECT id FROM actor_lookup), -- actor_id (direct, no records table) - (SELECT rkey FROM feedgen_parts), -- rkey (TEXT) + $1, -- actor_id (provided by caller) + $2, -- rkey (TEXT, provided by caller) $3::bytea, -- cid (embedded) $10, -- created_at (embedded) - (SELECT id FROM actor_lookup), -- owner_actor_id (same as actor_id for feed generators) + $1, -- owner_actor_id (same as actor_id for feed generators) $4, -- service_actor_id (passed as parameter, already ensured in executor) $5::text::content_mode, $6, -- name diff --git a/consumer/src/db/sql/list_block_upsert.sql b/consumer/src/db/sql/list_block_upsert.sql index 516b2d4b..80171fe2 100644 --- a/consumer/src/db/sql/list_block_upsert.sql +++ b/consumer/src/db/sql/list_block_upsert.sql @@ -1,14 +1,11 @@ -- Insert list_block with self-contained schema (no records table) --- Parameters: $1=at_uri, $2=actor_id, $3=list_uri, $4=cid(bytea), $5=rkey +-- Parameters: $1=actor_id, $2=list_uri, $3=cid(bytea), $4=rkey -- Note: created_at is derived from TID rkey -- NOTE: actor_id is provided by dispatcher after ensuring actor exists -WITH unused_params AS ( - SELECT $1::text as at_uri -- Type hint for unused parameter -), -list_parts AS ( +WITH list_parts AS ( SELECT - SPLIT_PART(SUBSTRING($3 FROM 6), '/', 1) as did, - tid_to_i64(SPLIT_PART(SUBSTRING($3 FROM 6), '/', 3)) as rkey + SPLIT_PART(SUBSTRING($2 FROM 6), '/', 1) as did, + tid_to_i64(SPLIT_PART(SUBSTRING($2 FROM 6), '/', 3)) as rkey ), list_lookup AS ( SELECT l.id @@ -18,9 +15,9 @@ list_lookup AS ( ) INSERT INTO list_blocks (actor_id, rkey, cid, list_id) SELECT - $2, -- actor_id (provided by dispatcher) - $5, -- rkey (INT8) - $4, -- cid (embedded, already bytea) + $1, -- actor_id (provided by dispatcher) + $4, -- rkey (INT8) + $3, -- cid (embedded, already bytea) (SELECT id FROM list_lookup) -- list_id (direct FK to lists.id) WHERE (SELECT id FROM list_lookup) IS NOT NULL -- Only insert if list exists ON CONFLICT (actor_id, rkey) DO NOTHING diff --git a/consumer/src/db/sql/list_upsert.sql b/consumer/src/db/sql/list_upsert.sql index 6fc6cc01..53cf6996 100644 --- a/consumer/src/db/sql/list_upsert.sql +++ b/consumer/src/db/sql/list_upsert.sql @@ -1,22 +1,19 @@ -- Insert/update list with self-contained schema (no records table) --- Parameters: $1=at_uri, $2=actor_id(INT4), $3=cid(bytea), $4=list_type, $5=name, $6=description, --- $7=description_facets, $8=avatar_cid, $9=rkey +-- Parameters: $1=actor_id(INT4), $2=cid(bytea), $3=list_type, $4=name, $5=description, +-- $6=description_facets, $7=avatar_cid, $8=rkey -- Note: created_at is derived from TID rkey -- Note: actor_id is provided by caller after calling get_actor_id (ensures zero sequence waste) -WITH unused_params AS ( - SELECT $1::text as at_uri -- Type hint for unused parameter -) INSERT INTO lists (actor_id, rkey, cid, owner_actor_id, list_type, name, description, description_facets, avatar_cid) SELECT - $2, -- actor_id (provided by caller) - $9, -- rkey (INT8) - $3::bytea, -- cid (embedded) - $2, -- owner_actor_id (same as actor_id for lists) - $4::text::list_type, - $5, -- name - $6, -- description - $7, -- description_facets - $8 -- avatar_cid + $1, -- actor_id (provided by caller) + $8, -- rkey (INT8) + $2::bytea, -- cid (embedded) + $1, -- owner_actor_id (same as actor_id for lists) + $3::text::list_type, + $4, -- name + $5, -- description + $6, -- description_facets + $7 -- avatar_cid ON CONFLICT (actor_id, rkey) DO UPDATE SET cid=EXCLUDED.cid, list_type=EXCLUDED.list_type, diff --git a/consumer/src/db/sql/postgate_upsert.sql b/consumer/src/db/sql/postgate_upsert.sql index 8c803541..8be11bcf 100644 --- a/consumer/src/db/sql/postgate_upsert.sql +++ b/consumer/src/db/sql/postgate_upsert.sql @@ -1,19 +1,11 @@ -- Insert/update postgate with CTE to resolve actor_id and post_id -- Self-contained schema: No records table, metadata embedded directly --- Parameters: $1=at_uri, $2=cid, $3=post_uri, $4=detached, $5=rules, $6=post_cid +-- Parameters: $1=actor_id, $2=rkey, $3=cid, $4=post_uri, $5=detached, $6=rules -- Note: created_at is derived from TID rkey -WITH postgate_parts AS ( +WITH post_parts AS ( SELECT - SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, - tid_to_i64(SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3)) as rkey -), -actor_lookup AS ( - SELECT id FROM actors WHERE did = (SELECT did FROM postgate_parts) -), -post_parts AS ( - SELECT - SPLIT_PART(SUBSTRING($3 FROM 6), '/', 1) as did, - tid_to_i64(SPLIT_PART(SUBSTRING($3 FROM 6), '/', 3)) as rkey + SPLIT_PART(SUBSTRING($4 FROM 6), '/', 1) as did, + tid_to_i64(SPLIT_PART(SUBSTRING($4 FROM 6), '/', 3)) as rkey ), post_lookup AS ( SELECT p.id @@ -24,11 +16,11 @@ post_lookup AS ( upserted_postgate AS ( INSERT INTO postgates (actor_id, rkey, cid, post_id, rules) SELECT - (SELECT id FROM actor_lookup), -- actor_id (direct, no records table) - (SELECT rkey FROM postgate_parts), -- rkey (INT8) - $2::bytea, -- cid (embedded) + $1, -- actor_id (provided by caller) + $2, -- rkey (INT8, provided by caller) + $3::bytea, -- cid (embedded) (SELECT id FROM post_lookup), -- post_id (direct FK to posts.id, NULL if not found) - $5::text[]::postgate_rule[] -- rules (cast array to enum array) + $6::text[]::postgate_rule[] -- rules (cast array to enum array) ON CONFLICT (actor_id, rkey) DO UPDATE SET cid=EXCLUDED.cid, post_id=EXCLUDED.post_id, @@ -46,7 +38,7 @@ detached_inserts AS ( SELECT (SELECT id FROM upserted_postgate), p.id - FROM UNNEST($4::text[]) AS uri + FROM UNNEST($5::text[]) AS uri CROSS JOIN LATERAL ( SELECT SPLIT_PART(SUBSTRING(uri FROM 6), '/', 1) as did, diff --git a/consumer/src/db/sql/starterpack_upsert.sql b/consumer/src/db/sql/starterpack_upsert.sql index d5e25aba..11c327a5 100644 --- a/consumer/src/db/sql/starterpack_upsert.sql +++ b/consumer/src/db/sql/starterpack_upsert.sql @@ -1,22 +1,22 @@ -- Insert/update starterpack with self-contained schema (no records table) --- Parameters: $1=at_uri, $2=actor_id(INT4), $3=cid(bytea), $4=record(unused), $5=name, $6=description, --- $7=description_facets, $8=list_id(INT8), $9=feeds(array of feed URIs), $10=rkey +-- Parameters: $1=actor_id(INT4), $2=cid(bytea), $3=record(unused), $4=name, $5=description, +-- $6=description_facets, $7=list_id(INT8), $8=feeds(array of feed URIs), $9=rkey -- Note: created_at is derived from TID rkey -- Note: actor_id is provided by caller after calling get_actor_id (ensures zero sequence waste) WITH unused_params AS ( - SELECT $1::text as at_uri, $4::jsonb as record -- Type hints for unused parameters + SELECT $3::jsonb as record -- Type hint for unused parameter ), upserted_starterpack AS ( INSERT INTO starterpacks (actor_id, rkey, cid, owner_actor_id, name, description, description_facets, list_id, status) SELECT - $2, -- actor_id (provided by caller) - $10, -- rkey (INT8) - $3::bytea, -- cid (embedded) - $2, -- owner_actor_id (same as actor_id for starterpacks) + $1, -- actor_id (provided by caller) + $9, -- rkey (INT8) + $2::bytea, -- cid (embedded) + $1, -- owner_actor_id (same as actor_id for starterpacks) + $4, $5, $6, - $7, - $8, -- list_id (already resolved, may be stub) + $7, -- list_id (already resolved, may be stub) 'complete'::record_status ON CONFLICT (actor_id, rkey) DO UPDATE SET cid=EXCLUDED.cid, @@ -34,14 +34,14 @@ deleted_feeds AS ( AND EXISTS (SELECT 1 FROM upserted_starterpack) ), -- Insert new feeds by resolving feed URIs to feedgen IDs (only if starterpack was actually inserted) --- Feeds are passed as array of URIs in $9, we enumerate them with position +-- Feeds are passed as array of URIs in $8, we enumerate them with position inserted_feeds AS ( INSERT INTO starterpack_feeds (starterpack_id, feed_id, position) SELECT (SELECT id FROM upserted_starterpack), fg.id, (row_number() OVER (ORDER BY ord) - 1)::smallint as position - FROM UNNEST($9::text[]) WITH ORDINALITY AS feed_uri_table(feed_uri, ord) + FROM UNNEST($8::text[]) WITH ORDINALITY AS feed_uri_table(feed_uri, ord) CROSS JOIN LATERAL ( SELECT SPLIT_PART(SUBSTRING(feed_uri FROM 6), '/', 1) as did, @@ -49,7 +49,7 @@ inserted_feeds AS ( ) parts INNER JOIN actors a ON a.did = parts.did INNER JOIN feedgens fg ON fg.actor_id = a.id AND fg.rkey = parts.rkey - WHERE $9 IS NOT NULL + WHERE $8 IS NOT NULL AND EXISTS (SELECT 1 FROM upserted_starterpack) ON CONFLICT DO NOTHING ) diff --git a/consumer/src/db/sql/threadgate_upsert.sql b/consumer/src/db/sql/threadgate_upsert.sql index 705cbf44..2683a465 100644 --- a/consumer/src/db/sql/threadgate_upsert.sql +++ b/consumer/src/db/sql/threadgate_upsert.sql @@ -1,23 +1,15 @@ -- Insert/update threadgate with CTE to resolve actor_id and post_id -- Self-contained schema: No records table, metadata embedded directly --- Parameters: $1=at_uri, $2=cid, $3=post_uri, $4=hidden_replies, $5=allow, --- $6=allowed_lists, $7=record(unused), $8=post_cid +-- Parameters: $1=actor_id, $2=rkey, $3=cid, $4=post_uri, $5=hidden_replies, $6=allow, +-- $7=allowed_lists, $8=record(unused) -- Note: created_at is derived from TID rkey WITH unused_params AS ( - SELECT $7::jsonb as record -- Type hint for unused parameter -), -threadgate_parts AS ( - SELECT - SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, - tid_to_i64(SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3)) as rkey -), -actor_lookup AS ( - SELECT id FROM actors WHERE did = (SELECT did FROM threadgate_parts) + SELECT $8::jsonb as record -- Type hint for unused parameter ), post_parts AS ( SELECT - SPLIT_PART(SUBSTRING($3 FROM 6), '/', 1) as did, - tid_to_i64(SPLIT_PART(SUBSTRING($3 FROM 6), '/', 3)) as rkey + SPLIT_PART(SUBSTRING($4 FROM 6), '/', 1) as did, + tid_to_i64(SPLIT_PART(SUBSTRING($4 FROM 6), '/', 3)) as rkey ), post_lookup AS ( SELECT p.id @@ -28,11 +20,11 @@ post_lookup AS ( upserted_threadgate AS ( INSERT INTO threadgates (actor_id, rkey, cid, post_id, allow) SELECT - (SELECT id FROM actor_lookup), -- actor_id (direct, no records table) - (SELECT rkey FROM threadgate_parts), -- rkey (INT8) - $2, -- cid (bytea - stored as digest) + $1, -- actor_id (provided by caller) + $2, -- rkey (INT8, provided by caller) + $3, -- cid (bytea - stored as digest) (SELECT id FROM post_lookup), -- post_id (direct FK to posts.id, NULL if not found) - $5::text[]::threadgate_rule[] -- allow (cast to enum array) + $6::text[]::threadgate_rule[] -- allow (cast to enum array) ON CONFLICT (actor_id, rkey) DO UPDATE SET cid=EXCLUDED.cid, post_id=EXCLUDED.post_id, @@ -54,7 +46,7 @@ hidden_inserts AS ( SELECT (SELECT id FROM upserted_threadgate), p.id - FROM UNNEST($4::text[]) AS uri + FROM UNNEST($5::text[]) AS uri CROSS JOIN LATERAL ( SELECT SPLIT_PART(SUBSTRING(uri FROM 6), '/', 1) as did, @@ -70,7 +62,7 @@ list_inserts AS ( SELECT (SELECT id FROM upserted_threadgate), l.id - FROM UNNEST($6::text[]) AS uri + FROM UNNEST($7::text[]) AS uri CROSS JOIN LATERAL ( SELECT SPLIT_PART(SUBSTRING(uri FROM 6), '/', 1) as did, @@ -78,7 +70,7 @@ list_inserts AS ( ) 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 $6 IS NOT NULL + WHERE $7 IS NOT NULL ON CONFLICT DO NOTHING ) -- Return the operation result (0 for insert, 1 for update)