/* * 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 crate::{ api::query::QueryResponseBuilder, registry::{ mapping::{RegistryGetResponse, RegistryQueryResponse, RegistrySetResponse}, query::RegistryQueryFilters, }, }; use common::{ Server, config::smtp::queue::{ArchivedQueueExpiry, QueueName}, ipc::QueueEvent, }; use jmap_proto::{error::set::SetError, object::registry::RegistryComparator, types::state::State}; use jmap_tools::{JsonPointer, JsonPointerItem, Key}; use registry::{ jmap::{IntoValue, JsonPointerPatch, RegistryJsonPatch}, schema::{ enums::{DeliveryErrorType, MessageFlag, RecipientFlag}, prelude::Property, structs::{ DeliveryError, QueueExpiry, QueueExpiryAttempts, QueueExpiryTtl, QueuedMessage, QueuedRecipient, RecipientStatus, ServerResponse, }, }, types::{datetime::UTCDateTime, ipaddr::IpAddr, map::Map}, }; use smtp::queue::{ self, ArchivedError, ArchivedErrorDetails, ArchivedMessage, ArchivedStatus, ErrorDetails, FROM_AUTHENTICATED, FROM_AUTOGENERATED, FROM_DSN, FROM_REPORT, FROM_UNAUTHENTICATED, FROM_UNAUTHENTICATED_DMARC, Message, MessageWrapper, RCPT_DSN_SENT, Schedule, Status, rcpt_spam_percentage, spool::SmtpSpool, }; use std::str::FromStr; use store::{ Deserialize, IterateParams, U64_LEN, ValueKey, ahash::AHashSet, registry::RegistryFilterOp, write::{AlignedBytes, Archive, QueueClass, ValueClass, key::DeserializeBigEndian, now}, }; use trc::AddContext; use types::{blob::BlobId, blob_hash::BlobHash, id::Id}; use utils::map::vec_map::VecMap; pub(crate) async fn queued_message_set( mut set: RegistrySetResponse<'_>, ) -> trc::Result> { // Fail all create operations set.fail_all_create("Queued messages cannot be created"); // Obtain tenant domains let tenant_domains = if let Some(tenant_id) = set.access_token.tenant_id() { Some(tenant_domains(set.server, tenant_id).await?) } else { None }; // Process update operations let mut refresh_queue = false; 'outer: for (id, value) in set.update.drain(..) { let queue_id = id.id(); let Some(archive) = set.server.read_message_archive(queue_id).await? else { set.response.not_updated.append(id, SetError::not_found()); continue; }; let archived_message = archive.to_unarchived::()?; // inbuxa: MT-5 if !tenant_domains .as_ref() .is_none_or(|domains| tenant_sees_archived(domains, archived_message.inner)) { set.response.not_updated.append(id, SetError::not_found()); continue; } // Process patches let mut message = map_message(archived_message.inner); message.next_retry = None; for (key, value) in value.into_expanded_object() { let ptr = match key { Key::Property(prop) => { JsonPointer::new(vec![JsonPointerItem::Key(Key::Property(prop))]) } Key::Borrowed(other) => JsonPointer::parse(other), Key::Owned(other) => JsonPointer::parse(&other), }; if let Err(err) = message.patch(JsonPointerPatch::new(&ptr).with_create(false), value) { set.response.not_updated.append(id, err.into()); continue 'outer; } } let set_next_retry = message.next_retry; // Process changes let mut has_changes = false; let mut modified_rcpts = AHashSet::new(); let mut queued_message = archived_message.deserialize()?; let prev_events = queued_message.next_events(); if queued_message.env_id.as_deref() != message.env_id.as_deref() { queued_message.env_id = message.env_id.as_deref().map(|v| v.into()); has_changes = true; } if queued_message.priority as i64 != message.priority { queued_message.priority = message.priority as i16; has_changes = true; } for (idx, rcpt) in queued_message.recipients.iter_mut().enumerate() { if !message .recipients .iter() .any(|(address, _)| address.as_str() == rcpt.address.as_ref()) { rcpt.status = Status::PermanentFailure(ErrorDetails { entity: "localhost".into(), details: queue::Error::Io("Delivery canceled.".into()), }); has_changes = true; modified_rcpts.insert(idx); } } for (address, rcpt) in message.recipients.into_iter() { let Some((idx, queued_rcpt)) = queued_message .recipients .iter_mut() .enumerate() .find(|(_, r)| r.address.as_ref() == address.as_str()) else { set.response.not_updated.append( id, SetError::invalid_properties() .with_description(format!("Recipient '{address}' does not exist")), ); continue 'outer; }; let mut changed = false; if rcpt.orcpt.as_deref() != queued_rcpt.orcpt.as_deref() { queued_rcpt.orcpt = rcpt.orcpt.as_deref().map(|v| v.into()); changed = true; } let expiry = match rcpt.expires { QueueExpiry::Ttl(ttl) => common::config::smtp::queue::QueueExpiry::Ttl( (ttl.expires_at.timestamp() as u64).saturating_sub(queued_message.created), ), QueueExpiry::Attempts(attempts) => { common::config::smtp::queue::QueueExpiry::Attempts( attempts.expires_attempts as u32, ) } }; if expiry != queued_rcpt.expires { queued_rcpt.expires = expiry; changed = true; } for (due, count, field) in [ (rcpt.retry_due, rcpt.retry_count, &mut queued_rcpt.retry), (rcpt.notify_due, rcpt.notify_count, &mut queued_rcpt.notify), ] { let schedule = Schedule { due: due.timestamp() as u64, inner: count as u32, }; if schedule != *field { *field = schedule; changed = true; } } if let Some(next_retry) = set_next_retry && !matches!(queued_rcpt.status, Status::PermanentFailure(_)) { let new_due = next_retry.timestamp() as u64; if queued_rcpt.retry.due != new_due { queued_rcpt.retry.due = new_due; changed = true; } } if matches!(rcpt.status, RecipientStatus::Scheduled) && !matches!(queued_rcpt.status, Status::Scheduled) { queued_rcpt.status = Status::Scheduled; changed = true; } if changed { has_changes = true; modified_rcpts.insert(idx); } } if has_changes { // Delete message if there are no pending deliveries let message = MessageWrapper::new(queued_message, queue_id, QueueName::default()); let is_success = if message.message.recipients.iter().any(|recipient| { matches!( recipient.status, Status::TemporaryFailure(_) | Status::Scheduled ) }) { message .save_registry_changes(set.server, prev_events, modified_rcpts) .await } else { message.remove_registry(set.server, prev_events).await }; if !is_success { set.response.not_updated.append( id, SetError::forbidden().with_description("Queue update operation failed"), ); continue; } refresh_queue = true; } set.response.updated.append(id, None); } if refresh_queue { let _ = set .server .inner .ipc .queue_tx .send(QueueEvent::Refresh) .await; } // Process destroy operations for id in set.destroy.drain(..) { let Some(message) = set.server.read_message(id.id(), QueueName::default()).await else { set.response.not_destroyed.append(id, SetError::not_found()); continue; }; // inbuxa: MT-5 if tenant_domains .as_ref() .is_none_or(|domains| tenant_sees(domains, &message.message)) { if message.remove(set.server, None).await { set.response.destroyed.push(id); } else { set.response.not_destroyed.append( id, SetError::forbidden().with_description("Queue delete operation failed"), ); } } else { set.response.not_destroyed.append(id, SetError::not_found()); } } Ok(set) } pub(crate) async fn queued_message_get( mut get: RegistryGetResponse<'_>, ) -> trc::Result> { let client_ids = get.ids.is_some(); let ids = if let Some(ids) = get.ids.take() { ids } else { queued_ids(get.server, get.server.core.jmap.get_max_objects) .await? .into_iter() .map(Id::from) .collect() }; // Obtain tenant domains let tenant_domains = if let Some(tenant_id) = get.access_token.tenant_id() { Some(tenant_domains(get.server, tenant_id).await?) } else { None }; for id in ids { let Some(message_archive) = get.server.read_message_archive(id.id()).await? else { if client_ids { get.not_found(id); } continue; }; let message_in = message_archive.unarchive::()?; // inbuxa: MT-5 if tenant_domains .as_ref() .is_none_or(|domains| tenant_sees_archived(domains, message_in)) { get.insert(id, map_message(message_in).into_value()); } else if client_ids { get.not_found(id); } } Ok(get) } pub(crate) async fn queued_message_query( mut req: RegistryQueryResponse<'_>, ) -> trc::Result { let mut due_from = 0u64; let mut due_to = u64::MAX; let mut queue_name = None; let mut filter_text = None; let mut filter_from = None; let mut filter_to = None; // Obtain tenant domains let tenant_domains = if let Some(tenant_id) = req.access_token.tenant_id() { Some(tenant_domains(req.server, tenant_id).await?) } else { None }; req.request .extract_filters(|property, op, value| match property { Property::Due => { if let Some(due) = value.as_str().and_then(|s| UTCDateTime::from_str(s).ok()) { let due = due.timestamp() as u64; let (from, to) = match op { RegistryFilterOp::Equal => (due, due), RegistryFilterOp::GreaterThan => (due + 1, u64::MAX), RegistryFilterOp::GreaterEqualThan => (due, u64::MAX), RegistryFilterOp::LowerThan => (0, due - 1), RegistryFilterOp::LowerEqualThan => (0, due), _ => return false, }; // Intersect with existing range due_from = due_from.max(from); due_to = due_to.min(to); due_from <= due_to } else { false } } Property::QueueName => { if let Some(value) = value.as_str().and_then(QueueName::new) { queue_name = Some(value); true } else { false } } Property::ReturnPath => { if let serde_json::Value::String(name) = value { filter_from = Some(name); true } else { false } } Property::To => { if let serde_json::Value::String(name) = value { filter_to = Some(name); true } else { false } } Property::Text => { if let serde_json::Value::String(name) = value { filter_text = Some(name); true } else { false } } _ => false, })?; if req .request .sort .as_ref() .and_then(|sort| sort.first()) .is_some_and(|comp| !matches!(comp.property, RegistryComparator::Property(Property::Due))) { return Err(trc::JmapEvent::UnsupportedSort .into_err() .details("Only sorting by 'due' is supported for queued messages".to_string())); } let params = req .request .extract_parameters(req.server.core.jmap.query_max_results, None)?; let has_filters = filter_text.is_some() || filter_from.is_some() || filter_to.is_some(); if has_filters || tenant_domains.is_some() { let from_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(0))); let to_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(u64::MAX))); let mut results = Vec::with_capacity(8); req.server .core .storage .data .iterate( IterateParams::new(from_key, to_key).ascending(), |key, value| { let message_ = as Deserialize>::deserialize(value) .add_context(|ctx| ctx.ctx(trc::Key::Key, key))?; let message = message_ .unarchive::() .add_context(|ctx| ctx.ctx(trc::Key::Key, key))?; if let Some(due) = message.next_delivery_event(queue_name) // inbuxa: MT-5 && tenant_domains .as_ref() .is_none_or(|domains| tenant_sees_archived(domains, message)) && (due_from..=due_to).contains(&due) && queue_name .as_ref() .is_none_or(|q| message.recipients.iter().any(|r| &r.queue == q)) && (!has_filters || (filter_text .as_ref() .map(|text| { message.return_path.contains(text) || message .recipients .iter() .any(|r| r.address().contains(text)) }) .unwrap_or_else(|| { filter_from .as_ref() .is_none_or(|from| message.return_path.contains(from)) && filter_to.as_ref().is_none_or(|to| { message .recipients .iter() .any(|r| r.address().contains(to)) }) }))) { results.push((key.deserialize_be_u64(0)?, due)); } Ok(true) }, ) .await .caused_by(trc::location!())?; // Build response let mut response = QueryResponseBuilder::new( results.len(), req.server.core.jmap.query_max_results, State::Initial, &req.request, ); if params.sort_ascending { results.sort_by_key(|(_, due)| *due); } else { results.sort_by_key(|(_, due)| u64::MAX - *due); } for (id, _) in results { if !response.add_id(id.into()) { break; } } Ok(response) } else { // Build response let mut response = QueryResponseBuilder::new( req.server.core.jmap.query_max_results, req.server.core.jmap.query_max_results, State::Initial, &req.request, ); let mut total = 0; if let Some(anchor) = req.request.anchor { let anchor_id = anchor.id(); if let Some(archive) = req.server.read_message_archive(anchor_id).await? && let Ok(archived) = archive.unarchive::() && let Some(anchor_due) = archived.next_delivery_event(queue_name) && anchor_due >= due_from && anchor_due <= due_to { if params.sort_ascending { due_from = anchor_due; } else { due_to = anchor_due; } } } let from_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( store::write::QueueEvent { due: due_from, queue_id: 0, queue_name: [0; 8], }, ))); let to_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( store::write::QueueEvent { due: due_to, queue_id: u64::MAX, queue_name: [u8::MAX; 8], }, ))); let mut seen_ids = AHashSet::with_capacity(8); req.server .store() .iterate( IterateParams::new(from_key, to_key) .set_ascending(params.sort_ascending) .no_values(), |key, _| { let id = key.deserialize_be_u64(U64_LEN)?; if queue_name.is_none_or(|queue_name| { queue_name.as_slice() == key.get(U64_LEN * 2..).unwrap_or_default() }) && seen_ids.insert(id) { total += 1; if response.response.total.is_some() { if !response.is_full() { response.add_id(id.into()); } Ok(true) } else { Ok(response.add_id(id.into())) } } else { Ok(true) } }, ) .await .caused_by(trc::location!())?; if response.response.total.is_some() { response.response.total = Some(total); } if let Some(limit) = response.response.limit && total < limit { response.response.limit = None; } Ok(response) } } // inbuxa: MT-5: a tenant sees mail to its domains, and its own people's sent mail async fn tenant_domains(server: &Server, tenant_id: u32) -> trc::Result> { inbuxa_features::tenancy::queue::tenant_domains(server.registry(), tenant_id).await } fn tenant_sees(domains: &AHashSet, message: &Message) -> bool { inbuxa_features::tenancy::queue::sees( domains, message.recipients.iter().map(|rcpt| rcpt.address()), &message.return_path, message.flags & FROM_AUTHENTICATED != 0, ) } fn tenant_sees_archived(domains: &AHashSet, message: &ArchivedMessage) -> bool { inbuxa_features::tenancy::queue::sees( domains, message.recipients.iter().map(|rcpt| rcpt.address()), &message.return_path, message.flags.to_native() & FROM_AUTHENTICATED != 0, ) } pub(crate) fn map_message(message_in: &ArchivedMessage) -> QueuedMessage { let mut message_out = QueuedMessage { blob_id: BlobId::new(BlobHash::from(&message_in.blob_hash), Default::default()), created_at: UTCDateTime::from_timestamp(message_in.created.to_native() as i64), env_id: message_in.env_id.as_ref().map(|v| v.to_string()), flags: Map::with_capacity(1), priority: message_in.priority.to_native() as i64, received_from_ip: IpAddr(message_in.received_from_ip.as_ipaddr()), received_via_port: message_in.received_via_port.to_native() as u64, recipients: VecMap::with_capacity(message_in.recipients.len()), return_path: if !message_in.return_path.is_empty() { message_in.return_path.to_string() } else { "<>".to_string() }, size: message_in.size.to_native(), next_retry: UTCDateTime::from_timestamp( message_in .next_delivery_event(None) .unwrap_or_else(now) .cast_signed(), ) .into(), next_notify: message_in .next_notify_event(None) .map(|ts| UTCDateTime::from_timestamp(ts.cast_signed())), }; // Parse flags let flags = message_in.flags.to_native(); for (bit, flag) in [ (FROM_AUTHENTICATED, MessageFlag::Authenticated), (FROM_UNAUTHENTICATED, MessageFlag::Unauthenticated), ( FROM_UNAUTHENTICATED_DMARC, MessageFlag::UnauthenticatedDmarc, ), (FROM_DSN, MessageFlag::Dsn), (FROM_REPORT, MessageFlag::Report), (FROM_AUTOGENERATED, MessageFlag::Autogenerated), ] { if flags & bit != 0 { message_out.flags.push(flag); } } // Parse recipients for rcpt_in in message_in.recipients.iter() { let mut rcpt_out = QueuedRecipient { expires: match &rcpt_in.expires { ArchivedQueueExpiry::Ttl(ttl) => QueueExpiry::Ttl(QueueExpiryTtl { expires_at: UTCDateTime::from_timestamp( message_in.created.to_native() as i64 + ttl.to_native() as i64, ), }), ArchivedQueueExpiry::Attempts(attempts) => { QueueExpiry::Attempts(QueueExpiryAttempts { expires_attempts: attempts.to_native() as u64, }) } }, flags: Default::default(), notify_count: rcpt_in.notify.inner.to_native() as u64, notify_due: UTCDateTime::from_timestamp(rcpt_in.notify.due.to_native() as i64), orcpt: rcpt_in.orcpt.as_ref().map(|v| v.to_string()), queue_name: rcpt_in.queue.as_str().to_string(), retry_count: rcpt_in.retry.inner.to_native() as u64, retry_due: UTCDateTime::from_timestamp(rcpt_in.retry.due.to_native() as i64), status: match &rcpt_in.status { ArchivedStatus::Scheduled => RecipientStatus::Scheduled, ArchivedStatus::Completed(status) => RecipientStatus::Completed(ServerResponse { response_code: (status.response.code.to_native() as u64).into(), response_enhanced: build_enhanced_code(&status.response.esc).into(), response_hostname: status.hostname.to_string().into(), response_message: status.response.message.to_string().into(), }), ArchivedStatus::TemporaryFailure(status) => { RecipientStatus::TemporaryFailure(map_error_details(status)) } ArchivedStatus::PermanentFailure(status) => { RecipientStatus::PermanentFailure(map_error_details(status)) } }, }; // Parse recipient flags let rcpt_flags = rcpt_in.flags.to_native(); for (bit, flag) in [(RCPT_DSN_SENT, RecipientFlag::DsnSent)] { if rcpt_flags & bit != 0 { rcpt_out.flags.push(flag); } } if rcpt_spam_percentage(rcpt_flags).is_some_and(|percentage| percentage >= 50) { rcpt_out.flags.push(RecipientFlag::SpamPayload); } message_out .recipients .append(rcpt_in.address.to_string(), rcpt_out); } message_out } fn map_error_details(err_in: &ArchivedErrorDetails) -> DeliveryError { let mut err_out = DeliveryError { response_hostname: err_in.entity.to_string().into(), ..Default::default() }; match &err_in.details { ArchivedError::DnsError(e) => { err_out.error_type = DeliveryErrorType::DnsError; err_out.error_message = e.to_string().into(); } ArchivedError::UnexpectedResponse(e) => { err_out.error_type = DeliveryErrorType::UnexpectedResponse; err_out.error_command = e.command.to_string().into(); err_out.response_code = (e.response.code.to_native() as u64).into(); err_out.response_enhanced = build_enhanced_code(&e.response.esc).into(); err_out.response_message = e.response.message.to_string().into(); } ArchivedError::ConnectionError(e) => { err_out.error_type = DeliveryErrorType::ConnectionError; err_out.error_message = e.to_string().into(); } ArchivedError::TlsError(e) => { err_out.error_type = DeliveryErrorType::TlsError; err_out.error_message = e.to_string().into(); } ArchivedError::DaneError(e) => { err_out.error_type = DeliveryErrorType::DaneError; err_out.error_message = e.to_string().into(); } ArchivedError::MtaStsError(e) => { err_out.error_type = DeliveryErrorType::MtaStsError; err_out.error_message = e.to_string().into(); } ArchivedError::RateLimited => { err_out.error_type = DeliveryErrorType::RateLimited; } ArchivedError::ConcurrencyLimited => { err_out.error_type = DeliveryErrorType::ConcurrencyLimited; } ArchivedError::Io(e) => { err_out.error_type = DeliveryErrorType::Io; err_out.error_message = e.to_string().into(); } } err_out } fn build_enhanced_code(esc: &[u8; 3]) -> String { format!("{}.{}.{}", esc[0], esc[1], esc[2]) } async fn queued_ids(server: &Server, max_results: usize) -> trc::Result> { let mut events = AHashSet::with_capacity(8); let from_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( store::write::QueueEvent { due: 0, queue_id: 0, queue_name: [0; 8], }, ))); let to_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( store::write::QueueEvent { due: u64::MAX, queue_id: u64::MAX, queue_name: [u8::MAX; 8], }, ))); server .store() .iterate( IterateParams::new(from_key, to_key).ascending().no_values(), |key, _| { events.insert(key.deserialize_be_u64(U64_LEN)?); Ok(events.len() < max_results) }, ) .await .caused_by(trc::location!()) .map(|_| events) }