From e60a80c8ce61dcfa01cae1b87a1cea6e14398bd2 Mon Sep 17 00:00:00 2001 From: Timothy Quilling Date: Fri, 14 Nov 2025 23:24:06 -0500 Subject: [PATCH] chore: normalize notifications --- .../database_writer/operations/executor.rs | 142 ++++++++++-------- consumer/src/db/mod.rs | 39 ++--- 2 files changed, 95 insertions(+), 86 deletions(-) diff --git a/consumer/src/database_writer/operations/executor.rs b/consumer/src/database_writer/operations/executor.rs index 23dcf2c1..275eefb4 100644 --- a/consumer/src/database_writer/operations/executor.rs +++ b/consumer/src/database_writer/operations/executor.rs @@ -360,8 +360,8 @@ pub async fn execute_operation( return Ok((vec![], vec![])); } - // Query recipient DID for allowlist and thread mute checks - // TODO: Refactor allowlist and is_thread_muted to use actor_ids directly + // Query recipient DID for allowlist check + // TODO: Refactor allowlist to use actor_ids directly let recipient_did: Option = conn .query_opt("SELECT did FROM actors WHERE id = $1", &[&recipient_actor_id]) .await? @@ -378,32 +378,18 @@ pub async fn execute_operation( } // Check if thread is muted (for reply, quote, and mention notifications) - // For normalized notifications, we need to reconstruct the subject URI for the mute check - if let (Some(subj_actor_id), Some(subj_type), Some(subj_rkey)) = - (subject_actor_id, subject_record_type.as_ref(), subject_rkey.as_ref()) + // Only posts can be thread roots, so only check for post subjects + if let (Some(subj_actor_id), Some("post"), Some(subj_rkey)) = + (subject_actor_id, subject_record_type.as_deref(), subject_rkey.as_ref()) { - // Get subject actor DID - if let Ok(Some(subject_did_row)) = conn - .query_opt("SELECT did FROM actors WHERE id = $1", &[&subj_actor_id]) - .await - { - let subject_did: String = subject_did_row.get(0); - let collection = match subj_type.as_str() { - "post" => "app.bsky.feed.post", - "like" => "app.bsky.feed.like", - "repost" => "app.bsky.feed.repost", - _ => "", - }; - let thread_root = format!("at://{}/{}/{}", subject_did, collection, subj_rkey); - - if db::is_thread_muted(conn, &recipient_did, &thread_root).await? { - tracing::debug!( - "Skipping notification for muted thread: recipient_actor_id={}, thread={}", - recipient_actor_id, - thread_root - ); - return Ok((vec![], vec![])); - } + if db::is_thread_muted(conn, recipient_actor_id, subj_actor_id, *subj_rkey).await? { + tracing::debug!( + "Skipping notification for muted thread: recipient_actor_id={}, root_post_actor_id={}, root_post_rkey={}", + recipient_actor_id, + subj_actor_id, + subj_rkey + ); + return Ok((vec![], vec![])); } } @@ -493,8 +479,60 @@ pub async fn execute_operation( ancestor_actor_id ); } else { - // Query ancestor DID for allowlist and thread mute checks - // TODO: Refactor allowlist and is_thread_muted to use actor_ids directly + // Parse reason_subject URI early for thread mute check and notification + // Format: at://did:plc:xyz/app.bsky.feed.post/rkey + let subject_rkey_str = reason_subject.split('/').next_back().unwrap_or(""); + let subject_rkey_i64 = match parakeet_db::tid_util::decode_tid(subject_rkey_str) { + Ok(rkey) => rkey, + Err(e) => { + tracing::warn!("Invalid subject TID in reason_subject: {} - {}", subject_rkey_str, e); + level += 1; + if let Some(parent) = ancestor_parent_uri { + current_uri = parent; + continue; + } else { + break; + } + } + }; + + // Get subject_actor_id from reason_subject DID + let subject_actor_id = if let Some(subject_did) = reason_subject + .strip_prefix("at://") + .and_then(|s| s.split('/').next()) + { + match conn.query_opt("SELECT id FROM actors WHERE did = $1", &[&subject_did]) + .await? + .and_then(|row| row.get::<_, Option>(0)) + { + Some(id) => id, + None => { + tracing::debug!( + "Skipping reply-chain notification at level {}: subject DID not found: {}", + level, + subject_did + ); + level += 1; + if let Some(parent) = ancestor_parent_uri { + current_uri = parent; + continue; + } else { + break; + } + } + } + } else { + tracing::warn!("Invalid reason_subject URI format: {}", reason_subject); + level += 1; + if let Some(parent) = ancestor_parent_uri { + current_uri = parent; + continue; + } else { + break; + } + }; + + // Query ancestor DID for allowlist check (TODO: refactor allowlist to use actor_ids) let ancestor_did: Option = conn .query_opt("SELECT did FROM actors WHERE id = $1", &[&ancestor_actor_id]) .await? @@ -522,13 +560,14 @@ pub async fn execute_operation( level, ancestor_actor_id ); - } else if db::is_thread_muted(conn, &ancestor_did, &reason_subject).await? { - // Skip if ancestor has muted this thread (TODO: refactor is_thread_muted to use actor_ids) + } else if db::is_thread_muted(conn, ancestor_actor_id, subject_actor_id, subject_rkey_i64).await? { + // Check if ancestor has muted this thread using actor_ids and i64 rkey tracing::debug!( - "Skipping reply-chain notification at level {}: actor_id={} muted thread {}", + "Skipping reply-chain notification at level {}: actor_id={} muted thread (root_post_actor_id={}, rkey={})", level, ancestor_actor_id, - reason_subject + subject_actor_id, + subject_rkey_i64 ); } else { // Create PostgreSQL notification for this ancestor @@ -548,16 +587,6 @@ pub async fn execute_operation( } }; - // Extract subject components from reason_subject (root or parent URI) - let subject_rkey_str = reason_subject.split('/').next_back().unwrap_or(""); - let subject_rkey_i64 = match parakeet_db::tid_util::decode_tid(subject_rkey_str) { - Ok(rkey) => Some(rkey), - Err(e) => { - tracing::warn!("Invalid subject TID: {} - {}", subject_rkey_str, e); - None - } - }; - // Parse CID string and extract digest let cid_digest = match Cid::try_from(cid.as_str()) { Ok(cid_obj) => { @@ -588,30 +617,17 @@ pub async fn execute_operation( } }; - // Get subject_actor_id from reason_subject URI - // We need to query the database to get the actor_id for the subject DID - let subject_actor_id = if let Some(subject_did) = reason_subject - .strip_prefix("at://") - .and_then(|s| s.split('/').next()) - { - conn.query_opt("SELECT id FROM actors WHERE did = $1", &[&subject_did]) - .await? - .and_then(|row| row.get::<_, Option>(0)) - } else { - None - }; - if let Err(e) = crate::external::pg_notifications::add_notification( pool, ancestor_actor_id, author_actor_id, - "post", // record_type (the reply is a post) - reply_rkey_i64, // record_rkey (i64) - &cid_digest, // record_cid (32-byte digest) - "reply", // reason - subject_actor_id, // subject_actor_id (root/parent post author) - Some("post"), // subject_record_type - subject_rkey_i64, // subject_rkey (Option) + "post", // record_type (the reply is a post) + reply_rkey_i64, // record_rkey (i64) + &cid_digest, // record_cid (32-byte digest) + "reply", // reason + Some(subject_actor_id), // subject_actor_id (root/parent post author) + Some("post"), // subject_record_type + Some(subject_rkey_i64), // subject_rkey (Option) created_at, ) .await diff --git a/consumer/src/db/mod.rs b/consumer/src/db/mod.rs index 66d0880d..6d050a10 100644 --- a/consumer/src/db/mod.rs +++ b/consumer/src/db/mod.rs @@ -96,36 +96,29 @@ pub async fn mark_stub_as_missing( Ok(()) } -/// Check if a thread is muted by a user +/// Check if a recipient has muted a thread +/// +/// # Arguments +/// +/// * `recipient_actor_id` - The actor_id of the potential recipient +/// * `root_post_actor_id` - The actor_id of the root post author +/// * `root_post_rkey` - The rkey (as i64) of the root post +/// +/// # Returns +/// +/// `true` if the recipient has muted this thread, `false` otherwise pub async fn is_thread_muted( conn: &C, - recipient_did: &str, - thread_root_uri: &str, + recipient_actor_id: i32, + root_post_actor_id: i32, + root_post_rkey: i64, ) -> Result { - // Parse the thread root URI to get DID and rkey - // Format: at://did:plc:xyz/app.bsky.feed.post/rkey - let uri_parts: Vec<&str> = thread_root_uri - .strip_prefix("at://") - .unwrap_or(thread_root_uri) - .split('/') - .collect(); - - if uri_parts.len() < 3 { - // Invalid URI format - not muted - return Ok(false); - } - - let root_did = uri_parts[0]; - let root_rkey = uri_parts[2]; - let row = conn .query_opt( "SELECT 1 FROM thread_mutes tm - INNER JOIN actors a ON tm.actor_id = a.id INNER JOIN posts p ON tm.root_post_id = p.id - INNER JOIN actors pa ON p.actor_id = pa.id - WHERE a.did = $1 AND pa.did = $2 AND p.rkey = tid_to_i64($3)", - &[&recipient_did, &root_did, &root_rkey], + WHERE tm.actor_id = $1 AND p.actor_id = $2 AND p.rkey = $3", + &[&recipient_actor_id, &root_post_actor_id, &root_post_rkey], ) .await?; Ok(row.is_some()) -- 2.51.2