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