use core::fmt;
use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
use std::time::Duration;
use reliar_core::{Classify, FailureKind, MessageType, Publisher, SerializedEnvelope};
use tokio::sync::Semaphore;
use tokio::task::{AbortHandle, Id as TaskId, JoinError, JoinSet};
use tokio::time::Instant;
use tokio_util::sync::CancellationToken;
use tracing::Instrument as _;
use crate::error::ConfigError;
use crate::metrics::{NoopMetrics, OutboxMetrics};
use crate::ordering::Ordering;
use crate::retry::{ExponentialBackoff, RetryPolicy};
use crate::settings::DispatcherSettings;
use crate::store::{
AcquireRequest, CompletedMessage, DeadReason, FailedMessage, FailureOutcome, MessageRef,
OutboxStore,
};
use crate::worker::WorkerId;
pub struct OutboxDispatcher<S, P, M = NoopMetrics, R = ExponentialBackoff> {
store: S,
publisher: Arc<P>,
metrics: Arc<M>,
retry: R,
settings: DispatcherSettings,
worker: WorkerId,
}
impl<S, P, M, R> fmt::Debug for OutboxDispatcher<S, P, M, R> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("OutboxDispatcher")
.field("worker", &self.worker)
.field("settings", &self.settings)
.finish_non_exhaustive()
}
}
impl<S, P> OutboxDispatcher<S, P>
where
S: OutboxStore,
P: Publisher,
{
pub fn builder(store: S, publisher: P) -> OutboxDispatcherBuilder<S, P> {
OutboxDispatcherBuilder {
store,
publisher,
metrics: NoopMetrics,
settings: DispatcherSettings::default(),
retry: DefaultRetry,
}
}
}
#[derive(Clone, Copy, Debug, Default)]
pub struct DefaultRetry;
#[must_use]
pub struct OutboxDispatcherBuilder<S, P, M = NoopMetrics, R = DefaultRetry> {
store: S,
publisher: P,
metrics: M,
settings: DispatcherSettings,
retry: R,
}
impl<S, P, M, R> fmt::Debug for OutboxDispatcherBuilder<S, P, M, R> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("OutboxDispatcherBuilder")
.field("settings", &self.settings)
.finish_non_exhaustive()
}
}
impl<S, P, M, R> OutboxDispatcherBuilder<S, P, M, R> {
pub fn settings(mut self, settings: DispatcherSettings) -> Self {
self.settings = settings;
self
}
pub const fn ordering(mut self, ordering: Ordering) -> Self {
self.settings.ordering = ordering;
self
}
pub fn worker_id(mut self, worker: WorkerId) -> Self {
self.settings.worker_id = Some(worker);
self
}
pub fn metrics<M2: OutboxMetrics>(self, metrics: M2) -> OutboxDispatcherBuilder<S, P, M2, R> {
OutboxDispatcherBuilder {
store: self.store,
publisher: self.publisher,
metrics,
settings: self.settings,
retry: self.retry,
}
}
pub fn retry_policy<R2: RetryPolicy>(self, policy: R2) -> OutboxDispatcherBuilder<S, P, M, R2> {
OutboxDispatcherBuilder {
store: self.store,
publisher: self.publisher,
metrics: self.metrics,
settings: self.settings,
retry: policy,
}
}
}
impl<S, P, M> OutboxDispatcherBuilder<S, P, M, DefaultRetry> {
pub fn build(self) -> Result<OutboxDispatcher<S, P, M, ExponentialBackoff>, ConfigError> {
validate_shared(&self.settings)?;
self.settings.retry.validate()?;
let worker = self.settings.worker_id.clone().unwrap_or_default();
Ok(OutboxDispatcher {
store: self.store,
publisher: Arc::new(self.publisher),
metrics: Arc::new(self.metrics),
retry: self.settings.retry,
settings: self.settings,
worker,
})
}
}
impl<S, P, M, R> OutboxDispatcherBuilder<S, P, M, R>
where
R: RetryPolicy + 'static,
{
pub fn build(self) -> Result<OutboxDispatcher<S, P, M, R>, ConfigError> {
validate_shared(&self.settings)?;
if self.settings.retry != ExponentialBackoff::default() {
return Err(ConfigError::RetryPolicyConflict);
}
if let Some(backoff) =
(&self.retry as &dyn core::any::Any).downcast_ref::<ExponentialBackoff>()
{
backoff.validate()?;
}
let worker = self.settings.worker_id.clone().unwrap_or_default();
Ok(OutboxDispatcher {
store: self.store,
publisher: Arc::new(self.publisher),
metrics: Arc::new(self.metrics),
retry: self.retry,
settings: self.settings,
worker,
})
}
}
fn validate_shared(settings: &DispatcherSettings) -> Result<(), ConfigError> {
if settings.max_in_flight == 0 {
return Err(ConfigError::ZeroInFlight);
}
if settings.batch_size == 0 {
return Err(ConfigError::ZeroBatchSize);
}
if settings.poll_interval.is_zero() {
return Err(ConfigError::ZeroPollInterval {
field: "poll_interval",
});
}
if settings.idle_poll_interval.is_zero() {
return Err(ConfigError::ZeroPollInterval {
field: "idle_poll_interval",
});
}
settings.ordering.validate()?;
if settings.lease <= settings.publish_timeout {
return Err(ConfigError::LeaseTooShort {
lease: settings.lease,
publish_timeout: settings.publish_timeout,
});
}
if settings.store_timeout >= settings.lease / 2 {
return Err(ConfigError::StoreTimeoutTooLong {
store_timeout: settings.store_timeout,
lease: settings.lease,
});
}
#[allow(
clippy::cast_possible_truncation,
reason = "max_in_flight is validated non-zero above; truncation only lowers the \
warning's sensitivity, never causes a false negative that matters"
)]
let max_in_flight = settings.max_in_flight as u32;
let expected_batch_duration = settings
.publish_timeout
.saturating_mul(settings.batch_size)
.checked_div(max_in_flight.max(1))
.unwrap_or(Duration::MAX);
if settings.lease <= expected_batch_duration {
tracing::warn!(
lease_ms = settings.lease.as_millis(),
expected_batch_duration_ms = expected_batch_duration.as_millis(),
"reliar.outbox: lease is not comfortably longer than batch_size × \
publish_timeout ÷ max_in_flight; a healthy batch may still outlive its lease \
(SRS §21.1)"
);
}
if settings.store_timeout > settings.drain_timeout {
tracing::warn!(
store_timeout_ms = settings.store_timeout.as_millis(),
drain_timeout_ms = settings.drain_timeout.as_millis(),
"reliar.outbox: store_timeout is longer than drain_timeout; a slow store call can \
still make shutdown wait past drain_timeout"
);
}
Ok(())
}
#[derive(Debug)]
#[non_exhaustive]
pub enum DispatchError<E> {
Configuration(ConfigError),
Store(E),
}
impl<E: fmt::Display> fmt::Display for DispatchError<E> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Configuration(err) => write!(f, "invalid dispatcher configuration: {err}"),
Self::Store(err) => write!(f, "permanent store error: {err}"),
}
}
}
impl<E: std::error::Error + 'static> std::error::Error for DispatchError<E> {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Self::Configuration(err) => Some(err),
Self::Store(err) => Some(err),
}
}
}
struct PendingRow {
message: MessageRef,
message_type: MessageType,
attempts_before: u32,
envelope: SerializedEnvelope,
}
struct OutstandingTask {
message: MessageRef,
started: Arc<AtomicBool>,
abort: AbortHandle,
}
impl<S, P, M, R> OutboxDispatcher<S, P, M, R>
where
S: OutboxStore + Send + Sync + 'static,
P: Publisher + Send + Sync + 'static,
M: OutboxMetrics + Send + Sync + 'static,
R: RetryPolicy + Send + Sync + 'static,
{
#[allow(
clippy::too_many_lines,
reason = "the claim/publish/drain state machine reads more clearly as one function than \
split across helpers that would each need most of this same local state"
)]
pub async fn run(self, cancel: CancellationToken) -> Result<(), DispatchError<S::Error>> {
let Self {
store,
publisher,
metrics,
retry,
settings,
worker,
} = self;
tracing::info!(
worker.id = %worker,
batch_size = settings.batch_size,
max_in_flight = settings.max_in_flight,
"reliar.outbox: dispatcher starting"
);
let semaphore = Arc::new(Semaphore::new(settings.max_in_flight));
let mut pending: VecDeque<PendingRow> = VecDeque::new();
let mut in_flight: JoinSet<PublishTaskOutcome> = JoinSet::new();
let mut outstanding: HashMap<TaskId, OutstandingTask> = HashMap::new();
let mut next_poll_at = Instant::now();
let lease_tick_period = (settings.lease / 2).max(Duration::from_millis(1));
let mut lease_ticker =
tokio::time::interval_at(Instant::now() + lease_tick_period, lease_tick_period);
lease_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut stats_ticker = (!settings.stats_interval.is_zero()).then(|| {
let mut interval = tokio::time::interval_at(
Instant::now() + settings.stats_interval,
settings.stats_interval,
);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval
});
let mut exit_reason: Option<DispatchError<S::Error>> = None;
let mut unwritten_complete: Vec<(Instant, CompletedMessage)> = Vec::new();
let mut unwritten_fail: Vec<(Instant, FailedMessage)> = Vec::new();
let outcome_retry_interval = settings
.poll_interval
.min(settings.lease / 4)
.max(Duration::from_millis(1));
let mut next_outcome_retry_at = Instant::now();
'run: loop {
let outstanding_count =
outstanding.len() + unwritten_complete.len() + unwritten_fail.len();
let has_capacity = outstanding_count < settings.max_in_flight;
let has_outstanding = outstanding_count > 0;
let has_unwritten = !unwritten_complete.is_empty() || !unwritten_fail.is_empty();
let lease_refs: Vec<MessageRef> = outstanding
.values()
.map(|task| task.message)
.chain(unwritten_complete.iter().map(|(_, c)| c.message))
.chain(unwritten_fail.iter().map(|(_, f)| f.message))
.collect();
tokio::select! {
biased;
() = cancel.cancelled() => {
break 'run;
}
Some(first) = in_flight.join_next_with_id(), if !in_flight.is_empty() => {
let mut results = vec![first];
while let Some(next) = in_flight.try_join_next_with_id() {
results.push(next);
}
record_publish_results(results, &mut outstanding, &mut unwritten_complete, &mut unwritten_fail, &retry, metrics.as_ref());
spawn_ready(&mut pending, &mut in_flight, &mut outstanding, &publisher, &metrics, &semaphore, settings.publish_timeout);
}
_ = lease_ticker.tick(), if has_outstanding => {
renew_leases(&store, &worker, &lease_refs, settings.lease, settings.store_timeout).await;
}
() = tick_optional(&mut stats_ticker) => {
report_stats(&store, metrics.as_ref(), settings.store_timeout).await;
}
retry_result = retry_unwritten_outcomes_when_due(
&mut unwritten_complete,
&mut unwritten_fail,
&store,
&worker,
settings.store_timeout,
settings.lease,
next_outcome_retry_at,
), if has_unwritten => {
let RetryOutcome { permanent, any_failed } = retry_result;
let still_unwritten =
!unwritten_complete.is_empty() || !unwritten_fail.is_empty();
next_outcome_retry_at = if any_failed && still_unwritten {
Instant::now() + outcome_retry_interval
} else {
Instant::now()
};
if let Some(err) = permanent {
exit_reason = Some(DispatchError::Store(err));
break 'run;
}
}
() = tokio::time::sleep_until(next_poll_at), if has_capacity => {
let capacity = settings.max_in_flight.saturating_sub(outstanding_count);
let capacity = u32::try_from(capacity).unwrap_or(u32::MAX);
let want = settings.batch_size.min(capacity);
let request = AcquireRequest::new(worker.clone())
.batch_size(want)
.lease(settings.lease)
.ordering(settings.ordering);
let claim_span = tracing::info_span!(
"reliar.outbox.claim",
worker.id = %worker,
batch.requested = want,
batch.claimed = tracing::field::Empty,
);
let acquired = tokio::time::timeout(settings.store_timeout, store.acquire(request).instrument(claim_span.clone())).await;
match acquired {
Ok(Ok(batch)) => {
claim_span.record("batch.claimed", batch.records.len());
metrics.claimed(batch.records.len());
if !batch.poisoned.is_empty() {
for poisoned in &batch.poisoned {
tracing::warn!(
message.id = %poisoned.id,
"reliar.outbox.dead: row undecodable, moved to dead by the store"
);
}
metrics.dead(batch.poisoned.len(), DeadReason::Undecodable);
}
if batch.records.is_empty() && batch.poisoned.is_empty() {
next_poll_at = Instant::now() + settings.idle_poll_interval;
} else {
next_poll_at = Instant::now() + settings.poll_interval;
for record in batch.records {
pending.push_back(PendingRow {
message: record.message_ref(),
message_type: record.envelope.message_type.clone(),
attempts_before: record.attempts,
envelope: record.envelope,
});
}
spawn_ready(&mut pending, &mut in_flight, &mut outstanding, &publisher, &metrics, &semaphore, settings.publish_timeout);
}
}
Ok(Err(err)) => {
if err.kind() == FailureKind::Permanent {
exit_reason = Some(DispatchError::Store(err));
break 'run;
}
tracing::error!(error = %err, "reliar.outbox.claim failed; backing off");
next_poll_at = Instant::now() + settings.idle_poll_interval;
}
Err(_elapsed) => {
tracing::error!(
store_timeout_ms = settings.store_timeout.as_millis(),
"reliar.outbox.claim timed out; backing off"
);
next_poll_at = Instant::now() + settings.idle_poll_interval;
}
}
}
}
}
let never_started: Vec<MessageRef> = pending.drain(..).map(|row| row.message).collect();
let not_yet_permitted: Vec<MessageRef> = {
let mut refs = Vec::new();
outstanding.retain(|_, task| {
if task.started.load(AtomicOrdering::Relaxed) {
true
} else {
task.abort.abort();
refs.push(task.message);
false
}
});
refs
};
let release_immediately: Vec<MessageRef> =
never_started.into_iter().chain(not_yet_permitted).collect();
if !release_immediately.is_empty() {
tracing::warn!(
count = release_immediately.len(),
"reliar.outbox: releasing rows that never started publishing — they will be \
reclaimed by the next owner (SRS §26.1)"
);
if let Err(err) = bounded(
settings.store_timeout,
store.release(&worker, &release_immediately),
)
.await
{
tracing::warn!(error = %err, "reliar.outbox: release of never-started rows failed; leases will expire naturally");
}
}
let drain_permanent_error = drain(
&mut in_flight,
&mut outstanding,
&mut unwritten_complete,
&mut unwritten_fail,
&store,
&worker,
&retry,
metrics.as_ref(),
settings.drain_timeout,
settings.store_timeout,
settings.lease,
)
.await;
if exit_reason.is_none() {
exit_reason = drain_permanent_error.map(DispatchError::Store);
}
tracing::info!(worker.id = %worker, "reliar.outbox: dispatcher stopped");
match exit_reason {
Some(err) => Err(err),
None => Ok(()),
}
}
}
fn spawn_ready<P, M>(
pending: &mut VecDeque<PendingRow>,
in_flight: &mut JoinSet<PublishTaskOutcome>,
outstanding: &mut HashMap<TaskId, OutstandingTask>,
publisher: &Arc<P>,
metrics: &Arc<M>,
semaphore: &Arc<Semaphore>,
publish_timeout: Duration,
) where
P: Publisher + Send + Sync + 'static,
M: OutboxMetrics + Send + Sync + 'static,
{
while let Some(row) = pending.pop_front() {
let message = row.message;
let started = Arc::new(AtomicBool::new(false));
let job = PublishJob {
message,
message_type: row.message_type,
attempts_before: row.attempts_before,
envelope: row.envelope,
publish_timeout,
};
let abort = in_flight.spawn(publish_one(
job,
Arc::clone(publisher),
Arc::clone(metrics),
Arc::clone(semaphore),
Arc::clone(&started),
));
outstanding.insert(
abort.id(),
OutstandingTask {
message,
started,
abort,
},
);
}
}
#[allow(
clippy::too_many_arguments,
reason = "each parameter is a distinct piece of the run loop's state or a distinct timeout \
budget (drain vs. store); bundling them would just move the count into a struct \
with the same number of fields"
)]
async fn drain<S, M, R>(
in_flight: &mut JoinSet<PublishTaskOutcome>,
outstanding: &mut HashMap<TaskId, OutstandingTask>,
unwritten_complete: &mut Vec<(Instant, CompletedMessage)>,
unwritten_fail: &mut Vec<(Instant, FailedMessage)>,
store: &S,
worker: &WorkerId,
retry: &R,
metrics: &M,
drain_timeout: Duration,
store_timeout: Duration,
lease: Duration,
) -> Option<S::Error>
where
S: OutboxStore,
M: OutboxMetrics,
R: RetryPolicy,
{
let drain_deadline = Instant::now() + drain_timeout;
let mut permanent_error = None;
'drain: loop {
if in_flight.is_empty() || Instant::now() >= drain_deadline || permanent_error.is_some() {
break 'drain;
}
tokio::select! {
biased;
Some(first) = in_flight.join_next_with_id() => {
let mut results = vec![first];
while let Some(next) = in_flight.try_join_next_with_id() {
results.push(next);
}
record_publish_results(results, outstanding, unwritten_complete, unwritten_fail, retry, metrics);
permanent_error = retry_unwritten_outcomes(unwritten_complete, unwritten_fail, store, worker, store_timeout, lease).await.permanent;
}
() = tokio::time::sleep_until(drain_deadline) => {
break 'drain;
}
}
}
if permanent_error.is_none() {
permanent_error = retry_unwritten_outcomes(
unwritten_complete,
unwritten_fail,
store,
worker,
store_timeout,
lease,
)
.await
.permanent;
}
if !unwritten_complete.is_empty() {
tracing::warn!(
count = unwritten_complete.len(),
"reliar.outbox: rows published successfully but not yet marked complete are left to \
their lease rather than released (SRS §23.2)"
);
}
let mut release_now: Vec<MessageRef> = outstanding.values().map(|task| task.message).collect();
release_now.extend(unwritten_fail.drain(..).map(|(_, failed)| failed.message));
if !release_now.is_empty() {
tracing::warn!(
count = release_now.len(),
"reliar.outbox: releasing rows whose publish failed or never resolved by the drain \
timeout — the next owner will retry them (SRS §22.1, §26.1)"
);
if let Err(err) = bounded(store_timeout, store.release(worker, &release_now)).await {
tracing::warn!(error = %err, "reliar.outbox: release on drain failed; leases will expire naturally");
}
}
permanent_error
}
struct RetryOutcome<E> {
permanent: Option<E>,
any_failed: bool,
}
async fn retry_unwritten_outcomes_when_due<S: OutboxStore>(
unwritten_complete: &mut Vec<(Instant, CompletedMessage)>,
unwritten_fail: &mut Vec<(Instant, FailedMessage)>,
store: &S,
worker: &WorkerId,
store_timeout: Duration,
lease: Duration,
due_at: Instant,
) -> RetryOutcome<S::Error> {
tokio::time::sleep_until(due_at).await;
retry_unwritten_outcomes(
unwritten_complete,
unwritten_fail,
store,
worker,
store_timeout,
lease,
)
.await
}
async fn retry_unwritten_outcomes<S: OutboxStore>(
unwritten_complete: &mut Vec<(Instant, CompletedMessage)>,
unwritten_fail: &mut Vec<(Instant, FailedMessage)>,
store: &S,
worker: &WorkerId,
store_timeout: Duration,
lease: Duration,
) -> RetryOutcome<S::Error> {
let mut any_failed = false;
if !unwritten_complete.is_empty() {
let items: Vec<CompletedMessage> =
unwritten_complete.iter().map(|(_, c)| c.clone()).collect();
let count = items.len();
match bounded(store_timeout, store.complete(worker, &items)).await {
Ok(_affected) => unwritten_complete.clear(),
Err(TimedOut::Store(err)) if err.kind() == FailureKind::Permanent => {
return RetryOutcome {
permanent: Some(err),
any_failed: true,
};
}
Err(err) => {
tracing::warn!(error = %err, count, "reliar.outbox.complete failed; retrying at the next scheduled attempt");
any_failed = true;
}
}
}
expire_past_lease(unwritten_complete, lease, "complete");
if !unwritten_fail.is_empty() {
let items: Vec<FailedMessage> = unwritten_fail.iter().map(|(_, f)| f.clone()).collect();
let count = items.len();
match bounded(store_timeout, store.fail(worker, &items)).await {
Ok(_affected) => unwritten_fail.clear(),
Err(TimedOut::Store(err)) if err.kind() == FailureKind::Permanent => {
return RetryOutcome {
permanent: Some(err),
any_failed: true,
};
}
Err(err) => {
tracing::warn!(error = %err, count, "reliar.outbox.fail failed; retrying at the next scheduled attempt");
any_failed = true;
}
}
}
expire_past_lease(unwritten_fail, lease, "fail");
RetryOutcome {
permanent: None,
any_failed,
}
}
fn expire_past_lease<T>(unwritten: &mut Vec<(Instant, T)>, lease: Duration, op: &'static str) {
let now = Instant::now();
let expired = unwritten
.iter()
.filter(|(since, _)| now.saturating_duration_since(*since) > lease)
.count();
if expired > 0 {
tracing::warn!(
op,
count = expired,
"reliar.outbox: giving up on an outcome write retried longer than the lease; the \
row's lease will lapse and another worker may reclaim it (SRS §22.1, M2)"
);
unwritten.retain(|(since, _)| now.saturating_duration_since(*since) <= lease);
}
}
async fn tick_optional(interval: &mut Option<tokio::time::Interval>) {
match interval {
Some(interval) => {
interval.tick().await;
}
None => std::future::pending::<()>().await,
}
}
async fn bounded<T, E, F>(budget: Duration, fut: F) -> Result<T, TimedOut<E>>
where
F: Future<Output = Result<T, E>>,
{
match tokio::time::timeout(budget, fut).await {
Ok(Ok(value)) => Ok(value),
Ok(Err(err)) => Err(TimedOut::Store(err)),
Err(_elapsed) => Err(TimedOut::Timeout(budget)),
}
}
enum TimedOut<E> {
Store(E),
Timeout(Duration),
}
impl<E: fmt::Display> fmt::Display for TimedOut<E> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Store(err) => write!(f, "{err}"),
Self::Timeout(budget) => write!(f, "timed out after {budget:?}"),
}
}
}
async fn renew_leases<S: OutboxStore>(
store: &S,
worker: &WorkerId,
outstanding: &[MessageRef],
lease: Duration,
store_timeout: Duration,
) {
match bounded(
store_timeout,
store.extend_lease(worker, outstanding, lease),
)
.await
{
Ok(affected) => {
let affected = usize::try_from(affected).unwrap_or(usize::MAX);
if affected < outstanding.len() {
tracing::warn!(
affected,
outstanding = outstanding.len(),
"reliar.outbox: extend_lease renewed fewer rows than outstanding; some rows' \
lease may already be reclaimed by another worker"
);
}
}
Err(err) => {
tracing::warn!(error = %err, "reliar.outbox: extend_lease failed; in-flight rows may lose their lease");
}
}
}
async fn report_stats<S: OutboxStore, M: OutboxMetrics>(
store: &S,
metrics: &M,
store_timeout: Duration,
) {
match bounded(store_timeout, store.stats()).await {
Ok(stats) => {
metrics.pending(stats.pending);
metrics.expired_pending(stats.expired_pending);
if let Some(lag) = stats.lag() {
metrics.oldest_pending_age(lag);
}
}
Err(err) => {
tracing::warn!(error = %err, "reliar.outbox: stats poll failed");
}
}
}
enum PublishTaskOutcome {
Published {
message: MessageRef,
message_type: MessageType,
},
Failed {
message: MessageRef,
attempts_before: u32,
kind: FailureKind,
error: String,
},
}
struct PublishJob {
message: MessageRef,
message_type: MessageType,
attempts_before: u32,
envelope: SerializedEnvelope,
publish_timeout: Duration,
}
#[tracing::instrument(
name = "reliar.outbox.publish",
skip_all,
fields(
message.id = %job.message.id,
message.type = %job.message_type,
attempt = job.attempts_before + 1
)
)]
async fn publish_one<P, M>(
job: PublishJob,
publisher: Arc<P>,
metrics: Arc<M>,
semaphore: Arc<Semaphore>,
started: Arc<AtomicBool>,
) -> PublishTaskOutcome
where
P: Publisher,
M: OutboxMetrics,
{
let PublishJob {
message,
message_type,
attempts_before,
envelope,
publish_timeout,
} = job;
let _permit = match semaphore.acquire_owned().await {
Ok(permit) => permit,
Err(_closed) => {
return PublishTaskOutcome::Failed {
message,
attempts_before,
kind: FailureKind::Transient,
error: "reliar-outbox: publish semaphore closed unexpectedly".to_string(),
};
}
};
started.store(true, AtomicOrdering::Relaxed);
let publish_started_at = Instant::now();
let outcome = tokio::time::timeout(publish_timeout, publisher.publish(&envelope)).await;
metrics.publish_duration(publish_started_at.elapsed(), &message_type);
match outcome {
Ok(Ok(())) => PublishTaskOutcome::Published {
message,
message_type,
},
Ok(Err(err)) => {
let kind = err.kind();
PublishTaskOutcome::Failed {
message,
attempts_before,
kind,
error: err.to_string(),
}
}
Err(_elapsed) => PublishTaskOutcome::Failed {
message,
attempts_before,
kind: FailureKind::Transient,
error: format!("reliar-outbox: publish timed out after {publish_timeout:?}"),
},
}
}
fn record_publish_results<M, R>(
results: Vec<Result<(TaskId, PublishTaskOutcome), JoinError>>,
outstanding: &mut HashMap<TaskId, OutstandingTask>,
unwritten_complete: &mut Vec<(Instant, CompletedMessage)>,
unwritten_fail: &mut Vec<(Instant, FailedMessage)>,
retry: &R,
metrics: &M,
) where
M: OutboxMetrics,
R: RetryPolicy,
{
for joined in results {
let outcome = match joined {
Ok((id, outcome)) => {
outstanding.remove(&id);
outcome
}
Err(join_error) => {
outstanding.remove(&join_error.id());
if join_error.is_panic() {
tracing::error!(
error = %join_error,
"reliar.outbox.publish task panicked; its lease will expire and the row will be reclaimed"
);
} else {
tracing::debug!(
error = %join_error,
"reliar.outbox.publish task was aborted before completing (shutdown); its lease will expire and the row will be reclaimed"
);
}
continue;
}
};
match outcome {
PublishTaskOutcome::Published {
message,
message_type,
} => {
metrics.published(1, &message_type);
unwritten_complete.push((Instant::now(), CompletedMessage::new(message)));
}
PublishTaskOutcome::Failed {
message,
attempts_before,
kind,
error,
} => {
let outcome = retry.next(attempts_before, kind);
match &outcome {
FailureOutcome::Retry { delay } => {
tracing::debug!(
message.id = %message.id,
attempt = attempts_before + 1,
delay_ms = delay.as_millis(),
"reliar.outbox.retry: scheduled"
);
metrics.retried(1, kind);
}
FailureOutcome::Dead { reason } => {
tracing::warn!(
message.id = %message.id,
attempt = attempts_before + 1,
reason = ?reason,
"reliar.outbox.dead: message moved to dead"
);
metrics.dead(1, *reason);
}
}
unwritten_fail.push((Instant::now(), FailedMessage::new(message, error, outcome)));
}
}
}
}