diff --git a/consumer/src/database_writer/mod.rs b/consumer/src/database_writer/mod.rs index 3361a0df..be7d7209 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, DatabaseOperation, PostStatsDeltas}; +pub use operations::{process_record_to_operations, DatabaseOperation}; 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 18a4246e..141035b8 100644 --- a/consumer/src/database_writer/operations/executor.rs +++ b/consumer/src/database_writer/operations/executor.rs @@ -117,12 +117,10 @@ pub fn describe_operation(op: &DatabaseOperation) -> (String, Option) { } } -/// Execute a single database operation and return any post stats deltas, actor stats deltas, and cache invalidations +/// Execute a single database operation and return cache invalidations /// /// This function will be called by the database writer for each operation. -/// It returns: -/// - 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) +/// It returns cache invalidations that should be sent via NOTIFY for real-time cache updates. /// /// Counts are maintained by database triggers. /// Notifications, fetch queue, and PDS cache use PostgreSQL/moka. @@ -132,7 +130,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)> { +) -> eyre::Result> { use crate::db; match op { @@ -180,7 +178,6 @@ pub async fn execute_operation( ); } if rows > 0 { - let deltas = Vec::new(); // Arrays updated directly by SQL let mut cache_invalidations = Vec::new(); // If this is a reply, increment parent post's reply count @@ -221,9 +218,9 @@ pub async fn execute_operation( } } - Ok((deltas, cache_invalidations)) + Ok(cache_invalidations) } else { - Ok((vec![], vec![])) + Ok(vec![]) } } DatabaseOperation::InsertLike { @@ -257,9 +254,9 @@ pub async fn execute_operation( // Invalidate post cache since like array changed let cache_invalidations = vec![cache::invalidate_post(&subject_uri)]; - Ok((vec![], cache_invalidations)) + Ok(cache_invalidations) } else { - Ok((vec![], vec![])) + Ok(vec![]) } } DatabaseOperation::InsertRepost { @@ -294,9 +291,9 @@ pub async fn execute_operation( // Invalidate post cache since repost array changed let cache_invalidations = vec![cache::invalidate_post(&subject_uri)]; - Ok((vec![], cache_invalidations)) + Ok(cache_invalidations) } else { - Ok((vec![], vec![])) + Ok(vec![]) } } DatabaseOperation::InsertFollow { @@ -309,9 +306,9 @@ pub async fn execute_operation( let rows = db::follow_insert(conn, rkey, actor_id, subject_actor_id, cid, record).await?; if rows > 0 { // Counts maintained by triggers - Ok((vec![], vec![])) + Ok(vec![]) } else { - Ok((vec![], vec![])) + Ok(vec![]) } } DatabaseOperation::InsertBlock { @@ -322,7 +319,7 @@ pub async fn execute_operation( record, } => { db::block_insert(conn, rkey, actor_id, subject_actor_id, cid, record).await?; - Ok((vec![], vec![])) // Blocks don't affect aggregates (yet) + Ok(vec![]) // Blocks don't affect aggregates (yet) } DatabaseOperation::InsertNotification { recipient_actor_id, @@ -338,7 +335,7 @@ pub async fn execute_operation( } => { // Skip self-notifications (author notifying themselves) if recipient_actor_id == author_actor_id { - return Ok((vec![], vec![])); + return Ok(vec![]); } // Query recipient DID for allowlist check @@ -346,12 +343,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![])); + return Ok(vec![]); }; // PostgreSQL-based notifications (only for allowlisted recipients) if !allowlist.cache.contains_did(&recipient_did) { - return Ok((vec![], vec![])); + return Ok(vec![]); } // Check if thread is muted (for reply, quote, and mention notifications) @@ -366,7 +363,7 @@ pub async fn execute_operation( subj_actor_id, subj_rkey ); - return Ok((vec![], vec![])); + return Ok(vec![]); } } @@ -403,7 +400,7 @@ pub async fn execute_operation( ); } - Ok((vec![], vec![])) // Notifications don't affect aggregates + Ok(vec![]) // Notifications don't affect aggregates } DatabaseOperation::NotifyReplyChain { reply_uri, @@ -645,7 +642,7 @@ pub async fn execute_operation( ); } - Ok((vec![], vec![])) // Notifications don't affect aggregates + Ok(vec![]) // Notifications don't affect aggregates } DatabaseOperation::UpsertActor { did, @@ -665,12 +662,12 @@ pub async fn execute_operation( timestamp, ) .await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::EnqueueFetch { uri, .. } => { // Enqueue to PostgreSQL fetch_queue table crate::db::fetch_queue::enqueue(conn, &uri).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::UpsertProfile { actor_id, cid, record } => { // Query DID from actor_id (needed by db functions) @@ -695,14 +692,14 @@ pub async fn execute_operation( crate::db::handle_resolution_queue::enqueue(pool, &did).await?; } - Ok((vec![], vec![])) + Ok(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![])) + Ok(vec![]) } DatabaseOperation::UpsertPostgate { rkey, @@ -711,7 +708,7 @@ pub async fn execute_operation( record, } => { db::postgate_upsert(conn, actor_id, rkey, cid, &record).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::UpsertThreadgate { rkey, @@ -720,7 +717,7 @@ pub async fn execute_operation( record, } => { db::threadgate_upsert(conn, actor_id, rkey, cid, record).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::UpsertList { rkey, @@ -732,9 +729,9 @@ pub async fn execute_operation( let inserted = db::list_upsert(conn, actor_id, &rkey, cid, record).await?; if inserted { // lists_count maintained by application - Ok((vec![], vec![])) + Ok(vec![]) } else { - Ok((vec![], vec![])) + Ok(vec![]) } } DatabaseOperation::InsertListItem { @@ -745,7 +742,7 @@ pub async fn execute_operation( record, } => { db::list_item_insert(conn, actor_id, rkey, subject_actor_id, cid, record).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::InsertListBlock { rkey, @@ -754,7 +751,7 @@ pub async fn execute_operation( record, } => { db::list_block_insert(conn, actor_id, rkey, cid, record).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::UpsertFeedGenerator { rkey, @@ -767,9 +764,9 @@ pub async fn execute_operation( let inserted = db::feedgen_upsert(conn, actor_id, &rkey, cid, service_actor_id, record).await?; if inserted { // feeds_count maintained by application - Ok((vec![], vec![])) + Ok(vec![]) } else { - Ok((vec![], vec![])) + Ok(vec![]) } } DatabaseOperation::UpsertStarterPack { @@ -781,9 +778,9 @@ pub async fn execute_operation( let inserted = db::starter_pack_upsert(conn, actor_id, rkey, cid, record).await?; if inserted { // starterpacks_count maintained by application - Ok((vec![], vec![])) + Ok(vec![]) } else { - Ok((vec![], vec![])) + Ok(vec![]) } } DatabaseOperation::InsertVerification { @@ -794,26 +791,26 @@ pub async fn execute_operation( record, } => { db::verification_insert(conn, actor_id, rkey, subject_actor_id, cid, record).await?; - Ok((vec![], vec![])) + Ok(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![])) + Ok(vec![]) } DatabaseOperation::UpsertNotificationDeclaration { actor_id, record } => { db::notif_decl_upsert(conn, actor_id, record).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::UpsertChatDeclaration { actor_id, record } => { db::chat_decl_upsert(conn, actor_id, record).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::UpsertBookmark { actor_id, rkey, record } => { db::bookmark_upsert(conn, actor_id, rkey, record).await?; - Ok((vec![], vec![])) + Ok(vec![]) } DatabaseOperation::DeleteRecord { actor_id, @@ -823,7 +820,6 @@ pub async fn execute_operation( } => { use crate::relay::types::CollectionType; - let deltas = Vec::new(); // Arrays updated directly let mut cache_invalidations = Vec::new(); // actor_id is already resolved in the resolution phase @@ -1095,7 +1091,7 @@ pub async fn execute_operation( // Note: No generic records table to delete from anymore - Ok((deltas, cache_invalidations)) + Ok(cache_invalidations) } DatabaseOperation::MaintainSelfLabels { actor_id, @@ -1107,7 +1103,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![])) + Ok(vec![]) } DatabaseOperation::MaintainPostgateDetaches { post_uri, @@ -1119,7 +1115,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![])) + Ok(vec![]) } DatabaseOperation::CachePdsMapping { did, @@ -1128,7 +1124,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![])) + Ok(vec![]) } DatabaseOperation::EnsureMinimumPostStats { post_uri: _, @@ -1145,7 +1141,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![])) + Ok(vec![]) } } } diff --git a/consumer/src/database_writer/operations/mod.rs b/consumer/src/database_writer/operations/mod.rs index f622c199..e93d8c6c 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::{DatabaseOperation, PostStatsDeltas}; +pub use types::DatabaseOperation; // 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 b448e566..65061037 100644 --- a/consumer/src/database_writer/operations/types.rs +++ b/consumer/src/database_writer/operations/types.rs @@ -251,16 +251,3 @@ pub enum DatabaseOperation { min_replies: u64, }, } - -/// Post stats deltas (DEPRECATED - now using array-only tracking) -/// This struct is kept for backwards compatibility during migration but is no longer used. -/// All functions should return empty HashMaps/Vecs for this type. -#[derive(Default, Debug, Clone)] -pub struct PostStatsDeltas { - // All fields removed - engagement is now tracked via arrays only -} - -// ActorStatsDeltas removed - all counts maintained by database triggers: -// - 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/consumer/src/database_writer/workers.rs b/consumer/src/database_writer/workers.rs index add008de..7142f629 100644 --- a/consumer/src/database_writer/workers.rs +++ b/consumer/src/database_writer/workers.rs @@ -71,7 +71,7 @@ fn spawn_bulk_worker( &allowlist, &pds_cache, ).await { - Ok(_stats_deltas) => { + Ok(()) => { // All counts maintained by database triggers } Err(e) => { @@ -105,8 +105,7 @@ async fn execute_bulk_operations( _source: super::EventSource, allowlist: &crate::db::Allowlist, pds_cache: &crate::external::pds_cache::PdsHostCache, -) -> eyre::Result> { - use std::collections::HashMap; +) -> eyre::Result<()> { let start = std::time::Instant::now(); let total_ops = bulk_ops.total_count(); @@ -132,12 +131,6 @@ async fn execute_bulk_operations( let txn = conn.transaction().await?; let mut total_inserted = 0; - let stats_deltas: HashMap<(i32, i64), super::PostStatsDeltas> = HashMap::new(); - - // Helper to accumulate deltas (DEPRECATED - arrays updated directly now) - let accumulate_delta = |_actor_id: i32, _rkey: i64, _delta: super::PostStatsDeltas| { - // No-op: engagement is now tracked via arrays, updated directly during operations - }; // Execute COPY for post likes (per-actor updates to respect decompression limits) if !bulk_ops.post_likes.is_empty() { @@ -148,8 +141,7 @@ async fn execute_bulk_operations( let op_elapsed = op_start.elapsed(); total_inserted += inserted.len(); - // Note: Like arrays are updated directly in bulk_copy, no deltas needed - // accumulate_delta is now a no-op + // Like arrays are updated directly in bulk_copy tracing::info!( count = count, @@ -225,8 +217,7 @@ async fn execute_bulk_operations( Ok(inserted) => { total_inserted += inserted.len(); - // Note: Repost arrays are updated directly in bulk_copy, no deltas needed - // accumulate_delta is now a no-op + // Repost arrays are updated directly in bulk_copy tracing::info!(inserted = inserted.len(), "Bulk COPY reposts completed"); } @@ -275,8 +266,7 @@ async fn execute_bulk_operations( // Generate deltas for posts for post_data in &bulk_ops.posts { - // Note: Reply/quote arrays are updated directly in bulk_copy, no deltas needed - // accumulate_delta is now a no-op + // Reply/quote arrays are updated directly in bulk_copy } // Execute child table inserts @@ -313,11 +303,7 @@ async fn execute_bulk_operations( allowlist, pds_cache, ).await { - Ok((op_deltas, _cache_invalidations)) => { - // Accumulate deltas from individual operations (DEPRECATED - no-op) - for ((actor_id, rkey), delta) in op_deltas { - accumulate_delta(actor_id, rkey, delta); - } + Ok(_cache_invalidations) => { // Actor counts maintained by triggers } Err(e) => { @@ -332,12 +318,11 @@ async fn execute_bulk_operations( repo = %repo, total_ops = total_ops, total_inserted = total_inserted, - stats_updates = stats_deltas.len(), duration_ms = elapsed.as_millis(), "Bulk COPY execution completed" ); - Ok(stats_deltas) + Ok(()) } /// Batch update post stats in the database @@ -1153,14 +1138,8 @@ fn spawn_worker(config: WorkerConfig) -> tokio::task::JoinHandle = std::collections::HashMap::new(); let mut event_cache_invalidations = Vec::new(); - // Helper to accumulate deltas (DEPRECATED - arrays updated directly now) - let accumulate_delta = |_actor_id: i32, _rkey: i64, _delta: super::PostStatsDeltas| { - // No-op: engagement is now tracked via arrays, updated directly during operations - }; - // Actor stats maintained by triggers let db_start = std::time::Instant::now(); @@ -1181,14 +1160,8 @@ fn spawn_worker(config: WorkerConfig) -> tokio::task::JoinHandle { - if !deltas.is_empty() { - affected += 1; - // Accumulate stats deltas (DEPRECATED - no-op) - for ((actor_id, rkey), delta) in deltas { - accumulate_delta(actor_id, rkey, delta); - } - } + Ok(cache_inv) => { + affected += 1; // Actor counts maintained by triggers event_cache_invalidations.extend(cache_inv); }