Journaling: capture at the queue, the built-in journal, retention
Phase 2 of the journaling spec. - A copy of each message is taken in MessageWrapper::queue, after DLP and transport rules, for every enabled journal that takes it (direction and scope: everyone, or accounts, groups, domains, tenants). If the copy can't be taken the message isn't queued (temporary failure). - The journal report: the envelope one field a line (sender, To, Cc, Bcc from the envelope, list members from their ORCPT, direction, held for review), then the queued message byte for byte as message/rfc822. - The built-in journal under J in the inbuxa subspace: one chain per node whose links name each entry by SHA-256, so entries can expire out of chain order; purge leaves a marker, and verify catches an entry changed or removed early and a report that doesn't match. - Retention per journal (30 to 3650 days); an entry keeps what it was written with. The daily maintenance purges what's due, keeping entries whose people a legal hold covers (deleted accounts a hold keeps too), and records the counts in the audit log. - inbuxa:Journal get/set, audited by the request layer. Permissions 680-683: administrators see and change journals; the Compliance Officer sees, searches and exports. Whoever changes journals may grant search and export without holding them, so officers can still be appointed. - Catalog entries (inbuxa:Journal, source "journal"); spec as-built notes. tests/src/system/journal.rs: validation, internal mail with a Bcc, outgoing into two journals, incoming over LMTP, the report and its original, tamper and early removal caught, hold-aware purge, retention changes leave entries alone, disabled and removed journals take nothing.
This commit is contained in:
@@ -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");
|
||||
}
|
||||
}
|
||||
@@ -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"])
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -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"));
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user