diff --git a/consumer/src/database_writer/locking.rs b/consumer/src/database_writer/locking.rs index a7c0c84b..a9dabd2b 100644 --- a/consumer/src/database_writer/locking.rs +++ b/consumer/src/database_writer/locking.rs @@ -11,11 +11,6 @@ use std::hash::{Hash, Hasher}; /// Uses DefaultHasher for fast, stable hashing within a process. /// Hash collisions are acceptable - they just cause false contention /// (performance impact only, not a correctness issue). -/// -/// # Example -/// ``` -/// let lock_id = did_to_lock_id("did:plc:abcdef123456"); -/// ``` pub fn did_to_lock_id(did: &str) -> i64 { let mut hasher = DefaultHasher::new(); did.hash(&mut hasher); @@ -30,11 +25,6 @@ pub fn did_to_lock_id(did: &str) -> i64 { /// Uses DefaultHasher for fast, stable hashing within a process. /// Hash collisions are acceptable - they just cause false contention /// (performance impact only, not a correctness issue). -/// -/// # Example -/// ``` -/// let lock_id = uri_to_lock_id("at://did:plc:abc/app.bsky.feed.post/123"); -/// ``` pub fn uri_to_lock_id(uri: &str) -> i64 { let mut hasher = DefaultHasher::new(); uri.hash(&mut hasher); diff --git a/consumer/src/db/operations/feed.rs b/consumer/src/db/operations/feed.rs index f6adab95..fe3947d5 100644 --- a/consumer/src/db/operations/feed.rs +++ b/consumer/src/db/operations/feed.rs @@ -9,7 +9,6 @@ use chrono::prelude::*; use deadpool_postgres::GenericClient; use eyre::{Context as _, OptionExt as _}; use ipld_core::cid::Cid; -use std::collections::HashSet; // Post functions @@ -469,37 +468,61 @@ async fn post_embed_record_insert( post_id: i64, embed: AppBskyEmbedRecord, post_created_at: DateTime, - is_backfill: bool, + _is_backfill: bool, ) -> Result { // Parse embedded record AT URI: at://did/collection/rkey let parts = embed.record.uri[5..].split('/').collect::>(); let collection = parts.get(1).copied().unwrap_or(""); // Check if this embed should be marked as detached (postgate logic) - // We need to look up the postgate by post_id + // We need to look up the postgate on the EMBEDDED post, not the quoting post let detached = if collection == "app.bsky.feed.post" { - let pg_data: Option<(DateTime, Vec, Vec)> = conn + // Query for postgate on the embedded post + // Check: + // 1. If the quoting post is in the specific detached list, OR + // 2. If there's a global disable rule and the quote was created after it became effective + let result: Option = conn .query_opt( - "SELECT pg.created_at, ARRAY[]::text[] as detached, pg.rules - FROM postgates pg - WHERE pg.post_id = $1", - &[&post_id], + "WITH embedded_parts AS ( + SELECT + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING($1 FROM 6), '/', 3) as rkey + ), + embedded_lookup AS ( + SELECT p.id + FROM embedded_parts ep + INNER JOIN actors a ON a.did = ep.did + INNER JOIN posts p ON p.actor_id = a.id AND p.rkey = tid_to_i64(ep.rkey) + ), + postgate_data AS ( + SELECT pg.created_at, pg.rules + FROM postgates pg + WHERE pg.post_id = (SELECT id FROM embedded_lookup) + ), + is_in_detached_list AS ( + SELECT EXISTS( + SELECT 1 + FROM postgates pg + INNER JOIN postgate_detached pd ON pd.postgate_id = pg.id + WHERE pg.post_id = (SELECT id FROM embedded_lookup) + AND pd.detached_post_id = $2 + ) as in_list + ) + SELECT + CASE + -- In specific detached list + WHEN (SELECT in_list FROM is_in_detached_list) THEN true + -- Global disable rule applies (created after postgate effective date) + WHEN (SELECT rules FROM postgate_data) @> ARRAY[$3::text]::postgate_rule[] + AND $4 > (SELECT created_at FROM postgate_data) THEN true + ELSE false + END", + &[&embed.record.uri, &post_id, &PG_DISABLE_RULE, &post_created_at], ) .await? - .map(|row| (row.get(0), row.get(1), row.get(2))); + .map(|row| row.get(0)); - if let Some((effective, _detached, rules)) = pg_data { - let rules: HashSet = HashSet::from_iter(rules); - let compare_date = if is_backfill { - post_created_at - } else { - Utc::now() - }; - - rules.contains(PG_DISABLE_RULE) && compare_date > effective - } else { - false - } + result.unwrap_or(false) } else { false }; diff --git a/consumer/src/db/workers.rs b/consumer/src/db/workers.rs index 91fa1156..274c3db8 100644 --- a/consumer/src/db/workers.rs +++ b/consumer/src/db/workers.rs @@ -14,18 +14,7 @@ use parakeet_db::types::ActorSyncState; /// This function uses PostgreSQL's unnest() to efficiently insert multiple actors /// in a single query. Existing actors are ignored (ON CONFLICT DO NOTHING). /// -/// # Arguments -/// * `conn` - Database connection (can be a transaction) -/// * `dids` - Slice of DID strings to ensure -/// -/// # Returns -/// Number of new actor records created (existing actors are not counted) -/// -/// # Example -/// ```no_run -/// let dids = vec!["did:plc:alice", "did:plc:bob"]; -/// let new_actors = bulk_ensure_actors(&txn, &dids).await?; -/// ``` +/// Returns the number of new actor records created (existing actors are not counted). pub async fn bulk_ensure_actors(conn: &C, dids: &[&str]) -> Result { if dids.is_empty() { return Ok(0); diff --git a/consumer/src/error.rs b/consumer/src/error.rs index d831c517..e0ccf045 100644 --- a/consumer/src/error.rs +++ b/consumer/src/error.rs @@ -10,26 +10,6 @@ //! - **Rich context**: Wrap errors at every layer with `.wrap_err()` or `.context()` //! - **Error boundaries**: Errors bubble to specific points (workers, handlers) where decisions are made //! - **Actionable errors**: Include enough context to diagnose without additional logs -//! -//! # Example -//! -//! ```rust -//! use crate::error::{Result, DbContextExt}; -//! -//! async fn insert_post(conn: &Client, uri: &str) -> Result { -//! let rkey = extract_rkey(uri) -//! .ok_or_else(|| eyre::eyre!("Invalid AT URI: missing rkey in {}", uri))?; -//! -//! conn.query_one( -//! "INSERT INTO posts (rkey) VALUES ($1) RETURNING id", -//! &[&rkey], -//! ) -//! .await -//! .wrap_err_with(|| format!("Failed to insert post with rkey {}", rkey))?; -//! -//! Ok(row.get(0)) -//! } -//! ``` /// Standard Result type used throughout the consumer /// diff --git a/consumer/src/utils.rs b/consumer/src/utils.rs index 8170759a..ceebfa55 100644 --- a/consumer/src/utils.rs +++ b/consumer/src/utils.rs @@ -55,12 +55,6 @@ pub fn at_uri_is_by(uri: &str, did: &str) -> bool { /// /// AT URIs have the format: `at://{did}/{collection}/{rkey}` /// This function extracts the DID component. -/// -/// # Example -/// ``` -/// let uri = "at://did:plc:abc123/app.bsky.feed.post/3jui2f5j4b52a"; -/// assert_eq!(extract_did_from_uri(uri), Some("did:plc:abc123")); -/// ``` pub fn extract_did_from_uri(uri: &str) -> Option<&str> { // AT URI format: at://{did}/{collection}/{rkey} // Split from the right: [rkey, collection, did, "at:/"] diff --git a/migrations/2025-11-04-060024_add_maintain_postgates_function/down.sql b/migrations/2025-11-04-060024_add_maintain_postgates_function/down.sql new file mode 100644 index 00000000..d9e94e24 --- /dev/null +++ b/migrations/2025-11-04-060024_add_maintain_postgates_function/down.sql @@ -0,0 +1,2 @@ +-- Remove maintain_postgates function +DROP FUNCTION IF EXISTS maintain_postgates(TEXT, TEXT[], TIMESTAMP); diff --git a/migrations/2025-11-04-060024_add_maintain_postgates_function/up.sql b/migrations/2025-11-04-060024_add_maintain_postgates_function/up.sql new file mode 100644 index 00000000..479f1818 --- /dev/null +++ b/migrations/2025-11-04-060024_add_maintain_postgates_function/up.sql @@ -0,0 +1,70 @@ +-- Create maintain_postgates function +-- Updates the 'detached' flag on existing post_embed_record entries when a postgate changes +-- +-- Parameters: +-- $1: post_uri - AT URI of the post that has the postgate (e.g., "at://did:plc:abc/app.bsky.feed.post/xyz") +-- $2: detached_uris - Array of post URIs that should be marked as detached +-- $3: disable_effective - Timestamp when the global disable rule became effective (nullable) +-- +-- Returns: Number of rows updated (as BIGINT for compatibility with conn.execute) +CREATE OR REPLACE FUNCTION maintain_postgates( + post_uri TEXT, + detached_uris TEXT[], + disable_effective TIMESTAMP +) RETURNS BIGINT AS $$ +DECLARE + target_post_id BIGINT; + rows_updated BIGINT; +BEGIN + -- Look up the post_id from the post URI + SELECT p.id INTO target_post_id + FROM posts p + INNER JOIN actors a ON p.actor_id = a.id + WHERE 'at://' || a.did || '/app.bsky.feed.post/' || i64_to_tid(p.rkey) = post_uri; + + -- If post doesn't exist, nothing to do + IF target_post_id IS NULL THEN + RETURN 0; + END IF; + + -- Update post_embed_record entries to mark quotes as detached + -- A quote post is detached if: + -- 1. It's in the detached_uris list, OR + -- 2. The global disable rule exists (disable_effective is not null) AND the quoting post was created after disable_effective + WITH detached_post_ids AS ( + -- Resolve the detached URIs to post IDs + SELECT p.id + FROM UNNEST(detached_uris) AS uri + CROSS JOIN LATERAL ( + SELECT + SPLIT_PART(SUBSTRING(uri FROM 6), '/', 1) as did, + SPLIT_PART(SUBSTRING(uri FROM 6), '/', 3) as rkey + ) parts + INNER JOIN actors a ON a.did = parts.did + INNER JOIN posts p ON p.actor_id = a.id AND p.rkey = tid_to_i64(parts.rkey) + ), + posts_to_update AS ( + -- Find all posts that quote the target post + SELECT per.post_id + FROM post_embed_record per + WHERE per.embedded_post_id = target_post_id + ) + UPDATE post_embed_record per + SET detached = ( + -- Mark as detached if either: + -- 1. The quoting post is in the specific detached list + per.post_id IN (SELECT id FROM detached_post_ids) + OR + -- 2. Global disable rule applies (quote created after disable_effective) + (disable_effective IS NOT NULL AND ( + SELECT p.created_at + FROM posts p + WHERE p.id = per.post_id + ) > disable_effective) + ) + WHERE per.embedded_post_id = target_post_id; + + GET DIAGNOSTICS rows_updated = ROW_COUNT; + RETURN rows_updated; +END; +$$ LANGUAGE plpgsql;