Journaling: outside archives, and Journal it in mail flow rules
Phase 3 of the journaling spec. - A journal's destination: builtIn (true for journals stored before) and archiveAddress, at least one. Reports to an archive are queued from the empty sender, one per address, flagged so they're never journaled. - A pending record per report. When the queue lets go of one without delivering it (refused, expired, deleted), it becomes its own entry in the built-in journal under the sending journals' retention, the journal's archiveFailures (count, last time, reason) goes up, and the audit log records it; if that can't be written it stays queued. - Journal it: a rule action naming a journal, on mail flow rules and beside a DLP rule's block, warn or hold. A journal whose scope chooses nobody takes only what rules send it. - The report lists recipients a rule added or redirected to under "Added by rule", by rule name. - A rule's route is cleared between messages in one SMTP session, with the new journal marks; a second message used to keep the first one's route. tests/src/system/journal.rs: destination validation, a rule-only journal fed by a rule that also adds a recipient, an unreachable archive's report kept in the built-in journal with the failure counted, a report delivered to an archive here and not journaled itself.
This commit is contained in:
@@ -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(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
{
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user