Journaling: outside archives, and Journal it in mail flow rules #116

Merged
jcoffey-dev merged 1 commits from feature/journal-archive into main 2026-09-29 04:33:59 +00:00
14 changed files with 893 additions and 73 deletions
+120
View File
@@ -0,0 +1,120 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! Reports on their way to an outside archive (JR-7). Keys, after `J`:
//!
//! - `o` + the report's queue id: what goes into the built-in journal if
//! the archive never takes the report, as JSON. Cleared once it's
//! delivered or kept.
//! - `w` + journal id (u32): how often that journal's archive didn't take a
//! report, and the last time and reason, for the console's warning.
use super::{FEATURE, Json, entries::Entry};
use serde::{Deserialize as SerdeDeserialize, Serialize as SerdeSerialize};
use store::{
SUBSPACE_INBUXA, Serialize, Store, ValueKey,
write::{AnyClass, BatchBuilder, ValueClass},
};
use trc::AddContext;
const KIND_PENDING: u8 = b'o';
const KIND_FAILURES: u8 = b'w';
/// A report queued to an archive.
#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)]
#[serde(rename_all = "camelCase")]
pub struct Pending {
pub address: String,
/// The entry, should the archive not take it: its own, with the
/// sending journals' retention, whatever else the built-in journal has.
pub entry: Entry,
}
/// How a journal's archive has been taking its reports.
#[derive(Debug, Clone, Default, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)]
#[serde(rename_all = "camelCase")]
pub struct Failures {
pub count: u64,
/// Seconds.
pub last_at: u64,
pub last_reason: String,
}
fn class(kind: u8, id: &[u8]) -> ValueClass {
let mut key = Vec::with_capacity(2 + id.len());
key.push(FEATURE);
key.push(kind);
key.extend_from_slice(id);
ValueClass::Any(AnyClass {
subspace: SUBSPACE_INBUXA,
key,
})
}
pub async fn set_pending(data: &Store, queue_id: u64, pending: &Pending) -> trc::Result<()> {
let mut batch = BatchBuilder::new();
batch.set(
class(KIND_PENDING, &queue_id.to_be_bytes()),
Json(pending).serialize()?,
);
data.write(batch.build_all())
.await
.caused_by(trc::location!())?;
Ok(())
}
pub async fn pending(data: &Store, queue_id: u64) -> trc::Result<Option<Pending>> {
Ok(data
.get_value::<Json<Pending>>(ValueKey::from(class(KIND_PENDING, &queue_id.to_be_bytes())))
.await
.caused_by(trc::location!())?
.map(|Json(pending)| pending))
}
pub async fn clear_pending(data: &Store, queue_id: u64) -> trc::Result<()> {
let mut batch = BatchBuilder::new();
batch.clear(class(KIND_PENDING, &queue_id.to_be_bytes()));
data.write(batch.build_all())
.await
.caused_by(trc::location!())?;
Ok(())
}
pub async fn failures(data: &Store, journal_id: u32) -> trc::Result<Failures> {
Ok(data
.get_value::<Json<Failures>>(ValueKey::from(class(
KIND_FAILURES,
&journal_id.to_be_bytes(),
)))
.await
.caused_by(trc::location!())?
.map(|Json(failures)| failures)
.unwrap_or_default())
}
/// Counts one report an archive didn't take, for each of `journals`.
pub async fn record_failure(
data: &Store,
journals: &[u32],
at: u64,
reason: &str,
) -> trc::Result<()> {
for journal_id in journals {
let mut failures = failures(data, *journal_id).await?;
failures.count += 1;
failures.last_at = at;
failures.last_reason = reason.chars().take(500).collect();
let mut batch = BatchBuilder::new();
batch.set(
class(KIND_FAILURES, &journal_id.to_be_bytes()),
Json(&failures).serialize()?,
);
data.write(batch.build_all())
.await
.caused_by(trc::location!())?;
}
Ok(())
}
+83 -2
View File
@@ -16,6 +16,7 @@
//! with `J`; journals are `j` + id (u32), as JSON. There are few, so they're
//! read whole.
pub mod archive;
pub mod entries;
pub mod report;
@@ -128,6 +129,12 @@ pub struct Journal {
/// How long an entry this journal writes is kept. An entry keeps the
/// retention it was written with (JR-12).
pub retention_days: u32,
/// Whether entries go into the built-in journal (JR-5).
#[serde(default = "yes")]
pub built_in: bool,
/// An outside archive's journal address, sent each report (JR-7).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub archive_address: Option<String>,
#[serde(default)]
pub created_by: String,
#[serde(default)]
@@ -164,13 +171,28 @@ impl Journal {
format!("Keep entries between {MIN_RETENTION_DAYS} and {MAX_RETENTION_DAYS} days."),
);
}
// Neither is a journal only rules send mail to (JR-10)
let chosen = self.scope.lists().iter().any(|list| !list.is_empty());
if self.scope.everyone == chosen {
if self.scope.everyone && chosen {
return invalid(
"scope",
"Journal everyone, or choose accounts, groups, domains or tenants; not both.",
);
}
if !self.built_in && self.archive_address.is_none() {
return invalid(
"builtIn",
"Keep entries in the built-in journal, send them to an archive, or both.",
);
}
if let Some(address) = &self.archive_address
&& !is_address(address)
{
return invalid(
"archiveAddress",
format!("\"{address}\" isn't an email address."),
);
}
if self.scope.lists().iter().any(|list| list.len() > MAX_LIST) {
return invalid("scope", format!("Choose at most {MAX_LIST} of each."));
}
@@ -179,6 +201,11 @@ impl Journal {
/// Whether this journal takes a message going `direction` with these
/// people here on either side.
/// Whether only rules send this journal mail (JR-10).
pub fn rules_only(&self) -> bool {
!self.scope.everyone && self.scope.lists().iter().all(|list| list.is_empty())
}
pub fn takes(&self, direction: Direction, members: &[Member]) -> bool {
self.enabled
&& self.direction.includes(direction)
@@ -186,6 +213,22 @@ impl Journal {
}
}
fn yes() -> bool {
true
}
/// An address an archive can be sent to: one `@`, something either side,
/// nothing that would break an envelope.
fn is_address(address: &str) -> bool {
address.len() <= 320
&& address.split_once('@').is_some_and(|(local, domain)| {
!local.is_empty() && domain.contains('.') && !domain.contains('@')
})
&& !address
.chars()
.any(|c| c.is_whitespace() || c.is_control() || matches!(c, '<' | '>' | ',' | ';'))
}
/// A value stored as JSON.
pub(crate) struct Json<T>(pub T);
@@ -347,6 +390,8 @@ mod tests {
direction: Direction::Any,
scope,
retention_days: 365,
built_in: true,
archive_address: None,
created_by: String::new(),
created_at: 0,
updated_at: 0,
@@ -372,7 +417,11 @@ mod tests {
.validate()
.is_ok()
);
assert!(journal(Scope::default()).validate().is_err());
// Nobody chosen: only rules send it mail
let rules_only = journal(Scope::default());
assert!(rules_only.validate().is_ok());
assert!(rules_only.rules_only());
assert!(!rules_only.takes(Direction::Any, &[member(3, vec![7])]));
let both = Scope {
everyone: true,
groups: vec![4],
@@ -381,6 +430,38 @@ mod tests {
assert_eq!(journal(both).validate().unwrap_err().property, "scope");
}
#[test]
fn destinations() {
let mut j = journal(Scope {
everyone: true,
..Default::default()
});
j.built_in = false;
assert_eq!(j.validate().unwrap_err().property, "builtIn");
j.archive_address = Some("[email protected]".into());
assert!(j.validate().is_ok());
for bad in [
"archive",
"a@b",
"a [email protected]",
"<[email protected]>",
"a@[email protected]",
] {
j.archive_address = Some(bad.into());
assert_eq!(
j.validate().unwrap_err().property,
"archiveAddress",
"{bad}"
);
}
// Stored before destinations existed: the built-in journal
let old: Journal = serde_json::from_str(
r#"{"name":"Old","direction":"any","scope":{"everyone":true},"retentionDays":30}"#,
)
.unwrap();
assert!(old.built_in && old.archive_address.is_none());
}
#[test]
fn retention_has_bounds() {
let mut j = journal(Scope {
+27 -1
View File
@@ -20,6 +20,8 @@ use sha2::{Digest, Sha256};
pub struct Recipient {
pub address: String,
pub orcpt: Option<String>,
/// The mail flow rule that added or redirected to it.
pub added_by: Option<String>,
}
/// What the queue knows about a message.
@@ -46,6 +48,8 @@ pub struct Fields {
pub bcc: Vec<String>,
/// A list's address, and its members among the recipients.
pub expanded: Vec<(String, Vec<String>)>,
/// A rule's name, and the recipients it added.
pub added: Vec<(String, Vec<String>)>,
}
/// One line's worth of a value: no line breaks, no control characters.
@@ -106,7 +110,12 @@ pub fn fields(envelope: &Envelope<'_>, original: &[u8]) -> Fields {
.as_deref()
.map(orcpt_address)
.filter(|via| !via.is_empty() && *via != address);
if header_to.contains(&address) {
if let Some(rule) = &rcpt.added_by {
match fields.added.iter_mut().find(|(name, _)| name == rule) {
Some((_, added)) => added.push(line(&rcpt.address)),
None => fields.added.push((line(rule), vec![line(&rcpt.address)])),
}
} else if header_to.contains(&address) {
fields.to.push(line(&rcpt.address));
} else if header_cc.contains(&address) {
fields.cc.push(line(&rcpt.address));
@@ -159,6 +168,9 @@ pub fn text(envelope: &Envelope<'_>, fields: &Fields) -> String {
for (list, members) in &fields.expanded {
field("Expanded", &format!("{list} -> {}", members.join(", ")));
}
for (rule, added) in &fields.added {
field("Added by rule", &format!("{rule} -> {}", added.join(", ")));
}
if envelope.held {
field("Held for review", "yes");
}
@@ -265,6 +277,7 @@ The figures.\r\n";
Recipient {
address: address.into(),
orcpt: orcpt.map(Into::into),
added_by: None,
}
}
@@ -337,6 +350,19 @@ The figures.\r\n";
assert_eq!(original(&report), Some(unterminated));
}
#[test]
fn rule_added_recipients_say_so() {
let mut copied = rcpt("[email protected]", None);
copied.added_by = Some("Copy finance".into());
let recipients = [rcpt("[email protected]", None), copied];
let env = envelope(&recipients);
let fields = fields(&env, ORIGINAL);
assert!(fields.bcc.is_empty(), "{fields:?}");
assert!(
text(&env, &fields).contains("Added by rule: Copy finance -> [email protected]\r\n")
);
}
#[test]
fn values_stay_on_one_line() {
let recipients = [rcpt("[email protected]", None)];
+73 -2
View File
@@ -95,6 +95,33 @@ pub(crate) mod jmap_ids {
}
}
/// One id in the same form.
pub(crate) mod jmap_id {
use serde::{Deserialize, Deserializer, Serializer, de::Error};
use std::str::FromStr;
use types::id::Id;
pub fn serialize<S: Serializer>(id: &u32, serializer: S) -> Result<S::Ok, S::Error> {
serializer.serialize_str(&Id::from(*id).to_string())
}
#[derive(Deserialize)]
#[serde(untagged)]
enum Either {
Text(String),
Number(u32),
}
pub fn deserialize<'de, D: Deserializer<'de>>(deserializer: D) -> Result<u32, D::Error> {
match Either::deserialize(deserializer)? {
Either::Number(n) => Ok(n),
Either::Text(text) => Id::from_str(&text)
.map(|id| id.document_id())
.map_err(|_| D::Error::custom(format!("\"{text}\" isn't an id"))),
}
}
}
/// A detector and the least it must find.
#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)]
#[serde(rename_all = "camelCase")]
@@ -222,6 +249,11 @@ pub enum Action {
Route {
queue: String,
},
/// Journaling spec, JR-10: a copy into this journal, whatever its scope.
Journal {
#[serde(with = "jmap_id")]
journal: u32,
},
// DLP actions
Block {
notice: String,
@@ -311,10 +343,16 @@ impl Rule {
if self.direction != Direction::Outgoing {
return Err(invalid("direction", "DLP rules check outgoing mail only."));
}
if dlp_actions != 1 || self.actions.len() != 1 {
// One of block, warn or hold; journaling may go with it
if dlp_actions != 1
|| self
.actions
.iter()
.any(|a| !a.is_dlp() && !matches!(a, Action::Journal { .. }))
{
return Err(invalid(
"actions",
"A DLP rule has exactly one action: block, warn or hold.",
"A DLP rule has exactly one action: block, warn or hold, and may also journal the message.",
));
}
}
@@ -478,6 +516,7 @@ fn validate_action(action: &Action) -> Result<(), String> {
}
Action::Refuse { text: t } => text(t, "refusal text"),
Action::Route { queue } => text(queue, "queue"),
Action::Journal { .. } => Ok(()),
Action::Block { notice } | Action::Warn { notice } | Action::Hold { notice, .. } => {
text(notice, "notice")
}
@@ -619,6 +658,38 @@ mod tests {
}
}
#[test]
fn journal_action_goes_with_either_kind() {
let hold = Action::Hold {
notice: "Held.".into(),
notify_sender: false,
};
let journal = Action::Journal { journal: 3 };
assert!(
rule(Kind::Dlp, vec![hold.clone(), journal.clone()])
.validate()
.is_ok()
);
assert!(rule(Kind::Dlp, vec![journal.clone()]).validate().is_err());
assert!(
rule(
Kind::Dlp,
vec![hold, Action::PrefixSubject { text: "x".into() }]
)
.validate()
.is_err()
);
assert!(
rule(Kind::Transport, vec![journal.clone()])
.validate()
.is_ok()
);
let json = serde_json::to_value(&journal).unwrap();
assert_eq!(json, serde_json::json!({"type": "journal", "journal": "d"}));
let back: Action = serde_json::from_value(json).unwrap();
assert_eq!(back, journal);
}
#[test]
fn wire_format() {
let json = r#"{"name":"Cards","kind":"dlp","direction":"outgoing",
@@ -31,6 +31,12 @@ pub enum JournalProperty {
Scope,
/// How long an entry is kept; each keeps what it was written with.
RetentionDays,
/// Whether entries go into the built-in journal.
BuiltIn,
/// An outside archive's journal address.
ArchiveAddress,
/// Reports the archive didn't take: how many, when and why last.
ArchiveFailures,
CreatedBy,
CreatedAt,
UpdatedAt,
@@ -59,6 +65,9 @@ impl Property for JournalProperty {
JournalProperty::Direction => "direction",
JournalProperty::Scope => "scope",
JournalProperty::RetentionDays => "retentionDays",
JournalProperty::BuiltIn => "builtIn",
JournalProperty::ArchiveAddress => "archiveAddress",
JournalProperty::ArchiveFailures => "archiveFailures",
JournalProperty::CreatedBy => "createdBy",
JournalProperty::CreatedAt => "createdAt",
JournalProperty::UpdatedAt => "updatedAt",
@@ -77,6 +86,9 @@ impl JournalProperty {
b"direction" => JournalProperty::Direction,
b"scope" => JournalProperty::Scope,
b"retentionDays" => JournalProperty::RetentionDays,
b"builtIn" => JournalProperty::BuiltIn,
b"archiveAddress" => JournalProperty::ArchiveAddress,
b"archiveFailures" => JournalProperty::ArchiveFailures,
b"createdBy" => JournalProperty::CreatedBy,
b"createdAt" => JournalProperty::CreatedAt,
b"updatedAt" => JournalProperty::UpdatedAt,
+56 -11
View File
@@ -11,7 +11,10 @@
//! 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 inbuxa_features::journal::{
self, Journal as Stored,
archive::{self, Failures},
};
use jmap_proto::{
error::set::SetError,
method::{
@@ -37,13 +40,22 @@ const ALL: &[P] = &[
P::Direction,
P::Scope,
P::RetentionDays,
P::BuiltIn,
P::ArchiveAddress,
P::ArchiveFailures,
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];
const SERVER_SET: &[P] = &[
P::Id,
P::ArchiveFailures,
P::CreatedBy,
P::CreatedAt,
P::UpdatedAt,
];
fn server_level(access_token: &AccessToken) -> trc::Result<()> {
if access_token.tenant_id().is_some() {
@@ -86,7 +98,7 @@ fn date(seconds: u64) -> JValue {
Value::Str(UTCDate::from_timestamp(seconds as i64).to_string().into())
}
fn to_value(journal: &Stored, properties: &[P]) -> JValue {
fn to_value(journal: &Stored, failures: &Failures, 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 {
@@ -94,6 +106,32 @@ fn to_value(journal: &Stored, properties: &[P]) -> JValue {
P::Id => Value::Element(JournalValue::Id(Id::from(journal.id))),
P::CreatedAt => date(journal.created_at),
P::UpdatedAt => date(journal.updated_at),
P::ArchiveAddress => journal
.archive_address
.as_ref()
.map_or(Value::Null, |a| Value::Str(a.clone().into())),
// JR-7: what the console warns about
P::ArchiveFailures => {
let mut out = Map::with_capacity(3);
out.insert_unchecked(Key::Borrowed("count"), Value::Number(failures.count.into()));
out.insert_unchecked(
Key::Borrowed("lastAt"),
if failures.count > 0 {
date(failures.last_at)
} else {
Value::Null
},
);
out.insert_unchecked(
Key::Borrowed("lastReason"),
if failures.count > 0 {
Value::Str(failures.last_reason.clone().into())
} else {
Value::Null
},
);
Value::Object(out)
}
other => json
.get(other.to_cow().as_ref())
.cloned()
@@ -162,21 +200,28 @@ pub async fn get(
not_found,
};
let journals = journal::all(server.store()).await?;
match ids {
None => {
response.list = journals
.iter()
.map(|journal| to_value(journal, &properties))
.collect()
}
let wanted: Vec<&Stored> = match ids {
None => journals.iter().collect(),
Some(ids) => {
let mut wanted = Vec::with_capacity(ids.len());
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)),
Some(journal) => wanted.push(journal),
None => response.push_not_found(id),
}
}
wanted
}
};
for journal in wanted {
let failures = if properties.contains(&P::ArchiveFailures) {
archive::failures(server.store(), journal.id).await?
} else {
Failures::default()
};
response
.list
.push(to_value(journal, &failures, &properties));
}
Ok(response)
}
+8
View File
@@ -102,6 +102,10 @@ pub struct SessionData {
pub dlp_refusal: Option<DlpRefusal>,
// inbuxa: a mail flow rule's route for this message
pub mailflow_queue: Option<String>,
// inbuxa: journaling (JR-3, JR-10): journals rules sent this message
// to, and recipients rules added, by rule name
pub journal_marks: Vec<u32>,
pub journal_added: Vec<(String, String)>,
}
/// inbuxa: a DATA refusal by DLP rules: blocked, or a warning the sender
@@ -189,6 +193,8 @@ impl SessionData {
dlp_override: None,
dlp_refusal: None,
mailflow_queue: None,
journal_marks: Vec::new(),
journal_added: Vec::new(),
}
}
}
@@ -315,6 +321,8 @@ impl SessionData {
dlp_override: None,
dlp_refusal: None,
mailflow_queue: None,
journal_marks: Vec::new(),
journal_added: Vec::new(),
}
}
}
+15 -3
View File
@@ -763,19 +763,27 @@ impl<T: SessionStream> Session<T> {
}
for change in envelope {
match change {
super::mailflow::EnvelopeChange::AddRecipient(address) => {
super::mailflow::EnvelopeChange::AddRecipient(address, rule) => {
if !self
.data
.rcpt_to
.iter()
.any(|r| r.address_lcase.eq_ignore_ascii_case(&address))
{
self.data.journal_added.push((address.to_lowercase(), rule));
self.data.rcpt_to.push(SessionAddress::new(address));
}
}
super::mailflow::EnvelopeChange::Redirect(addresses) => {
super::mailflow::EnvelopeChange::Redirect(addresses, rule) => {
self.data.journal_added = addresses
.iter()
.map(|a| (a.to_lowercase(), rule.clone()))
.collect();
self.data.rcpt_to = addresses.into_iter().map(SessionAddress::new).collect();
}
super::mailflow::EnvelopeChange::Journal(journal) => {
self.data.journal_marks.push(journal);
}
super::mailflow::EnvelopeChange::Route(queue) => {
self.data.mailflow_queue = Some(queue);
}
@@ -882,7 +890,11 @@ impl<T: SessionStream> Session<T> {
.with_dkim_signers(dkim_signers)
.with_original_raw_message(original_message)
.with_original_authenticated_message(auth_message)
.with_metadata(metadata),
.with_metadata(metadata)
.with_journal(
std::mem::take(&mut self.data.journal_marks),
std::mem::take(&mut self.data.journal_added),
),
)
.await
{
+28 -5
View File
@@ -63,9 +63,12 @@ pub struct HeldDraft {
/// What a transport rule changes about where a message goes.
pub enum EnvelopeChange {
AddRecipient(String),
Redirect(Vec<String>),
/// An address, and the rule that added it.
AddRecipient(String, String),
Redirect(Vec<String>, String),
Route(String),
/// Journaling spec, JR-10: a journal the message goes to.
Journal(u32),
}
/// `[override: reason]` at the start of a subject: the reason, and the
@@ -402,11 +405,17 @@ impl<T: SessionStream> Session<T> {
rewrite::prefix_subject(now, text)
}
RuleAction::AddRecipient { address } => {
changes.push(EnvelopeChange::AddRecipient(address.clone()));
changes.push(EnvelopeChange::AddRecipient(
address.clone(),
matched.name.clone(),
));
None
}
RuleAction::Redirect { addresses } => {
changes.push(EnvelopeChange::Redirect(addresses.clone()));
changes.push(EnvelopeChange::Redirect(
addresses.clone(),
matched.name.clone(),
));
None
}
RuleAction::Route { queue } => {
@@ -423,7 +432,8 @@ impl<T: SessionStream> Session<T> {
}
RuleAction::Block { .. }
| RuleAction::Warn { .. }
| RuleAction::Hold { .. } => None,
| RuleAction::Hold { .. }
| RuleAction::Journal { .. } => None,
};
if next.is_some() {
current = next;
@@ -450,6 +460,19 @@ impl<T: SessionStream> Session<T> {
.await;
}
}
// JR-10: journals any matched rule sends the message to,
// DLP rules included
for matched in &outcome.matched {
for action in &matched.actions {
if let RuleAction::Journal { journal } = action
&& !changes
.iter()
.any(|c| matches!(c, EnvelopeChange::Journal(j) if j == journal))
{
changes.push(EnvelopeChange::Journal(*journal));
}
}
}
match hold {
Some(draft) => Checked::Hold {
draft,
+4
View File
@@ -494,6 +494,10 @@ impl<T: AsyncWrite + AsyncRead + Unpin> Session<T> {
self.data.delivery_by = 0;
self.data.future_release = 0;
self.data.rcpt_oks = 0;
// inbuxa: what mail flow rules decided was for the last message only
self.data.mailflow_queue = None;
self.data.journal_marks.clear();
self.data.journal_added.clear();
}
pub fn reset_tls(&mut self) {
+165 -44
View File
@@ -4,16 +4,22 @@
* 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.
//! inbuxa: journaling (journaling spec, JR-1 to 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; reports to outside archives,
//! and what happens when an archive doesn't take one.
use crate::queue::{FROM_AUTHENTICATED, FROM_AUTOGENERATED, FROM_DSN, FROM_REPORT, Message};
use crate::queue::{
FROM_AUTHENTICATED, FROM_AUTOGENERATED, FROM_DSN, FROM_REPORT, Message, MessageSource, Status,
spool::{QueueParams, SmtpSpool},
};
use common::Server;
use inbuxa_features::{
audit::{Action, Actor, Outcome, Record, Target},
hold::Member,
journal::{
self, Direction,
archive::{self, Pending},
entries::{self, Entry},
report::{self, Envelope, Recipient},
},
@@ -22,6 +28,15 @@ use inbuxa_features::{
use store::write::{BatchBuilder, BlobLink, BlobOp, now};
use types::blob_hash::BlobHash;
/// What mail flow rules decided about a message at DATA (JR-3, JR-10).
#[derive(Debug, Clone, Default)]
pub struct Hints {
/// Journals a rule sent it to.
pub marks: Vec<u32>,
/// Recipients a rule added (lowercase), and the rule's name.
pub added: Vec<(String, String)>,
}
/// 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).
@@ -35,6 +50,7 @@ pub async fn capture(
queue_id: u64,
message: &Message,
raw: &[u8],
hints: &Hints,
) -> trc::Result<()> {
if message.flags & (FROM_JOURNAL | FROM_REPORT) != 0 {
return Ok(());
@@ -77,9 +93,10 @@ pub async fn capture(
}
}
let direction = Direction::of(sender_local, any_remote, any_local);
// A journal takes it through its scope, or because a rule sent it there
let taken: Vec<&journal::Journal> = journals
.iter()
.filter(|j| j.takes(direction, &members))
.filter(|j| j.takes(direction, &members) || hints.marks.contains(&j.id))
.collect();
if taken.is_empty() {
return Ok(());
@@ -98,6 +115,11 @@ pub async fn capture(
.map(|rcpt| Recipient {
address: rcpt.address.to_string(),
orcpt: rcpt.orcpt.as_deref().map(Into::into),
added_by: hints
.added
.iter()
.find(|(address, _)| address.eq_ignore_ascii_case(&rcpt.address))
.map(|(_, rule)| rule.clone()),
})
.collect();
let envelope = Envelope {
@@ -112,48 +134,147 @@ pub async fn capture(
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,
let hash = BlobHash::generate(&bytes);
let entry_for = |journals: &[&journal::Journal]| {
let retention_days = journals
.iter()
.map(|j| j.retention_days)
.max()
.unwrap_or_default();
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.clone(),
message_id: fields.message_id.clone(),
accounts: members.iter().map(|m| m.account).collect(),
tenants: tenants.clone(),
journals: journals.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?;
// The built-in journal: one entry, however many journals keep it there
let built_in: Vec<&journal::Journal> = taken.iter().copied().filter(|j| j.built_in).collect();
if !built_in.is_empty() {
// The report's blob, reserved until the entry links it
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?;
entries::append(
server.store(),
server.core.network.node_id,
&entry_for(&built_in),
)
.await?;
}
// Outside archives: one report per address (JR-4, JR-7)
let mut addresses: Vec<(String, Vec<&journal::Journal>)> = Vec::new();
for journal in &taken {
if let Some(address) = &journal.archive_address {
let address = address.to_lowercase();
match addresses.iter_mut().find(|(a, _)| *a == address) {
Some((_, journals)) => journals.push(journal),
None => addresses.push((address, vec![journal])),
}
}
}
for (address, journals) in addresses {
// From nobody: an archive's refusal comes back to no one, and the
// queue's own record of it is what counts (settle, below)
let mut report = server.new_message("", MessageSource::Autogenerated, 0);
report.message.flags |= FROM_JOURNAL;
report.add_expanded_recipient(&address, server).await;
let pending = Pending {
address,
entry: entry_for(&journals),
};
archive::set_pending(server.store(), report.queue_id, &pending).await?;
let report_id = report.queue_id;
// Boxed: queueing the report comes back through this function
let queued = Box::pin(report.queue(QueueParams::new(&bytes, 0, server))).await;
if !queued {
archive::clear_pending(server.store(), report_id).await?;
return Err(trc::StoreEvent::UnexpectedError
.into_err()
.details("Failed to queue a journal report"));
}
}
Ok(())
}
/// JR-7: a journal report is leaving the queue. Delivered, its pending
/// record goes; not delivered (refused, expired, or deleted from the
/// queue), it goes into the built-in journal instead, and the journals
/// that sent it count a failure. An error means nothing was settled, and
/// the report must stay queued.
pub async fn settle(server: &Server, queue_id: u64, message: &Message) -> trc::Result<()> {
let store = server.store();
let Some(pending) = archive::pending(store, queue_id).await? else {
return Ok(());
};
let delivered = !message.recipients.is_empty()
&& message
.recipients
.iter()
.all(|rcpt| matches!(rcpt.status, Status::Completed(_)));
if !delivered {
let reason = if message
.recipients
.iter()
.any(|rcpt| matches!(rcpt.status, Status::PermanentFailure(_)))
{
"the archive refused it"
} else {
"it wasn't delivered before leaving the queue"
};
entries::append(store, server.core.network.node_id, &pending.entry).await?;
let at = now();
archive::record_failure(store, &pending.entry.journals, at, reason).await?;
server
.audit_note(Record {
at: at * 1000,
actor: Actor::system("Journal"),
via: None,
remote_ip: None,
action: Action::Create,
target: Target {
kind: "inbuxa:JournalEntry".into(),
id: Some(format!("{:x}", pending.entry.queue_id)),
name: None,
account_id: None,
tenant_id: None,
},
changes: vec![],
details: Some(format!(
"A journal report to {} wasn't delivered ({reason}); kept in the built-in journal",
pending.address
)),
reason: None,
outcome: Outcome::success(),
})
.await;
}
archive::clear_pending(store, queue_id).await
}
+39
View File
@@ -370,6 +370,8 @@ pub(crate) struct QueueParams<'x, 'y> {
pub session_id: u64,
pub server: &'y Server,
pub train_spam: Option<(bool, String)>,
// inbuxa: journaling, JR-3, JR-10
pub journal: crate::queue::journal::Hints,
}
impl MessageWrapper {
@@ -390,6 +392,7 @@ impl MessageWrapper {
server,
train_spam,
metadata,
journal,
..
} = params;
let event = self.message.queued_event();
@@ -462,6 +465,7 @@ impl MessageWrapper {
self.queue_id,
&self.message,
message.as_ref(),
&journal,
)
.await
{
@@ -813,6 +817,20 @@ impl MessageWrapper {
}
pub async fn remove(self, server: &Server, prev_event: Option<u64>) -> bool {
// inbuxa: journaling, JR-7: a journal report the archive never took
// goes into the built-in journal before it leaves the queue
if self.message.flags & crate::queue::journal::FROM_JOURNAL != 0
&& let Err(err) =
crate::queue::journal::settle(server, self.queue_id, &self.message).await
{
trc::error!(
err.details("Failed to settle a journal report; it stays queued.")
.span_id(self.span_id)
.caused_by(trc::location!())
);
return false;
}
let mut batch = BatchBuilder::new();
if let Some(prev_event) = prev_event {
@@ -987,6 +1005,19 @@ impl MessageWrapper {
server: &Server,
prev_events: AHashMap<QueueName, u64>,
) -> bool {
// inbuxa: journaling, JR-7, as in `remove`
if self.message.flags & crate::queue::journal::FROM_JOURNAL != 0
&& let Err(err) =
crate::queue::journal::settle(server, self.queue_id, &self.message).await
{
trc::error!(
err.details("Failed to settle a journal report; it stays queued.")
.span_id(self.span_id)
.caused_by(trc::location!())
);
return false;
}
let mut batch = BatchBuilder::new();
for (queue_name, due) in prev_events {
@@ -1142,9 +1173,17 @@ impl<'x, 'y> QueueParams<'x, 'y> {
original_raw_message: None,
original_authenticated_message: None,
metadata: Vec::new(),
journal: Default::default(),
}
}
/// inbuxa: journals mail flow rules sent the message to, and the
/// recipients they added, by rule name.
pub fn with_journal(mut self, marks: Vec<u32>, added: Vec<(String, String)>) -> Self {
self.journal = crate::queue::journal::Hints { marks, added };
self
}
pub fn with_train_spam(mut self, train_spam: Option<(bool, String)>) -> Self {
self.train_spam = train_spam;
self
+22 -2
View File
@@ -290,8 +290,28 @@ fills in what it left open:
- **`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.
Phase 3 (`feature/journal-archive`):
- **Destinations** are two properties of a journal: `builtIn` (true for
journals stored before phase 3) and `archiveAddress`. At least one.
- **Journals only rules use**: a journal whose scope chooses nobody takes
only what a **Journal it** action sends it. The action goes on mail flow
rules, and on a DLP rule beside its block, warn or hold (a blocked
message isn't queued, so it isn't journaled).
- **Reports to an archive** are queued from the empty sender, so a refusal
comes back to no one; a pending record per report says what to keep.
When the queue lets go of a report without delivering it (refused,
expired, or deleted from the queue), the report becomes its own entry in
the built-in journal under the sending journals' retention, even when
another journal already kept the message there, the journal's
`archiveFailures` (count, last time, reason) goes up, and the audit log
records it. If that can't be written, the report stays queued.
- **Added by rule** lists recipients a transport rule added or redirected
to, by rule name, instead of counting them as Bcc.
- A rule's route (and now its journal marks) is cleared between messages
in one SMTP session; before, a second message in the same session kept
the first one's route.
## Known gaps
+241 -3
View File
@@ -20,6 +20,7 @@ use inbuxa_features::journal::{
};
use registry::schema::structs::{Expression, MtaStageAuth};
use serde_json::{Value, json};
use std::str::FromStr;
use store::{Deserialize, write::BatchBuilder};
const USING: &[&str] = &[
@@ -27,6 +28,7 @@ const USING: &[&str] = &[
"urn:ietf:params:jmap:mail",
"urn:ietf:params:jmap:submission",
"urn:inbuxa:jmap",
"urn:inbuxa:jmap:registry",
];
async fn call(account: &Account, method: &str, mut arguments: Value) -> (String, Value) {
@@ -171,8 +173,11 @@ pub async fn test(test: &mut TestServer) {
"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},
"none": {"name": "Nowhere", "enabled": true, "direction": "any",
"scope": {"everyone": true}, "retentionDays": 365, "builtIn": false},
"badaddr": {"name": "Bad archive", "enabled": true, "direction": "any",
"scope": {"everyone": true}, "retentionDays": 365,
"archiveAddress": "not an address"},
"server": {"name": "Mine", "enabled": true, "direction": "any",
"scope": {"everyone": true}, "retentionDays": 365,
"createdBy": "me"},
@@ -183,7 +188,7 @@ pub async fn test(test: &mut TestServer) {
}}),
)
.await;
for refused in ["short", "both", "none", "server"] {
for refused in ["short", "both", "none", "badaddr", "server"] {
assert_eq!(
response["notCreated"][refused]["type"], "invalidProperties",
"{refused}: {response}"
@@ -197,6 +202,14 @@ pub async fn test(test: &mut TestServer) {
response["notCreated"]["both"]["properties"],
json!(["scope"])
);
assert_eq!(
response["notCreated"]["none"]["properties"],
json!(["builtIn"])
);
assert_eq!(
response["notCreated"]["badaddr"]["properties"],
json!(["archiveAddress"])
);
let everything = response["created"]["all"]["id"]
.as_str()
.unwrap_or_else(|| panic!("{response}"))
@@ -445,6 +458,230 @@ pub async fn test(test: &mut TestServer) {
);
}
/// Phase 3: journals only rules send mail to, recipients a rule added,
/// reports sent to an outside archive, and what happens when the archive
/// doesn't take one.
pub async fn archive(test: &mut TestServer) {
println!("Running journal archive tests...");
let admin = test.account("[email protected]");
let sender = admin
.create_user_account(
"[email protected]",
"archive-sender-secret-7201",
"Archive sender",
&[],
vec![],
)
.await;
let vault = admin
.create_user_account(
"[email protected]",
"journal-vault-secret-7202",
"Journal vault",
&[],
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": "Archive drafts"}}}),
)
.await;
let mailbox = response["created"]["m"]["id"].as_str().unwrap().to_string();
let (_, response) = call(
&admin,
"inbuxa:Journal/set",
json!({"create": {
"rules": {"name": "Only what rules send", "enabled": true, "direction": "any",
"scope": {}, "retentionDays": 30},
"local": {"name": "To the vault", "enabled": true, "direction": "outgoing",
"scope": {"accounts": [sender.id_string()]}, "retentionDays": 30,
"builtIn": false, "archiveAddress": "[email protected]"},
"remote": {"name": "To an outside archive", "enabled": true, "direction": "internal",
"scope": {"accounts": [sender.id_string()]}, "retentionDays": 30,
"builtIn": false, "archiveAddress": "[email protected]"}
}}),
)
.await;
let id = |name: &str| {
response["created"][name]["id"]
.as_str()
.unwrap_or_else(|| panic!("{name}: {response}"))
.to_string()
};
let (rules_only, local, remote) = (id("rules"), id("local"), id("remote"));
let number = |id: &str| types::id::Id::from_str(id).unwrap().document_id();
let (_, response) = call(
&admin,
"inbuxa:MailRule/set",
json!({"create": {"r": {
"name": "Copy and journal", "kind": "transport", "direction": "outgoing",
"conditions": [{"type": "words", "words": ["journal-me"]}],
"actions": [
{"type": "addRecipient", "address": "[email protected]"},
{"type": "journal", "journal": rules_only.clone()}
]
}}}),
)
.await;
let rule = response["created"]["r"]["id"]
.as_str()
.unwrap_or_else(|| panic!("{response}"))
.to_string();
inbuxa_features::journal::invalidate();
// A rule sends it to a journal whose scope takes nobody, and says who
// it added
let response = send(
&sender,
&identity,
&mailbox,
&["[email protected]"],
&["[email protected]"],
"Marked journal-me",
)
.await;
assert!(response["created"].get("s").is_some(), "{response}");
let entry = all_entries(test)
.await
.into_iter()
.map(|(_, e)| e)
.find(|e| e.subject == "Marked journal-me" && e.journals.contains(&number(&rules_only)))
.expect("journaled by the rule");
assert_eq!(entry.journals, vec![number(&rules_only)], "{entry:?}");
let text = String::from_utf8_lossy(&report_of(test, &entry).await).into_owned();
assert!(
text.contains("Added by rule: Copy and journal -> [email protected]\r\n"),
"{text}"
);
assert!(!text.contains("Bcc:"), "{text}");
// The same message went to the outside archive, which can't be reached
// from here: once it leaves the queue (given up on, or deleted), it's
// kept in the built-in journal
let fallback = |entries: &[(EntryId, Entry)]| {
entries
.iter()
.any(|(_, e)| e.subject == "Marked journal-me" && e.journals == vec![number(&remote)])
};
let mut deleted = false;
for _ in 0..100 {
if fallback(&all_entries(test).await) {
break;
}
let (_, response) = call(&admin, "x:QueuedMessage/get", json!({"ids": null})).await;
if let Some(queued) = response["list"]
.as_array()
.unwrap()
.iter()
.find(|m| m.to_string().contains("[email protected]"))
{
let queued_id = queued["id"].as_str().unwrap().to_string();
let (_, response) = call(
&admin,
"x:QueuedMessage/set",
json!({"destroy": [queued_id.clone()]}),
)
.await;
deleted = response["destroyed"] == json!([queued_id]);
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
let kept: Vec<Entry> = all_entries(test)
.await
.into_iter()
.map(|(_, e)| e)
.filter(|e| e.subject == "Marked journal-me")
.collect();
assert_eq!(kept.len(), 2, "{kept:?}");
assert!(kept.iter().any(|e| e.journals == vec![number(&remote)]));
let (_, response) = call(
&admin,
"inbuxa:Journal/get",
json!({"ids": [remote.clone()]}),
)
.await;
let failures = &response["list"][0]["archiveFailures"];
assert_eq!(failures["count"], 1, "{response}");
assert_eq!(
failures["lastReason"],
if deleted {
"it wasn't delivered before leaving the queue"
} else {
"the archive refused it"
},
"{response}"
);
let (_, response) = call(
&admin,
"inbuxa:Journal/get",
json!({"ids": [local.clone()]}),
)
.await;
assert_eq!(response["list"][0]["archiveFailures"]["count"], 0);
// Delivered to an archive here: the report arrives, and nothing goes
// into the built-in journal for that journal
let response = send(
&sender,
&identity,
&mailbox,
&["[email protected]"],
&["[email protected]"],
"To the vault",
)
.await;
assert!(response["created"].get("s").is_some(), "{response}");
let mut arrived = Vec::new();
for _ in 0..100 {
let (_, response) = call(
&vault,
"Email/query",
json!({"filter": {"subject": "Journal report: To the vault"}}),
)
.await;
arrived = response["ids"].as_array().cloned().unwrap_or_default();
if !arrived.is_empty() {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
assert_eq!(arrived.len(), 1, "the report arrived");
assert!(entry_for(test, "To the vault").await.is_none());
assert!(
all_entries(test)
.await
.iter()
.all(|(_, e)| !e.subject.starts_with("Journal report")),
"reports aren't journaled"
);
let (_, response) = call(
&admin,
"inbuxa:Journal/get",
json!({"ids": [local.clone()]}),
)
.await;
assert_eq!(response["list"][0]["archiveFailures"]["count"], 0);
call(&admin, "inbuxa:MailRule/set", json!({"destroy": [rule]})).await;
call(
&admin,
"inbuxa:Journal/set",
json!({"destroy": [rules_only, local, remote]}),
)
.await;
inbuxa_features::journal::invalidate();
}
struct Raw(Vec<u8>);
impl Deserialize for Raw {
@@ -465,6 +702,7 @@ pub async fn journal_tests() {
let admin = test.create_admin_account("[email protected]").await;
test.insert_account(admin);
self::test(&mut test).await;
self::archive(&mut test).await;
if test.is_reset() {
test.temp_dir.delete();
}