Files
inbuxa-server/crates/services/src/task_manager/maintenance.rs
T
jcoffey-dev 441ad0b18e
ci / fork-checks (pull_request) Successful in 41s
ci / build (pull_request) Successful in 8m2s
Journaling: capture at the queue, the built-in journal, retention
Phase 2 of the journaling spec.

- A copy of each message is taken in MessageWrapper::queue, after DLP and
  transport rules, for every enabled journal that takes it (direction and
  scope: everyone, or accounts, groups, domains, tenants). If the copy
  can't be taken the message isn't queued (temporary failure).
- The journal report: the envelope one field a line (sender, To, Cc, Bcc
  from the envelope, list members from their ORCPT, direction, held for
  review), then the queued message byte for byte as message/rfc822.
- The built-in journal under J in the inbuxa subspace: one chain per node
  whose links name each entry by SHA-256, so entries can expire out of
  chain order; purge leaves a marker, and verify catches an entry changed
  or removed early and a report that doesn't match.
- Retention per journal (30 to 3650 days); an entry keeps what it was
  written with. The daily maintenance purges what's due, keeping entries
  whose people a legal hold covers (deleted accounts a hold keeps too),
  and records the counts in the audit log.
- inbuxa:Journal get/set, audited by the request layer. Permissions
  680-683: administrators see and change journals; the Compliance Officer
  sees, searches and exports. Whoever changes journals may grant search and
  export without holding them, so officers can still be appointed.
- Catalog entries (inbuxa:Journal, source "journal"); spec as-built notes.

tests/src/system/journal.rs: validation, internal mail with a Bcc,
outgoing into two journals, incoming over LMTP, the report and its
original, tamper and early removal caught, hold-aware purge, retention
changes leave entries alone, disabled and removed journals take nothing.
2026-09-28 20:46:04 -07:00

676 lines
24 KiB
Rust

/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use std::time::Instant;
use crate::task_manager::{
TaskResult,
index::{reindex_account, reindex_telemetry},
};
use common::{
KV_ACME, KV_GREYLIST, KV_LOCK_DAV, KV_LOCK_QUEUE_MESSAGE, KV_LOCK_TASK, KV_OAUTH,
KV_QUOTA_BLOB, KV_RATE_LIMIT_AUTH, KV_RATE_LIMIT_CONTACT, KV_RATE_LIMIT_HTTP_ANONYMOUS,
KV_RATE_LIMIT_HTTP_AUTHENTICATED, KV_RATE_LIMIT_IMAP, KV_RATE_LIMIT_LOITER, KV_RATE_LIMIT_RCPT,
KV_RATE_LIMIT_SCAN, KV_RATE_LIMIT_SMTP, KV_SIEVE_ID, Server,
storage::index::ObjectIndexBuilder,
};
use email::{
cache::MessageCacheFetch,
message::{delete::EmailDeletion, ingest::EmailIngest, metadata::MessageData},
sieve::SieveScript,
};
use groupware::{
calendar::{Calendar, CalendarEvent, CalendarEventNotification},
contact::{AddressBook, ContactCard},
file::FileNode,
};
use registry::{
schema::{
enums::{TaskAccountMaintenanceType, TaskStoreMaintenanceType, TaskTenantMaintenanceType},
prelude::{Object, ObjectInner, ObjectType, Property},
structs::{
Task, TaskAccountMaintenance, TaskStatus, TaskStoreMaintenance, TaskTenantMaintenance,
},
},
types::EnumImpl,
};
use smtp::reporting::index::ExternalReportIndex;
use store::{
Serialize, ValueKey,
rand::{self},
registry::{RegistryFilter, RegistryQuery},
roaring::RoaringBitmap,
write::{AlignedBytes, Archive, Archiver, BatchBuilder, RegistryClass, ValueClass, now},
};
use trc::{AddContext, StoreEvent};
use types::{
collection::Collection,
field::{EmailField, MailboxField},
id::Id,
};
pub(crate) trait MaintenanceTask: Sync + Send {
fn store_maintenance(
&self,
task: &TaskStoreMaintenance,
) -> impl Future<Output = TaskResult> + Send;
fn account_maintenance(
&self,
task: &TaskAccountMaintenance,
) -> impl Future<Output = TaskResult> + Send;
fn tenant_maintenance(
&self,
task: &TaskTenantMaintenance,
) -> impl Future<Output = TaskResult> + Send;
}
impl MaintenanceTask for Server {
async fn store_maintenance(&self, task: &TaskStoreMaintenance) -> TaskResult {
match store_maintenance(self, task).await {
Ok(result) => result,
Err(err) => {
let result = TaskResult::temporary(err.to_string());
trc::error!(err.details("Failed to perform store maintenance task"));
result
}
}
}
async fn account_maintenance(&self, task: &TaskAccountMaintenance) -> TaskResult {
match account_maintenance(self, task).await {
Ok(result) => result,
Err(err) => {
let result = TaskResult::temporary(err.to_string());
trc::error!(
err.account_id(task.account_id.document_id())
.details("Failed to perform account maintenance task")
);
result
}
}
}
async fn tenant_maintenance(&self, task: &TaskTenantMaintenance) -> TaskResult {
match tenant_maintenance(self, task).await {
Ok(result) => result,
Err(err) => {
let result = TaskResult::temporary(err.to_string());
trc::error!(err.details("Failed to perform tenant maintenance task"));
result
}
}
}
}
async fn store_maintenance(
server: &Server,
task: &TaskStoreMaintenance,
) -> trc::Result<TaskResult> {
match task.maintenance_type {
TaskStoreMaintenanceType::ReindexAccounts
| TaskStoreMaintenanceType::PurgeAccounts
| TaskStoreMaintenanceType::ResetUserQuotas => {
let mut batch = BatchBuilder::new();
let now = now() as i64;
let maintenance_type = match task.maintenance_type {
TaskStoreMaintenanceType::ReindexAccounts => TaskAccountMaintenanceType::Reindex,
TaskStoreMaintenanceType::PurgeAccounts => TaskAccountMaintenanceType::Purge,
TaskStoreMaintenanceType::ResetUserQuotas => {
TaskAccountMaintenanceType::RecalculateQuota
}
_ => unreachable!(),
};
for account_id in server
.registry()
.query::<RoaringBitmap>(RegistryQuery::new(ObjectType::Account))
.await?
{
#[cfg(feature = "test_mode")]
let status = TaskStatus::at(now);
#[cfg(not(feature = "test_mode"))]
let status =
TaskStatus::at(now + rand::RngExt::random_range(&mut rand::rng(), 0..=300));
batch.schedule_task(Task::AccountMaintenance(TaskAccountMaintenance {
account_id: account_id.into(),
maintenance_type,
status,
}));
if batch.is_large_batch() {
server.core.storage.data.write(batch.build_all()).await?;
server.notify_task_queue();
batch = BatchBuilder::new();
}
}
if !batch.is_empty() {
server.core.storage.data.write(batch.build_all()).await?;
server.notify_task_queue();
}
}
TaskStoreMaintenanceType::ReindexTelemetry => {
reindex_telemetry(server).await?;
}
TaskStoreMaintenanceType::PurgeData => {
// Delete expired external reports
let now = now();
let mut batch = BatchBuilder::new();
for object in [
ObjectType::DmarcExternalReport,
ObjectType::TlsExternalReport,
ObjectType::ArfExternalReport,
] {
let ids = server
.registry()
.query::<Vec<Id>>(RegistryQuery::new(object).filter(RegistryFilter::less_than(
Property::ExpiresAt,
now,
false,
)))
.await?;
let object_id = object.to_id();
for id in ids {
let item_id = id.id();
if let Some(report) = server
.store()
.get_value::<Object>(ValueKey::from(ValueClass::Registry(
RegistryClass::Item { object_id, item_id },
)))
.await?
{
match &report.inner {
ObjectInner::DmarcExternalReport(report) => {
report.write_ops(&mut batch, item_id, false);
}
ObjectInner::TlsExternalReport(report) => {
report.write_ops(&mut batch, item_id, false);
}
ObjectInner::ArfExternalReport(report) => {
report.write_ops(&mut batch, item_id, false);
}
_ => {}
}
if batch.is_large_batch() {
server.store().write(batch.build_all()).await?;
batch = BatchBuilder::new();
}
}
}
}
if !batch.is_empty() {
server.store().write(batch.build_all()).await?;
}
let started = Instant::now();
server
.store()
.purge_store()
.await
.caused_by(trc::location!())?;
server
.in_memory_store()
.purge_in_memory_store()
.await
.caused_by(trc::location!())?;
server
.registry()
.purge_dead_nodes()
.await
.caused_by(trc::location!())?;
// inbuxa: UD-13: archived items past their deadline go
inbuxa_features::undelete::records::remove_expired(
&server.core.storage.data,
server.registry(),
)
.await
.caused_by(trc::location!())?;
use common::telemetry::metrics::store::MetricsStore;
// inbuxa: MON-17, MON-38: history past its retention goes; a
// failure leaves it for the next run
let retention = common::telemetry::metrics::store::retention(server).await;
if let Some(keep) = retention.hold_metrics_for
&& !server.metrics_store().is_none()
&& let Err(err) = server
.metrics_store()
.purge_metrics(keep.into_inner())
.await
{
trc::error!(err.details("Failed to purge metric history"));
}
if let Some(keep) = retention.hold_traces_for
&& !server.tracing_store().is_none()
{
use common::telemetry::tracers::store::TracingStore;
if let Err(err) = server
.tracing_store()
.purge_spans(keep.into_inner(), Some(server.search_store()))
.await
{
trc::error!(err.details("Failed to purge trace history"));
}
}
// inbuxa: AL-7: locks' grants reach folders the server made on
// its own (a Sieve fileinto :create)
if let Err(err) = email::inbuxa_lock::reconcile_all(server).await {
trc::error!(err.details("Failed to re-apply account locks"));
}
// inbuxa: personal-data catalog: the inventory's daily look for
// a change, and snapshots past the audit log's retention go
if let Err(err) = server
.inventory_snapshot(inbuxa_features::privacy::snapshot::Trigger::Daily)
.await
{
trc::error!(err.details("Failed to record an inventory snapshot"));
}
if let Err(err) = server.purge_inventory_snapshots().await {
trc::error!(err.details("Failed to purge inventory snapshots"));
}
// inbuxa: personal-data catalog, D2: bans past their period go
if let Err(err) = server.purge_expired_blocked_ips().await {
trc::error!(err.details("Failed to purge expired IP bans"));
}
// inbuxa: DLP, §2.6: mail nobody reviewed in time goes back
if let Err(err) = smtp::queue::held::expire(server).await {
trc::error!(err.details("Failed to return unreviewed held mail"));
}
// inbuxa: journaling, JR-13: entries past their retention go,
// except those a legal hold keeps
if let Err(err) = purge_journal(server).await {
trc::error!(err.details("Failed to purge journal entries"));
}
// inbuxa: AU-7: audit records past their retention go; a
// failure leaves them for the next run
if let Err(err) = server.audit_purge().await {
trc::error!(err.details("Failed to purge audit records"));
}
trc::event!(
Store(StoreEvent::DataStorePurged),
Elapsed = started.elapsed()
);
}
TaskStoreMaintenanceType::PurgeBlob => {
if let Some(shard_index) = task.shard_index {
server
.store()
.purge_blobs(server.blob_store().clone(), shard_index as u8)
.await
.caused_by(trc::location!())?;
} else {
let mut batch = BatchBuilder::new();
let now = now() as i64;
for shard_index in 0..=u8::MAX {
batch.schedule_task(Task::StoreMaintenance(TaskStoreMaintenance {
maintenance_type: TaskStoreMaintenanceType::PurgeBlob,
shard_index: Some(shard_index as u64),
status: TaskStatus::at(now),
}));
if batch.is_large_batch() {
server.core.storage.data.write(batch.build_all()).await?;
server.notify_task_queue();
batch = BatchBuilder::new();
}
}
if !batch.is_empty() {
server.core.storage.data.write(batch.build_all()).await?;
server.notify_task_queue();
}
}
}
TaskStoreMaintenanceType::RemoveGreylist
| TaskStoreMaintenanceType::RemoveLockQueueMessage
| TaskStoreMaintenanceType::RemoveLockTask
| TaskStoreMaintenanceType::RemoveLockDav
| TaskStoreMaintenanceType::RemoveSieveId
| TaskStoreMaintenanceType::ResetRateLimiters
| TaskStoreMaintenanceType::ResetBlobQuotas
| TaskStoreMaintenanceType::RemoveAuthTokens => {
#[cfg(feature = "test_mode")]
if let Some(test_var) = task.shard_index {
use crate::task_manager::TaskFailureType;
// Simulate success for testing purposes
match test_var {
0 => {
return Ok(TaskResult::Success(vec![]));
}
1 => {
return Ok(TaskResult::temporary(
"Simulated temporary failure".to_string(),
));
}
2 => {
return Ok(TaskResult::permanent("Simulated permanent failure"));
}
retry => {
return Ok(TaskResult::Failure {
typ: TaskFailureType::Retry(retry),
message: "Simulated retry failure".to_string(),
max_attempts: None,
});
}
}
}
let prefixes = match task.maintenance_type {
TaskStoreMaintenanceType::RemoveGreylist => &[KV_GREYLIST][..],
TaskStoreMaintenanceType::RemoveLockQueueMessage => &[KV_LOCK_QUEUE_MESSAGE][..],
TaskStoreMaintenanceType::RemoveLockTask => &[KV_LOCK_TASK][..],
TaskStoreMaintenanceType::RemoveLockDav => &[KV_LOCK_DAV][..],
TaskStoreMaintenanceType::RemoveSieveId => &[KV_SIEVE_ID][..],
TaskStoreMaintenanceType::ResetRateLimiters => &[
KV_RATE_LIMIT_RCPT,
KV_RATE_LIMIT_SCAN,
KV_RATE_LIMIT_LOITER,
KV_RATE_LIMIT_AUTH,
KV_RATE_LIMIT_SMTP,
KV_RATE_LIMIT_CONTACT,
KV_RATE_LIMIT_HTTP_AUTHENTICATED,
KV_RATE_LIMIT_HTTP_ANONYMOUS,
KV_RATE_LIMIT_IMAP,
][..],
TaskStoreMaintenanceType::ResetBlobQuotas => &[KV_QUOTA_BLOB][..],
TaskStoreMaintenanceType::RemoveAuthTokens => &[KV_ACME, KV_OAUTH][..],
_ => unreachable!(),
};
for &prefix in prefixes {
server
.in_memory_store()
.key_delete_prefix(&[prefix])
.await?;
}
}
// inbuxa: MT-21: recompute every tenant's usage
TaskStoreMaintenanceType::ResetTenantQuotas => {
for tenant_id in inbuxa_features::tenancy::quota::all_tenants(server.registry()).await?
{
recalculate_tenant_quota(server, tenant_id).await?;
}
}
}
Ok(TaskResult::Success(vec![]))
}
/// inbuxa: journaling, JR-13: removes journal entries past their
/// retention, keeping any whose sender or recipients a legal hold covers
/// (deleted accounts a hold keeps included), and records how many went.
async fn purge_journal(server: &Server) -> trc::Result<()> {
use inbuxa_features::audit::{Action, Actor, Outcome, Record, Target};
let mut held = server.held_accounts().await?;
if !held.is_empty() {
for (account_id, kept) in
inbuxa_features::undelete::data::kept_accounts(server.store()).await?
{
if server.is_kept_held(account_id, &kept).await? {
held.insert(account_id);
}
}
}
let at = store::write::now();
let purged = inbuxa_features::journal::entries::purge(server.store(), at, |entry| {
entry.accounts.iter().any(|account| held.contains(account))
})
.await?;
if purged.removed > 0 || purged.kept_for_hold > 0 {
server
.audit_note(Record {
at: at * 1000,
actor: Actor::system("Journal"),
via: None,
remote_ip: None,
action: Action::Destroy,
target: Target {
kind: "inbuxa:JournalEntry".into(),
id: None,
name: None,
account_id: None,
tenant_id: None,
},
changes: vec![],
details: Some(format!(
"{} past their retention removed; {} kept for a legal hold",
purged.removed, purged.kept_for_hold
)),
reason: None,
outcome: Outcome::success(),
})
.await;
}
Ok(())
}
async fn account_maintenance(
server: &Server,
task: &TaskAccountMaintenance,
) -> trc::Result<TaskResult> {
match task.maintenance_type {
TaskAccountMaintenanceType::Purge => {
server.purge_account(task.account_id.document_id()).await?;
}
TaskAccountMaintenanceType::Reindex => {
reindex_account(server, task.account_id.document_id()).await?;
}
TaskAccountMaintenanceType::RecalculateImapUid => {
reset_imap_uids(server, task.account_id.document_id()).await?;
}
TaskAccountMaintenanceType::RecalculateQuota => {
recalculate_quota(server, task.account_id.document_id()).await?;
}
}
Ok(TaskResult::Success(vec![]))
}
async fn tenant_maintenance(
server: &Server,
task: &TaskTenantMaintenance,
) -> trc::Result<TaskResult> {
match task.maintenance_type {
TaskTenantMaintenanceType::RecalculateQuota => {
recalculate_tenant_quota(server, task.tenant_id.document_id()).await?;
}
}
Ok(TaskResult::Success(vec![]))
}
async fn recalculate_quota(server: &Server, account_id: u32) -> trc::Result<()> {
let mut quota = 0;
for collection in [
Collection::Email,
Collection::Calendar,
Collection::CalendarEvent,
Collection::CalendarEventNotification,
Collection::AddressBook,
Collection::ContactCard,
Collection::FileNode,
Collection::SieveScript,
] {
server
.archives(account_id, collection, &(), |_, archive| {
match collection {
Collection::Email => {
quota += archive.unarchive::<MessageData>()?.size.to_native() as i64;
}
Collection::Calendar => {
quota += archive.unarchive::<Calendar>()?.size() as i64;
}
Collection::CalendarEvent => {
quota += archive.unarchive::<CalendarEvent>()?.size() as i64;
}
Collection::CalendarEventNotification => {
quota += archive.unarchive::<CalendarEventNotification>()?.size() as i64;
}
Collection::AddressBook => {
quota += archive.unarchive::<AddressBook>()?.size() as i64;
}
Collection::ContactCard => {
quota += archive.unarchive::<ContactCard>()?.size() as i64;
}
Collection::FileNode => {
quota += archive.unarchive::<FileNode>()?.size() as i64;
}
Collection::SieveScript => {
quota += u32::from(archive.unarchive::<SieveScript>()?.size) as i64;
}
_ => {}
}
Ok(true)
})
.await
.caused_by(trc::location!())?;
}
let mut batch = BatchBuilder::new();
batch
.with_account_id(account_id)
.clear(ValueClass::Quota)
.add(ValueClass::Quota, quota);
server
.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())
.map(|_| ())
}
// inbuxa: MT-21: recompute a tenant's usage from its members'
async fn recalculate_tenant_quota(server: &Server, tenant_id: u32) -> trc::Result<()> {
inbuxa_features::tenancy::quota::recalculate(
&server.core.storage.data,
server.registry(),
tenant_id,
)
.await
.map(|_| ())
}
async fn reset_imap_uids(server: &Server, account_id: u32) -> trc::Result<(u32, u32)> {
let mut mailbox_count = 0;
let mut email_count = 0;
let cache = server
.get_cached_messages(account_id)
.await
.caused_by(trc::location!())?;
for &mailbox_id in cache.mailboxes.index.keys() {
let mailbox = server
.store()
.get_value::<Archive<AlignedBytes>>(ValueKey::archive(
account_id,
Collection::Mailbox,
mailbox_id,
))
.await
.caused_by(trc::location!())?
.ok_or_else(|| trc::ImapEvent::Error.into_err().caused_by(trc::location!()))?
.into_deserialized::<email::mailbox::Mailbox>()
.caused_by(trc::location!())?;
let mut new_mailbox = mailbox.inner.clone();
new_mailbox.uid_validity = rand::random::<u32>();
let mut batch = BatchBuilder::new();
batch
.with_account_id(account_id)
.with_collection(Collection::Mailbox)
.with_document(mailbox_id)
.custom(
ObjectIndexBuilder::new()
.with_current(mailbox)
.with_changes(new_mailbox),
)
.caused_by(trc::location!())?
.clear(MailboxField::UidCounter);
server
.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
mailbox_count += 1;
}
// Reset all UIDs
for message_id in cache.emails.items.iter().map(|i| i.document_id) {
let data = server
.store()
.get_value::<Archive<AlignedBytes>>(ValueKey::archive(
account_id,
Collection::Email,
message_id,
))
.await
.caused_by(trc::location!())?;
let data_ = if let Some(data) = data {
data
} else {
continue;
};
let data = data_
.to_unarchived::<MessageData>()
.caused_by(trc::location!())?;
let mut new_data = data
.deserialize::<MessageData>()
.caused_by(trc::location!())?;
let ids = server
.assign_email_ids(
account_id,
new_data.mailboxes.iter().map(|m| m.mailbox_id),
false,
)
.await
.caused_by(trc::location!())?;
for (uid_mailbox, uid) in new_data.mailboxes.iter_mut().zip(ids) {
uid_mailbox.uid = uid;
}
// Prepare write batch
let mut batch = BatchBuilder::new();
batch
.with_account_id(account_id)
.with_collection(Collection::Email)
.with_document(message_id)
.assert_value(ValueClass::Property(EmailField::Archive.into()), &data)
.set(
EmailField::Archive,
Archiver::new(new_data)
.serialize()
.caused_by(trc::location!())?,
);
server
.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
email_count += 1;
}
Ok((mailbox_count, email_count))
}