Compare commits

..
6 Commits
Author SHA1 Message Date
jcoffey-dev eea96e8674 Security to-do list: accepted items, kept on the server
ci / fork-checks (pull_request) Successful in 19s
ci / build (pull_request) Successful in 8m5s
2026-09-28 21:22:43 -07:00
jcoffey-dev 4c5583e725 Merge pull request 'Journaling: capture at the queue, the built-in journal, retention' (#115) from feature/journal-capture into main
ci / fork-checks (push) Successful in 46s
ci / build (push) Canceled after 22m4s
2026-09-29 04:11:40 +00:00
jcoffey-dev 441ad0b18e Journaling: capture at the queue, the built-in journal, retention
ci / fork-checks (pull_request) Successful in 41s
ci / build (pull_request) Successful in 8m2s
Phase 2 of the journaling spec.

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

tests/src/system/journal.rs: validation, internal mail with a Bcc,
outgoing into two journals, incoming over LMTP, the report and its
original, tamper and early removal caught, hold-aware purge, retention
changes leave entries alone, disabled and removed journals take nothing.
2026-09-28 20:46:04 -07:00
jcoffey-dev 792ff9d1ee Merge pull request 'Spec: journaling' (#113) from spec/journaling into main
ci / fork-checks (push) Successful in 48s
ci / build (push) Canceled after 54m15s
2026-09-29 03:17:22 +00:00
jcoffey-dev af49e94d97 Journaling spec: approved, with the answers
ci / fork-checks (pull_request) Successful in 32s
ci / build (pull_request) Successful in 3m56s
2026-09-28 20:13:15 -07:00
jcoffey-dev 37c00b609c Spec: journaling
ci / fork-checks (pull_request) Successful in 47s
ci / build (pull_request) Successful in 23m6s
2026-09-28 19:45:37 -07:00
33 changed files with 3147 additions and 6 deletions
+15
View File
@@ -165,6 +165,14 @@ impl AccessToken {
mut requested_permissions: Permissions,
) -> Result<(), Vec<Permission>> {
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 {
@@ -310,6 +318,13 @@ impl Default for DefaultPermissions {
| Permission::SysSecurityAccept => {
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
@@ -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`].
@@ -53,6 +53,8 @@ const ADMIN_GRANTS: &[Permission] = &[
Permission::SysDlpPolicyUpdate,
Permission::SysDlpReviewGet,
Permission::SysDlpReviewUpdate,
Permission::SysJournalGet,
Permission::SysJournalUpdate,
Permission::SysSecurityAccept,
];
@@ -63,6 +65,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
+689
View File
@@ -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<String>,
pub subject: String,
pub message_id: String,
/// The people here on either side, whose holds keep the entry.
pub accounts: Vec<u32>,
pub tenants: Vec<u32>,
/// The journals that took it.
pub journals: Vec<u32>,
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<BlobHash> {
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<u8> {
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<Self> {
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<u8>);
impl Deserialize for Raw {
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
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<ValueClass> {
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<ValueClass> {
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<Vec<u64>> {
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<Vec<u8>> {
(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<Option<Head>> {
data.get_value::<Head>(key(KIND_HEAD, &[node]))
.await
.caused_by(trc::location!())
}
async fn floor(data: &Store, node: u64) -> trc::Result<Floor> {
Ok(data
.get_value::<Json<Floor>>(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<Vec<u64>> {
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<EntryId> {
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<Option<Entry>> {
Ok(data
.get_value::<Json<Entry>>(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<Vec<(EntryId, Entry)>> {
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<Purged> {
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<u64> = 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::<Raw>(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<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
}
/// 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<Vec<ChainReport>> {
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::<Link>::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::<Raw>(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::<Entry>::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::<Raw>(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");
}
}
+431
View File
@@ -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<u32>,
#[serde(default, with = "jmap_ids")]
pub groups: Vec<u32>,
#[serde(default, with = "jmap_ids")]
pub domains: Vec<u32>,
#[serde(default, with = "jmap_ids")]
pub tenants: Vec<u32>,
}
impl Scope {
fn lists(&self) -> [&Vec<u32>; 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<String>) -> 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<T>(pub T);
impl<T: SerdeSerialize> Serialize for Json<T> {
fn serialize(&self) -> trc::Result<Vec<u8>> {
serde_json::to_vec(&self.0).map_err(|err| {
trc::StoreEvent::UnexpectedError
.into_err()
.details("Failed to serialize a journal record")
.reason(err)
})
}
}
impl<T: DeserializeOwned + Sync + Send> Deserialize for Json<T> {
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
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<ValueClass> {
ValueKey::from(class(id))
}
pub async fn get(data: &Store, id: u32) -> trc::Result<Option<Journal>> {
Ok(data
.get_value::<Json<Journal>>(key(id))
.await
.caused_by(trc::location!())?
.map(|Json(journal)| journal))
}
/// Every journal, oldest first.
pub async fn all(data: &Store) -> trc::Result<Vec<Journal>> {
let mut journals = Vec::new();
data.iterate(IterateParams::new(key(0), key(u32::MAX)), |_, value| {
if let Ok(Json(journal)) = Json::<Journal>::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<u32> {
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<Vec<Journal>>)>;
static CACHE: RwLock<Cached> = 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<Arc<Vec<Journal>>> {
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::<Vec<_>>(),
);
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<u32>) -> 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"])
);
}
}
+359
View File
@@ -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<String>,
}
/// 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<String>,
pub cc: Vec<String>,
/// Envelope recipients in neither To nor Cc, nor reached through a list.
pub bcc: Vec<String>,
/// A list's address, and its members among the recipients.
pub expanded: Vec<(String, Vec<String>)>,
}
/// 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::<String>()
.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<String> {
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<u8>, 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<u8> = 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: <journal.{:x}.{}@{}>\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: [email protected]\r\n\
To: Bank <[email protected]>\r\n\
Cc: [email protected]\r\n\
Subject: Q3 figures\r\n\
Message-ID: <[email protected]>\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: "[email protected]",
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("[email protected]", None),
rcpt("[email protected]", Some("rfc822;[email protected]")),
rcpt("[email protected]", None),
rcpt("[email protected]", Some("[email protected]")),
rcpt("[email protected]", Some("rfc822;[email protected]")),
];
let fields = fields(&envelope(&recipients), ORIGINAL);
assert_eq!(fields.subject, "Q3 figures");
assert_eq!(fields.message_id, "<[email protected]>");
assert_eq!(fields.to, vec!["[email protected]"]);
assert_eq!(fields.cc, vec!["[email protected]"]);
assert_eq!(fields.bcc, vec!["[email protected]"]);
assert_eq!(
fields.expanded,
vec![(
"[email protected]".to_string(),
vec![
"[email protected]".to_string(),
"[email protected]".to_string()
]
)]
);
}
#[test]
fn report_carries_the_original_untouched() {
let recipients = [
rcpt("[email protected]", None),
rcpt("[email protected]", None),
];
let (report, _) = build(
&envelope(&recipients),
ORIGINAL,
"[email protected]",
"mx.example.com",
);
let text = String::from_utf8_lossy(&report);
assert!(text.contains("Sender: [email protected]\r\n"));
assert!(text.contains("Bcc: [email protected]\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,
"[email protected]",
"mx.example.com",
);
assert_eq!(original(&report), Some(unterminated));
}
#[test]
fn values_stay_on_one_line() {
let recipients = [rcpt("[email protected]", None)];
let mut env = envelope(&recipients);
env.sender = "[email protected]\r\nBcc: [email protected]";
env.held = true;
let body = text(&env, &Fields::default());
assert_eq!(body.matches("\r\n").count(), body.lines().count());
assert!(body.contains("Sender: [email protected] Bcc: [email protected]\r\n"));
assert!(body.contains("Held for review: yes\r\n"));
}
#[test]
fn an_empty_sender_is_shown_as_such() {
let recipients = [rcpt("[email protected]", None)];
let mut env = envelope(&recipients);
env.sender = "";
assert!(text(&env, &Fields::default()).starts_with("Sender: <>\r\n"));
}
}
+1
View File
@@ -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;
+1 -1
View File
@@ -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;
@@ -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<Self> {
// 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<Self> {
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<Self, Self::Err> {
JournalProperty::parse(s).ok_or(())
}
}
impl Element for JournalValue {
type Property = JournalProperty;
fn try_parse<P>(key: &Key<'_, Self::Property>, value: &str) -> Option<Self> {
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<String>,
}
impl<'de> DeserializeArguments<'de> for JournalSetArguments {
fn deserialize_argument<A>(&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::<serde::de::IgnoredAny>()?;
}
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<Id> for JournalValue {
fn from(id: Id) -> Self {
JournalValue::Id(id)
}
}
impl JmapObjectId for JournalValue {
fn as_id(&self) -> Option<Id> {
match self {
JournalValue::Id(id) => Some(*id),
}
}
fn as_any_id(&self) -> Option<AnyId> {
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<Id> {
None
}
fn as_any_id(&self) -> Option<AnyId> {
None
}
fn as_id_ref(&self) -> Option<&str> {
None
}
fn try_set_id(&mut self, _: AnyId) -> bool {
false
}
}
+1
View File
@@ -31,6 +31,7 @@ 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_security_acceptance; // inbuxa: accepted security to-do items
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
+3
View File
@@ -91,6 +91,9 @@ impl Response<'_> {
GetResponseMethod::SecurityAcceptance(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)
}
@@ -56,6 +56,7 @@ impl Response<'_> {
GetRequestMethod::LegalHold(request) => request.resolve_references(self)?,
GetRequestMethod::MailRule(request) => request.resolve_references(self)?,
GetRequestMethod::SecurityAcceptance(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)?,
@@ -135,6 +136,9 @@ impl Response<'_> {
SetRequestMethod::SecurityAcceptance(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)?
}
+9 -1
View File
@@ -71,6 +71,8 @@ pub enum MethodObject {
// inbuxa: accepted security to-do items
SecurityAcceptance,
HeldMessage,
// inbuxa: journaling
Journal,
TenantProtocolPolicy,
}
@@ -112,7 +114,8 @@ impl MethodObject {
| MethodObject::HoldExport
| MethodObject::MailRule
| MethodObject::SecurityAcceptance
| MethodObject::HeldMessage => Capability::Inbuxa,
| MethodObject::HeldMessage
| MethodObject::Journal => Capability::Inbuxa,
MethodObject::ProtocolPolicy => Capability::Inbuxa,
MethodObject::TenantProtocolPolicy => Capability::Inbuxa,
}
@@ -312,6 +315,8 @@ impl MethodName {
(MethodFunction::Set, MethodObject::MailRule) => "inbuxa:MailRule/set",
(MethodFunction::Get, MethodObject::SecurityAcceptance) => "inbuxa:SecurityAcceptance/get",
(MethodFunction::Set, MethodObject::SecurityAcceptance) => "inbuxa:SecurityAcceptance/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",
@@ -472,6 +477,8 @@ impl MethodName {
"inbuxa:MailRule/set" => (MethodObject::MailRule, MethodFunction::Set),
"inbuxa:SecurityAcceptance/get" => (MethodObject::SecurityAcceptance, MethodFunction::Get),
"inbuxa:SecurityAcceptance/set" => (MethodObject::SecurityAcceptance, 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),
@@ -547,6 +554,7 @@ impl Display for MethodObject {
MethodObject::LegalHold => "inbuxa:LegalHold",
MethodObject::MailRule => "inbuxa:MailRule",
MethodObject::SecurityAcceptance => "inbuxa:SecurityAcceptance",
MethodObject::Journal => "inbuxa:Journal",
MethodObject::HeldMessage => "inbuxa:HeldMessage",
MethodObject::HoldExport => "inbuxa:HoldExport",
MethodObject::ProtocolPolicy => "inbuxa:ProtocolPolicy",
+2
View File
@@ -126,6 +126,7 @@ pub enum GetRequestMethod {
LegalHold(Box<GetRequest<crate::object::inbuxa_legal_hold::LegalHold>>),
MailRule(Box<GetRequest<crate::object::inbuxa_mail_rule::MailRule>>),
SecurityAcceptance(Box<GetRequest<crate::object::inbuxa_security_acceptance::SecurityAcceptance>>),
Journal(Box<GetRequest<crate::object::inbuxa_journal::Journal>>),
HeldMessage(Box<GetRequest<crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<GetRequest<crate::object::inbuxa_hold_export::HoldExport>>),
ProtocolPolicy(Box<GetRequest<crate::object::inbuxa_protocol_policy::ProtocolPolicy>>),
@@ -167,6 +168,7 @@ pub enum SetRequestMethod<'x> {
SecurityAcceptance(
Box<SetRequest<'x, crate::object::inbuxa_security_acceptance::SecurityAcceptance>>,
),
Journal(Box<SetRequest<'x, crate::object::inbuxa_journal::Journal>>),
HeldMessage(Box<SetRequest<'x, crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<SetRequest<'x, crate::object::inbuxa_hold_export::HoldExport>>),
ProtocolPolicy(Box<SetRequest<'x, crate::object::inbuxa_protocol_policy::ProtocolPolicy>>),
+15
View File
@@ -668,6 +668,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)),
+15
View File
@@ -113,6 +113,7 @@ pub enum GetResponseMethod {
LegalHold(GetResponse<crate::object::inbuxa_legal_hold::LegalHold>),
MailRule(GetResponse<crate::object::inbuxa_mail_rule::MailRule>),
SecurityAcceptance(GetResponse<crate::object::inbuxa_security_acceptance::SecurityAcceptance>),
Journal(GetResponse<crate::object::inbuxa_journal::Journal>),
HeldMessage(GetResponse<crate::object::inbuxa_held_message::HeldMessage>),
HoldExport(GetResponse<crate::object::inbuxa_hold_export::HoldExport>),
ProtocolPolicy(GetResponse<crate::object::inbuxa_protocol_policy::ProtocolPolicy>),
@@ -154,6 +155,7 @@ pub enum SetResponseMethod {
SecurityAcceptance(
Box<SetResponse<crate::object::inbuxa_security_acceptance::SecurityAcceptance>>,
),
Journal(Box<SetResponse<crate::object::inbuxa_journal::Journal>>),
HeldMessage(Box<SetResponse<crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<SetResponse<crate::object::inbuxa_hold_export::HoldExport>>),
Explanation(Box<SetResponse<crate::object::inbuxa_explanation::Explanation>>),
@@ -862,6 +864,19 @@ impl<'x> From<SetResponse<crate::object::inbuxa_mail_rule::MailRule>> for Respon
}
}
// inbuxa: journaling
impl<'x> From<GetResponse<crate::object::inbuxa_journal::Journal>> for ResponseMethod<'x> {
fn from(value: GetResponse<crate::object::inbuxa_journal::Journal>) -> Self {
ResponseMethod::Get(GetResponseMethod::Journal(value))
}
}
impl<'x> From<SetResponse<crate::object::inbuxa_journal::Journal>> for ResponseMethod<'x> {
fn from(value: SetResponse<crate::object::inbuxa_journal::Journal>) -> Self {
ResponseMethod::Set(SetResponseMethod::Journal(Box::new(value)))
}
}
impl<'x> From<GetResponse<crate::object::inbuxa_legal_hold::LegalHold>> for ResponseMethod<'x> {
fn from(value: GetResponse<crate::object::inbuxa_legal_hold::LegalHold>) -> Self {
ResponseMethod::Get(GetResponseMethod::LegalHold(value))
+11
View File
@@ -116,6 +116,8 @@ impl JmapAuthorization for AccessToken {
Permission::SysDlpPolicyGet
}
}
// inbuxa: journaling (JR-18)
GetRequestMethod::Journal(_) => Permission::SysJournalGet,
GetRequestMethod::HoldExport(_) => Permission::SysLegalHoldExport,
// inbuxa: accepted security items are read by whoever may
// see the server's security settings
@@ -291,6 +293,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: accepting a security to-do item, or removing
// an acceptance; nothing is ever edited
SetRequestMethod::SecurityAcceptance(s) => {
@@ -470,6 +480,7 @@ impl JmapAuthorization for AccessToken {
| MethodObject::MailRule
| MethodObject::SecurityAcceptance
| MethodObject::HeldMessage
| MethodObject::Journal
| MethodObject::ProtocolPolicy
| MethodObject::TenantProtocolPolicy => Permission::JmapEmailChanges,
// inbuxa: x:MaskedEmail/changes reads what /get reads
+24
View File
@@ -288,6 +288,9 @@ impl RequestHandler for Server {
SetResponseMethod::SecurityAcceptance(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);
}
@@ -522,6 +525,11 @@ impl RequestHandler for Server {
.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)?;
@@ -1018,6 +1026,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)
+1
View File
@@ -432,6 +432,7 @@ impl IntermediateChangesResponse {
| MethodObject::HoldExport
| MethodObject::MailRule
| MethodObject::SecurityAcceptance
| MethodObject::Journal
| MethodObject::HeldMessage
| MethodObject::ProtocolPolicy
| MethodObject::TenantProtocolPolicy
+273
View File
@@ -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<serde_json::Map<String, serde_json::Value>, SetError<P>> {
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<String, serde_json::Value>) -> Result<Stored, SetError<P>> {
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::<P>().unwrap_or(P::Name);
SetError::invalid_properties()
.with_property(property)
.with_description(invalid.reason)
})?;
Ok(journal)
}
fn journal_id(id: Id) -> Option<u32> {
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<Journal>,
) -> trc::Result<GetResponse<Journal>> {
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<SetResponse<Journal>> {
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(&current) {
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)
}
+1
View File
@@ -12,6 +12,7 @@ pub mod account_lock;
pub mod legal_hold;
pub mod mail_rule;
pub mod security_acceptance;
pub mod journal;
pub mod held_message;
pub mod dlp_settings;
pub mod hold_export;
+6 -1
View File
@@ -1755,8 +1755,13 @@ pub enum Permission {
SysDlpPolicyUpdate = 677,
SysDlpReviewGet = 678,
SysDlpReviewUpdate = 679,
// inbuxa: journaling
SysJournalGet = 680,
SysJournalUpdate = 681,
SysJournalSearch = 682,
SysJournalExport = 683,
// inbuxa: the security to-do list, accepting an item
SysSecurityAccept = 680,
SysSecurityAccept = 684,
SysAccountGet = 219,
SysAccountCreate = 220,
SysAccountUpdate = 221,
+14 -2
View File
@@ -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"sysSecurityAccept" => Permission::SysSecurityAccept,
b"sysAccountGet" => Permission::SysAccountGet,
b"sysAccountCreate" => Permission::SysAccountCreate,
@@ -7794,6 +7798,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::SysSecurityAccept => "sysSecurityAccept",
Permission::SysAccountGet => "sysAccountGet",
Permission::SysAccountCreate => "sysAccountCreate",
@@ -8484,7 +8492,11 @@ impl EnumImpl for Permission {
677 => Some(Permission::SysDlpPolicyUpdate),
678 => Some(Permission::SysDlpReviewGet),
679 => Some(Permission::SysDlpReviewUpdate),
680 => Some(Permission::SysSecurityAccept),
680 => Some(Permission::SysJournalGet),
681 => Some(Permission::SysJournalUpdate),
682 => Some(Permission::SysJournalSearch),
683 => Some(Permission::SysJournalExport),
684 => Some(Permission::SysSecurityAccept),
219 => Some(Permission::SysAccountGet),
220 => Some(Permission::SysAccountCreate),
221 => Some(Permission::SysAccountUpdate),
@@ -8929,7 +8941,7 @@ impl EnumImpl for Permission {
}
}
const COUNT: usize = 681;
const COUNT: usize = 685;
}
impl serde::Serialize for Permission {
@@ -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,
+159
View File
@@ -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<Member> = 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<Recipient> = 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<u32> = 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(())
}
+1
View File
@@ -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;
+21
View File
@@ -2,6 +2,8 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
*
* 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,
+324
View File
@@ -0,0 +1,324 @@
# Feature spec: journaling
Status: **approved 2026-09-28**, with the answers under [Settled](#settled).
Not a rebuild of an upstream feature, so it has no line in SPEC.md §4's table.
Rule IDs: **JR-**.
## Provenance
Written for the record SPEC.md §3 rule 3 asks for. Sources, and nothing else:
| Source | License | Used for |
|---|---|---|
| This repository at `94a3a76` (2026-09-28): `crates/smtp/src/inbound/data.rs`, `inbound/rcpt.rs`, `queue/spool.rs`, `outbound/delivery.rs`, `crates/common/src/network/mta.rs`, `crates/features/src/{hold,audit,mailflow,undelete}`, `crates/store/src/write/{mod,blob}.rs`, `crates/jmap/src/inbuxa/hold_export.rs` | AGPL-3.0-only | Where every message passes, what the envelope holds, how holds keep blobs, how the audit chain and hold export work |
| `inbuxa-drafts/queue/journaling.md` | Own | What John asked for, and the gaps to settle |
| The DLP and mail flow rules spec, the audit-hold-lock spec, the personal-data catalog spec | Own | Conditions, the audit log, legal holds, roles, the catalog check |
| RFC 5321, RFC 3461 (DSN, ORCPT), RFC 2046 (`message/rfc822`), RFC 5322 | Public | The envelope, the original recipient of an expanded list, the report's shape |
No Enterprise-only file or snippet was used, and no third-party journaling
product or report format was consulted: the journal report below is our own
layout of the SMTP envelope around the untouched message.
## What it is
A **journal** is a copy of each message the server handles, captured in
transit with its **envelope** (the real sender and every recipient, including
Bcc and the members of lists), kept where nobody can change or remove it
until its retention ends, or sent to an outside archive. It sits beside two
things that exist:
- **Legal hold** keeps what's in chosen mailboxes, including what their owners
delete. It starts when a hold is placed and can't see Bcc or what was sent
from a mailbox that no longer exists.
- **The audit log** records what people and the server did, never the mail.
A journal answers the question neither can: *what went through, to whom,
from the day it was turned on*.
**Out of scope**: journaling mail stored before it's turned on, files,
calendar and contacts, IMAP APPEND (a mail app saving to its own Sent folder
sends nothing), and mail a mail app sends through another server.
Nothing in code, docs, UI text or output claims the product meets a legal or
regulatory standard. The pages say what's captured, where it's kept and for
how long.
## 1. What exists today
Checked by reading the code at `94a3a76`:
| Need | Today |
|---|---|
| One place all mail passes | `MessageWrapper::queue` (`queue/spool.rs` ~L375). SMTP, JMAP submission (`jmap/src/submission/set.rs` builds a local session and runs `queue_message`), inbound mail, Sieve redirects and vacation replies, and DSNs all queue through it. Local and remote delivery both start from the queue. |
| The envelope | At queue time: `mail_from`, every `rcpt_to` (Bcc included), the authenticated account with its groups and tenant, the queue id. Lists are **already expanded** at RCPT (`rcpt_resolve` → `RcptResolution::Expand`, `inbound/rcpt.rs`); the list address survives as each member's ORCPT (`dsn_info`). |
| A copy out | Sieve at DATA, milters and MTA hooks can send one, but all run **before** DLP and transport rules, so they miss recipients the rules add, and a Sieve copy carries no envelope. |
| Keeping a blob nobody can delete | No "undeletable" flag. Blobs are content-addressed (can't be edited); a `BlobLink::Temporary { until }` keeps one until `until`. Legal hold uses `until` = year 9999. |
| A record nobody can quietly change | The audit log's per-node SHA-256 chain (`features/src/audit/log.rs`): each entry carries `prev`, the head is asserted on append, purge leaves a floor hash, `verify` walks it. |
| Export | Hold export (LH-12): a ZIP of `.eml` files, `manifest.csv` with a SHA-256 per file, `manifest.sha256`, capped at 2 GiB. |
| Conditions by sender, recipient, group, tenant | The mail flow engine (`features/src/mailflow/engine.rs`), at DATA. |
## 2. Design
### 2.1 Where the copy is taken (JR-1, JR-2)
**JR-1.** The journal is taken in `MessageWrapper::queue`, after the message
is spooled, behind one `// inbuxa:` marked block. That's after DLP and
transport rules, so the envelope is the one the message actually leaves or
arrives with, and it covers every path that queues mail.
**JR-2.** What isn't journaled: journal reports themselves (they carry a
queue flag, so a report to an outside archive can't journal itself), and the
server's own DMARC and TLS reports. DSNs and Sieve redirects and vacation
replies are journaled (question 4). A message **refused** at DATA (DLP block,
a transport rule's refusal) was never accepted and isn't journaled; the
audit log already records it. A message **held** for DLP review is journaled
when it's queued, which is when it's held, with the hold noted in the entry.
### 2.2 The journal report (JR-3, JR-4)
**JR-3.** Each copy is a **journal report**: a new message whose first part
is `text/plain`, one field a line:
```
Sender: alice@example.com
Signed in as: alice@example.com
Subject: Q3 figures
Message-ID: <…>
Queue ID: 1a2b3c…
Received: 2026-09-28T14:03:11Z
Direction: outgoing
To: bank@elsewhere.example
Cc: bob@example.com
Bcc: carol@example.com
Expanded: finance@example.com -> dan@example.com, erin@example.com
Held for review: yes
```
and whose second part is the message as queued, **byte for byte**, as
`message/rfc822`. `Bcc:` lists envelope recipients that aren't in the
message's To or Cc headers. `Expanded:` groups the members of a list under
the list address, from their ORCPT. Recipients a transport rule added say so
(`Added by rule: <name>`). The field names are fixed English (they're a
record, not interface text), so a script can read them.
**JR-4.** One report per queued message, with the whole envelope, whatever
the scope matched on (§2.4). A message to 40 recipients is one report, not
40.
### 2.3 Where reports go (JR-5 to JR-8)
Each journal has a **destination** (question 1):
**JR-5. The built-in journal.** Records under a new prefix `J` in
`SUBSPACE_INBUXA`: queue id, received time, direction, sender, recipients,
the tenant(s), which journal matched, the report's blob hash and size, its
SHA-256, and the time it may be purged. The report blob is kept by a
`BlobLink::Temporary { until }` set to the end of its retention. There is no
JMAP `set` or `destroy` for entries: nothing in the product changes or
removes one before its time.
**JR-6. The chain.** Each entry carries the SHA-256 of the entry before it,
one chain per node, the same construction as the audit log (and its code,
generalized rather than copied). The console's **Check the journal** walks
it, and every blob's hash against its entry, and says what it found. Someone
with the server's disks can still remove data, and the chain is how that
shows; the docs say exactly that, and don't say it can't happen.
**JR-7. An outside archive.** The report is queued to an address (the
archive's journal mailbox) like any mail, with the queue's retries. A report
the archive refuses permanently, or can't take within the queue's limit, goes
into the built-in journal instead and raises a warning on the Overview
(question 6). Delivery is by the ordinary queue, so TLS and routing settings
apply; a queue route can be chosen for it.
**JR-8. Both**: the built-in journal and an outside archive.
### 2.4 Which mail: journals and their scope (JR-9 to JR-11)
**JR-9.** A **journal** is a named object (`inbuxa:Journal`): on or off, a
destination, a retention, and a scope. The scope is who: **everyone**, or
senders and recipients in chosen **accounts, groups, domains or tenants**,
and which **direction**: outgoing, incoming, internal, any. A message is
journaled once per journal whose scope any sender or recipient is in; two
journals with the same destination never write the same message twice.
**JR-10.** **By what's in it**: a new mail flow rule action, **Journal it**,
names a journal. The rule's conditions (detectors, words, attachments,
headers) decide; the copy is still taken at queue time (the rule only marks
the message). This is the rule-based journaling the queue note called
premium; here it's one more action, not a separate tier (question 2).
**JR-11.** Scope is evaluated from the envelope and directory membership at
queue time (`Server::member_of`), no message parsing, so journaling
everything costs a lookup per recipient and one blob write per message.
### 2.5 Retention and legal hold (JR-12 to JR-14)
**JR-12.** Each journal has a retention in days (question 3). An entry keeps
the retention it was written with: shortening a journal's retention applies
to new entries only, so nobody can empty the journal by editing a number.
Lengthening it applies to new entries too, and the console says so.
**JR-13.** Purge runs in the daily maintenance, removes entries past their
time and drops their blob link, and leaves a floor hash so the chain still
verifies, as the audit log does. An entry whose sender or any recipient is
under a **legal hold** isn't purged while the hold lasts (`holds_on`, read
uncached, as holds are everywhere).
**JR-14.** Deleting an account doesn't remove its journal entries; they end
with their retention (question 7). The privacy catalog says so.
### 2.6 Search, reading, export (JR-15 to JR-17)
**JR-15.** **Management › Compliance › Journal**: search by sender,
recipient, date range, direction, subject words (from the report's header
fields, not the body: no full-text index of the journal in this version).
Results list the envelope; **Read…** opens the report.
**JR-16.** **Export** a search as a ZIP in the hold export's shape: the
reports as `.eml`, `manifest.csv` with the envelope columns and a SHA-256
per file, `manifest.sha256`, the same 2 GiB cap. Export runs as a task and
the result is a blob owned by the person who asked for it.
**JR-17.** Every search, read and export is in the audit log, with who and
the search terms; so is every change to a journal.
### 2.7 Permissions (JR-18)
**JR-18.** New permissions after the DLP set (680 onward):
`sysJournalGet` / `sysJournalUpdate` (see and change journals),
`sysJournalSearch` (search and read entries), `sysJournalExport`. Superuser
only, by default. The officer grant audience adds Get, Search and Export to
the Compliance Officer; administrators configure journals but don't read
them unless granted Search (question 5). Journals are server-level, with a
tenant scope, as DLP rules are; nobody in a tenant reaches them.
### 2.8 Privacy catalog
New objects get catalog entries (`resources/privacy/catalog.toml`):
`inbuxa:Journal` (none), `inbuxa:JournalEntry` (mail content and envelope,
kept for the journal's retention, access audited, not erased with the
account). `privacy-check.py` enforces it.
### 2.9 Mixed versions, clusters, rollback
- Entries and journals live in the shared data store; report blobs in the
blob store. A **node-local** blob store (FileSystem, or RocksDB/SQLite as
the blob store) on a cluster means a node's journal lives on that node;
the console warns when journaling is on and the blob store isn't shared.
- During a rolling upgrade a node on the old version doesn't journal. The
console says so when nodes report different versions; the release notes
say to turn journals on after every node is upgraded.
- Rollback: the new prefix and the queue flag are ignored by an older
version; nothing in the queue's archived format changes (the "journal
report" flag rides in the existing message flags if one is free, else in
a side key by queue id; checked in phase 2 before writing code).
### 2.10 Cost
One extra blob per journaled message (the report wraps the original, so it
doesn't share its hash), plus one small record. The console shows the
journal's size and growth per day on the journal page, from the entries.
## 3. Console
- **Management › Compliance › Journaling**: journals (name, scope,
destination, retention, on/off), described in words like DLP rules
("Journal all mail to and from Finance into the built-in journal, kept
7 years"); **Check the journal**.
- **Management › Compliance › Journal**: search, read, export.
- The mail flow rule editor gains **Journal it**.
- The Overview warns about undelivered outside reports (JR-7) and a
node-local blob store (§2.9).
## 4. Webmail
Nothing. People aren't told a message was journaled, as they aren't told
about legal hold; the docs say journaling exists and what it captures.
## 5. Tests
Unit: the report's fields (Bcc computed from headers, list expansion from
ORCPT, rule-added recipients), scope matching, retention arithmetic, the
chain. Integration (`tests/src/system/`): SMTP and JMAP sends, inbound
mail, internal mail, a list and a Bcc recipient, a DLP-held message, a
Sieve redirect; the report equals the queued bytes; no `set`/`destroy`;
shortening retention doesn't touch existing entries; a hold stops purge;
account deletion leaves entries; an outside archive that refuses falls back
to the built-in journal; export manifest hashes; audit records for search,
read, export.
## 6. Phases
1. This spec, approved.
2. Capture at the queue, the report, the built-in journal with its chain,
retention and purge, holds; `inbuxa:Journal` and `inbuxa:JournalEntry`;
catalog entries; tests.
3. Outside archive and the fallback; **Journal it** in mail flow rules.
4. Search, read and export (a task), audit records.
5. Console pages; docs; a row in `inbuxa-drafts/divergence-log.md`.
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
server, never reaches the queue.
- Mail stored before journaling is on isn't journaled (legal hold covers
mailboxes).
- Search reads envelope and header fields, not bodies.
- Group accounts (`GroupAccount`) resolve as one account, not members; their
mail is journaled under the group's address.
## Settled
John, 2026-09-28, all seven as recommended:
1. **Destinations**: the built-in journal, an outside archive by address, or
both, per journal (JR-5, JR-7, JR-8).
2. **Scope**: everyone, or chosen accounts, groups, domains and tenants by
direction, plus a **Journal it** rule action; no standard/premium split
(JR-9, JR-10).
3. **Retention**: no default; 30 days to 10 years, picked when a journal is
turned on; existing entries keep theirs (JR-12).
4. **Which mail**: everything queued, including DSNs, Sieve redirects and
vacation replies, except DMARC/TLS reports and journal reports (JR-2).
5. **Who reads it**: administrators configure; Compliance Officers search,
read and export; administrators read only if granted Search (JR-18).
6. **An outside archive that won't take a report**: kept in the built-in
journal, with a warning (JR-7).
7. **Deleted accounts**: journal entries stay until their retention ends,
and the catalog says so (JR-14).
+25
View File
@@ -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:SecurityAcceptance"]
file = "inbuxa_security_acceptance.rs"
default = "none"
@@ -419,6 +431,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"]
Binary file not shown.
+1 -1
View File
@@ -1 +1 @@
-Qr56jJ7QowRD0wYXigExwHILIn1MXCgO4TBSAUHk8w
0Ay1tm9k_D94V7sLFfK7pIxwt3QdfYK_VBNs-rLGjhs
+471
View File
@@ -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::<Vec<_>>(),
"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::<Vec<_>>()
}
}}}),
)
.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<u8> {
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("[email protected]");
let sender = admin
.create_user_account(
"[email protected]",
"journal-sender-secret-7101",
"Journal sender",
&[],
vec![],
)
.await;
let other = admin
.create_user_account(
"[email protected]",
"journal-other-secret-7102",
"Journal other",
&[],
vec![],
)
.await;
let (_, response) = call(
&sender,
"Identity/set",
json!({"create": {"i": {"name": "Sender", "email": "[email protected]"}}}),
)
.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,
&["[email protected]"],
&["[email protected]"],
"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"], "[email protected]");
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,
&["[email protected]"],
&["[email protected]", "[email protected]"],
"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, "[email protected]");
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: [email protected]\r\n"),
"{text}"
);
assert!(
text.contains("Bcc: [email protected]\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,
&["[email protected]"],
&["[email protected]"],
"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(
"[email protected]",
&["[email protected]"],
"From: [email protected]\r\nTo: [email protected]\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::<u64>() >= 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::<Raw>(key.clone())
.await
.unwrap()
.expect("stored entry")
.0;
let mut forged: Entry = serde_json::from_slice(&stored).unwrap();
forged.recipients = vec!["[email protected]".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::<u64>() >= 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,
&["[email protected]"],
&["[email protected]"],
"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<u8>);
impl Deserialize for Raw {
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
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("[email protected]").await;
test.insert_account(admin);
self::test(&mut test).await;
if test.is_reset() {
test.temp_dir.delete();
}
}
+1
View File
@@ -16,6 +16,7 @@ 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 security_acceptances; // inbuxa: accepted security to-do items
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