A new method, inbuxa:Explanation/set, asks the node's local model for a short plain-words reading of one thing an administrator is looking at: a failed recipient in the queue, a Classify verdict, a log line or trace event, or one setting with its saved value. The server builds the prompt itself from stored data and the registry schema, never from text the console sends, and grounds SMTP replies in RFC 3463 and RFC 5321. What the model is never shown: secrets (including ones nested inside a setting, like an AI model's HTTP auth), raw protocol events, and the contents of any other event. A tag name that doesn't have a tag's shape is refused before a model is asked. Calls share the AI gate with spam classification, but mail always keeps its slot, and Explain has its own hourly count per account and its own on/off switch in inbuxa:AiLimits. The permission is sysAiExplain, superuser only; tenant administrators can't use it. The session carries an aiExplain flag so a console knows when to offer the button. An install whose roles were stored before the permission existed gets it added once, at start-up, to the roles that are administrators' alone, not the User role their defaults share with every account. An operator who removes it later isn't overruled. Tests: unit tests in inbuxa-features and jmap, and ai_explain_tests (run with --ignored) covering the acceptance tests and the upgrade.
778 lines
28 KiB
Rust
778 lines
28 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 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<RegistrySetResponse<'_>> {
|
|
// 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::<Message>()?;
|
|
// 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<RegistryGetResponse<'_>> {
|
|
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::<Message>()?;
|
|
// 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<QueryResponseBuilder> {
|
|
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_ = <Archive<AlignedBytes> as Deserialize>::deserialize(value)
|
|
.add_context(|ctx| ctx.ctx(trc::Key::Key, key))?;
|
|
let message = message_
|
|
.unarchive::<queue::Message>()
|
|
.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::<Message>()
|
|
&& 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<AHashSet<String>> {
|
|
inbuxa_features::tenancy::queue::tenant_domains(server.registry(), tenant_id).await
|
|
}
|
|
|
|
fn tenant_sees(domains: &AHashSet<String>, 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<String>, 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<AHashSet<u64>> {
|
|
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)
|
|
}
|