diff --git a/consumer/src/db/operations/feed/helpers.rs b/consumer/src/db/operations/feed/helpers.rs index 8e285f8f..2e2a4899 100644 --- a/consumer/src/db/operations/feed/helpers.rs +++ b/consumer/src/db/operations/feed/helpers.rs @@ -324,12 +324,19 @@ pub(crate) async fn get_repost_id( conn: &C, actor_id: i32, rkey: &str, - _cid_str: &str, + cid_str: &str, ) -> Result<(i64, bool)> { // Convert rkey (TID string) to INT8 for database lookup let rkey_i64 = parakeet_db::models::tid_to_i64(rkey) .wrap_err_with(|| format!("Invalid TID encoding in rkey: {}", rkey))?; + // Parse the CID string to get the digest + let cid = ipld_core::cid::Cid::try_from(cid_str) + .wrap_err_with(|| format!("Invalid CID format: {}", cid_str))?; + let cid_bytes = cid.to_bytes(); + let cid_digest = parakeet_db::cid_util::cid_to_digest(&cid_bytes) + .ok_or_eyre("CID must be valid AT Protocol CID")?; + // Acquire advisory lock using actor_id and rkey to prevent concurrent access races let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("reposts", actor_id, rkey_i64); conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) @@ -338,15 +345,15 @@ pub(crate) async fn get_repost_id( // Use CTE to SELECT first, then conditionally INSERT only if not found // This prevents unnecessary sequence consumption when repost stub already exists // Still uses ON CONFLICT for race condition safety between concurrent transactions - // Note: CID is synthetic, generated from actor_id + rkey + // Note: CID is from the via field's StrongRef (real CID from referencing record) let row = conn .query_one( "WITH existing AS ( SELECT id FROM reposts WHERE actor_id = $1 AND rkey = $2 ), inserted AS ( - INSERT INTO reposts (actor_id, rkey, post_id, via_repost_id, status) - SELECT $1, $2, NULL, NULL, 'stub'::repost_status + INSERT INTO reposts (actor_id, rkey, cid, post_id, via_repost_id, status) + SELECT $1, $2, $3, NULL, NULL, 'stub'::repost_status WHERE NOT EXISTS (SELECT 1 FROM existing) ON CONFLICT (actor_id, rkey) DO UPDATE SET actor_id = EXCLUDED.actor_id RETURNING id @@ -354,7 +361,7 @@ pub(crate) async fn get_repost_id( SELECT COALESCE((SELECT id FROM existing), (SELECT id FROM inserted)) as id, (SELECT id FROM existing) IS NULL as was_created", - &[&actor_id, &rkey_i64], + &[&actor_id, &rkey_i64, &cid_digest], ) .await?; diff --git a/consumer/src/db/operations/feed/repost.rs b/consumer/src/db/operations/feed/repost.rs index a392eeba..1cc40234 100644 --- a/consumer/src/db/operations/feed/repost.rs +++ b/consumer/src/db/operations/feed/repost.rs @@ -30,11 +30,15 @@ pub async fn repost_insert( conn: &C, rkey: i64, actor_id: i32, - _cid: Cid, + cid: Cid, rec: AppBskyFeedRepost, via_repost_id: Option, source: crate::database_writer::EventSource, ) -> Result { + let cid_bytes = cid.to_bytes(); + let cid_digest = parakeet_db::cid_util::cid_to_digest(&cid_bytes) + .expect("CID must be valid AT Protocol CID"); + let subject_uri = &rec.subject.uri; let subject_cid_str = rec.subject.cid.to_string(); @@ -57,17 +61,18 @@ pub async fn repost_insert( // If a stub exists, upgrade it to complete status and fill in the data // If it's a new repost, insert with status='complete' // Note: created_at is derived from TID rkey - // Note: CID is synthetic, generated from actor_id + rkey + // Note: CID is real (from database) - reposts can be referenced via like.via field let rows = conn .execute( - "INSERT INTO reposts (actor_id, rkey, post_id, via_repost_id, status) - VALUES ($1, $2, $3, $4, 'complete'::repost_status) + "INSERT INTO reposts (actor_id, rkey, cid, post_id, via_repost_id, status) + VALUES ($1, $2, $3, $4, $5, 'complete'::repost_status) ON CONFLICT (actor_id, rkey) DO UPDATE SET + cid = EXCLUDED.cid, post_id = EXCLUDED.post_id, via_repost_id = EXCLUDED.via_repost_id, status = 'complete'::repost_status WHERE reposts.status = 'stub'", - &[&actor_id, &rkey, &post_id, &via_repost_id], + &[&actor_id, &rkey, &cid_digest, &post_id, &via_repost_id], ) .await?; diff --git a/migrations/2025-11-16-014047_drop_reposts_cid/down.sql b/migrations/2025-11-16-014047_drop_reposts_cid/down.sql deleted file mode 100644 index 9d493f38..00000000 --- a/migrations/2025-11-16-014047_drop_reposts_cid/down.sql +++ /dev/null @@ -1,4 +0,0 @@ --- Rollback: Add cid column back --- NOTE: On rollback, you would need to re-index from the relay to populate real CIDs --- or accept synthetic CIDs (which is what we're using anyway) -ALTER TABLE reposts ADD COLUMN cid BYTEA NOT NULL DEFAULT E'\\x0000000000000000000000000000000000000000000000000000000000000000'::bytea; diff --git a/migrations/2025-11-16-014047_drop_reposts_cid/up.sql b/migrations/2025-11-16-014047_drop_reposts_cid/up.sql deleted file mode 100644 index 533b4ae8..00000000 --- a/migrations/2025-11-16-014047_drop_reposts_cid/up.sql +++ /dev/null @@ -1,4 +0,0 @@ --- Drop the cid column from reposts table --- Repost CIDs are now generated synthetically from (actor_id, rkey) --- This enables massive compression gains (14-23x) when using TimescaleDB -ALTER TABLE reposts DROP COLUMN cid; diff --git a/parakeet-db/src/models.rs b/parakeet-db/src/models.rs index 964d9ec1..15bdfbbe 100644 --- a/parakeet-db/src/models.rs +++ b/parakeet-db/src/models.rs @@ -299,8 +299,8 @@ pub struct Repost { pub post_id: Option, // FK to posts (nullable for pending linkage) pub via_repost_id: Option, // FK to reposts (repost context: "repost via repost") pub status: RepostStatus, // ENUM: complete | stub + pub cid: Vec, // Real CID from database (reposts can be referenced via like.via field) // Note: created_at derived from TID rkey via created_at() method - // Note: CID is synthetic, generated from actor_id + rkey } #[derive(Clone, Debug, Queryable, Selectable, Identifiable)] diff --git a/parakeet-db/src/schema.rs b/parakeet-db/src/schema.rs index 539fcd58..00572066 100644 --- a/parakeet-db/src/schema.rs +++ b/parakeet-db/src/schema.rs @@ -616,6 +616,7 @@ diesel::table! { post_id -> Nullable, via_repost_id -> Nullable, status -> RepostStatus, + cid -> Bytea, } } diff --git a/parakeet/src/db/bookmarks.rs b/parakeet/src/db/bookmarks.rs index 7dca47d5..2c39588c 100644 --- a/parakeet/src/db/bookmarks.rs +++ b/parakeet/src/db/bookmarks.rs @@ -62,18 +62,18 @@ pub async fn get_user_bookmarks( created_at: chrono::DateTime, #[diesel(sql_type = Text)] did: String, - #[diesel(sql_type = Integer)] - actor_id: i32, #[diesel(sql_type = BigInt)] rkey: i64, + #[diesel(sql_type = diesel::sql_types::Binary)] + cid: Vec, } // Use .bind() for cursor parameter to prevent SQL injection diesel::sql_query( "SELECT tid_timestamp(b.rkey) as created_at, a.did, - p.actor_id, - p.rkey + p.rkey, + p.cid FROM bookmarks b INNER JOIN posts p ON b.subject_id = p.id INNER JOIN actors a ON p.actor_id = a.id @@ -94,8 +94,9 @@ pub async fn get_user_bookmarks( .map(|r| { let encoded_rkey = parakeet_db::tid_util::encode_tid(r.rkey); let subject_uri = format!("at://{}/app.bsky.feed.post/{}", r.did, encoded_rkey); - // Generate synthetic CID from actor_id and rkey - let cid_str = parakeet_db::cid_util::post_cid_string(r.actor_id, r.rkey); + // Convert real CID from database to string + let cid_str = parakeet_db::cid_util::digest_to_record_cid_string(&r.cid) + .unwrap_or_else(|| String::from("bafyrei_invalid_cid")); (r.created_at, subject_uri, cid_str) }) .collect() diff --git a/parakeet/src/db/states.rs b/parakeet/src/db/states.rs index dff58e5a..bee40e71 100644 --- a/parakeet/src/db/states.rs +++ b/parakeet/src/db/states.rs @@ -96,6 +96,8 @@ pub struct PostStateRet { pub actor_id: i32, #[diesel(sql_type = diesel::sql_types::BigInt)] pub rkey: i64, + #[diesel(sql_type = diesel::sql_types::Binary)] + pub cid: Vec, #[diesel(sql_type = Nullable)] pub like_rkey: Option, #[diesel(sql_type = Nullable)] @@ -115,9 +117,10 @@ impl PostStateRet { format!("at://{}/app.bsky.feed.post/{}", self.did, encoded_rkey) } - /// Get the CID as a string (generated synthetically) + /// Get the CID as a string (from database) pub fn cid_string(&self) -> String { - parakeet_db::cid_util::post_cid_string(self.actor_id, self.rkey) + parakeet_db::cid_util::digest_to_record_cid_string(&self.cid) + .unwrap_or_else(|| String::from("bafyrei_invalid_cid")) } /// Get the like rkey as a base32 TID string if present diff --git a/parakeet/src/loaders/post.rs b/parakeet/src/loaders/post.rs index 5274eb39..352b950f 100644 --- a/parakeet/src/loaders/post.rs +++ b/parakeet/src/loaders/post.rs @@ -445,36 +445,39 @@ impl BatchFn for PostLoader { let encoded_rkey = parakeet_db::tid_util::encode_tid(p.rkey); let at_uri = format!("at://{}/app.bsky.feed.post/{}", p.author_did, encoded_rkey); - // Generate synthetic CID from actor_id and rkey - let cid_str = parakeet_db::cid_util::post_cid_string(p.actor_id, p.rkey); + // Convert real CID from database to string + let cid_str = parakeet_db::cid_util::digest_to_record_cid_string(&p.cid) + .unwrap_or_else(|| String::from("bafyrei_invalid_cid")); - // Encode parent URI and CID if present (generate synthetic CID) + // Encode parent URI and CID if present (use real CID from database) let (parent_uri, parent_cid) = - if let (Some(parent_did), Some(parent_actor_id), Some(parent_rkey)) = - (&p.parent_did, p.parent_actor_id, p.parent_rkey) + if let (Some(parent_did), Some(parent_rkey), Some(parent_cid_bytes)) = + (&p.parent_did, p.parent_rkey, &p.parent_cid_bytes) { let encoded_parent_rkey = parakeet_db::tid_util::encode_tid(parent_rkey); let parent_uri_str = format!( "at://{}/app.bsky.feed.post/{}", parent_did, encoded_parent_rkey ); - // Generate synthetic CID for parent post - let parent_cid_str = parakeet_db::cid_util::post_cid_string(parent_actor_id, parent_rkey); + // Convert real CID from database to string + let parent_cid_str = parakeet_db::cid_util::digest_to_record_cid_string(parent_cid_bytes) + .unwrap_or_else(|| String::from("bafyrei_invalid_cid")); (Some(parent_uri_str), Some(parent_cid_str)) } else { (None, None) }; - // Encode root URI and CID if present (generate synthetic CID) + // Encode root URI and CID if present (use real CID from database) let (root_uri, root_cid) = - if let (Some(root_did), Some(root_actor_id), Some(root_rkey)) = - (&p.root_did, p.root_actor_id, p.root_rkey) + if let (Some(root_did), Some(root_rkey), Some(root_cid_bytes)) = + (&p.root_did, p.root_rkey, &p.root_cid_bytes) { let encoded_root_rkey = parakeet_db::tid_util::encode_tid(root_rkey); let root_uri_str = format!("at://{}/app.bsky.feed.post/{}", root_did, encoded_root_rkey); - // Generate synthetic CID for root post - let root_cid_str = parakeet_db::cid_util::post_cid_string(root_actor_id, root_rkey); + // Convert real CID from database to string + let root_cid_str = parakeet_db::cid_util::digest_to_record_cid_string(root_cid_bytes) + .unwrap_or_else(|| String::from("bafyrei_invalid_cid")); (Some(root_uri_str), Some(root_cid_str)) } else { (None, None) diff --git a/parakeet/src/sql/post_state.sql b/parakeet/src/sql/post_state.sql index 283be901..c461d744 100644 --- a/parakeet/src/sql/post_state.sql +++ b/parakeet/src/sql/post_state.sql @@ -10,6 +10,7 @@ FROM ( a.did, p.actor_id, p.rkey, + p.cid, l.rkey as like_rkey, r.rkey as repost_rkey, b.actor_id IS NOT NULL as bookmarked, diff --git a/parakeet/src/xrpc/app_bsky/notification/mod.rs b/parakeet/src/xrpc/app_bsky/notification/mod.rs index 019e2482..5b6f8272 100644 --- a/parakeet/src/xrpc/app_bsky/notification/mod.rs +++ b/parakeet/src/xrpc/app_bsky/notification/mod.rs @@ -294,20 +294,9 @@ pub async fn list_notifications( notifications.push(Notification { uri, - cid: match notif.record_type { - parakeet_db::types::NotificationRecordType::Like => { - parakeet_db::cid_util::like_cid_string(notif.author_actor_id, notif.record_rkey) - } - parakeet_db::types::NotificationRecordType::Repost => { - parakeet_db::cid_util::repost_cid_string(notif.author_actor_id, notif.record_rkey) - } - parakeet_db::types::NotificationRecordType::Follow => { - parakeet_db::cid_util::follow_cid_string(notif.author_actor_id, notif.record_rkey) - } - parakeet_db::types::NotificationRecordType::Post => { - parakeet_db::cid_util::post_cid_string(notif.author_actor_id, notif.record_rkey) - } - }, + // Use the real CID from database (already stored for all record types) + cid: parakeet_db::cid_util::digest_to_record_cid_string(¬if.record_cid) + .unwrap_or_else(|| String::from("bafyrei_invalid_cid")), author, reason: notif.reason.to_string(), reason_subject, diff --git a/parakeet/src/xrpc/com_atproto/repo.rs b/parakeet/src/xrpc/com_atproto/repo.rs index 33da4135..6b2c77cc 100644 --- a/parakeet/src/xrpc/com_atproto/repo.rs +++ b/parakeet/src/xrpc/com_atproto/repo.rs @@ -9,8 +9,8 @@ use serde_json::Value; #[derive(diesel::QueryableByName)] struct PostRecord { - #[diesel(sql_type = diesel::sql_types::Integer)] - actor_id: i32, + #[diesel(sql_type = diesel::sql_types::Binary)] + cid: Vec, #[diesel(sql_type = diesel::sql_types::Jsonb)] record: Value, } @@ -89,10 +89,10 @@ pub async fn get_record( let rkey_bigint = parakeet_db::tid_util::decode_tid(&query.rkey) .map_err(|_| Error::invalid_request(Some("Invalid rkey format".to_string())))?; - // Use raw SQL to get actor_id and record (CID generated synthetically) + // Use raw SQL to get CID and record from posts table let result: PostRecord = diesel_async::RunQueryDsl::get_result( diesel::sql_query( - "SELECT r.actor_id, record + "SELECT posts.cid, record FROM posts INNER JOIN records r ON posts.record_id = r.id INNER JOIN actors a ON r.actor_id = a.id @@ -107,8 +107,9 @@ pub async fn get_record( ) .await?; - // Generate synthetic CID from actor_id and rkey - let cid_str = parakeet_db::cid_util::post_cid_string(result.actor_id, rkey_bigint); + // Convert real CID from database to string + let cid_str = parakeet_db::cid_util::digest_to_record_cid_string(&result.cid) + .unwrap_or_else(|| String::from("bafyrei_invalid_cid")); (cid_str, result.record) }