diff --git a/consumer/src/database_writer/operations/executor.rs b/consumer/src/database_writer/operations/executor.rs index 4d1ccec7..3e285814 100644 --- a/consumer/src/database_writer/operations/executor.rs +++ b/consumer/src/database_writer/operations/executor.rs @@ -190,10 +190,11 @@ pub async fn execute_operation( // Resolve parent URI to natural key if let Ok(Some((parent_actor_id, parent_rkey))) = crate::db::id_resolution::resolve_post_uri(conn, parent_uri).await { // Append to parent's reply arrays (array-only tracking) + // Arrays are correlated (actor_id, rkey) pairs sorted by rkey for better columnstore compression let _ = conn.execute( "UPDATE posts - SET reply_actor_ids = COALESCE(reply_actor_ids, ARRAY[]::integer[]) || $3, - reply_rkeys = COALESCE(reply_rkeys, ARRAY[]::bigint[]) || $4 + SET reply_actor_ids = (SELECT array_agg(actor_id ORDER BY rkey) FROM unnest(COALESCE(reply_actor_ids, ARRAY[]::integer[]) || $3, COALESCE(reply_rkeys, ARRAY[]::bigint[]) || $4) AS t(actor_id, rkey)), + reply_rkeys = (SELECT array_agg(rkey ORDER BY rkey) FROM unnest(COALESCE(reply_actor_ids, ARRAY[]::integer[]) || $3, COALESCE(reply_rkeys, ARRAY[]::bigint[]) || $4) AS t(actor_id, rkey)) WHERE (actor_id, rkey) = ($1, $2) AND NOT ($3 = ANY(COALESCE(reply_actor_ids, ARRAY[])))", &[&parent_actor_id, &parent_rkey, &actor_id, &rkey], @@ -208,10 +209,11 @@ pub async fn execute_operation( // Resolve quoted URI to natural key if let Ok(Some((quoted_actor_id, quoted_rkey))) = crate::db::id_resolution::resolve_post_uri(conn, quoted_uri).await { // Append to quoted post's quote arrays (array-only tracking) + // Arrays are correlated (actor_id, rkey) pairs sorted by rkey for better columnstore compression let _ = conn.execute( "UPDATE posts - SET quote_actor_ids = COALESCE(quote_actor_ids, ARRAY[]::integer[]) || $3, - quote_rkeys = COALESCE(quote_rkeys, ARRAY[]::bigint[]) || $4 + SET quote_actor_ids = (SELECT array_agg(actor_id ORDER BY rkey) FROM unnest(COALESCE(quote_actor_ids, ARRAY[]::integer[]) || $3, COALESCE(quote_rkeys, ARRAY[]::bigint[]) || $4) AS t(actor_id, rkey)), + quote_rkeys = (SELECT array_agg(rkey ORDER BY rkey) FROM unnest(COALESCE(quote_actor_ids, ARRAY[]::integer[]) || $3, COALESCE(quote_rkeys, ARRAY[]::bigint[]) || $4) AS t(actor_id, rkey)) WHERE (actor_id, rkey) = ($1, $2) AND NOT ($3 = ANY(COALESCE(quote_actor_ids, ARRAY[])))", &["ed_actor_id, "ed_rkey, &actor_id, &rkey], diff --git a/consumer/src/db/operations/feed/post.rs b/consumer/src/db/operations/feed/post.rs index a814f4b0..371db391 100644 --- a/consumer/src/db/operations/feed/post.rs +++ b/consumer/src/db/operations/feed/post.rs @@ -266,6 +266,14 @@ pub async fn post_insert( // Mentions, facets, and embeds are now stored directly in the posts table via composite types // No need for separate insert operations + // Append post rkey to actor's post_rkeys array (sorted for better columnstore compression) + let _ = conn.execute( + "UPDATE actors + SET post_rkeys = (SELECT array_agg(elem ORDER BY elem) FROM unnest(COALESCE(post_rkeys, ARRAY[]::bigint[]) || $2::bigint) elem) + WHERE id = $1", + &[&actor_id, &rkey], + ).await; + Ok(1) } else { Ok(0) diff --git a/consumer/src/db/operations/feed/repost.rs b/consumer/src/db/operations/feed/repost.rs index 273c76a9..4b3d02e3 100644 --- a/consumer/src/db/operations/feed/repost.rs +++ b/consumer/src/db/operations/feed/repost.rs @@ -86,6 +86,16 @@ pub async fn repost_insert( // Note: We no longer enqueue posts for fetching when reposts are created. // With natural keys, we don't create post stubs, so there's nothing to enqueue. + // Append repost rkey to actor's repost_rkeys array (sorted for better columnstore compression) + if rows > 0 { + let _ = conn.execute( + "UPDATE actors + SET repost_rkeys = (SELECT array_agg(elem ORDER BY elem) FROM unnest(COALESCE(repost_rkeys, ARRAY[]::bigint[]) || $2::bigint) elem) + WHERE id = $1", + &[&actor_id, &rkey], + ).await; + } + Ok(rows) }