Files
inbuxa-server/crates/jmap/src/registry/mapping/queued_message.rs
T
jcoffey-dev f44382fb09
ci / fork-checks (pull_request) Successful in 33s
ci / build (pull_request) Successful in 9m58s
DLP phase 3: hold for review
The hold action now holds (dlp-and-mail-flow-rules spec, §2.6), where
until now it blocked.

- At DATA a hold decision queues the message with its release a century
  off (the queue's future-release mechanism, so the stored format is
  unchanged and an older node just never sends it), transport rules
  still applied, and replies 250 Held for review. A review record under
  R/h + queue id keeps the sender, recipients, subject, size, rules and
  detector counts. The sender is told when the rule asks.
- smtp/queue/held.rs: release (each recipient due now, its next notice
  as far off as it was, its lifetime counted from the release), reject
  (removed from the queue, the sender told, with the reviewer's note),
  and expiry: the daily clean-up rejects what nobody reviewed in 7 days,
  recorded as the server's doing.
- inbuxa:HeldMessage get/set: the review queue, sysDlpReviewGet to list
  and read (preview, 64 KB of text, only when asked for and recorded as
  blobAccess), sysDlpReviewUpdate to release or reject, a reason
  required and audited by the request layer; no create or destroy;
  server-level only.
- Guards: Emails > Queue refuses to change or delete held mail; the
  sender can't unsend it.
- Privacy catalog entry for inbuxa:HeldMessage; spec §2.6 as built.

Tests: mail_rules_tests gains the whole flow (held and listed with
counts, sender notified and nothing delivered, queue and unsend
refused, preview recorded, reject needs a reason and tells the sender
the note, release delivers, expiry returns it, decisions audited with
reasons). smtp inbound, system_tests (after one BlobNotFound in
antispam, the known flake, then clean), features and common unit tests.
2026-09-28 18:33:12 -07:00

795 lines
29 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;
/// 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<'_>> {
// 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();
// 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;
};
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(..) {
// 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;
};
// 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)
}