Merge pull request 'DLP phase 3: hold for review' (#108) from feature/dlp-hold into main
ci / fork-checks (push) Successful in 19s
ci / build (push) Canceled after 4m22s

This commit was merged in pull request #108.
This commit is contained in:
2026-09-29 01:43:44 +00:00
25 changed files with 1561 additions and 57 deletions
+183
View File
@@ -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<u32>,
#[serde(default)]
pub tenant_id: Option<u32>,
pub recipients: Vec<String>,
pub subject: String,
pub size: u64,
pub rules: Vec<HeldRule>,
/// 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>(T);
impl<T: SerdeSerialize> Serialize for Json<T> {
fn serialize(&self) -> trc::Result<Vec<u8>> {
serde_json::to_vec(&self.0).map_err(|err| {
trc::StoreEvent::UnexpectedError
.into_err()
.details("Failed to serialize held message")
.reason(err)
})
}
}
impl<T: DeserializeOwned + Sync + Send> Deserialize for Json<T> {
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
serde_json::from_slice(bytes).map(Json).map_err(|err| {
trc::StoreEvent::DataCorruption
.into_err()
.details("Invalid 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<ValueClass> {
ValueKey::from(class(queue_id))
}
pub async fn get(data: &Store, queue_id: u64) -> trc::Result<Option<Held>> {
Ok(data
.get_value::<Json<Held>>(key(queue_id))
.await
.caused_by(trc::location!())?
.map(|Json(held)| held))
}
pub async fn is_held(data: &Store, queue_id: u64) -> trc::Result<bool> {
get(data, queue_id).await.map(|held| held.is_some())
}
/// Every held message, oldest first.
pub async fn all(data: &Store) -> trc::Result<Vec<Held>> {
let mut held = Vec::new();
data.iterate(IterateParams::new(key(0), key(u64::MAX)), |_, value| {
if let Ok(Json(record)) = Json::<Held>::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: "[email protected]".into(),
account_id: Some(7),
tenant_id: None,
recipients: vec!["[email protected]".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::<Held>(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);
}
}
+1
View File
@@ -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;
@@ -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<Self> {
// 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<Self> {
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<Self, Self::Err> {
HeldMessageProperty::parse(s).ok_or(())
}
}
impl Element for HeldMessageValue {
type Property = HeldMessageProperty;
fn try_parse<P>(key: &Key<'_, Self::Property>, value: &str) -> Option<Self> {
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<String>,
}
impl<'de> DeserializeArguments<'de> for HeldMessageSetArguments {
fn deserialize_argument<A>(&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::<serde::de::IgnoredAny>()?;
}
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<Id> for HeldMessageValue {
fn from(id: Id) -> Self {
HeldMessageValue::Id(id)
}
}
impl JmapObjectId for HeldMessageValue {
fn as_id(&self) -> Option<Id> {
match self {
HeldMessageValue::Id(id) => Some(*id),
}
}
fn as_any_id(&self) -> Option<AnyId> {
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<Id> {
None
}
fn as_any_id(&self) -> Option<AnyId> {
None
}
fn as_id_ref(&self) -> Option<&str> {
None
}
fn try_set_id(&mut self, _: AnyId) -> bool {
false
}
}
+1
View File
@@ -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
+3
View File
@@ -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)
}
@@ -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)?
}
+8 -1
View File
@@ -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",
+2
View File
@@ -124,6 +124,7 @@ pub enum GetRequestMethod {
AccountLock(Box<GetRequest<crate::object::inbuxa_account_lock::AccountLock>>),
LegalHold(Box<GetRequest<crate::object::inbuxa_legal_hold::LegalHold>>),
MailRule(Box<GetRequest<crate::object::inbuxa_mail_rule::MailRule>>),
HeldMessage(Box<GetRequest<crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<GetRequest<crate::object::inbuxa_hold_export::HoldExport>>),
ProtocolPolicy(Box<GetRequest<crate::object::inbuxa_protocol_policy::ProtocolPolicy>>),
TenantProtocolPolicy(
@@ -160,6 +161,7 @@ pub enum SetRequestMethod<'x> {
AccountLock(Box<SetRequest<'x, crate::object::inbuxa_account_lock::AccountLock>>),
LegalHold(Box<SetRequest<'x, crate::object::inbuxa_legal_hold::LegalHold>>),
MailRule(Box<SetRequest<'x, crate::object::inbuxa_mail_rule::MailRule>>),
HeldMessage(Box<SetRequest<'x, crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<SetRequest<'x, crate::object::inbuxa_hold_export::HoldExport>>),
ProtocolPolicy(Box<SetRequest<'x, crate::object::inbuxa_protocol_policy::ProtocolPolicy>>),
TenantProtocolPolicy(
+15
View File
@@ -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)),
+14
View File
@@ -111,6 +111,7 @@ pub enum GetResponseMethod {
AccountLock(GetResponse<crate::object::inbuxa_account_lock::AccountLock>),
LegalHold(GetResponse<crate::object::inbuxa_legal_hold::LegalHold>),
MailRule(GetResponse<crate::object::inbuxa_mail_rule::MailRule>),
HeldMessage(GetResponse<crate::object::inbuxa_held_message::HeldMessage>),
HoldExport(GetResponse<crate::object::inbuxa_hold_export::HoldExport>),
ProtocolPolicy(GetResponse<crate::object::inbuxa_protocol_policy::ProtocolPolicy>),
TenantProtocolPolicy(
@@ -147,6 +148,7 @@ pub enum SetResponseMethod {
AccountLock(Box<SetResponse<crate::object::inbuxa_account_lock::AccountLock>>),
LegalHold(Box<SetResponse<crate::object::inbuxa_legal_hold::LegalHold>>),
MailRule(Box<SetResponse<crate::object::inbuxa_mail_rule::MailRule>>),
HeldMessage(Box<SetResponse<crate::object::inbuxa_held_message::HeldMessage>>),
HoldExport(Box<SetResponse<crate::object::inbuxa_hold_export::HoldExport>>),
Explanation(Box<SetResponse<crate::object::inbuxa_explanation::Explanation>>),
ProtocolPolicy(Box<SetResponse<crate::object::inbuxa_protocol_policy::ProtocolPolicy>>),
@@ -801,6 +803,18 @@ impl<'x> From<SetResponse<crate::object::inbuxa_account_lock::AccountLock>> for
}
// inbuxa: legal hold
impl<'x> From<GetResponse<crate::object::inbuxa_held_message::HeldMessage>> for ResponseMethod<'x> {
fn from(value: GetResponse<crate::object::inbuxa_held_message::HeldMessage>) -> Self {
ResponseMethod::Get(GetResponseMethod::HeldMessage(value))
}
}
impl<'x> From<SetResponse<crate::object::inbuxa_held_message::HeldMessage>> for ResponseMethod<'x> {
fn from(value: SetResponse<crate::object::inbuxa_held_message::HeldMessage>) -> Self {
ResponseMethod::Set(SetResponseMethod::HeldMessage(Box::new(value)))
}
}
impl<'x> From<GetResponse<crate::object::inbuxa_mail_rule::MailRule>> for ResponseMethod<'x> {
fn from(value: GetResponse<crate::object::inbuxa_mail_rule::MailRule>) -> Self {
ResponseMethod::Get(GetResponseMethod::MailRule(value))
+11
View File
@@ -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
+24
View File
@@ -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();
+1
View File
@@ -430,6 +430,7 @@ impl IntermediateChangesResponse {
| MethodObject::LegalHold
| MethodObject::HoldExport
| MethodObject::MailRule
| MethodObject::HeldMessage
| MethodObject::ProtocolPolicy
| MethodObject::TenantProtocolPolicy
| MethodObject::Registry(_) => unreachable!(),
+325
View File
@@ -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<Option<String>> {
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<HeldMessage>,
) -> trc::Result<GetResponse<HeldMessage>> {
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<P> {
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<SetResponse<HeldMessage>> {
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::<String>());
}
}
(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)
}
+1
View File
@@ -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;
@@ -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<Property> {
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<RegistrySetResponse<'_>> {
@@ -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;
+11
View File
@@ -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
@@ -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 {
+42 -10
View File
@@ -740,19 +740,36 @@ impl<T: SessionStream> Session<T> {
// 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 } => {
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)) {
if !self
.data
.rcpt_to
.iter()
.any(|r| r.address_lcase.eq_ignore_ascii_case(&address))
{
self.data.rcpt_to.push(SessionAddress::new(address));
}
}
@@ -764,12 +781,6 @@ impl<T: SessionStream> Session<T> {
}
}
}
}
super::mailflow::Checked::Refuse(reply, refusal) => {
self.data.dlp_refusal = refusal;
return reply.into();
}
}
// Build message
let mail_from = self.data.mail_from.clone().unwrap();
@@ -850,6 +861,19 @@ impl<T: SessionStream> Session<T> {
.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::<Vec<_>>(),
message.message.size,
)
});
if message
.queue(
QueueParams::new(raw_message, self.data.session_id, &self.server)
@@ -864,9 +888,17 @@ impl<T: SessionStream> Session<T> {
{
self.state = State::Accepted(queue_id);
self.data.messages_sent += 1;
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()
}
+128 -24
View File
@@ -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<u8>, Option<DlpRefusal>),
/// Queue it held for review (§2.6), with any transport changes.
Hold {
draft: HeldDraft,
message: Option<Vec<u8>>,
envelope: Vec<EnvelopeChange>,
},
}
/// What the review record will say, once the message has a queue id.
pub struct HeldDraft {
pub subject: String,
pub rules: Vec<HeldRule>,
pub counts: Vec<(String, usize)>,
pub notify_sender: bool,
}
/// What a transport rule changes about where a message goes.
@@ -219,6 +234,7 @@ impl<T: SessionStream> Session<T> {
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<T: SessionStream> Session<T> {
.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(
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)),
),
Decision::Pass => {
);
}
// 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<Vec<u8>> =
@@ -344,12 +389,18 @@ impl<T: SessionStream> Session<T> {
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<T: SessionStream> Session<T> {
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<T: SessionStream> Session<T> {
.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<String>,
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<T: SessionStream> Session<T> {
.collect::<Vec<_>>()
.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",
+177
View File
@@ -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<bool> {
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::<Message>()?.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<bool> {
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::<Vec<_>>()
.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::<String>, 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<usize> {
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)
}
+3
View File
@@ -2,6 +2,8 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
*
* 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;
+13 -2
View File
@@ -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
+15
View File
@@ -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"
+324
View File
@@ -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("[email protected]");
let sender = admin
.create_user_account(
"[email protected]",
"hold-sender-secret-7703",
"Hold sender",
&[],
vec![],
)
.await;
let (_, response) = call(
&sender,
"Identity/set",
json!({"create": {"i": {"name": "Sender", "email": "[email protected]"}}}),
)
.await;
let identity = response["created"]["i"]["id"].as_str().unwrap().to_string();
let (_, response) = call(
&sender,
"Mailbox/set",
json!({"create": {"m": {"name": "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,
&["[email protected]"],
"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"], "[email protected]");
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,
&["[email protected]"],
"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,
&["[email protected]"],
"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<String> {
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();
}