From 59ab6d91ab3b786edc5eba65adbc58b80cb570e0 Mon Sep 17 00:00:00 2001 From: mdecimus Date: Mon, 20 Oct 2025 12:10:36 +0200 Subject: [PATCH] Sharded cluster node roles --- crates/common/src/config/network.rs | 22 ++++++++++---- crates/common/src/ipc.rs | 5 +++- crates/email/src/message/delete.rs | 17 +++++++++-- crates/http/src/management/stores.rs | 7 +++-- crates/services/src/housekeeper/mod.rs | 38 +++++++++++++++++------- crates/services/src/task_manager/mod.rs | 39 ++++++++++++++++++++----- crates/trc/src/event/description.rs | 2 ++ crates/trc/src/event/level.rs | 1 + crates/trc/src/lib.rs | 1 + crates/trc/src/serializers/binary.rs | 2 ++ 10 files changed, 105 insertions(+), 29 deletions(-) diff --git a/crates/common/src/config/network.rs b/crates/common/src/config/network.rs index d8c9921d..ce55c937 100644 --- a/crates/common/src/config/network.rs +++ b/crates/common/src/config/network.rs @@ -4,11 +4,9 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -use std::time::Duration; - use crate::expr::{if_block::IfBlock, tokenizer::TokenMap}; use ahash::AHashSet; - +use std::time::Duration; use utils::config::{Config, Rate, utils::ParseValue}; use super::*; @@ -356,8 +354,8 @@ impl AsnGeoLookupConfig { } impl ClusterRole { - pub fn is_enabled(&self) -> bool { - matches!(self, ClusterRole::Enabled) + pub fn is_enabled_or_sharded(&self) -> bool { + matches!(self, ClusterRole::Enabled | ClusterRole::Sharded { .. }) } pub fn is_enabled_for_account(&self, account_id: u32) -> bool { @@ -370,4 +368,18 @@ impl ClusterRole { } => (account_id % total_shards) == *shard_id, } } + + pub fn is_enabled_for_hash(&self, item: &impl std::hash::Hash) -> bool { + match self { + ClusterRole::Enabled => true, + ClusterRole::Disabled => false, + ClusterRole::Sharded { + shard_id, + total_shards, + } => { + (ahash::RandomState::new().hash_one(item) % (*total_shards as u64)) + == (*shard_id as u64) + } + } + } } diff --git a/crates/common/src/ipc.rs b/crates/common/src/ipc.rs index eedd8c61..b827d60b 100644 --- a/crates/common/src/ipc.rs +++ b/crates/common/src/ipc.rs @@ -41,7 +41,10 @@ pub enum PurgeType { store: InMemoryStore, prefix: Option>, }, - Account(Option), + Account { + account_id: Option, + use_roles: bool, + }, } #[derive(Debug)] diff --git a/crates/email/src/message/delete.rs b/crates/email/src/message/delete.rs index 60dfbd08..0b35a2e4 100644 --- a/crates/email/src/message/delete.rs +++ b/crates/email/src/message/delete.rs @@ -32,7 +32,7 @@ pub trait EmailDeletion: Sync + Send { document_ids: RoaringBitmap, ) -> impl Future> + Send; - fn purge_accounts(&self) -> impl Future + Send; + fn purge_accounts(&self, use_roles: bool) -> impl Future + Send; fn purge_account(&self, account_id: u32) -> impl Future + Send; @@ -99,10 +99,21 @@ impl EmailDeletion for Server { Ok(not_destroyed) } - async fn purge_accounts(&self) { + async fn purge_accounts(&self, use_roles: bool) { if let Ok(Some(account_ids)) = self.get_document_ids(u32::MAX, Collection::Principal).await { - let mut account_ids: Vec = account_ids.into_iter().collect(); + let mut account_ids: Vec = account_ids + .into_iter() + .filter(|id| { + !use_roles + || self + .core + .network + .roles + .purge_accounts + .is_enabled_for_account(*id) + }) + .collect(); // Shuffle account ids account_ids.shuffle(&mut store::rand::rng()); diff --git a/crates/http/src/management/stores.rs b/crates/http/src/management/stores.rs index 12804b21..cd77cd6d 100644 --- a/crates/http/src/management/stores.rs +++ b/crates/http/src/management/stores.rs @@ -212,8 +212,11 @@ impl ManageStore for Server { None }; - self.housekeeper_request(HousekeeperEvent::Purge(PurgeType::Account(account_id))) - .await + self.housekeeper_request(HousekeeperEvent::Purge(PurgeType::Account { + account_id, + use_roles: false, + })) + .await } (Some("reindex"), id, None, &Method::GET) => { // Validate the access token diff --git a/crates/services/src/housekeeper/mod.rs b/crates/services/src/housekeeper/mod.rs index 7a5f1590..e604cb7f 100644 --- a/crates/services/src/housekeeper/mod.rs +++ b/crates/services/src/housekeeper/mod.rs @@ -78,9 +78,10 @@ pub fn spawn_housekeeper(inner: Arc, mut rx: mpsc::Receiver, mut rx: mpsc::Receiver, mut rx: mpsc::Receiver, mut rx: mpsc::Receiver { queue.schedule( @@ -343,7 +344,15 @@ pub fn spawn_housekeeper(inner: Arc, mut rx: mpsc::Receiver { @@ -436,7 +445,13 @@ pub fn spawn_housekeeper(inner: Arc, mut rx: mpsc::Receiver // SPDX-License-Identifier: LicenseRef-SEL @@ -665,7 +680,7 @@ impl Purge for Server { .into(), ), PurgeType::Lookup { .. } => ("in-memory-prefix", None), - PurgeType::Account(_) => ("account", None), + PurgeType::Account { .. } => ("account", None), }; if let Some(lock_name) = &lock_name { match self @@ -750,11 +765,14 @@ impl Purge for Server { trc::error!(err.details("Failed to purge in-memory store")); } } - PurgeType::Account(account_id) => { + PurgeType::Account { + account_id, + use_roles, + } => { if let Some(account_id) = account_id { self.purge_account(account_id).await; } else { - self.purge_accounts().await; + self.purge_accounts(use_roles).await; } } } diff --git a/crates/services/src/task_manager/mod.rs b/crates/services/src/task_manager/mod.rs index d2cc145d..5950b844 100644 --- a/crates/services/src/task_manager/mod.rs +++ b/crates/services/src/task_manager/mod.rs @@ -95,14 +95,14 @@ pub fn spawn_task_manager(inner: Arc) { for mut rx_index in [rx_index_1, rx_index_2, rx_index_3, rx_index_4] { let inner = inner.clone(); let server_instance = server_instance.clone(); - let todo = "shard tasks based on config"; tokio::spawn(async move { - while let Some(mut task) = rx_index.recv().await { + while let Some(task) = rx_index.recv().await { let server = inner.build_server(); + // Lock task if server.try_lock_task(&task).await { - let success = match &mut task.action { + let success = match &task.action { TaskAction::Index { hash } => { server .fts_index(task.account_id, task.document_id, hash) @@ -262,7 +262,7 @@ impl TaskQueueManager for Server { .map_err(|err| { trc::error!( err.caused_by(trc::location!()) - .details("Failed to iterate over index emails") + .details("Failed to iterate over task queue.") ); }); @@ -279,12 +279,35 @@ impl TaskQueueManager for Server { tasks.shuffle(&mut rand::rng()); } + // Dispatch tasks + let roles = &self.core.network.roles; for event in tasks { let tx = match &event.action { - TaskAction::Index { .. } => &ipc.tx_fts, - TaskAction::BayesTrain { .. } => &ipc.tx_bayes, - TaskAction::SendAlarm { .. } => &ipc.tx_alarm, - TaskAction::SendImip => &ipc.tx_imip, + TaskAction::Index { .. } if roles.fts_indexing.is_enabled_for_hash(&event) => { + &ipc.tx_fts + } + TaskAction::BayesTrain { .. } + if roles.bayes_training.is_enabled_for_hash(&event) => + { + &ipc.tx_bayes + } + TaskAction::SendAlarm { .. } + if roles.calendar_alerts.is_enabled_for_hash(&event) => + { + &ipc.tx_alarm + } + TaskAction::SendImip if roles.imip_processing.is_enabled_for_hash(&event) => { + &ipc.tx_imip + } + _ => { + trc::event!( + TaskQueue(TaskQueueEvent::TaskIgnored), + AccountId = event.account_id, + DocumentId = event.document_id, + ); + + continue; + } }; if tx.send(event).await.is_err() { trc::event!( diff --git a/crates/trc/src/event/description.rs b/crates/trc/src/event/description.rs index e310edee..04ce2ad1 100644 --- a/crates/trc/src/event/description.rs +++ b/crates/trc/src/event/description.rs @@ -195,6 +195,7 @@ impl TaskQueueEvent { TaskQueueEvent::TaskLocked => "Task is locked by another process", TaskQueueEvent::BlobNotFound => "Blob not found for task", TaskQueueEvent::MetadataNotFound => "Metadata not found for task", + TaskQueueEvent::TaskIgnored => "Task ignored based on current server roles", } } @@ -204,6 +205,7 @@ impl TaskQueueEvent { TaskQueueEvent::TaskLocked => "The task id is locked by another process", TaskQueueEvent::BlobNotFound => "The requested blob was not found for task", TaskQueueEvent::MetadataNotFound => "The metadata was not found for task", + TaskQueueEvent::TaskIgnored => "The task was ignored based on the current server roles", } } } diff --git a/crates/trc/src/event/level.rs b/crates/trc/src/event/level.rs index 59576005..459f67a1 100644 --- a/crates/trc/src/event/level.rs +++ b/crates/trc/src/event/level.rs @@ -380,6 +380,7 @@ impl EventType { TaskQueueEvent::BlobNotFound | TaskQueueEvent::TaskAcquired | TaskQueueEvent::TaskLocked + | TaskQueueEvent::TaskIgnored | TaskQueueEvent::MetadataNotFound => Level::Debug, }, EventType::Dmarc(_) => Level::Debug, diff --git a/crates/trc/src/lib.rs b/crates/trc/src/lib.rs index 4be4fb5d..35f290d1 100644 --- a/crates/trc/src/lib.rs +++ b/crates/trc/src/lib.rs @@ -239,6 +239,7 @@ pub enum HousekeeperEvent { pub enum TaskQueueEvent { TaskAcquired, TaskLocked, + TaskIgnored, BlobNotFound, MetadataNotFound, } diff --git a/crates/trc/src/serializers/binary.rs b/crates/trc/src/serializers/binary.rs index c0795bbe..442fba1b 100644 --- a/crates/trc/src/serializers/binary.rs +++ b/crates/trc/src/serializers/binary.rs @@ -893,6 +893,7 @@ impl EventType { EventType::Calendar(CalendarEvent::ItipMessageSent) => 583, EventType::Calendar(CalendarEvent::ItipMessageReceived) => 584, EventType::Calendar(CalendarEvent::ItipMessageError) => 585, + EventType::TaskQueue(TaskQueueEvent::TaskIgnored) => 586, } } @@ -1524,6 +1525,7 @@ impl EventType { 583 => Some(EventType::Calendar(CalendarEvent::ItipMessageSent)), 584 => Some(EventType::Calendar(CalendarEvent::ItipMessageReceived)), 585 => Some(EventType::Calendar(CalendarEvent::ItipMessageError)), + 586 => Some(EventType::TaskQueue(TaskQueueEvent::TaskIgnored)), _ => None, } }