Files
inbuxa-server/crates/email/src/sieve/ingest.rs
T
jcoffey-dev 447229f871
ci / fork-checks (pull_request) Successful in 1m4s
ci / build (pull_request) Successful in 8m47s
Lock accounts: keep receiving mail, no sign-in, hand to delegates
A locked account can't sign in (it fails as a wrong password does), its
sessions end on every node, refresh tokens stop working, and its Sieve
scripts forward and reply to nothing. Mail keeps arriving.

Delegates get real ACL grants on the account's mailboxes, calendars,
address books and files at read, organize or full, with the rights they
replaced restored on unlock. Folders made later are granted after the
create and in a daily sweep. Organize delegates can't destroy; send-as
needs organize or full. The JMAP session marks delegated accounts in
urn:inbuxa:jmap.

New inbuxa:AccountLock object with get/set, permissions 665-668, and a
Compliance > Locked Accounts entry in the schema. Lock, unlock and
delegate changes need a reason and are audited; delegate access and
writes are audited too (audit-hold-lock spec AL-1 to AL-12).
2026-09-27 14:46:06 -07:00

820 lines
33 KiB
Rust

/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::{ActiveScript, SeenIdHash, SieveScript};
use crate::{
cache::{MessageCacheFetch, mailbox::MailboxCacheAccess},
mailbox::{INBOX_ID, TRASH_ID, manage::MailboxFnc},
message::{
delivery::{AutogeneratedMessage, IngestRecipient},
ingest::{EmailIngest, IngestEmail, IngestSource, IngestedEmail},
},
};
use common::{
Server, auth::AccessToken, config::mailstore::spamfilter::spam_status,
scripts::plugins::PluginContext,
};
use mail_builder::headers::date::Date;
use mail_parser::{HeaderName, MessageParser};
use sieve::{Envelope, Event, Input, Mailbox, Recipient, Sieve};
use std::{borrow::Cow, sync::Arc};
use std::{future::Future, str::FromStr};
use store::{
Deserialize, Serialize, ValueKey,
ahash::AHashMap,
dispatch::lookup::KeyValue,
write::{
AlignedBytes, Archive, ArchiveVersion, Archiver, BatchBuilder, BlobLink, BlobOp, ValueClass,
},
};
use trc::{AddContext, SieveEvent, SmtpEvent};
use types::{
blob_hash::BlobHash,
collection::Collection,
field::{PrincipalField, SieveField},
id::Id,
keyword::Keyword,
special_use::SpecialUse,
};
struct SieveMessage<'x> {
pub raw_message: Cow<'x, [u8]>,
pub file_into: Vec<u32>,
pub did_file_into: bool,
pub flags: Vec<Keyword>,
}
pub trait SieveScriptIngest: Sync + Send {
#[allow(clippy::too_many_arguments)]
fn sieve_script_ingest(
&self,
access_token: &AccessToken,
blob_hash: &BlobHash,
raw_message: &[u8],
envelope_from: &str,
envelope_from_authenticated: bool,
envelope_to: &IngestRecipient,
session_id: u64,
active_script: ActiveScript,
autogenerated: &mut Vec<AutogeneratedMessage>,
) -> impl Future<Output = trc::Result<IngestedEmail>> + Send;
fn sieve_script_get_active_id(
&self,
account_id: u32,
) -> impl Future<Output = trc::Result<Option<u32>>> + Send;
fn sieve_script_get_active(
&self,
account_id: u32,
) -> impl Future<Output = trc::Result<Option<ActiveScript>>> + Send;
fn sieve_script_get_by_name(
&self,
account_id: u32,
name: &str,
) -> impl Future<Output = trc::Result<Option<Sieve>>> + Send;
fn sieve_script_compile(
&self,
account_id: u32,
document_id: u32,
) -> impl Future<Output = trc::Result<Option<CompiledScript>>> + Send;
}
impl SieveScriptIngest for Server {
#[allow(clippy::blocks_in_conditions)]
async fn sieve_script_ingest(
&self,
access_token: &AccessToken,
blob_hash: &BlobHash,
raw_message: &[u8],
envelope_from: &str,
envelope_from_authenticated: bool,
envelope_to: &IngestRecipient,
session_id: u64,
active_script: ActiveScript,
autogenerated: &mut Vec<AutogeneratedMessage>,
) -> trc::Result<IngestedEmail> {
// Parse message
let message = if let Some(message) = MessageParser::new().parse(raw_message) {
message
} else {
return Err(
trc::EventType::MessageIngest(trc::MessageIngestEvent::Error)
.ctx(trc::Key::Code, 550)
.ctx(trc::Key::Reason, "Failed to parse e-mail message."),
);
};
let received_headers = message
.headers()
.iter()
.filter(|header| matches!(header.name, HeaderName::Received))
.count();
// Obtain mailboxIds
let account_id = access_token.account_id();
let mut cache = self
.get_cached_messages(account_id)
.await
.caused_by(trc::location!())?;
// Create Sieve instance
let orcpt = envelope_to.orcpt_parameter();
let mut instance = self.core.sieve.untrusted_runtime.filter_parsed(message);
// Set account name and email
let account_info = self.account(account_id).await.caused_by(trc::location!())?;
let mail_from = account_info.name().to_string();
instance.set_user_full_name(
account_info
.description()
.unwrap_or_else(|| account_info.name()),
);
instance.set_user_address(&mail_from);
// Set envelope
instance.set_envelope(Envelope::From, envelope_from);
instance.set_envelope(Envelope::To, envelope_to.address.as_str());
if let Some(orcpt) = &orcpt {
instance.set_envelope(Envelope::Orcpt, orcpt.as_str());
}
instance.set_spam_status(spam_status(envelope_to.spam_percentage));
let mut input = Input::script(
active_script.script_name.to_string(),
active_script.script.clone(),
);
let mut do_discard = false;
let mut do_deliver = false;
let mut do_redirect = false;
let mut reject_reason = None;
let mut messages: Vec<SieveMessage> = vec![SieveMessage {
raw_message: raw_message.into(),
file_into: Vec::new(),
flags: Vec::new(),
did_file_into: false,
}];
let mut ingested_message = IngestedEmail {
document_id: 0,
thread_id: 0,
change_id: u64::MAX,
blob_id: Default::default(),
size: raw_message.len(),
imap_uids: Vec::new(),
};
let mut checked_ids: AHashMap<SeenIdHash, bool> = AHashMap::new();
while let Some(event) = instance.run(input) {
match event {
Ok(event) => match event {
Event::IncludeScript { name, .. } => match &name {
sieve::Script::Personal(name_) => {
if let Ok(Some(script)) =
self.sieve_script_get_by_name(account_id, name_).await
{
input = Input::script(name, script);
} else {
input = false.into();
}
}
sieve::Script::Global(name_) => {
if let Some(script) = self.get_untrusted_sieve_script(name_, session_id)
{
input = Input::script(name, script.clone());
} else {
input = false.into();
}
}
},
Event::MailboxExists {
mailboxes,
special_use,
} => {
if !mailboxes.is_empty() {
let mut special_use_ids = Vec::with_capacity(special_use.len());
for role in special_use.iter().map(|v| SpecialUse::parse_use_attr(v)) {
special_use_ids.push(match role {
Some(SpecialUse::Inbox) => INBOX_ID,
Some(SpecialUse::Trash) => TRASH_ID,
Some(role) => cache
.mailbox_by_role(&role)
.map(|m| m.document_id)
.unwrap_or(u32::MAX),
None => u32::MAX,
});
}
let mut result = true;
for mailbox in mailboxes {
match mailbox {
Mailbox::Name(name) => {
if !matches!(
cache.mailbox_by_path(&name),
Some(item) if special_use_ids.is_empty() ||
special_use_ids.contains(&item.document_id)
) {
result = false;
break;
}
}
Mailbox::Id(id) => {
if !matches!(Id::from_str(&id), Ok(id) if
cache.has_mailbox_id(&id.document_id()) &&
(special_use_ids.is_empty() ||
special_use_ids.contains(&id.document_id())))
{
result = false;
break;
}
}
}
}
input = result.into();
} else if !special_use.is_empty() {
let mut result = true;
for role in special_use.iter().map(|v| SpecialUse::parse_use_attr(v)) {
match role {
Some(SpecialUse::Inbox | SpecialUse::Trash) => {}
Some(other) if cache.mailbox_by_role(&other).is_some() => {}
_ => {
result = false;
break;
}
}
}
input = result.into();
} else {
input = false.into();
}
}
Event::DuplicateId { id, expiry, last } => {
let id_hash = SeenIdHash::new(
account_id,
active_script.version.hash().unwrap_or_default(),
&id,
);
if let Some(result) = checked_ids.get(&id_hash) {
input = (*result).into();
} else {
let exists = self
.in_memory_store()
.key_exists(id_hash.key())
.await
.caused_by(trc::location!())?;
if !exists || last {
self.in_memory_store()
.key_set(KeyValue::new(id_hash.key(), vec![]).expires(expiry))
.await
.caused_by(trc::location!())?;
}
checked_ids.insert(id_hash, exists);
input = exists.into();
}
}
Event::Discard => {
do_discard = true;
input = true.into();
}
// inbuxa: AL-4: a locked account answers no sender, so a
// rejection is kept instead; sieve has already cleared
// the implicit keep, so it is filed here
Event::Reject { .. } if access_token.is_locked() => {
if let Some(message) = messages.get_mut(0)
&& !message.file_into.contains(&INBOX_ID)
{
message.file_into.push(INBOX_ID);
}
do_deliver = true;
input = true.into();
}
Event::Reject { reason, .. } => {
reject_reason = reason.into();
do_discard = true;
input = true.into();
}
Event::Keep { flags, message_id } => {
if let Some(message) = messages.get_mut(message_id) {
message.flags = flags.into_iter().map(Keyword::from).collect();
if !message.file_into.contains(&INBOX_ID) {
message.file_into.push(INBOX_ID);
}
do_deliver = true;
} else {
trc::event!(
Sieve(SieveEvent::UnexpectedError),
Details = "Unknown message id.",
MessageId = message_id,
SpanId = session_id
);
}
input = true.into();
}
Event::FileInto {
folder,
flags,
mailbox_id,
special_use,
create,
message_id,
} => {
let mut target_id = u32::MAX;
// Find mailbox by Id
if let Some(mailbox_id) = mailbox_id.and_then(|m| Id::from_str(&m).ok()) {
let mailbox_id = mailbox_id.document_id();
if cache.has_mailbox_id(&mailbox_id) {
target_id = mailbox_id;
}
}
// Find mailbox by role
if target_id == u32::MAX
&& let Some(special_use) =
special_use.as_deref().and_then(SpecialUse::parse_use_attr)
{
match special_use {
SpecialUse::Inbox => {
target_id = INBOX_ID;
}
SpecialUse::Trash => {
target_id = TRASH_ID;
}
role => {
if let Some(item) = cache.mailbox_by_role(&role) {
target_id = item.document_id;
}
}
}
}
// Find mailbox by name
if target_id == u32::MAX {
if !create {
if let Some(m) = cache.mailbox_by_path(&folder) {
target_id = m.document_id;
}
} else if let Some(document_id) = self
.mailbox_create_path(account_id, &folder)
.await
.caused_by(trc::location!())?
{
cache = self
.get_cached_messages(account_id)
.await
.caused_by(trc::location!())?;
target_id = document_id;
}
}
// Default to Inbox
if target_id == u32::MAX {
target_id = INBOX_ID;
}
if let Some(message) = messages.get_mut(message_id) {
message.flags = flags.into_iter().map(Keyword::from).collect();
if !message.file_into.contains(&target_id) {
message.file_into.push(target_id);
}
message.did_file_into = true;
do_deliver = true;
} else {
trc::event!(
Sieve(SieveEvent::UnexpectedError),
Details = "Unknown message id.",
MessageId = message_id,
SpanId = session_id
);
}
input = true.into();
}
// inbuxa: AL-4: a locked account sends nothing on its
// own: no redirect, vacation reply or notification. An
// unsent redirect leaves the message to be kept.
Event::SendMessage { .. } if access_token.is_locked() => {
trc::event!(
Sieve(SieveEvent::ActionReject),
Details = "Account is locked: nothing is sent",
SpanId = session_id
);
input = true.into();
}
Event::SendMessage {
recipient,
message_id,
..
} => {
input = true.into();
if let Some(message) = messages.get(message_id) {
if received_headers >= self.core.sieve.max_received_headers {
trc::event!(
Smtp(SmtpEvent::LoopDetected),
From = mail_from.clone(),
Total = received_headers,
Limit = self.core.sieve.max_received_headers,
SpanId = session_id,
);
continue;
}
let recipients: Vec<String> = match recipient {
Recipient::Address(rcpt) => vec![rcpt],
Recipient::Group(rcpts) => rcpts,
Recipient::List(_) => {
// Not yet implemented
continue;
}
};
if message.raw_message.len() <= self.core.email.mail_max_size {
trc::event!(
Sieve(SieveEvent::SendMessage),
From = mail_from.clone(),
To = recipients
.iter()
.map(|r| trc::Value::String(r.as_str().into()))
.collect::<Vec<_>>(),
Size = message.raw_message.len(),
SpanId = session_id
);
let mut raw_message =
Vec::with_capacity(160 + message.raw_message.len());
write_received_header(
&mut raw_message,
&self.core.network.server_name,
session_id,
);
raw_message.extend_from_slice(message.raw_message.as_ref());
autogenerated.push(AutogeneratedMessage {
sender_address: mail_from.clone(),
recipients,
message: raw_message,
});
do_redirect = true;
} else {
trc::event!(
Sieve(SieveEvent::MessageTooLarge),
From = mail_from.clone(),
To = recipients
.iter()
.map(|r| trc::Value::String(r.as_str().into()))
.collect::<Vec<_>>(),
Size = message.raw_message.len(),
Limit = self.core.email.mail_max_size,
SpanId = session_id,
);
}
} else {
trc::event!(
Sieve(SieveEvent::UnexpectedError),
Details = "Unknown message id.",
MessageId = message_id,
SpanId = session_id
);
continue;
}
}
Event::ListContains { .. }
| Event::Notify { .. }
| Event::SetEnvelope { .. } => {
// Not allowed
input = false.into();
}
Event::Function { id, arguments } => {
input = self
.core
.run_plugin(
id,
PluginContext {
session_id,
server: self,
message: instance.message(),
modifications: &mut Vec::new(),
access_token: access_token.into(),
arguments,
},
)
.await;
}
Event::CreatedMessage { message, .. } => {
messages.push(SieveMessage {
raw_message: message.into(),
file_into: Vec::new(),
flags: Vec::new(),
did_file_into: false,
});
input = true.into();
}
},
#[cfg(feature = "test_mode")]
Err(sieve::runtime::RuntimeError::ScriptErrorMessage(err)) => {
panic!("Sieve test failed: {}", err);
}
Err(err) => {
trc::event!(
Sieve(SieveEvent::RuntimeError),
Reason = err.to_string(),
SpanId = session_id
);
input = true.into();
}
}
}
// Fail-safe, no discard and no keep seen, assume that something went wrong and file anyway.
if !do_deliver && !do_discard && !do_redirect {
messages[0].file_into.push(INBOX_ID);
}
// Deliver messages
let mut last_temp_error = None;
let mut has_delivered = false;
for (message_id, sieve_message) in messages.into_iter().enumerate() {
if !sieve_message.file_into.is_empty() {
// Parse message if needed
let (blob_hash, message) = if message_id == 0 && !instance.has_message_changed() {
(blob_hash.into(), instance.take_message())
} else if let Some(message) =
MessageParser::new().parse(sieve_message.raw_message.as_ref())
{
(None, message)
} else {
trc::event!(
Sieve(SieveEvent::UnexpectedError),
Details = "Failed to parse Sieve generated message.",
SpanId = session_id
);
continue;
};
// Deliver message
match self
.email_ingest(IngestEmail {
raw_message: &sieve_message.raw_message,
blob_hash,
message: message.into(),
access_token,
mailbox_ids: sieve_message.file_into,
keywords: sieve_message.flags,
received_at: None,
source: IngestSource::Smtp {
deliver_to: envelope_to.address.as_str(),
is_sender_authenticated: envelope_from_authenticated,
is_spam: envelope_to.is_spam() && !sieve_message.did_file_into,
},
session_id,
})
.await
{
Ok(ingested_message_) => {
has_delivered = true;
ingested_message = ingested_message_;
}
Err(err) => {
last_temp_error = err.into();
}
}
}
}
if let Some(reject_reason) = reject_reason {
Err(
trc::EventType::MessageIngest(trc::MessageIngestEvent::Error)
.ctx(trc::Key::Code, 571)
.ctx(trc::Key::Reason, reject_reason),
)
} else if has_delivered || last_temp_error.is_none() {
Ok(ingested_message)
} else {
// There were problems during delivery
#[allow(clippy::unnecessary_unwrap)]
Err(last_temp_error.unwrap())
}
}
async fn sieve_script_get_active_id(&self, account_id: u32) -> trc::Result<Option<u32>> {
self.store()
.get_value::<u32>(ValueKey {
account_id,
collection: Collection::Principal.into(),
document_id: 0,
class: ValueClass::Property(PrincipalField::ActiveScriptId.into()),
})
.await
.caused_by(trc::location!())
}
async fn sieve_script_get_active(&self, account_id: u32) -> trc::Result<Option<ActiveScript>> {
// Find the currently active script
if let Some(document_id) = self
.store()
.get_value::<u32>(ValueKey {
account_id,
collection: Collection::Principal.into(),
document_id: 0,
class: ValueClass::Property(PrincipalField::ActiveScriptId.into()),
})
.await
.caused_by(trc::location!())?
{
if let Some(script) = self.sieve_script_compile(account_id, document_id).await? {
Ok(Some(ActiveScript {
document_id,
script: Arc::new(script.script),
script_name: script.name,
version: script.version,
}))
} else {
Ok(None)
}
} else {
Ok(None)
}
}
async fn sieve_script_get_by_name(
&self,
account_id: u32,
name: &str,
) -> trc::Result<Option<Sieve>> {
// Find the script by name
if let Some(document_id) = self
.document_ids_matching(
account_id,
Collection::SieveScript,
SieveField::Name,
name.to_lowercase().as_bytes(),
)
.await
.caused_by(trc::location!())?
.min()
{
self.sieve_script_compile(account_id, document_id)
.await
.map(|script| script.map(|s| s.script))
} else {
Ok(None)
}
}
#[allow(clippy::blocks_in_conditions)]
async fn sieve_script_compile(
&self,
account_id: u32,
document_id: u32,
) -> trc::Result<Option<CompiledScript>> {
// Obtain script object
let Some(script_object) = self
.store()
.get_value::<Archive<AlignedBytes>>(ValueKey::archive(
account_id,
Collection::SieveScript,
document_id,
))
.await?
else {
return Ok(None);
};
// Obtain the sieve script length
let version = script_object.version;
let unarchived_script = script_object
.unarchive::<SieveScript>()
.caused_by(trc::location!())?;
let script_offset = u32::from(unarchived_script.size) as usize;
// Obtain the sieve script blob
let script_bytes = self
.core
.storage
.blob
.get_blob(unarchived_script.blob_hash.0.as_ref(), 0..usize::MAX)
.await
.caused_by(trc::location!())?
.ok_or_else(|| {
trc::StoreEvent::NotFound
.into_err()
.caused_by(trc::location!())
.document_id(document_id)
})?;
// Obtain the precompiled script
if let Some(script) = script_bytes.get(script_offset..).and_then(|bytes| {
<Archive<AlignedBytes> as Deserialize>::deserialize(bytes)
.ok()?
.deserialize::<Sieve>()
.ok()
}) {
Ok(Some(CompiledScript {
script,
name: unarchived_script.name.as_str().into(),
version,
}))
} else {
// Deserialization failed, probably because the script compiler version changed
match self.core.sieve.untrusted_compiler.compile(
script_bytes.get(0..script_offset).ok_or_else(|| {
trc::StoreEvent::NotFound
.into_err()
.caused_by(trc::location!())
.document_id(document_id)
})?,
) {
Ok(sieve) => {
// Store updated compiled sieve script
let sieve = Archiver::new(sieve).untrusted();
let compiled_bytes = sieve.serialize().caused_by(trc::location!())?;
let mut updated_sieve_bytes =
Vec::with_capacity(script_offset + compiled_bytes.len());
updated_sieve_bytes.extend_from_slice(&script_bytes[0..script_offset]);
updated_sieve_bytes.extend_from_slice(&compiled_bytes);
// Store updated blob
let (new_blob_hash, new_blob_hold) = self
.put_temporary_blob(account_id, &updated_sieve_bytes, 60)
.await?;
let mut new_script_object =
rkyv::deserialize(unarchived_script).caused_by(trc::location!())?;
let blob_hash =
std::mem::replace(&mut new_script_object.blob_hash, new_blob_hash.clone());
let new_archive = Archiver::new(new_script_object);
// Update script object
let mut batch = BatchBuilder::new();
batch
.with_account_id(account_id)
.with_collection(Collection::SieveScript)
.with_document(document_id)
.assert_value(SieveField::Archive, &script_object)
.set(
SieveField::Archive,
new_archive.serialize().caused_by(trc::location!())?,
)
.clear(BlobOp::Link {
hash: blob_hash,
to: BlobLink::Document,
})
.set(
BlobOp::Link {
hash: new_blob_hash,
to: BlobLink::Document,
},
Vec::new(),
)
.clear(new_blob_hold);
self.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
Ok(Some(CompiledScript {
script: sieve.into_inner(),
name: new_archive.into_inner().name,
version,
}))
}
Err(error) => Err(trc::StoreEvent::UnexpectedError
.caused_by(trc::location!())
.reason(error)
.details("Failed to compile Sieve script")),
}
}
}
}
fn write_received_header(buf: &mut Vec<u8>, hostname: &str, id: u64) {
buf.extend_from_slice(b"Received: from localhost (localhost [127.0.0.1])\r\n\tby ");
buf.extend_from_slice(hostname.as_bytes());
buf.extend_from_slice(concat!(" (", types::brand!(), " SMTP) with LMTP id ").as_bytes());
buf.extend_from_slice(format!("{id:X}").as_bytes());
buf.extend_from_slice(b";\r\n\t");
buf.extend_from_slice(Date::now().to_rfc822().as_bytes());
buf.extend_from_slice(b"\r\n");
}
pub struct CompiledScript {
pub script: Sieve,
pub name: String,
pub version: ArchiveVersion,
}