use std::fmt::Debug;
use std::sync::Arc;
use crate::core::StoreError;
use crate::journal::{JournalStore, Record, RecordKind};
use async_trait::async_trait;
use super::{Delivered, PushMessage, PushNamespace, PushRegistration, PushStore, PushTransport};
#[async_trait]
pub trait Projection: Send + Sync + Debug {
async fn messages(&self, record: &Record) -> Result<Vec<PushMessage>, 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 parked: usize,
pub unserved: usize,
pub saturated: bool,
}
impl PushSweepReport {
#[must_use]
pub const fn needs_attention(&self) -> bool {
self.saturated || self.parked > 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,
max_in_flight: usize,
}
#[derive(Debug, Default, Clone, Copy)]
struct Progress {
records: usize,
deliveries: usize,
retries: usize,
completed: usize,
parked: usize,
}
impl DeliveryWorker {
pub const DEFAULT_MAX_ATTEMPTS: u32 = 32;
pub const MAX_RETRY_AFTER: u64 = 3_600;
pub const DEFAULT_MAX_IN_FLIGHT: usize = 16;
}
const DRAIN_PAGE: usize = 256;
impl DeliveryWorker {
#[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,
max_in_flight: Self::DEFAULT_MAX_IN_FLIGHT,
}
}
#[must_use]
pub const fn max_in_flight(mut self, n: usize) -> Self {
assert!(
n > 0,
"a push concurrency of zero contacts no receiver and reports a \
quiet plane; schedule no worker instead"
);
self.max_in_flight = n;
self
}
#[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> {
use futures_util::StreamExt as _;
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()
};
let outcomes: Vec<_> = futures_util::stream::iter(
batch
.rows
.into_iter()
.take(limit)
.map(|registration| self.deliver_one(registration, at)),
)
.buffer_unordered(self.max_in_flight)
.collect()
.await;
for outcome in outcomes {
let progress = outcome?;
report.records += progress.records;
report.deliveries += progress.deliveries;
report.retries += progress.retries;
report.completed += progress.completed;
report.parked += progress.parked;
}
Ok(report)
}
async fn deliver_one(
&self,
registration: PushRegistration,
at: u64,
) -> Result<Progress, crate::core::StoreError> {
let mut progress = Progress::default();
let mut attempts = registration.attempts;
let mut records = self
.journal
.read(registration.config.task, registration.next_seq)
.await?;
records.truncate(DRAIN_PAGE);
if records.is_empty() && self.cleanup_acknowledged_terminal(®istration).await? {
progress.completed += 1;
return Ok(progress);
}
for record in records {
let messages = match self.projection.messages(&record).await {
Ok(messages) => messages,
Err(error) => {
self.give_up_or_retry(
®istration,
at,
attempts,
&Failure::transient(error.to_string()),
&mut progress,
)
.await?;
break;
}
};
let mut failed = None;
for message in messages {
match self
.transport
.deliver(®istration.config, &message, at)
.await
{
Ok(Delivered::Accepted) => progress.deliveries += 1,
Ok(other) => {
failed = Some(Failure {
error: format!("receiver outcome: {other:?}"),
permanent: other.is_permanent(),
retry_after: other.retry_after(),
});
}
Err(error) => {
failed = Some(Failure {
permanent: error.is_permanent(),
error: error.to_string(),
retry_after: None,
});
}
}
if failed.is_some() {
break;
}
}
if let Some(failure) = failed {
self.give_up_or_retry(®istration, at, attempts, &failure, &mut progress)
.await?;
break;
}
self.store
.advance(
registration.config.task,
®istration.config.id,
record.body.seq.saturating_add(1),
)
.await?;
attempts = 0;
progress.records += 1;
if self.projection.terminal(&record) {
self.store
.delete(registration.config.task, ®istration.config.id)
.await?;
progress.completed += 1;
break;
}
}
Ok(progress)
}
async fn give_up_or_retry(
&self,
registration: &PushRegistration,
at: u64,
attempts: u32,
failure: &Failure,
report: &mut Progress,
) -> Result<(), crate::core::StoreError> {
let exhausted = attempts.saturating_add(1) >= self.max_attempts;
if failure.permanent || exhausted {
let reason = if failure.permanent {
"this destination answered permanently"
} 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 = %failure.error,
"parking a push registration: {reason} — its cursor is kept, so \
`unpark` resumes at the first record the receiver never took"
);
self.store
.park(
registration.config.task,
®istration.config.id,
&failure.error,
)
.await?;
report.parked += 1;
return Ok(());
}
self.store
.retry(
registration.config.task,
®istration.config.id,
Self::next_attempt_at(registration, at, attempts, failure.retry_after),
&failure.error,
)
.await?;
report.retries += 1;
Ok(())
}
fn next_attempt_at(
registration: &PushRegistration,
at: u64,
attempts: u32,
advice: Option<u64>,
) -> u64 {
if let Some(seconds) = advice {
return at.saturating_add(seconds.clamp(1, Self::MAX_RETRY_AFTER));
}
let window = 1u64 << attempts.min(8);
let half = window / 2;
let offset = spread(registration, attempts) % (half.saturating_add(1));
at.saturating_add(half.saturating_add(offset).max(1))
}
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)
}
}
struct Failure {
error: String,
permanent: bool,
retry_after: Option<u64>,
}
impl Failure {
fn transient(error: String) -> Self {
Self {
error,
permanent: false,
retry_after: None,
}
}
}
fn spread(registration: &PushRegistration, attempts: u32) -> u64 {
use sha2::Digest as _;
let mut hasher = sha2::Sha256::new();
hasher.update(registration.config.task.to_string().as_bytes());
hasher.update([0x1f]);
hasher.update(registration.config.id.as_bytes());
hasher.update([0x1f]);
hasher.update(attempts.to_be_bytes());
let digest = hasher.finalize();
u64::from_be_bytes(digest[..8].try_into().expect("SHA-256 is 32 bytes"))
}