Report reschedules keep the task queue readable

Setting deliverAt on an internal DMARC or TLS report wrote the new task
queue row with the report's object type (0x21, 0x6e) instead of the task
type (7, 8), and left the task row at its old due. The task manager's scan
failed on that row with store.data-corruption ("Failed to iterate over task
queue"), and because the error ended the whole scan, every task due after
the row stopped running on every node.

- reschedule_ops writes the new queue row through schedule_task_with_id, so
  it carries the task type and the task row gets the new due. It removes
  the row the task is actually queued under (the task's due, which differs
  from deliverAt once the task has been retried) and any row an earlier
  reschedule left at deliverAt.
- x:DmarcInternalReport/set and x:TlsInternalReport/set lock the report's
  task while they move it, as x:Task/set does, refuse while the report is
  being sent, release the locks however the request ends, and wake the task
  manager.
- The task manager logs a queue row it can't read (id, due, key, value) and
  skips it instead of ending the scan. It then repairs the row from its task:
  the row is rewritten with the task's type, and a row with no task behind
  it is removed. A row holding a report's object type for a report task is
  what the old reschedule wrote: the task is moved to that row's time, as
  the reschedule intended, and its old queue row is removed. Stores that
  already hold such a row recover on their own once it comes due.
- x:Task/query with a type filter skips an unreadable row instead of
  failing.

Test: smtp::reporting::reschedule (RocksDB and PostgreSQL). It fails on
main: x:Task/get shows the old due, and with that check removed, neither
report nor a later task ever runs.

(cherry picked from commit 1a7859a8cc)
This commit is contained in:
2026-09-24 18:35:50 -07:00
parent 4e6c8b916e
commit 0b4aa9c084
6 changed files with 602 additions and 33 deletions
+63 -5
View File
@@ -2,6 +2,8 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]> * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* *
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/ */
use crate::{ use crate::{
@@ -15,22 +17,42 @@ use jmap_proto::{error::set::SetError, types::state::State};
use jmap_tools::{Key, Value}; use jmap_tools::{Key, Value};
use registry::{ use registry::{
jmap::IntoValue, jmap::IntoValue,
schema::prelude::{Object, ObjectInner, ObjectType, Property}, schema::{
prelude::{Object, ObjectInner, ObjectType, Property},
structs::Task,
},
types::{EnumImpl, datetime::UTCDateTime}, types::{EnumImpl, datetime::UTCDateTime},
}; };
use services::task_manager::lock::TaskLockManager;
use smtp::reporting::index::{ExternalReportIndex, InternalReportIndex}; use smtp::reporting::index::{ExternalReportIndex, InternalReportIndex};
use std::str::FromStr; use std::str::FromStr;
use store::{ use store::{
U64_LEN, ValueKey, U64_LEN, ValueKey,
registry::{RegistryFilter, RegistryFilterValue, RegistryQuery}, registry::{RegistryFilter, RegistryFilterValue, RegistryQuery},
write::{BatchBuilder, RegistryClass, ValueClass, key::KeySerializer}, write::{BatchBuilder, RegistryClass, TaskQueueClass, ValueClass, key::KeySerializer},
}; };
use trc::AddContext; use trc::AddContext;
use types::id::Id; use types::id::Id;
pub(crate) async fn report_set( pub(crate) async fn report_set(
mut set: RegistrySetResponse<'_>, set: RegistrySetResponse<'_>,
) -> trc::Result<RegistrySetResponse<'_>> { ) -> trc::Result<RegistrySetResponse<'_>> {
// inbuxa: task locks taken to reschedule reports are released however
// the request ends; a held lock is renewed, so a leaked one would keep
// the report's task from ever running
let server = set.server;
let mut locked_tasks = Vec::new();
let result = report_set_locked(set, &mut locked_tasks).await;
for task_id in locked_tasks {
server.remove_index_lock(task_id).await;
}
result
}
async fn report_set_locked<'x>(
mut set: RegistrySetResponse<'x>,
locked_tasks: &mut Vec<u64>,
) -> trc::Result<RegistrySetResponse<'x>> {
let object_id = set.object_type.to_id(); let object_id = set.object_type.to_id();
// Reports cannot be created // Reports cannot be created
@@ -89,12 +111,45 @@ pub(crate) async fn report_set(
.get_value::<Object>(ValueKey::from(key.clone())) .get_value::<Object>(ValueKey::from(key.clone()))
.await? .await?
{ {
// inbuxa: the report's task shares its id. Hold the task
// while its queue rows move, as x:Task/set does, and move the
// row the task is actually queued under
if !set.server.try_lock_task(item_id).await {
set.response.not_updated.append(
id,
SetError::forbidden().with_description(
"The report is being sent and cannot be rescheduled".to_string(),
),
);
continue;
}
locked_tasks.push(item_id);
let queued = set
.server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id: item_id },
)))
.await?;
match &mut report_obj.inner { match &mut report_obj.inner {
ObjectInner::DmarcInternalReport(report) => { ObjectInner::DmarcInternalReport(report) => {
report.reschedule_ops(&mut batch, item_id, report_obj.revision, deliver_at); report.reschedule_ops(
&mut batch,
item_id,
report_obj.revision,
deliver_at,
queued.as_ref(),
);
} }
ObjectInner::TlsInternalReport(report) => { ObjectInner::TlsInternalReport(report) => {
report.reschedule_ops(&mut batch, item_id, report_obj.revision, deliver_at); report.reschedule_ops(
&mut batch,
item_id,
report_obj.revision,
deliver_at,
queued.as_ref(),
);
} }
_ => {} _ => {}
} }
@@ -156,6 +211,9 @@ pub(crate) async fn report_set(
.write(batch.build_all()) .write(batch.build_all())
.await .await
.caused_by(trc::location!())?; .caused_by(trc::location!())?;
// inbuxa: a rescheduled report may now be due sooner than the task
// manager's next scan
set.server.notify_task_queue();
} }
Ok(set) Ok(set)
+4 -9
View File
@@ -463,15 +463,10 @@ pub(crate) async fn task_query(
.set_values(typ.is_some()), .set_values(typ.is_some()),
|key, value| { |key, value| {
if let Some(typ) = typ { if let Some(typ) = typ {
let task_type = // inbuxa: a row whose type can't be read matches no type
TaskType::from_id(value.deserialize_be_u16(0)?).ok_or_else(|| { // filter; the task manager logs and repairs it
trc::StoreEvent::DataCorruption let task_type = value.deserialize_be_u16(0).ok().and_then(TaskType::from_id);
.into_err() if task_type != Some(typ) {
.ctx(trc::Key::Key, key.to_vec())
.ctx(trc::Key::Value, value.to_vec())
.caused_by(trc::location!())
})?;
if task_type != typ {
return Ok(true); return Ok(true);
} }
} }
+133 -6
View File
@@ -27,6 +27,7 @@ use common::network::limiter::ConcurrencyLimiter;
use common::network::{ServerInstance, TcpAcceptor}; use common::network::{ServerInstance, TcpAcceptor};
use common::{Inner, Server}; use common::{Inner, Server};
use registry::schema::enums::TaskType; use registry::schema::enums::TaskType;
use registry::schema::prelude::ObjectType;
use registry::schema::structs::{ use registry::schema::structs::{
Task, TaskManager, TaskRetryStrategy, TaskStatus, TaskStatusFailed, TaskStatusRetry, Task, TaskManager, TaskRetryStrategy, TaskStatus, TaskStatusFailed, TaskStatusRetry,
}; };
@@ -337,6 +338,7 @@ impl TaskQueueManager for Server {
// Retrieve tasks pending to be processed // Retrieve tasks pending to be processed
let mut tasks = Vec::new(); let mut tasks = Vec::new();
let mut unreadable = Vec::new();
let now = Instant::now(); let now = Instant::now();
let mut next_event = None; let mut next_event = None;
let roles = &self.core.network.roles; let roles = &self.core.network.roles;
@@ -351,12 +353,21 @@ impl TaskQueueManager for Server {
let task_id = key.deserialize_be_u64(U64_LEN)?; let task_id = key.deserialize_be_u64(U64_LEN)?;
if task_due <= now_timestamp { if task_due <= now_timestamp {
let task_type_idx = value.deserialize_be_u16(0)?; // inbuxa: a row whose task type can't be read is
let task_type = TaskType::from_id(task_type_idx).ok_or_else(|| { // set aside, not allowed to end the scan: every
trc::StoreEvent::DataCorruption // task due after it would wait behind it
.caused_by(trc::location!()) let Some((task_type_idx, task_type)) = value
.ctx(trc::Key::Value, value) .deserialize_be_u16(0)
})?; .ok()
.and_then(|idx| TaskType::from_id(idx).map(|typ| (idx, typ)))
else {
unreadable.push(UnreadableDueRow {
due: task_due,
id: task_id,
value: value.to_vec(),
});
return Ok(true);
};
let enabled = match task_type { let enabled = match task_type {
TaskType::IndexDocument TaskType::IndexDocument
| TaskType::UnindexDocument | TaskType::UnindexDocument
@@ -446,6 +457,11 @@ impl TaskQueueManager for Server {
); );
}); });
if !unreadable.is_empty() && repair_due_rows(self, unreadable).await {
// Look again at once for the rows that were rewritten
self.notify_task_queue();
}
if !tasks.is_empty() { if !tasks.is_empty() {
trc::event!( trc::event!(
TaskManager(TaskManagerEvent::TaskAcquired), TaskManager(TaskManagerEvent::TaskAcquired),
@@ -686,3 +702,114 @@ impl TaskResult {
) )
} }
} }
/// inbuxa: a task queue row whose task type could not be read.
struct UnreadableDueRow {
due: u64,
id: u64,
value: Vec<u8>,
}
/// inbuxa: logs each unreadable queue row and repairs it from the task it
/// schedules. The task row says what the task is, so the queue row is
/// rewritten with that task's type; a row with no task behind it is removed.
///
/// Rescheduling an internal DMARC or TLS report wrote the report's object
/// type into the queue row instead of the task type. Such a row is the time
/// an administrator chose, so the task is moved to it as the reschedule
/// meant to do: the task row takes that due, and a queue row left at the
/// task's previous due is removed. Returns whether any row was repaired.
async fn repair_due_rows(server: &Server, rows: Vec<UnreadableDueRow>) -> bool {
let mut repaired = false;
for row in rows {
let UnreadableDueRow { due, id, value } = row;
trc::error!(
trc::StoreEvent::DataCorruption
.into_err()
.id(id)
.ctx(trc::Key::Due, trc::Value::Timestamp(due))
.ctx(
trc::Key::Key,
[due.to_be_bytes(), id.to_be_bytes()].concat()
)
.ctx(trc::Key::Value, value.clone())
.details("Unreadable task queue row skipped")
.caused_by(trc::location!())
);
let task_key = ValueClass::TaskQueue(TaskQueueClass::Task { id });
let due_key = ValueClass::TaskQueue(TaskQueueClass::Due { id, due });
let task = match server
.store()
.get_value::<Task>(ValueKey::from(task_key.clone()))
.await
{
Ok(task) => task,
Err(err) => {
trc::error!(
err.id(id)
.details("Failed to read the task of an unreadable queue row.")
.caused_by(trc::location!())
);
continue;
}
};
let mut batch = BatchBuilder::new();
let action = if let Some(mut task) = task {
let task_type = task.object_type();
batch.assert_value(task_key.clone(), AssertValue::Some);
if rescheduled_report_type(&value) == Some(task_type) {
let old_due = task.due_timestamp();
if old_due != due {
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id,
due: old_due,
}));
}
task.set_status(TaskStatus::at(due as i64));
}
batch
.set(due_key, task_type.to_id().serialize())
.set(task_key, task.to_pickled_vec());
"Rewrote the queue row from its task."
} else {
batch.clear(due_key);
"Removed a queue row with no task."
};
match server.store().write(batch.build_all()).await {
Ok(_) => {
repaired = true;
trc::event!(
TaskManager(TaskManagerEvent::TaskIgnored),
Id = id,
Due = trc::Value::Timestamp(due),
Reason = action,
);
}
Err(err) if err.matches(trc::EventType::Store(trc::StoreEvent::AssertValueFailed)) => {
// The task went away meanwhile; the next scan looks again
}
Err(err) => {
trc::error!(
err.id(id)
.details("Failed to repair an unreadable queue row.")
.caused_by(trc::location!())
);
}
}
}
repaired
}
/// inbuxa: the task type a report reschedule meant, when a queue row holds
/// an internal report's object type (the value that reschedule wrote).
fn rescheduled_report_type(value: &[u8]) -> Option<TaskType> {
let id = u16::from_be_bytes(value.get(..2)?.try_into().ok()?);
match ObjectType::from_id(id)? {
ObjectType::DmarcInternalReport => Some(TaskType::DmarcReport),
ObjectType::TlsInternalReport => Some(TaskType::TlsReport),
_ => None,
}
}
+29 -13
View File
@@ -2,6 +2,8 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]> * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* *
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/ */
use registry::{ use registry::{
@@ -40,35 +42,49 @@ pub trait InternalReportIndex: ObjectImpl {
fn primary_key(&self) -> ValueClass; fn primary_key(&self) -> ValueClass;
/// Moves the report's delivery, and its queued task, to `at`.
///
/// inbuxa: the new queue row carries the task's type, as
/// `schedule_task_with_id` writes it, and the task row gets the new due
/// too. `queued` is the task as stored: its due, not the report's
/// `deliverAt`, is the queue row that exists (they differ once the task
/// has been retried).
fn reschedule_ops( fn reschedule_ops(
&mut self, &mut self,
batch: &mut BatchBuilder, batch: &mut BatchBuilder,
item_id: u64, item_id: u64,
revision: u64, revision: u64,
at: UTCDateTime, at: UTCDateTime,
queued: Option<&Task>,
) { ) {
let current_deliver_at = self.deliver_at(); let current_deliver_at = self.deliver_at();
let current_due = current_deliver_at.timestamp() as u64;
let queued_due = queued.map_or(current_due, |task| task.due_timestamp());
let new_due = at.timestamp() as u64;
if current_deliver_at != at { if current_deliver_at != at || queued_due != new_due {
let object = Self::OBJECT; let object = Self::OBJECT;
let object_id = object.to_id(); let object_id = object.to_id();
let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id }); let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id });
self.set_deliver_at(at); self.set_deliver_at(at);
batch.assert_value(key.clone(), AssertValue::Hash(revision));
if queued_due != new_due {
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: queued_due,
}));
}
// A row an earlier reschedule left at the report's deliverAt
if current_due != new_due && current_due != queued_due {
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: current_due,
}));
}
batch batch
.assert_value(key.clone(), AssertValue::Hash(revision)) .schedule_task_with_id(item_id, self.task(item_id))
.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: current_deliver_at.timestamp() as u64,
}))
.set(
ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: at.timestamp() as u64,
}),
object_id.serialize(),
)
.set(key, self.to_pickled_vec()); .set(key, self.to_pickled_vec());
} }
} }
+3
View File
@@ -2,9 +2,12 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]> * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* *
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/ */
pub mod analyze; pub mod analyze;
pub mod dmarc; pub mod dmarc;
pub mod reschedule; // inbuxa: report reschedules and unreadable queue rows
pub mod scheduler; pub mod scheduler;
pub mod tls; pub mod tls;
+370
View File
@@ -0,0 +1,370 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! Rescheduling an internal DMARC or TLS report over JMAP moves its task: the
//! task runs at the new time, x:Task/get shows the new due, and tasks due
//! after it still run. A task queue row whose type can't be read is logged
//! and repaired rather than stopping every task due after it, including the
//! rows an earlier reschedule wrote with the report's object type.
use crate::utils::server::{TestServer, TestServerBuilder};
use common::{
Server,
config::smtp::report::AggregateFrequency,
ipc::{DmarcEvent, PolicyType, TlsEvent},
};
use mail_auth::{
common::parse::TxtRecordParser,
dmarc::Dmarc,
mta_sts::TlsRpt,
report::{ActionDisposition, DmarcResult, Record},
};
use registry::{
schema::{
enums::{TaskStoreMaintenanceType, TaskType},
prelude::{ObjectType, Property},
structs::{
DmarcInternalReport, DmarcReportSettings, Expression, Task, TaskStatus,
TaskStoreMaintenance, TlsInternalReport, TlsReportSettings,
},
},
types::{EnumImpl, ObjectImpl, datetime::UTCDateTime},
};
use serde_json::json;
use smtp::reporting::{index::InternalReportIndex, send::MtaReportSend};
use std::{
sync::Arc,
time::{Duration, Instant},
};
use store::{
SerializeInfallible, ValueKey,
write::{BatchBuilder, RegistryClass, TaskQueueClass, ValueClass, now},
};
use types::id::Id;
use utils::snowflake::SnowflakeIdGenerator;
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn report_reschedule() {
let mut test = TestServerBuilder::new("smtp_report_reschedule")
.await
.with_http_listener(19057)
.await
.capture_queue()
.build()
.await;
let admin = test.account("admin");
admin
.registry_create_object(TlsReportSettings {
max_report_size: Expression {
else_: "1024".into(),
..Default::default()
},
..Default::default()
})
.await;
admin
.registry_create_object(DmarcReportSettings {
aggregate_max_report_size: Expression {
else_: "1024".into(),
..Default::default()
},
..Default::default()
})
.await;
admin.reload_settings().await;
test.reload_core();
test.expect_reload_settings().await;
let admin = test.account("admin");
// A daily DMARC and TLS report, due a day from now
schedule_dmarc(&test, "foobar.org").await;
schedule_tls(&test, "foobar.org").await;
let dmarc_id = wait_for_report::<DmarcInternalReport>(&test, "foobar.org").await;
let tls_id = wait_for_report::<TlsInternalReport>(&test, "foobar.org").await;
// Reschedule both to a few seconds from now, with a task due after them
let at = now() + 3;
let later = marker_task(&test.server, at + 3).await;
for (object, id, task_type) in [
(
ObjectType::DmarcInternalReport,
dmarc_id,
TaskType::DmarcReport,
),
(ObjectType::TlsInternalReport, tls_id, TaskType::TlsReport),
] {
admin
.registry_update_object(
object,
id,
json!({
Property::DeliverAt: UTCDateTime::from_timestamp(at as i64),
}),
)
.await;
// x:Task/get shows the new due, and the queue row carries the task's
// type. Upstream wrote the report's object type there and left the
// task at its old due
let task = admin.registry_get::<Task>(id).await;
assert_eq!(task.object_type(), task_type);
assert_eq!(
task.due_timestamp(),
at,
"{object:?} task due not moved: {task:?}"
);
assert_eq!(
queue_row(&test.server, id.id(), at).await,
Some(task_type.to_id().serialize()),
"{object:?} queue row"
);
}
// Both reports go out at the new time, and the later task still runs
wait_until_run(&test.server, &[dmarc_id.id(), tls_id.id(), later]).await;
assert!(now() >= at, "the reports went out before their new time");
assert!(
admin
.registry_get_all::<DmarcInternalReport>()
.await
.is_empty()
);
assert!(
admin
.registry_get_all::<TlsInternalReport>()
.await
.is_empty()
);
// Rows an earlier reschedule may have left in a store: one with the
// report's object type and the task left at its old due, and one that
// is unreadable and has no task behind it. Neither may hold back a task
// due after them.
schedule_dmarc(&test, "foobar.net").await;
let dmarc_id = wait_for_report::<DmarcInternalReport>(&test, "foobar.net").await;
let at = now() + 2;
let old_due = old_style_reschedule(&test.server, dmarc_id.id(), at).await;
let orphan = SnowflakeIdGenerator::global_id().unwrap();
let mut batch = BatchBuilder::new();
batch.set(
ValueClass::TaskQueue(TaskQueueClass::Due {
id: orphan,
due: at,
}),
vec![0xff, 0xff],
);
test.server.store().write(batch.build_all()).await.unwrap();
let later = marker_task(&test.server, at + 2).await;
wait_until_run(&test.server, &[dmarc_id.id(), later]).await;
assert!(
admin
.registry_get_all::<DmarcInternalReport>()
.await
.is_empty()
);
assert_eq!(queue_row(&test.server, orphan, at).await, None);
assert_eq!(queue_row(&test.server, dmarc_id.id(), at).await, None);
assert_eq!(queue_row(&test.server, dmarc_id.id(), old_due).await, None);
// x:Task/query by type skips an unreadable row rather than failing
let mut batch = BatchBuilder::new();
let due = now() + 3600;
batch.set(
ValueClass::TaskQueue(TaskQueueClass::Due { id: orphan, due }),
vec![0xff, 0xff],
);
test.server.store().write(batch.build_all()).await.unwrap();
admin
.registry_query_ids(
ObjectType::Task,
vec![(Property::Type, TaskType::DmarcReport.as_str())],
Vec::<&str>::new(),
)
.await;
let mut batch = BatchBuilder::new();
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: orphan,
due,
}));
test.server.store().write(batch.build_all()).await.unwrap();
if test.is_reset() {
test.temp_dir.delete();
}
}
async fn schedule_dmarc(test: &TestServer, domain: &str) {
test.server
.schedule_report(DmarcEvent {
domain: domain.to_string(),
report_record: Record::new()
.with_source_ip("192.168.1.2".parse().unwrap())
.with_action_disposition(ActionDisposition::Pass)
.with_dmarc_dkim_result(DmarcResult::Pass)
.with_dmarc_spf_result(DmarcResult::Fail)
.with_envelope_from("[email protected]")
.with_envelope_to("[email protected]")
.with_header_from("[email protected]"),
dmarc_record: Arc::new(
Dmarc::parse(format!("v=DMARC1; p=reject; rua=mailto:reports@{domain}").as_bytes())
.unwrap(),
),
interval: AggregateFrequency::Daily,
span_id: 0,
})
.await;
}
async fn schedule_tls(test: &TestServer, domain: &str) {
test.server
.schedule_report(TlsEvent {
domain: domain.to_string(),
policy: PolicyType::None,
failure: None,
tls_record: Arc::new(
TlsRpt::parse(format!("v=TLSRPTv1;rua=mailto:reports@{domain}").as_bytes())
.unwrap(),
),
interval: AggregateFrequency::Daily,
span_id: 0,
})
.await;
}
trait ReportDomain: ObjectImpl {
fn report_domain(&self) -> &str;
}
impl ReportDomain for DmarcInternalReport {
fn report_domain(&self) -> &str {
&self.domain
}
}
impl ReportDomain for TlsInternalReport {
fn report_domain(&self) -> &str {
&self.domain
}
}
async fn wait_for_report<T: ReportDomain>(test: &TestServer, domain: &str) -> Id {
let admin = test.account("admin");
for _ in 0..100 {
if let Some((id, _)) = admin
.registry_get_all::<T>()
.await
.into_iter()
.find(|(_, report)| report.report_domain() == domain)
{
return id;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
panic!("No {} for {domain}", T::OBJECT.as_str());
}
/// A task that succeeds when it runs, due at `due`.
async fn marker_task(server: &Server, due: u64) -> u64 {
let id = SnowflakeIdGenerator::global_id().unwrap();
let mut batch = BatchBuilder::new();
batch.schedule_task_with_id(
id,
Task::StoreMaintenance(TaskStoreMaintenance {
maintenance_type: TaskStoreMaintenanceType::RemoveLockDav,
shard_index: Some(0),
status: TaskStatus::at(due as i64),
}),
);
server.store().write(batch.build_all()).await.unwrap();
server.notify_task_queue();
id
}
/// What the reschedule before this fix wrote: the report's object type in
/// the new queue row, and the task row left at its old due. Returns that
/// old due.
async fn old_style_reschedule(server: &Server, item_id: u64, at: u64) -> u64 {
let object_id = ObjectType::DmarcInternalReport.to_id();
let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id });
let mut report = server
.store()
.get_value::<DmarcInternalReport>(ValueKey::from(key.clone()))
.await
.unwrap()
.unwrap();
let old_due = report.deliver_at().timestamp() as u64;
report.set_deliver_at(UTCDateTime::from_timestamp(at as i64));
let mut batch = BatchBuilder::new();
batch
.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: old_due,
}))
.set(
ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: at,
}),
object_id.serialize(),
)
.set(key, report.to_pickled_vec());
server.store().write(batch.build_all()).await.unwrap();
server.notify_task_queue();
old_due
}
struct RawValue(Vec<u8>);
impl store::Deserialize for RawValue {
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
Ok(RawValue(bytes.to_vec()))
}
}
async fn queue_row(server: &Server, id: u64, due: u64) -> Option<Vec<u8>> {
server
.store()
.get_value::<RawValue>(ValueKey::from(ValueClass::TaskQueue(TaskQueueClass::Due {
id,
due,
})))
.await
.unwrap()
.map(|raw| raw.0)
}
async fn task_exists(server: &Server, id: u64) -> bool {
server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id },
)))
.await
.unwrap()
.is_some()
}
async fn wait_until_run(server: &Server, ids: &[u64]) {
let started = Instant::now();
loop {
let mut pending = Vec::new();
for id in ids {
if task_exists(server, *id).await {
pending.push(*id);
}
}
if pending.is_empty() {
return;
}
if started.elapsed() > Duration::from_secs(30) {
panic!("tasks {pending:?} never ran");
}
tokio::time::sleep(Duration::from_millis(200)).await;
}
}