use std::future::Future;
use std::sync::Arc;
use crate::bus::source::{MessageSource, ReceivedMessage};
use crate::bus::{FailureAction, MessageRouter, RunOptions, TransportError, TransportErrorKind};
use crate::bus::{Message, MessageKind};
pub async fn run_source<R, S, I>(
router: Arc<R>,
mut source: S,
options: RunOptions<I>,
) -> Result<(), TransportError>
where
R: MessageRouter,
S: MessageSource,
I: Send,
{
let service = router.consumer_group();
let transport = source.transport_name();
loop {
let Some(received) = recv_next(&mut source, service, transport).await? else {
break;
};
if let Some(error) = received.decode_error() {
let action = options.failure_policy.resolve(error);
record_transport_failure(service, transport, error.kind(), action);
let kind = received.message().kind;
match action {
FailureAction::Nack => {
let reason = error.to_string();
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::NACK,
crate::telemetry::transport_outcome::NACK,
|| received.nack(&reason),
)
.await?;
}
FailureAction::DeadLetter => {
let reason = error.to_string();
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::DEAD_LETTER,
crate::telemetry::transport_outcome::DEAD_LETTER,
|| received.dead_letter(&reason),
)
.await?;
}
FailureAction::Park => {
let reason = error.to_string();
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::PARK,
crate::telemetry::transport_outcome::PARK,
|| received.park(&reason),
)
.await?;
}
FailureAction::LogAndAck => {
eprintln!("[bus::runner] dropping undecodable message after permanent failure: {error}");
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::ACK,
crate::telemetry::transport_outcome::LOG_AND_ACK,
|| received.ack(),
)
.await?;
}
FailureAction::Stop => return Err(TransportError::permanent(error.to_string())),
}
continue;
}
if !router.handles(received.message().kind, received.message().name()) {
let kind = received.message().kind;
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::ACK,
crate::telemetry::transport_outcome::IGNORED,
|| received.ack(),
)
.await?;
continue;
}
let kind = received.message().kind;
match dispatch(
router.as_ref(),
&options,
received.message(),
received.ordered_delivery(),
)
.await
{
Ok(()) => {
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::ACK,
crate::telemetry::transport_outcome::ACK,
|| received.ack(),
)
.await?;
}
Err(error) if error.should_retain_and_stop() => {
record_transport_failure(
service,
transport,
error.kind(),
crate::telemetry::transport_outcome::NACK,
);
let reason = error.to_string();
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::NACK,
crate::telemetry::transport_outcome::NACK,
|| received.nack(&reason),
)
.await?;
return Err(error);
}
Err(error) => match options.failure_policy.resolve(&error) {
action @ FailureAction::Nack => {
record_transport_failure(service, transport, error.kind(), action);
let reason = error.to_string();
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::NACK,
crate::telemetry::transport_outcome::NACK,
|| received.nack(&reason),
)
.await?;
}
action @ FailureAction::DeadLetter => {
record_transport_failure(service, transport, error.kind(), action);
let reason = error.to_string();
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::DEAD_LETTER,
crate::telemetry::transport_outcome::DEAD_LETTER,
|| received.dead_letter(&reason),
)
.await?;
}
action @ FailureAction::Park => {
record_transport_failure(service, transport, error.kind(), action);
let reason = error.to_string();
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::PARK,
crate::telemetry::transport_outcome::PARK,
|| received.park(&reason),
)
.await?;
}
FailureAction::LogAndAck => {
record_transport_failure(
service,
transport,
error.kind(),
FailureAction::LogAndAck,
);
eprintln!(
"[bus::runner] dropping message '{}' after permanent failure: {error}",
received.message().name()
);
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::ACK,
crate::telemetry::transport_outcome::LOG_AND_ACK,
|| received.ack(),
)
.await?;
}
FailureAction::Stop => {
record_transport_failure(service, transport, error.kind(), FailureAction::Stop);
return Err(error);
}
},
}
}
Ok(())
}
async fn settle_and_record<F, Fut>(
service: Option<&str>,
transport: &str,
kind: MessageKind,
settle_action: &'static str,
outcome: &'static str,
settle: F,
) -> Result<(), TransportError>
where
F: FnOnce() -> Fut,
Fut: Future<Output = Result<(), TransportError>>,
{
match settle().await {
Ok(()) => {
record_transport_message(service, transport, kind, outcome);
Ok(())
}
Err(error) => {
record_transport_failure(
service,
transport,
error.kind(),
crate::telemetry::settle_failure_action(settle_action),
);
Err(error)
}
}
}
async fn recv_next<S: MessageSource>(
source: &mut S,
service: Option<&str>,
transport: &str,
) -> Result<Option<S::Received>, TransportError> {
match source.recv().await {
Ok(received) => Ok(received),
Err(error) => {
record_transport_failure(
service,
transport,
error.kind(),
crate::telemetry::failure_action::RECV_ERROR,
);
Err(error)
}
}
}
async fn dispatch<R: MessageRouter, I>(
router: &R,
options: &RunOptions<I>,
message: &Message,
ordered: Option<&crate::bus::OrderedDelivery>,
) -> Result<(), TransportError> {
#[cfg(feature = "otel")]
{
use tracing::Instrument as _;
let span = transport_receive_span(message);
crate::trace_context::set_span_parent_from_metadata_if_no_current_span(
&span,
&message.metadata,
);
return async {
options
.validate_message_id(message)
.map_err(|err| TransportError::permanent(err.to_string()).with_source(err))?;
router.dispatch_ordered(message, ordered).await
}
.instrument(span)
.await;
}
#[cfg(not(feature = "otel"))]
{
options
.validate_message_id(message)
.map_err(|err| TransportError::permanent(err.to_string()).with_source(err))?;
router.dispatch_ordered(message, ordered).await
}
}
#[cfg(feature = "otel")]
fn transport_receive_span(message: &Message) -> tracing::Span {
crate::telemetry::transport_receive_span(message)
}
fn record_transport_message(
service: Option<&str>,
transport: &str,
kind: MessageKind,
outcome: &str,
) {
#[cfg(feature = "metrics")]
crate::metrics::record_transport_message(service, transport, kind, outcome);
#[cfg(not(feature = "metrics"))]
let _ = (service, transport, kind, outcome);
}
fn record_transport_failure<A>(
service: Option<&str>,
transport: &str,
kind: TransportErrorKind,
action: A,
) where
A: IntoFailureActionLabel,
{
#[cfg(feature = "metrics")]
crate::metrics::record_transport_failure(
service,
transport,
crate::telemetry::transport_failure_class(kind),
action.into_failure_action_label(),
);
#[cfg(not(feature = "metrics"))]
{
let _ = (service, transport, kind);
let _ = action.into_failure_action_label();
}
}
trait IntoFailureActionLabel {
fn into_failure_action_label(self) -> &'static str;
}
impl IntoFailureActionLabel for FailureAction {
fn into_failure_action_label(self) -> &'static str {
crate::telemetry::failure_action_label(self)
}
}
impl IntoFailureActionLabel for &'static str {
fn into_failure_action_label(self) -> &'static str {
self
}
}