/* * 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 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, pub did_file_into: bool, pub flags: Vec, } 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, ) -> impl Future> + Send; fn sieve_script_get_active_id( &self, account_id: u32, ) -> impl Future>> + Send; fn sieve_script_get_active( &self, account_id: u32, ) -> impl Future>> + Send; fn sieve_script_get_by_name( &self, account_id: u32, name: &str, ) -> impl Future>> + Send; fn sieve_script_compile( &self, account_id: u32, document_id: u32, ) -> impl Future>> + 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, ) -> trc::Result { // 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 = 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 = 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(); } 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(); } 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 = 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::>(), 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::>(), 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> { self.store() .get_value::(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> { // Find the currently active script if let Some(document_id) = self .store() .get_value::(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> { // 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> { // Obtain script object let Some(script_object) = self .store() .get_value::>(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::() .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| { as Deserialize>::deserialize(bytes) .ok()? .deserialize::() .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, 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, }