From 2ebd09a935b5418cb126748567f1113004ae360c Mon Sep 17 00:00:00 2001 From: Timothy Quilling Date: Fri, 31 Oct 2025 01:50:50 -0400 Subject: [PATCH] feat: database hacks -- insanity, consumer --- Cargo.lock | 1 + consumer/src/db/gates.rs | 13 ++- consumer/src/db/operations/feed.rs | 5 +- consumer/src/db/operations/notification.rs | 36 +++++--- consumer/src/db/sql/block_upsert.sql | 26 +++++- consumer/src/db/sql/bookmarks_upsert.sql | 30 ++++++- consumer/src/db/sql/feedgen_upsert.sql | 62 ++++++++++--- consumer/src/db/sql/label_copy_upsert.sql | 8 +- consumer/src/db/sql/label_defs_upsert.sql | 9 +- consumer/src/db/sql/label_service_upsert.sql | 37 ++++++-- consumer/src/db/sql/list_block_upsert.sql | 43 ++++++++- consumer/src/db/sql/list_item_upsert.sql | 43 ++++++++- consumer/src/db/sql/list_upsert.sql | 46 +++++++--- consumer/src/db/sql/post_insert.sql | 52 +++++------ consumer/src/db/sql/postgate_upsert.sql | 45 ++++++++-- consumer/src/db/sql/profile_upsert.sql | 91 ++++++++++++++++---- consumer/src/db/sql/starterpack_upsert.sql | 59 ++++++++++--- consumer/src/db/sql/status_upsert.sql | 66 +++++++++++--- consumer/src/db/sql/threadgate_upsert.sql | 46 ++++++++-- parakeet-db/Cargo.toml | 1 + parakeet-db/src/actor_cache.rs | 6 +- parakeet-db/src/at_uri_util.rs | 1 - parakeet-db/src/cid_util.rs | 1 - parakeet-db/src/types.rs | 35 ++++++++ 24 files changed, 611 insertions(+), 151 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index bcc016b8..5fbe485d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2672,6 +2672,7 @@ name = "parakeet-db" version = "0.1.0" dependencies = [ "bloomfilter", + "bytes", "chrono", "cid", "diesel", diff --git a/consumer/src/db/gates.rs b/consumer/src/db/gates.rs index 190fdebf..84873cef 100644 --- a/consumer/src/db/gates.rs +++ b/consumer/src/db/gates.rs @@ -40,9 +40,20 @@ pub async fn post_enforce_threadgate( let allow: HashSet = HashSet::from_iter(allow); if allow.contains(THREADGATE_RULE_FOLLOWER) || allow.contains(THREADGATE_RULE_FOLLOWING) { + // Check follow relationships directly from follows table + // following: root_author follows post_author + // followed: post_author follows root_author let profile_state: Option<(bool, bool)> = conn .query_opt( - "SELECT following IS NOT NULL, followed IS NOT NULL FROM profile_states WHERE did=$1 AND subject=$2", + "WITH root_actor AS ( + SELECT id FROM actors WHERE did = $1 + ), + post_actor AS ( + SELECT id FROM actors WHERE did = $2 + ) + SELECT + EXISTS(SELECT 1 FROM follows WHERE actor_id = (SELECT id FROM root_actor) AND subject_actor_id = (SELECT id FROM post_actor)) as following, + EXISTS(SELECT 1 FROM follows WHERE actor_id = (SELECT id FROM post_actor) AND subject_actor_id = (SELECT id FROM root_actor)) as followed", &[&root_author, &post_author], ) .await? diff --git a/consumer/src/db/operations/feed.rs b/consumer/src/db/operations/feed.rs index 29c8676d..82891f38 100644 --- a/consumer/src/db/operations/feed.rs +++ b/consumer/src/db/operations/feed.rs @@ -122,7 +122,7 @@ pub async fn post_delete( .query_opt( "WITH deleted_post AS ( DELETE FROM posts WHERE at_uri=$1 - RETURNING actor_id, parent_record_id, (SELECT uri FROM post_embed_record WHERE post_uri = $1 LIMIT 1) as embed_uri + RETURNING record_id, parent_record_id, (SELECT uri FROM post_embed_record WHERE post_uri = $1 LIMIT 1) as embed_uri ) SELECT a.did, @@ -131,7 +131,8 @@ pub async fn post_delete( ELSE NULL END as parent_uri, dp.embed_uri FROM deleted_post dp - INNER JOIN actors a ON dp.actor_id = a.id + INNER JOIN records r ON dp.record_id = r.id + INNER JOIN actors a ON r.actor_id = a.id LEFT JOIN records pr ON dp.parent_record_id = pr.id LEFT JOIN actors pa ON pr.actor_id = pa.id", &[&at_uri], diff --git a/consumer/src/db/operations/notification.rs b/consumer/src/db/operations/notification.rs index aad5e80d..745ebfc1 100644 --- a/consumer/src/db/operations/notification.rs +++ b/consumer/src/db/operations/notification.rs @@ -14,8 +14,8 @@ pub async fn notification_insert( author_did: &str, reason: &str, reason_subject: Option<&str>, - cid: &str, - created_at: DateTime, + _cid: &str, // Unused - kept for API compatibility + _created_at: DateTime, // Unused - kept for API compatibility ) -> PgExecResult { // Use CTE to atomically ensure author actor exists and resolve IDs before inserting notification conn.execute( @@ -27,9 +27,6 @@ pub async fn notification_insert( recipient_lookup AS ( SELECT id FROM actors WHERE did = $2 ), - author_lookup AS ( - SELECT id FROM actors WHERE did = $3 - ), record_parts AS ( SELECT SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, @@ -43,25 +40,36 @@ pub async fn notification_insert( INNER JOIN records r ON r.actor_id = a.id AND r.collection::text = rp.collection AND r.rkey = rp.rkey + ), + reason_subject_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($5 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($5 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($5 FROM 6), '/', 3) as rkey + WHERE $5 IS NOT NULL + ), + reason_subject_lookup AS ( + SELECT r.id + FROM reason_subject_parts rsp + INNER JOIN actors a ON a.did = rsp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rsp.collection + AND r.rkey = rsp.rkey ) - INSERT INTO notifications (record_id, recipient_actor_id, author_actor_id, reason, reason_subject, cid_digest, created_at) + INSERT INTO notifications (record_id, recipient_actor_id, reason, reason_subject_record_id, is_read) SELECT (SELECT id FROM record_lookup), (SELECT id FROM recipient_lookup), - (SELECT id FROM author_lookup), - $4, - $5, - DECODE(SUBSTRING($6 FROM 9), 'hex'), - $7 - ON CONFLICT (record_id, recipient_actor_id) DO NOTHING", + $4::notification_reason, + (SELECT id FROM reason_subject_lookup), + FALSE + ON CONFLICT (recipient_actor_id, record_id) DO NOTHING", &[ &uri, &recipient_did, &author_did, &reason, &reason_subject, - &cid, - &created_at.naive_utc(), ], ) .await diff --git a/consumer/src/db/sql/block_upsert.sql b/consumer/src/db/sql/block_upsert.sql index 5d002cb3..b262c29c 100644 --- a/consumer/src/db/sql/block_upsert.sql +++ b/consumer/src/db/sql/block_upsert.sql @@ -1,5 +1,6 @@ --- Insert block with CTE to ensure both blocker and blocked actors exist +-- Insert block with CTE to resolve record_id, actor_id, and subject_actor_id -- This prevents foreign key violations when indexing blocks +-- Parameters: $1=at_uri, $2=blocker_did, $3=blocked_did, $4=created_at WITH ensure_blocker AS ( INSERT INTO actors (did, status, sync_state, last_indexed) VALUES ($2, 'active', 'partial', NOW()) @@ -15,7 +16,24 @@ blocker_lookup AS ( ), blocked_lookup AS ( SELECT id FROM actors WHERE did = $3 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint ) -INSERT INTO blocks (rkey, actor_id, subject_actor_id, created_at) -SELECT $1, (SELECT id FROM blocker_lookup), (SELECT id FROM blocked_lookup), $4 -ON CONFLICT (actor_id, rkey) DO NOTHING +INSERT INTO blocks (actor_id, subject_actor_id, record_id) +SELECT + (SELECT id FROM blocker_lookup), + (SELECT id FROM blocked_lookup), + (SELECT id FROM record_lookup) +ON CONFLICT (actor_id, subject_actor_id) DO NOTHING diff --git a/consumer/src/db/sql/bookmarks_upsert.sql b/consumer/src/db/sql/bookmarks_upsert.sql index 2d06cd2c..149c27a8 100644 --- a/consumer/src/db/sql/bookmarks_upsert.sql +++ b/consumer/src/db/sql/bookmarks_upsert.sql @@ -1,4 +1,26 @@ -INSERT INTO bookmarks (did, rkey, subject, subject_type, tags, created_at) -VALUES ($1, $2, $3, $4, $5, $6) --- Bookmarks are immutable - have unique (did, rkey) key -ON CONFLICT (did, rkey) DO NOTHING \ No newline at end of file +-- Insert bookmark with CTE to resolve actor_id and subject_record_id +-- Parameters: $1=did, $2=rkey, $3=subject_uri, $4=subject_type, $5=tags(unused), $6=created_at +WITH actor_lookup AS ( + SELECT id FROM actors WHERE did = $1 +), +subject_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 +), +subject_lookup AS ( + SELECT r.id + FROM subject_parts sp + INNER JOIN actors a ON a.did = sp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = sp.collection + AND r.rkey = sp.rkey::bigint +) +INSERT INTO bookmarks (actor_id, subject_record_id, created_at) +SELECT + (SELECT id FROM actor_lookup), + (SELECT id FROM subject_lookup), + $6 +-- Bookmarks are immutable - have unique (actor_id, subject_record_id) key +ON CONFLICT (actor_id, subject_record_id) DO NOTHING \ No newline at end of file diff --git a/consumer/src/db/sql/feedgen_upsert.sql b/consumer/src/db/sql/feedgen_upsert.sql index dcdb5c3b..a717f84f 100644 --- a/consumer/src/db/sql/feedgen_upsert.sql +++ b/consumer/src/db/sql/feedgen_upsert.sql @@ -1,12 +1,50 @@ -INSERT INTO feedgens (at_uri, owner, cid, service_did, content_mode, name, description, description_facets, avatar_cid, - created_at) -VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) -ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, - service_did=EXCLUDED.service_did, - content_mode=EXCLUDED.content_mode, - name=EXCLUDED.name, - description=EXCLUDED.description, - description_facets=EXCLUDED.description_facets, - avatar_cid=EXCLUDED.avatar_cid, - indexed_at=NOW() -RETURNING XMAX::text::int \ No newline at end of file +-- Insert/update feedgen with CTE to resolve record_id, owner_actor_id, and service_actor_id +-- Parameters: $1=at_uri, $2=owner, $3=cid, $4=service_did, $5=content_mode, $6=name, +-- $7=description, $8=description_facets, $9=avatar_cid, $10=created_at +WITH owner_lookup AS ( + SELECT id FROM actors WHERE did = $2 +), +service_ensure AS ( + INSERT INTO actors (did, status, sync_state, last_indexed) + VALUES ($4, 'active', 'partial', NOW()) + ON CONFLICT (did) DO NOTHING +), +service_lookup AS ( + SELECT id FROM actors WHERE did = $4 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint +) +INSERT INTO feedgens ( + record_id, owner_actor_id, service_actor_id, content_mode, name, + description, description_facets, avatar_cid, accepts_interactions +) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM owner_lookup), + (SELECT id FROM service_lookup), + $5::content_mode, + $6, + $7, + $8, + $9, + true -- Default accepts_interactions to true +ON CONFLICT (record_id) DO UPDATE SET + service_actor_id=EXCLUDED.service_actor_id, + content_mode=EXCLUDED.content_mode, + name=EXCLUDED.name, + description=EXCLUDED.description, + description_facets=EXCLUDED.description_facets, + avatar_cid=EXCLUDED.avatar_cid +RETURNING XMAX::text::int diff --git a/consumer/src/db/sql/label_copy_upsert.sql b/consumer/src/db/sql/label_copy_upsert.sql index 82458906..b553aad3 100644 --- a/consumer/src/db/sql/label_copy_upsert.sql +++ b/consumer/src/db/sql/label_copy_upsert.sql @@ -1,8 +1,10 @@ +-- Copy labels from label_tmp staging table to labels table +-- Uses DISTINCT ON to ensure unique (labeler_actor_id, label, uri) combinations INSERT INTO labels -SELECT DISTINCT on (labeler, label, uri) * +SELECT DISTINCT on (labeler_actor_id, label, uri) * FROM label_tmp -ON CONFLICT (labeler, label, uri) DO UPDATE +ON CONFLICT (labeler_actor_id, label, uri) DO UPDATE SET negated=EXCLUDED.negated, expires=EXCLUDED.expires, sig=EXCLUDED.sig, - created_at=excluded.created_at \ No newline at end of file + created_at=excluded.created_at diff --git a/consumer/src/db/sql/label_defs_upsert.sql b/consumer/src/db/sql/label_defs_upsert.sql index 5cd2c49f..e9341ac9 100644 --- a/consumer/src/db/sql/label_defs_upsert.sql +++ b/consumer/src/db/sql/label_defs_upsert.sql @@ -1,9 +1,12 @@ -INSERT INTO labeler_defs (labeler, label_identifier, severity, blurs, default_setting, adult_only, locales) +-- Insert/update labeler definition +-- Parameters: $1=labeler_actor_id, $2=label_identifier, $3=severity, $4=blurs, +-- $5=default_setting, $6=adult_only, $7=locales +INSERT INTO labeler_defs (labeler_actor_id, label_identifier, severity, blurs, default_setting, adult_only, locales) VALUES ($1, $2, $3, $4, $5, $6, $7) -ON CONFLICT (labeler, label_identifier) DO UPDATE +ON CONFLICT (labeler_actor_id, label_identifier) DO UPDATE SET severity=EXCLUDED.severity, blurs=EXCLUDED.blurs, default_setting=EXCLUDED.default_setting, adult_only=EXCLUDED.adult_only, locales=EXCLUDED.locales, - indexed_at=NOW() \ No newline at end of file + indexed_at=NOW() diff --git a/consumer/src/db/sql/label_service_upsert.sql b/consumer/src/db/sql/label_service_upsert.sql index cf126607..39071a43 100644 --- a/consumer/src/db/sql/label_service_upsert.sql +++ b/consumer/src/db/sql/label_service_upsert.sql @@ -1,7 +1,30 @@ -INSERT INTO labelers (did, cid, reasons, subject_types, subject_collections) -VALUES ($1, $2, $3, $4, $5) -ON CONFLICT (did) DO UPDATE SET cid=EXCLUDED.cid, - reasons=EXCLUDED.reasons, - subject_types=EXCLUDED.subject_types, - subject_collections=EXCLUDED.subject_collections, - indexed_at=NOW() \ No newline at end of file +-- Insert/update labeler service with CTE to resolve record_id and actor_id +-- Parameters: $1=did, $2=cid, $3=reasons, $4=subject_types, $5=subject_collections +WITH actor_lookup AS ( + SELECT id FROM actors WHERE did = $1 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING(('at://' || $1 || '/app.bsky.labeler.service/self') FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING(('at://' || $1 || '/app.bsky.labeler.service/self') FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING(('at://' || $1 || '/app.bsky.labeler.service/self') FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey +) +INSERT INTO labelers (record_id, actor_id, reasons, subject_types, subject_collections) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM actor_lookup), + $3, + $4, + $5 +ON CONFLICT (record_id) DO UPDATE SET + reasons=EXCLUDED.reasons, + subject_types=EXCLUDED.subject_types, + subject_collections=EXCLUDED.subject_collections diff --git a/consumer/src/db/sql/list_block_upsert.sql b/consumer/src/db/sql/list_block_upsert.sql index 20a521f7..e346a958 100644 --- a/consumer/src/db/sql/list_block_upsert.sql +++ b/consumer/src/db/sql/list_block_upsert.sql @@ -1,5 +1,40 @@ --- Insert list_block +-- Insert list_block with CTE to resolve record_id, actor_id, and list_record_id -- Actor existence is ensured by batch writer before this operation -INSERT INTO list_blocks (at_uri, did, list_uri, created_at) -VALUES ($1, $2, $3, $4) -ON CONFLICT DO NOTHING +-- Parameters: $1=at_uri, $2=did, $3=list_uri, $4=created_at +WITH actor_lookup AS ( + SELECT id FROM actors WHERE did = $2 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint +), +list_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 +), +list_lookup AS ( + SELECT r.id + FROM list_parts lp + INNER JOIN actors a ON a.did = lp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = lp.collection + AND r.rkey = lp.rkey::bigint +) +INSERT INTO list_blocks (record_id, actor_id, list_record_id) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM actor_lookup), + (SELECT id FROM list_lookup) +ON CONFLICT (record_id) DO NOTHING diff --git a/consumer/src/db/sql/list_item_upsert.sql b/consumer/src/db/sql/list_item_upsert.sql index cf20d28b..9661a89d 100644 --- a/consumer/src/db/sql/list_item_upsert.sql +++ b/consumer/src/db/sql/list_item_upsert.sql @@ -1,10 +1,45 @@ --- Insert list_item with CTE to ensure subject actor exists +-- Insert list_item with CTE to resolve record_id, list_record_id, and subject_actor_id -- This prevents foreign key violations when indexing list items +-- Parameters: $1=at_uri, $2=list_uri, $3=subject, $4=created_at WITH ensure_subject AS ( INSERT INTO actors (did, status, sync_state, last_indexed) VALUES ($3, 'active', 'partial', NOW()) ON CONFLICT (did) DO NOTHING +), +subject_lookup AS ( + SELECT id FROM actors WHERE did = $3 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint +), +list_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($2 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($2 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($2 FROM 6), '/', 3) as rkey +), +list_lookup AS ( + SELECT r.id + FROM list_parts lp + INNER JOIN actors a ON a.did = lp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = lp.collection + AND r.rkey = lp.rkey::bigint ) -INSERT INTO list_items (at_uri, list_uri, subject, created_at) -VALUES ($1, $2, $3, $4) -ON CONFLICT DO NOTHING +INSERT INTO list_items (record_id, list_record_id, subject_actor_id) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM list_lookup), + (SELECT id FROM subject_lookup) +ON CONFLICT (record_id) DO NOTHING diff --git a/consumer/src/db/sql/list_upsert.sql b/consumer/src/db/sql/list_upsert.sql index af850932..dd8a0e51 100644 --- a/consumer/src/db/sql/list_upsert.sql +++ b/consumer/src/db/sql/list_upsert.sql @@ -1,10 +1,36 @@ -INSERT INTO lists (at_uri, owner, cid, list_type, name, description, description_facets, avatar_cid, created_at) -VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) -ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, - list_type=EXCLUDED.list_type, - name=EXCLUDED.name, - description=EXCLUDED.description, - description_facets=EXCLUDED.description_facets, - avatar_cid=EXCLUDED.avatar_cid, - indexed_at=NOW() -RETURNING XMAX::text::int \ No newline at end of file +-- Insert/update list with CTE to resolve record_id and owner_actor_id +-- Parameters: $1=at_uri, $2=owner, $3=cid, $4=list_type, $5=name, $6=description, +-- $7=description_facets, $8=avatar_cid, $9=created_at +WITH owner_lookup AS ( + SELECT id FROM actors WHERE did = $2 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint +) +INSERT INTO lists (record_id, owner_actor_id, list_type, name, description, description_facets, avatar_cid) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM owner_lookup), + $4::list_type, + $5, + $6, + $7, + $8 +ON CONFLICT (record_id) DO UPDATE SET + list_type=EXCLUDED.list_type, + name=EXCLUDED.name, + description=EXCLUDED.description, + description_facets=EXCLUDED.description_facets, + avatar_cid=EXCLUDED.avatar_cid +RETURNING XMAX::text::int diff --git a/consumer/src/db/sql/post_insert.sql b/consumer/src/db/sql/post_insert.sql index 763b5a40..d5109193 100644 --- a/consumer/src/db/sql/post_insert.sql +++ b/consumer/src/db/sql/post_insert.sql @@ -1,13 +1,21 @@ --- Insert post with CTE to ensure author exists and resolve parent/root record IDs +-- Insert post with CTE to resolve record_id and parent/root record IDs -- This prevents foreign key violations when indexing posts -WITH actor_ensure AS ( - INSERT INTO actors (did, status, sync_state, last_indexed) - VALUES ($2, 'active', 'partial', NOW()) - ON CONFLICT (did) DO NOTHING - RETURNING id +-- Parameters: $1=at_uri, $2=repo(did), $3=cid, $4=record, $5=content, $6=facets, +-- $7=languages, $8=tags, $9=parent_uri, $10=root_uri, $11=embed, +-- $12=embed_subtype, $13=mentions, $14=violates_threadgate, $15=created_at +WITH record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey ), -actor_id_lookup AS ( - SELECT id FROM actors WHERE did = $2 +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint ), parent_parts AS ( SELECT @@ -22,7 +30,7 @@ parent_lookup AS ( INNER JOIN actors a ON a.did = pp.did INNER JOIN records r ON r.actor_id = a.id AND r.collection::text = pp.collection - AND r.rkey = pp.rkey + AND r.rkey = pp.rkey::bigint ), root_parts AS ( SELECT @@ -37,28 +45,22 @@ root_lookup AS ( INNER JOIN actors a ON a.did = rp.did INNER JOIN records r ON r.actor_id = a.id AND r.collection::text = rp.collection - AND r.rkey = rp.rkey + AND r.rkey = rp.rkey::bigint ) INSERT INTO posts ( - at_uri, cid_digest, actor_id, record, content, facets, languages, tags, - parent_record_id, root_record_id, embed, embed_subtype, mentions, - violates_threadgate, created_at + record_id, content, language_code, tags, + parent_record_id, root_record_id, embed_type, embed_subtype, + violates_threadgate ) SELECT - $1, -- at_uri - DECODE(SUBSTRING($3 FROM 9), 'hex'), -- cid_digest (strip 0x01711220 header, convert hex string to bytea) - (SELECT id FROM actor_id_lookup), -- actor_id - $4, -- record + (SELECT id FROM record_lookup), -- record_id $5, -- content - $6, -- facets - $7, -- languages + CASE WHEN array_length($7::text[], 1) > 0 THEN ($7::text[])[1]::language_code ELSE NULL END, -- language_code (first language or NULL) $8, -- tags (SELECT id FROM parent_lookup), -- parent_record_id (SELECT id FROM root_lookup), -- root_record_id - $11, -- embed - $12, -- embed_subtype - $13, -- mentions - $14, -- violates_threadgate - $15 -- created_at + $11::embed_type, -- embed_type + $12::embed_type, -- embed_subtype + $14 -- violates_threadgate -- Posts are immutable - cannot be edited in AT Protocol -ON CONFLICT (at_uri) DO NOTHING \ No newline at end of file +ON CONFLICT (record_id) 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 index bdbab7fe..022294b0 100644 --- a/consumer/src/db/sql/postgate_upsert.sql +++ b/consumer/src/db/sql/postgate_upsert.sql @@ -1,7 +1,38 @@ -INSERT INTO postgates (at_uri, cid, post_uri, detached, rules, created_at) -VALUES ($1, $2, $3, $4, $5, $6) -ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, - post_uri=EXCLUDED.post_uri, - detached=EXCLUDED.detached, - rules=EXCLUDED.rules, - indexed_at=NOW() \ No newline at end of file +-- Insert/update postgate with CTE to resolve record_id and post_record_id +-- Parameters: $1=at_uri, $2=cid, $3=post_uri, $4=detached, $5=rules, $6=created_at +WITH record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint +), +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 ( + SELECT r.id + FROM post_parts pp + INNER JOIN actors a ON a.did = pp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = pp.collection + AND r.rkey = pp.rkey::bigint +) +INSERT INTO postgates (record_id, post_record_id, rules) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM post_lookup), + $5 +ON CONFLICT (record_id) DO UPDATE SET + post_record_id=EXCLUDED.post_record_id, + rules=EXCLUDED.rules diff --git a/consumer/src/db/sql/profile_upsert.sql b/consumer/src/db/sql/profile_upsert.sql index 80615215..b7ff7b4b 100644 --- a/consumer/src/db/sql/profile_upsert.sql +++ b/consumer/src/db/sql/profile_upsert.sql @@ -1,15 +1,76 @@ -INSERT INTO profiles (did, cid, avatar_cid, banner_cid, display_name, description, pinned_uri, pinned_cid, - joined_sp_uri, joined_sp_cid, pronouns, website, created_at) -VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13) -ON CONFLICT (did) DO UPDATE SET cid=EXCLUDED.cid, - avatar_cid=EXCLUDED.avatar_cid, - banner_cid=EXCLUDED.banner_cid, - display_name=EXCLUDED.display_name, - description=EXCLUDED.description, - pinned_uri=EXCLUDED.pinned_uri, - pinned_cid=EXCLUDED.pinned_cid, - joined_sp_uri=EXCLUDED.joined_sp_cid, - joined_sp_cid=EXCLUDED.joined_sp_cid, - pronouns=EXCLUDED.pronouns, - website=EXCLUDED.website, - indexed_at=NOW() \ No newline at end of file +-- Insert/update profile with CTE to resolve actor_id and pinned/joined_sp record IDs +-- Parameters: $1=did, $2=cid, $3=avatar_cid, $4=banner_cid, $5=display_name, $6=description, +-- $7=pinned_uri, $8=pinned_cid, $9=joined_sp_uri, $10=joined_sp_cid, +-- $11=pronouns, $12=website, $13=created_at +WITH actor_lookup AS ( + SELECT id FROM actors WHERE did = $1 +), +pinned_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($7 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($7 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($7 FROM 6), '/', 3) as rkey + WHERE $7 IS NOT NULL +), +pinned_lookup AS ( + SELECT r.id + FROM pinned_parts pp + INNER JOIN actors a ON a.did = pp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = pp.collection + AND r.rkey = pp.rkey::bigint +), +joined_sp_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($9 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($9 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($9 FROM 6), '/', 3) as rkey + WHERE $9 IS NOT NULL +), +joined_sp_lookup AS ( + SELECT r.id + FROM joined_sp_parts jsp + INNER JOIN actors a ON a.did = jsp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = jsp.collection + AND r.rkey = jsp.rkey::bigint +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING(('at://' || $1 || '/app.bsky.actor.profile/self') FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING(('at://' || $1 || '/app.bsky.actor.profile/self') FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING(('at://' || $1 || '/app.bsky.actor.profile/self') FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey +) +INSERT INTO profiles ( + actor_id, record_id, avatar_cid, banner_cid, display_name, description, + pinned_record_id, joined_sp_record_id, pronouns, website +) +SELECT + (SELECT id FROM actor_lookup), + (SELECT id FROM record_lookup), + $3, -- avatar_cid + $4, -- banner_cid + $5, -- display_name + $6, -- description + (SELECT id FROM pinned_lookup), + (SELECT id FROM joined_sp_lookup), + $11, -- pronouns + $12 -- website +ON CONFLICT (actor_id) DO UPDATE SET + record_id=EXCLUDED.record_id, + avatar_cid=EXCLUDED.avatar_cid, + banner_cid=EXCLUDED.banner_cid, + display_name=EXCLUDED.display_name, + description=EXCLUDED.description, + pinned_record_id=EXCLUDED.pinned_record_id, + joined_sp_record_id=EXCLUDED.joined_sp_record_id, + pronouns=EXCLUDED.pronouns, + website=EXCLUDED.website diff --git a/consumer/src/db/sql/starterpack_upsert.sql b/consumer/src/db/sql/starterpack_upsert.sql index 325f32ca..58c57b0c 100644 --- a/consumer/src/db/sql/starterpack_upsert.sql +++ b/consumer/src/db/sql/starterpack_upsert.sql @@ -1,11 +1,48 @@ -INSERT INTO starterpacks (at_uri, owner, cid, record, name, description, description_facets, list, feeds, created_at) -VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) -ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, - record=EXCLUDED.record, - name=EXCLUDED.name, - description=EXCLUDED.description, - description_facets=EXCLUDED.description_facets, - list=EXCLUDED.list, - feeds=EXCLUDED.feeds, - indexed_at=NOW() -RETURNING XMAX::text::int \ No newline at end of file +-- Insert/update starterpack with CTE to resolve record_id, owner_actor_id, and list_record_id +-- Parameters: $1=at_uri, $2=owner, $3=cid, $4=record(unused), $5=name, $6=description, +-- $7=description_facets, $8=list, $9=feeds(unused), $10=created_at +WITH owner_lookup AS ( + SELECT id FROM actors WHERE did = $2 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint +), +list_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($8 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($8 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($8 FROM 6), '/', 3) as rkey +), +list_lookup AS ( + SELECT r.id + FROM list_parts lp + INNER JOIN actors a ON a.did = lp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = lp.collection + AND r.rkey = lp.rkey::bigint +) +INSERT INTO starterpacks (record_id, owner_actor_id, name, description, description_facets, list_record_id) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM owner_lookup), + $5, + $6, + $7, + (SELECT id FROM list_lookup) +ON CONFLICT (record_id) DO UPDATE SET + name=EXCLUDED.name, + description=EXCLUDED.description, + description_facets=EXCLUDED.description_facets, + list_record_id=EXCLUDED.list_record_id +RETURNING XMAX::text::int diff --git a/consumer/src/db/sql/status_upsert.sql b/consumer/src/db/sql/status_upsert.sql index 71012475..3afb976e 100644 --- a/consumer/src/db/sql/status_upsert.sql +++ b/consumer/src/db/sql/status_upsert.sql @@ -1,12 +1,54 @@ -INSERT INTO statuses (did, status, duration, record, embed_uri, embed_title, embed_description, thumb_mime_type, - thumb_cid, created_at) -VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) -ON CONFLICT (did) DO UPDATE SET status=EXCLUDED.status, - duration=EXCLUDED.duration, - record=EXCLUDED.record, - embed_uri=EXCLUDED.embed_uri, - embed_title=EXCLUDED.embed_title, - embed_description=EXCLUDED.embed_description, - thumb_mime_type=EXCLUDED.thumb_mime_type, - thumb_cid=EXCLUDED.thumb_cid, - indexed_at=NOW() \ No newline at end of file +-- Insert/update status with CTE to resolve record_id and embed_record_id +-- Parameters: $1=did, $2=status, $3=duration, $4=record(unused), $5=embed_uri, +-- $6=embed_title(unused), $7=embed_description(unused), +-- $8=thumb_mime_type, $9=thumb_cid, $10=created_at +WITH actor_lookup AS ( + SELECT id FROM actors WHERE did = $1 +), +record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING(('at://' || $1 || '/chat.bsky.actor.status/self') FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING(('at://' || $1 || '/chat.bsky.actor.status/self') FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING(('at://' || $1 || '/chat.bsky.actor.status/self') FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey +), +embed_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($5 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($5 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($5 FROM 6), '/', 3) as rkey + WHERE $5 IS NOT NULL +), +embed_lookup AS ( + SELECT r.id + FROM embed_parts ep + INNER JOIN actors a ON a.did = ep.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = ep.collection + AND r.rkey = ep.rkey::bigint +) +INSERT INTO statuses ( + record_id, actor_id, status, duration, embed_record_id, + thumb_mime_type, thumb_cid +) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM actor_lookup), + $2::status_type, + $3, + (SELECT id FROM embed_lookup), + $8::image_mime_type, + $9 +ON CONFLICT (record_id) DO UPDATE SET + status=EXCLUDED.status, + duration=EXCLUDED.duration, + embed_record_id=EXCLUDED.embed_record_id, + thumb_mime_type=EXCLUDED.thumb_mime_type, + thumb_cid=EXCLUDED.thumb_cid diff --git a/consumer/src/db/sql/threadgate_upsert.sql b/consumer/src/db/sql/threadgate_upsert.sql index 75596404..dd2c7833 100644 --- a/consumer/src/db/sql/threadgate_upsert.sql +++ b/consumer/src/db/sql/threadgate_upsert.sql @@ -1,8 +1,38 @@ -INSERT INTO threadgates (at_uri, cid, post_uri, hidden_replies, allow, allowed_lists, record, created_at) -VALUES ($1, $2, $3, $4, $5, $6, $7, $8) -ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, - hidden_replies=EXCLUDED.hidden_replies, - allow=EXCLUDED.allow, - allowed_lists=EXCLUDED.allowed_lists, - record=EXCLUDED.record, - indexed_at=NOW() \ No newline at end of file +-- Insert/update threadgate with CTE to resolve record_id and post_record_id +-- Parameters: $1=at_uri, $2=cid, $3=post_uri, $4=hidden_replies, $5=allow, +-- $6=allowed_lists, $7=record(unused), $8=created_at +WITH record_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 2) as collection, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey +), +record_lookup AS ( + SELECT r.id + FROM record_parts rp + INNER JOIN actors a ON a.did = rp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = rp.collection + AND r.rkey = rp.rkey::bigint +), +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 ( + SELECT r.id + FROM post_parts pp + INNER JOIN actors a ON a.did = pp.did + INNER JOIN records r ON r.actor_id = a.id + AND r.collection::text = pp.collection + AND r.rkey = pp.rkey::bigint +) +INSERT INTO threadgates (record_id, post_record_id, allow) +SELECT + (SELECT id FROM record_lookup), + (SELECT id FROM post_lookup), + $5 +ON CONFLICT (record_id) DO UPDATE SET + allow=EXCLUDED.allow diff --git a/parakeet-db/Cargo.toml b/parakeet-db/Cargo.toml index fa59ef20..aaf7b78b 100644 --- a/parakeet-db/Cargo.toml +++ b/parakeet-db/Cargo.toml @@ -5,6 +5,7 @@ edition = "2021" [dependencies] bloomfilter = "1.0.14" +bytes = "1" chrono = { version = "0.4.39", features = ["serde"] } cid = "0.11" diesel = { version = "2.2.6", features = ["chrono", "serde_json", "postgres"] } diff --git a/parakeet-db/src/actor_cache.rs b/parakeet-db/src/actor_cache.rs index 71bb4b02..b9acc9ca 100644 --- a/parakeet-db/src/actor_cache.rs +++ b/parakeet-db/src/actor_cache.rs @@ -173,7 +173,7 @@ impl CachedActorStore { shard .get(did) - .map(|actor| AllowlistStatus::from(actor.sync_state.clone())) + .map(|actor| AllowlistStatus::from(actor.sync_state)) .unwrap_or(AllowlistStatus::NotAllowed) } @@ -214,7 +214,7 @@ impl CachedActorStore { results[idx] = shard .get(did) - .map(|actor| AllowlistStatus::from(actor.sync_state.clone())) + .map(|actor| AllowlistStatus::from(actor.sync_state)) .unwrap_or(AllowlistStatus::NotAllowed); } @@ -250,7 +250,7 @@ impl CachedActorStore { let shard = self.shards[shard_idx].read().expect("Failed to read shard"); if let Some(actor) = shard.get(did) { - let status = AllowlistStatus::from(actor.sync_state.clone()); + let status = AllowlistStatus::from(actor.sync_state); if status.is_fully_allowed() { return true; } diff --git a/parakeet-db/src/at_uri_util.rs b/parakeet-db/src/at_uri_util.rs index 4693f869..5c7e9bea 100644 --- a/parakeet-db/src/at_uri_util.rs +++ b/parakeet-db/src/at_uri_util.rs @@ -8,7 +8,6 @@ /// - rkey (TEXT) /// /// These utilities help convert between the formats. - /// Parse an AT URI into its components /// /// # Arguments diff --git a/parakeet-db/src/cid_util.rs b/parakeet-db/src/cid_util.rs index 8b65d2c6..b084ca69 100644 --- a/parakeet-db/src/cid_util.rs +++ b/parakeet-db/src/cid_util.rs @@ -23,7 +23,6 @@ /// (blobs use raw codec, records use dag-cbor codec). /// /// Reference: https://atproto.com/specs/data-model - /// CIDv1 header for blob CIDs (images, videos, etc.) /// 0x01 (CIDv1) + 0x55 (raw codec) + 0x1220 (SHA-256 multihash header) pub const BLOB_CID_HEADER: [u8; 4] = [0x01, 0x55, 0x12, 0x20]; diff --git a/parakeet-db/src/types.rs b/parakeet-db/src/types.rs index 3cd4cee8..38783091 100644 --- a/parakeet-db/src/types.rs +++ b/parakeet-db/src/types.rs @@ -96,6 +96,41 @@ macro_rules! diesel_enum { Ok(diesel::serialize::IsNull::No) } } + + // postgres-types impls for consumer (uses deadpool-postgres/tokio-postgres) + impl postgres_types::ToSql for $rust_name { + fn to_sql( + &self, + _ty: &postgres_types::Type, + out: &mut postgres_types::private::BytesMut, + ) -> Result> { + use bytes::BufMut; + let s = self.to_string(); + out.put_slice(s.as_bytes()); + Ok(postgres_types::IsNull::No) + } + + fn accepts(ty: &postgres_types::Type) -> bool { + ty.name() == stringify!($pg_type) || ty.name() == "text" + } + + postgres_types::to_sql_checked!(); + } + + impl<'a> postgres_types::FromSql<'a> for $rust_name { + fn from_sql( + _ty: &postgres_types::Type, + raw: &'a [u8], + ) -> Result> { + use std::str::FromStr; + let s = std::str::from_utf8(raw)?; + Self::from_str(s).map_err(|e| e.into()) + } + + fn accepts(ty: &postgres_types::Type) -> bool { + ty.name() == stringify!($pg_type) || ty.name() == "text" + } + } }; } -- 2.51.2