diff --git a/.gitea/workflows/publish.yml b/.gitea/workflows/publish.yml index f087bd5..5df8087 100644 --- a/.gitea/workflows/publish.yml +++ b/.gitea/workflows/publish.yml @@ -14,8 +14,11 @@ # 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 # reports a different version from its tag. -# * the tag must be on main, so an image never describes code that was never -# reviewed onto the default branch. +# * the tag must be on main or on a release/* branch, so an image never +# 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 # (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 exit 1 fi - git merge-base --is-ancestor "$(git rev-parse "${TAG}^{commit}")" origin/main \ - || { echo "$TAG is not on main" >&2; exit 1; } + commit="$(git rev-parse "${TAG}^{commit}")" + 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" diff --git a/crates/jmap/src/registry/mapping/report.rs b/crates/jmap/src/registry/mapping/report.rs index da01a22..d8a6fcb 100644 --- a/crates/jmap/src/registry/mapping/report.rs +++ b/crates/jmap/src/registry/mapping/report.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use crate::{ @@ -15,22 +17,42 @@ use jmap_proto::{error::set::SetError, types::state::State}; use jmap_tools::{Key, Value}; use registry::{ jmap::IntoValue, - schema::prelude::{Object, ObjectInner, ObjectType, Property}, + schema::{ + prelude::{Object, ObjectInner, ObjectType, Property}, + structs::Task, + }, types::{EnumImpl, datetime::UTCDateTime}, }; +use services::task_manager::lock::TaskLockManager; use smtp::reporting::index::{ExternalReportIndex, InternalReportIndex}; use std::str::FromStr; use store::{ U64_LEN, ValueKey, registry::{RegistryFilter, RegistryFilterValue, RegistryQuery}, - write::{BatchBuilder, RegistryClass, ValueClass, key::KeySerializer}, + write::{BatchBuilder, RegistryClass, TaskQueueClass, ValueClass, key::KeySerializer}, }; use trc::AddContext; use types::id::Id; pub(crate) async fn report_set( - mut set: RegistrySetResponse<'_>, + set: RegistrySetResponse<'_>, ) -> trc::Result> { + // 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, +) -> trc::Result> { let object_id = set.object_type.to_id(); // Reports cannot be created @@ -89,12 +111,45 @@ pub(crate) async fn report_set( .get_value::(ValueKey::from(key.clone())) .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::(ValueKey::from(ValueClass::TaskQueue( + TaskQueueClass::Task { id: item_id }, + ))) + .await?; + match &mut report_obj.inner { 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) => { - 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()) .await .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) diff --git a/crates/jmap/src/registry/mapping/task.rs b/crates/jmap/src/registry/mapping/task.rs index 229dd16..3a6cf10 100644 --- a/crates/jmap/src/registry/mapping/task.rs +++ b/crates/jmap/src/registry/mapping/task.rs @@ -463,15 +463,10 @@ pub(crate) async fn task_query( .set_values(typ.is_some()), |key, value| { if let Some(typ) = typ { - let task_type = - TaskType::from_id(value.deserialize_be_u16(0)?).ok_or_else(|| { - trc::StoreEvent::DataCorruption - .into_err() - .ctx(trc::Key::Key, key.to_vec()) - .ctx(trc::Key::Value, value.to_vec()) - .caused_by(trc::location!()) - })?; - if task_type != typ { + // inbuxa: a row whose type can't be read matches no type + // filter; the task manager logs and repairs it + let task_type = value.deserialize_be_u16(0).ok().and_then(TaskType::from_id); + if task_type != Some(typ) { return Ok(true); } } diff --git a/crates/services/src/task_manager/manager.rs b/crates/services/src/task_manager/manager.rs index 180aeac..e3c9331 100644 --- a/crates/services/src/task_manager/manager.rs +++ b/crates/services/src/task_manager/manager.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use crate::task_manager::acme::AcmeTask; @@ -27,6 +29,7 @@ use common::network::limiter::ConcurrencyLimiter; use common::network::{ServerInstance, TcpAcceptor}; use common::{Inner, Server}; use registry::schema::enums::TaskType; +use registry::schema::prelude::ObjectType; use registry::schema::structs::{ Task, TaskManager, TaskRetryStrategy, TaskStatus, TaskStatusFailed, TaskStatusRetry, }; @@ -337,6 +340,7 @@ impl TaskQueueManager for Server { // Retrieve tasks pending to be processed let mut tasks = Vec::new(); + let mut unreadable = Vec::new(); let now = Instant::now(); let mut next_event = None; let roles = &self.core.network.roles; @@ -351,12 +355,21 @@ impl TaskQueueManager for Server { let task_id = key.deserialize_be_u64(U64_LEN)?; if task_due <= now_timestamp { - let task_type_idx = value.deserialize_be_u16(0)?; - let task_type = TaskType::from_id(task_type_idx).ok_or_else(|| { - trc::StoreEvent::DataCorruption - .caused_by(trc::location!()) - .ctx(trc::Key::Value, value) - })?; + // inbuxa: a row whose task type can't be read is + // set aside, not allowed to end the scan: every + // task due after it would wait behind it + let Some((task_type_idx, task_type)) = 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 { TaskType::IndexDocument | 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() { trc::event!( 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, +} + +/// 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) -> 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::(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 { + 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, + } +} diff --git a/crates/smtp/src/reporting/index.rs b/crates/smtp/src/reporting/index.rs index 83cc993..0d92347 100644 --- a/crates/smtp/src/reporting/index.rs +++ b/crates/smtp/src/reporting/index.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use registry::{ @@ -40,35 +42,49 @@ pub trait InternalReportIndex: ObjectImpl { 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( &mut self, batch: &mut BatchBuilder, item_id: u64, revision: u64, at: UTCDateTime, + queued: Option<&Task>, ) { 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_id = object.to_id(); let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id }); self.set_deliver_at(at); - batch - .assert_value(key.clone(), AssertValue::Hash(revision)) - .clear(ValueClass::TaskQueue(TaskQueueClass::Due { + batch.assert_value(key.clone(), AssertValue::Hash(revision)); + if queued_due != new_due { + batch.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(), - ) + 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 + .schedule_task_with_id(item_id, self.task(item_id)) .set(key, self.to_pickled_vec()); } } diff --git a/crates/types/src/branding.rs b/crates/types/src/branding.rs index 624ef22..e564d7c 100644 --- a/crates/types/src/branding.rs +++ b/crates/types/src/branding.rs @@ -81,7 +81,7 @@ fn legacy_setting(name: &str, is_set: impl Fn(&str) -> bool) -> Option { #[macro_export] macro_rules! brand_version { () => { - "2026.9.24.3" + "2026.9.24.4" }; } diff --git a/tests/src/smtp/reporting/mod.rs b/tests/src/smtp/reporting/mod.rs index 0972ccc..1aa3f3f 100644 --- a/tests/src/smtp/reporting/mod.rs +++ b/tests/src/smtp/reporting/mod.rs @@ -2,9 +2,12 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ pub mod analyze; pub mod dmarc; +pub mod reschedule; // inbuxa: report reschedules and unreadable queue rows pub mod scheduler; pub mod tls; diff --git a/tests/src/smtp/reporting/reschedule.rs b/tests/src/smtp/reporting/reschedule.rs new file mode 100644 index 0000000..675b5c1 --- /dev/null +++ b/tests/src/smtp/reporting/reschedule.rs @@ -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::(&test, "foobar.org").await; + let tls_id = wait_for_report::(&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::(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::() + .await + .is_empty() + ); + assert!( + admin + .registry_get_all::() + .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::(&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::() + .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("hello@example.org") + .with_envelope_to("other@example.org") + .with_header_from("bye@example.org"), + 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(test: &TestServer, domain: &str) -> Id { + let admin = test.account("admin"); + for _ in 0..100 { + if let Some((id, _)) = admin + .registry_get_all::() + .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::(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); + +impl store::Deserialize for RawValue { + fn deserialize(bytes: &[u8]) -> trc::Result { + Ok(RawValue(bytes.to_vec())) + } +} + +async fn queue_row(server: &Server, id: u64, due: u64) -> Option> { + server + .store() + .get_value::(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::(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; + } +}