diff --git a/consumer/src/database_writer/locking.rs b/consumer/src/database_writer/locking.rs index 21efa27f..7b146528 100644 --- a/consumer/src/database_writer/locking.rs +++ b/consumer/src/database_writer/locking.rs @@ -96,6 +96,33 @@ pub fn actor_record_lock_str(table_name: &str, actor_id: i32, rkey: &str) -> (i3 (table_id, key_id) } +/// Acquire a single advisory lock +/// +/// This is a simple wrapper around the common pattern of: +/// ```ignore +/// conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]).await?; +/// ``` +/// +/// The lock is automatically released when the transaction commits/rolls back. +/// +/// # Example +/// ```ignore +/// let (table_id, key_id) = table_record_lock("actors", did); +/// acquire_lock(conn, table_id, key_id).await?; +/// ``` +pub async fn acquire_lock( + conn: &C, + table_id: i32, + key_id: i32, +) -> eyre::Result<()> { + conn.execute( + "SELECT pg_advisory_xact_lock($1, $2)", + &[&table_id, &key_id], + ) + .await?; + Ok(()) +} + /// Acquire transaction-level advisory locks for a set of DIDs /// /// Locks are acquired in sorted DID order to prevent deadlocks. diff --git a/consumer/src/db/actor.rs b/consumer/src/db/actor.rs index 3aacbf75..1d9ae061 100644 --- a/consumer/src/db/actor.rs +++ b/consumer/src/db/actor.rs @@ -16,8 +16,7 @@ pub async fn actor_upsert( ) -> Result { // Acquire advisory lock on DID to prevent races let (table_id, key_id) = crate::database_writer::locking::table_record_lock("actors", did); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Allow allowlist states (synced, dirty, processing) to flow freely // Allow upgrading from partial to allowlist states @@ -327,8 +326,7 @@ pub async fn ensure_actor_id( // This ensures only one transaction at a time can create/access this specific actor // Prevents sequence consumption from concurrent inserts or rollback scenarios let (table_id, key_id) = crate::database_writer::locking::table_record_lock("actors", did); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Use CTE to SELECT first, then conditionally INSERT only if not found // This prevents unnecessary sequence consumption when actor already exists @@ -449,8 +447,7 @@ pub async fn ensure_actor_id_with_cache( // Slow path: Not in cache, do database operation // Acquire per-DID advisory lock to prevent race conditions let (table_id, key_id) = crate::database_writer::locking::table_record_lock("actors", did); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Modified CTE that also returns sync_state for allowlist check and was_created flag let row = match (status, handle) { diff --git a/consumer/src/db/operations/feed/helpers.rs b/consumer/src/db/operations/feed/helpers.rs index 6cee11cf..fea48523 100644 --- a/consumer/src/db/operations/feed/helpers.rs +++ b/consumer/src/db/operations/feed/helpers.rs @@ -93,8 +93,7 @@ pub(super) async fn get_feedgen_id( ) -> Result<(i64, bool)> { // Acquire advisory lock on feedgen URI to prevent concurrent access races let (table_id, key_id) = crate::database_writer::locking::table_record_lock("feedgens", at_uri); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Parse the CID string to get the digest let cid = ipld_core::cid::Cid::try_from(cid_str) @@ -163,8 +162,7 @@ pub async fn ensure_list_id(conn: &C, at_uri: &str) -> Result< pub(super) async fn get_list_id(conn: &C, at_uri: &str) -> Result<(i64, bool)> { // Acquire advisory lock on list URI to prevent concurrent access races let (table_id, key_id) = crate::database_writer::locking::table_record_lock("lists", at_uri); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Extract owner DID and rkey from AT URI let did = parakeet_db::at_uri_util::extract_did(at_uri) diff --git a/consumer/src/db/operations/feed/like.rs b/consumer/src/db/operations/feed/like.rs index 9398072f..e60c5bba 100644 --- a/consumer/src/db/operations/feed/like.rs +++ b/consumer/src/db/operations/feed/like.rs @@ -51,8 +51,7 @@ pub async fn like_insert( // Acquire advisory lock on record to prevent deadlocks let lock_start = std::time::Instant::now(); let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("likes", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; let lock_ms = lock_start.elapsed().as_millis(); // Note: via_repost_id is already resolved in the reference extraction phase (workers.rs) diff --git a/consumer/src/db/operations/feed/post.rs b/consumer/src/db/operations/feed/post.rs index fc7d3f1f..0d35d5c1 100644 --- a/consumer/src/db/operations/feed/post.rs +++ b/consumer/src/db/operations/feed/post.rs @@ -76,8 +76,7 @@ pub async fn post_insert( // Acquire advisory lock using actor_id and rkey to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("posts", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Get parent and root post natural keys (create stubs if necessary) // First get the actor IDs, then get the post natural keys diff --git a/consumer/src/db/operations/feed/repost.rs b/consumer/src/db/operations/feed/repost.rs index 9baa51f2..273c76a9 100644 --- a/consumer/src/db/operations/feed/repost.rs +++ b/consumer/src/db/operations/feed/repost.rs @@ -44,8 +44,7 @@ pub async fn repost_insert( // Acquire advisory lock on record to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("reposts", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Note: via_repost_id is already resolved in the reference extraction phase (workers.rs) // This ensures the FK constraint is satisfied before we reach this point diff --git a/consumer/src/db/operations/graph.rs b/consumer/src/db/operations/graph.rs index 394a3371..6115764b 100644 --- a/consumer/src/db/operations/graph.rs +++ b/consumer/src/db/operations/graph.rs @@ -20,8 +20,7 @@ pub async fn follow_insert( ) -> Result { // Acquire advisory lock on record to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("follows", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Insert follow - simple insert on primary key (actor_id, rkey) // Note: CID is synthetic, generated from actor_id + rkey @@ -69,8 +68,7 @@ pub async fn block_insert( ) -> Result { // Acquire advisory lock on record to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("blocks", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Insert block - simple insert on primary key (actor_id, rkey) // Note: CID is synthetic, generated from actor_id + rkey @@ -114,8 +112,7 @@ pub async fn list_upsert( // Acquire table-scoped advisory lock on lists table for this record 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?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; conn.query_one( include_str!("../sql/list_upsert.sql"), @@ -162,8 +159,7 @@ pub async fn list_block_insert( // Acquire advisory lock on record to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("list_blocks", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; conn.execute( include_str!("../sql/list_block_upsert.sql"), @@ -210,8 +206,7 @@ pub async fn list_item_insert( // Acquire advisory lock on record to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("list_items", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Resolve list_id by creating stub list if needed let list_id = super::feed::ensure_list_id(conn, &rec.list).await?; @@ -262,8 +257,7 @@ pub async fn verification_insert( // Acquire advisory lock on record to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("verification", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Insert verification - actors already resolved (no SELECT subquery needed) // Note: created_at is derived from TID rkey diff --git a/consumer/src/db/operations/labeler.rs b/consumer/src/db/operations/labeler.rs index b3f96283..d82bb1ad 100644 --- a/consumer/src/db/operations/labeler.rs +++ b/consumer/src/db/operations/labeler.rs @@ -24,8 +24,7 @@ pub async fn ensure_labeler_stub( // Labeler URI is always at://{did}/app.bsky.labeler.service/self let labeler_uri = format!("at://{}/app.bsky.labeler.service/self", did); let (table_id, key_id) = crate::database_writer::locking::table_record_lock("labelers", &labeler_uri); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Parse the CID string to get the digest let cid = ipld_core::cid::Cid::try_from(cid_str) @@ -75,8 +74,7 @@ pub async fn labeler_upsert( // Acquire table-scoped advisory lock on labelers table for this actor // This prevents concurrent labeler updates for the same actor let (table_id, key_id) = crate::database_writer::locking::table_record_lock("labelers", did); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; let cid_bytes = cid.to_bytes(); let cid_digest = parakeet_db::cid_util::cid_to_digest(&cid_bytes) diff --git a/consumer/src/db/operations/starter_pack.rs b/consumer/src/db/operations/starter_pack.rs index 0630f9c1..727f9697 100644 --- a/consumer/src/db/operations/starter_pack.rs +++ b/consumer/src/db/operations/starter_pack.rs @@ -24,8 +24,7 @@ pub async fn starter_pack_upsert( // Acquire advisory lock on record to prevent deadlocks let (table_id, key_id) = crate::database_writer::locking::actor_record_lock("starterpacks", actor_id, rkey); - conn.execute("SELECT pg_advisory_xact_lock($1, $2)", &[&table_id, &key_id]) - .await?; + crate::database_writer::locking::acquire_lock(conn, table_id, key_id).await?; // Resolve list_id by creating stub list if needed let list_id = super::feed::ensure_list_id(conn, &rec.list).await?;