diff --git a/crates/common/src/auth/permissions.rs b/crates/common/src/auth/permissions.rs index 04a98be..a111b77 100644 --- a/crates/common/src/auth/permissions.rs +++ b/crates/common/src/auth/permissions.rs @@ -165,6 +165,14 @@ impl AccessToken { mut requested_permissions: Permissions, ) -> Result<(), Vec> { requested_permissions.difference(self.permissions_bits()); + // inbuxa: journaling, JR-18: whoever sets up journals may give + // others (or, through a role, themselves) the reading of them, + // which administrators don't hold by default; the role change is + // in the audit log + if self.has_permission(Permission::SysJournalUpdate) { + requested_permissions.clear(Permission::SysJournalSearch as usize); + requested_permissions.clear(Permission::SysJournalExport as usize); + } if requested_permissions.is_empty() { Ok(()) } else { @@ -307,6 +315,13 @@ impl Default for DefaultPermissions { | Permission::SysDlpReviewUpdate => { default.superuser.push(permission); } + // inbuxa: journals are the server's; administrators set them + // up but read what's journaled only if granted it + // (journaling spec, JR-18, settled answer 5) + Permission::SysJournalGet | Permission::SysJournalUpdate => { + default.superuser.push(permission); + } + Permission::SysJournalSearch | Permission::SysJournalExport => {} // inbuxa: AL-12: tenant administrators lock and delegate // within their tenant Permission::SysAccountLockGet diff --git a/crates/common/src/manager/compliance_roles.rs b/crates/common/src/manager/compliance_roles.rs index 3b53f5f..21f65a1 100644 --- a/crates/common/src/manager/compliance_roles.rs +++ b/crates/common/src/manager/compliance_roles.rs @@ -69,6 +69,10 @@ const OFFICER: &[Permission] = &[ Permission::SysDlpPolicyGet, Permission::SysDlpReviewGet, Permission::SysDlpReviewUpdate, + // journaling spec, JR-18: see journals, search and export them + Permission::SysJournalGet, + Permission::SysJournalSearch, + Permission::SysJournalExport, ]; /// What a tenant's officer holds besides [`READS`]. diff --git a/crates/common/src/manager/granted_permissions.rs b/crates/common/src/manager/granted_permissions.rs index fac45e4..7d653a9 100644 --- a/crates/common/src/manager/granted_permissions.rs +++ b/crates/common/src/manager/granted_permissions.rs @@ -52,6 +52,8 @@ const ADMIN_GRANTS: &[Permission] = &[ Permission::SysDlpPolicyUpdate, Permission::SysDlpReviewGet, Permission::SysDlpReviewUpdate, + Permission::SysJournalGet, + Permission::SysJournalUpdate, ]; /// Granted to the server-level Compliance Officer role once it exists: @@ -61,6 +63,10 @@ const OFFICER_GRANTS: &[Permission] = &[ Permission::SysDlpPolicyGet, Permission::SysDlpReviewGet, Permission::SysDlpReviewUpdate, + // journaling spec, JR-18: see journals, search and export them + Permission::SysJournalGet, + Permission::SysJournalSearch, + Permission::SysJournalExport, ]; /// Granted to the default tenant administrator roles: reading and exporting diff --git a/crates/features/src/journal/entries.rs b/crates/features/src/journal/entries.rs new file mode 100644 index 0000000..ebdf3c9 --- /dev/null +++ b/crates/features/src/journal/entries.rs @@ -0,0 +1,689 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! The built-in journal (JR-5, JR-6, JR-13). Keys, after `J`: +//! +//! - `e` + node + seq: a chain link: its seq, the hash of the link before +//! it, and the SHA-256 of its entry. One chain per node, as the audit log +//! keeps (AU-6), but a link names its entry by hash instead of holding it, +//! so an entry can go at the end of its own retention without breaking +//! the chain: entries don't expire in chain order. +//! - `c` + node + seq: the entry, as JSON; its bytes are what the link's +//! hash names. +//! - `p` + node + seq: when an entry past its retention was purged. A link +//! whose entry is gone without this marker is a broken chain. +//! - `t` + time + node + seq: the time index, for search. +//! - `x` + expiry + node + seq: the expiry index, for purge. +//! - `h` + node: the chain's head: its hash, then its seq as the last eight +//! bytes, which each append asserts. +//! - `f` + node: where the chain starts after purged links at its start +//! were cleared, and the hash the first kept link names. +//! +//! The report itself is a blob, kept by a temporary link that lasts until +//! its entry is purged. Nothing here changes or removes an entry before +//! its time; nothing in JMAP can. + +use super::{Direction, FEATURE, Json}; +use crate::hold::HELD_UNTIL; +use serde::{Deserialize as SerdeDeserialize, Serialize as SerdeSerialize}; +use sha2::{Digest, Sha256}; +use std::fmt; +use store::{ + BlobStore, Deserialize, IterateParams, SUBSPACE_INBUXA, Serialize, Store, ValueKey, + write::{AnyClass, BatchBuilder, BlobLink, BlobOp, ValueClass, assert::AssertValue}, +}; +use tokio::sync::Mutex; +use trc::AddContext; +use types::blob_hash::BlobHash; + +const KIND_LINK: u8 = b'e'; +const KIND_CONTENT: u8 = b'c'; +const KIND_PURGED: u8 = b'p'; +const KIND_TIME: u8 = b't'; +const KIND_EXPIRY: u8 = b'x'; +const KIND_HEAD: u8 = b'h'; +const KIND_FLOOR: u8 = b'f'; + +const APPEND_ATTEMPTS: usize = 5; +/// Entries purged per batch. +const PURGE_BATCH: usize = 100; + +/// Where one entry sits: its node's chain and its place in it. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct EntryId { + pub node: u64, + pub seq: u64, +} + +impl EntryId { + /// As one number, for JMAP ids: the node in the top 16 bits. + pub fn to_u64(&self) -> u64 { + (self.node << 48) | (self.seq & ((1 << 48) - 1)) + } + + pub fn from_u64(id: u64) -> Self { + EntryId { + node: id >> 48, + seq: id & ((1 << 48) - 1), + } + } +} + +impl fmt::Display for EntryId { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}-{}", self.node, self.seq) + } +} + +/// One journaled message (JR-5). +#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)] +#[serde(rename_all = "camelCase")] +pub struct Entry { + pub queue_id: u64, + /// Seconds. + pub at: u64, + pub direction: Direction, + pub sender: String, + pub authenticated: bool, + pub recipients: Vec, + pub subject: String, + pub message_id: String, + /// The people here on either side, whose holds keep the entry. + pub accounts: Vec, + pub tenants: Vec, + /// The journals that took it. + pub journals: Vec, + pub held: bool, + /// The report's blob, hex. + pub blob: String, + pub size: u64, + /// SHA-256 of the report, hex. + pub sha256: String, + /// Seconds. + pub expires_at: u64, +} + +impl Entry { + pub fn blob_hash(&self) -> Option { + let bytes = unhex(&self.blob)?; + BlobHash::try_from_hash_slice(&bytes).ok() + } +} + +#[derive(Debug, Clone, SerdeSerialize, SerdeDeserialize)] +#[serde(rename_all = "camelCase")] +struct Link { + seq: u64, + prev: String, + content: String, +} + +#[derive(Debug, Clone, Default, PartialEq, SerdeSerialize, SerdeDeserialize)] +struct Floor { + seq: u64, + prev: String, +} + +#[derive(Debug, Clone, Default, PartialEq)] +struct Head { + seq: u64, + hash: String, +} + +impl Head { + fn to_bytes(&self) -> Vec { + let mut bytes = self.hash.as_bytes().to_vec(); + bytes.extend_from_slice(&self.seq.to_be_bytes()); + bytes + } +} + +impl Deserialize for Head { + fn deserialize(bytes: &[u8]) -> trc::Result { + let split = bytes.len().checked_sub(8).ok_or_else(|| { + trc::StoreEvent::DataCorruption + .into_err() + .details("Invalid journal chain head") + })?; + Ok(Head { + seq: u64::from_be_bytes(bytes[split..].try_into().unwrap()), + hash: String::from_utf8_lossy(&bytes[..split]).into_owned(), + }) + } +} + +struct Raw(Vec); + +impl Deserialize for Raw { + fn deserialize(bytes: &[u8]) -> trc::Result { + Ok(Raw(bytes.to_vec())) + } +} + +fn class(kind: u8, parts: &[u64]) -> ValueClass { + let mut key = Vec::with_capacity(2 + parts.len() * 8); + key.push(FEATURE); + key.push(kind); + for part in parts { + key.extend_from_slice(&part.to_be_bytes()); + } + ValueClass::Any(AnyClass { + subspace: SUBSPACE_INBUXA, + key, + }) +} + +fn key(kind: u8, parts: &[u64]) -> ValueKey { + ValueKey::from(class(kind, parts)) +} + +/// Where an entry's content is kept, for tests that check tampering shows. +pub fn content_key(id: EntryId) -> ValueKey { + key(KIND_CONTENT, &[id.node, id.seq]) +} + +/// The numbers after the kind byte, from the key's tail. +fn parse_key(key: &[u8], kind: u8, parts: usize) -> Option> { + let len = 2 + parts * 8; + let tail = key.get(key.len().checked_sub(len)?..)?; + (tail[0] == FEATURE && tail[1] == kind).then_some(())?; + Some( + tail[2..] + .chunks_exact(8) + .map(|chunk| u64::from_be_bytes(chunk.try_into().unwrap())) + .collect(), + ) +} + +pub fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|b| format!("{b:02x}")).collect() +} + +fn unhex(value: &str) -> Option> { + (value.len() % 2 == 0).then_some(())?; + (0..value.len()) + .step_by(2) + .map(|i| u8::from_str_radix(value.get(i..i + 2)?, 16).ok()) + .collect() +} + +pub fn sha256(bytes: &[u8]) -> String { + hex(&Sha256::digest(bytes)) +} + +async fn head(data: &Store, node: u64) -> trc::Result> { + data.get_value::(key(KIND_HEAD, &[node])) + .await + .caused_by(trc::location!()) +} + +async fn floor(data: &Store, node: u64) -> trc::Result { + Ok(data + .get_value::>(key(KIND_FLOOR, &[node])) + .await + .caused_by(trc::location!())? + .map(|Json(floor)| floor) + .unwrap_or(Floor { + seq: 1, + prev: String::new(), + })) +} + +async fn nodes(data: &Store) -> trc::Result> { + let mut nodes = Vec::new(); + data.iterate( + IterateParams::new(key(KIND_HEAD, &[0]), key(KIND_HEAD, &[u64::MAX])).no_values(), + |key, _| { + if let Some(parts) = parse_key(key, KIND_HEAD, 1) { + nodes.push(parts[0]); + } + Ok(true) + }, + ) + .await + .caused_by(trc::location!())?; + Ok(nodes) +} + +/// Lines up this process's appends; the store's assert settles the rest. +static APPENDING: Mutex<()> = Mutex::const_new(()); + +/// Adds an entry to this node's chain, and links its report's blob (already +/// written) until the entry is purged. An error means nothing was written. +pub async fn append(data: &Store, node: u64, entry: &Entry) -> trc::Result { + let blob = entry.blob_hash().ok_or_else(|| { + trc::StoreEvent::UnexpectedError + .into_err() + .details("Journal entry without a blob") + })?; + let content = Json(entry).serialize()?; + let content_hash = sha256(&content); + let _appending = APPENDING.lock().await; + let mut attempt = 0; + loop { + attempt += 1; + let current = head(data, node).await?; + let (seq, prev) = current + .as_ref() + .map_or((1, String::new()), |head| (head.seq + 1, head.hash.clone())); + let link = Json(&Link { + seq, + prev, + content: content_hash.clone(), + }) + .serialize()?; + let new_head = Head { + seq, + hash: sha256(&link), + }; + + let mut batch = BatchBuilder::new(); + batch.assert_value( + class(KIND_HEAD, &[node]), + current.map_or(AssertValue::None, |head| AssertValue::U64(head.seq)), + ); + batch + .set(class(KIND_LINK, &[node, seq]), link) + .set(class(KIND_CONTENT, &[node, seq]), content.clone()) + .set(class(KIND_TIME, &[entry.at, node, seq]), vec![]) + .set(class(KIND_EXPIRY, &[entry.expires_at, node, seq]), vec![]) + .set(class(KIND_HEAD, &[node]), new_head.to_bytes()) + .set( + BlobOp::Link { + hash: blob.clone(), + to: BlobLink::Temporary { until: HELD_UNTIL }, + }, + vec![], + ) + .set(BlobOp::Commit { hash: blob.clone() }, vec![]); + match data.write(batch.build_all()).await { + Ok(_) => return Ok(EntryId { node, seq }), + Err(err) + if attempt < APPEND_ATTEMPTS + && matches!( + err.as_ref(), + trc::EventType::Store(trc::StoreEvent::AssertValueFailed) + ) => {} + Err(err) => return Err(err.caused_by(trc::location!())), + } + } +} + +/// One entry, unless it was purged. +pub async fn get(data: &Store, id: EntryId) -> trc::Result> { + Ok(data + .get_value::>(key(KIND_CONTENT, &[id.node, id.seq])) + .await + .caused_by(trc::location!())? + .map(|Json(entry)| entry)) +} + +/// Entries written in `[after, before)` (seconds), newest first, up to +/// `limit`. +pub async fn list( + data: &Store, + after: u64, + before: u64, + limit: usize, +) -> trc::Result> { + let mut ids = Vec::new(); + data.iterate( + IterateParams::new( + key(KIND_TIME, &[after, 0, 0]), + key(KIND_TIME, &[before.saturating_sub(1), u64::MAX, u64::MAX]), + ) + .descending() + .no_values(), + |key, _| { + if let Some(parts) = parse_key(key, KIND_TIME, 3) { + ids.push(EntryId { + node: parts[1], + seq: parts[2], + }); + } + Ok(ids.len() < limit) + }, + ) + .await + .caused_by(trc::location!())?; + let mut out = Vec::with_capacity(ids.len()); + for id in ids { + if let Some(entry) = get(data, id).await? { + out.push((id, entry)); + } + } + Ok(out) +} + +/// What a purge did. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct Purged { + pub removed: usize, + /// Past their time, kept for a legal hold. + pub kept_for_hold: usize, +} + +/// Removes entries past their retention (JR-13), except those `held` keeps: +/// the entry, its indexes and its blob's link go; the chain link stays, +/// with a purge marker. Then each chain's start moves past purged links. +pub async fn purge( + data: &Store, + now: u64, + held: impl Fn(&Entry) -> bool + Sync + Send, +) -> trc::Result { + let mut due = Vec::new(); + data.iterate( + IterateParams::new( + key(KIND_EXPIRY, &[0, 0, 0]), + key(KIND_EXPIRY, &[now, u64::MAX, u64::MAX]), + ) + .ascending() + .no_values(), + |key, _| { + if let Some(parts) = parse_key(key, KIND_EXPIRY, 3) { + due.push(( + parts[0], + EntryId { + node: parts[1], + seq: parts[2], + }, + )); + } + Ok(due.len() < 100_000) + }, + ) + .await + .caused_by(trc::location!())?; + + let mut purged = Purged::default(); + for chunk in due.chunks(PURGE_BATCH) { + let mut batch = BatchBuilder::new(); + for (expires_at, id) in chunk { + let parts = [id.node, id.seq]; + let Some(entry) = get(data, *id).await? else { + // Its entry is already gone: only the index is left + batch.clear(class(KIND_EXPIRY, &[*expires_at, id.node, id.seq])); + continue; + }; + if held(&entry) { + purged.kept_for_hold += 1; + continue; + } + batch + .clear(class(KIND_CONTENT, &parts)) + .clear(class(KIND_TIME, &[entry.at, id.node, id.seq])) + .clear(class(KIND_EXPIRY, &[*expires_at, id.node, id.seq])) + .set(class(KIND_PURGED, &parts), now.to_be_bytes().to_vec()); + if let Some(blob) = entry.blob_hash() { + batch.clear(BlobOp::Link { + hash: blob, + to: BlobLink::Temporary { until: HELD_UNTIL }, + }); + } + purged.removed += 1; + } + if !batch.is_empty() { + data.write(batch.build_all()) + .await + .caused_by(trc::location!())?; + } + } + + for node in nodes(data).await? { + advance_floor(data, node).await?; + } + Ok(purged) +} + +/// Clears the purged links at the start of a node's chain, recording where +/// it now starts and the hash that start names. +async fn advance_floor(data: &Store, node: u64) -> trc::Result<()> { + let start = floor(data, node).await?; + let mut cleared: Vec = Vec::new(); + let mut next = start.clone(); + let mut purged_seqs = Vec::new(); + data.iterate( + IterateParams::new( + key(KIND_PURGED, &[node, start.seq]), + key(KIND_PURGED, &[node, u64::MAX]), + ) + .ascending() + .no_values(), + |key, _| { + if let Some(parts) = parse_key(key, KIND_PURGED, 2) { + purged_seqs.push(parts[1]); + } + Ok(purged_seqs.len() < 100_000) + }, + ) + .await + .caused_by(trc::location!())?; + for seq in purged_seqs { + if seq != next.seq { + break; + } + let Some(Raw(link)) = data + .get_value::(key(KIND_LINK, &[node, seq])) + .await + .caused_by(trc::location!())? + else { + break; + }; + next = Floor { + seq: seq + 1, + prev: sha256(&link), + }; + cleared.push(seq); + } + if cleared.is_empty() { + return Ok(()); + } + // The floor moves first: a run cut short leaves links before it, which + // the next run clears, never a chain that looks broken + let mut batch = BatchBuilder::new(); + batch.set(class(KIND_FLOOR, &[node]), Json(&next).serialize()?); + data.write(batch.build_all()) + .await + .caused_by(trc::location!())?; + for chunk in cleared.chunks(PURGE_BATCH) { + let mut batch = BatchBuilder::new(); + for seq in chunk { + batch + .clear(class(KIND_LINK, &[node, *seq])) + .clear(class(KIND_PURGED, &[node, *seq])); + } + data.write(batch.build_all()) + .await + .caused_by(trc::location!())?; + } + Ok(()) +} + +/// One node's chain, as [`verify`] found it. +#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize)] +#[serde(rename_all = "camelCase")] +pub struct ChainReport { + pub node: u64, + pub entries: u64, + pub purged: u64, + pub first_seq: u64, + pub last_seq: u64, + #[serde(skip_serializing_if = "Option::is_none")] + pub broken_at: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub reason: Option, +} + +/// Rechecks every node's chain (JR-6): each link names the hash of the one +/// before it, seqs run without gaps, the head matches the last link, each +/// entry hashes to what its link names or was purged, and, with `blobs`, +/// each report is there and hashes to what its entry names. +pub async fn verify(data: &Store, blobs: Option<&BlobStore>) -> trc::Result> { + let mut reports = Vec::new(); + for node in nodes(data).await? { + let start = floor(data, node).await?; + let head = head(data, node).await?.unwrap_or_default(); + let mut report = ChainReport { + node, + entries: 0, + purged: 0, + first_seq: start.seq, + last_seq: start.seq.saturating_sub(1), + broken_at: None, + reason: None, + }; + let mut links = Vec::new(); + data.iterate( + IterateParams::new( + key(KIND_LINK, &[node, start.seq]), + key(KIND_LINK, &[node, u64::MAX]), + ) + .ascending(), + |key, value| { + if let Some(parts) = parse_key(key, KIND_LINK, 2) { + links.push((parts[1], value.to_vec())); + } + Ok(true) + }, + ) + .await + .caused_by(trc::location!())?; + + let mut expected_seq = start.seq; + let mut expected_prev = start.prev.clone(); + for (seq, bytes) in links { + let broken = |report: &mut ChainReport, reason: &str| { + report.broken_at = Some(EntryId { node, seq }.to_string()); + report.reason = Some(reason.to_string()); + }; + let Ok(Json(link)) = Json::::deserialize(&bytes) else { + broken(&mut report, "The link can't be read."); + break; + }; + if seq != expected_seq || link.seq != seq { + report.broken_at = Some(EntryId { node, seq }.to_string()); + report.reason = Some(format!( + "Entry {expected_seq} is missing; the next one found is {seq}." + )); + break; + } + if link.prev != expected_prev { + broken( + &mut report, + "The link doesn't follow from the one before it: one of them was changed.", + ); + break; + } + match data + .get_value::(key(KIND_CONTENT, &[node, seq])) + .await + .caused_by(trc::location!())? + { + Some(Raw(content)) => { + if sha256(&content) != link.content { + broken(&mut report, "The entry was changed after it was written."); + break; + } + if let Some(blobs) = blobs { + let Ok(Json(entry)) = Json::::deserialize(&content) else { + broken(&mut report, "The entry can't be read."); + break; + }; + let report_bytes = match entry.blob_hash() { + Some(hash) => blobs + .get_blob(hash.as_slice(), 0..usize::MAX) + .await + .caused_by(trc::location!())?, + None => None, + }; + match report_bytes { + Some(bytes) if sha256(&bytes) == entry.sha256 => {} + Some(_) => { + broken(&mut report, "The report doesn't match its entry."); + break; + } + None => { + broken(&mut report, "The report is missing."); + break; + } + } + } + report.entries += 1; + } + None => { + if data + .get_value::(key(KIND_PURGED, &[node, seq])) + .await + .caused_by(trc::location!())? + .is_none() + { + broken(&mut report, "The entry was removed before its time."); + break; + } + report.purged += 1; + } + } + expected_prev = sha256(&bytes); + expected_seq = seq + 1; + report.last_seq = seq; + } + + if report.broken_at.is_none() + && (head.seq != report.last_seq + || (report.last_seq >= report.first_seq && head.hash != expected_prev)) + { + report.broken_at = Some( + EntryId { + node, + seq: report.last_seq, + } + .to_string(), + ); + report.reason = Some( + "The chain's recorded end doesn't match its last link: entries were removed \ + or changed at the end." + .into(), + ); + } + reports.push(report); + } + Ok(reports) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn keys_read_back() { + let ValueClass::Any(any) = class(KIND_EXPIRY, &[5, 3, 9]) else { + panic!() + }; + assert_eq!(parse_key(&any.key, KIND_EXPIRY, 3), Some(vec![5, 3, 9])); + let mut with_subspace = vec![SUBSPACE_INBUXA]; + with_subspace.extend_from_slice(&any.key); + assert_eq!( + parse_key(&with_subspace, KIND_EXPIRY, 3), + Some(vec![5, 3, 9]) + ); + assert_eq!(parse_key(&any.key, KIND_TIME, 3), None); + } + + #[test] + fn hex_round_trips() { + let bytes = [0u8, 1, 0xab, 0xff]; + assert_eq!(unhex(&hex(&bytes)), Some(bytes.to_vec())); + assert_eq!(unhex("abc"), None); + assert_eq!(unhex("zz"), None); + } + + #[test] + fn ids_read_back() { + let id = EntryId { node: 3, seq: 77 }; + assert_eq!(EntryId::from_u64(id.to_u64()), id); + assert_eq!(id.to_string(), "3-77"); + } +} diff --git a/crates/features/src/journal/mod.rs b/crates/features/src/journal/mod.rs new file mode 100644 index 0000000..f0c2e35 --- /dev/null +++ b/crates/features/src/journal/mod.rs @@ -0,0 +1,431 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Journaling (journaling spec, JR-1 to JR-18): a copy of each message the +//! server queues, with its envelope, kept where nothing in the product +//! changes or removes it before its retention ends. +//! +//! - this module: journals, what makes one valid, and where they're kept; +//! - [`report`]: the journal report around the untouched message (JR-3); +//! - [`entries`]: the built-in journal and its chain (JR-5, JR-6, JR-13). +//! +//! Kept in the fork's subspace (`store::SUBSPACE_INBUXA`). Every key starts +//! with `J`; journals are `j` + id (u32), as JSON. There are few, so they're +//! read whole. + +pub mod entries; +pub mod report; + +use crate::{hold::Member, mailflow::rules::jmap_ids}; +use serde::{Deserialize as SerdeDeserialize, Serialize as SerdeSerialize, de::DeserializeOwned}; +use std::{ + sync::{Arc, RwLock}, + time::{Duration, Instant}, +}; +use store::{ + Deserialize, IterateParams, SUBSPACE_INBUXA, Serialize, Store, ValueKey, + write::{AnyClass, BatchBuilder, ValueClass, assert::AssertValue}, +}; +use trc::AddContext; + +pub(crate) const FEATURE: u8 = b'J'; +const KIND_JOURNAL: u8 = b'j'; +const CREATE_ATTEMPTS: usize = 5; + +/// Retention a journal may be given, in days (settled answer 3). +pub const MIN_RETENTION_DAYS: u32 = 30; +pub const MAX_RETENTION_DAYS: u32 = 3650; +/// Most entries in one scope list. +const MAX_LIST: usize = 5_000; + +/// Which way a message goes, from this server's side (JR-9). +#[derive(Debug, Clone, Copy, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)] +#[serde(rename_all = "camelCase")] +pub enum Direction { + /// From someone here to at least one recipient elsewhere. + Outgoing, + /// From elsewhere to someone here. + Incoming, + /// From someone here, to people here only. + Internal, + Any, +} + +impl Direction { + pub fn as_str(&self) -> &'static str { + match self { + Direction::Outgoing => "outgoing", + Direction::Incoming => "incoming", + Direction::Internal => "internal", + Direction::Any => "any", + } + } + + /// A message's direction: `Any` is never one. + pub fn of(sender_local: bool, any_remote: bool, any_local: bool) -> Direction { + match (sender_local, any_remote) { + (true, true) => Direction::Outgoing, + (true, false) => Direction::Internal, + (false, _) if any_local => Direction::Incoming, + // Nobody here on either side: relayed mail counts as outgoing + (false, _) => Direction::Outgoing, + } + } + + fn includes(&self, direction: Direction) -> bool { + *self == Direction::Any || *self == direction + } +} + +/// Whose mail a journal takes (JR-9): everyone, or people reached through +/// their account, domain, group or tenant. Ids are in the JMAP form. +#[derive(Debug, Clone, Default, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)] +#[serde(rename_all = "camelCase")] +pub struct Scope { + #[serde(default)] + pub everyone: bool, + #[serde(default, with = "jmap_ids")] + pub accounts: Vec, + #[serde(default, with = "jmap_ids")] + pub groups: Vec, + #[serde(default, with = "jmap_ids")] + pub domains: Vec, + #[serde(default, with = "jmap_ids")] + pub tenants: Vec, +} + +impl Scope { + fn lists(&self) -> [&Vec; 4] { + [&self.accounts, &self.groups, &self.domains, &self.tenants] + } + + /// Whether this scope reaches one person here. + pub fn covers(&self, member: &Member) -> bool { + self.everyone + || self.accounts.contains(&member.account) + || member.domains.iter().any(|d| self.domains.contains(d)) + || member.groups.iter().any(|g| self.groups.contains(g)) + || member.tenant.is_some_and(|t| self.tenants.contains(&t)) + } +} + +/// A journal (JR-9): what it takes, and how long its entries are kept. +#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)] +#[serde(rename_all = "camelCase")] +pub struct Journal { + #[serde(default)] + pub id: u32, + pub name: String, + #[serde(default)] + pub description: String, + #[serde(default)] + pub enabled: bool, + pub direction: Direction, + pub scope: Scope, + /// How long an entry this journal writes is kept. An entry keeps the + /// retention it was written with (JR-12). + pub retention_days: u32, + #[serde(default)] + pub created_by: String, + #[serde(default)] + pub created_at: u64, + #[serde(default)] + pub updated_at: u64, +} + +/// Why a journal was refused: the property, and what to do. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Invalid { + pub property: &'static str, + pub reason: String, +} + +fn invalid(property: &'static str, reason: impl Into) -> Result<(), Invalid> { + Err(Invalid { + property, + reason: reason.into(), + }) +} + +impl Journal { + pub fn validate(&self) -> Result<(), Invalid> { + if self.name.trim().is_empty() { + return invalid("name", "Give the journal a name."); + } + if self.name.len() > 200 || self.description.len() > 2_000 { + return invalid("name", "The name or description is too long."); + } + if !(MIN_RETENTION_DAYS..=MAX_RETENTION_DAYS).contains(&self.retention_days) { + return invalid( + "retentionDays", + format!("Keep entries between {MIN_RETENTION_DAYS} and {MAX_RETENTION_DAYS} days."), + ); + } + let chosen = self.scope.lists().iter().any(|list| !list.is_empty()); + if self.scope.everyone == chosen { + return invalid( + "scope", + "Journal everyone, or choose accounts, groups, domains or tenants; not both.", + ); + } + if self.scope.lists().iter().any(|list| list.len() > MAX_LIST) { + return invalid("scope", format!("Choose at most {MAX_LIST} of each.")); + } + Ok(()) + } + + /// Whether this journal takes a message going `direction` with these + /// people here on either side. + pub fn takes(&self, direction: Direction, members: &[Member]) -> bool { + self.enabled + && self.direction.includes(direction) + && (self.scope.everyone || members.iter().any(|m| self.scope.covers(m))) + } +} + +/// A value stored as JSON. +pub(crate) struct Json(pub T); + +impl Serialize for Json { + fn serialize(&self) -> trc::Result> { + serde_json::to_vec(&self.0).map_err(|err| { + trc::StoreEvent::UnexpectedError + .into_err() + .details("Failed to serialize a journal record") + .reason(err) + }) + } +} + +impl Deserialize for Json { + fn deserialize(bytes: &[u8]) -> trc::Result { + serde_json::from_slice(bytes).map(Json).map_err(|err| { + trc::StoreEvent::DataCorruption + .into_err() + .details("Invalid journal record") + .reason(err) + }) + } +} + +fn class(id: u32) -> ValueClass { + let mut key = Vec::with_capacity(6); + key.push(FEATURE); + key.push(KIND_JOURNAL); + key.extend_from_slice(&id.to_be_bytes()); + ValueClass::Any(AnyClass { + subspace: SUBSPACE_INBUXA, + key, + }) +} + +fn key(id: u32) -> ValueKey { + ValueKey::from(class(id)) +} + +pub async fn get(data: &Store, id: u32) -> trc::Result> { + Ok(data + .get_value::>(key(id)) + .await + .caused_by(trc::location!())? + .map(|Json(journal)| journal)) +} + +/// Every journal, oldest first. +pub async fn all(data: &Store) -> trc::Result> { + let mut journals = Vec::new(); + data.iterate(IterateParams::new(key(0), key(u32::MAX)), |_, value| { + if let Ok(Json(journal)) = Json::::deserialize(value) { + journals.push(journal); + } + Ok(true) + }) + .await + .caused_by(trc::location!())?; + journals.sort_by_key(|journal| journal.id); + Ok(journals) +} + +/// Writes a new journal under the next free id, which it returns. +pub async fn create(data: &Store, journal: &Journal) -> trc::Result { + let mut attempt = 0; + loop { + attempt += 1; + let id = all(data).await?.iter().map(|j| j.id).max().unwrap_or(0) + 1; + let stored = Journal { + id, + ..journal.clone() + }; + let mut batch = BatchBuilder::new(); + batch.assert_value(class(id), AssertValue::None); + batch.set(class(id), Json(&stored).serialize()?); + match data.write(batch.build_all()).await { + Ok(_) => { + invalidate(); + return Ok(id); + } + Err(err) + if attempt < CREATE_ATTEMPTS + && matches!( + err.as_ref(), + trc::EventType::Store(trc::StoreEvent::AssertValueFailed) + ) => {} + Err(err) => return Err(err.caused_by(trc::location!())), + } + } +} + +/// Replaces a stored journal (same id). +pub async fn update(data: &Store, journal: &Journal) -> trc::Result<()> { + let mut batch = BatchBuilder::new(); + batch.set(class(journal.id), Json(journal).serialize()?); + data.write(batch.build_all()) + .await + .caused_by(trc::location!())?; + invalidate(); + Ok(()) +} + +/// Removes a journal. Its entries stay, each until its own time. +pub async fn delete(data: &Store, id: u32) -> trc::Result<()> { + let mut batch = BatchBuilder::new(); + batch.clear(class(id)); + data.write(batch.build_all()) + .await + .caused_by(trc::location!())?; + invalidate(); + Ok(()) +} + +/// How long a node keeps its copy of the journals before reading them again. +pub const TTL: Duration = Duration::from_secs(30); + +type Cached = Option<(Instant, Arc>)>; +static CACHE: RwLock = RwLock::new(None); + +/// Forgets this node's copy, so the next message reads the journals again. +pub fn invalidate() { + if let Ok(mut cache) = CACHE.write() { + *cache = None; + } +} + +/// The enabled journals, from this node's copy (refreshed every [`TTL`]). +pub async fn enabled(data: &Store) -> trc::Result>> { + if let Ok(cache) = CACHE.read() + && let Some((at, journals)) = cache.as_ref() + && at.elapsed() < TTL + { + return Ok(journals.clone()); + } + let journals = Arc::new( + all(data) + .await? + .into_iter() + .filter(|journal| journal.enabled) + .collect::>(), + ); + if let Ok(mut cache) = CACHE.write() { + *cache = Some((Instant::now(), journals.clone())); + } + Ok(journals) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn journal(scope: Scope) -> Journal { + Journal { + id: 1, + name: "Finance".into(), + description: String::new(), + enabled: true, + direction: Direction::Any, + scope, + retention_days: 365, + created_by: String::new(), + created_at: 0, + updated_at: 0, + } + } + + fn member(account: u32, groups: Vec) -> Member { + Member { + account, + domains: vec![1], + groups, + tenant: None, + } + } + + #[test] + fn scope_is_everyone_or_chosen() { + assert!( + journal(Scope { + everyone: true, + ..Default::default() + }) + .validate() + .is_ok() + ); + assert!(journal(Scope::default()).validate().is_err()); + let both = Scope { + everyone: true, + groups: vec![4], + ..Default::default() + }; + assert_eq!(journal(both).validate().unwrap_err().property, "scope"); + } + + #[test] + fn retention_has_bounds() { + let mut j = journal(Scope { + everyone: true, + ..Default::default() + }); + j.retention_days = 29; + assert_eq!(j.validate().unwrap_err().property, "retentionDays"); + j.retention_days = 3651; + assert!(j.validate().is_err()); + j.retention_days = 3650; + assert!(j.validate().is_ok()); + } + + #[test] + fn takes_by_direction_and_member() { + let mut j = journal(Scope { + groups: vec![7], + ..Default::default() + }); + assert!(j.takes(Direction::Outgoing, &[member(3, vec![7])])); + assert!(!j.takes(Direction::Outgoing, &[member(3, vec![8])])); + assert!(!j.takes(Direction::Outgoing, &[])); + j.direction = Direction::Incoming; + assert!(!j.takes(Direction::Outgoing, &[member(3, vec![7])])); + j.enabled = false; + assert!(!j.takes(Direction::Incoming, &[member(3, vec![7])])); + } + + #[test] + fn directions() { + assert_eq!(Direction::of(true, true, true), Direction::Outgoing); + assert_eq!(Direction::of(true, false, true), Direction::Internal); + assert_eq!(Direction::of(false, false, true), Direction::Incoming); + assert_eq!(Direction::of(false, true, true), Direction::Incoming); + } + + #[test] + fn scope_ids_are_jmap_ids() { + let scope: Scope = serde_json::from_str(r#"{"groups":["b"],"tenants":[7]}"#).unwrap(); + assert_eq!(scope.groups, vec![1]); + assert_eq!(scope.tenants, vec![7]); + assert_eq!( + serde_json::to_value(&scope).unwrap()["tenants"], + serde_json::json!(["h"]) + ); + } +} diff --git a/crates/features/src/journal/report.rs b/crates/features/src/journal/report.rs new file mode 100644 index 0000000..c8c3da8 --- /dev/null +++ b/crates/features/src/journal/report.rs @@ -0,0 +1,359 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! The journal report (JR-3, JR-4): a message whose first part lists the +//! envelope, one field a line, and whose second part is the message as it +//! was queued, byte for byte, as `message/rfc822`. Field names are fixed +//! English: a report is a record, and scripts read it. + +use super::Direction; +use mail_builder::headers::{Header, date::Date, text::Text}; +use mail_parser::MessageParser; +use sha2::{Digest, Sha256}; + +/// One envelope recipient, with the address it was given as (a list's, for +/// the list's members). +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Recipient { + pub address: String, + pub orcpt: Option, +} + +/// What the queue knows about a message. +#[derive(Debug, Clone)] +pub struct Envelope<'x> { + pub sender: &'x str, + pub authenticated: bool, + pub recipients: &'x [Recipient], + pub queue_id: u64, + /// Seconds. + pub received: u64, + pub direction: Direction, + pub held: bool, +} + +/// What a report says, besides the envelope's own fields. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct Fields { + pub subject: String, + pub message_id: String, + pub to: Vec, + pub cc: Vec, + /// Envelope recipients in neither To nor Cc, nor reached through a list. + pub bcc: Vec, + /// A list's address, and its members among the recipients. + pub expanded: Vec<(String, Vec)>, +} + +/// One line's worth of a value: no line breaks, no control characters. +fn line(value: &str) -> String { + value + .chars() + .map(|c| if c.is_control() { ' ' } else { c }) + .collect::() + .trim() + .to_string() +} + +/// The address an ORCPT names, without its `rfc822;` type. +fn orcpt_address(orcpt: &str) -> String { + let orcpt = orcpt.trim(); + let bare = match orcpt.split_once(';') { + Some((kind, address)) if kind.eq_ignore_ascii_case("rfc822") => address, + _ => orcpt, + }; + bare.trim().to_lowercase() +} + +/// Sorts the envelope's recipients by how they were addressed. +pub fn fields(envelope: &Envelope<'_>, original: &[u8]) -> Fields { + let parsed = MessageParser::default().parse_headers(original); + let headed = |which: Option<&mail_parser::Address<'_>>| -> Vec { + which + .map(|list| { + list.iter() + .filter_map(|addr| addr.address()) + .map(|address| address.to_lowercase()) + .collect() + }) + .unwrap_or_default() + }; + let (subject, message_id, header_to, header_cc) = match &parsed { + Some(message) => ( + message.subject().map(line).unwrap_or_default(), + message + .message_id() + .map(|id| format!("<{}>", line(id))) + .unwrap_or_default(), + headed(message.to()), + headed(message.cc()), + ), + None => Default::default(), + }; + + let mut fields = Fields { + subject, + message_id, + ..Default::default() + }; + for rcpt in envelope.recipients { + let address = rcpt.address.to_lowercase(); + let via = rcpt + .orcpt + .as_deref() + .map(orcpt_address) + .filter(|via| !via.is_empty() && *via != address); + if header_to.contains(&address) { + fields.to.push(line(&rcpt.address)); + } else if header_cc.contains(&address) { + fields.cc.push(line(&rcpt.address)); + } else if let Some(via) = via { + match fields.expanded.iter_mut().find(|(list, _)| *list == via) { + Some((_, members)) => members.push(line(&rcpt.address)), + None => fields + .expanded + .push((line(&via), vec![line(&rcpt.address)])), + } + } else { + fields.bcc.push(line(&rcpt.address)); + } + } + fields +} + +/// The report's first part. +pub fn text(envelope: &Envelope<'_>, fields: &Fields) -> String { + let mut out = String::new(); + let mut field = |name: &str, value: &str| { + if !value.is_empty() { + out.push_str(name); + out.push_str(": "); + out.push_str(value); + out.push_str("\r\n"); + } + }; + let sender = if envelope.sender.is_empty() { + "<>".to_string() + } else { + line(envelope.sender) + }; + field("Sender", &sender); + field( + "Authenticated", + if envelope.authenticated { "yes" } else { "no" }, + ); + field("Subject", &fields.subject); + field("Message-ID", &fields.message_id); + field("Queue ID", &format!("{:x}", envelope.queue_id)); + field( + "Received", + &mail_parser::DateTime::from_timestamp(envelope.received as i64).to_rfc3339(), + ); + field("Direction", envelope.direction.as_str()); + field("To", &fields.to.join(", ")); + field("Cc", &fields.cc.join(", ")); + field("Bcc", &fields.bcc.join(", ")); + for (list, members) in &fields.expanded { + field("Expanded", &format!("{list} -> {}", members.join(", "))); + } + if envelope.held { + field("Held for review", "yes"); + } + out +} + +fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|b| format!("{b:02x}")).collect() +} + +/// Whether a message can travel as 8bit: no NULs, no line past 998 bytes. +fn fits_8bit(message: &[u8]) -> bool { + !message.contains(&0) && message.split(|b| *b == b'\n').all(|l| l.len() <= 998) +} + +/// The whole report: headers, the fields, then the original untouched. +/// `from` is the address the report is from; `host` names the server in its +/// Message-ID. +pub fn build( + envelope: &Envelope<'_>, + original: &[u8], + from: &str, + host: &str, +) -> (Vec, Fields) { + let fields = fields(envelope, original); + let body = text(envelope, &fields); + // A boundary that can't occur in the original + let mut boundary = format!("journal-{}", &hex(&Sha256::digest(original))[..32]); + while original + .windows(boundary.len()) + .any(|window| window == boundary.as_bytes()) + { + boundary.push('x'); + } + + let mut out: Vec = Vec::with_capacity(original.len() + body.len() + 1024); + out.extend_from_slice(format!("From: Journal <{}>\r\n", line(from)).as_bytes()); + out.extend_from_slice(b"Date: "); + out.extend_from_slice(Date::new(envelope.received as i64).to_rfc822().as_bytes()); + out.extend_from_slice(b"\r\n"); + out.extend_from_slice(b"Subject: "); + let subject = if fields.subject.is_empty() { + "Journal report".to_string() + } else { + format!("Journal report: {}", fields.subject) + }; + Text::new(subject).write_header(&mut out, "Subject: ".len()); + out.extend_from_slice( + format!( + "Message-ID: \r\n", + envelope.queue_id, + envelope.received, + line(host) + ) + .as_bytes(), + ); + out.extend_from_slice(format!("X-Inbuxa-Journal: {:x}\r\n", envelope.queue_id).as_bytes()); + out.extend_from_slice(b"MIME-Version: 1.0\r\n"); + out.extend_from_slice( + format!("Content-Type: multipart/mixed; boundary=\"{boundary}\"\r\n\r\n").as_bytes(), + ); + out.extend_from_slice(format!("--{boundary}\r\n").as_bytes()); + out.extend_from_slice( + b"Content-Type: text/plain; charset=utf-8\r\nContent-Transfer-Encoding: 8bit\r\n\r\n", + ); + out.extend_from_slice(body.as_bytes()); + out.extend_from_slice(format!("\r\n--{boundary}\r\n").as_bytes()); + out.extend_from_slice(b"Content-Type: message/rfc822\r\n"); + out.extend_from_slice(b"Content-Disposition: attachment; filename=\"original.eml\"\r\n"); + out.extend_from_slice(if fits_8bit(original) { + b"Content-Transfer-Encoding: 8bit\r\n\r\n".as_slice() + } else { + b"Content-Transfer-Encoding: binary\r\n\r\n".as_slice() + }); + out.extend_from_slice(original); + // The line break before a boundary belongs to the boundary: the + // original keeps its own last one + out.extend_from_slice(format!("\r\n--{boundary}--\r\n").as_bytes()); + (out, fields) +} + +/// Where the original starts and ends inside a report [`build`] made. +pub fn original(report: &[u8]) -> Option<&[u8]> { + let parsed = MessageParser::default().parse(report)?; + let part = parsed.attachment(0)?; + let start = part.raw_body_offset() as usize; + let end = part.raw_end_offset() as usize; + report.get(start..end) +} + +#[cfg(test)] +mod tests { + use super::*; + + const ORIGINAL: &[u8] = b"From: alice@example.com\r\n\ +To: Bank \r\n\ +Cc: bob@example.com\r\n\ +Subject: Q3 figures\r\n\ +Message-ID: \r\n\ +\r\n\ +The figures.\r\n"; + + fn rcpt(address: &str, orcpt: Option<&str>) -> Recipient { + Recipient { + address: address.into(), + orcpt: orcpt.map(Into::into), + } + } + + fn envelope(recipients: &[Recipient]) -> Envelope<'_> { + Envelope { + sender: "alice@example.com", + authenticated: true, + recipients, + queue_id: 0x1a2b, + received: 1_790_000_000, + direction: Direction::Outgoing, + held: false, + } + } + + #[test] + fn recipients_sorted_by_how_they_were_addressed() { + let recipients = [ + rcpt("pay@bank.example", None), + rcpt("Bob@example.com", Some("rfc822;bob@example.com")), + rcpt("carol@example.com", None), + rcpt("dan@example.com", Some("finance@example.com")), + rcpt("erin@example.com", Some("rfc822;Finance@example.com")), + ]; + let fields = fields(&envelope(&recipients), ORIGINAL); + assert_eq!(fields.subject, "Q3 figures"); + assert_eq!(fields.message_id, ""); + assert_eq!(fields.to, vec!["pay@bank.example"]); + assert_eq!(fields.cc, vec!["Bob@example.com"]); + assert_eq!(fields.bcc, vec!["carol@example.com"]); + assert_eq!( + fields.expanded, + vec![( + "finance@example.com".to_string(), + vec![ + "dan@example.com".to_string(), + "erin@example.com".to_string() + ] + )] + ); + } + + #[test] + fn report_carries_the_original_untouched() { + let recipients = [ + rcpt("pay@bank.example", None), + rcpt("carol@example.com", None), + ]; + let (report, _) = build( + &envelope(&recipients), + ORIGINAL, + "postmaster@example.com", + "mx.example.com", + ); + let text = String::from_utf8_lossy(&report); + assert!(text.contains("Sender: alice@example.com\r\n")); + assert!(text.contains("Bcc: carol@example.com\r\n")); + assert!(text.contains("Queue ID: 1a2b\r\n")); + assert!(text.contains("Direction: outgoing\r\n")); + assert!(text.contains("Subject: Journal report: Q3 figures\r\n")); + assert!(!text.contains("Held for review")); + assert_eq!(original(&report), Some(ORIGINAL)); + let unterminated = &ORIGINAL[..ORIGINAL.len() - 2]; + let (report, _) = build( + &envelope(&recipients), + unterminated, + "postmaster@example.com", + "mx.example.com", + ); + assert_eq!(original(&report), Some(unterminated)); + } + + #[test] + fn values_stay_on_one_line() { + let recipients = [rcpt("x@example.com", None)]; + let mut env = envelope(&recipients); + env.sender = "evil@example.com\r\nBcc: nobody@example.com"; + env.held = true; + let body = text(&env, &Fields::default()); + assert_eq!(body.matches("\r\n").count(), body.lines().count()); + assert!(body.contains("Sender: evil@example.com Bcc: nobody@example.com\r\n")); + assert!(body.contains("Held for review: yes\r\n")); + } + + #[test] + fn an_empty_sender_is_shown_as_such() { + let recipients = [rcpt("x@example.com", None)]; + let mut env = envelope(&recipients); + env.sender = ""; + assert!(text(&env, &Fields::default()).starts_with("Sender: <>\r\n")); + } +} diff --git a/crates/features/src/lib.rs b/crates/features/src/lib.rs index 9df288c..1351d41 100644 --- a/crates/features/src/lib.rs +++ b/crates/features/src/lib.rs @@ -22,6 +22,7 @@ pub mod ai; pub mod audit; pub mod branding; pub mod hold; +pub mod journal; pub mod lock; pub mod mailflow; pub mod masked_email; diff --git a/crates/features/src/mailflow/rules.rs b/crates/features/src/mailflow/rules.rs index 1816523..a94174b 100644 --- a/crates/features/src/mailflow/rules.rs +++ b/crates/features/src/mailflow/rules.rs @@ -62,7 +62,7 @@ fn one() -> u32 { /// Group and tenant ids in the JMAP form clients use (`"b"`, `"c"`…), held /// as numbers for matching. Plain numbers are read too. -mod jmap_ids { +pub(crate) mod jmap_ids { use serde::{Deserialize, Deserializer, Serializer, de::Error, ser::SerializeSeq}; use std::str::FromStr; use types::id::Id; diff --git a/crates/jmap-proto/src/object/inbuxa_journal.rs b/crates/jmap-proto/src/object/inbuxa_journal.rs new file mode 100644 index 0000000..5b79069 --- /dev/null +++ b/crates/jmap-proto/src/object/inbuxa_journal.rs @@ -0,0 +1,205 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! `inbuxa:Journal/get` and `/set` under `urn:inbuxa:jmap`: journals +//! (journaling spec, JR-9, JR-12). What a journal has taken stays when the +//! journal changes or goes; each entry keeps its own retention. + +use crate::{ + object::{AnyId, JmapObject, JmapObjectId}, + request::deserialize::DeserializeArguments, +}; +use jmap_tools::{Element, Key, Property}; +use std::{borrow::Cow, str::FromStr}; +use types::id::Id; + +#[derive(Debug, Clone, Default)] +pub struct Journal; + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum JournalProperty { + Id, + Name, + Description, + Enabled, + /// `outgoing`, `incoming`, `internal` or `any`. + Direction, + /// Everyone, or chosen accounts, groups, domains and tenants. + Scope, + /// How long an entry is kept; each keeps what it was written with. + RetentionDays, + CreatedBy, + CreatedAt, + UpdatedAt, +} + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum JournalValue { + Id(Id), +} + +impl Property for JournalProperty { + fn try_parse(parent: Option<&Key<'_, Self>>, value: &str) -> Option { + // Keys inside the scope stay plain keys + match parent { + None => JournalProperty::parse(value), + Some(_) => None, + } + } + + fn to_cow(&self) -> Cow<'static, str> { + match self { + JournalProperty::Id => "id", + JournalProperty::Name => "name", + JournalProperty::Description => "description", + JournalProperty::Enabled => "enabled", + JournalProperty::Direction => "direction", + JournalProperty::Scope => "scope", + JournalProperty::RetentionDays => "retentionDays", + JournalProperty::CreatedBy => "createdBy", + JournalProperty::CreatedAt => "createdAt", + JournalProperty::UpdatedAt => "updatedAt", + } + .into() + } +} + +impl JournalProperty { + fn parse(value: &str) -> Option { + hashify::tiny_map!(value.as_bytes(), + b"id" => JournalProperty::Id, + b"name" => JournalProperty::Name, + b"description" => JournalProperty::Description, + b"enabled" => JournalProperty::Enabled, + b"direction" => JournalProperty::Direction, + b"scope" => JournalProperty::Scope, + b"retentionDays" => JournalProperty::RetentionDays, + b"createdBy" => JournalProperty::CreatedBy, + b"createdAt" => JournalProperty::CreatedAt, + b"updatedAt" => JournalProperty::UpdatedAt, + ) + } +} + +impl FromStr for JournalProperty { + type Err = (); + + fn from_str(s: &str) -> Result { + JournalProperty::parse(s).ok_or(()) + } +} + +impl Element for JournalValue { + type Property = JournalProperty; + + fn try_parse

(key: &Key<'_, Self::Property>, value: &str) -> Option { + match key { + Key::Property(JournalProperty::Id) => Id::from_str(value).ok().map(JournalValue::Id), + _ => None, + } + } + + fn to_cow(&self) -> Cow<'static, str> { + match self { + JournalValue::Id(id) => id.to_string().into(), + } + } +} + +/// The set call's own argument: why, for the audit log. +#[derive(Debug, Clone, Default)] +pub struct JournalSetArguments { + pub reason: Option, +} + +impl<'de> DeserializeArguments<'de> for JournalSetArguments { + fn deserialize_argument(&mut self, key: &str, map: &mut A) -> Result<(), A::Error> + where + A: serde::de::MapAccess<'de>, + { + if key == "reason" { + self.reason = map.next_value()?; + } else { + let _ = map.next_value::()?; + } + Ok(()) + } +} + +impl JmapObject for Journal { + type Property = JournalProperty; + + type Element = JournalValue; + + type Id = Id; + + type Filter = (); + + type Comparator = (); + + type GetArguments = (); + + type SetArguments<'de> = JournalSetArguments; + + type QueryArguments = (); + + type CopyArguments = (); + + type ParseArguments = (); + + const ID_PROPERTY: Self::Property = JournalProperty::Id; +} + +impl From for JournalValue { + fn from(id: Id) -> Self { + JournalValue::Id(id) + } +} + +impl JmapObjectId for JournalValue { + fn as_id(&self) -> Option { + match self { + JournalValue::Id(id) => Some(*id), + } + } + + fn as_any_id(&self) -> Option { + match self { + JournalValue::Id(id) => Some(AnyId::Id(*id)), + } + } + + fn as_id_ref(&self) -> Option<&str> { + None + } + + fn try_set_id(&mut self, new_id: AnyId) -> bool { + if let AnyId::Id(id) = new_id { + *self = JournalValue::Id(id); + true + } else { + false + } + } +} + +impl JmapObjectId for JournalProperty { + fn as_id(&self) -> Option { + None + } + + fn as_any_id(&self) -> Option { + None + } + + fn as_id_ref(&self) -> Option<&str> { + None + } + + fn try_set_id(&mut self, _: AnyId) -> bool { + false + } +} diff --git a/crates/jmap-proto/src/object/mod.rs b/crates/jmap-proto/src/object/mod.rs index 003cf17..7d91f9a 100644 --- a/crates/jmap-proto/src/object/mod.rs +++ b/crates/jmap-proto/src/object/mod.rs @@ -30,6 +30,7 @@ pub mod inbuxa_inventory_snapshot; // inbuxa: personal-data catalog pub mod inbuxa_audit; // inbuxa: the audit log pub mod inbuxa_legal_hold; // inbuxa: legal hold pub mod inbuxa_mail_rule; // inbuxa: DLP and mail flow rules +pub mod inbuxa_journal; // inbuxa: journaling pub mod inbuxa_held_message; // inbuxa: mail held for review pub mod inbuxa_hold_export; // inbuxa: legal hold exports pub mod inbuxa_explanation; // inbuxa: "Explain this" with the local model diff --git a/crates/jmap-proto/src/references/eval.rs b/crates/jmap-proto/src/references/eval.rs index 216d826..08a22b4 100644 --- a/crates/jmap-proto/src/references/eval.rs +++ b/crates/jmap-proto/src/references/eval.rs @@ -88,6 +88,9 @@ impl Response<'_> { GetResponseMethod::MailRule(response) => { response.eval_jptr(path, &mut results) } + GetResponseMethod::Journal(response) => { + response.eval_jptr(path, &mut results) + } GetResponseMethod::HeldMessage(response) => { response.eval_jptr(path, &mut results) } diff --git a/crates/jmap-proto/src/references/resolve.rs b/crates/jmap-proto/src/references/resolve.rs index b4a3ed3..dd69481 100644 --- a/crates/jmap-proto/src/references/resolve.rs +++ b/crates/jmap-proto/src/references/resolve.rs @@ -55,6 +55,7 @@ impl Response<'_> { GetRequestMethod::AccountLock(request) => request.resolve_references(self)?, GetRequestMethod::LegalHold(request) => request.resolve_references(self)?, GetRequestMethod::MailRule(request) => request.resolve_references(self)?, + GetRequestMethod::Journal(request) => request.resolve_references(self)?, GetRequestMethod::HeldMessage(request) => request.resolve_references(self)?, GetRequestMethod::HoldExport(request) => request.resolve_references(self)?, GetRequestMethod::ProtocolPolicy(request) => request.resolve_references(self)?, @@ -131,6 +132,9 @@ impl Response<'_> { SetRequestMethod::MailRule(request) => { request.resolve_references(self, 1, false)? } + SetRequestMethod::Journal(request) => { + request.resolve_references(self, 1, false)? + } SetRequestMethod::HeldMessage(request) => { request.resolve_references(self, 1, false)? } diff --git a/crates/jmap-proto/src/request/method.rs b/crates/jmap-proto/src/request/method.rs index 99d5258..b42d028 100644 --- a/crates/jmap-proto/src/request/method.rs +++ b/crates/jmap-proto/src/request/method.rs @@ -69,6 +69,8 @@ pub enum MethodObject { // inbuxa: DLP and mail flow rules MailRule, HeldMessage, + // inbuxa: journaling + Journal, TenantProtocolPolicy, } @@ -109,7 +111,8 @@ impl MethodObject { | MethodObject::LegalHold | MethodObject::HoldExport | MethodObject::MailRule - | MethodObject::HeldMessage => Capability::Inbuxa, + | MethodObject::HeldMessage + | MethodObject::Journal => Capability::Inbuxa, MethodObject::ProtocolPolicy => Capability::Inbuxa, MethodObject::TenantProtocolPolicy => Capability::Inbuxa, } @@ -307,6 +310,8 @@ impl MethodName { (MethodFunction::Set, MethodObject::LegalHold) => "inbuxa:LegalHold/set", (MethodFunction::Get, MethodObject::MailRule) => "inbuxa:MailRule/get", (MethodFunction::Set, MethodObject::MailRule) => "inbuxa:MailRule/set", + (MethodFunction::Get, MethodObject::Journal) => "inbuxa:Journal/get", + (MethodFunction::Set, MethodObject::Journal) => "inbuxa:Journal/set", (MethodFunction::Get, MethodObject::HeldMessage) => "inbuxa:HeldMessage/get", (MethodFunction::Set, MethodObject::HeldMessage) => "inbuxa:HeldMessage/set", (MethodFunction::Get, MethodObject::HoldExport) => "inbuxa:HoldExport/get", @@ -465,6 +470,8 @@ impl MethodName { "inbuxa:LegalHold/set" => (MethodObject::LegalHold, MethodFunction::Set), "inbuxa:MailRule/get" => (MethodObject::MailRule, MethodFunction::Get), "inbuxa:MailRule/set" => (MethodObject::MailRule, MethodFunction::Set), + "inbuxa:Journal/get" => (MethodObject::Journal, MethodFunction::Get), + "inbuxa:Journal/set" => (MethodObject::Journal, MethodFunction::Set), "inbuxa:HeldMessage/get" => (MethodObject::HeldMessage, MethodFunction::Get), "inbuxa:HeldMessage/set" => (MethodObject::HeldMessage, MethodFunction::Set), "inbuxa:HoldExport/get" => (MethodObject::HoldExport, MethodFunction::Get), @@ -539,6 +546,7 @@ impl Display for MethodObject { MethodObject::AccountLock => "inbuxa:AccountLock", MethodObject::LegalHold => "inbuxa:LegalHold", MethodObject::MailRule => "inbuxa:MailRule", + MethodObject::Journal => "inbuxa:Journal", MethodObject::HeldMessage => "inbuxa:HeldMessage", MethodObject::HoldExport => "inbuxa:HoldExport", MethodObject::ProtocolPolicy => "inbuxa:ProtocolPolicy", diff --git a/crates/jmap-proto/src/request/mod.rs b/crates/jmap-proto/src/request/mod.rs index 0923c14..b2ad5ef 100644 --- a/crates/jmap-proto/src/request/mod.rs +++ b/crates/jmap-proto/src/request/mod.rs @@ -125,6 +125,7 @@ pub enum GetRequestMethod { AccountLock(Box>), LegalHold(Box>), MailRule(Box>), + Journal(Box>), HeldMessage(Box>), HoldExport(Box>), ProtocolPolicy(Box>), @@ -163,6 +164,7 @@ pub enum SetRequestMethod<'x> { AccountLock(Box>), LegalHold(Box>), MailRule(Box>), + Journal(Box>), HeldMessage(Box>), HoldExport(Box>), ProtocolPolicy(Box>), diff --git a/crates/jmap-proto/src/request/parser.rs b/crates/jmap-proto/src/request/parser.rs index 4d067db..7c0f1b0 100644 --- a/crates/jmap-proto/src/request/parser.rs +++ b/crates/jmap-proto/src/request/parser.rs @@ -653,6 +653,21 @@ impl<'de> Visitor<'de> for CallVisitor { return Err(de::Error::invalid_length(1, &self)); } }, + // inbuxa: journaling + (MethodFunction::Get, MethodObject::Journal) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Get(GetRequestMethod::Journal(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, + (MethodFunction::Set, MethodObject::Journal) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Set(SetRequestMethod::Journal(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, // inbuxa: legal hold (MethodFunction::Get, MethodObject::LegalHold) => match seq.next_element() { Ok(Some(value)) => RequestMethod::Get(GetRequestMethod::LegalHold(value)), diff --git a/crates/jmap-proto/src/response/mod.rs b/crates/jmap-proto/src/response/mod.rs index 6320886..d9e9e74 100644 --- a/crates/jmap-proto/src/response/mod.rs +++ b/crates/jmap-proto/src/response/mod.rs @@ -112,6 +112,7 @@ pub enum GetResponseMethod { AccountLock(GetResponse), LegalHold(GetResponse), MailRule(GetResponse), + Journal(GetResponse), HeldMessage(GetResponse), HoldExport(GetResponse), ProtocolPolicy(GetResponse), @@ -150,6 +151,7 @@ pub enum SetResponseMethod { AccountLock(Box>), LegalHold(Box>), MailRule(Box>), + Journal(Box>), HeldMessage(Box>), HoldExport(Box>), Explanation(Box>), @@ -841,6 +843,19 @@ impl<'x> From> for Respon } } +// inbuxa: journaling +impl<'x> From> for ResponseMethod<'x> { + fn from(value: GetResponse) -> Self { + ResponseMethod::Get(GetResponseMethod::Journal(value)) + } +} + +impl<'x> From> for ResponseMethod<'x> { + fn from(value: SetResponse) -> Self { + ResponseMethod::Set(SetResponseMethod::Journal(Box::new(value))) + } +} + impl<'x> From> for ResponseMethod<'x> { fn from(value: GetResponse) -> Self { ResponseMethod::Get(GetResponseMethod::LegalHold(value)) diff --git a/crates/jmap/src/api/auth.rs b/crates/jmap/src/api/auth.rs index 83af147..73b7e36 100644 --- a/crates/jmap/src/api/auth.rs +++ b/crates/jmap/src/api/auth.rs @@ -116,6 +116,8 @@ impl JmapAuthorization for AccessToken { Permission::SysDlpPolicyGet } } + // inbuxa: journaling (JR-18) + GetRequestMethod::Journal(_) => Permission::SysJournalGet, GetRequestMethod::HoldExport(_) => Permission::SysLegalHoldExport, // inbuxa: legacy protocols off. It takes listeners away and // puts them back, so it takes the listener's permissions @@ -288,6 +290,14 @@ impl JmapAuthorization for AccessToken { .details("You are not authorized to change mail rules")) } } + // inbuxa: journaling (JR-18) + SetRequestMethod::Journal(s) => validate_set( + s, + self, + Permission::SysJournalUpdate, + Permission::SysJournalUpdate, + Permission::SysJournalUpdate, + ), // inbuxa: LH-12, exporting held data SetRequestMethod::HoldExport(s) => validate_set( s, @@ -451,6 +461,7 @@ impl JmapAuthorization for AccessToken { | MethodObject::HoldExport | MethodObject::MailRule | MethodObject::HeldMessage + | MethodObject::Journal | MethodObject::ProtocolPolicy | MethodObject::TenantProtocolPolicy => Permission::JmapEmailChanges, // inbuxa: x:MaskedEmail/changes reads what /get reads diff --git a/crates/jmap/src/api/request.rs b/crates/jmap/src/api/request.rs index 3dbc901..d627d51 100644 --- a/crates/jmap/src/api/request.rs +++ b/crates/jmap/src/api/request.rs @@ -285,6 +285,9 @@ impl RequestHandler for Server { SetResponseMethod::MailRule(set_response) => { set_response.update_created_ids(&mut response); } + SetResponseMethod::Journal(set_response) => { + set_response.update_created_ids(&mut response); + } SetResponseMethod::HeldMessage(set_response) => { set_response.update_created_ids(&mut response); } @@ -512,6 +515,11 @@ impl RequestHandler for Server { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; crate::inbuxa::mail_rule::get(self, access_token, *req).await?.into() } + // inbuxa: journaling + GetRequestMethod::Journal(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + crate::inbuxa::journal::get(self, access_token, *req).await?.into() + } // inbuxa: the audit log (AU-9) GetRequestMethod::AuditEvent(mut req) => { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; @@ -985,6 +993,22 @@ impl RequestHandler for Server { .await? .into() } + SetRequestMethod::Journal(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + let reason = req.arguments.reason.clone(); + crate::inbuxa::audit::recorded( + self, + access_token, + session, + &method_name.obj.to_string(), + None, + reason, + *req, + |req| Box::pin(crate::inbuxa::journal::set(self, access_token, req)), + ) + .await? + .into() + } SetRequestMethod::AuditExport(mut req) => { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; crate::inbuxa::audit_log::export_set(self, access_token, session, *req) diff --git a/crates/jmap/src/changes/get.rs b/crates/jmap/src/changes/get.rs index 343ef12..5370409 100644 --- a/crates/jmap/src/changes/get.rs +++ b/crates/jmap/src/changes/get.rs @@ -431,6 +431,7 @@ impl IntermediateChangesResponse { | MethodObject::LegalHold | MethodObject::HoldExport | MethodObject::MailRule + | MethodObject::Journal | MethodObject::HeldMessage | MethodObject::ProtocolPolicy | MethodObject::TenantProtocolPolicy diff --git a/crates/jmap/src/inbuxa/journal.rs b/crates/jmap/src/inbuxa/journal.rs new file mode 100644 index 0000000..2497235 --- /dev/null +++ b/crates/jmap/src/inbuxa/journal.rs @@ -0,0 +1,273 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! `inbuxa:Journal` (journaling spec, JR-9, JR-12, JR-18): journals, seen +//! with `sysJournalGet` and changed with `sysJournalUpdate`, which the +//! request layer checks. Journals are the server's: nobody in a tenant +//! reaches them. The request layer records every change in the audit log. +//! Changing or removing a journal never touches what it has taken. + +use common::{Server, auth::AccessToken}; +use inbuxa_features::journal::{self, Journal as Stored}; +use jmap_proto::{ + error::set::SetError, + method::{ + get::{GetRequest, GetResponse}, + set::{SetRequest, SetResponse}, + }, + object::inbuxa_journal::{Journal, JournalProperty as P, JournalValue}, + request::IntoValid, + types::date::UTCDate, +}; +use jmap_tools::{Key, Map, Property, Value}; +use std::borrow::Cow; +use store::write::now; +use types::id::Id; + +type JValue = Value<'static, P, JournalValue>; + +const ALL: &[P] = &[ + P::Id, + P::Name, + P::Description, + P::Enabled, + P::Direction, + P::Scope, + P::RetentionDays, + P::CreatedBy, + P::CreatedAt, + P::UpdatedAt, +]; + +/// Properties the server sets; a client that sends them is refused. +const SERVER_SET: &[P] = &[P::Id, P::CreatedBy, P::CreatedAt, P::UpdatedAt]; + +fn server_level(access_token: &AccessToken) -> trc::Result<()> { + if access_token.tenant_id().is_some() { + Err(trc::JmapEvent::Forbidden + .into_err() + .details("Journals are the server's.")) + } else { + Ok(()) + } +} + +fn json_to_value(json: serde_json::Value) -> JValue { + match json { + serde_json::Value::Null => Value::Null, + serde_json::Value::Bool(b) => Value::Bool(b), + serde_json::Value::Number(n) => { + if let Some(n) = n.as_u64() { + Value::Number(n.into()) + } else if let Some(n) = n.as_i64() { + Value::Number(n.into()) + } else { + Value::Number(n.as_f64().unwrap_or_default().into()) + } + } + serde_json::Value::String(s) => Value::Str(Cow::Owned(s)), + serde_json::Value::Array(items) => { + Value::Array(items.into_iter().map(json_to_value).collect()) + } + serde_json::Value::Object(map) => { + let mut out = Map::with_capacity(map.len()); + for (key, value) in map { + out.insert_unchecked(Key::Owned(key), json_to_value(value)); + } + Value::Object(out) + } + } +} + +fn date(seconds: u64) -> JValue { + Value::Str(UTCDate::from_timestamp(seconds as i64).to_string().into()) +} + +fn to_value(journal: &Stored, properties: &[P]) -> JValue { + let json = serde_json::to_value(journal).unwrap_or_default(); + let mut out = Map::with_capacity(properties.len()); + for property in properties { + let value = match property { + P::Id => Value::Element(JournalValue::Id(Id::from(journal.id))), + P::CreatedAt => date(journal.created_at), + P::UpdatedAt => date(journal.updated_at), + other => json + .get(other.to_cow().as_ref()) + .cloned() + .map_or(Value::Null, json_to_value), + }; + out.insert_unchecked(Key::Property(property.clone()), value); + } + Value::Object(out) +} + +/// A journal as sent: its JSON object, top-level keys only those a client +/// may set. +fn client_json( + value: Value<'_, P, JournalValue>, +) -> Result, SetError

> { + let mut map = serde_json::Map::new(); + for (key, value) in value.into_expanded_object() { + match &key { + Key::Property(p) if SERVER_SET.contains(p) => { + return Err(SetError::invalid_properties() + .with_property(p.clone()) + .with_description("The server sets this.")); + } + Key::Property(p) => { + map.insert(p.to_cow().into_owned(), value.into()); + } + _ => { + return Err(SetError::invalid_properties().with_property(key.clone().into_owned())); + } + } + } + Ok(map) +} + +fn parse(json: serde_json::Map) -> Result> { + let journal: Stored = + serde_json::from_value(serde_json::Value::Object(json)).map_err(|err| { + SetError::invalid_properties().with_description(format!("Not a valid journal: {err}")) + })?; + journal.validate().map_err(|invalid| { + let property = invalid.property.parse::

().unwrap_or(P::Name); + SetError::invalid_properties() + .with_property(property) + .with_description(invalid.reason) + })?; + Ok(journal) +} + +fn journal_id(id: Id) -> Option { + u32::try_from(id.id()).ok() +} + +/// `inbuxa:Journal/get`: every journal, oldest first. +pub async fn get( + server: &Server, + access_token: &AccessToken, + mut request: GetRequest, +) -> trc::Result> { + server_level(access_token)?; + let properties = request.unwrap_properties(ALL); + let (ids, not_found) = request.unwrap_ids(server.core.jmap.get_max_objects)?; + let mut response = GetResponse { + account_id: request.account_id.into(), + state: None, + list: Vec::new(), + not_found, + }; + let journals = journal::all(server.store()).await?; + match ids { + None => { + response.list = journals + .iter() + .map(|journal| to_value(journal, &properties)) + .collect() + } + Some(ids) => { + for id in ids { + match journal_id(id).and_then(|id| journals.iter().find(|j| j.id == id)) { + Some(journal) => response.list.push(to_value(journal, &properties)), + None => response.push_not_found(id), + } + } + } + } + Ok(response) +} + +/// `inbuxa:Journal/set`: create, change or remove journals. +pub async fn set( + server: &Server, + access_token: &AccessToken, + mut request: SetRequest<'_, Journal>, +) -> trc::Result> { + server_level(access_token)?; + let mut response = SetResponse::from_request(&request, server.core.jmap.set_max_objects)?; + let data = server.store(); + let actor = server.audit_actor(access_token).await; + + for (client_id, value) in request.unwrap_create() { + let stored = match client_json(value).and_then(parse) { + Ok(stored) => stored, + Err(error) => { + response.not_created.append(client_id, error); + continue; + } + }; + let at = now(); + let stored = Stored { + created_by: actor.name.clone(), + created_at: at, + updated_at: at, + ..stored + }; + let id = journal::create(data, &stored).await?; + let mut out = Map::with_capacity(1); + out.insert_unchecked( + Key::Property(P::Id), + Value::Element(JournalValue::Id(Id::from(id))), + ); + response.created.insert(client_id, Value::Object(out)); + } + + for (id, value) in request.unwrap_update().into_valid() { + let Some(current) = (match journal_id(id) { + Some(journal_id) => journal::get(data, journal_id).await?, + None => None, + }) else { + response.not_updated.append(id, SetError::not_found()); + continue; + }; + // The stored journal, with each property sent replacing its own + let mut json = match serde_json::to_value(¤t) { + Ok(serde_json::Value::Object(map)) => map, + _ => serde_json::Map::new(), + }; + let changes = match client_json(value) { + Ok(changes) => changes, + Err(error) => { + response.not_updated.append(id, error); + continue; + } + }; + json.extend(changes); + let next = match parse(json) { + Ok(next) => next, + Err(error) => { + response.not_updated.append(id, error); + continue; + } + }; + let next = Stored { + id: current.id, + created_by: current.created_by.clone(), + created_at: current.created_at, + updated_at: now(), + ..next + }; + if next != current { + journal::update(data, &next).await?; + } + response.updated.append(id, None); + } + + for id in request.unwrap_destroy().into_valid() { + let Some(current) = (match journal_id(id) { + Some(journal_id) => journal::get(data, journal_id).await?, + None => None, + }) else { + response.not_destroyed.append(id, SetError::not_found()); + continue; + }; + journal::delete(data, current.id).await?; + response.destroyed.push(id); + } + + Ok(response) +} diff --git a/crates/jmap/src/inbuxa/mod.rs b/crates/jmap/src/inbuxa/mod.rs index 9e978f1..5d2ab1b 100644 --- a/crates/jmap/src/inbuxa/mod.rs +++ b/crates/jmap/src/inbuxa/mod.rs @@ -11,6 +11,7 @@ pub mod access; pub mod account_lock; pub mod legal_hold; pub mod mail_rule; +pub mod journal; pub mod held_message; pub mod dlp_settings; pub mod hold_export; diff --git a/crates/registry/src/schema/enums.rs b/crates/registry/src/schema/enums.rs index f578905..e0d5c98 100644 --- a/crates/registry/src/schema/enums.rs +++ b/crates/registry/src/schema/enums.rs @@ -1755,6 +1755,11 @@ pub enum Permission { SysDlpPolicyUpdate = 677, SysDlpReviewGet = 678, SysDlpReviewUpdate = 679, + // inbuxa: journaling + SysJournalGet = 680, + SysJournalUpdate = 681, + SysJournalSearch = 682, + SysJournalExport = 683, SysAccountGet = 219, SysAccountCreate = 220, SysAccountUpdate = 221, diff --git a/crates/registry/src/schema/enums_impl.rs b/crates/registry/src/schema/enums_impl.rs index ed72ffe..f8cba76 100644 --- a/crates/registry/src/schema/enums_impl.rs +++ b/crates/registry/src/schema/enums_impl.rs @@ -7097,6 +7097,10 @@ impl EnumImpl for Permission { b"sysDlpPolicyUpdate" => Permission::SysDlpPolicyUpdate, b"sysDlpReviewGet" => Permission::SysDlpReviewGet, b"sysDlpReviewUpdate" => Permission::SysDlpReviewUpdate, + b"sysJournalGet" => Permission::SysJournalGet, + b"sysJournalUpdate" => Permission::SysJournalUpdate, + b"sysJournalSearch" => Permission::SysJournalSearch, + b"sysJournalExport" => Permission::SysJournalExport, b"sysAccountGet" => Permission::SysAccountGet, b"sysAccountCreate" => Permission::SysAccountCreate, b"sysAccountUpdate" => Permission::SysAccountUpdate, @@ -7793,6 +7797,10 @@ impl EnumImpl for Permission { Permission::SysDlpPolicyUpdate => "sysDlpPolicyUpdate", Permission::SysDlpReviewGet => "sysDlpReviewGet", Permission::SysDlpReviewUpdate => "sysDlpReviewUpdate", + Permission::SysJournalGet => "sysJournalGet", + Permission::SysJournalUpdate => "sysJournalUpdate", + Permission::SysJournalSearch => "sysJournalSearch", + Permission::SysJournalExport => "sysJournalExport", Permission::SysAccountGet => "sysAccountGet", Permission::SysAccountCreate => "sysAccountCreate", Permission::SysAccountUpdate => "sysAccountUpdate", @@ -8482,6 +8490,10 @@ impl EnumImpl for Permission { 677 => Some(Permission::SysDlpPolicyUpdate), 678 => Some(Permission::SysDlpReviewGet), 679 => Some(Permission::SysDlpReviewUpdate), + 680 => Some(Permission::SysJournalGet), + 681 => Some(Permission::SysJournalUpdate), + 682 => Some(Permission::SysJournalSearch), + 683 => Some(Permission::SysJournalExport), 219 => Some(Permission::SysAccountGet), 220 => Some(Permission::SysAccountCreate), 221 => Some(Permission::SysAccountUpdate), @@ -8926,7 +8938,7 @@ impl EnumImpl for Permission { } } - const COUNT: usize = 680; + const COUNT: usize = 684; } impl serde::Serialize for Permission { diff --git a/crates/services/src/task_manager/maintenance.rs b/crates/services/src/task_manager/maintenance.rs index b8e2414..3f404e9 100644 --- a/crates/services/src/task_manager/maintenance.rs +++ b/crates/services/src/task_manager/maintenance.rs @@ -291,6 +291,12 @@ async fn store_maintenance( 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 { @@ -409,6 +415,54 @@ async fn store_maintenance( 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, diff --git a/crates/smtp/src/queue/journal.rs b/crates/smtp/src/queue/journal.rs new file mode 100644 index 0000000..10280c3 --- /dev/null +++ b/crates/smtp/src/queue/journal.rs @@ -0,0 +1,159 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! inbuxa: journaling (journaling spec, JR-1 to JR-5, JR-11): the copy +//! taken as a message is queued, after DLP and transport rules, so it has +//! the envelope the message actually leaves or arrives with. + +use crate::queue::{FROM_AUTHENTICATED, FROM_AUTOGENERATED, FROM_DSN, FROM_REPORT, Message}; +use common::Server; +use inbuxa_features::{ + hold::Member, + journal::{ + self, Direction, + entries::{self, Entry}, + report::{self, Envelope, Recipient}, + }, + mailflow::held::HOLD_SECONDS, +}; +use store::write::{BatchBuilder, BlobLink, BlobOp, now}; +use types::blob_hash::BlobHash; + +/// Marks a journal report the server queued itself, so it's never +/// journaled (JR-2). Free in the message flags (the MAIL parameters use +/// the low bits, the sources bits 32 to 37). +pub const FROM_JOURNAL: u64 = 1 << 48; + +/// Journals `message`, whose queued bytes are `raw`, into every enabled +/// journal that takes it. An error means it may not have been journaled, +/// and the caller must not queue it. +pub async fn capture( + server: &Server, + queue_id: u64, + message: &Message, + raw: &[u8], +) -> trc::Result<()> { + if message.flags & (FROM_JOURNAL | FROM_REPORT) != 0 { + return Ok(()); + } + let journals = journal::enabled(server.store()).await?; + if journals.is_empty() { + return Ok(()); + } + + // Who's here on either side, and which way it goes + let mut members: Vec = Vec::new(); + let mut sender_local = message.flags & FROM_AUTHENTICATED != 0 + || (message.return_path.is_empty() && message.flags & (FROM_DSN | FROM_AUTOGENERATED) != 0); + if !message.return_path.is_empty() + && let Some(id) = server + .account_id_from_email(&message.return_path, false) + .await? + { + sender_local = true; + if let Some(member) = server.member_of(id).await { + members.push(member); + } + } + let (mut any_local, mut any_remote) = (false, false); + for rcpt in &message.recipients { + let address = rcpt.address.to_lowercase(); + let domain = address.rsplit_once('@').map_or("", |(_, d)| d); + let local_domain = server.domain(domain).await.ok().flatten().is_some(); + match server.account_id_from_email(&address, false).await? { + Some(id) => { + any_local = true; + if !members.iter().any(|m| m.account == id) + && let Some(member) = server.member_of(id).await + { + members.push(member); + } + } + None if local_domain => any_local = true, + None => any_remote = true, + } + } + let direction = Direction::of(sender_local, any_remote, any_local); + let taken: Vec<&journal::Journal> = journals + .iter() + .filter(|j| j.takes(direction, &members)) + .collect(); + if taken.is_empty() { + return Ok(()); + } + + // DLP holds a message by putting its release a century off + let at = now(); + let held = !message.recipients.is_empty() + && message + .recipients + .iter() + .all(|rcpt| rcpt.retry.due >= at + HOLD_SECONDS / 2); + let recipients: Vec = message + .recipients + .iter() + .map(|rcpt| Recipient { + address: rcpt.address.to_string(), + orcpt: rcpt.orcpt.as_deref().map(Into::into), + }) + .collect(); + let envelope = Envelope { + sender: &message.return_path, + authenticated: message.flags & FROM_AUTHENTICATED != 0, + recipients: &recipients, + queue_id, + received: message.created, + direction, + held, + }; + let host = server.core.network.server_name.as_str(); + let (bytes, fields) = report::build(&envelope, raw, &format!("postmaster@{host}"), host); + + // The report's blob, reserved until the entry links it + let hash = BlobHash::generate(&bytes); + let mut batch = BatchBuilder::new(); + batch.set( + BlobOp::Link { + hash: hash.clone(), + to: BlobLink::Temporary { until: at + 120 }, + }, + vec![], + ); + server.store().write(batch.build_all()).await?; + server + .blob_store() + .put_blob(hash.as_slice(), &bytes, server.core.email.compression) + .await?; + + let retention_days = taken + .iter() + .map(|j| j.retention_days) + .max() + .unwrap_or_default(); + let mut tenants: Vec = members.iter().filter_map(|m| m.tenant).collect(); + tenants.sort_unstable(); + tenants.dedup(); + let entry = Entry { + queue_id, + at, + direction, + sender: message.return_path.to_string(), + authenticated: envelope.authenticated, + recipients: recipients.iter().map(|r| r.address.clone()).collect(), + subject: fields.subject, + message_id: fields.message_id, + accounts: members.iter().map(|m| m.account).collect(), + tenants, + journals: taken.iter().map(|j| j.id).collect(), + held, + blob: entries::hex(hash.as_slice()), + size: bytes.len() as u64, + sha256: entries::sha256(&bytes), + expires_at: at + u64::from(retention_days) * 86_400, + }; + entries::append(server.store(), server.core.network.node_id, &entry).await?; + Ok(()) +} diff --git a/crates/smtp/src/queue/mod.rs b/crates/smtp/src/queue/mod.rs index 3e717bd..d31323c 100644 --- a/crates/smtp/src/queue/mod.rs +++ b/crates/smtp/src/queue/mod.rs @@ -24,6 +24,7 @@ use utils::DomainPart; pub mod dsn; pub mod held; // inbuxa: mail held for review +pub mod journal; // inbuxa: journaling pub mod manager; pub mod quota; pub mod spool; diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index c3e80ce..131aecb 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use super::{ @@ -453,6 +455,25 @@ impl MessageWrapper { return false; } + // inbuxa: journaling, JR-1: the copy is taken before the message is + // queued; if it can't be, the message isn't queued either + if let Err(err) = crate::queue::journal::capture( + server, + self.queue_id, + &self.message, + message.as_ref(), + ) + .await + { + trc::error!( + err.details("Failed to journal a message.") + .span_id(session_id) + .caused_by(trc::location!()) + ); + + return false; + } + trc::event!( Queue(event), SpanId = session_id, diff --git a/docs/spec/features/journaling.md b/docs/spec/features/journaling.md index d667f20..2754199 100644 --- a/docs/spec/features/journaling.md +++ b/docs/spec/features/journaling.md @@ -261,6 +261,38 @@ read, export. Each phase is its own PR with tests; releases as John decides. Like DLP, it stays out of production until John says. +## As built + +Phase 2 (`feature/journal-capture`), where it differs from the design or +fills in what it left open: + +- **The chain** is the journal's own (`crates/features/src/journal/ + entries.rs`), not the audit log's code shared. Entries expire out of chain + order (each keeps its journal's retention, and holds keep some longer), so + a link names its entry by SHA-256 instead of holding it: purging removes + the entry, its indexes and its report's blob link, and writes a purge + marker; the link stays. An entry missing without a marker is a broken + chain. Purged links at a chain's start are cleared and a floor recorded, + as the audit log does. +- **If the copy can't be taken**, the message isn't queued: the sender gets + a temporary failure and tries again. Nothing leaves unjournaled. +- **The report** says `Authenticated: yes|no` instead of the signed-in + account (the queue doesn't keep which account it was). `Added by rule` + comes with **Journal it** in phase 3. A recipient given with an ORCPT + that names another address counts as expanded from that address. +- **Journal reports** the server queues carry message flag bit 48 + (`FROM_JOURNAL`); an older version ignores the bit. +- **Permissions 680–683**: administrators get `sysJournalGet`/`Update`; the + Compliance Officer gets `Get`, `Search` and `Export`. So that an + administrator can still appoint an officer (and grant reading as settled + answer 5 describes), whoever holds `sysJournalUpdate` may grant `Search` + and `Export` without holding them; the role change is in the audit log. +- **`inbuxa:JournalEntry`** (get, query) and **Check the journal** over + JMAP come in phase 4 with search, so every read is audited from the first + version that allows one. Phase 2 has `inbuxa:Journal` only. +- **Outside archives** (a journal's destination) come in phase 3; every + journal writes to the built-in journal until then. + ## Known gaps - A message a person saves to Sent over IMAP, or sends through another diff --git a/resources/privacy/catalog.toml b/resources/privacy/catalog.toml index b6a6b9a..7461873 100644 --- a/resources/privacy/catalog.toml +++ b/resources/privacy/catalog.toml @@ -113,6 +113,18 @@ exceptions = ["contact", "content"] actions = ["contact", "content"] createdBy = ["identifier"] +[object."inbuxa:Journal"] +file = "inbuxa_journal.rs" +default = "none" +whose = ["administrator"] +where = ["data-store"] +scope = "server" +retention = "unbounded" +[object."inbuxa:Journal".properties] +name = ["content"] +description = ["content"] +createdBy = ["identifier"] + [object."inbuxa:LegalHold"] file = "inbuxa_legal_hold.rs" default = "none" @@ -408,6 +420,19 @@ captures = ["x:Email.maxMaskedAddresses"] leaves_host = false written_by = ["crates/features/src/masked_email/data.rs"] +# Journaling (journaling spec, JR-5, JR-14): a copy of each message a +# journal takes, with its envelope, kept for the journal's retention even +# after the account is deleted, and longer while a legal hold covers +# someone on it. +[source."journal"] +categories = ["content", "identifier", "contact", "metadata"] +whose = ["holder", "correspondent"] +where = ["data-store", "blob-store"] +scope = "server" +retention = { setting = "inbuxa:Journal.retentionDays" } +leaves_host = false +written_by = ["crates/features/src/journal/entries.rs", "crates/smtp/src/queue/journal.rs"] + [source."outbound-reports"] categories = ["network", "identifier", "content"] whose = ["correspondent"] diff --git a/resources/schema/schema.json.gz b/resources/schema/schema.json.gz index 85a0a52..91d98fc 100644 Binary files a/resources/schema/schema.json.gz and b/resources/schema/schema.json.gz differ diff --git a/resources/schema/schema.json.sha256 b/resources/schema/schema.json.sha256 index dfb1703..8c25d1c 100644 --- a/resources/schema/schema.json.sha256 +++ b/resources/schema/schema.json.sha256 @@ -1 +1 @@ -bLr7o49yK6CtWpENWq5zjgKrbCSzAXW3DZfNAEpnXEE \ No newline at end of file +JGXu3u5TIsHK6jNasyQy1wVxK-2Ku-PSa5NgiJP61q0 \ No newline at end of file diff --git a/tests/src/system/journal.rs b/tests/src/system/journal.rs new file mode 100644 index 0000000..3377d52 --- /dev/null +++ b/tests/src/system/journal.rs @@ -0,0 +1,471 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Journaling (journaling spec, phase 2): journals over JMAP, the copy +//! taken as mail is queued with its whole envelope, the report around the +//! untouched message, retention, purge, and a chain that shows tampering. + +use crate::utils::{ + account::Account, + server::{TestServer, TestServerBuilder}, + smtp::SmtpConnection, +}; +use inbuxa_features::journal::{ + Direction, + entries::{self, Entry, EntryId}, + report, +}; +use registry::schema::structs::{Expression, MtaStageAuth}; +use serde_json::{Value, json}; +use store::{Deserialize, write::BatchBuilder}; + +const USING: &[&str] = &[ + "urn:ietf:params:jmap:core", + "urn:ietf:params:jmap:mail", + "urn:ietf:params:jmap:submission", + "urn:inbuxa:jmap", +]; + +async fn call(account: &Account, method: &str, mut arguments: Value) -> (String, Value) { + if arguments.get("accountId").is_none() { + arguments["accountId"] = account.id_string().into(); + } + let response = account + .jmap_request(USING, json!([[method, arguments, "0"]])) + .await; + let call = response + .0 + .pointer("/methodResponses/0") + .cloned() + .unwrap_or_else(|| panic!("{method}: {}", response.0)); + ( + call[0].as_str().unwrap_or_default().to_string(), + call[1].clone(), + ) +} + +/// Sends a message whose headers name `to`, to the envelope `rcpt_to`. +async fn send( + sender: &Account, + identity: &str, + mailbox: &str, + to: &[&str], + rcpt_to: &[&str], + subject: &str, +) -> Value { + let (_, response) = call( + sender, + "Email/set", + json!({"create": {"e": { + "mailboxIds": {mailbox: true}, + "from": [{"email": sender.name()}], + "to": to.iter().map(|a| json!({"email": a})).collect::>(), + "subject": subject, + "bodyValues": {"b": {"value": "The body."}}, + "textBody": [{"partId": "b", "type": "text/plain"}] + }}}), + ) + .await; + let email = response["created"]["e"]["id"] + .as_str() + .unwrap_or_else(|| panic!("draft: {response}")) + .to_string(); + call( + sender, + "EmailSubmission/set", + json!({"create": {"s": { + "emailId": email, + "identityId": identity, + "envelope": { + "mailFrom": {"email": sender.name()}, + "rcptTo": rcpt_to.iter().map(|a| json!({"email": a})).collect::>() + } + }}}), + ) + .await + .1 +} + +async fn all_entries(test: &TestServer) -> Vec<(EntryId, Entry)> { + entries::list(test.server.store(), 0, u64::MAX, 10_000) + .await + .unwrap() +} + +async fn entry_for(test: &TestServer, subject: &str) -> Option<(EntryId, Entry)> { + all_entries(test) + .await + .into_iter() + .find(|(_, e)| e.subject == subject) +} + +async fn report_of(test: &TestServer, entry: &Entry) -> Vec { + let hash = entry.blob_hash().expect("blob hash"); + test.server + .blob_store() + .get_blob(hash.as_slice(), 0..usize::MAX) + .await + .unwrap() + .expect("report blob") +} + +pub async fn test(test: &mut TestServer) { + println!("Running journaling tests..."); + let admin = test.account("admin@example.com"); + let sender = admin + .create_user_account( + "journal-sender@example.com", + "journal-sender-secret-7101", + "Journal sender", + &[], + vec![], + ) + .await; + let other = admin + .create_user_account( + "journal-other@example.com", + "journal-other-secret-7102", + "Journal other", + &[], + vec![], + ) + .await; + let (_, response) = call( + &sender, + "Identity/set", + json!({"create": {"i": {"name": "Sender", "email": "journal-sender@example.com"}}}), + ) + .await; + let identity = response["created"]["i"]["id"].as_str().unwrap().to_string(); + let (_, response) = call( + &sender, + "Mailbox/set", + json!({"create": {"m": {"name": "Journal drafts"}}}), + ) + .await; + let mailbox = response["created"]["m"]["id"].as_str().unwrap().to_string(); + + // Nothing is journaled while there are no journals + let response = send( + &sender, + &identity, + &mailbox, + &["journal-other@example.com"], + &["journal-other@example.com"], + "Before any journal", + ) + .await; + assert!(response["created"].get("s").is_some(), "{response}"); + assert!(all_entries(test).await.is_empty()); + + // Journals: checked when written, the server's own properties refused + let (_, response) = call( + &admin, + "inbuxa:Journal/set", + json!({"create": { + "short": {"name": "Short", "enabled": true, "direction": "any", + "scope": {"everyone": true}, "retentionDays": 29}, + "both": {"name": "Both", "enabled": true, "direction": "any", + "scope": {"everyone": true, "accounts": [sender.id_string()]}, + "retentionDays": 365}, + "none": {"name": "None", "enabled": true, "direction": "any", + "scope": {}, "retentionDays": 365}, + "server": {"name": "Mine", "enabled": true, "direction": "any", + "scope": {"everyone": true}, "retentionDays": 365, + "createdBy": "me"}, + "all": {"name": "Everything", "enabled": true, "direction": "any", + "scope": {"everyone": true}, "retentionDays": 365}, + "out": {"name": "Sender's outgoing", "enabled": true, "direction": "outgoing", + "scope": {"accounts": [sender.id_string()]}, "retentionDays": 3650} + }}), + ) + .await; + for refused in ["short", "both", "none", "server"] { + assert_eq!( + response["notCreated"][refused]["type"], "invalidProperties", + "{refused}: {response}" + ); + } + assert_eq!( + response["notCreated"]["short"]["properties"], + json!(["retentionDays"]) + ); + assert_eq!( + response["notCreated"]["both"]["properties"], + json!(["scope"]) + ); + let everything = response["created"]["all"]["id"] + .as_str() + .unwrap_or_else(|| panic!("{response}")) + .to_string(); + let outgoing = response["created"]["out"]["id"] + .as_str() + .unwrap() + .to_string(); + let (_, response) = call(&admin, "inbuxa:Journal/get", json!({"ids": null})).await; + let list = response["list"].as_array().unwrap(); + assert_eq!(list.len(), 2, "{response}"); + assert_eq!(list[0]["name"], "Everything"); + assert_eq!(list[0]["createdBy"], "admin@example.com"); + assert_eq!(list[1]["scope"]["accounts"], json!([sender.id_string()])); + // Each node reads journals again within 30 seconds; this one at once + inbuxa_features::journal::invalidate(); + + // Internal mail with a Bcc recipient: one entry, the whole envelope + let response = send( + &sender, + &identity, + &mailbox, + &["journal-sender@example.com"], + &["journal-sender@example.com", "journal-other@example.com"], + "Internal with Bcc", + ) + .await; + assert!(response["created"].get("s").is_some(), "{response}"); + let (_, entry) = entry_for(test, "Internal with Bcc") + .await + .expect("journaled"); + assert_eq!(entry.direction, Direction::Internal); + assert_eq!(entry.sender, "journal-sender@example.com"); + assert!(entry.authenticated); + assert_eq!(entry.recipients.len(), 2, "{entry:?}"); + assert_eq!( + entry.journals.len(), + 1, + "internal isn't outgoing: {entry:?}" + ); + assert!(!entry.held); + assert_eq!(entry.expires_at, entry.at + 365 * 86_400); + let bytes = report_of(test, &entry).await; + assert_eq!(entries::sha256(&bytes), entry.sha256); + let text = String::from_utf8_lossy(&bytes); + assert!(text.contains("Direction: internal\r\n"), "{text}"); + assert!( + text.contains("To: journal-sender@example.com\r\n"), + "{text}" + ); + assert!( + text.contains("Bcc: journal-other@example.com\r\n"), + "{text}" + ); + let original = report::original(&bytes).expect("original part"); + let original = String::from_utf8_lossy(original); + assert!( + original.contains("Subject: Internal with Bcc"), + "{original}" + ); + assert!(original.contains("The body."), "{original}"); + assert!(!original.contains("Bcc:"), "the original is as sent"); + + // Outgoing: both journals take it, and it's kept for the longer + let response = send( + &sender, + &identity, + &mailbox, + &["someone@elsewhere.org"], + &["someone@elsewhere.org"], + "Leaving", + ) + .await; + assert!(response["created"].get("s").is_some(), "{response}"); + let (leaving_id, entry) = entry_for(test, "Leaving").await.expect("journaled"); + assert_eq!(entry.direction, Direction::Outgoing); + assert_eq!(entry.journals.len(), 2, "{entry:?}"); + assert_eq!(entry.expires_at, entry.at + 3650 * 86_400); + + // Incoming from outside + admin + .registry_create_object(MtaStageAuth { + require: Expression { + else_: "false".to_string(), + ..Default::default() + }, + ..Default::default() + }) + .await; + let mut lmtp = SmtpConnection::connect().await; + lmtp.ingest( + "someone@elsewhere.org", + &["journal-other@example.com"], + "From: someone@elsewhere.org\r\nTo: journal-other@example.com\r\nSubject: Arriving\r\n\r\nHi.\r\n", + ) + .await; + let (_, entry) = entry_for(test, "Arriving").await.expect("journaled"); + assert_eq!(entry.direction, Direction::Incoming); + assert!(!entry.authenticated); + assert_eq!(entry.accounts, vec![other.id().document_id()]); + + // The chain checks out, reports included + let store = test.server.store(); + let blobs = test.server.blob_store(); + let reports = entries::verify(store, Some(blobs)).await.unwrap(); + assert!(reports.iter().all(|r| r.broken_at.is_none()), "{reports:?}"); + let journaled = all_entries(test).await.len() as u64; + assert!(reports.iter().map(|r| r.entries).sum::() >= journaled); + + // An entry changed in the store shows; put back, it checks out again + let key = entries::content_key(leaving_id); + let stored = store + .get_value::(key.clone()) + .await + .unwrap() + .expect("stored entry") + .0; + let mut forged: Entry = serde_json::from_slice(&stored).unwrap(); + forged.recipients = vec!["nobody@elsewhere.org".into()]; + let mut batch = BatchBuilder::new(); + batch.set(key.class.clone(), serde_json::to_vec(&forged).unwrap()); + store.write(batch.build_all()).await.unwrap(); + let reports = entries::verify(store, None).await.unwrap(); + let broken = reports + .iter() + .find(|r| r.broken_at.is_some()) + .expect("broken"); + assert_eq!( + broken.broken_at.as_deref(), + Some(leaving_id.to_string().as_str()) + ); + assert!( + broken + .reason + .as_deref() + .unwrap_or_default() + .contains("changed") + ); + let mut batch = BatchBuilder::new(); + batch.set(key.class.clone(), stored.clone()); + store.write(batch.build_all()).await.unwrap(); + assert!( + entries::verify(store, None) + .await + .unwrap() + .iter() + .all(|r| r.broken_at.is_none()) + ); + + // An entry removed without a purge shows too + let mut batch = BatchBuilder::new(); + batch.clear(key.class.clone()); + store.write(batch.build_all()).await.unwrap(); + let reports = entries::verify(store, None).await.unwrap(); + assert!( + reports.iter().any(|r| r + .reason + .as_deref() + .unwrap_or_default() + .contains("before its time")), + "{reports:?}" + ); + let mut batch = BatchBuilder::new(); + batch.set(key.class.clone(), stored); + store.write(batch.build_all()).await.unwrap(); + + // Retention: nothing is due yet; a year on, what's kept for a hold + // stays, the rest goes, and the chain still checks out + let now = store::write::now(); + let purged = entries::purge(store, now, |_| false).await.unwrap(); + assert_eq!(purged.removed, 0); + let sender_id = sender.id().document_id(); + let later = now + 400 * 86_400; + let purged = entries::purge(store, later, |e| e.accounts.contains(&sender_id)) + .await + .unwrap(); + assert!(purged.removed >= 1, "{purged:?}"); + assert!(purged.kept_for_hold >= 1, "{purged:?}"); + assert!(entry_for(test, "Arriving").await.is_none(), "purged"); + assert!(entry_for(test, "Internal with Bcc").await.is_some(), "held"); + assert!(entry_for(test, "Leaving").await.is_some(), "ten years"); + let reports = entries::verify(store, Some(blobs)).await.unwrap(); + assert!(reports.iter().all(|r| r.broken_at.is_none()), "{reports:?}"); + assert!(reports.iter().map(|r| r.purged).sum::() >= 1); + + // Once the hold is gone the held entry goes too + let purged = entries::purge(store, later, |_| false).await.unwrap(); + assert!(purged.removed >= 1, "{purged:?}"); + assert!(entry_for(test, "Internal with Bcc").await.is_none()); + assert!( + entries::verify(store, Some(blobs)) + .await + .unwrap() + .iter() + .all(|r| r.broken_at.is_none()) + ); + + // Changing a journal's retention doesn't touch what it has taken + let before = entry_for(test, "Leaving").await.unwrap().1.expires_at; + let (_, response) = call( + &admin, + "inbuxa:Journal/set", + json!({"update": {outgoing.clone(): {"retentionDays": 30}}}), + ) + .await; + assert!(response["updated"].get(&outgoing).is_some(), "{response}"); + assert_eq!( + entry_for(test, "Leaving").await.unwrap().1.expires_at, + before + ); + + // Journals turned off or removed take nothing more; entries stay + let (_, response) = call( + &admin, + "inbuxa:Journal/set", + json!({"update": {everything.clone(): {"enabled": false}}, "destroy": [outgoing]}), + ) + .await; + assert!(response["updated"].get(&everything).is_some(), "{response}"); + assert_eq!(response["destroyed"].as_array().map(|d| d.len()), Some(1)); + inbuxa_features::journal::invalidate(); + let count = all_entries(test).await.len(); + let response = send( + &sender, + &identity, + &mailbox, + &["journal-other@example.com"], + &["journal-other@example.com"], + "After the journals", + ) + .await; + assert!(response["created"].get("s").is_some(), "{response}"); + assert_eq!(all_entries(test).await.len(), count); + assert!(entry_for(test, "Leaving").await.is_some()); + + // Every change to a journal is in the audit log + let (_, response) = call( + &admin, + "inbuxa:AuditEvent/query", + json!({"filter": {"targetKind": "inbuxa:Journal"}}), + ) + .await; + assert!( + response["ids"].as_array().map_or(0, |ids| ids.len()) >= 4, + "{response}" + ); +} + +struct Raw(Vec); + +impl Deserialize for Raw { + fn deserialize(bytes: &[u8]) -> trc::Result { + Ok(Raw(bytes.to_vec())) + } +} + +#[ignore] +#[tokio::test(flavor = "multi_thread")] +pub async fn journal_tests() { + let mut test = TestServerBuilder::new("journal_tests") + .await + .with_default_listeners() + .await + .build() + .await; + let admin = test.create_admin_account("admin@example.com").await; + test.insert_account(admin); + self::test(&mut test).await; + if test.is_reset() { + test.temp_dir.delete(); + } +} diff --git a/tests/src/system/mod.rs b/tests/src/system/mod.rs index 4c51b14..b181c64 100644 --- a/tests/src/system/mod.rs +++ b/tests/src/system/mod.rs @@ -15,6 +15,7 @@ pub mod account_lock; // inbuxa: account lock with delegation pub mod legal_hold; // inbuxa: legal hold pub mod compliance; // inbuxa: the compliance roles pub mod mail_rules; // inbuxa: DLP and mail flow rules +pub mod journal; // inbuxa: journaling pub mod audit; // inbuxa: the audit log pub mod authorization; pub mod auto_reload; // inbuxa: registry writes apply at once