1use std::collections::HashMap;
7use std::future::Future;
8use std::pin::Pin;
9use std::sync::{Arc, Mutex};
10use std::task::{Context, Poll};
11use std::time::Duration;
12
13use async_trait::async_trait;
14use tokio::sync::{AcquireError, OwnedSemaphorePermit, Semaphore, mpsc, oneshot};
15use tokio_util::sync::CancellationToken;
16use tower::Service;
17
18use camel_component_api::UriConfig;
19use camel_component_api::parse_uri;
20use camel_component_api::{BoxProcessor, CamelError, Exchange};
21use camel_component_api::{
22 Component, ComponentMetadata, Consumer, ConsumerContext, ConsumerStartupMode, Endpoint,
23 ProducerContext,
24};
25use tracing::{debug, error, info, warn};
26
27type DirectSender = mpsc::Sender<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>;
35type DirectRegistry = Arc<Mutex<HashMap<String, DirectSender>>>;
36type AcquirePermitFut =
37 Pin<Box<dyn Future<Output = Result<OwnedSemaphorePermit, AcquireError>> + Send>>;
38
39fn validate_name(name: &str) -> Result<(), CamelError> {
45 if name.trim().is_empty() {
46 return Err(CamelError::InvalidUri(
47 "direct: endpoint name must not be empty".to_string(),
48 ));
49 }
50 if name.contains(char::is_whitespace) {
51 return Err(CamelError::InvalidUri(
52 "direct: endpoint name must not contain whitespace".to_string(),
53 ));
54 }
55 Ok(())
56}
57
58#[derive(Debug, Clone, UriConfig)]
68#[uri_scheme = "direct"]
69#[uri_config(
70 skip_impl,
71 metadata(
72 scheme = "direct",
73 description = "Synchronous in-memory direct invocation between routes",
74 producer,
75 consumer
76 ),
77 crate = "camel_component_api"
78)]
79pub struct DirectConfig {
80 pub name: String,
82 #[uri_param(
84 name = "timeout_ms",
85 default = "30000",
86 desc = "Producer call timeout in milliseconds"
87 )]
88 pub timeout_ms: Option<u64>,
89 #[uri_param(
91 name = "failIfNoConsumers",
92 default = "true",
93 desc = "Fail if no consumer registered for the name"
94 )]
95 pub fail_if_no_consumers: Option<bool>,
96}
97
98impl DirectConfig {
99 pub fn from_uri(uri: &str) -> Result<Self, CamelError> {
100 let parts = parse_uri(uri)?;
101 if parts.scheme != "direct" {
102 return Err(CamelError::InvalidUri(format!(
103 "invalid scheme '{}', expected 'direct'",
104 parts.scheme
105 )));
106 }
107
108 let parse_bool = |name: &str, value: &str| -> Result<bool, CamelError> {
109 match value.to_ascii_lowercase().as_str() {
110 "true" | "1" | "yes" => Ok(true),
111 "false" | "0" | "no" => Ok(false),
112 _ => Err(CamelError::InvalidUri(format!(
113 "invalid value for {}: invalid boolean value: '{}'",
114 name, value
115 ))),
116 }
117 };
118
119 let timeout_ms = parts
120 .params
121 .get("timeout_ms")
122 .map(|v| {
123 v.parse::<u64>().map_err(|e| {
124 CamelError::InvalidUri(format!("invalid value for timeout_ms: {}", e))
125 })
126 })
127 .transpose()?;
128
129 if parts.params.contains_key("block") {
130 return Err(CamelError::InvalidUri("block is not supported".into()));
131 }
132
133 let fail_if_no_consumers = parts
134 .params
135 .get("fail_if_no_consumers")
136 .or_else(|| parts.params.get("failIfNoConsumers"))
137 .map(|v| parse_bool("fail_if_no_consumers", v))
138 .transpose()?;
139
140 if parts.params.contains_key("exchange_pattern")
141 || parts.params.contains_key("exchangePattern")
142 {
143 return Err(CamelError::InvalidUri(
144 "exchange_pattern is not supported".into(),
145 ));
146 }
147
148 Ok(Self {
149 name: parts.path,
150 timeout_ms,
151 fail_if_no_consumers,
152 })
153 }
154}
155
156pub struct DirectComponent {
168 registry: DirectRegistry,
169}
170
171impl DirectComponent {
172 pub fn new() -> Self {
173 Self {
174 registry: Arc::new(Mutex::new(HashMap::new())),
175 }
176 }
177}
178
179impl Default for DirectComponent {
180 fn default() -> Self {
181 Self::new()
182 }
183}
184
185impl Component for DirectComponent {
186 fn scheme(&self) -> &str {
187 "direct"
188 }
189
190 fn metadata(&self) -> ComponentMetadata {
191 DirectConfig::metadata()
192 }
193
194 fn create_endpoint(
195 &self,
196 uri: &str,
197 _ctx: &dyn camel_component_api::ComponentContext,
198 ) -> Result<Box<dyn Endpoint>, CamelError> {
199 let config = DirectConfig::from_uri(uri)?;
200 validate_name(&config.name)?;
201 let name = config.name.clone();
202 debug!(endpoint_name = %name, "direct endpoint created");
203 Ok(Box::new(DirectEndpoint {
204 uri: uri.to_string(),
205 config,
206 registry: Arc::clone(&self.registry),
207 }))
208 }
209}
210
211struct DirectEndpoint {
216 uri: String,
217 config: DirectConfig,
218 registry: DirectRegistry,
219}
220
221impl Endpoint for DirectEndpoint {
222 fn uri(&self) -> &str {
223 &self.uri
224 }
225
226 fn create_consumer(
227 &self,
228 rt: Arc<dyn camel_component_api::RuntimeObservability>,
229 ) -> Result<Box<dyn Consumer>, CamelError> {
230 Ok(Box::new(DirectConsumer::new(
231 self.config.name.clone(),
232 Arc::clone(&self.registry),
233 rt,
234 )))
235 }
236
237 fn create_producer(
238 &self,
239 _rt: Arc<dyn camel_component_api::RuntimeObservability>,
240 _ctx: &ProducerContext,
241 ) -> Result<BoxProcessor, CamelError> {
242 Ok(BoxProcessor::new(DirectProducer {
243 name: self.config.name.clone(),
244 registry: Arc::clone(&self.registry),
245 config: self.config.clone(),
246 semaphore: Arc::new(Semaphore::new(1)),
247 pending_permit: None,
248 acquire_fut: None,
249 fail_if_no_consumers: self.config.fail_if_no_consumers,
250 }))
251 }
252}
253
254struct DirectConsumer {
261 name: String,
262 registry: DirectRegistry,
263 cancel: Option<CancellationToken>,
264 runtime: Arc<dyn camel_component_api::RuntimeObservability>,
265}
266
267impl DirectConsumer {
268 fn new(
269 name: String,
270 registry: DirectRegistry,
271 runtime: Arc<dyn camel_component_api::RuntimeObservability>,
272 ) -> Self {
273 Self {
274 name,
275 registry,
276 cancel: None,
277 runtime,
278 }
279 }
280}
281
282#[async_trait]
283impl Consumer for DirectConsumer {
284 fn startup_mode(&self) -> ConsumerStartupMode {
285 ConsumerStartupMode::Explicit
286 }
287
288 async fn start(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
289 let (tx, mut rx) =
291 mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(32);
292
293 {
295 let mut reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
296 if let Some(existing) = reg.get(&self.name)
297 && !existing.is_closed()
298 {
299 return Err(CamelError::EndpointCreationFailed(format!(
300 "direct endpoint '{}' already has a registered consumer",
301 self.name
302 )));
303 }
304 reg.insert(self.name.clone(), tx);
305 }
306
307 context.mark_ready();
308
309 let name = self.name.clone();
310 let registry = Arc::clone(&self.registry);
311 let cancel = context.cancel_token();
312 let cancel_clone = cancel.clone();
313 let route_id = context.route_id().to_owned();
314 let runtime = Arc::clone(&self.runtime);
315
316 info!(endpoint_name = %self.name, "direct consumer started");
317
318 self.cancel = Some(cancel);
319
320 loop {
321 tokio::select! {
322 _ = cancel_clone.cancelled() => {
323 debug!(endpoint_name = %name, "direct consumer received cancellation");
324 break;
325 }
326 msg = rx.recv() => {
327 match msg {
328 Some((exchange, reply_tx)) => {
329 debug!(
330 endpoint_name = %name,
331 exchange_id = %exchange.correlation_id,
332 "direct consumer received exchange"
333 );
334 let result = context.send_and_wait(exchange).await;
335 if let Err(ref err) = result {
336 runtime
341 .metrics()
342 .increment_errors(&route_id, "b-prime:direct:send-and-wait");
343 error!(
345 endpoint_name = %name,
346 error = %err,
347 "direct consumer pipeline error"
348 );
349 }
350 let _ = reply_tx.send(result);
351 }
352 None => break,
353 }
354 }
355 }
356 }
357
358 {
360 let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
361 reg.remove(&name);
362 }
363
364 debug!(endpoint_name = %name, "direct consumer stopped");
365 Ok(())
366 }
367
368 async fn stop(&mut self) -> Result<(), CamelError> {
369 if let Some(cancel) = self.cancel.take() {
371 cancel.cancel();
372 }
373
374 let mut reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
375 reg.remove(&self.name);
376
377 debug!(endpoint_name = %self.name, "direct consumer stopped");
378 Ok(())
379 }
380}
381
382struct DirectProducer {
389 name: String,
390 registry: DirectRegistry,
391 config: DirectConfig,
392 semaphore: Arc<Semaphore>,
393 pending_permit: Option<OwnedSemaphorePermit>,
394 acquire_fut: Option<AcquirePermitFut>,
395 fail_if_no_consumers: Option<bool>,
396}
397
398impl Clone for DirectProducer {
399 fn clone(&self) -> Self {
400 Self {
401 name: self.name.clone(),
402 registry: self.registry.clone(),
403 config: self.config.clone(),
404 semaphore: self.semaphore.clone(),
405 pending_permit: None,
406 acquire_fut: None,
407 fail_if_no_consumers: self.fail_if_no_consumers,
408 }
409 }
410}
411
412impl Service<Exchange> for DirectProducer {
413 type Response = Exchange;
414 type Error = CamelError;
415 type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
416
417 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
418 if self.pending_permit.is_some() {
420 return Poll::Ready(Ok(()));
421 }
422
423 {
425 let reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
426 match reg.get(&self.name) {
427 None => {
428 if self.fail_if_no_consumers != Some(false) {
429 return Poll::Ready(Err(CamelError::EndpointCreationFailed(format!(
430 "direct endpoint '{}' not registered",
431 self.name
432 ))));
433 }
434 }
435 Some(sender) if sender.is_closed() => {
436 return Poll::Ready(Err(CamelError::EndpointCreationFailed(format!(
437 "direct endpoint '{}' channel closed",
438 self.name
439 ))));
440 }
441 Some(_) => {}
442 }
443 }
444
445 let fut = self
447 .acquire_fut
448 .get_or_insert_with(|| Box::pin(Arc::clone(&self.semaphore).acquire_owned()));
449 match fut.as_mut().poll(cx) {
450 Poll::Ready(Ok(permit)) => {
451 self.acquire_fut = None;
452 self.pending_permit = Some(permit);
453 Poll::Ready(Ok(()))
454 }
455 Poll::Ready(Err(_)) => Poll::Ready(Err(CamelError::ChannelClosed)),
456 Poll::Pending => Poll::Pending,
457 }
458 }
459
460 fn call(&mut self, exchange: Exchange) -> Self::Future {
461 let _permit = match self.pending_permit.take() {
462 Some(p) => p,
463 None => {
464 return Box::pin(async {
465 Err(CamelError::ProcessorError(
466 "call() invoked without poll_ready()".into(),
467 ))
468 });
469 }
470 };
471
472 let name = self.name.clone();
473 let registry = Arc::clone(&self.registry);
474 let timeout = Duration::from_millis(self.config.timeout_ms.unwrap_or(30_000));
475 let exchange_id = exchange.correlation_id.clone();
476
477 debug!(
478 endpoint_name = %name,
479 exchange_id = %exchange_id,
480 "direct producer call entry"
481 );
482
483 Box::pin(async move {
484 tokio::time::timeout(timeout, async {
485 let sender = {
486 let reg = registry.lock().unwrap_or_else(|e| e.into_inner());
487 reg.get(&name)
488 .ok_or_else(|| {
489 let err = CamelError::EndpointCreationFailed(format!(
490 "no consumer registered for direct:{name}"
491 ));
492 warn!(endpoint_name = %name, error = %err, "direct send failed");
495 err
496 })?
497 .clone()
498 };
499
500 let (reply_tx, reply_rx) = oneshot::channel();
501 sender.send((exchange, reply_tx)).await.map_err(|err| {
502 warn!(endpoint_name = %name, error = %err, "direct send failed");
505 CamelError::ChannelClosed
506 })?;
507
508 let result = reply_rx.await.map_err(|err| {
509 warn!(endpoint_name = %name, error = %err, "direct send failed");
512 CamelError::ChannelClosed
513 })?;
514
515 debug!(endpoint_name = %name, "direct message sent");
516 result
517 })
518 .await
519 .map_err(|_| CamelError::ProcessorError(format!("direct:{name} call timed out")))?
520 })
521 }
522}
523
524#[cfg(test)]
529mod tests {
530 use std::sync::Mutex;
531
532 use camel_api::MetricsCollector;
533 use camel_component_api::HealthCheckRegistry;
534 use std::time::Duration;
535 fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
536 std::sync::Arc::new(NoOpComponentContext)
540 }
541
542 use super::*;
543 use camel_component_api::ExchangeEnvelope;
544 use camel_component_api::Message;
545 use camel_component_api::NoOpComponentContext;
546 use camel_component_api::RuntimeObservability;
547 use camel_component_api::StartupSignal;
548 use std::task::RawWakerVTable;
549 use tower::ServiceExt;
550
551 struct RecordingMetrics {
556 errors: Arc<Mutex<Vec<(String, String)>>>,
557 }
558
559 impl MetricsCollector for RecordingMetrics {
560 fn record_exchange_duration(&self, _: &str, _: Duration) {}
561 fn increment_errors(&self, route_id: &str, error_type: &str) {
562 self.errors
563 .lock()
564 .unwrap()
565 .push((route_id.to_string(), error_type.to_string()));
566 }
567 fn increment_exchanges(&self, _: &str) {}
568 fn set_queue_depth(&self, _: &str, _: usize) {}
569 fn record_circuit_breaker_change(&self, _: &str, _: &str, _: &str) {}
570 }
571
572 struct RecordingRuntime {
573 metrics_collector: Arc<RecordingMetrics>,
574 }
575
576 impl RecordingRuntime {
577 fn new(errors: Arc<Mutex<Vec<(String, String)>>>) -> Self {
578 Self {
579 metrics_collector: Arc::new(RecordingMetrics { errors }),
580 }
581 }
582 }
583
584 impl RuntimeObservability for RecordingRuntime {
585 fn metrics(&self) -> Arc<dyn MetricsCollector> {
586 self.metrics_collector.clone() as Arc<dyn MetricsCollector>
587 }
588 fn health(&self) -> Arc<dyn HealthCheckRegistry> {
589 panic!("RecordingRuntime::health not used in this test")
590 }
591 }
592
593 fn noop_waker() -> std::task::Waker {
594 const VTABLE: RawWakerVTable = RawWakerVTable::new(|_| RAW, |_| {}, |_| {}, |_| {});
595 const RAW: std::task::RawWaker = std::task::RawWaker::new(std::ptr::null(), &VTABLE);
596 unsafe { std::task::Waker::from_raw(RAW) }
597 }
598
599 fn test_producer_ctx() -> ProducerContext {
600 ProducerContext::new()
601 }
602
603 #[test]
604 fn test_direct_component_scheme() {
605 let component = DirectComponent::new();
606 assert_eq!(component.scheme(), "direct");
607 }
608
609 #[test]
610 fn test_direct_component_default() {
611 let component = DirectComponent::default();
612 assert_eq!(component.scheme(), "direct");
613 }
614
615 #[test]
616 fn test_direct_config_from_uri() {
617 let config = DirectConfig::from_uri("direct:orders").unwrap();
618 assert_eq!(config.name, "orders");
619 }
620
621 #[test]
622 fn rejects_block_param() {
623 let result = DirectConfig::from_uri("direct:foo?block=true");
624 assert!(result.is_err());
625 assert!(
626 result
627 .unwrap_err()
628 .to_string()
629 .contains("block is not supported"),
630 "expected block rejection"
631 );
632 }
633
634 #[test]
635 fn rejects_exchange_pattern_snake_case() {
636 let result = DirectConfig::from_uri("direct:foo?exchange_pattern=InOnly");
637 assert!(result.is_err());
638 assert!(
639 result
640 .unwrap_err()
641 .to_string()
642 .contains("exchange_pattern is not supported"),
643 "expected exchange_pattern rejection"
644 );
645 }
646
647 #[test]
648 fn rejects_exchange_pattern_camel_case() {
649 let result = DirectConfig::from_uri("direct:foo?exchangePattern=InOnly");
650 assert!(result.is_err());
651 assert!(
652 result
653 .unwrap_err()
654 .to_string()
655 .contains("exchange_pattern is not supported"),
656 "expected exchangePattern rejection"
657 );
658 }
659
660 #[test]
661 fn test_direct_endpoint_uri() {
662 let component = DirectComponent::new();
663 let endpoint = component
664 .create_endpoint("direct:uri-check", &NoOpComponentContext)
665 .unwrap();
666 assert_eq!(endpoint.uri(), "direct:uri-check");
667 }
668
669 #[test]
670 fn test_direct_creates_endpoint() {
671 let component = DirectComponent::new();
672 let endpoint = component.create_endpoint("direct:foo", &NoOpComponentContext);
673 assert!(endpoint.is_ok());
674 }
675
676 #[test]
677 fn test_direct_wrong_scheme() {
678 let component = DirectComponent::new();
679 let result = component.create_endpoint("timer:tick", &NoOpComponentContext);
680 assert!(result.is_err());
681 }
682
683 #[test]
684 fn test_direct_endpoint_creates_consumer() {
685 let component = DirectComponent::new();
686 let endpoint = component
687 .create_endpoint("direct:foo", &NoOpComponentContext)
688 .unwrap();
689 assert!(endpoint.create_consumer(rt()).is_ok());
690 }
691
692 #[test]
693 fn test_direct_endpoint_creates_producer() {
694 let ctx = test_producer_ctx();
695 let component = DirectComponent::new();
696 let endpoint = component
697 .create_endpoint("direct:foo", &NoOpComponentContext)
698 .unwrap();
699 assert!(endpoint.create_producer(rt(), &ctx).is_ok());
700 }
701
702 #[test]
703 fn test_direct_empty_name_rejected() {
704 let component = DirectComponent::new();
705 match component.create_endpoint("direct:", &NoOpComponentContext) {
706 Err(e) => assert!(
707 e.to_string().contains("must not be empty"),
708 "unexpected error: {e}"
709 ),
710 Ok(_) => panic!("expected error for empty name"),
711 }
712 }
713
714 #[tokio::test]
715 async fn test_direct_producer_no_consumer_registered() {
716 let ctx = test_producer_ctx();
717 let component = DirectComponent::new();
718 let endpoint = component
719 .create_endpoint("direct:missing", &NoOpComponentContext)
720 .unwrap();
721 let producer = endpoint.create_producer(rt(), &ctx).unwrap();
722
723 let exchange = Exchange::new(Message::new("test"));
724 let result = producer.oneshot(exchange).await;
725 assert!(result.is_err());
726 }
727
728 #[tokio::test]
729 async fn test_direct_duplicate_consumer_returns_error() {
730 let component = DirectComponent::new();
731 let endpoint = component
732 .create_endpoint("direct:dup", &NoOpComponentContext)
733 .unwrap();
734
735 let mut consumer_a = endpoint.create_consumer(rt()).unwrap();
736 let mut consumer_b = endpoint.create_consumer(rt()).unwrap();
737
738 let (route_tx_a, _route_rx_a) = mpsc::channel::<ExchangeEnvelope>(16);
739 let token_a = tokio_util::sync::CancellationToken::new();
740 let ctx_a = ConsumerContext::new(
741 route_tx_a,
742 token_a.clone(),
743 "direct-test-route-a".to_string(),
744 );
745 let handle_a = tokio::spawn(async move {
746 consumer_a.start(ctx_a).await.unwrap();
747 });
748 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
749
750 let (route_tx_b, _route_rx_b) = mpsc::channel::<ExchangeEnvelope>(16);
751 let ctx_b = ConsumerContext::new(
752 route_tx_b,
753 tokio_util::sync::CancellationToken::new(),
754 "direct-test-route-b".to_string(),
755 );
756 let result = consumer_b.start(ctx_b).await;
757
758 assert!(matches!(
759 result,
760 Err(CamelError::EndpointCreationFailed(msg))
761 if msg.contains("already has a registered consumer")
762 ));
763
764 token_a.cancel();
765 handle_a.await.unwrap();
766 }
767
768 #[tokio::test]
769 async fn test_direct_producer_consumer_roundtrip() {
770 let component = DirectComponent::new();
771
772 let consumer_endpoint = component
774 .create_endpoint("direct:test", &NoOpComponentContext)
775 .unwrap();
776 let mut consumer = consumer_endpoint.create_consumer(rt()).unwrap();
777
778 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
780 let ctx = ConsumerContext::new(
781 route_tx,
782 tokio_util::sync::CancellationToken::new(),
783 "direct-test-route".to_string(),
784 );
785
786 tokio::spawn(async move {
788 consumer.start(ctx).await.unwrap();
789 });
790
791 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
793
794 tokio::spawn(async move {
796 while let Some(envelope) = route_rx.recv().await {
797 let ExchangeEnvelope { exchange, reply_tx } = envelope;
798 if let Some(tx) = reply_tx {
799 let _ = tx.send(Ok(exchange));
800 }
801 }
802 });
803
804 let ctx = test_producer_ctx();
806 let producer_endpoint = component
807 .create_endpoint("direct:test", &NoOpComponentContext)
808 .unwrap();
809 let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
810
811 let exchange = Exchange::new(Message::new("hello direct"));
812 let result = producer.oneshot(exchange).await;
813
814 assert!(result.is_ok());
815 let reply = result.unwrap();
816 assert_eq!(reply.input.body.as_text(), Some("hello direct"));
817 }
818
819 #[tokio::test]
820 async fn test_direct_propagates_error_when_no_handler() {
821 let component = DirectComponent::new();
822
823 let consumer_endpoint = component
824 .create_endpoint("direct:err-test", &NoOpComponentContext)
825 .unwrap();
826 let mut consumer = consumer_endpoint.create_consumer(rt()).unwrap();
827
828 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
829 let ctx = ConsumerContext::new(
830 route_tx,
831 tokio_util::sync::CancellationToken::new(),
832 "direct-test-route".to_string(),
833 );
834
835 tokio::spawn(async move {
836 consumer.start(ctx).await.unwrap();
837 });
838
839 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
840
841 tokio::spawn(async move {
843 while let Some(envelope) = route_rx.recv().await {
844 if let Some(tx) = envelope.reply_tx {
845 let _ = tx.send(Err(CamelError::ProcessorError("subroute failed".into())));
846 }
847 }
848 });
849
850 let ctx = test_producer_ctx();
851 let producer_endpoint = component
852 .create_endpoint("direct:err-test", &NoOpComponentContext)
853 .unwrap();
854 let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
855
856 let exchange = Exchange::new(Message::new("test"));
857 let result = producer.oneshot(exchange).await;
858 assert!(result.is_err());
859 assert!(matches!(result.unwrap_err(), CamelError::ProcessorError(_)));
860 }
861
862 #[tokio::test]
863 async fn test_direct_consumer_stop_unregisters() {
864 let component = DirectComponent::new();
865 let endpoint = component
866 .create_endpoint("direct:cleanup", &NoOpComponentContext)
867 .unwrap();
868
869 let mut consumer = endpoint.create_consumer(rt()).unwrap();
871
872 let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
873 let ctx = ConsumerContext::new(
874 route_tx,
875 tokio_util::sync::CancellationToken::new(),
876 "direct-test-route".to_string(),
877 );
878
879 let handle = tokio::spawn(async move {
881 consumer.start(ctx).await.unwrap();
882 });
883
884 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
885
886 {
888 let reg = component.registry.lock().unwrap_or_else(|e| e.into_inner());
889 assert!(reg.contains_key("cleanup"));
890 }
891
892 let mut stop_consumer = DirectConsumer {
894 name: "cleanup".to_string(),
895 registry: Arc::clone(&component.registry),
896 cancel: None,
897 runtime: rt(),
898 };
899 stop_consumer.stop().await.unwrap();
900
901 {
903 let reg = component.registry.lock().unwrap_or_else(|e| e.into_inner());
904 assert!(!reg.contains_key("cleanup"));
905 }
906
907 handle.abort();
908 }
909
910 #[tokio::test]
911 async fn test_direct_consumer_respects_cancellation() {
912 use tokio_util::sync::CancellationToken;
913
914 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
915 let token = CancellationToken::new();
916 let (tx, _rx) = mpsc::channel(16);
917 let ctx = ConsumerContext::new(tx, token.clone(), "direct-test-route".to_string());
918
919 let mut consumer = DirectConsumer {
920 name: "cancel-test".to_string(),
921 registry: registry.clone(),
922 cancel: None,
923 runtime: rt(),
924 };
925
926 let handle = tokio::spawn(async move {
927 consumer.start(ctx).await.unwrap();
928 });
929
930 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
931 assert!(
932 registry
933 .lock()
934 .unwrap_or_else(|e| e.into_inner())
935 .contains_key("cancel-test")
936 );
937
938 token.cancel();
939 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
943
944 assert!(
946 !registry
947 .lock()
948 .unwrap_or_else(|e| e.into_inner())
949 .contains_key("cancel-test")
950 );
951
952 let _ = handle.await;
953 }
954
955 #[tokio::test]
956 async fn test_direct_consumer_stop_missing_entry_is_ok() {
957 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
958 let mut consumer = DirectConsumer {
959 name: "never-registered".to_string(),
960 registry,
961 cancel: None,
962 runtime: rt(),
963 };
964 let result = consumer.stop().await;
965 assert!(result.is_ok());
966 }
967
968 #[test]
969 fn test_poll_ready_endpoint_not_registered() {
970 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
971 let producer = DirectProducer {
972 name: "missing".to_string(),
973 registry,
974 config: DirectConfig {
975 name: "missing".to_string(),
976 timeout_ms: None,
977 fail_if_no_consumers: None,
978 },
979 semaphore: Arc::new(Semaphore::new(1)),
980 pending_permit: None,
981 acquire_fut: None,
982 fail_if_no_consumers: None,
983 };
984 let waker = noop_waker();
985 let mut cx = Context::from_waker(&waker);
986 let mut producer = producer;
987 let result = producer.poll_ready(&mut cx);
988 assert!(matches!(
989 result,
990 Poll::Ready(Err(CamelError::EndpointCreationFailed(_)))
991 ));
992 }
993
994 #[test]
995 fn test_poll_ready_endpoint_registered() {
996 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
997 let (tx, _rx) =
998 mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(1);
999 registry.lock().unwrap().insert("active".to_string(), tx);
1000 let producer = DirectProducer {
1001 name: "active".to_string(),
1002 registry,
1003 config: DirectConfig {
1004 name: "active".to_string(),
1005 timeout_ms: None,
1006 fail_if_no_consumers: None,
1007 },
1008 semaphore: Arc::new(Semaphore::new(1)),
1009 pending_permit: None,
1010 acquire_fut: None,
1011 fail_if_no_consumers: None,
1012 };
1013 let waker = noop_waker();
1014 let mut cx = Context::from_waker(&waker);
1015 let mut producer = producer;
1016 let result = producer.poll_ready(&mut cx);
1017 assert!(matches!(result, Poll::Ready(Ok(()))));
1018 }
1019
1020 #[test]
1021 fn test_poll_ready_allows_missing_consumer_when_fail_if_no_consumers_false() {
1022 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
1023 let producer = DirectProducer {
1024 name: "missing-ok".to_string(),
1025 registry,
1026 config: DirectConfig {
1027 name: "missing-ok".to_string(),
1028 timeout_ms: None,
1029 fail_if_no_consumers: Some(false),
1030 },
1031 semaphore: Arc::new(Semaphore::new(1)),
1032 pending_permit: None,
1033 acquire_fut: None,
1034 fail_if_no_consumers: Some(false),
1035 };
1036
1037 let waker = noop_waker();
1038 let mut cx = Context::from_waker(&waker);
1039 let mut producer = producer;
1040 let result = producer.poll_ready(&mut cx);
1041 assert!(matches!(result, Poll::Ready(Ok(()))));
1042 }
1043
1044 #[test]
1045 fn test_poll_ready_channel_closed() {
1046 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
1047 let (tx, rx) =
1048 mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(1);
1049 drop(rx);
1050 registry.lock().unwrap().insert("closed".to_string(), tx);
1051 let producer = DirectProducer {
1052 name: "closed".to_string(),
1053 registry,
1054 config: DirectConfig {
1055 name: "closed".to_string(),
1056 timeout_ms: None,
1057 fail_if_no_consumers: None,
1058 },
1059 semaphore: Arc::new(Semaphore::new(1)),
1060 pending_permit: None,
1061 acquire_fut: None,
1062 fail_if_no_consumers: None,
1063 };
1064 let waker = noop_waker();
1065 let mut cx = Context::from_waker(&waker);
1066 let mut producer = producer;
1067 let result = producer.poll_ready(&mut cx);
1068 assert!(matches!(
1069 result,
1070 Poll::Ready(Err(CamelError::EndpointCreationFailed(_)))
1071 ));
1072 }
1073
1074 #[tokio::test]
1075 async fn test_direct_stop_cancels_loop() {
1076 use tokio_util::sync::CancellationToken;
1077
1078 let component = DirectComponent::new();
1079 let endpoint = component
1080 .create_endpoint("direct:stop-test", &NoOpComponentContext)
1081 .unwrap();
1082 let mut consumer = endpoint.create_consumer(rt()).unwrap();
1083
1084 let token = CancellationToken::new();
1085 let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1086 let ctx = ConsumerContext::new(route_tx, token.clone(), "direct-test-route".to_string());
1087
1088 let handle = tokio::spawn(async move {
1091 consumer.start(ctx).await.unwrap();
1092 });
1093
1094 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1095 assert!(
1096 component
1097 .registry
1098 .lock()
1099 .unwrap_or_else(|e| e.into_inner())
1100 .contains_key("stop-test")
1101 );
1102
1103 let mut stop_consumer = DirectConsumer {
1105 name: "stop-test".to_string(),
1106 registry: Arc::clone(&component.registry),
1107 cancel: Some(token.clone()),
1108 runtime: rt(),
1109 };
1110 stop_consumer.stop().await.unwrap();
1111
1112 let result = tokio::time::timeout(std::time::Duration::from_secs(2), handle).await;
1114 assert!(result.is_ok(), "Consumer loop did not stop within 2s");
1115
1116 assert!(
1118 !component
1119 .registry
1120 .lock()
1121 .unwrap_or_else(|e| e.into_inner())
1122 .contains_key("stop-test")
1123 );
1124 }
1125
1126 #[tokio::test]
1127 async fn test_direct_producer_timeout() {
1128 let component = DirectComponent::new();
1129 let endpoint = component
1130 .create_endpoint("direct:timeout-test", &NoOpComponentContext)
1131 .unwrap();
1132 let mut consumer = endpoint.create_consumer(rt()).unwrap();
1133
1134 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1136 let token = tokio_util::sync::CancellationToken::new();
1137 let ctx = ConsumerContext::new(route_tx, token.clone(), "direct-test-route".to_string());
1138 tokio::spawn(async move {
1139 consumer.start(ctx).await.unwrap();
1140 });
1141
1142 tokio::spawn(async move {
1144 let mut held_reply: Vec<oneshot::Sender<Result<Exchange, CamelError>>> = Vec::new();
1145 while let Some(envelope) = route_rx.recv().await {
1146 held_reply.push(envelope.reply_tx.unwrap());
1147 }
1148 drop(held_reply);
1149 });
1150
1151 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1152
1153 let _ = test_producer_ctx();
1155 let _producer_endpoint = component
1156 .create_endpoint("direct:timeout-test", &NoOpComponentContext)
1157 .unwrap();
1158 let producer = DirectProducer {
1159 name: "timeout-test".to_string(),
1160 registry: Arc::clone(&component.registry),
1161 config: DirectConfig {
1162 name: "timeout-test".to_string(),
1163 timeout_ms: Some(100), fail_if_no_consumers: None,
1165 },
1166 semaphore: Arc::new(Semaphore::new(1)),
1167 pending_permit: None,
1168 acquire_fut: None,
1169 fail_if_no_consumers: None,
1170 };
1171
1172 let exchange = Exchange::new(Message::new("test"));
1173 let mut svc = producer;
1174 let _ = svc.poll_ready(&mut Context::from_waker(&noop_waker()));
1175 let result = svc.call(exchange).await;
1176 assert!(result.is_err(), "Expected timeout error");
1177 assert!(
1178 result.unwrap_err().to_string().contains("timed out"),
1179 "Expected timeout message"
1180 );
1181
1182 token.cancel();
1183 }
1184
1185 #[tokio::test]
1186 async fn test_send_and_wait_error_increments_errors_metric() {
1187 let errors: Arc<Mutex<Vec<(String, String)>>> = Arc::new(Mutex::new(Vec::new()));
1188 let runtime = Arc::new(RecordingRuntime::new(Arc::clone(&errors)));
1189
1190 let component = DirectComponent::new();
1191 let endpoint = component
1192 .create_endpoint("direct:metrics-error", &NoOpComponentContext)
1193 .unwrap();
1194 let mut consumer = endpoint.create_consumer(runtime).unwrap();
1195
1196 let cancel = tokio_util::sync::CancellationToken::new();
1198 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1199 let ctx = ConsumerContext::new(route_tx, cancel.clone(), "test-route-id".to_string());
1200
1201 tokio::spawn(async move {
1202 consumer.start(ctx).await.unwrap();
1203 });
1204
1205 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1206
1207 tokio::spawn(async move {
1209 while let Some(envelope) = route_rx.recv().await {
1210 if let Some(tx) = envelope.reply_tx {
1211 let _ = tx.send(Err(CamelError::ProcessorError(
1212 "pipeline failure".to_string(),
1213 )));
1214 }
1215 }
1216 });
1217
1218 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1219
1220 let ctx = test_producer_ctx();
1223 let producer_endpoint = component
1224 .create_endpoint("direct:metrics-error", &NoOpComponentContext)
1225 .unwrap();
1226 let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
1227
1228 let exchange = Exchange::new(Message::new("test"));
1229 let _result = producer.oneshot(exchange).await;
1232
1233 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1234
1235 cancel.cancel();
1236
1237 let recorded = errors.lock().unwrap();
1239 assert_eq!(
1240 recorded.len(),
1241 1,
1242 "expected 1 increment_errors call, got {}: {:?}",
1243 recorded.len(),
1244 *recorded
1245 );
1246 assert_eq!(recorded[0].0, "test-route-id");
1247 assert_eq!(recorded[0].1, "b-prime:direct:send-and-wait");
1248 }
1249
1250 #[test]
1251 fn test_empty_endpoint_name_rejected() {
1252 let result = DirectConfig::from_uri("direct:");
1253 if let Ok(config) = result {
1255 assert!(
1256 validate_name(&config.name).is_err(),
1257 "expected validation error for empty name"
1258 );
1259 }
1260 let component = DirectComponent::new();
1262 let result = component.create_endpoint("direct:", &NoOpComponentContext);
1263 assert!(result.is_err(), "empty endpoint name must be rejected");
1264 }
1265
1266 #[test]
1267 fn test_whitespace_endpoint_name_rejected() {
1268 let result = DirectConfig::from_uri("direct:my endpoint");
1269 if let Ok(config) = result {
1270 assert!(
1271 validate_name(&config.name).is_err(),
1272 "expected validation error for whitespace in name"
1273 );
1274 }
1275 let component = DirectComponent::new();
1276 let result = component.create_endpoint("direct:my endpoint", &NoOpComponentContext);
1277 assert!(result.is_err(), "whitespace endpoint name must be rejected");
1278 }
1279
1280 #[test]
1281 fn test_valid_endpoint_name_accepted() {
1282 let component = DirectComponent::new();
1283 let result = component.create_endpoint("direct:my-endpoint", &NoOpComponentContext);
1284 assert!(result.is_ok(), "valid endpoint name should be accepted");
1285 }
1286
1287 #[test]
1288 fn test_direct_consumer_startup_mode_is_explicit() {
1289 let component = DirectComponent::new();
1290 let endpoint = component
1291 .create_endpoint("direct:ready-check", &NoOpComponentContext)
1292 .unwrap();
1293 let consumer = endpoint.create_consumer(rt()).unwrap();
1294 assert_eq!(
1295 consumer.startup_mode(),
1296 ConsumerStartupMode::Explicit,
1297 "DirectConsumer must opt into Explicit startup"
1298 );
1299 }
1300
1301 #[tokio::test]
1302 async fn test_direct_consumer_marks_ready_after_registration() {
1303 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
1304
1305 let mut consumer = DirectConsumer {
1306 name: "ready-probe-direct".into(),
1307 registry: registry.clone(),
1308 cancel: None,
1309 runtime: rt(),
1310 };
1311
1312 let (tx, _rx) = mpsc::channel::<ExchangeEnvelope>(16);
1313 let token = CancellationToken::new();
1314 let ctx = ConsumerContext::new(tx, token.clone(), "ready-probe-route".to_string());
1315
1316 let (signal, startup_rx) = StartupSignal::pair();
1317 let ctx = ctx.with_startup(signal);
1318
1319 tokio::spawn(async move {
1320 let _ = consumer.start(ctx).await;
1321 });
1322
1323 let result = tokio::time::timeout(Duration::from_secs(2), startup_rx.await_ready())
1324 .await
1325 .expect("DirectConsumer must call ctx.mark_ready() after registration");
1326 assert!(
1327 result.is_ok(),
1328 "mark_ready must resolve Ok after registration"
1329 );
1330
1331 {
1332 let reg = registry.lock().unwrap_or_else(|e| e.into_inner());
1333 assert!(
1334 reg.contains_key("ready-probe-direct"),
1335 "registry must contain consumer name after mark_ready resolves"
1336 );
1337 }
1338
1339 token.cancel();
1340 }
1341}