diff --git a/consumer/src/database_writer/operations/executor.rs b/consumer/src/database_writer/operations/executor.rs index 1a4baefa..af61cc03 100644 --- a/consumer/src/database_writer/operations/executor.rs +++ b/consumer/src/database_writer/operations/executor.rs @@ -788,7 +788,7 @@ pub async fn execute_operation( cid, record, } => { - let inserted = db::list_upsert(conn, actor_id, rkey, cid, record).await?; + let inserted = db::list_upsert(conn, actor_id, &rkey, cid, record).await?; if inserted { let actor_deltas = actor_list_delta(actor_id, 1); Ok((vec![], actor_deltas, vec![])) @@ -1078,7 +1078,10 @@ pub async fn execute_operation( db::threadgate_delete(conn, actor_id, rkey).await?; } CollectionType::BskyList => { - let rows = db::list_delete(conn, actor_id, rkey).await?; + // Lists use arbitrary string rkeys - extract from at_uri + let rkey_str = parakeet_db::at_uri_util::extract_rkey(&at_uri) + .ok_or_else(|| eyre::eyre!("Invalid at_uri: missing rkey"))?; + let rows = db::list_delete(conn, actor_id, rkey_str).await?; if rows > 0 { actor_deltas.extend(actor_list_delta(actor_id, -1)); } diff --git a/consumer/src/database_writer/operations/handlers/list.rs b/consumer/src/database_writer/operations/handlers/list.rs index f73da55b..05007fe2 100644 --- a/consumer/src/database_writer/operations/handlers/list.rs +++ b/consumer/src/database_writer/operations/handlers/list.rs @@ -11,20 +11,19 @@ pub fn handle_list( let mut operations = Vec::new(); let mut cache_invalidations = Vec::new(); - // Validate TID timestamp is within acceptable range - if let Err(e) = crate::database_writer::validate_tid_timestamp(&ctx.rkey) { - tracing::warn!("Invalid list TID timestamp: {}", e); - return (vec![], vec![]); + // Lists can use both TID and arbitrary string rkeys (like "nfb", "bblock") + // If it's a TID, validate the timestamp is within acceptable range + if parakeet_db::models::tid_to_i64(&ctx.rkey).is_ok() { + if let Err(e) = crate::database_writer::validate_tid_timestamp(&ctx.rkey) { + tracing::warn!("Invalid list TID timestamp: {}", e); + return (vec![], vec![]); + } } - // Convert TID string to i64 (safe because validation passed above) - let rkey_i64 = parakeet_db::models::tid_to_i64(&ctx.rkey) - .expect("TID validation passed, conversion should succeed"); - let labels = record.labels.clone(); operations.push(DatabaseOperation::UpsertList { - rkey: rkey_i64, + rkey: ctx.rkey.clone(), actor_id: ctx.actor_id, cid: ctx.cid, record, diff --git a/consumer/src/database_writer/operations/types.rs b/consumer/src/database_writer/operations/types.rs index 3c960f3d..4fc7344a 100644 --- a/consumer/src/database_writer/operations/types.rs +++ b/consumer/src/database_writer/operations/types.rs @@ -133,7 +133,7 @@ pub enum DatabaseOperation { // Lists UpsertList { - rkey: i64, // TID converted to i64 + rkey: String, // Can be TID or arbitrary string like "nfb", "bblock" actor_id: i32, cid: Cid, record: AppBskyGraphList, diff --git a/consumer/src/db/gates/queries.rs b/consumer/src/db/gates/queries.rs index 03a86ffe..f4d602ea 100644 --- a/consumer/src/db/gates/queries.rs +++ b/consumer/src/db/gates/queries.rs @@ -90,7 +90,7 @@ pub async fn check_list_membership( SELECT l.id 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 = tid_to_i64(lp.rkey) + INNER JOIN lists l ON l.actor_id = a.id AND l.rkey = lp.rkey ) SELECT count(*) FROM list_items li diff --git a/consumer/src/db/operations/feed/helpers.rs b/consumer/src/db/operations/feed/helpers.rs index 5c967533..721c81cb 100644 --- a/consumer/src/db/operations/feed/helpers.rs +++ b/consumer/src/db/operations/feed/helpers.rs @@ -188,18 +188,8 @@ pub(super) async fn get_list_id(conn: &C, at_uri: &str) -> Res let rkey = parakeet_db::at_uri_util::extract_rkey(at_uri) .ok_or_else(|| eyre::eyre!("Invalid AT URI: missing rkey in {}", at_uri))?; - // Try to convert rkey to TID - // Some old lists (like Blockenheimer) use non-TID rkeys like "bblock" - // We don't support these, so log a warning and skip them - let rkey_i64 = parakeet_db::models::tid_to_i64(rkey) - .inspect_err(|_e| { - tracing::warn!( - at_uri = %at_uri, - rkey = %rkey, - "Skipping list with non-TID rkey (not supported)" - ); - }) - .wrap_err_with(|| format!("List uses non-TID rkey which is not supported: {}", at_uri))?; + // Lists can use both TID rkeys (like "3jzfcijpj2z2a") and arbitrary string rkeys (like "nfb", "bblock") + // We support both now that lists.rkey is text // Get/create actor_id for owner let (actor_id, _, _) = get_actor_id(conn, did).await?; @@ -209,7 +199,7 @@ pub(super) async fn get_list_id(conn: &C, at_uri: &str) -> Res // Still uses ON CONFLICT for race condition safety between concurrent transactions // For stubs, owner_actor_id equals actor_id, list_type and name are NULL // Use zero CID as placeholder (will be updated when real list is fetched) - // Note: created_at is derived from TID rkey + // Note: created_at is derived from TID rkey if present, or set to epoch for non-TID rkeys let zero_cid = vec![0_u8; 32]; let row = conn .query_one( @@ -226,7 +216,7 @@ pub(super) async fn get_list_id(conn: &C, at_uri: &str) -> Res 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, &zero_cid], + &[&actor_id, &rkey, &zero_cid], ) .await?; diff --git a/consumer/src/db/operations/graph.rs b/consumer/src/db/operations/graph.rs index 60f44ece..394a3371 100644 --- a/consumer/src/db/operations/graph.rs +++ b/consumer/src/db/operations/graph.rs @@ -100,7 +100,7 @@ pub async fn block_delete(conn: &C, rkey: i64, actor_id: i32) pub async fn list_upsert( conn: &C, actor_id: i32, - rkey: i64, + rkey: &str, cid: Cid, rec: AppBskyGraphList, ) -> Result { @@ -113,7 +113,7 @@ pub async fn list_upsert( let avatar = blob_to_cid_bytes(rec.avatar); // Acquire table-scoped advisory lock on lists table for this record - let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("lists", actor_id, rkey); + let (table_id, key_id) = crate::database_writer::locking::actor_record_lock_str("lists", actor_id, rkey); conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) .await?; @@ -135,7 +135,7 @@ pub async fn list_upsert( .map(|r| r.get::<_, i32>(0) == 0) } -pub async fn list_delete(conn: &C, actor_id: i32, rkey: i64) -> Result { +pub async fn list_delete(conn: &C, actor_id: i32, rkey: &str) -> Result { conn.execute( "DELETE FROM lists WHERE actor_id = $1 diff --git a/consumer/src/db/record_exists/mod.rs b/consumer/src/db/record_exists/mod.rs index 3f9b9070..343778f3 100644 --- a/consumer/src/db/record_exists/mod.rs +++ b/consumer/src/db/record_exists/mod.rs @@ -70,10 +70,8 @@ pub async fn record_exists(conn: &PgObject, at_uri: &str) -> Result { "app.bsky.feed.generator" => queries::feedgen_exists(&***conn, did, rkey).await?, // Lists: Check lists table directly by (actor_id, rkey) - "app.bsky.graph.list" => { - let rkey_i64 = decode_tid(rkey).map_err(|e| eyre::eyre!("Invalid TID: {}", e))?; - queries::list_exists(&***conn, did, rkey_i64).await? - } + // Lists can use both TID and arbitrary string rkeys + "app.bsky.graph.list" => queries::list_exists(&***conn, did, rkey).await?, // List items: Check list_items table directly by (actor_id, rkey) "app.bsky.graph.listitem" => { diff --git a/consumer/src/db/record_exists/queries.rs b/consumer/src/db/record_exists/queries.rs index 652583ba..f84e3312 100644 --- a/consumer/src/db/record_exists/queries.rs +++ b/consumer/src/db/record_exists/queries.rs @@ -122,7 +122,7 @@ pub async fn feedgen_exists( /// - complete: Return true (already have data, skip) /// - deleted: Return true (soft-deleted, skip) /// - forbidden: Return true (access denied, skip) -pub async fn list_exists(conn: &C, did: &str, rkey: i64) -> QueryResult { +pub async fn list_exists(conn: &C, did: &str, rkey: &str) -> QueryResult { let row = conn .query_opt( "SELECT 1 FROM lists l diff --git a/migrations/2025-12-04-233127_change_lists_rkey_to_text/down.sql b/migrations/2025-12-04-233127_change_lists_rkey_to_text/down.sql new file mode 100644 index 00000000..d564ec34 --- /dev/null +++ b/migrations/2025-12-04-233127_change_lists_rkey_to_text/down.sql @@ -0,0 +1,8 @@ +-- Revert lists.rkey from text back to bigint +-- WARNING: This will fail if any lists have non-TID rkeys (like "nfb", "bblock") +-- Only use this if you're certain all lists use TID rkeys + +-- Convert text rkeys back to bigint (will fail on non-TID rkeys) +ALTER TABLE lists ALTER COLUMN rkey TYPE bigint USING tid_to_i64(rkey); + +-- Constraints are already correct (the ALTER TYPE doesn't change them) diff --git a/migrations/2025-12-04-233127_change_lists_rkey_to_text/up.sql b/migrations/2025-12-04-233127_change_lists_rkey_to_text/up.sql new file mode 100644 index 00000000..6f67ea37 --- /dev/null +++ b/migrations/2025-12-04-233127_change_lists_rkey_to_text/up.sql @@ -0,0 +1,16 @@ +-- Change lists.rkey from bigint (TID) to text (allows both TIDs and arbitrary strings) +-- Lists in AT Protocol can use arbitrary string rkeys like "nfb", "bblock", etc. +-- This matches the feedgens table which also uses text for rkey + +-- Convert existing TID bigints to base32 string representation +ALTER TABLE lists ALTER COLUMN rkey TYPE text USING i64_to_tid(rkey); + +-- Update unique constraint to work with text type +ALTER TABLE lists DROP CONSTRAINT lists_actor_id_rkey_key; +ALTER TABLE lists ADD CONSTRAINT lists_actor_id_rkey_key UNIQUE (actor_id, rkey); + +-- Update foreign key constraints that reference (actor_id, rkey) +-- Note: threadgate_allowed_lists references lists via list_id (synthetic), not natural key +-- list_items references lists via list_id (synthetic), not natural key +-- list_blocks references lists via list_id (synthetic), not natural key +-- So no FK updates needed diff --git a/parakeet-db/src/models.rs b/parakeet-db/src/models.rs index 16148f49..e491fd11 100644 --- a/parakeet-db/src/models.rs +++ b/parakeet-db/src/models.rs @@ -523,7 +523,7 @@ pub struct FeedGen { pub struct List { pub id: i64, // PK: Surrogate key pub actor_id: i32, // FK to actors - pub rkey: i64, // TID as INT8 + pub rkey: String, // Can be TID or arbitrary string like "nfb", "bblock" pub cid: Vec, // 32-byte CID digest pub owner_actor_id: i32, // FK to actors (denormalized for filtering) pub list_type: Option, // ENUM: curatelist | modlist | referencelist (NULL for stubs) @@ -532,7 +532,7 @@ pub struct List { pub description_facets: Option, pub avatar_cid: Option>, // 32-byte CID pub status: RecordStatus, // ENUM: complete | stub | deleted | forbidden - // Note: created_at derived from TID rkey via created_at() method + // Note: created_at derived from TID rkey if present, or epoch for non-TID rkeys } #[derive(Debug, Queryable, Selectable, Identifiable)] @@ -846,8 +846,7 @@ impl_tid_rkey!(Repost); impl_tid_rkey!(Follow); impl_tid_rkey!(Block); impl_tid_rkey!(Bookmark); -// NOTE: FeedGen intentionally excluded - uses arbitrary String rkeys, not TIDs -impl_tid_rkey!(List); +// NOTE: FeedGen and List intentionally excluded - use arbitrary String rkeys, not TIDs impl_tid_rkey!(ListItem); impl_tid_rkey!(ListBlock); impl_tid_rkey!(StarterPack); @@ -868,9 +867,19 @@ impl_tid_created_at!(Repost); impl_tid_created_at!(Follow); impl_tid_created_at!(Block); impl_tid_created_at!(Bookmark); -impl_tid_created_at!(List); impl_tid_created_at!(ListItem); impl_tid_created_at!(ListBlock); +// NOTE: List uses String rkey and has custom created_at() implementation impl_tid_created_at!(StarterPack); impl_tid_created_at!(Threadgate); impl_tid_created_at!(Verification); + +// Custom created_at() implementation for List (String rkey can be TID or arbitrary) +impl List { + /// Get created_at timestamp derived from rkey if it's a TID, or epoch for non-TID rkeys + pub fn created_at(&self) -> DateTime { + crate::models::tid_to_i64(&self.rkey) + .map(crate::tid_util::tid_to_datetime) + .unwrap_or_else(|_| DateTime::::UNIX_EPOCH) + } +} diff --git a/parakeet-db/src/schema.rs b/parakeet-db/src/schema.rs index 7deb974f..251824a9 100644 --- a/parakeet-db/src/schema.rs +++ b/parakeet-db/src/schema.rs @@ -420,7 +420,7 @@ diesel::table! { lists (id) { id -> Int8, actor_id -> Int4, - rkey -> Int8, + rkey -> Text, cid -> Bytea, owner_actor_id -> Int4, list_type -> Nullable, diff --git a/parakeet/src/loaders/list.rs b/parakeet/src/loaders/list.rs index 48627892..6f986ed8 100644 --- a/parakeet/src/loaders/list.rs +++ b/parakeet/src/loaders/list.rs @@ -37,7 +37,7 @@ pub fn build_lists_batch_query() -> &'static str { FROM lists l INNER JOIN actors owner ON l.actor_id = owner.id LEFT JOIN list_items li ON li.list_id = l.id - WHERE 'at://' || owner.did || '/app.bsky.graph.list/' || i64_to_tid(l.rkey) = ANY($1) + WHERE 'at://' || owner.did || '/app.bsky.graph.list/' || l.rkey = ANY($1) GROUP BY l.id, l.actor_id, l.rkey, l.cid, l.owner_actor_id, l.list_type, l.name, l.description, l.description_facets, l.avatar_cid, l.status, owner.did" } @@ -57,8 +57,8 @@ impl BatchFn for ListLoader { id: i64, #[diesel(sql_type = diesel::sql_types::Integer)] actor_id: i32, - #[diesel(sql_type = diesel::sql_types::BigInt)] - rkey: i64, + #[diesel(sql_type = diesel::sql_types::Text)] + rkey: String, #[diesel(sql_type = diesel::sql_types::Binary)] cid: Vec, #[diesel(sql_type = diesel::sql_types::Timestamptz)] @@ -96,15 +96,14 @@ impl BatchFn for ListLoader { }); HashMap::from_iter(lists_with_computed.into_iter().map(|l| { - // Encode TID using Rust utility function - let encoded_rkey = parakeet_db::tid_util::encode_tid(l.rkey); - let at_uri = format!("at://{}/app.bsky.graph.list/{}", l.owner_did, encoded_rkey); + // Rkey is already in string form (either encoded TID or arbitrary string like "nfb") + let at_uri = format!("at://{}/app.bsky.graph.list/{}", l.owner_did, l.rkey); let enriched = EnrichedList { list: models::List { id: l.id, actor_id: l.actor_id, - rkey: l.rkey, + rkey: l.rkey.clone(), cid: l.cid.clone(), owner_actor_id: l.owner_actor_id, list_type: l.list_type,