use std::fmt::Debug;
use std::sync::Arc;
use async_trait::async_trait;
use serde_json::Value;
use crate::core::StoreError;
use crate::journal::{JournalStore, Record, RecordKind};
use super::{Delivered, PushNamespace, PushRegistration, PushStore, PushTransport};
#[async_trait]
pub trait Projection: Send + Sync + Debug {
async fn payloads(&self, record: &Record) -> Result<Vec<Value>, StoreError>;
fn terminal(&self, record: &Record) -> bool {
matches!(record.kind(), RecordKind::RunSealed { .. })
}
fn namespace(&self) -> PushNamespace;
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct PushSweepReport {
pub registrations: usize,
pub records: usize,
pub deliveries: usize,
pub retries: usize,
pub completed: usize,
pub abandoned: usize,
pub unserved: usize,
pub saturated: bool,
}
impl PushSweepReport {
#[must_use]
pub const fn needs_attention(&self) -> bool {
self.saturated || self.abandoned > 0
}
}
#[derive(Debug, Clone)]
pub struct DeliveryWorker {
journal: Arc<dyn JournalStore>,
store: Arc<dyn PushStore>,
transport: Arc<dyn PushTransport>,
projection: Arc<dyn Projection>,
max_attempts: u32,
}
impl DeliveryWorker {
pub const DEFAULT_MAX_ATTEMPTS: u32 = 32;
#[must_use]
pub fn new(
journal: Arc<dyn JournalStore>,
store: Arc<dyn PushStore>,
transport: Arc<dyn PushTransport>,
projection: Arc<dyn Projection>,
) -> Self {
Self {
journal,
store,
transport,
projection,
max_attempts: Self::DEFAULT_MAX_ATTEMPTS,
}
}
#[must_use]
pub const fn max_attempts(mut self, attempts: u32) -> Self {
assert!(
attempts > 0,
"a push retry ceiling of zero abandons every receiver on its first \
hiccup; configure no push instead"
);
self.max_attempts = attempts;
self
}
pub async fn run_once(
&self,
at: u64,
limit: usize,
) -> Result<PushSweepReport, crate::core::StoreError> {
let batch = self
.store
.due_in(at, limit.saturating_add(1), self.projection.namespace())
.await?;
let saturated = batch.rows.len() > limit;
let mut report = PushSweepReport {
registrations: batch.rows.len().min(limit),
unserved: batch.unserved,
saturated,
..PushSweepReport::default()
};
for registration in batch.rows.into_iter().take(limit) {
let mut attempts = registration.attempts;
let records = self
.journal
.read(registration.config.task, registration.next_seq)
.await?;
if records.is_empty() && self.cleanup_acknowledged_terminal(®istration).await? {
report.completed += 1;
continue;
}
for record in records {
let payloads = match self.projection.payloads(&record).await {
Ok(payloads) => payloads,
Err(error) => {
self.give_up_or_retry(
®istration,
at,
attempts,
&error.to_string(),
false,
&mut report,
)
.await?;
break;
}
};
let mut failed = None;
for payload in payloads {
match self.transport.deliver(®istration.config, &payload).await {
Ok(Delivered::Accepted) => report.deliveries += 1,
Ok(other) => failed = Some((format!("receiver outcome: {other:?}"), false)),
Err(error) => failed = Some((error.to_string(), error.is_permanent())),
}
if failed.is_some() {
break;
}
}
if let Some((error, permanent)) = failed {
self.give_up_or_retry(
®istration,
at,
attempts,
&error,
permanent,
&mut report,
)
.await?;
break;
}
self.store
.advance(
registration.config.task,
®istration.config.id,
record.body.seq.saturating_add(1),
)
.await?;
attempts = 0;
report.records += 1;
if self.projection.terminal(&record) {
self.store
.delete(registration.config.task, ®istration.config.id)
.await?;
report.completed += 1;
break;
}
}
}
Ok(report)
}
async fn give_up_or_retry(
&self,
registration: &PushRegistration,
at: u64,
attempts: u32,
error: &str,
permanent: bool,
report: &mut PushSweepReport,
) -> Result<(), crate::core::StoreError> {
let exhausted = attempts.saturating_add(1) >= self.max_attempts;
if permanent || exhausted {
let reason = if permanent {
"the deployment no longer permits this destination"
} else {
"the receiver did not answer within the retry ceiling"
};
tracing::warn!(
task = %registration.config.task,
config = %registration.config.id,
url = %registration.config.url,
attempts = attempts.saturating_add(1),
%error,
"abandoning a push registration: {reason}"
);
self.store
.delete(registration.config.task, ®istration.config.id)
.await?;
report.abandoned += 1;
return Ok(());
}
let exponent = attempts.min(8);
self.store
.retry(
registration.config.task,
®istration.config.id,
at.saturating_add(1u64 << exponent),
error,
)
.await?;
report.retries += 1;
Ok(())
}
async fn cleanup_acknowledged_terminal(
&self,
registration: &PushRegistration,
) -> Result<bool, crate::core::StoreError> {
if registration.next_seq <= 1 {
return Ok(false);
}
let previous = self
.journal
.read(
registration.config.task,
registration.next_seq.saturating_sub(1),
)
.await?;
let completed = previous.last().is_some_and(|record| {
record.body.seq.saturating_add(1) == registration.next_seq
&& self.projection.terminal(record)
});
if completed {
self.store
.delete(registration.config.task, ®istration.config.id)
.await?;
}
Ok(completed)
}
}