diff --git a/crates/features/src/journal/archive.rs b/crates/features/src/journal/archive.rs new file mode 100644 index 0000000..d6b6f8d --- /dev/null +++ b/crates/features/src/journal/archive.rs @@ -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> { + Ok(data + .get_value::>(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 { + Ok(data + .get_value::>(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(()) +} diff --git a/crates/features/src/journal/mod.rs b/crates/features/src/journal/mod.rs index f0c2e35..43ce152 100644 --- a/crates/features/src/journal/mod.rs +++ b/crates/features/src/journal/mod.rs @@ -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, #[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(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("journal@archive.example".into()); + assert!(j.validate().is_ok()); + for bad in [ + "archive", + "a@b", + "a b@c.example", + "", + "a@b@c.example", + ] { + 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 { diff --git a/crates/features/src/journal/report.rs b/crates/features/src/journal/report.rs index c8c3da8..a24ad92 100644 --- a/crates/features/src/journal/report.rs +++ b/crates/features/src/journal/report.rs @@ -20,6 +20,8 @@ use sha2::{Digest, Sha256}; pub struct Recipient { pub address: String, pub orcpt: Option, + /// The mail flow rule that added or redirected to it. + pub added_by: Option, } /// What the queue knows about a message. @@ -46,6 +48,8 @@ pub struct Fields { pub bcc: Vec, /// A list's address, and its members among the recipients. pub expanded: Vec<(String, Vec)>, + /// A rule's name, and the recipients it added. + pub added: Vec<(String, Vec)>, } /// 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("archive@example.com", None); + copied.added_by = Some("Copy finance".into()); + let recipients = [rcpt("pay@bank.example", 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 -> archive@example.com\r\n") + ); + } + #[test] fn values_stay_on_one_line() { let recipients = [rcpt("x@example.com", None)]; diff --git a/crates/features/src/mailflow/rules.rs b/crates/features/src/mailflow/rules.rs index a94174b..e2d8a86 100644 --- a/crates/features/src/mailflow/rules.rs +++ b/crates/features/src/mailflow/rules.rs @@ -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(id: &u32, serializer: S) -> Result { + 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 { + 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", diff --git a/crates/jmap-proto/src/object/inbuxa_journal.rs b/crates/jmap-proto/src/object/inbuxa_journal.rs index 5b79069..7b5b0d5 100644 --- a/crates/jmap-proto/src/object/inbuxa_journal.rs +++ b/crates/jmap-proto/src/object/inbuxa_journal.rs @@ -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, diff --git a/crates/jmap/src/inbuxa/journal.rs b/crates/jmap/src/inbuxa/journal.rs index 2497235..858f63d 100644 --- a/crates/jmap/src/inbuxa/journal.rs +++ b/crates/jmap/src/inbuxa/journal.rs @@ -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) } diff --git a/crates/smtp/src/core/mod.rs b/crates/smtp/src/core/mod.rs index dedad18..96c6ccb 100644 --- a/crates/smtp/src/core/mod.rs +++ b/crates/smtp/src/core/mod.rs @@ -102,6 +102,10 @@ pub struct SessionData { pub dlp_refusal: Option, // inbuxa: a mail flow rule's route for this message pub mailflow_queue: Option, + // inbuxa: journaling (JR-3, JR-10): journals rules sent this message + // to, and recipients rules added, by rule name + pub journal_marks: Vec, + 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(), } } } diff --git a/crates/smtp/src/inbound/data.rs b/crates/smtp/src/inbound/data.rs index b6a7280..ed0c5e0 100644 --- a/crates/smtp/src/inbound/data.rs +++ b/crates/smtp/src/inbound/data.rs @@ -763,19 +763,27 @@ impl Session { } 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 Session { .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 { diff --git a/crates/smtp/src/inbound/mailflow.rs b/crates/smtp/src/inbound/mailflow.rs index 36789e9..966aaee 100644 --- a/crates/smtp/src/inbound/mailflow.rs +++ b/crates/smtp/src/inbound/mailflow.rs @@ -63,9 +63,12 @@ pub struct HeldDraft { /// What a transport rule changes about where a message goes. pub enum EnvelopeChange { - AddRecipient(String), - Redirect(Vec), + /// An address, and the rule that added it. + AddRecipient(String, String), + Redirect(Vec, 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 Session { 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 Session { } RuleAction::Block { .. } | RuleAction::Warn { .. } - | RuleAction::Hold { .. } => None, + | RuleAction::Hold { .. } + | RuleAction::Journal { .. } => None, }; if next.is_some() { current = next; @@ -450,6 +460,19 @@ impl Session { .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, diff --git a/crates/smtp/src/inbound/session.rs b/crates/smtp/src/inbound/session.rs index 3319f57..84ece5c 100644 --- a/crates/smtp/src/inbound/session.rs +++ b/crates/smtp/src/inbound/session.rs @@ -494,6 +494,10 @@ impl Session { 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) { diff --git a/crates/smtp/src/queue/journal.rs b/crates/smtp/src/queue/journal.rs index 10280c3..37ec934 100644 --- a/crates/smtp/src/queue/journal.rs +++ b/crates/smtp/src/queue/journal.rs @@ -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, + /// 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 = 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 +} diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index 131aecb..88a4f06 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -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) -> 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, ) -> 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, 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 diff --git a/docs/spec/features/journaling.md b/docs/spec/features/journaling.md index 2754199..4d844e1 100644 --- a/docs/spec/features/journaling.md +++ b/docs/spec/features/journaling.md @@ -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 diff --git a/tests/src/system/journal.rs b/tests/src/system/journal.rs index 3377d52..4c9ac4f 100644 --- a/tests/src/system/journal.rs +++ b/tests/src/system/journal.rs @@ -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("admin@example.com"); + let sender = admin + .create_user_account( + "archive-sender@example.com", + "archive-sender-secret-7201", + "Archive sender", + &[], + vec![], + ) + .await; + let vault = admin + .create_user_account( + "journal-vault@example.com", + "journal-vault-secret-7202", + "Journal vault", + &[], + vec![], + ) + .await; + let (_, response) = call( + &sender, + "Identity/set", + json!({"create": {"i": {"name": "Sender", "email": "archive-sender@example.com"}}}), + ) + .await; + let identity = response["created"]["i"]["id"].as_str().unwrap().to_string(); + let (_, response) = call( + &sender, + "Mailbox/set", + json!({"create": {"m": {"name": "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": "journal-vault@example.com"}, + "remote": {"name": "To an outside archive", "enabled": true, "direction": "internal", + "scope": {"accounts": [sender.id_string()]}, "retentionDays": 30, + "builtIn": false, "archiveAddress": "vault@elsewhere.org"} + }}), + ) + .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": "journal-vault@example.com"}, + {"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, + &["archive-sender@example.com"], + &["archive-sender@example.com"], + "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 -> journal-vault@example.com\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("vault@elsewhere.org")) + { + 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 = 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, + &["someone@elsewhere.org"], + &["someone@elsewhere.org"], + "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); impl Deserialize for Raw { @@ -465,6 +702,7 @@ pub async fn journal_tests() { let admin = test.create_admin_account("admin@example.com").await; test.insert_account(admin); self::test(&mut test).await; + self::archive(&mut test).await; if test.is_reset() { test.temp_dir.delete(); }