diff --git a/consumer/src/database_writer/mod.rs b/consumer/src/database_writer/mod.rs index 884323e3..3361a0df 100644 --- a/consumer/src/database_writer/mod.rs +++ b/consumer/src/database_writer/mod.rs @@ -23,7 +23,7 @@ pub mod workers; // acquire_did_locks is used internally by workers, not exposed pub use bulk_types::{BulkOperations, UnresolvedRecord}; -pub use operations::{process_record_to_operations, ActorStatsDeltas, DatabaseOperation, PostStatsDeltas}; +pub use operations::{process_record_to_operations, DatabaseOperation, PostStatsDeltas}; pub use reference_extraction::{extract_references, RecordReferences}; pub use timestamp::{validate_record_timestamp_with_tid, validate_tid_timestamp}; pub use workers::spawn_database_writer; diff --git a/consumer/src/database_writer/operations/executor.rs b/consumer/src/database_writer/operations/executor.rs index 70e79536..4d1ccec7 100644 --- a/consumer/src/database_writer/operations/executor.rs +++ b/consumer/src/database_writer/operations/executor.rs @@ -7,46 +7,16 @@ //! - `execute_operation()` - Single operation execution with aggregate delta generation //! - `describe_operation()` - Logging helper for operation names and URIs //! -//! ## Aggregate Deltas +//! ## Aggregate Deltas (Phase 8: DEPRECATED) //! -//! Operations return aggregate deltas ONLY when they actually affect rows in the database. -//! This makes stats idempotent - duplicate events won't double-count. +//! Actor stats deltas have been removed. Counts are now maintained automatically by database triggers: +//! - followers_count, following_count (from followers/following arrays) +//! - posts_count (from post_rkeys array) +//! Operations still return PostStatsDeltas for compatibility, but they are empty (arrays updated directly). -use super::{cache, ActorStatsDeltas, DatabaseOperation}; -use eyre::WrapErr as _; +use super::{cache, DatabaseOperation}; use ipld_core::cid::Cid; -// Actor stats delta helpers for denormalized stats in actors table - -/// Helper to create follow stats deltas (affects both follower and followee) -/// Returns deltas for both: the actor who followed (+1 following) and the target actor (+1 followers) -fn follow_deltas(actor_id: i32, subject_actor_id: i32, delta: i32) -> Vec<(i32, ActorStatsDeltas)> { - vec![ - (actor_id, ActorStatsDeltas { following_delta: delta, ..Default::default() }), - (subject_actor_id, ActorStatsDeltas { followers_delta: delta, ..Default::default() }), - ] -} - -/// Helper to create post delta for actor -fn actor_post_delta(actor_id: i32, delta: i32) -> Vec<(i32, ActorStatsDeltas)> { - vec![(actor_id, ActorStatsDeltas { posts_delta: delta, ..Default::default() })] -} - -/// Helper to create list delta for actor -fn actor_list_delta(actor_id: i32, delta: i16) -> Vec<(i32, ActorStatsDeltas)> { - vec![(actor_id, ActorStatsDeltas { lists_delta: delta, ..Default::default() })] -} - -/// Helper to create feedgen delta for actor -fn actor_feedgen_delta(actor_id: i32, delta: i16) -> Vec<(i32, ActorStatsDeltas)> { - vec![(actor_id, ActorStatsDeltas { feeds_delta: delta, ..Default::default() })] -} - -/// Helper to create starterpack delta for actor -fn actor_starterpack_delta(actor_id: i32, delta: i16) -> Vec<(i32, ActorStatsDeltas)> { - vec![(actor_id, ActorStatsDeltas { starterpacks_delta: delta, ..Default::default() })] -} - pub fn describe_operation(op: &DatabaseOperation) -> (String, Option) { match op { DatabaseOperation::InsertPost { rkey, actor_id, .. } => { @@ -153,10 +123,10 @@ pub fn describe_operation(op: &DatabaseOperation) -> (String, Option) { /// /// This function will be called by the database writer for each operation. /// It returns: -/// - Vec of PostStatsDeltas for affected posts (actor_id, rkey) -> deltas -/// - Vec of ActorStatsDeltas for affected actors (actor_id) -> deltas +/// - Vec of PostStatsDeltas for affected posts (DEPRECATED - always empty, kept for compatibility) /// - Cache invalidations that should be sent via NOTIFY (for real-time cache updates) /// +/// Phase 8: ActorStatsDeltas removed - counts maintained by database triggers automatically. /// Redis has been eliminated! Notifications, fetch queue, and PDS cache now use PostgreSQL/moka. pub async fn execute_operation( pool: &deadpool_postgres::Pool, @@ -164,7 +134,7 @@ pub async fn execute_operation( op: DatabaseOperation, allowlist: &crate::db::Allowlist, pds_cache: &crate::external::pds_cache::PdsHostCache, -) -> eyre::Result<(Vec<((i32, i64), super::PostStatsDeltas)>, Vec<(i32, super::ActorStatsDeltas)>, Vec)> { +) -> eyre::Result<(Vec<((i32, i64), super::PostStatsDeltas)>, Vec)> { use crate::db; match op { @@ -212,8 +182,7 @@ pub async fn execute_operation( ); } if rows > 0 { - let mut deltas = Vec::new(); - let actor_deltas = actor_post_delta(actor_id, 1); + let deltas = Vec::new(); // Phase 8: Empty - arrays updated directly by SQL let mut cache_invalidations = Vec::new(); // If this is a reply, increment parent post's reply count @@ -252,9 +221,9 @@ pub async fn execute_operation( } } - Ok((deltas, actor_deltas, cache_invalidations)) + Ok((deltas, cache_invalidations)) } else { - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } DatabaseOperation::InsertLike { @@ -288,9 +257,9 @@ pub async fn execute_operation( // Invalidate post cache since like array changed let cache_invalidations = vec![cache::invalidate_post(&subject_uri)]; - Ok((vec![], vec![], cache_invalidations)) + Ok((vec![], cache_invalidations)) } else { - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } DatabaseOperation::InsertRepost { @@ -325,9 +294,9 @@ pub async fn execute_operation( // Invalidate post cache since repost array changed let cache_invalidations = vec![cache::invalidate_post(&subject_uri)]; - Ok((vec![], vec![], cache_invalidations)) + Ok((vec![], cache_invalidations)) } else { - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } DatabaseOperation::InsertFollow { @@ -339,10 +308,10 @@ pub async fn execute_operation( } => { let rows = db::follow_insert(conn, rkey, actor_id, subject_actor_id, cid, record).await?; if rows > 0 { - let actor_deltas = follow_deltas(actor_id, subject_actor_id, 1); - Ok((vec![], actor_deltas, vec![])) + // Phase 8: Counts maintained by triggers (followers/following arrays) + Ok((vec![], vec![])) } else { - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } DatabaseOperation::InsertBlock { @@ -353,7 +322,7 @@ pub async fn execute_operation( record, } => { db::block_insert(conn, rkey, actor_id, subject_actor_id, cid, record).await?; - Ok((vec![], vec![], vec![])) // Blocks don't affect aggregates (yet) + Ok((vec![], vec![])) // Blocks don't affect aggregates (yet) } DatabaseOperation::InsertNotification { recipient_actor_id, @@ -369,7 +338,7 @@ pub async fn execute_operation( } => { // Skip self-notifications (author notifying themselves) if recipient_actor_id == author_actor_id { - return Ok((vec![], vec![], vec![])); + return Ok((vec![], vec![])); } // Query recipient DID for allowlist check @@ -378,12 +347,12 @@ pub async fn execute_operation( let Some(recipient_did) = recipient_did else { tracing::debug!("Skipping notification - recipient actor_id not found: {}", recipient_actor_id); - return Ok((vec![], vec![], vec![])); + return Ok((vec![], vec![])); }; // PostgreSQL-based notifications (only for allowlisted recipients) if !allowlist.cache.contains_did(&recipient_did) { - return Ok((vec![], vec![], vec![])); + return Ok((vec![], vec![])); } // Check if thread is muted (for reply, quote, and mention notifications) @@ -398,7 +367,7 @@ pub async fn execute_operation( subj_actor_id, subj_rkey ); - return Ok((vec![], vec![], vec![])); + return Ok((vec![], vec![])); } } @@ -435,7 +404,7 @@ pub async fn execute_operation( ); } - Ok((vec![], vec![], vec![])) // Notifications don't affect aggregates + Ok((vec![], vec![])) // Notifications don't affect aggregates } DatabaseOperation::NotifyReplyChain { reply_uri, @@ -677,7 +646,7 @@ pub async fn execute_operation( ); } - Ok((vec![], vec![], vec![])) // Notifications don't affect aggregates + Ok((vec![], vec![])) // Notifications don't affect aggregates } DatabaseOperation::UpsertActor { did, @@ -697,12 +666,12 @@ pub async fn execute_operation( timestamp, ) .await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::EnqueueFetch { uri, .. } => { // Enqueue to PostgreSQL fetch_queue table crate::db::fetch_queue::enqueue(conn, &uri).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertProfile { actor_id, cid, record } => { // Query DID from actor_id (needed by db functions) @@ -727,14 +696,14 @@ pub async fn execute_operation( crate::db::handle_resolution_queue::enqueue(pool, &did).await?; } - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertStatus { actor_id, cid, record } => { // Query DID from actor_id (needed by db functions) let did = db::actor_did_from_id(conn, actor_id).await?; db::status_upsert(conn, actor_id, &did, cid, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertPostgate { rkey, @@ -743,7 +712,7 @@ pub async fn execute_operation( record, } => { db::postgate_upsert(conn, actor_id, rkey, cid, &record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertThreadgate { rkey, @@ -752,7 +721,7 @@ pub async fn execute_operation( record, } => { db::threadgate_upsert(conn, actor_id, rkey, cid, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertList { rkey, @@ -763,10 +732,10 @@ pub async fn execute_operation( // Lists use arbitrary string rkeys (not TID-based like posts) 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![])) + // Phase 8: lists_count maintained by application (not from array) + Ok((vec![], vec![])) } else { - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } DatabaseOperation::InsertListItem { @@ -778,7 +747,7 @@ pub async fn execute_operation( } => { db::list_item_insert(conn, actor_id, rkey, subject_actor_id, cid, record).await?; // TODO: Could track list item_count here - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::InsertListBlock { rkey, @@ -787,7 +756,7 @@ pub async fn execute_operation( record, } => { db::list_block_insert(conn, actor_id, rkey, cid, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertFeedGenerator { rkey, @@ -799,10 +768,10 @@ pub async fn execute_operation( // Service actor ID already resolved during reference extraction let inserted = db::feedgen_upsert(conn, actor_id, &rkey, cid, service_actor_id, record).await?; if inserted { - let actor_deltas = actor_feedgen_delta(actor_id, 1); - Ok((vec![], actor_deltas, vec![])) + // Phase 8: feeds_count maintained by application (not from array) + Ok((vec![], vec![])) } else { - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } DatabaseOperation::UpsertStarterPack { @@ -813,10 +782,10 @@ pub async fn execute_operation( } => { let inserted = db::starter_pack_upsert(conn, actor_id, rkey, cid, record).await?; if inserted { - let actor_deltas = actor_starterpack_delta(actor_id, 1); - Ok((vec![], actor_deltas, vec![])) + // Phase 8: starterpacks_count maintained by application (not from array) + Ok((vec![], vec![])) } else { - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } DatabaseOperation::InsertVerification { @@ -827,26 +796,26 @@ pub async fn execute_operation( record, } => { db::verification_insert(conn, actor_id, rkey, subject_actor_id, cid, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertLabeler { actor_id, cid, record } => { // Query DID from actor_id (needed by db functions) let did = db::actor_did_from_id(conn, actor_id).await?; db::labeler_upsert(conn, actor_id, &did, cid, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertNotificationDeclaration { actor_id, record } => { db::notif_decl_upsert(conn, actor_id, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertChatDeclaration { actor_id, record } => { db::chat_decl_upsert(conn, actor_id, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::UpsertBookmark { actor_id, rkey, record } => { db::bookmark_upsert(conn, actor_id, rkey, record).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::DeleteRecord { actor_id, @@ -856,8 +825,7 @@ pub async fn execute_operation( } => { use crate::relay::types::CollectionType; - let mut deltas = Vec::new(); - let mut actor_deltas = Vec::new(); + let deltas = Vec::new(); // Phase 8: Empty - arrays updated directly let mut cache_invalidations = Vec::new(); // actor_id is already resolved in the resolution phase @@ -955,10 +923,7 @@ pub async fn execute_operation( CollectionType::BskyFollow => { // Returns subject (target DID) if deleted if let Some(target_did) = db::follow_delete(conn, rkey, actor_id).await? { - // Resolve target DID to actor_id for deltas - if let Ok(Some(subject_actor_id)) = db::actor_id_from_did(conn, &target_did).await { - actor_deltas.extend(follow_deltas(actor_id, subject_actor_id, -1)); - } + // Phase 8: Counts maintained by triggers (followers/following arrays) // Remove PostgreSQL notification (recipient = followed user) if allowlist.cache.contains_did(&target_did) { @@ -989,8 +954,7 @@ pub async fn execute_operation( if let Some((parent_key, embedded_key)) = db::post_delete(conn, actor_id, rkey).await? { - // Decrement actor's post count - actor_deltas.extend(actor_post_delta(actor_id, -1)); + // Phase 8: posts_count maintained by triggers (post_rkeys array) // Remove from parent's reply arrays if this was a reply if let Some((parent_actor_id, parent_rkey)) = parent_key { @@ -1080,7 +1044,7 @@ pub async fn execute_operation( .ok_or_else(|| eyre::eyre!("Invalid at_uri: missing rkey"))?; let rows = db::feedgen_delete(conn, actor_id, rkey_str).await?; if rows > 0 { - actor_deltas.extend(actor_feedgen_delta(actor_id, -1)); + // Phase 8: feeds_count maintained by application (not from array) } } CollectionType::BskyFeedPostgate => { @@ -1096,7 +1060,7 @@ pub async fn execute_operation( let rows = db::list_delete(conn, actor_id, rkey_str).await?; if rows > 0 { - actor_deltas.extend(actor_list_delta(actor_id, -1)); + // Phase 8: lists_count maintained by application (not from array) } } CollectionType::BskyListBlock => { @@ -1108,7 +1072,7 @@ pub async fn execute_operation( CollectionType::BskyStarterPack => { let rows = db::starter_pack_delete(conn, actor_id, rkey).await?; if rows > 0 { - actor_deltas.extend(actor_starterpack_delta(actor_id, -1)); + // Phase 8: starterpacks_count maintained by application (not from array) } } CollectionType::BskyVerification => { @@ -1133,7 +1097,7 @@ pub async fn execute_operation( // Note: No generic records table to delete from anymore - Ok((deltas, actor_deltas, cache_invalidations)) + Ok((deltas, cache_invalidations)) } DatabaseOperation::MaintainSelfLabels { actor_id, @@ -1145,7 +1109,7 @@ pub async fn execute_operation( let did = db::actor_did_from_id(conn, actor_id).await?; db::maintain_self_labels(conn, &did, cid, &at_uri, labels).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::MaintainPostgateDetaches { post_uri, @@ -1157,7 +1121,7 @@ pub async fn execute_operation( // - Uses direct database lookups for IDs // - Does single efficient UPDATE db::postgate_maintain_detaches_cached(conn, &post_uri, &detached_uris, disable_effective).await?; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::CachePdsMapping { did, @@ -1166,7 +1130,7 @@ pub async fn execute_operation( } => { // Cache PDS mapping in moka-based cache (in-memory, TTL: 24h) pds_cache.set(did, host).await; - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } DatabaseOperation::EnsureMinimumPostStats { post_uri: _, @@ -1183,7 +1147,7 @@ pub async fn execute_operation( // // For now, this operation is a no-op to avoid breaking callers. // TODO: Remove this operation entirely once callers are updated. - Ok((vec![], vec![], vec![])) + Ok((vec![], vec![])) } } } diff --git a/consumer/src/database_writer/operations/mod.rs b/consumer/src/database_writer/operations/mod.rs index 109b4e97..f622c199 100644 --- a/consumer/src/database_writer/operations/mod.rs +++ b/consumer/src/database_writer/operations/mod.rs @@ -122,7 +122,7 @@ pub mod processor; pub mod types; pub use processor::process_record_to_operations; -pub use types::{ActorStatsDeltas, DatabaseOperation, PostStatsDeltas}; +pub use types::{DatabaseOperation, PostStatsDeltas}; // SelfLabels available via types::SelfLabels if needed, but most code imports from lexica directly // Notification helpers re-exported for convenience pub use notifications::*; diff --git a/consumer/src/database_writer/operations/types.rs b/consumer/src/database_writer/operations/types.rs index 724a8afc..28da528a 100644 --- a/consumer/src/database_writer/operations/types.rs +++ b/consumer/src/database_writer/operations/types.rs @@ -260,15 +260,8 @@ pub struct PostStatsDeltas { // All fields removed - engagement is now tracked via arrays only } -/// Actor stats deltas for batch updating denormalized stats columns in actors table -/// Accumulates deltas per actor during operation processing -/// NULL = 0 in database (better TimescaleDB compression) -#[derive(Default, Debug, Clone)] -pub struct ActorStatsDeltas { - pub followers_delta: i32, // Changed by Follow insert/delete - pub following_delta: i32, // Changed by Follow insert/delete - pub posts_delta: i32, // Changed by Post insert/delete - pub lists_delta: i16, // Changed by List insert/delete - pub feeds_delta: i16, // Changed by Feedgen insert/delete - pub starterpacks_delta: i16, // Changed by Starterpack insert/delete -} +// Phase 8: ActorStatsDeltas removed - all counts now maintained by database triggers +// Triggers automatically maintain: +// - followers_count, following_count (from followers/following arrays) +// - posts_count (from post_rkeys array) +// - lists_count, feeds_count, starterpacks_count (maintained by application) diff --git a/migrations/2025-12-08-043658_add_count_and_notify_triggers/down.sql b/migrations/2025-12-08-043658_add_count_and_notify_triggers/down.sql new file mode 100644 index 00000000..ec70b43b --- /dev/null +++ b/migrations/2025-12-08-043658_add_count_and_notify_triggers/down.sql @@ -0,0 +1,32 @@ +-- Rollback Phase 8: Remove Count Columns and Triggers + +-- Drop NOTIFY triggers +DROP TRIGGER IF EXISTS actor_notify_trigger ON actors; +DROP TRIGGER IF EXISTS post_notify_trigger ON posts; +DROP TRIGGER IF EXISTS feedgen_notify_trigger ON feedgens; +DROP TRIGGER IF EXISTS list_notify_trigger ON lists; +DROP TRIGGER IF EXISTS starterpack_notify_trigger ON starterpacks; + +-- Drop count maintenance triggers +DROP TRIGGER IF EXISTS post_counts_trigger ON posts; +DROP TRIGGER IF EXISTS actor_counts_trigger ON actors; + +-- Drop trigger functions +DROP FUNCTION IF EXISTS notify_actor_change(); +DROP FUNCTION IF EXISTS notify_post_change(); +DROP FUNCTION IF EXISTS notify_feedgen_change(); +DROP FUNCTION IF EXISTS notify_list_change(); +DROP FUNCTION IF EXISTS notify_starterpack_change(); +DROP FUNCTION IF EXISTS update_post_counts(); +DROP FUNCTION IF EXISTS update_actor_counts(); + +-- Drop indexes +DROP INDEX IF EXISTS idx_posts_like_count_desc; +DROP INDEX IF EXISTS idx_posts_repost_count_desc; + +-- Drop count columns from posts +ALTER TABLE posts + DROP COLUMN IF EXISTS like_count, + DROP COLUMN IF EXISTS repost_count, + DROP COLUMN IF EXISTS reply_count, + DROP COLUMN IF EXISTS quote_count; diff --git a/migrations/2025-12-08-043658_add_count_and_notify_triggers/up.sql b/migrations/2025-12-08-043658_add_count_and_notify_triggers/up.sql new file mode 100644 index 00000000..568ec6cd --- /dev/null +++ b/migrations/2025-12-08-043658_add_count_and_notify_triggers/up.sql @@ -0,0 +1,166 @@ +-- Phase 8: Add Count Columns and Triggers for Automatic Maintenance +-- +-- This migration: +-- 1. Adds count columns to posts table +-- 2. Creates triggers to maintain counts automatically from arrays +-- 3. Creates NOTIFY triggers for cache invalidation +-- 4. Creates triggers to maintain actor counts +-- +-- Benefits: +-- - Eliminates delta generation code in application +-- - Atomic count updates (no drift) +-- - Automatic cache invalidation via NOTIFY +-- - Counts always accurate + +-- ============================================================================= +-- STEP 1: Add count columns to posts table +-- ============================================================================= + +ALTER TABLE posts + ADD COLUMN like_count INTEGER, + ADD COLUMN repost_count INTEGER, + ADD COLUMN reply_count INTEGER, + ADD COLUMN quote_count INTEGER; + +-- Backfill counts from existing arrays +UPDATE posts SET + like_count = COALESCE(array_length(like_actor_ids, 1), 0), + repost_count = COALESCE(array_length(repost_actor_ids, 1), 0), + reply_count = COALESCE(array_length(reply_actor_ids, 1), 0), + quote_count = COALESCE(array_length(quote_actor_ids, 1), 0); + +-- Make counts NOT NULL after backfill +ALTER TABLE posts + ALTER COLUMN like_count SET NOT NULL, + ALTER COLUMN like_count SET DEFAULT 0, + ALTER COLUMN repost_count SET NOT NULL, + ALTER COLUMN repost_count SET DEFAULT 0, + ALTER COLUMN reply_count SET NOT NULL, + ALTER COLUMN reply_count SET DEFAULT 0, + ALTER COLUMN quote_count SET NOT NULL, + ALTER COLUMN quote_count SET DEFAULT 0; + +-- Add indexes for sorting/filtering by counts +CREATE INDEX idx_posts_like_count_desc ON posts (like_count DESC NULLS LAST) WHERE like_count > 0; +CREATE INDEX idx_posts_repost_count_desc ON posts (repost_count DESC NULLS LAST) WHERE repost_count > 0; + +-- ============================================================================= +-- STEP 2: Create trigger functions for post count maintenance +-- ============================================================================= + +-- Trigger to maintain post engagement counts automatically +CREATE OR REPLACE FUNCTION update_post_counts() +RETURNS TRIGGER AS $$ +BEGIN + NEW.like_count := COALESCE(array_length(NEW.like_actor_ids, 1), 0); + NEW.repost_count := COALESCE(array_length(NEW.repost_actor_ids, 1), 0); + NEW.reply_count := COALESCE(array_length(NEW.reply_actor_ids, 1), 0); + NEW.quote_count := COALESCE(array_length(NEW.quote_actor_ids, 1), 0); + RETURN NEW; +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER post_counts_trigger + BEFORE INSERT OR UPDATE OF like_actor_ids, repost_actor_ids, reply_actor_ids, quote_actor_ids ON posts + FOR EACH ROW + EXECUTE FUNCTION update_post_counts(); + +-- ============================================================================= +-- STEP 3: Create trigger functions for actor count maintenance +-- ============================================================================= + +-- Trigger to maintain actor aggregate counts automatically +CREATE OR REPLACE FUNCTION update_actor_counts() +RETURNS TRIGGER AS $$ +BEGIN + -- Social graph counts + NEW.followers_count := COALESCE(array_length(NEW.followers, 1), 0); + NEW.following_count := COALESCE(array_length(NEW.following, 1), 0); + + -- Content counts (from rkey arrays) + -- Note: Only post_rkeys and repost_rkeys exist as arrays + -- lists_count, feeds_count, starterpacks_count are maintained separately by application + NEW.posts_count := COALESCE(array_length(NEW.post_rkeys, 1), 0); + + RETURN NEW; +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER actor_counts_trigger + BEFORE INSERT OR UPDATE OF followers, following, post_rkeys ON actors + FOR EACH ROW + EXECUTE FUNCTION update_actor_counts(); + +-- ============================================================================= +-- STEP 4: Create NOTIFY triggers for cache invalidation +-- ============================================================================= + +-- Actor cache invalidation +CREATE OR REPLACE FUNCTION notify_actor_change() +RETURNS TRIGGER AS $$ +BEGIN + PERFORM pg_notify('cache_invalidate', 'actor:' || COALESCE(NEW.id, OLD.id)::text); + RETURN COALESCE(NEW, OLD); +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER actor_notify_trigger + AFTER INSERT OR UPDATE OR DELETE ON actors + FOR EACH ROW + EXECUTE FUNCTION notify_actor_change(); + +-- Post cache invalidation +CREATE OR REPLACE FUNCTION notify_post_change() +RETURNS TRIGGER AS $$ +BEGIN + PERFORM pg_notify('cache_invalidate', 'post:' || COALESCE(NEW.actor_id, OLD.actor_id)::text || ':' || COALESCE(NEW.rkey, OLD.rkey)::text); + RETURN COALESCE(NEW, OLD); +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER post_notify_trigger + AFTER INSERT OR UPDATE OR DELETE ON posts + FOR EACH ROW + EXECUTE FUNCTION notify_post_change(); + +-- FeedGen cache invalidation +CREATE OR REPLACE FUNCTION notify_feedgen_change() +RETURNS TRIGGER AS $$ +BEGIN + PERFORM pg_notify('cache_invalidate', 'feedgen:' || COALESCE(NEW.actor_id, OLD.actor_id)::text || ':' || COALESCE(NEW.rkey, OLD.rkey)::text); + RETURN COALESCE(NEW, OLD); +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER feedgen_notify_trigger + AFTER INSERT OR UPDATE OR DELETE ON feedgens + FOR EACH ROW + EXECUTE FUNCTION notify_feedgen_change(); + +-- List cache invalidation +CREATE OR REPLACE FUNCTION notify_list_change() +RETURNS TRIGGER AS $$ +BEGIN + PERFORM pg_notify('cache_invalidate', 'list:' || COALESCE(NEW.actor_id, OLD.actor_id)::text || ':' || COALESCE(NEW.rkey, OLD.rkey)::text); + RETURN COALESCE(NEW, OLD); +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER list_notify_trigger + AFTER INSERT OR UPDATE OR DELETE ON lists + FOR EACH ROW + EXECUTE FUNCTION notify_list_change(); + +-- StarterPack cache invalidation +CREATE OR REPLACE FUNCTION notify_starterpack_change() +RETURNS TRIGGER AS $$ +BEGIN + PERFORM pg_notify('cache_invalidate', 'starterpack:' || COALESCE(NEW.actor_id, OLD.actor_id)::text || ':' || COALESCE(NEW.rkey, OLD.rkey)::text); + RETURN COALESCE(NEW, OLD); +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER starterpack_notify_trigger + AFTER INSERT OR UPDATE OR DELETE ON starterpacks + FOR EACH ROW + EXECUTE FUNCTION notify_starterpack_change(); diff --git a/parakeet-db/src/models.rs b/parakeet-db/src/models.rs index cd801629..9c02d378 100644 --- a/parakeet-db/src/models.rs +++ b/parakeet-db/src/models.rs @@ -258,6 +258,11 @@ pub struct Post { pub postgate_detached_rkeys: Option>>, // Detached embed rkeys (parallel array) // Phase 7: Labels array (from labels table, denormalized) pub labels: Option>>, // Labels applied to this post + // Phase 8: Engagement counts (maintained automatically by triggers) + pub like_count: i32, // Count of likes (maintained by update_post_counts trigger) + pub repost_count: i32, // Count of reposts (maintained by update_post_counts trigger) + pub reply_count: i32, // Count of replies (maintained by update_post_counts trigger) + pub quote_count: i32, // Count of quote posts (maintained by update_post_counts trigger) // Note: created_at derived from TID rkey via created_at() method } @@ -325,28 +330,6 @@ impl Post { .map(|v| v.len()) .unwrap_or(0) } - - // Count helpers (compute from array lengths, replacing old count columns) - - /// Get like count from array length - pub fn like_count(&self) -> usize { - self.like_actor_ids.as_ref().map_or(0, |v| v.len()) - } - - /// Get reply count from array length - pub fn reply_count(&self) -> usize { - self.reply_actor_ids.as_ref().map_or(0, |v| v.len()) - } - - /// Get quote count from array length - pub fn quote_count(&self) -> usize { - self.quote_actor_ids.as_ref().map_or(0, |v| v.len()) - } - - /// Get repost count from array length - pub fn repost_count(&self) -> usize { - self.repost_actor_ids.as_ref().map_or(0, |v| v.len()) - } } /// Represents a single like extracted from post arrays @@ -377,10 +360,10 @@ pub struct PostStats { impl PostStats { pub fn from_post(post: &Post) -> Self { Self { - likes: post.like_count() as i32, - replies: post.reply_count() as i32, - reposts: post.repost_count() as i32, - quotes: post.quote_count() as i32, + likes: post.like_count, + replies: post.reply_count, + reposts: post.repost_count, + quotes: post.quote_count, } } diff --git a/parakeet-db/src/schema.rs b/parakeet-db/src/schema.rs index a043f87e..0bcbf3ef 100644 --- a/parakeet-db/src/schema.rs +++ b/parakeet-db/src/schema.rs @@ -497,6 +497,10 @@ diesel::table! { postgate_detached_actor_ids -> Nullable>>, postgate_detached_rkeys -> Nullable>>, labels -> Nullable>>, + like_count -> Int4, + repost_count -> Int4, + reply_count -> Int4, + quote_count -> Int4, } } diff --git a/parakeet/src/loaders/post.rs b/parakeet/src/loaders/post.rs index 940c373a..e3ba648d 100644 --- a/parakeet/src/loaders/post.rs +++ b/parakeet/src/loaders/post.rs @@ -147,6 +147,15 @@ pub struct PostWithComputed { pub repost_actor_ids: Option>>, #[diesel(sql_type = diesel::sql_types::Nullable>>)] pub repost_rkeys: Option>>, + // Phase 8: Engagement counts (maintained by triggers) + #[diesel(sql_type = diesel::sql_types::Integer)] + pub like_count: i32, + #[diesel(sql_type = diesel::sql_types::Integer)] + pub repost_count: i32, + #[diesel(sql_type = diesel::sql_types::Integer)] + pub reply_count: i32, + #[diesel(sql_type = diesel::sql_types::Integer)] + pub quote_count: i32, } /// Build SQL query for batch loading posts with computed fields @@ -205,7 +214,12 @@ pub fn build_posts_batch_query() -> &'static str { p.quote_actor_ids, p.quote_rkeys, p.repost_actor_ids, - p.repost_rkeys + p.repost_rkeys, + -- Phase 8: Engagement counts (maintained by triggers) + p.like_count, + p.repost_count, + p.reply_count, + p.quote_count FROM posts p INNER JOIN unnest($1::integer[], $2::bigint[]) AS lookup(lookup_actor_id, lookup_rkey) ON p.actor_id = lookup.lookup_actor_id AND p.rkey = lookup.lookup_rkey @@ -938,6 +952,11 @@ impl BatchFn for PostLoader { postgate_detached_rkeys: None, // Phase 7: Labels (not loaded for hydration - will be loaded separately if needed) labels: None, + // Phase 8: Engagement counts (loaded from database, maintained by triggers) + like_count: p.like_count, + repost_count: p.repost_count, + reply_count: p.reply_count, + quote_count: p.quote_count, }; // Encode TIDs using Rust utility functions @@ -1279,12 +1298,12 @@ impl BatchFn for PostLoader { &created_at, ); - // Create PostStats from array lengths using helper methods + // Create PostStats from count columns (maintained by triggers) let stats = parakeet_db::models::PostStats { - likes: post.like_count() as i32, - replies: post.reply_count() as i32, - reposts: post.repost_count() as i32, - quotes: post.quote_count() as i32, + likes: post.like_count, + replies: post.reply_count, + reposts: post.repost_count, + quotes: post.quote_count, }; let hydrated_post = HydratedPost {