diff --git a/crates/features/src/mailflow/held.rs b/crates/features/src/mailflow/held.rs new file mode 100644 index 0000000..7e85ae1 --- /dev/null +++ b/crates/features/src/mailflow/held.rs @@ -0,0 +1,183 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Mail held for review (dlp-and-mail-flow-rules spec, §2.6). +//! +//! A held message is queued as any other, but released [`HOLD_SECONDS`] +//! from now, the queue's own future-release mechanism: nothing about the +//! queue's stored format changes, so a node on an older version reads it +//! and simply never sends it. Beside it, a review record under `R` `h` + +//! queue id (u64) says why it's held, for the review queue. +//! +//! A reviewer releases it (it's rescheduled from the queue's settings and +//! delivered) or rejects it (it's removed, and the sender told). Unreviewed +//! mail is rejected after [`KEEP_DAYS`]. + +use serde::{Deserialize as SerdeDeserialize, Serialize as SerdeSerialize, de::DeserializeOwned}; +use store::{ + Deserialize, IterateParams, SUBSPACE_INBUXA, Serialize, Store, ValueKey, + write::{AnyClass, BatchBuilder, ValueClass}, +}; +use trc::AddContext; + +const FEATURE: u8 = b'R'; +const KIND_HELD: u8 = b'h'; + +/// How far off a held message's release is set: a century, so it never +/// comes due on its own. +pub const HOLD_SECONDS: u64 = 100 * 365 * 24 * 60 * 60; + +/// How long unreviewed mail waits before it's rejected (settled answer 5). +pub const KEEP_DAYS: u64 = 7; + +/// A rule that held the message, with its notice. +#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)] +pub struct HeldRule { + pub name: String, + pub notice: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, SerdeSerialize, SerdeDeserialize)] +#[serde(rename_all = "camelCase")] +pub struct Held { + pub queue_id: u64, + pub sender: String, + #[serde(default)] + pub account_id: Option, + #[serde(default)] + pub tenant_id: Option, + pub recipients: Vec, + pub subject: String, + pub size: u64, + pub rules: Vec, + /// Each detector that counted, and its count. + #[serde(default)] + pub counts: Vec<(String, usize)>, + /// Seconds since the epoch. + pub held_at: u64, + pub expires_at: u64, +} + +impl Held { + pub fn is_expired(&self, now: u64) -> bool { + now >= self.expires_at + } +} + +struct Json(T); + +impl Serialize for Json { + fn serialize(&self) -> trc::Result> { + serde_json::to_vec(&self.0).map_err(|err| { + trc::StoreEvent::UnexpectedError + .into_err() + .details("Failed to serialize held message") + .reason(err) + }) + } +} + +impl Deserialize for Json { + fn deserialize(bytes: &[u8]) -> trc::Result { + serde_json::from_slice(bytes).map(Json).map_err(|err| { + trc::StoreEvent::DataCorruption + .into_err() + .details("Invalid held message") + .reason(err) + }) + } +} + +fn class(queue_id: u64) -> ValueClass { + let mut key = Vec::with_capacity(10); + key.push(FEATURE); + key.push(KIND_HELD); + key.extend_from_slice(&queue_id.to_be_bytes()); + ValueClass::Any(AnyClass { + subspace: SUBSPACE_INBUXA, + key, + }) +} + +fn key(queue_id: u64) -> ValueKey { + ValueKey::from(class(queue_id)) +} + +pub async fn get(data: &Store, queue_id: u64) -> trc::Result> { + Ok(data + .get_value::>(key(queue_id)) + .await + .caused_by(trc::location!())? + .map(|Json(held)| held)) +} + +pub async fn is_held(data: &Store, queue_id: u64) -> trc::Result { + get(data, queue_id).await.map(|held| held.is_some()) +} + +/// Every held message, oldest first. +pub async fn all(data: &Store) -> trc::Result> { + let mut held = Vec::new(); + data.iterate(IterateParams::new(key(0), key(u64::MAX)), |_, value| { + if let Ok(Json(record)) = Json::::deserialize(value) { + held.push(record); + } + Ok(true) + }) + .await + .caused_by(trc::location!())?; + held.sort_by_key(|h| (h.held_at, h.queue_id)); + Ok(held) +} + +pub async fn create(data: &Store, held: &Held) -> trc::Result<()> { + let mut batch = BatchBuilder::new(); + batch.set(class(held.queue_id), Json(held).serialize()?); + data.write(batch.build_all()) + .await + .caused_by(trc::location!())?; + Ok(()) +} + +pub async fn delete(data: &Store, queue_id: u64) -> trc::Result<()> { + let mut batch = BatchBuilder::new(); + batch.clear(class(queue_id)); + data.write(batch.build_all()) + .await + .caused_by(trc::location!())?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn wire_format_and_expiry() { + let held = Held { + queue_id: 42, + sender: "dana@example.com".into(), + account_id: Some(7), + tenant_id: None, + recipients: vec!["x@elsewhere.org".into()], + subject: "Numbers".into(), + size: 900, + rules: vec![HeldRule { + name: "Cards".into(), + notice: "Held for review".into(), + }], + counts: vec![("payment-card".into(), 5)], + held_at: 1_000, + expires_at: 1_000 + KEEP_DAYS * 86_400, + }; + let json = serde_json::to_value(&held).unwrap(); + assert_eq!(json["heldAt"], 1_000); + assert_eq!(serde_json::from_value::(json).unwrap(), held); + assert!(!held.is_expired(1_000 + KEEP_DAYS * 86_400 - 1)); + assert!(held.is_expired(1_000 + KEEP_DAYS * 86_400)); + assert!(HOLD_SECONDS > 90 * 365 * 86_400); + } +} diff --git a/crates/features/src/mailflow/mod.rs b/crates/features/src/mailflow/mod.rs index d07dc1f..6104c52 100644 --- a/crates/features/src/mailflow/mod.rs +++ b/crates/features/src/mailflow/mod.rs @@ -26,6 +26,7 @@ pub mod cache; pub mod detectors; pub mod engine; pub mod extract; +pub mod held; pub mod rewrite; pub mod rules; pub mod words; diff --git a/crates/jmap-proto/src/object/inbuxa_held_message.rs b/crates/jmap-proto/src/object/inbuxa_held_message.rs new file mode 100644 index 0000000..1f7436a --- /dev/null +++ b/crates/jmap-proto/src/object/inbuxa_held_message.rs @@ -0,0 +1,213 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! `inbuxa:HeldMessage/get` and `/set` under `urn:inbuxa:jmap`: mail held +//! for review (dlp-and-mail-flow-rules spec, §2.6). Get lists it; `preview` +//! (the text, only when asked for) is recorded as access to someone's mail. +//! Set only updates: `{"decision": "release"}`, or `"reject"` with an +//! optional `note` for the sender. The call's `reason` goes into the audit +//! log and is required. + +use crate::{ + object::{AnyId, JmapObject, JmapObjectId}, + request::deserialize::DeserializeArguments, +}; +use jmap_tools::{Element, Key, Property}; +use std::{borrow::Cow, str::FromStr}; +use types::id::Id; + +#[derive(Debug, Clone, Default)] +pub struct HeldMessage; + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum HeldMessageProperty { + Id, + Sender, + Recipients, + Subject, + Size, + Rules, + Counts, + HeldAt, + ExpiresAt, + Preview, + Decision, + Note, +} + +#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum HeldMessageValue { + Id(Id), +} + +impl Property for HeldMessageProperty { + fn try_parse(parent: Option<&Key<'_, Self>>, value: &str) -> Option { + // Keys inside rules and counts stay plain keys + match parent { + None => HeldMessageProperty::parse(value), + Some(_) => None, + } + } + + fn to_cow(&self) -> Cow<'static, str> { + match self { + HeldMessageProperty::Id => "id", + HeldMessageProperty::Sender => "sender", + HeldMessageProperty::Recipients => "recipients", + HeldMessageProperty::Subject => "subject", + HeldMessageProperty::Size => "size", + HeldMessageProperty::Rules => "rules", + HeldMessageProperty::Counts => "counts", + HeldMessageProperty::HeldAt => "heldAt", + HeldMessageProperty::ExpiresAt => "expiresAt", + HeldMessageProperty::Preview => "preview", + HeldMessageProperty::Decision => "decision", + HeldMessageProperty::Note => "note", + } + .into() + } +} + +impl HeldMessageProperty { + fn parse(value: &str) -> Option { + hashify::tiny_map!(value.as_bytes(), + b"id" => HeldMessageProperty::Id, + b"sender" => HeldMessageProperty::Sender, + b"recipients" => HeldMessageProperty::Recipients, + b"subject" => HeldMessageProperty::Subject, + b"size" => HeldMessageProperty::Size, + b"rules" => HeldMessageProperty::Rules, + b"counts" => HeldMessageProperty::Counts, + b"heldAt" => HeldMessageProperty::HeldAt, + b"expiresAt" => HeldMessageProperty::ExpiresAt, + b"preview" => HeldMessageProperty::Preview, + b"decision" => HeldMessageProperty::Decision, + b"note" => HeldMessageProperty::Note, + ) + } +} + +impl FromStr for HeldMessageProperty { + type Err = (); + + fn from_str(s: &str) -> Result { + HeldMessageProperty::parse(s).ok_or(()) + } +} + +impl Element for HeldMessageValue { + type Property = HeldMessageProperty; + + fn try_parse

(key: &Key<'_, Self::Property>, value: &str) -> Option { + match key { + Key::Property(HeldMessageProperty::Id) => { + Id::from_str(value).ok().map(HeldMessageValue::Id) + } + _ => None, + } + } + + fn to_cow(&self) -> Cow<'static, str> { + match self { + HeldMessageValue::Id(id) => id.to_string().into(), + } + } +} + +/// The set call's own argument: why, for the audit log (required). +#[derive(Debug, Clone, Default)] +pub struct HeldMessageSetArguments { + pub reason: Option, +} + +impl<'de> DeserializeArguments<'de> for HeldMessageSetArguments { + fn deserialize_argument(&mut self, key: &str, map: &mut A) -> Result<(), A::Error> + where + A: serde::de::MapAccess<'de>, + { + if key == "reason" { + self.reason = map.next_value()?; + } else { + let _ = map.next_value::()?; + } + Ok(()) + } +} + +impl JmapObject for HeldMessage { + type Property = HeldMessageProperty; + + type Element = HeldMessageValue; + + type Id = Id; + + type Filter = (); + + type Comparator = (); + + type GetArguments = (); + + type SetArguments<'de> = HeldMessageSetArguments; + + type QueryArguments = (); + + type CopyArguments = (); + + type ParseArguments = (); + + const ID_PROPERTY: Self::Property = HeldMessageProperty::Id; +} + +impl From for HeldMessageValue { + fn from(id: Id) -> Self { + HeldMessageValue::Id(id) + } +} + +impl JmapObjectId for HeldMessageValue { + fn as_id(&self) -> Option { + match self { + HeldMessageValue::Id(id) => Some(*id), + } + } + + fn as_any_id(&self) -> Option { + match self { + HeldMessageValue::Id(id) => Some(AnyId::Id(*id)), + } + } + + fn as_id_ref(&self) -> Option<&str> { + None + } + + fn try_set_id(&mut self, new_id: AnyId) -> bool { + if let AnyId::Id(id) = new_id { + *self = HeldMessageValue::Id(id); + true + } else { + false + } + } +} + +impl JmapObjectId for HeldMessageProperty { + fn as_id(&self) -> Option { + None + } + + fn as_any_id(&self) -> Option { + None + } + + fn as_id_ref(&self) -> Option<&str> { + None + } + + fn try_set_id(&mut self, _: AnyId) -> bool { + false + } +} diff --git a/crates/jmap-proto/src/object/mod.rs b/crates/jmap-proto/src/object/mod.rs index a85bb6f..a0053ca 100644 --- a/crates/jmap-proto/src/object/mod.rs +++ b/crates/jmap-proto/src/object/mod.rs @@ -29,6 +29,7 @@ pub mod inbuxa_inventory_snapshot; // inbuxa: personal-data catalog pub mod inbuxa_audit; // inbuxa: the audit log pub mod inbuxa_legal_hold; // inbuxa: legal hold pub mod inbuxa_mail_rule; // inbuxa: DLP and mail flow rules +pub mod inbuxa_held_message; // inbuxa: mail held for review pub mod inbuxa_hold_export; // inbuxa: legal hold exports pub mod inbuxa_explanation; // inbuxa: "Explain this" with the local model pub mod inbuxa_protocol_policy; // inbuxa: legacy protocols off diff --git a/crates/jmap-proto/src/references/eval.rs b/crates/jmap-proto/src/references/eval.rs index b2678aa..fb45299 100644 --- a/crates/jmap-proto/src/references/eval.rs +++ b/crates/jmap-proto/src/references/eval.rs @@ -85,6 +85,9 @@ impl Response<'_> { GetResponseMethod::MailRule(response) => { response.eval_jptr(path, &mut results) } + GetResponseMethod::HeldMessage(response) => { + response.eval_jptr(path, &mut results) + } GetResponseMethod::HoldExport(response) => { response.eval_jptr(path, &mut results) } diff --git a/crates/jmap-proto/src/references/resolve.rs b/crates/jmap-proto/src/references/resolve.rs index a447be1..ac1a860 100644 --- a/crates/jmap-proto/src/references/resolve.rs +++ b/crates/jmap-proto/src/references/resolve.rs @@ -54,6 +54,7 @@ impl Response<'_> { GetRequestMethod::AccountLock(request) => request.resolve_references(self)?, GetRequestMethod::LegalHold(request) => request.resolve_references(self)?, GetRequestMethod::MailRule(request) => request.resolve_references(self)?, + GetRequestMethod::HeldMessage(request) => request.resolve_references(self)?, GetRequestMethod::HoldExport(request) => request.resolve_references(self)?, GetRequestMethod::ProtocolPolicy(request) => request.resolve_references(self)?, GetRequestMethod::TenantProtocolPolicy(request) => { @@ -126,6 +127,9 @@ impl Response<'_> { SetRequestMethod::MailRule(request) => { request.resolve_references(self, 1, false)? } + SetRequestMethod::HeldMessage(request) => { + request.resolve_references(self, 1, false)? + } SetRequestMethod::HoldExport(request) => { request.resolve_references(self, 1, false)? } diff --git a/crates/jmap-proto/src/request/method.rs b/crates/jmap-proto/src/request/method.rs index 8a91a58..3431827 100644 --- a/crates/jmap-proto/src/request/method.rs +++ b/crates/jmap-proto/src/request/method.rs @@ -67,6 +67,7 @@ pub enum MethodObject { ProtocolPolicy, // inbuxa: DLP and mail flow rules MailRule, + HeldMessage, TenantProtocolPolicy, } @@ -105,7 +106,8 @@ impl MethodObject { | MethodObject::AccountLock | MethodObject::LegalHold | MethodObject::HoldExport - | MethodObject::MailRule => Capability::Inbuxa, + | MethodObject::MailRule + | MethodObject::HeldMessage => Capability::Inbuxa, MethodObject::ProtocolPolicy => Capability::Inbuxa, MethodObject::TenantProtocolPolicy => Capability::Inbuxa, } @@ -301,6 +303,8 @@ impl MethodName { (MethodFunction::Set, MethodObject::LegalHold) => "inbuxa:LegalHold/set", (MethodFunction::Get, MethodObject::MailRule) => "inbuxa:MailRule/get", (MethodFunction::Set, MethodObject::MailRule) => "inbuxa:MailRule/set", + (MethodFunction::Get, MethodObject::HeldMessage) => "inbuxa:HeldMessage/get", + (MethodFunction::Set, MethodObject::HeldMessage) => "inbuxa:HeldMessage/set", (MethodFunction::Get, MethodObject::HoldExport) => "inbuxa:HoldExport/get", (MethodFunction::Set, MethodObject::HoldExport) => "inbuxa:HoldExport/set", (MethodFunction::Set, MethodObject::AuditVerification) => { @@ -455,6 +459,8 @@ impl MethodName { "inbuxa:LegalHold/set" => (MethodObject::LegalHold, MethodFunction::Set), "inbuxa:MailRule/get" => (MethodObject::MailRule, MethodFunction::Get), "inbuxa:MailRule/set" => (MethodObject::MailRule, MethodFunction::Set), + "inbuxa:HeldMessage/get" => (MethodObject::HeldMessage, MethodFunction::Get), + "inbuxa:HeldMessage/set" => (MethodObject::HeldMessage, MethodFunction::Set), "inbuxa:HoldExport/get" => (MethodObject::HoldExport, MethodFunction::Get), "inbuxa:HoldExport/set" => (MethodObject::HoldExport, MethodFunction::Set), "inbuxa:AuditVerification/set" => (MethodObject::AuditVerification, MethodFunction::Set), @@ -526,6 +532,7 @@ impl Display for MethodObject { MethodObject::AccountLock => "inbuxa:AccountLock", MethodObject::LegalHold => "inbuxa:LegalHold", MethodObject::MailRule => "inbuxa:MailRule", + MethodObject::HeldMessage => "inbuxa:HeldMessage", MethodObject::HoldExport => "inbuxa:HoldExport", MethodObject::ProtocolPolicy => "inbuxa:ProtocolPolicy", MethodObject::TenantProtocolPolicy => "inbuxa:TenantProtocolPolicy", diff --git a/crates/jmap-proto/src/request/mod.rs b/crates/jmap-proto/src/request/mod.rs index 1d5f0e1..c4ff028 100644 --- a/crates/jmap-proto/src/request/mod.rs +++ b/crates/jmap-proto/src/request/mod.rs @@ -124,6 +124,7 @@ pub enum GetRequestMethod { AccountLock(Box>), LegalHold(Box>), MailRule(Box>), + HeldMessage(Box>), HoldExport(Box>), ProtocolPolicy(Box>), TenantProtocolPolicy( @@ -160,6 +161,7 @@ pub enum SetRequestMethod<'x> { AccountLock(Box>), LegalHold(Box>), MailRule(Box>), + HeldMessage(Box>), HoldExport(Box>), ProtocolPolicy(Box>), TenantProtocolPolicy( diff --git a/crates/jmap-proto/src/request/parser.rs b/crates/jmap-proto/src/request/parser.rs index 28512b9..a811598 100644 --- a/crates/jmap-proto/src/request/parser.rs +++ b/crates/jmap-proto/src/request/parser.rs @@ -609,6 +609,21 @@ impl<'de> Visitor<'de> for CallVisitor { return Err(de::Error::invalid_length(1, &self)); } }, + // inbuxa: mail held for review + (MethodFunction::Get, MethodObject::HeldMessage) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Get(GetRequestMethod::HeldMessage(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, + (MethodFunction::Set, MethodObject::HeldMessage) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Set(SetRequestMethod::HeldMessage(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, // inbuxa: DLP and mail flow rules (MethodFunction::Get, MethodObject::MailRule) => match seq.next_element() { Ok(Some(value)) => RequestMethod::Get(GetRequestMethod::MailRule(value)), diff --git a/crates/jmap-proto/src/response/mod.rs b/crates/jmap-proto/src/response/mod.rs index 67eaaf8..cfc4c44 100644 --- a/crates/jmap-proto/src/response/mod.rs +++ b/crates/jmap-proto/src/response/mod.rs @@ -111,6 +111,7 @@ pub enum GetResponseMethod { AccountLock(GetResponse), LegalHold(GetResponse), MailRule(GetResponse), + HeldMessage(GetResponse), HoldExport(GetResponse), ProtocolPolicy(GetResponse), TenantProtocolPolicy( @@ -147,6 +148,7 @@ pub enum SetResponseMethod { AccountLock(Box>), LegalHold(Box>), MailRule(Box>), + HeldMessage(Box>), HoldExport(Box>), Explanation(Box>), ProtocolPolicy(Box>), @@ -801,6 +803,18 @@ impl<'x> From> for } // inbuxa: legal hold +impl<'x> From> for ResponseMethod<'x> { + fn from(value: GetResponse) -> Self { + ResponseMethod::Get(GetResponseMethod::HeldMessage(value)) + } +} + +impl<'x> From> for ResponseMethod<'x> { + fn from(value: SetResponse) -> Self { + ResponseMethod::Set(SetResponseMethod::HeldMessage(Box::new(value))) + } +} + impl<'x> From> for ResponseMethod<'x> { fn from(value: GetResponse) -> Self { ResponseMethod::Get(GetResponseMethod::MailRule(value)) diff --git a/crates/jmap/src/api/auth.rs b/crates/jmap/src/api/auth.rs index 27d808c..724a450 100644 --- a/crates/jmap/src/api/auth.rs +++ b/crates/jmap/src/api/auth.rs @@ -106,6 +106,8 @@ impl JmapAuthorization for AccessToken { // inbuxa: DLP and mail flow rules share an object; either // permission reaches it, and the handler shows each kind // only to those who may see it + // inbuxa: mail held for review (§2.8) + GetRequestMethod::HeldMessage(_) => Permission::SysDlpReviewGet, GetRequestMethod::MailRule(_) => { if self.has_permission(Permission::SysMailRuleGet) { Permission::SysMailRuleGet @@ -257,6 +259,14 @@ impl JmapAuthorization for AccessToken { Permission::SysLegalHoldUpdate, Permission::SysLegalHoldUpdate, ), + // inbuxa: releasing or rejecting held mail (§2.8) + SetRequestMethod::HeldMessage(s) => validate_set( + s, + self, + Permission::SysDlpReviewUpdate, + Permission::SysDlpReviewUpdate, + Permission::SysDlpReviewUpdate, + ), // inbuxa: DLP and mail flow rules: either change // permission gets in; the handler checks each rule's kind SetRequestMethod::MailRule(_) => { @@ -431,6 +441,7 @@ impl JmapAuthorization for AccessToken { | MethodObject::LegalHold | MethodObject::HoldExport | MethodObject::MailRule + | MethodObject::HeldMessage | MethodObject::ProtocolPolicy | MethodObject::TenantProtocolPolicy => Permission::JmapEmailChanges, // inbuxa: x:MaskedEmail/changes reads what /get reads diff --git a/crates/jmap/src/api/request.rs b/crates/jmap/src/api/request.rs index 903bb61..513bc9b 100644 --- a/crates/jmap/src/api/request.rs +++ b/crates/jmap/src/api/request.rs @@ -282,6 +282,9 @@ impl RequestHandler for Server { SetResponseMethod::MailRule(set_response) => { set_response.update_created_ids(&mut response); } + SetResponseMethod::HeldMessage(set_response) => { + set_response.update_created_ids(&mut response); + } SetResponseMethod::HoldExport(set_response) => { set_response.update_created_ids(&mut response); } @@ -489,6 +492,11 @@ impl RequestHandler for Server { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; crate::inbuxa::legal_hold::get(self, *req).await?.into() } + // inbuxa: mail held for review + GetRequestMethod::HeldMessage(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + crate::inbuxa::held_message::get(self, access_token, *req).await?.into() + } // inbuxa: DLP and mail flow rules GetRequestMethod::MailRule(mut req) => { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; @@ -918,6 +926,22 @@ impl RequestHandler for Server { .await? .into() } + SetRequestMethod::HeldMessage(mut req) => { + resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; + let reason = req.arguments.reason.clone(); + crate::inbuxa::audit::recorded( + self, + access_token, + session, + &method_name.obj.to_string(), + None, + reason, + *req, + |req| Box::pin(crate::inbuxa::held_message::set(self, access_token, req)), + ) + .await? + .into() + } SetRequestMethod::MailRule(mut req) => { resolve_account_id(&mut req.account_id, method_name.obj, access_token)?; let reason = req.arguments.reason.clone(); diff --git a/crates/jmap/src/changes/get.rs b/crates/jmap/src/changes/get.rs index 2888c77..f3c2af9 100644 --- a/crates/jmap/src/changes/get.rs +++ b/crates/jmap/src/changes/get.rs @@ -430,6 +430,7 @@ impl IntermediateChangesResponse { | MethodObject::LegalHold | MethodObject::HoldExport | MethodObject::MailRule + | MethodObject::HeldMessage | MethodObject::ProtocolPolicy | MethodObject::TenantProtocolPolicy | MethodObject::Registry(_) => unreachable!(), diff --git a/crates/jmap/src/inbuxa/held_message.rs b/crates/jmap/src/inbuxa/held_message.rs new file mode 100644 index 0000000..4130ed1 --- /dev/null +++ b/crates/jmap/src/inbuxa/held_message.rs @@ -0,0 +1,325 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! `inbuxa:HeldMessage` (dlp-and-mail-flow-rules spec, §2.6, §2.8): the +//! review queue. `sysDlpReviewGet` lists held mail and reads it; +//! `sysDlpReviewUpdate` releases or rejects it, with a reason the request +//! layer records. Reading a held message's text is recorded as access to +//! the sender's mail. Nobody in a tenant reaches this (settled answer 3). + +use common::{Server, auth::AccessToken, config::smtp::queue::QueueName}; +use inbuxa_features::{ + audit::{Action, Outcome, Record, Target}, + mailflow::held::{self, Held}, +}; +use jmap_proto::{ + error::set::SetError, + method::{ + get::{GetRequest, GetResponse}, + set::{SetRequest, SetResponse}, + }, + object::inbuxa_held_message::{ + HeldMessage, HeldMessageProperty as P, HeldMessageSetArguments, HeldMessageValue, + }, + request::IntoValid, + types::date::UTCDate, +}; +use jmap_tools::{Key, Map, Value}; +use mail_parser::{MessageParser, MimeHeaders, PartType}; +use smtp::queue::spool::SmtpSpool; +use std::borrow::Cow; +use types::id::Id; + +type HValue = Value<'static, P, HeldMessageValue>; + +const ALL: &[P] = &[ + P::Id, + P::Sender, + P::Recipients, + P::Subject, + P::Size, + P::Rules, + P::Counts, + P::HeldAt, + P::ExpiresAt, +]; + +/// How much of a held message's text a preview shows. +const PREVIEW_LIMIT: usize = 64 * 1024; + +fn server_level(access_token: &AccessToken) -> trc::Result<()> { + if access_token.tenant_id().is_some() { + Err(trc::JmapEvent::Forbidden + .into_err() + .details("Held mail is the server's to review.")) + } else { + Ok(()) + } +} + +fn date(seconds: u64) -> HValue { + Value::Str(UTCDate::from_timestamp(seconds as i64).to_string().into()) +} + +fn text(s: &str) -> HValue { + Value::Str(Cow::Owned(s.to_string())) +} + +/// The text a reviewer reads: the subject, each body as text, and the +/// attachments' names; at most [`PREVIEW_LIMIT`]. +async fn preview(server: &Server, queue_id: u64) -> trc::Result> { + let Some(message) = server.read_message(queue_id, QueueName::default()).await else { + return Ok(None); + }; + let Some(raw) = server + .blob_store() + .get_blob(message.message.blob_hash.as_slice(), 0..usize::MAX) + .await? + else { + return Ok(None); + }; + let Some(parsed) = MessageParser::new().parse(&raw) else { + return Ok(Some( + String::from_utf8_lossy(&raw[..raw.len().min(PREVIEW_LIMIT)]).into_owned(), + )); + }; + let mut out = String::new(); + for part in parsed.text_bodies() { + match &part.body { + PartType::Text(text) => out.push_str(text), + PartType::Html(html) => out.push_str(&mail_parser::decoders::html::html_to_text(html)), + _ => {} + } + out.push_str("\n\n"); + } + let attachments: Vec<&str> = parsed + .attachments() + .filter_map(|a| a.attachment_name()) + .collect(); + if !attachments.is_empty() { + out.push_str(&format!("Attachments: {}\n", attachments.join(", "))); + } + if out.len() > PREVIEW_LIMIT { + let mut cut = PREVIEW_LIMIT; + while !out.is_char_boundary(cut) { + cut -= 1; + } + out.truncate(cut); + } + Ok(Some(out)) +} + +fn to_value(record: &Held, properties: &[P], preview: Option<&str>) -> HValue { + let mut out = Map::with_capacity(properties.len()); + for property in properties { + let value = match property { + P::Id => Value::Element(HeldMessageValue::Id(Id::from(record.queue_id))), + P::Sender => text(&record.sender), + P::Recipients => Value::Array(record.recipients.iter().map(|r| text(r)).collect()), + P::Subject => text(&record.subject), + P::Size => Value::Number(record.size.into()), + P::Rules => Value::Array( + record + .rules + .iter() + .map(|rule| { + let mut map = Map::with_capacity(2); + map.insert_unchecked(Key::Borrowed("name"), text(&rule.name)); + map.insert_unchecked(Key::Borrowed("notice"), text(&rule.notice)); + Value::Object(map) + }) + .collect(), + ), + P::Counts => Value::Array( + record + .counts + .iter() + .map(|(detector, count)| { + let mut map = Map::with_capacity(2); + map.insert_unchecked(Key::Borrowed("detector"), text(detector)); + map.insert_unchecked( + Key::Borrowed("count"), + Value::Number((*count as u64).into()), + ); + Value::Object(map) + }) + .collect(), + ), + P::HeldAt => date(record.held_at), + P::ExpiresAt => date(record.expires_at), + P::Preview => preview.map_or(Value::Null, text), + P::Decision | P::Note => Value::Null, + }; + out.insert_unchecked(Key::Property(property.clone()), value); + } + Value::Object(out) +} + +/// `inbuxa:HeldMessage/get`: held mail, oldest first. +pub async fn get( + server: &Server, + access_token: &AccessToken, + mut request: GetRequest, +) -> trc::Result> { + server_level(access_token)?; + let properties = request.unwrap_properties(ALL); + let (ids, not_found) = request.unwrap_ids(server.core.jmap.get_max_objects)?; + let mut response = GetResponse { + account_id: request.account_id.into(), + state: None, + list: Vec::new(), + not_found, + }; + let all = held::all(server.store()).await?; + let wanted: Vec<&Held> = match &ids { + None => all.iter().collect(), + Some(ids) => { + let mut found = Vec::new(); + for id in ids { + match all.iter().find(|h| h.queue_id == id.id()) { + Some(record) => found.push(record), + None => response.push_not_found(*id), + } + } + found + } + }; + let with_preview = properties.contains(&P::Preview); + for record in wanted { + let text = if with_preview { + let text = preview(server, record.queue_id).await?; + // Reading someone's mail is recorded, as any access is + server + .audit_note(Record { + at: store::write::now() * 1000, + actor: server.audit_actor(access_token).await, + via: access_token.origin().cloned(), + remote_ip: None, + action: Action::BlobAccess, + target: Target { + kind: "inbuxa:HeldMessage".into(), + id: Some(Id::from(record.queue_id).to_string()), + name: Some(record.subject.clone()), + account_id: record.account_id, + tenant_id: record.tenant_id, + }, + changes: vec![], + details: Some(format!( + "Read a message held for review, from {}", + record.sender + )), + reason: None, + outcome: Outcome::success(), + }) + .await; + text + } else { + None + }; + response + .list + .push(to_value(record, &properties, text.as_deref())); + } + Ok(response) +} + +fn invalid(property: P, why: &str) -> SetError

{ + SetError::invalid_properties() + .with_property(property) + .with_description(why.to_string()) +} + +/// `inbuxa:HeldMessage/set`: update with `decision` release or reject (and +/// an optional `note` for the sender). There is no create or destroy. +pub async fn set( + server: &Server, + access_token: &AccessToken, + mut request: SetRequest<'_, HeldMessage>, +) -> trc::Result> { + server_level(access_token)?; + let mut response = SetResponse::from_request(&request, server.core.jmap.set_max_objects)?; + let arguments: HeldMessageSetArguments = std::mem::take(&mut request.arguments); + let has_reason = arguments + .reason + .as_deref() + .is_some_and(|r| !r.trim().is_empty()); + + for (client_id, _) in request.unwrap_create() { + response.not_created.append( + client_id, + SetError::forbidden().with_description("Mail is held by DLP rules, not created."), + ); + } + + 'update: for (id, value) in request.unwrap_update().into_valid() { + let Some(record) = held::get(server.store(), id.id()).await? else { + response.not_updated.append(id, SetError::not_found()); + continue; + }; + if !has_reason { + response.not_updated.append( + id, + SetError::invalid_properties().with_description( + "Say why: a reason is required and is kept in the audit log.", + ), + ); + continue; + } + let mut decision = None; + let mut note = None; + for (key, value) in value.into_expanded_object() { + match (&key, value) { + (Key::Property(P::Decision), Value::Str(s)) if s == "release" || s == "reject" => { + decision = Some(s.to_string()); + } + (Key::Property(P::Note), Value::Str(s)) => { + let s = s.trim(); + if !s.is_empty() { + note = Some(s.chars().take(1000).collect::()); + } + } + (Key::Property(P::Note), Value::Null) => {} + _ => { + response.not_updated.append( + id, + invalid( + P::Decision, + "Send decision: \"release\" or \"reject\", and an optional note.", + ), + ); + continue 'update; + } + } + } + let done = match decision.as_deref() { + Some("release") => smtp::queue::held::release(server, record.queue_id).await?, + Some("reject") => smtp::queue::held::reject(server, &record, note.as_deref()).await?, + _ => { + response + .not_updated + .append(id, invalid(P::Decision, "Say release or reject.")); + continue; + } + }; + if done { + response.updated.append(id, None); + } else { + response.not_updated.append( + id, + SetError::not_found().with_description("The message is no longer in the queue."), + ); + } + } + + for id in request.unwrap_destroy().into_valid() { + response.not_destroyed.append( + id, + SetError::forbidden().with_description("Release or reject it instead."), + ); + } + + Ok(response) +} diff --git a/crates/jmap/src/inbuxa/mod.rs b/crates/jmap/src/inbuxa/mod.rs index a5a2638..5756737 100644 --- a/crates/jmap/src/inbuxa/mod.rs +++ b/crates/jmap/src/inbuxa/mod.rs @@ -11,6 +11,7 @@ pub mod access; pub mod account_lock; pub mod legal_hold; pub mod mail_rule; +pub mod held_message; pub mod hold_export; pub mod hold_export_api; pub mod audit; diff --git a/crates/jmap/src/registry/mapping/queued_message.rs b/crates/jmap/src/registry/mapping/queued_message.rs index c8c8c26..c016898 100644 --- a/crates/jmap/src/registry/mapping/queued_message.rs +++ b/crates/jmap/src/registry/mapping/queued_message.rs @@ -49,6 +49,12 @@ use trc::AddContext; use types::{blob::BlobId, blob_hash::BlobHash, id::Id}; use utils::map::vec_map::VecMap; +/// inbuxa: held mail is the review queue's to decide. +fn held_refusal() -> SetError { + SetError::forbidden() + .with_description("This message is held for review: release or reject it under Compliance, Held mail.") +} + pub(crate) async fn queued_message_set( mut set: RegistrySetResponse<'_>, ) -> trc::Result> { @@ -66,6 +72,12 @@ pub(crate) async fn queued_message_set( let mut refresh_queue = false; 'outer: for (id, value) in set.update.drain(..) { let queue_id = id.id(); + // inbuxa: held mail is released or rejected by review, not here + // (dlp-and-mail-flow-rules spec, §2.6) + if inbuxa_features::mailflow::held::is_held(set.server.store(), queue_id).await? { + set.response.not_updated.append(id, held_refusal()); + continue; + } let Some(archive) = set.server.read_message_archive(queue_id).await? else { set.response.not_updated.append(id, SetError::not_found()); continue; @@ -238,6 +250,11 @@ pub(crate) async fn queued_message_set( // Process destroy operations for id in set.destroy.drain(..) { + // inbuxa: §2.6, as above + if inbuxa_features::mailflow::held::is_held(set.server.store(), id.id()).await? { + set.response.not_destroyed.append(id, held_refusal()); + continue; + } let Some(message) = set.server.read_message(id.id(), QueueName::default()).await else { set.response.not_destroyed.append(id, SetError::not_found()); continue; diff --git a/crates/jmap/src/submission/set.rs b/crates/jmap/src/submission/set.rs index 7c1a6fc..055ab30 100644 --- a/crates/jmap/src/submission/set.rs +++ b/crates/jmap/src/submission/set.rs @@ -214,6 +214,17 @@ impl EmailSubmissionSet for Server { } match undo_status { + // inbuxa: held for review: the review decides, not an unsend + // (dlp-and-mail-flow-rules spec, §2.6) + Some(email_submission::UndoStatus::Canceled) + if inbuxa_features::mailflow::held::is_held(self.store(), queue_id).await? => + { + response.not_updated.append( + id, + SetError::new(SetErrorType::CannotUnsend) + .with_description("The message is held for review and can't be unsent."), + ); + } Some(email_submission::UndoStatus::Canceled) => { if let Some(queue_message) = self.read_message(queue_id, QueueName::default()).await diff --git a/crates/services/src/task_manager/maintenance.rs b/crates/services/src/task_manager/maintenance.rs index 9169b5a..b8e2414 100644 --- a/crates/services/src/task_manager/maintenance.rs +++ b/crates/services/src/task_manager/maintenance.rs @@ -286,6 +286,11 @@ async fn store_maintenance( trc::error!(err.details("Failed to purge expired IP bans")); } + // inbuxa: DLP, §2.6: mail nobody reviewed in time goes back + if let Err(err) = smtp::queue::held::expire(server).await { + trc::error!(err.details("Failed to return unreviewed held mail")); + } + // inbuxa: AU-7: audit records past their retention go; a // failure leaves them for the next run if let Err(err) = server.audit_purge().await { diff --git a/crates/smtp/src/inbound/data.rs b/crates/smtp/src/inbound/data.rs index 4d13fc0..b6a7280 100644 --- a/crates/smtp/src/inbound/data.rs +++ b/crates/smtp/src/inbound/data.rs @@ -740,35 +740,46 @@ impl Session { // inbuxa: DLP (dlp-and-mail-flow-rules spec, §2.1): after the system // script, before headers and signing - match self + let mut held_draft = None; + let (message, envelope) = match self .check_mail_rules(edited_message.as_deref().unwrap_or(raw_message.as_slice())) .await { - super::mailflow::Checked::Accept => {} - super::mailflow::Checked::Changed { message, envelope } => { - if let Some(message) = message { - edited_message = Some(message); - } - for change in envelope { - match change { - super::mailflow::EnvelopeChange::AddRecipient(address) => { - if !self.data.rcpt_to.iter().any(|r| r.address_lcase.eq_ignore_ascii_case(&address)) { - self.data.rcpt_to.push(SessionAddress::new(address)); - } - } - super::mailflow::EnvelopeChange::Redirect(addresses) => { - self.data.rcpt_to = addresses.into_iter().map(SessionAddress::new).collect(); - } - super::mailflow::EnvelopeChange::Route(queue) => { - self.data.mailflow_queue = Some(queue); - } - } - } + super::mailflow::Checked::Accept => (None, Vec::new()), + super::mailflow::Checked::Changed { message, envelope } => (message, envelope), + // §2.6: queued, but not due for a century; a reviewer releases it + super::mailflow::Checked::Hold { draft, message, envelope } => { + self.data.future_release = inbuxa_features::mailflow::held::HOLD_SECONDS; + held_draft = Some(draft); + (message, envelope) } super::mailflow::Checked::Refuse(reply, refusal) => { self.data.dlp_refusal = refusal; return reply.into(); } + }; + if let Some(message) = message { + edited_message = Some(message); + } + for change in envelope { + match change { + super::mailflow::EnvelopeChange::AddRecipient(address) => { + if !self + .data + .rcpt_to + .iter() + .any(|r| r.address_lcase.eq_ignore_ascii_case(&address)) + { + self.data.rcpt_to.push(SessionAddress::new(address)); + } + } + super::mailflow::EnvelopeChange::Redirect(addresses) => { + self.data.rcpt_to = addresses.into_iter().map(SessionAddress::new).collect(); + } + super::mailflow::EnvelopeChange::Route(queue) => { + self.data.mailflow_queue = Some(queue); + } + } } // Build message @@ -850,6 +861,19 @@ impl Session { .server .eval_signers(&ac.dkim.sign, self, self.data.session_id) .await; + // inbuxa: §2.6, who the held message is from and to + let held_envelope = held_draft.as_ref().map(|_| { + ( + message.message.return_path.to_string(), + message + .message + .recipients + .iter() + .map(|r| r.address.to_string()) + .collect::>(), + message.message.size, + ) + }); if message .queue( QueueParams::new(raw_message, self.data.session_id, &self.server) @@ -864,9 +888,17 @@ impl Session { { self.state = State::Accepted(queue_id); self.data.messages_sent += 1; - format!("250 2.0.0 Message queued with id {queue_id:x}.\r\n") - .into_bytes() - .into() + if let (Some(draft), Some((sender, recipients, size))) = (held_draft, held_envelope) + { + self.record_held(queue_id, draft, sender, recipients, size).await; + format!("250 2.0.0 Held for review, id {queue_id:x}.\r\n") + .into_bytes() + .into() + } else { + format!("250 2.0.0 Message queued with id {queue_id:x}.\r\n") + .into_bytes() + .into() + } } else { (b"451 4.3.5 Unable to accept message at this time.\r\n"[..]).into() } diff --git a/crates/smtp/src/inbound/mailflow.rs b/crates/smtp/src/inbound/mailflow.rs index a84704b..4088f92 100644 --- a/crates/smtp/src/inbound/mailflow.rs +++ b/crates/smtp/src/inbound/mailflow.rs @@ -21,6 +21,7 @@ use inbuxa_features::{ Attachment, Content, Decision, Envelope, Outcome as RulesOutcome, Recipient, RuleRef, }, extract::{self, Extracted, Limits}, + held::{self, Held, HeldRule, KEEP_DAYS}, rewrite, rules::{Action as RuleAction, Kind}, }, @@ -44,6 +45,20 @@ pub enum Checked { }, /// Refuse, with this SMTP reply, and for a JMAP submission, why. Refuse(Vec, Option), + /// Queue it held for review (§2.6), with any transport changes. + Hold { + draft: HeldDraft, + message: Option>, + envelope: Vec, + }, +} + +/// What the review record will say, once the message has a queue id. +pub struct HeldDraft { + pub subject: String, + pub rules: Vec, + pub counts: Vec<(String, usize)>, + pub notify_sender: bool, } /// What a transport rule changes about where a message goes. @@ -219,6 +234,7 @@ impl Session { None }; let checked_subject = tag.as_ref().map_or(subject, |(_, rest)| rest.as_str()); + let held_subject = checked_subject.to_string(); let mut content = Content { subject: checked_subject, @@ -319,22 +335,51 @@ impl Session { .await; } - match decision { - Decision::Block(rules) | Decision::Hold { rules, .. } => { - // Hold for review is phase 3: until then a hold rule blocks, - // rather than let the message through unreviewed + let hold = match decision { + Decision::Block(rules) => { let refusal = refusal(true, &rules); - Checked::Refuse(format!("550 5.7.1 {}\r\n", notices(&rules)).into_bytes(), Some(refusal)) + return Checked::Refuse( + format!("550 5.7.1 {}\r\n", notices(&rules)).into_bytes(), + Some(refusal), + ); } - Decision::Warn(rules) => Checked::Refuse( - format!( - "550 5.7.1 {} To send anyway, start the subject with [override: your reason]\r\n", - notices(&rules) - ) - .into_bytes(), - Some(refusal(false, &rules)), - ), - Decision::Pass => { + Decision::Warn(rules) => { + return Checked::Refuse( + format!( + "550 5.7.1 {} To send anyway, start the subject with [override: your reason]\r\n", + notices(&rules) + ) + .into_bytes(), + Some(refusal(false, &rules)), + ); + } + // Accepted and queued, but not sent until a reviewer says so + // (§2.6); the transport rules still apply, so what's released + // is what would have gone out + Decision::Hold { + rules, + notify_sender, + } => Some(HeldDraft { + subject: held_subject, + rules: rules + .iter() + .map(|r| HeldRule { + name: r.name.clone(), + notice: r.notice.clone(), + }) + .collect(), + counts: outcome + .matched + .iter() + .filter(|m| m.kind == Kind::Dlp) + .flat_map(|m| m.counts.iter().cloned()) + .collect(), + notify_sender, + }), + Decision::Pass => None, + }; + { + { // The tag was an instruction to the server, not part of the // subject: it doesn't go out let mut current: Option> = @@ -344,12 +389,18 @@ impl Session { for action in &matched.actions { let now = current.as_deref().unwrap_or(message); let next = match action { - RuleAction::AddDisclaimer { text, html, position } => { - rewrite::add_disclaimer(now, text, html.as_deref(), *position) + RuleAction::AddDisclaimer { + text, + html, + position, + } => rewrite::add_disclaimer(now, text, html.as_deref(), *position), + RuleAction::AddHeader { name, value } => { + Some(rewrite::add_header(now, name, value)) } - RuleAction::AddHeader { name, value } => Some(rewrite::add_header(now, name, value)), RuleAction::RemoveHeader { name } => rewrite::remove_header(now, name), - RuleAction::PrefixSubject { text } => rewrite::prefix_subject(now, text), + RuleAction::PrefixSubject { text } => { + rewrite::prefix_subject(now, text) + } RuleAction::AddRecipient { address } => { changes.push(EnvelopeChange::AddRecipient(address.clone())); None @@ -363,13 +414,16 @@ impl Session { None } RuleAction::Refuse { text } => { - self.record_transport(&sender, &matched.name, "refused", &domains).await; + self.record_transport(&sender, &matched.name, "refused", &domains) + .await; return Checked::Refuse( format!("550 5.7.1 {}\r\n", reply_text(text)).into_bytes(), None, ); } - RuleAction::Block { .. } | RuleAction::Warn { .. } | RuleAction::Hold { .. } => None, + RuleAction::Block { .. } + | RuleAction::Warn { .. } + | RuleAction::Hold { .. } => None, }; if next.is_some() { current = next; @@ -381,25 +435,76 @@ impl Session { .actions .iter() .filter_map(|a| match a { - RuleAction::AddRecipient { address } => Some(format!("copied to {address}")), - RuleAction::Redirect { addresses } => Some(format!("redirected to {}", addresses.join(", "))), + RuleAction::AddRecipient { address } => { + Some(format!("copied to {address}")) + } + RuleAction::Redirect { addresses } => { + Some(format!("redirected to {}", addresses.join(", "))) + } RuleAction::Route { queue } => Some(format!("routed through {queue}")), _ => None, }) .collect(); if !routed.is_empty() { - self.record_transport(&sender, &matched.name, &routed.join(", "), &domains).await; + self.record_transport(&sender, &matched.name, &routed.join(", "), &domains) + .await; } } - if current.is_none() && changes.is_empty() { - Checked::Accept - } else { - Checked::Changed { message: current, envelope: changes } + match hold { + Some(draft) => Checked::Hold { + draft, + message: current, + envelope: changes, + }, + None if current.is_none() && changes.is_empty() => Checked::Accept, + None => Checked::Changed { + message: current, + envelope: changes, + }, } } } } + /// Writes the review record for a message just queued held (§2.6), + /// and tells the sender when the rule asks. A failure to write it is + /// logged: the message stays held, never sent unreviewed. + pub async fn record_held( + &self, + queue_id: u64, + draft: HeldDraft, + sender: String, + recipients: Vec, + size: u64, + ) { + let at = store::write::now(); + let account = self.data.authenticated_as.as_ref(); + let record = Held { + queue_id, + sender, + account_id: account.map(|a| a.account_id), + tenant_id: account.and_then(|a| a.account.id_tenant), + recipients, + subject: draft.subject, + size, + rules: draft.rules, + counts: draft.counts, + held_at: at, + expires_at: at + KEEP_DAYS * 86_400, + }; + if let Err(err) = held::create(self.server.store(), &record).await { + trc::error!( + err.span_id(self.data.session_id) + .caused_by(trc::location!()) + .details("Failed to write the review record of a held message") + ); + return; + } + if draft.notify_sender { + crate::queue::held::notify_held(&self.server, &record).await; + } + } + /// A transport rule that refused a message or changed where it goes /// (§2.7): who sent it (or the server, for incoming mail), the rule, /// what it did. @@ -485,9 +590,8 @@ impl Session { .collect::>() .join("; "); let (what, outcome, reason) = match decision { - Decision::Block(_) | Decision::Hold { .. } => { - ("blocked", Outcome::refused("inbuxa:dlpBlocked", None), None) - } + Decision::Hold { .. } => ("held for review", Outcome::success(), None), + Decision::Block(_) => ("blocked", Outcome::refused("inbuxa:dlpBlocked", None), None), Decision::Warn(_) => ("warned", Outcome::refused("inbuxa:dlpWarning", None), None), Decision::Pass => ( "sent after a warning", diff --git a/crates/smtp/src/queue/held.rs b/crates/smtp/src/queue/held.rs new file mode 100644 index 0000000..291c1bd --- /dev/null +++ b/crates/smtp/src/queue/held.rs @@ -0,0 +1,177 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! inbuxa: mail held for review (dlp-and-mail-flow-rules spec, §2.6): +//! releasing it, rejecting it, rejecting what nobody reviewed in time, and +//! the notices the sender gets. +//! +//! A held message sits in the queue with its release [`HOLD_SECONDS`] off. +//! Releasing it undoes exactly that: each recipient due now, its next +//! notice as far from now as it was from its retry, its lifetime counted +//! from the release. + +use crate::{ + queue::{Message, MessageWrapper, Status, spool::SmtpSpool}, + reporting::send::MtaReportSend, +}; +use common::{ + Server, + config::smtp::queue::{QueueExpiry, QueueName}, + ipc::QueueEvent, +}; +use inbuxa_features::{ + audit::{Action, Actor, Outcome, Record, Target}, + mailflow::held::{self, HOLD_SECONDS, Held, KEEP_DAYS}, +}; +use mail_builder::{ + MessageBuilder, + headers::{HeaderType, address::Address}, +}; +use store::{ahash::AHashSet, write::now}; + +/// Puts a held message back on its way. False when it's no longer queued. +pub async fn release(server: &Server, queue_id: u64) -> trc::Result { + let Some(archive) = server.read_message_archive(queue_id).await? else { + held::delete(server.store(), queue_id).await?; + return Ok(false); + }; + let mut message: Message = archive.to_unarchived::()?.deserialize()?; + let prev_events = message.next_events(); + let at = now(); + let mut modified = AHashSet::new(); + for (idx, rcpt) in message.recipients.iter_mut().enumerate() { + if !matches!(rcpt.status, Status::Scheduled | Status::TemporaryFailure(_)) { + continue; + } + let notify_gap = rcpt.notify.due.saturating_sub(rcpt.retry.due); + rcpt.retry.due = at; + rcpt.notify.due = at + notify_gap; + if let QueueExpiry::Ttl(ttl) = rcpt.expires { + rcpt.expires = QueueExpiry::Ttl( + ttl.saturating_sub(HOLD_SECONDS) + at.saturating_sub(message.created), + ); + } + modified.insert(idx); + } + let saved = MessageWrapper::new(message, queue_id, QueueName::default()) + .save_registry_changes(server, prev_events, modified) + .await; + held::delete(server.store(), queue_id).await?; + let _ = server.inner.ipc.queue_tx.send(QueueEvent::Refresh).await; + Ok(saved) +} + +/// Takes a held message out of the queue and tells its sender, with the +/// reviewer's note if there is one. False when it's no longer queued. +pub async fn reject(server: &Server, record: &Held, note: Option<&str>) -> trc::Result { + let removed = match server + .read_message(record.queue_id, QueueName::default()) + .await + { + Some(message) => message.remove(server, None).await, + None => false, + }; + held::delete(server.store(), record.queue_id).await?; + let mut text = format!( + "Your message \"{}\" to {} was held for review under this server's rules, and wasn't sent.\r\n", + record.subject, + record.recipients.join(", ") + ); + match note { + Some(note) => text.push_str(&format!("\r\nThe reviewer's note: {note}\r\n")), + None => text.push_str(&format!( + "\r\nNobody reviewed it within {KEEP_DAYS} days, so it was returned.\r\n" + )), + } + notify( + server, + record, + &format!("Not sent: {}", record.subject), + text, + ) + .await; + let _ = server.inner.ipc.queue_tx.send(QueueEvent::Refresh).await; + Ok(removed) +} + +/// Tells the sender their message is held (when the rule asks). +pub async fn notify_held(server: &Server, record: &Held) { + let notices = record + .rules + .iter() + .map(|r| r.notice.as_str()) + .collect::>() + .join(" "); + let text = format!( + "Your message \"{}\" to {} is held for review under this server's rules: {notices}\r\n\r\n\ + It will be sent if a reviewer releases it, and returned otherwise within {KEEP_DAYS} days.\r\n", + record.subject, + record.recipients.join(", "), + ); + notify( + server, + record, + &format!("Held for review: {}", record.subject), + text, + ) + .await; +} + +async fn notify(server: &Server, record: &Held, subject: &str, text: String) { + let domain = record + .sender + .rsplit_once('@') + .map_or("localhost", |(_, d)| d); + let from = format!("postmaster@{domain}"); + let message = MessageBuilder::new() + .from(Address::new_address(Some("Mail review"), from.clone())) + .to(Address::new_address(None::, record.sender.clone())) + .subject(subject) + .header("Auto-Submitted", HeaderType::Text("auto-replied".into())) + .text_body(text) + .write_to_vec() + .unwrap_or_default(); + server + .send_autogenerated(from, [record.sender.as_str()].into_iter(), message, None, 0) + .await; +} + +/// Rejects every held message nobody reviewed in time (§2.6), each +/// recorded as the server's doing. Returns how many. +pub async fn expire(server: &Server) -> trc::Result { + let at = now(); + let mut count = 0; + for record in held::all(server.store()).await? { + if !record.is_expired(at) { + continue; + } + reject(server, &record, None).await?; + count += 1; + server + .audit_note(Record { + at: at * 1000, + actor: Actor::system("DLP"), + via: None, + remote_ip: None, + action: Action::Destroy, + target: Target { + kind: "inbuxa:HeldMessage".into(), + id: Some(record.queue_id.to_string()), + name: Some(record.subject.clone()), + account_id: record.account_id, + tenant_id: record.tenant_id, + }, + changes: vec![], + details: Some(format!( + "Rejected: nobody reviewed it within {KEEP_DAYS} days; the sender was told" + )), + reason: None, + outcome: Outcome::success(), + }) + .await; + } + Ok(count) +} diff --git a/crates/smtp/src/queue/mod.rs b/crates/smtp/src/queue/mod.rs index bacc6c9..3e717bd 100644 --- a/crates/smtp/src/queue/mod.rs +++ b/crates/smtp/src/queue/mod.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use common::{ @@ -21,6 +23,7 @@ use types::blob_hash::BlobHash; use utils::DomainPart; pub mod dsn; +pub mod held; // inbuxa: mail held for review pub mod manager; pub mod quota; pub mod spool; diff --git a/docs/spec/features/dlp-and-mail-flow-rules.md b/docs/spec/features/dlp-and-mail-flow-rules.md index 00489dd..79cbb4a 100644 --- a/docs/spec/features/dlp-and-mail-flow-rules.md +++ b/docs/spec/features/dlp-and-mail-flow-rules.md @@ -293,8 +293,19 @@ once it's held (the webmail says so). Held messages count against no one's quota. Each held message and each decision is in the audit log. -**Until phase 3** a hold rule blocks, with its notice, rather than let the -message through unreviewed. +**As built (phase 3).** Holding uses the queue's own future-release +mechanism: the message is queued with its release a century off, every +recipient's retry, notice and expiry pushed with it, so the stored format +doesn't change. Release puts each recipient due now, keeps the gap to its +next notice, and counts its lifetime from the release. The review record +(`inbuxa:HeldMessage`, under `R` `h` + queue id) holds the sender, +recipients, subject, size, rules and counts. Transport rules still apply to +held mail, so what's released is what would have gone out. The daily +clean-up rejects what's past its 7 days (recorded as the server's doing). +`preview` returns the text (64 KB) only when asked for, and each read is +recorded as `blobAccess`. Emails › Queue refuses to change or delete held +mail, and the sender can't unsend it. The 7 days is a constant for now; a +setting comes with the console page. ### 2.7 What's recorded diff --git a/resources/privacy/catalog.toml b/resources/privacy/catalog.toml index 524954f..e26adf1 100644 --- a/resources/privacy/catalog.toml +++ b/resources/privacy/catalog.toml @@ -79,6 +79,21 @@ lockedAt = ["metadata"] lockedBy = ["identifier"] delegates = ["identifier"] +[object."inbuxa:HeldMessage"] +file = "inbuxa_held_message.rs" +default = "none" +whose = ["holder", "correspondent"] +where = ["data-store", "blob-store"] +scope = "server" +retention = "object-life" +[object."inbuxa:HeldMessage".properties] +sender = ["identifier", "contact"] +recipients = ["identifier", "contact"] +subject = ["content"] +preview = ["content"] +note = ["content"] +counts = ["metadata"] + [object."inbuxa:MailRule"] file = "inbuxa_mail_rule.rs" default = "none" diff --git a/tests/src/system/mail_rules.rs b/tests/src/system/mail_rules.rs index 3eae538..842d8f5 100644 --- a/tests/src/system/mail_rules.rs +++ b/tests/src/system/mail_rules.rs @@ -773,6 +773,329 @@ pub async fn transport(test: &mut TestServer) { call(&admin, "inbuxa:MailRule/set", json!({"destroy": transport})).await; } +/// Hold for review (§2.6): held mail waits, shows in the review queue, +/// can't be sent around the review, and a reviewer releases or rejects it. +pub async fn hold(test: &mut TestServer) { + println!("Running held mail tests..."); + let admin = test.account("admin@example.com"); + let sender = admin + .create_user_account( + "hold-sender@example.com", + "hold-sender-secret-7703", + "Hold sender", + &[], + vec![], + ) + .await; + let (_, response) = call( + &sender, + "Identity/set", + json!({"create": {"i": {"name": "Sender", "email": "hold-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": "Hold drafts"}}}), + ) + .await; + let mailbox = response["created"]["m"]["id"].as_str().unwrap().to_string(); + let (_, response) = call( + &admin, + "inbuxa:MailRule/set", + json!({"create": {"h": { + "name": "Hold cards", "kind": "dlp", "direction": "outgoing", + "conditions": [{"type": "words", "words": ["hold-me"]}, + {"type": "detected", "detectors": [{"id": "payment-card"}]}], + "actions": [{"type": "hold", "notice": "Card numbers are reviewed first.", "notifySender": true}] + }}}), + ) + .await; + let rule = response["created"]["h"]["id"] + .as_str() + .unwrap_or_else(|| panic!("{response}")) + .to_string(); + let body = "hold-me: card 4242 4242 4242 4242"; + + // Accepted, held, listed + let response = submit( + &sender, + &identity, + &mailbox, + &["hold-sender@example.com"], + "Held one", + body, + None, + ) + .await; + let submission = response["created"]["s"]["id"] + .as_str() + .unwrap_or_else(|| panic!("held, not refused: {response}")) + .to_string(); + let (_, response) = call(&admin, "inbuxa:HeldMessage/get", json!({"ids": null})).await; + let list = response["list"] + .as_array() + .unwrap_or_else(|| panic!("{response}")); + assert_eq!(list.len(), 1, "{response}"); + let first = list[0]["id"].as_str().unwrap().to_string(); + assert_eq!(list[0]["sender"], "hold-sender@example.com"); + assert_eq!(list[0]["subject"], "Held one"); + assert_eq!(list[0]["rules"][0]["name"], "Hold cards"); + assert_eq!( + list[0]["counts"], + json!([{"detector": "words", "count": 1}, {"detector": "payment-card", "count": 1}]) + ); + assert!( + list[0].get("preview").is_none(), + "no preview unless asked for" + ); + + // The sender is told, and it isn't delivered + let notice = received( + &sender, + "Held for review: Held one", + "X-Flow", + Some(&mailbox), + ) + .await; + assert!( + notice[0].1.contains("Card numbers are reviewed first."), + "{notice:?}" + ); + assert!( + received_now(&sender, "Held one", &mailbox) + .await + .iter() + .all(|s| s.starts_with("Held for review")) + ); + + // Not around the review: not from the queue, not by unsending + let (_, response) = call( + &admin, + "x:QueuedMessage/set", + json!({"update": {first.as_str(): {"nextRetry": "2026-01-01T00:00:00Z"}}}), + ) + .await; + assert!( + response["notUpdated"][first.as_str()]["description"] + .as_str() + .unwrap_or_default() + .contains("held for review"), + "{response}" + ); + let (_, response) = call(&admin, "x:QueuedMessage/set", json!({"destroy": [first]})).await; + assert!( + response["notDestroyed"].get(first.as_str()).is_some(), + "{response}" + ); + let (_, response) = call( + &sender, + "EmailSubmission/set", + json!({"update": {submission.as_str(): {"undoStatus": "canceled"}}}), + ) + .await; + assert_eq!( + response["notUpdated"][submission.as_str()]["type"], + "cannotUnsend", + "{response}" + ); + + // Reading it is recorded + let (_, response) = call( + &admin, + "inbuxa:HeldMessage/get", + json!({"ids": [first], "properties": ["id", "preview"]}), + ) + .await; + assert!( + response["list"][0]["preview"] + .as_str() + .unwrap_or_default() + .contains("4242 4242"), + "{response}" + ); + let (_, response) = call( + &admin, + "inbuxa:AuditEvent/query", + json!({"filter": {"targetKind": "inbuxa:HeldMessage", "action": "blobAccess"}}), + ) + .await; + assert_eq!( + response["ids"].as_array().map(|i| i.len()), + Some(1), + "{response}" + ); + + // Rejected: a reason is required; the sender gets the note + let (_, response) = call( + &admin, + "inbuxa:HeldMessage/set", + json!({"update": {first.as_str(): {"decision": "reject"}}}), + ) + .await; + assert!( + response["notUpdated"].get(first.as_str()).is_some(), + "no reason: {response}" + ); + let (_, response) = call( + &admin, + "inbuxa:HeldMessage/set", + json!({"reason": "Card data may not leave by mail", "update": {first.as_str(): {"decision": "reject", "note": "Use the payments portal."}}}), + ) + .await; + assert!( + response["updated"].get(first.as_str()).is_some(), + "{response}" + ); + let notice = received(&sender, "Not sent: Held one", "X-Flow", Some(&mailbox)).await; + assert!( + notice[0].1.contains("Use the payments portal."), + "{notice:?}" + ); + let (_, response) = call(&admin, "x:QueuedMessage/get", json!({"ids": [first]})).await; + assert_eq!( + response["notFound"][0], + first.as_str(), + "gone from the queue: {response}" + ); + + // Released: delivered + let response = submit( + &sender, + &identity, + &mailbox, + &["hold-sender@example.com"], + "Held two", + body, + None, + ) + .await; + assert!(response["created"].get("s").is_some(), "{response}"); + let (_, response) = call(&admin, "inbuxa:HeldMessage/get", json!({"ids": null})).await; + let second = response["list"][0]["id"] + .as_str() + .unwrap_or_else(|| panic!("{response}")) + .to_string(); + let (_, response) = call( + &admin, + "inbuxa:HeldMessage/set", + json!({"reason": "Finance approved", "update": {second.as_str(): {"decision": "release"}}}), + ) + .await; + assert!( + response["updated"].get(second.as_str()).is_some(), + "{response}" + ); + let mut delivered = Vec::new(); + for _ in 0..50 { + delivered = received_now(&sender, "Held two", &mailbox).await; + if delivered.iter().any(|s| s == "Held two") { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + } + assert!(delivered.iter().any(|s| s == "Held two"), "{delivered:?}"); + let (_, response) = call(&admin, "inbuxa:HeldMessage/get", json!({"ids": null})).await; + assert_eq!( + response["list"].as_array().map(|l| l.len()), + Some(0), + "{response}" + ); + + // Both decisions are in the audit log, with their reasons + let (_, response) = call( + &admin, + "inbuxa:AuditEvent/query", + json!({"filter": {"targetKind": "inbuxa:HeldMessage", "action": "update"}}), + ) + .await; + let ids = response["ids"].clone(); + let (_, response) = call(&admin, "inbuxa:AuditEvent/get", json!({"ids": ids})).await; + let reasons: Vec<&str> = response["list"] + .as_array() + .unwrap() + .iter() + .filter_map(|e| e["reason"].as_str()) + .collect(); + assert!( + reasons.contains(&"Card data may not leave by mail") + && reasons.contains(&"Finance approved"), + "{response}" + ); + + // Unreviewed: returned by the daily clean-up once its time is up + let response = submit( + &sender, + &identity, + &mailbox, + &["hold-sender@example.com"], + "Held three", + body, + None, + ) + .await; + assert!(response["created"].get("s").is_some(), "{response}"); + let store = test.server.store(); + let mut record = inbuxa_features::mailflow::held::all(store) + .await + .unwrap() + .pop() + .expect("held"); + record.expires_at = 0; + inbuxa_features::mailflow::held::create(store, &record) + .await + .unwrap(); + assert_eq!(smtp::queue::held::expire(&test.server).await.unwrap(), 1); + assert!( + inbuxa_features::mailflow::held::all(store) + .await + .unwrap() + .is_empty() + ); + let notice = received(&sender, "Not sent: Held three", "X-Flow", Some(&mailbox)).await; + assert!( + notice[0].1.contains("Nobody reviewed it within 7 days"), + "{notice:?}" + ); + let (_, response) = call( + &admin, + "inbuxa:AuditEvent/query", + json!({"filter": {"targetKind": "inbuxa:HeldMessage", "action": "destroy"}}), + ) + .await; + assert_eq!( + response["ids"].as_array().map(|i| i.len()), + Some(1), + "expiry recorded: {response}" + ); + + call(&admin, "inbuxa:MailRule/set", json!({"destroy": [rule]})).await; +} + +/// Subjects in `account` matching `text` right now, not in `drafts`. +async fn received_now(account: &Account, text: &str, drafts: &str) -> Vec { + let (_, response) = call( + account, + "Email/query", + json!({"filter": {"subject": text, "inMailboxOtherThan": [drafts]}}), + ) + .await; + let ids = response["ids"].clone(); + let (_, response) = call( + account, + "Email/get", + json!({"ids": ids, "properties": ["subject"]}), + ) + .await; + response["list"] + .as_array() + .unwrap() + .iter() + .map(|e| e["subject"].as_str().unwrap_or_default().to_string()) + .collect() +} + #[ignore] #[tokio::test(flavor = "multi_thread")] pub async fn mail_rules_tests() { @@ -787,6 +1110,7 @@ pub async fn mail_rules_tests() { self::test(&mut test).await; self::dlp(&mut test).await; self::transport(&mut test).await; + self::hold(&mut test).await; if test.is_reset() { test.temp_dir.delete(); }