use std::time::Duration;
use tokio::sync::watch;
use tracing::{debug, warn};
use crate::dal::DAL;
use crate::database::universal_types::UniversalTimestamp;
use crate::delivery::WakeHandle;
use crate::error::ValidationError;
use crate::models::delivery_outbox::DeliveryState;
#[derive(Debug, Clone)]
pub struct SweeperConfig {
pub sweep_interval: Duration,
pub stuck_threshold: Duration,
pub batch_limit: i64,
}
impl Default for SweeperConfig {
fn default() -> Self {
Self {
sweep_interval: Duration::from_secs(30),
stuck_threshold: Duration::from_secs(60),
batch_limit: 256,
}
}
}
pub struct DeliverySweeper {
dal: DAL,
wake: WakeHandle,
config: SweeperConfig,
}
impl DeliverySweeper {
pub fn new(dal: DAL, wake: WakeHandle) -> Self {
Self::with_config(dal, wake, SweeperConfig::default())
}
pub fn with_config(dal: DAL, wake: WakeHandle, config: SweeperConfig) -> Self {
Self { dal, wake, config }
}
pub async fn run(self, mut shutdown: watch::Receiver<bool>) {
let mut ticker = tokio::time::interval(self.config.sweep_interval);
ticker.tick().await;
loop {
tokio::select! {
_ = ticker.tick() => {
match self.sweep_once().await {
Ok(reset) => debug!(reset, "delivery sweep complete"),
Err(e) => warn!(error = %e, "delivery sweep failed"),
}
}
res = shutdown.changed() => {
if res.is_err() || *shutdown.borrow() {
break;
}
}
}
}
debug!("delivery sweeper stopped");
}
pub async fn sweep_once(&self) -> Result<usize, ValidationError> {
let cutoff_chrono = chrono::Utc::now()
- chrono::Duration::from_std(self.config.stuck_threshold)
.unwrap_or_else(|_| chrono::Duration::seconds(60));
let cutoff = UniversalTimestamp(cutoff_chrono);
let stuck = self
.dal
.delivery_outbox()
.list_stuck(cutoff, self.config.batch_limit)
.await?;
let mut reset = 0usize;
for row in &stuck {
if row.state() == Some(DeliveryState::Delivered) {
match self.dal.delivery_outbox().reset_to_pending(row.id).await {
Ok(()) => reset += 1,
Err(e) => debug!(
id = row.id,
error = %e,
"sweep reset skipped (state changed concurrently)"
),
}
}
}
if !stuck.is_empty() {
self.wake.wake();
}
metrics::counter!("cloacina_delivery_outbox_sweep_runs_total").increment(1);
if reset > 0 {
metrics::counter!("cloacina_delivery_outbox_sweep_redeliveries_total")
.increment(reset as u64);
}
match self.dal.delivery_outbox().count_open().await {
Ok(n) => metrics::gauge!("cloacina_delivery_outbox_open").set(n as f64),
Err(e) => debug!(error = %e, "count_open failed during sweep"),
}
Ok(reset)
}
}
#[cfg(all(test, feature = "sqlite"))]
mod tests {
use super::*;
use crate::database::Database;
use crate::delivery::{DeliveryOutcome, DeliveryRelay, DeliverySink};
use crate::models::delivery_outbox::{DeliveryOutbox, NewDeliveryOutbox};
use async_trait::async_trait;
use std::sync::{Arc, Mutex};
async fn unique_dal() -> DAL {
let url = format!(
"file:delivery_sweeper_test_{}?mode=memory&cache=shared",
uuid::Uuid::new_v4()
);
let db = Database::new(&url, "", 5);
db.run_migrations()
.await
.expect("migrations should succeed");
DAL::new(db)
}
fn work(recipient: &str) -> NewDeliveryOutbox {
NewDeliveryOutbox {
recipient: recipient.to_string(),
kind: "work".to_string(),
tenant_id: None,
payload: b"x".to_vec(),
}
}
struct NullSink;
#[async_trait]
impl DeliverySink for NullSink {
async fn deliver(
&self,
_row: &DeliveryOutbox,
) -> Result<DeliveryOutcome, crate::delivery::DeliveryError> {
Ok(DeliveryOutcome::NoRoute)
}
}
fn make_sweeper(dal: DAL, threshold: Duration) -> DeliverySweeper {
let relay = DeliveryRelay::new(dal.clone(), Arc::new(NullSink));
let wake = relay.wake_handle();
DeliverySweeper::with_config(
dal,
wake,
SweeperConfig {
sweep_interval: Duration::from_secs(30),
stuck_threshold: threshold,
batch_limit: 256,
},
)
}
#[tokio::test]
async fn sweep_resets_stuck_delivered_rows() {
let dal = unique_dal().await;
let row = dal
.delivery_outbox()
.enqueue(work("agent:1"))
.await
.unwrap();
dal.delivery_outbox().mark_delivered(row.id).await.unwrap();
let sweeper = make_sweeper(dal.clone(), Duration::ZERO);
let reset = sweeper.sweep_once().await.unwrap();
assert_eq!(reset, 1);
let pending = dal.delivery_outbox().list_pending(10).await.unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, row.id);
}
#[tokio::test]
async fn sweep_skips_fresh_rows() {
let dal = unique_dal().await;
let row = dal
.delivery_outbox()
.enqueue(work("agent:1"))
.await
.unwrap();
dal.delivery_outbox().mark_delivered(row.id).await.unwrap();
let sweeper = make_sweeper(dal.clone(), Duration::from_secs(3600));
let reset = sweeper.sweep_once().await.unwrap();
assert_eq!(reset, 0, "fresh row should not be considered stuck");
let open = dal
.delivery_outbox()
.list_open_for_recipient("agent:1", 10)
.await
.unwrap();
assert_eq!(open.len(), 1);
assert_eq!(open[0].delivery_state, "delivered");
}
#[tokio::test]
async fn second_sweep_is_idempotent() {
let dal = unique_dal().await;
let row = dal
.delivery_outbox()
.enqueue(work("agent:1"))
.await
.unwrap();
dal.delivery_outbox().mark_delivered(row.id).await.unwrap();
let sweeper = make_sweeper(dal.clone(), Duration::ZERO);
assert_eq!(sweeper.sweep_once().await.unwrap(), 1);
assert_eq!(sweeper.sweep_once().await.unwrap(), 0);
assert_eq!(
dal.delivery_outbox().list_pending(10).await.unwrap().len(),
1
);
}
#[tokio::test]
async fn concurrent_sweepers_are_race_safe() {
let dal = unique_dal().await;
let _r = dal
.delivery_outbox()
.enqueue(work("agent:1"))
.await
.unwrap();
dal.delivery_outbox().mark_delivered(_r.id).await.unwrap();
let a = make_sweeper(dal.clone(), Duration::ZERO);
let b = make_sweeper(dal.clone(), Duration::ZERO);
let (ra, rb) = tokio::join!(a.sweep_once(), b.sweep_once());
let total = ra.unwrap() + rb.unwrap();
assert_eq!(total, 1, "exactly one sweeper should win the reset CAS");
assert_eq!(
dal.delivery_outbox().list_pending(10).await.unwrap().len(),
1
);
let _ = Mutex::new(()); }
}