5 Commits
Author SHA1 Message Date
jcoffey-dev f398d95062 Merge pull request 'Hotfix 2026.9.24.4: report reschedules keep the task queue readable' (#49) from hotfix/2026.9.24.4 into release/2026.9.24.4
publish / version (push) Successful in 41s
publish / publish (push) Successful in 56m8s
publish / release (push) Successful in 6s
publish / binaries (push) Successful in 49s
2026-09-25 01:43:57 +00:00
jcoffey-dev 81ec917d74 Mark task_manager/manager.rs as modified for the AGPL notice
ci / fork-checks (pull_request) Successful in 22s
ci / build (pull_request) Successful in 3m37s
2026-09-24 18:39:42 -07:00
jcoffey-dev e2ab26ad19 Release 2026.9.24.4
ci / fork-checks (pull_request) Failing after 29s
ci / build (pull_request) Canceled after 3m15s
2026-09-24 18:35:50 -07:00
jcoffey-dev 0b4aa9c084 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)
2026-09-24 18:35:50 -07:00
jcoffey-dev 4e6c8b916e Publish: accept tags on release/* branches for hotfix releases
The publish workflow only built a tag whose commit is on main. That keeps
every image tied to reviewed code, but it means production can only get a
fix together with everything that has landed on main since its release.

A tag on a release/* branch is now accepted too. A hotfix branch starts at
an earlier release tag, takes fixes through pull requests into it (so the
code is still reviewed and CI-tested before it is tagged), bumps
brand_version! and is tagged there. The tag must still equal
v<brand_version!>, and the step prints which branch it was found on.

A tag runs the workflow file from its own commit, so a hotfix branch that
starts before this change needs this commit cherry-picked onto it before
its tag is pushed.

(cherry picked from commit 4b85113262)
2026-09-24 18:31:52 -07:00
8 changed files with 617 additions and 38 deletions
+12 -4
View File
@@ -14,8 +14,11 @@
# crates/types/src/branding.rs, not Cargo.toml, and the image is tagged # crates/types/src/branding.rs, not Cargo.toml, and the image is tagged
# with it, so a tag beside an unbumped macro would publish an image that # with it, so a tag beside an unbumped macro would publish an image that
# reports a different version from its tag. # reports a different version from its tag.
# * the tag must be on main, so an image never describes code that was never # * the tag must be on main or on a release/* branch, so an image never
# reviewed onto the default branch. # describes code that was never reviewed onto one of them. A release/*
# branch carries a hotfix: it starts at an earlier release tag, takes
# fixes through pull requests into it, and is tagged there, so production
# can get a fix without everything that has landed on main since.
# #
# :latest moves with every published tag: tags are cut by the weekly release # :latest moves with every published tag: tags are cut by the weekly release
# (or by hand for a real release); there are no prerelease tags here. # (or by hand for a real release); there are no prerelease tags here.
@@ -57,8 +60,13 @@ jobs:
echo "Refusing to publish an image that would report the wrong version." >&2 echo "Refusing to publish an image that would report the wrong version." >&2
exit 1 exit 1
fi fi
git merge-base --is-ancestor "$(git rev-parse "${TAG}^{commit}")" origin/main \ commit="$(git rev-parse "${TAG}^{commit}")"
|| { echo "$TAG is not on main" >&2; exit 1; } on=""
for ref in origin/main $(git for-each-ref --format='%(refname:short)' 'refs/remotes/origin/release/*'); do
if git merge-base --is-ancestor "$commit" "$ref"; then on="$ref"; break; fi
done
[ -n "$on" ] || { echo "$TAG is not on main or a release/* branch" >&2; exit 1; }
echo "$TAG is on $on"
echo "version=$V" >> "$GITHUB_OUTPUT" echo "version=$V" >> "$GITHUB_OUTPUT"
echo "version $V" echo "version $V"
+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);
} }
} }
+135 -6
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::task_manager::acme::AcmeTask; use crate::task_manager::acme::AcmeTask;
@@ -27,6 +29,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 +340,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 +355,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 +459,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 +704,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 batch.assert_value(key.clone(), AssertValue::Hash(revision));
.assert_value(key.clone(), AssertValue::Hash(revision)) if queued_due != new_due {
.clear(ValueClass::TaskQueue(TaskQueueClass::Due { batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id, id: item_id,
due: current_deliver_at.timestamp() as u64, due: queued_due,
})) }));
.set( }
ValueClass::TaskQueue(TaskQueueClass::Due { // A row an earlier reschedule left at the report's deliverAt
id: item_id, if current_due != new_due && current_due != queued_due {
due: at.timestamp() as u64, batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
}), id: item_id,
object_id.serialize(), due: current_due,
) }));
}
batch
.schedule_task_with_id(item_id, self.task(item_id))
.set(key, self.to_pickled_vec()); .set(key, self.to_pickled_vec());
} }
} }
+1 -1
View File
@@ -81,7 +81,7 @@ fn legacy_setting(name: &str, is_set: impl Fn(&str) -> bool) -> Option<String> {
#[macro_export] #[macro_export]
macro_rules! brand_version { macro_rules! brand_version {
() => { () => {
"2026.9.24.3" "2026.9.24.4"
}; };
} }
+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;
}
}