use std::future::Future;
use std::pin::Pin;
use crate::bus::MessagePublisher;
use crate::outbox::{OutboxMessage, OutboxPublishHook};
use crate::repository::RepositoryError;
use super::outbox_dispatch::{publish_and_settle, SettleOutcome};
use super::OutboxStore;
pub struct BusOutboxPublishHook<S, P> {
store: S,
publisher: P,
max_attempts: u32,
service_name: Option<String>,
}
impl<S, P> BusOutboxPublishHook<S, P> {
pub fn new(store: S, publisher: P, max_attempts: u32) -> Self {
Self {
store,
publisher,
max_attempts,
service_name: None,
}
}
pub fn with_service(mut self, service_name: Option<String>) -> Self {
self.service_name = service_name;
self
}
}
impl<S, P> OutboxPublishHook for BusOutboxPublishHook<S, P>
where
S: OutboxStore,
P: MessagePublisher,
{
fn publish_claimed<'a>(
&'a self,
claimed: Vec<OutboxMessage>,
) -> Pin<Box<dyn Future<Output = Result<(), RepositoryError>> + Send + 'a>> {
Box::pin(async move {
let settled = publish_and_settle(
&self.store,
&self.publisher,
claimed,
self.max_attempts,
std::num::NonZeroUsize::MIN,
)
.await?;
self.record_outbox_outcomes(&settled).await;
Ok(())
})
}
}
impl<S, P> BusOutboxPublishHook<S, P>
where
S: OutboxStore,
{
async fn record_outbox_outcomes(&self, settled: &SettleOutcome) {
#[cfg(feature = "metrics")]
{
let service = self.service_name.as_deref();
crate::metrics::record_outbox_messages(
service,
crate::telemetry::outbox_outcome::PUBLISHED,
settled.published,
);
crate::metrics::record_outbox_messages(
service,
crate::telemetry::outbox_outcome::RELEASED,
settled.released,
);
crate::metrics::record_outbox_messages(
service,
crate::telemetry::outbox_outcome::FAILED,
settled.failed,
);
super::outbox_dispatch::record_backlog_gauges(&self.store, service).await;
}
#[cfg(not(feature = "metrics"))]
let _ = settled;
}
}