use std::future::Future;
use std::sync::Arc;
use super::source::{MessageSource, ReceivedMessage};
use super::{FailureAction, MessageRouter, RunOptions, TransportError, TransportErrorKind};
use super::{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()).await {
Ok(()) => {
settle_and_record(
service,
transport,
kind,
crate::telemetry::transport_outcome::ACK,
crate::telemetry::transport_outcome::ACK,
|| received.ack(),
)
.await?;
}
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,
) -> 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(message).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(message).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
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bus::{FailurePolicy, Handlers, MessageKind};
use std::collections::VecDeque;
use std::future::Future;
use std::sync::Mutex;
fn block_on<F: Future>(future: F) -> F::Output {
use std::ptr;
use std::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
const VTABLE: RawWakerVTable = RawWakerVTable::new(
|_| RawWaker::new(ptr::null(), &VTABLE),
|_| {},
|_| {},
|_| {},
);
let raw = RawWaker::new(ptr::null(), &VTABLE);
let waker = unsafe { Waker::from_raw(raw) };
let mut cx = Context::from_waker(&waker);
let mut future = std::pin::pin!(future);
loop {
if let Poll::Ready(output) = future.as_mut().poll(&mut cx) {
return output;
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum Event {
Handled(String),
Ack,
Nack(String),
DeadLetter(String),
Park(String),
}
struct Recorder {
events: Mutex<Vec<Event>>,
}
impl Recorder {
fn new() -> Arc<Self> {
Arc::new(Self {
events: Mutex::new(Vec::new()),
})
}
fn push(&self, event: Event) {
self.events.lock().unwrap().push(event);
}
fn events(&self) -> Vec<Event> {
self.events.lock().unwrap().clone()
}
}
struct FakeReceived {
message: Message,
recorder: Arc<Recorder>,
settle_ok: bool,
decode_error: Option<TransportError>,
}
impl FakeReceived {
fn settle(self, event: Event) -> Result<(), TransportError> {
self.recorder.push(event);
if self.settle_ok {
Ok(())
} else {
Err(TransportError::retryable("settle failed"))
}
}
}
impl ReceivedMessage for FakeReceived {
fn message(&self) -> &Message {
&self.message
}
fn decode_error(&self) -> Option<&TransportError> {
self.decode_error.as_ref()
}
async fn ack(self) -> Result<(), TransportError> {
self.settle(Event::Ack)
}
async fn nack(self, reason: &str) -> Result<(), TransportError> {
self.settle(Event::Nack(reason.to_string()))
}
async fn dead_letter(self, reason: &str) -> Result<(), TransportError> {
self.settle(Event::DeadLetter(reason.to_string()))
}
async fn park(self, reason: &str) -> Result<(), TransportError> {
self.settle(Event::Park(reason.to_string()))
}
}
struct FakeSource {
queue: VecDeque<Message>,
recorder: Arc<Recorder>,
settle_ok: bool,
recv_error: bool,
decode_error: bool,
}
impl MessageSource for FakeSource {
type Received = FakeReceived;
async fn recv(&mut self) -> Result<Option<FakeReceived>, TransportError> {
if self.recv_error {
return Err(TransportError::retryable("recv failed"));
}
let decode_error = self.decode_error;
Ok(self.queue.pop_front().map(|message| FakeReceived {
message,
recorder: self.recorder.clone(),
settle_ok: self.settle_ok,
decode_error: decode_error
.then(|| TransportError::permanent("corrupt row: name failed to decode")),
}))
}
}
fn event_message(name: &str, id: Option<&str>) -> Message {
let mut message = Message::new(name, MessageKind::Event, b"{}".to_vec());
if let Some(id) = id {
message = message.with_id(id);
}
message
}
fn router(recorder: &Arc<Recorder>) -> Arc<Handlers> {
let ok = recorder.clone();
let retryable = recorder.clone();
let permanent = recorder.clone();
Arc::new(
Handlers::new()
.on_event("ok", move |msg: &Message| {
let ok = ok.clone();
let name = msg.name().to_string();
async move {
ok.push(Event::Handled(name));
Ok(())
}
})
.on_event("retryable", move |msg: &Message| {
let retryable = retryable.clone();
let name = msg.name().to_string();
async move {
retryable.push(Event::Handled(name));
Err(TransportError::retryable("infra"))
}
})
.on_event("permanent", move |msg: &Message| {
let permanent = permanent.clone();
let name = msg.name().to_string();
async move {
permanent.push(Event::Handled(name));
Err(TransportError::permanent("nope"))
}
}),
)
}
struct RunResult {
outcome: Result<(), TransportError>,
events: Vec<Event>,
}
fn run<I: Send>(messages: Vec<Message>, options: RunOptions<I>) -> RunResult {
run_with(messages, options, true, false)
}
fn run_with<I: Send>(
messages: Vec<Message>,
options: RunOptions<I>,
settle_ok: bool,
recv_error: bool,
) -> RunResult {
let recorder = Recorder::new();
let svc = router(&recorder);
let source = FakeSource {
queue: messages.into_iter().collect(),
recorder: recorder.clone(),
settle_ok,
recv_error,
decode_error: false,
};
let outcome = block_on(run_source(svc, source, options));
RunResult {
outcome,
events: recorder.events(),
}
}
#[test]
fn success_dispatches_then_acks_in_order() {
let result = run(vec![event_message("ok", None)], RunOptions::idempotent());
assert!(result.outcome.is_ok());
assert_eq!(
result.events,
vec![Event::Handled("ok".to_string()), Event::Ack]
);
}
#[test]
fn processes_every_message_then_stops_on_none() {
let result = run(
vec![event_message("ok", None), event_message("ok", None)],
RunOptions::idempotent(),
);
assert!(result.outcome.is_ok());
assert_eq!(
result.events,
vec![
Event::Handled("ok".to_string()),
Event::Ack,
Event::Handled("ok".to_string()),
Event::Ack,
]
);
}
#[cfg(feature = "metrics")]
#[test]
fn metrics_record_success_and_failure_settlement_outcomes() {
let _guard = crate::metrics::lock_for_tests();
crate::metrics::reset_for_tests();
let recorder = Recorder::new();
let svc = Arc::new(
router(&recorder)
.as_ref()
.clone()
.named("metrics-runner-settlement"),
);
let source = FakeSource {
queue: vec![event_message("ok", None), event_message("retryable", None)]
.into_iter()
.collect(),
recorder,
settle_ok: true,
recv_error: false,
decode_error: false,
};
let outcome = block_on(run_source(svc, source, RunOptions::idempotent()));
assert!(outcome.is_ok());
let text = crate::metrics::prometheus_text();
assert!(
text.contains(
"distributed_transport_messages_total{service=\"metrics-runner-settlement\",transport=\"unknown\",message_kind=\"event\",outcome=\"ack\"} 1"
),
"metrics should include the ack outcome:\n{text}"
);
assert!(
text.contains(
"distributed_transport_messages_total{service=\"metrics-runner-settlement\",transport=\"unknown\",message_kind=\"event\",outcome=\"nack\"} 1"
),
"metrics should include the nack outcome:\n{text}"
);
assert!(
text.contains(
"distributed_transport_failures_total{service=\"metrics-runner-settlement\",transport=\"unknown\",failure_class=\"retryable\",action=\"nack\"} 1"
),
"metrics should include the retryable failure action:\n{text}"
);
}
#[cfg(feature = "metrics")]
#[test]
fn metrics_record_settle_failures_before_propagating() {
let _guard = crate::metrics::lock_for_tests();
crate::metrics::reset_for_tests();
let recorder = Recorder::new();
let svc = Arc::new(
router(&recorder)
.as_ref()
.clone()
.named("metrics-runner-settle-failure"),
);
let source = FakeSource {
queue: vec![event_message("ok", None)].into_iter().collect(),
recorder,
settle_ok: false,
recv_error: false,
decode_error: false,
};
let outcome = block_on(run_source(svc, source, RunOptions::idempotent()));
assert!(outcome
.expect_err("settle error should propagate")
.is_retryable());
let text = crate::metrics::prometheus_text();
assert!(
text.contains(
"distributed_transport_failures_total{service=\"metrics-runner-settle-failure\",transport=\"unknown\",failure_class=\"retryable\",action=\"settle_ack\"} 1"
),
"metrics should include the settle failure:\n{text}"
);
assert!(
!text.contains(
"distributed_transport_messages_total{service=\"metrics-runner-settle-failure\",transport=\"unknown\",message_kind=\"event\",outcome=\"ack\"} 1"
),
"settle failure should not record an ack outcome:\n{text}"
);
}
#[test]
fn retryable_failure_nacks_without_acking() {
let result = run(
vec![event_message("retryable", None)],
RunOptions::idempotent(),
);
assert!(result.outcome.is_ok());
assert_eq!(
result.events.first(),
Some(&Event::Handled("retryable".to_string()))
);
assert!(matches!(result.events.get(1), Some(Event::Nack(_))));
assert!(!result.events.contains(&Event::Ack));
}
#[test]
fn permanent_failure_dead_letters_under_default_policy() {
let result = run(
vec![event_message("permanent", None)],
RunOptions::idempotent(),
);
assert!(result.outcome.is_ok());
assert_eq!(
result.events.first(),
Some(&Event::Handled("permanent".to_string()))
);
assert!(matches!(result.events.get(1), Some(Event::DeadLetter(_))));
}
#[test]
fn permanent_failure_parks_under_park_policy() {
let result = run(
vec![event_message("permanent", None)],
RunOptions::idempotent().with_failure_policy(FailurePolicy::Park),
);
assert!(result.outcome.is_ok());
assert!(matches!(result.events.get(1), Some(Event::Park(_))));
}
#[test]
fn permanent_failure_logs_and_acks_under_log_and_ack_policy() {
let result = run(
vec![event_message("permanent", None)],
RunOptions::idempotent().with_failure_policy(FailurePolicy::LogAndAck),
);
assert!(result.outcome.is_ok());
assert_eq!(result.events.get(1), Some(&Event::Ack));
}
#[test]
fn permanent_failure_nacks_under_retry_policy() {
let result = run(
vec![event_message("permanent", None)],
RunOptions::idempotent().with_failure_policy(FailurePolicy::Retry),
);
assert!(result.outcome.is_ok());
assert!(matches!(result.events.get(1), Some(Event::Nack(_))));
}
#[test]
fn stop_policy_returns_error_without_settling() {
let result = run(
vec![event_message("permanent", None), event_message("ok", None)],
RunOptions::idempotent().with_failure_policy(FailurePolicy::Stop),
);
let err = result
.outcome
.expect_err("stop policy should surface the error");
assert!(err.is_permanent());
assert_eq!(result.events, vec![Event::Handled("permanent".to_string())]);
}
#[test]
fn inbox_mode_rejects_message_without_stable_id_before_dispatch() {
let result = run(vec![event_message("ok", None)], RunOptions::inbox(()));
assert!(result.outcome.is_ok());
assert_eq!(result.events.len(), 1);
match &result.events[0] {
Event::DeadLetter(reason) => {
assert!(reason.contains("stable message id is required but missing"))
}
other => panic!("expected dead-letter, got {other:?}"),
}
assert!(!result.events.iter().any(|e| matches!(e, Event::Handled(_))));
}
#[test]
fn inbox_mode_dispatches_when_stable_id_is_present() {
let result = run(
vec![event_message("ok", Some("evt-1"))],
RunOptions::inbox(()),
);
assert!(result.outcome.is_ok());
assert_eq!(
result.events,
vec![Event::Handled("ok".to_string()), Event::Ack]
);
}
#[test]
fn recv_error_propagates_and_is_not_swallowed() {
let result = run_with(
vec![event_message("ok", None)],
RunOptions::idempotent(),
true,
true,
);
let err = result.outcome.expect_err("recv error should propagate");
assert!(err.is_retryable());
assert!(result.events.is_empty());
}
#[test]
fn settle_error_propagates_and_is_not_swallowed() {
let result = run_with(
vec![event_message("ok", None)],
RunOptions::idempotent(),
false,
false,
);
let err = result.outcome.expect_err("settle error should propagate");
assert!(err.is_retryable());
assert_eq!(
result.events,
vec![Event::Handled("ok".to_string()), Event::Ack]
);
}
#[test]
fn settle_error_on_failure_path_propagates() {
let result = run_with(
vec![event_message("retryable", None)],
RunOptions::idempotent(),
false,
false,
);
let err = result
.outcome
.expect_err("nack settle error should propagate");
assert!(err.is_retryable());
assert_eq!(
result.events.first(),
Some(&Event::Handled("retryable".to_string()))
);
assert!(matches!(result.events.get(1), Some(Event::Nack(_))));
}
#[test]
fn unhandled_message_is_acked_and_ignored() {
let result = run(
vec![event_message("unrelated", None), event_message("ok", None)],
RunOptions::idempotent(),
);
assert!(result.outcome.is_ok());
assert_eq!(
result.events,
vec![Event::Ack, Event::Handled("ok".to_string()), Event::Ack]
);
}
fn run_decode_error<I: Send>(messages: Vec<Message>, options: RunOptions<I>) -> RunResult {
let recorder = Recorder::new();
let svc = router(&recorder);
let source = FakeSource {
queue: messages.into_iter().collect(),
recorder: recorder.clone(),
settle_ok: true,
recv_error: false,
decode_error: true,
};
let outcome = block_on(run_source(svc, source, options));
RunResult {
outcome,
events: recorder.events(),
}
}
#[test]
fn corrupt_row_dead_letters_under_default_policy_not_acked_and_ignored() {
let result = run_decode_error(vec![event_message("", None)], RunOptions::idempotent());
assert!(result.outcome.is_ok());
assert_eq!(result.events.len(), 1);
match &result.events[0] {
Event::DeadLetter(reason) => assert!(reason.contains("corrupt row")),
other => panic!("expected dead-letter, got {other:?}"),
}
assert!(!result.events.iter().any(|e| matches!(e, Event::Handled(_))));
assert!(!result.events.contains(&Event::Ack));
}
#[test]
fn corrupt_row_parks_under_park_policy() {
let result = run_decode_error(
vec![event_message("", None)],
RunOptions::idempotent().with_failure_policy(FailurePolicy::Park),
);
assert!(result.outcome.is_ok());
assert!(matches!(result.events.first(), Some(Event::Park(_))));
}
#[test]
fn corrupt_row_stops_under_stop_policy_with_permanent_error() {
let result = run_decode_error(
vec![event_message("", None)],
RunOptions::idempotent().with_failure_policy(FailurePolicy::Stop),
);
let err = result
.outcome
.expect_err("stop policy surfaces the decode error");
assert!(err.is_permanent());
assert!(
result.events.is_empty(),
"stop does not settle the corrupt row"
);
}
#[test]
fn run_source_future_is_send() {
fn assert_send<T: Send>(_: &T) {}
let recorder = Recorder::new();
let svc = router(&recorder);
let source = FakeSource {
queue: VecDeque::new(),
recorder,
settle_ok: true,
recv_error: false,
decode_error: false,
};
let future = run_source(svc, source, RunOptions::idempotent());
assert_send(&future);
assert!(block_on(future).is_ok());
}
struct DefaultReceived {
message: Message,
recorder: Arc<Recorder>,
}
impl ReceivedMessage for DefaultReceived {
fn message(&self) -> &Message {
&self.message
}
async fn ack(self) -> Result<(), TransportError> {
self.recorder.push(Event::Ack);
Ok(())
}
async fn nack(self, reason: &str) -> Result<(), TransportError> {
self.recorder.push(Event::Nack(reason.to_string()));
Ok(())
}
}
#[test]
fn default_dead_letter_and_park_degrade_to_nack() {
let recorder = Recorder::new();
let dl = DefaultReceived {
message: event_message("ok", None),
recorder: recorder.clone(),
};
block_on(dl.dead_letter("boom")).unwrap();
let park = DefaultReceived {
message: event_message("ok", None),
recorder: recorder.clone(),
};
block_on(park.park("hold")).unwrap();
assert_eq!(
recorder.events(),
vec![
Event::Nack("boom".to_string()),
Event::Nack("hold".to_string())
]
);
}
}