use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use async_trait::async_trait;
use tokio::sync::Notify;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
struct ImmediateLifetimeConsumer {
started: Arc<AtomicBool>,
}
#[async_trait]
impl Consumer for ImmediateLifetimeConsumer {
async fn start(&mut self, ctx: ConsumerContext) -> Result<(), CamelError> {
self.started.store(true, Ordering::SeqCst);
ctx.cancelled().await;
Ok(())
}
async fn stop(&mut self) -> Result<(), CamelError> {
Ok(())
}
}
struct ExplicitReadyConsumer {
ready_after: Duration,
started: Arc<AtomicBool>,
entered: Arc<Notify>,
mark_ready_called: Arc<AtomicBool>,
}
#[async_trait]
impl Consumer for ExplicitReadyConsumer {
async fn start(&mut self, ctx: ConsumerContext) -> Result<(), CamelError> {
self.started.store(true, Ordering::SeqCst);
self.entered.notify_one();
tokio::time::sleep(self.ready_after).await;
ctx.mark_ready();
self.mark_ready_called.store(true, Ordering::SeqCst);
ctx.cancelled().await;
Ok(())
}
async fn stop(&mut self) -> Result<(), CamelError> {
Ok(())
}
fn startup_mode(&self) -> ConsumerStartupMode {
ConsumerStartupMode::Explicit
}
}
struct ExplicitBindFailConsumer;
#[async_trait]
impl Consumer for ExplicitBindFailConsumer {
async fn start(&mut self, _ctx: ConsumerContext) -> Result<(), CamelError> {
Err(CamelError::RouteError(
"simulated bind failure: Address already in use".to_string(),
))
}
async fn stop(&mut self) -> Result<(), CamelError> {
Ok(())
}
fn startup_mode(&self) -> ConsumerStartupMode {
ConsumerStartupMode::Explicit
}
}
struct ExplicitOkNoMarkReadyConsumer;
#[async_trait]
impl Consumer for ExplicitOkNoMarkReadyConsumer {
async fn start(&mut self, _ctx: ConsumerContext) -> Result<(), CamelError> {
Ok(())
}
async fn stop(&mut self) -> Result<(), CamelError> {
Ok(())
}
fn startup_mode(&self) -> ConsumerStartupMode {
ConsumerStartupMode::Explicit
}
}
#[tokio::test]
async fn spawn_consumer_task_immediate_consumer_returns_resolved_receiver() {
let (tx, _rx) = mpsc::channel(1);
let cancel = CancellationToken::new();
let ctx = ConsumerContext::new(tx, cancel.clone(), "immediate-route".to_string());
let started = Arc::new(AtomicBool::new(false));
let consumer = ImmediateLifetimeConsumer {
started: Arc::clone(&started),
};
let (handle, startup_rx) = spawn_consumer_task(
"immediate-route".to_string(),
Box::new(consumer),
ctx,
None,
None,
false,
);
let result = tokio::time::timeout(Duration::from_millis(200), startup_rx.await_ready())
.await
.expect("immediate startup receiver must not block");
assert!(result.is_ok(), "immediate receiver resolves Ok");
for _ in 0..50 {
if started.load(Ordering::SeqCst) {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert!(
started.load(Ordering::SeqCst),
"immediate consumer should have entered start()"
);
cancel.cancel();
let _ = tokio::time::timeout(Duration::from_secs(2), handle).await;
}
#[tokio::test]
async fn spawn_consumer_task_explicit_consumer_waits_for_mark_ready() {
let (tx, _rx) = mpsc::channel(1);
let cancel = CancellationToken::new();
let ctx = ConsumerContext::new(tx, cancel.clone(), "explicit-route".to_string());
let started = Arc::new(AtomicBool::new(false));
let mark_ready_called = Arc::new(AtomicBool::new(false));
let entered = Arc::new(Notify::new());
let consumer = ExplicitReadyConsumer {
ready_after: Duration::from_millis(40),
started: Arc::clone(&started),
entered: Arc::clone(&entered),
mark_ready_called: Arc::clone(&mark_ready_called),
};
let (handle, startup_rx) = spawn_consumer_task(
"explicit-route".to_string(),
Box::new(consumer),
ctx,
None,
None,
false,
);
entered.notified().await;
assert!(
started.load(Ordering::SeqCst),
"consumer should have entered start() before ready signal"
);
assert!(
!mark_ready_called.load(Ordering::SeqCst),
"mark_ready must not have fired yet — receiver should still be pending"
);
let result = tokio::time::timeout(Duration::from_secs(2), startup_rx.await_ready())
.await
.expect("explicit receiver must resolve after mark_ready");
assert!(
result.is_ok(),
"explicit receiver resolves Ok after mark_ready"
);
assert!(
mark_ready_called.load(Ordering::SeqCst),
"mark_ready must have been called before receiver resolved"
);
cancel.cancel();
let _ = tokio::time::timeout(Duration::from_secs(2), handle).await;
}
#[tokio::test]
async fn spawn_consumer_task_explicit_consumer_start_error_propagates() {
let (tx, _rx) = mpsc::channel(1);
let cancel = CancellationToken::new();
let ctx = ConsumerContext::new(tx, cancel, "explicit-fail-route".to_string());
let (handle, startup_rx) = spawn_consumer_task(
"explicit-fail-route".to_string(),
Box::new(ExplicitBindFailConsumer),
ctx,
None,
None,
false,
);
let result = tokio::time::timeout(Duration::from_secs(2), startup_rx.await_ready())
.await
.expect("explicit receiver must resolve (with Err) on start failure");
let err = result.expect_err("start() Err must propagate to receiver");
match err {
CamelError::RouteError(msg) => assert!(
msg.contains("simulated bind failure"),
"error message must carry the consumer's failure: got {msg}"
),
other => panic!("expected RouteError, got {other:?}"),
}
handle
.await
.expect("consumer task must join after start failure");
}
#[tokio::test]
async fn spawn_consumer_task_explicit_consumer_ok_without_mark_ready_does_not_hang_controller() {
let (tx, _rx) = mpsc::channel(1);
let cancel = CancellationToken::new();
let ctx = ConsumerContext::new(tx, cancel.clone(), "explicit-no-mark-ready".to_string());
let (handle, startup_rx) = spawn_consumer_task(
"explicit-no-mark-ready".to_string(),
Box::new(ExplicitOkNoMarkReadyConsumer),
ctx,
None,
None,
false,
);
let result = tokio::time::timeout(Duration::from_secs(2), startup_rx.await_ready())
.await
.expect("defensive fallback must resolve receiver (no controller hang)");
assert!(
result.is_ok(),
"defensive fallback must resolve receiver as Ok"
);
cancel.cancel();
let _ = tokio::time::timeout(Duration::from_secs(2), handle).await;
}
struct ExplicitDeferredMarkReadyConsumer {
bg_handle: Option<tokio::task::JoinHandle<Result<(), CamelError>>>,
}
#[async_trait]
impl Consumer for ExplicitDeferredMarkReadyConsumer {
async fn start(&mut self, ctx: ConsumerContext) -> Result<(), CamelError> {
let signal = ctx.startup_signal();
let cancel = ctx.cancel_token();
self.bg_handle = Some(tokio::spawn(async move {
tokio::select! {
_ = cancel.cancelled() => return Ok(()),
_ = tokio::time::sleep(Duration::from_millis(50)) => {}
}
signal.mark_ready();
cancel.cancelled().await;
Ok(())
}));
Ok(())
}
async fn stop(&mut self) -> Result<(), CamelError> {
Ok(())
}
fn startup_mode(&self) -> ConsumerStartupMode {
ConsumerStartupMode::Explicit
}
fn background_task_handle(
&mut self,
) -> Option<tokio::task::JoinHandle<Result<(), CamelError>>> {
self.bg_handle.take()
}
}
#[tokio::test]
async fn spawn_consumer_task_explicit_deferred_mark_ready_not_defeated_by_fallback() {
let (tx, _rx) = mpsc::channel(1);
let cancel = CancellationToken::new();
let ctx = ConsumerContext::new(tx, cancel.clone(), "deferred-ready".to_string());
let consumer = ExplicitDeferredMarkReadyConsumer { bg_handle: None };
let (handle, startup_rx) = spawn_consumer_task(
"deferred-ready".to_string(),
Box::new(consumer),
ctx,
None,
None,
false,
);
let start = std::time::Instant::now();
let result = tokio::time::timeout(Duration::from_secs(2), startup_rx.await_ready())
.await
.expect("bg task must call mark_ready");
let elapsed = start.elapsed();
assert!(result.is_ok(), "deferred mark_ready must resolve Ok");
assert!(
elapsed >= Duration::from_millis(20),
"receiver resolved in {elapsed:?} — defensive fallback likely fired prematurely"
);
cancel.cancel();
let _ = tokio::time::timeout(Duration::from_secs(2), handle).await;
}