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::task::JoinHandle;
16use tokio_util::sync::CancellationToken;
17use tower::Service;
18
19use camel_api::component_metadata::{
20 ComponentCapabilities, ComponentMetadata, OptionKind, UriOption,
21};
22use camel_component_api::parse_uri;
23use camel_component_api::{BoxProcessor, CamelError, Exchange};
24use camel_component_api::{Component, Consumer, ConsumerContext, Endpoint, ProducerContext};
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)]
68pub struct DirectConfig {
69 pub name: String,
71 pub timeout_ms: Option<u64>,
73 pub block: Option<bool>,
76 pub fail_if_no_consumers: Option<bool>,
78 pub exchange_pattern: Option<String>,
80}
81
82impl DirectConfig {
83 pub fn from_uri(uri: &str) -> Result<Self, CamelError> {
84 let parts = parse_uri(uri)?;
85 if parts.scheme != "direct" {
86 return Err(CamelError::InvalidUri(format!(
87 "invalid scheme '{}', expected 'direct'",
88 parts.scheme
89 )));
90 }
91
92 let parse_bool = |name: &str, value: &str| -> Result<bool, CamelError> {
93 match value.to_ascii_lowercase().as_str() {
94 "true" | "1" | "yes" => Ok(true),
95 "false" | "0" | "no" => Ok(false),
96 _ => Err(CamelError::InvalidUri(format!(
97 "invalid value for {}: invalid boolean value: '{}'",
98 name, value
99 ))),
100 }
101 };
102
103 let timeout_ms = parts
104 .params
105 .get("timeout_ms")
106 .map(|v| {
107 v.parse::<u64>().map_err(|e| {
108 CamelError::InvalidUri(format!("invalid value for timeout_ms: {}", e))
109 })
110 })
111 .transpose()?;
112
113 let block = parts
114 .params
115 .get("block")
116 .map(|v| parse_bool("block", v))
117 .transpose()?;
118
119 let fail_if_no_consumers = parts
120 .params
121 .get("fail_if_no_consumers")
122 .or_else(|| parts.params.get("failIfNoConsumers"))
123 .map(|v| parse_bool("fail_if_no_consumers", v))
124 .transpose()?;
125
126 let exchange_pattern = parts
127 .params
128 .get("exchange_pattern")
129 .or_else(|| parts.params.get("exchangePattern"))
130 .cloned();
131
132 Ok(Self {
133 name: parts.path,
134 timeout_ms,
135 block,
136 fail_if_no_consumers,
137 exchange_pattern,
138 })
139 }
140}
141
142pub struct DirectComponent {
154 registry: DirectRegistry,
155}
156
157impl DirectComponent {
158 pub fn new() -> Self {
159 Self {
160 registry: Arc::new(Mutex::new(HashMap::new())),
161 }
162 }
163}
164
165impl Default for DirectComponent {
166 fn default() -> Self {
167 Self::new()
168 }
169}
170
171impl Component for DirectComponent {
172 fn scheme(&self) -> &str {
173 "direct"
174 }
175
176 fn metadata(&self) -> ComponentMetadata {
177 ComponentMetadata {
178 scheme: "direct".to_string(),
179 version: env!("CARGO_PKG_VERSION").to_string(),
180 description: "Synchronous in-memory direct invocation between routes".to_string(),
181 uri_syntax: "direct:name".to_string(),
182 capabilities: ComponentCapabilities {
183 supports_consumer: true,
184 supports_producer: true,
185 ..Default::default()
186 },
187 uri_options: vec![
188 UriOption::new(
189 "timeout_ms",
190 "Producer call timeout in milliseconds",
191 OptionKind::Int,
192 )
193 .with_default("30000"),
194 UriOption::new(
195 "failIfNoConsumers",
196 "Fail if no consumer registered for the name",
197 OptionKind::Bool,
198 )
199 .with_default("true"),
200 ],
201 ..ComponentMetadata::minimal("direct")
202 }
203 }
204
205 fn create_endpoint(
206 &self,
207 uri: &str,
208 _ctx: &dyn camel_component_api::ComponentContext,
209 ) -> Result<Box<dyn Endpoint>, CamelError> {
210 let config = DirectConfig::from_uri(uri)?;
211 validate_name(&config.name)?;
212 let name = config.name.clone();
213 debug!(endpoint_name = %name, "direct endpoint created");
214 Ok(Box::new(DirectEndpoint {
215 uri: uri.to_string(),
216 config,
217 registry: Arc::clone(&self.registry),
218 }))
219 }
220}
221
222struct DirectEndpoint {
227 uri: String,
228 config: DirectConfig,
229 registry: DirectRegistry,
230}
231
232impl Endpoint for DirectEndpoint {
233 fn uri(&self) -> &str {
234 &self.uri
235 }
236
237 fn create_consumer(
238 &self,
239 rt: Arc<dyn camel_component_api::RuntimeObservability>,
240 ) -> Result<Box<dyn Consumer>, CamelError> {
241 Ok(Box::new(DirectConsumer::new(
242 self.config.name.clone(),
243 Arc::clone(&self.registry),
244 rt,
245 )))
246 }
247
248 fn create_producer(
249 &self,
250 _rt: Arc<dyn camel_component_api::RuntimeObservability>,
251 _ctx: &ProducerContext,
252 ) -> Result<BoxProcessor, CamelError> {
253 Ok(BoxProcessor::new(DirectProducer {
254 name: self.config.name.clone(),
255 registry: Arc::clone(&self.registry),
256 config: self.config.clone(),
257 semaphore: Arc::new(Semaphore::new(1)),
258 pending_permit: None,
259 acquire_fut: None,
260 fail_if_no_consumers: self.config.fail_if_no_consumers,
261 }))
262 }
263}
264
265struct DirectConsumer {
272 name: String,
273 registry: DirectRegistry,
274 cancel: Option<CancellationToken>,
275 handle: Option<JoinHandle<Result<(), CamelError>>>,
276 runtime: Arc<dyn camel_component_api::RuntimeObservability>,
277}
278
279impl DirectConsumer {
280 fn new(
281 name: String,
282 registry: DirectRegistry,
283 runtime: Arc<dyn camel_component_api::RuntimeObservability>,
284 ) -> Self {
285 Self {
286 name,
287 registry,
288 cancel: None,
289 handle: None,
290 runtime,
291 }
292 }
293}
294
295#[async_trait]
296impl Consumer for DirectConsumer {
297 async fn start(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
298 let (tx, mut rx) =
300 mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(32);
301
302 {
304 let mut reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
305 if let Some(existing) = reg.get(&self.name)
306 && !existing.is_closed()
307 {
308 return Err(CamelError::EndpointCreationFailed(format!(
309 "direct endpoint '{}' already has a registered consumer",
310 self.name
311 )));
312 }
313 reg.insert(self.name.clone(), tx);
314 }
315
316 let name = self.name.clone();
317 let registry = Arc::clone(&self.registry);
318 let cancel = context.cancel_token();
319 let cancel_clone = cancel.clone();
320 let route_id = context.route_id().to_owned();
321 let runtime = Arc::clone(&self.runtime);
322
323 info!(endpoint_name = %self.name, "direct consumer started");
324
325 let handle = tokio::spawn(async move {
327 loop {
328 tokio::select! {
329 _ = cancel_clone.cancelled() => {
330 debug!(endpoint_name = %name, "direct consumer received cancellation");
331 break;
332 }
333 msg = rx.recv() => {
334 match msg {
335 Some((exchange, reply_tx)) => {
336 debug!(
337 endpoint_name = %name,
338 exchange_id = %exchange.correlation_id,
339 "direct consumer received exchange"
340 );
341 let result = context.send_and_wait(exchange).await;
342 if let Err(ref err) = result {
343 runtime
348 .metrics()
349 .increment_errors(&route_id, "b-prime:direct:send-and-wait");
350 error!(
352 endpoint_name = %name,
353 error = %err,
354 "direct consumer pipeline error"
355 );
356 }
357 let _ = reply_tx.send(result);
358 }
359 None => break,
360 }
361 }
362 }
363 }
364
365 {
367 let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
368 reg.remove(&name);
369 }
370
371 debug!(endpoint_name = %name, "direct consumer stopped");
372 Ok(())
373 });
374
375 self.cancel = Some(cancel);
376 self.handle = Some(handle);
377 Ok(())
378 }
379
380 async fn stop(&mut self) -> Result<(), CamelError> {
381 if let Some(cancel) = self.cancel.take() {
383 cancel.cancel();
384 }
385
386 if let Some(mut h) = self.handle.take() {
388 if tokio::time::timeout(Duration::from_secs(5), &mut h)
389 .await
390 .is_err()
391 {
392 h.abort();
393 let _ = h.await;
394 warn!(endpoint_name = %self.name, "consumer task did not stop in 5s; aborted");
395 let mut reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
397 reg.remove(&self.name);
398 }
399 } else {
400 let mut reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
402 reg.remove(&self.name);
403 }
404
405 debug!(endpoint_name = %self.name, "direct consumer stopped");
406 Ok(())
407 }
408
409 fn background_task_handle(&mut self) -> Option<JoinHandle<Result<(), CamelError>>> {
410 self.handle.take()
415 }
416}
417
418struct DirectProducer {
425 name: String,
426 registry: DirectRegistry,
427 config: DirectConfig,
428 semaphore: Arc<Semaphore>,
429 pending_permit: Option<OwnedSemaphorePermit>,
430 acquire_fut: Option<AcquirePermitFut>,
431 fail_if_no_consumers: Option<bool>,
432}
433
434impl Clone for DirectProducer {
435 fn clone(&self) -> Self {
436 Self {
437 name: self.name.clone(),
438 registry: self.registry.clone(),
439 config: self.config.clone(),
440 semaphore: self.semaphore.clone(),
441 pending_permit: None,
442 acquire_fut: None,
443 fail_if_no_consumers: self.fail_if_no_consumers,
444 }
445 }
446}
447
448impl Service<Exchange> for DirectProducer {
449 type Response = Exchange;
450 type Error = CamelError;
451 type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
452
453 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
454 if self.pending_permit.is_some() {
456 return Poll::Ready(Ok(()));
457 }
458
459 {
461 let reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
462 match reg.get(&self.name) {
463 None => {
464 if self.fail_if_no_consumers != Some(false) {
465 return Poll::Ready(Err(CamelError::EndpointCreationFailed(format!(
466 "direct endpoint '{}' not registered",
467 self.name
468 ))));
469 }
470 }
471 Some(sender) if sender.is_closed() => {
472 return Poll::Ready(Err(CamelError::EndpointCreationFailed(format!(
473 "direct endpoint '{}' channel closed",
474 self.name
475 ))));
476 }
477 Some(_) => {}
478 }
479 }
480
481 let fut = self
483 .acquire_fut
484 .get_or_insert_with(|| Box::pin(Arc::clone(&self.semaphore).acquire_owned()));
485 match fut.as_mut().poll(cx) {
486 Poll::Ready(Ok(permit)) => {
487 self.acquire_fut = None;
488 self.pending_permit = Some(permit);
489 Poll::Ready(Ok(()))
490 }
491 Poll::Ready(Err(_)) => Poll::Ready(Err(CamelError::ChannelClosed)),
492 Poll::Pending => Poll::Pending,
493 }
494 }
495
496 fn call(&mut self, exchange: Exchange) -> Self::Future {
497 let _permit = match self.pending_permit.take() {
498 Some(p) => p,
499 None => {
500 return Box::pin(async {
501 Err(CamelError::ProcessorError(
502 "call() invoked without poll_ready()".into(),
503 ))
504 });
505 }
506 };
507
508 let name = self.name.clone();
509 let registry = Arc::clone(&self.registry);
510 let timeout = Duration::from_millis(self.config.timeout_ms.unwrap_or(30_000));
511 let exchange_id = exchange.correlation_id.clone();
512
513 debug!(
514 endpoint_name = %name,
515 exchange_id = %exchange_id,
516 "direct producer call entry"
517 );
518
519 Box::pin(async move {
520 tokio::time::timeout(timeout, async {
521 let sender = {
522 let reg = registry.lock().unwrap_or_else(|e| e.into_inner());
523 reg.get(&name)
524 .ok_or_else(|| {
525 let err = CamelError::EndpointCreationFailed(format!(
526 "no consumer registered for direct:{name}"
527 ));
528 warn!(endpoint_name = %name, error = %err, "direct send failed");
531 err
532 })?
533 .clone()
534 };
535
536 let (reply_tx, reply_rx) = oneshot::channel();
537 sender.send((exchange, reply_tx)).await.map_err(|err| {
538 warn!(endpoint_name = %name, error = %err, "direct send failed");
541 CamelError::ChannelClosed
542 })?;
543
544 let result = reply_rx.await.map_err(|err| {
545 warn!(endpoint_name = %name, error = %err, "direct send failed");
548 CamelError::ChannelClosed
549 })?;
550
551 debug!(endpoint_name = %name, "direct message sent");
552 result
553 })
554 .await
555 .map_err(|_| CamelError::ProcessorError(format!("direct:{name} call timed out")))?
556 })
557 }
558}
559
560#[cfg(test)]
565mod tests {
566 use std::sync::Mutex;
567
568 use camel_api::MetricsCollector;
569 use camel_component_api::HealthCheckRegistry;
570 use std::time::Duration;
571 fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
572 std::sync::Arc::new(NoOpComponentContext)
576 }
577
578 use super::*;
579 use camel_component_api::ExchangeEnvelope;
580 use camel_component_api::Message;
581 use camel_component_api::NoOpComponentContext;
582 use camel_component_api::RuntimeObservability;
583 use std::task::RawWakerVTable;
584 use tower::ServiceExt;
585
586 struct RecordingMetrics {
591 errors: Arc<Mutex<Vec<(String, String)>>>,
592 }
593
594 impl MetricsCollector for RecordingMetrics {
595 fn record_exchange_duration(&self, _: &str, _: Duration) {}
596 fn increment_errors(&self, route_id: &str, error_type: &str) {
597 self.errors
598 .lock()
599 .unwrap()
600 .push((route_id.to_string(), error_type.to_string()));
601 }
602 fn increment_exchanges(&self, _: &str) {}
603 fn set_queue_depth(&self, _: &str, _: usize) {}
604 fn record_circuit_breaker_change(&self, _: &str, _: &str, _: &str) {}
605 }
606
607 struct RecordingRuntime {
608 metrics_collector: Arc<RecordingMetrics>,
609 }
610
611 impl RecordingRuntime {
612 fn new(errors: Arc<Mutex<Vec<(String, String)>>>) -> Self {
613 Self {
614 metrics_collector: Arc::new(RecordingMetrics { errors }),
615 }
616 }
617 }
618
619 impl RuntimeObservability for RecordingRuntime {
620 fn metrics(&self) -> Arc<dyn MetricsCollector> {
621 self.metrics_collector.clone() as Arc<dyn MetricsCollector>
622 }
623 fn health(&self) -> Arc<dyn HealthCheckRegistry> {
624 panic!("RecordingRuntime::health not used in this test")
625 }
626 }
627
628 fn noop_waker() -> std::task::Waker {
629 const VTABLE: RawWakerVTable = RawWakerVTable::new(|_| RAW, |_| {}, |_| {}, |_| {});
630 const RAW: std::task::RawWaker = std::task::RawWaker::new(std::ptr::null(), &VTABLE);
631 unsafe { std::task::Waker::from_raw(RAW) }
632 }
633
634 fn test_producer_ctx() -> ProducerContext {
635 ProducerContext::new()
636 }
637
638 #[test]
639 fn test_direct_component_scheme() {
640 let component = DirectComponent::new();
641 assert_eq!(component.scheme(), "direct");
642 }
643
644 #[test]
645 fn test_direct_component_default() {
646 let component = DirectComponent::default();
647 assert_eq!(component.scheme(), "direct");
648 }
649
650 #[test]
651 fn test_direct_config_from_uri() {
652 let config = DirectConfig::from_uri("direct:orders").unwrap();
653 assert_eq!(config.name, "orders");
654 }
655
656 #[test]
657 fn test_direct_endpoint_uri() {
658 let component = DirectComponent::new();
659 let endpoint = component
660 .create_endpoint("direct:uri-check", &NoOpComponentContext)
661 .unwrap();
662 assert_eq!(endpoint.uri(), "direct:uri-check");
663 }
664
665 #[test]
666 fn test_direct_creates_endpoint() {
667 let component = DirectComponent::new();
668 let endpoint = component.create_endpoint("direct:foo", &NoOpComponentContext);
669 assert!(endpoint.is_ok());
670 }
671
672 #[test]
673 fn test_direct_wrong_scheme() {
674 let component = DirectComponent::new();
675 let result = component.create_endpoint("timer:tick", &NoOpComponentContext);
676 assert!(result.is_err());
677 }
678
679 #[test]
680 fn test_direct_endpoint_creates_consumer() {
681 let component = DirectComponent::new();
682 let endpoint = component
683 .create_endpoint("direct:foo", &NoOpComponentContext)
684 .unwrap();
685 assert!(endpoint.create_consumer(rt()).is_ok());
686 }
687
688 #[test]
689 fn test_direct_endpoint_creates_producer() {
690 let ctx = test_producer_ctx();
691 let component = DirectComponent::new();
692 let endpoint = component
693 .create_endpoint("direct:foo", &NoOpComponentContext)
694 .unwrap();
695 assert!(endpoint.create_producer(rt(), &ctx).is_ok());
696 }
697
698 #[test]
699 fn test_direct_empty_name_rejected() {
700 let component = DirectComponent::new();
701 match component.create_endpoint("direct:", &NoOpComponentContext) {
702 Err(e) => assert!(
703 e.to_string().contains("must not be empty"),
704 "unexpected error: {e}"
705 ),
706 Ok(_) => panic!("expected error for empty name"),
707 }
708 }
709
710 #[tokio::test]
711 async fn test_direct_producer_no_consumer_registered() {
712 let ctx = test_producer_ctx();
713 let component = DirectComponent::new();
714 let endpoint = component
715 .create_endpoint("direct:missing", &NoOpComponentContext)
716 .unwrap();
717 let producer = endpoint.create_producer(rt(), &ctx).unwrap();
718
719 let exchange = Exchange::new(Message::new("test"));
720 let result = producer.oneshot(exchange).await;
721 assert!(result.is_err());
722 }
723
724 #[tokio::test]
725 async fn test_direct_duplicate_consumer_returns_error() {
726 let component = DirectComponent::new();
727 let endpoint = component
728 .create_endpoint("direct:dup", &NoOpComponentContext)
729 .unwrap();
730
731 let mut consumer_a = endpoint.create_consumer(rt()).unwrap();
732 let mut consumer_b = endpoint.create_consumer(rt()).unwrap();
733
734 let (route_tx_a, _route_rx_a) = mpsc::channel::<ExchangeEnvelope>(16);
735 let ctx_a = ConsumerContext::new(
736 route_tx_a,
737 tokio_util::sync::CancellationToken::new(),
738 "direct-test-route-a".to_string(),
739 );
740 consumer_a.start(ctx_a).await.unwrap();
741
742 let (route_tx_b, _route_rx_b) = mpsc::channel::<ExchangeEnvelope>(16);
743 let ctx_b = ConsumerContext::new(
744 route_tx_b,
745 tokio_util::sync::CancellationToken::new(),
746 "direct-test-route-b".to_string(),
747 );
748 let result = consumer_b.start(ctx_b).await;
749
750 assert!(matches!(
751 result,
752 Err(CamelError::EndpointCreationFailed(msg))
753 if msg.contains("already has a registered consumer")
754 ));
755
756 consumer_a.stop().await.unwrap();
757 }
758
759 #[tokio::test]
760 async fn test_direct_producer_consumer_roundtrip() {
761 let component = DirectComponent::new();
762
763 let consumer_endpoint = component
765 .create_endpoint("direct:test", &NoOpComponentContext)
766 .unwrap();
767 let mut consumer = consumer_endpoint.create_consumer(rt()).unwrap();
768
769 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
771 let ctx = ConsumerContext::new(
772 route_tx,
773 tokio_util::sync::CancellationToken::new(),
774 "direct-test-route".to_string(),
775 );
776
777 tokio::spawn(async move {
779 consumer.start(ctx).await.unwrap();
780 });
781
782 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
784
785 tokio::spawn(async move {
787 while let Some(envelope) = route_rx.recv().await {
788 let ExchangeEnvelope { exchange, reply_tx } = envelope;
789 if let Some(tx) = reply_tx {
790 let _ = tx.send(Ok(exchange));
791 }
792 }
793 });
794
795 let ctx = test_producer_ctx();
797 let producer_endpoint = component
798 .create_endpoint("direct:test", &NoOpComponentContext)
799 .unwrap();
800 let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
801
802 let exchange = Exchange::new(Message::new("hello direct"));
803 let result = producer.oneshot(exchange).await;
804
805 assert!(result.is_ok());
806 let reply = result.unwrap();
807 assert_eq!(reply.input.body.as_text(), Some("hello direct"));
808 }
809
810 #[tokio::test]
811 async fn test_direct_propagates_error_when_no_handler() {
812 let component = DirectComponent::new();
813
814 let consumer_endpoint = component
815 .create_endpoint("direct:err-test", &NoOpComponentContext)
816 .unwrap();
817 let mut consumer = consumer_endpoint.create_consumer(rt()).unwrap();
818
819 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
820 let ctx = ConsumerContext::new(
821 route_tx,
822 tokio_util::sync::CancellationToken::new(),
823 "direct-test-route".to_string(),
824 );
825
826 tokio::spawn(async move {
827 consumer.start(ctx).await.unwrap();
828 });
829
830 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
831
832 tokio::spawn(async move {
834 while let Some(envelope) = route_rx.recv().await {
835 if let Some(tx) = envelope.reply_tx {
836 let _ = tx.send(Err(CamelError::ProcessorError("subroute failed".into())));
837 }
838 }
839 });
840
841 let ctx = test_producer_ctx();
842 let producer_endpoint = component
843 .create_endpoint("direct:err-test", &NoOpComponentContext)
844 .unwrap();
845 let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
846
847 let exchange = Exchange::new(Message::new("test"));
848 let result = producer.oneshot(exchange).await;
849 assert!(result.is_err());
850 assert!(matches!(result.unwrap_err(), CamelError::ProcessorError(_)));
851 }
852
853 #[tokio::test]
854 async fn test_direct_consumer_stop_unregisters() {
855 let component = DirectComponent::new();
856 let endpoint = component
857 .create_endpoint("direct:cleanup", &NoOpComponentContext)
858 .unwrap();
859
860 let mut consumer = endpoint.create_consumer(rt()).unwrap();
862
863 let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
864 let ctx = ConsumerContext::new(
865 route_tx,
866 tokio_util::sync::CancellationToken::new(),
867 "direct-test-route".to_string(),
868 );
869
870 let handle = tokio::spawn(async move {
872 consumer.start(ctx).await.unwrap();
873 });
874
875 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
876
877 {
879 let reg = component.registry.lock().unwrap_or_else(|e| e.into_inner());
880 assert!(reg.contains_key("cleanup"));
881 }
882
883 let mut stop_consumer = DirectConsumer {
885 name: "cleanup".to_string(),
886 registry: Arc::clone(&component.registry),
887 cancel: None,
888 handle: None,
889 runtime: rt(),
890 };
891 stop_consumer.stop().await.unwrap();
892
893 {
895 let reg = component.registry.lock().unwrap_or_else(|e| e.into_inner());
896 assert!(!reg.contains_key("cleanup"));
897 }
898
899 handle.abort();
900 }
901
902 #[tokio::test]
903 async fn test_direct_consumer_respects_cancellation() {
904 use tokio_util::sync::CancellationToken;
905
906 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
907 let token = CancellationToken::new();
908 let (tx, _rx) = mpsc::channel(16);
909 let ctx = ConsumerContext::new(tx, token.clone(), "direct-test-route".to_string());
910
911 let mut consumer = DirectConsumer {
912 name: "cancel-test".to_string(),
913 registry: registry.clone(),
914 cancel: None,
915 handle: None,
916 runtime: rt(),
917 };
918
919 let handle = tokio::spawn(async move {
920 consumer.start(ctx).await.unwrap();
921 });
922
923 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
924 assert!(
925 registry
926 .lock()
927 .unwrap_or_else(|e| e.into_inner())
928 .contains_key("cancel-test")
929 );
930
931 token.cancel();
932 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
936
937 assert!(
939 !registry
940 .lock()
941 .unwrap_or_else(|e| e.into_inner())
942 .contains_key("cancel-test")
943 );
944
945 let _ = handle.await;
946 }
947
948 #[tokio::test]
949 async fn test_direct_consumer_stop_missing_entry_is_ok() {
950 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
951 let mut consumer = DirectConsumer {
952 name: "never-registered".to_string(),
953 registry,
954 cancel: None,
955 handle: None,
956 runtime: rt(),
957 };
958 let result = consumer.stop().await;
959 assert!(result.is_ok());
960 }
961
962 #[test]
963 fn test_poll_ready_endpoint_not_registered() {
964 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
965 let producer = DirectProducer {
966 name: "missing".to_string(),
967 registry,
968 config: DirectConfig {
969 name: "missing".to_string(),
970 timeout_ms: None,
971 block: None,
972 fail_if_no_consumers: None,
973 exchange_pattern: None,
974 },
975 semaphore: Arc::new(Semaphore::new(1)),
976 pending_permit: None,
977 acquire_fut: None,
978 fail_if_no_consumers: None,
979 };
980 let waker = noop_waker();
981 let mut cx = Context::from_waker(&waker);
982 let mut producer = producer;
983 let result = producer.poll_ready(&mut cx);
984 assert!(matches!(
985 result,
986 Poll::Ready(Err(CamelError::EndpointCreationFailed(_)))
987 ));
988 }
989
990 #[test]
991 fn test_poll_ready_endpoint_registered() {
992 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
993 let (tx, _rx) =
994 mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(1);
995 registry.lock().unwrap().insert("active".to_string(), tx);
996 let producer = DirectProducer {
997 name: "active".to_string(),
998 registry,
999 config: DirectConfig {
1000 name: "active".to_string(),
1001 timeout_ms: None,
1002 block: None,
1003 fail_if_no_consumers: None,
1004 exchange_pattern: None,
1005 },
1006 semaphore: Arc::new(Semaphore::new(1)),
1007 pending_permit: None,
1008 acquire_fut: None,
1009 fail_if_no_consumers: None,
1010 };
1011 let waker = noop_waker();
1012 let mut cx = Context::from_waker(&waker);
1013 let mut producer = producer;
1014 let result = producer.poll_ready(&mut cx);
1015 assert!(matches!(result, Poll::Ready(Ok(()))));
1016 }
1017
1018 #[test]
1019 fn test_poll_ready_allows_missing_consumer_when_fail_if_no_consumers_false() {
1020 let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
1021 let producer = DirectProducer {
1022 name: "missing-ok".to_string(),
1023 registry,
1024 config: DirectConfig {
1025 name: "missing-ok".to_string(),
1026 timeout_ms: None,
1027 block: None,
1028 fail_if_no_consumers: Some(false),
1029 exchange_pattern: None,
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 block: None,
1058 fail_if_no_consumers: None,
1059 exchange_pattern: None,
1060 },
1061 semaphore: Arc::new(Semaphore::new(1)),
1062 pending_permit: None,
1063 acquire_fut: None,
1064 fail_if_no_consumers: None,
1065 };
1066 let waker = noop_waker();
1067 let mut cx = Context::from_waker(&waker);
1068 let mut producer = producer;
1069 let result = producer.poll_ready(&mut cx);
1070 assert!(matches!(
1071 result,
1072 Poll::Ready(Err(CamelError::EndpointCreationFailed(_)))
1073 ));
1074 }
1075
1076 #[tokio::test]
1077 async fn test_direct_stop_cancels_loop() {
1078 use tokio_util::sync::CancellationToken;
1079
1080 let component = DirectComponent::new();
1081 let endpoint = component
1082 .create_endpoint("direct:stop-test", &NoOpComponentContext)
1083 .unwrap();
1084 let mut consumer = endpoint.create_consumer(rt()).unwrap();
1085
1086 let token = CancellationToken::new();
1087 let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1088 let ctx = ConsumerContext::new(route_tx, token.clone(), "direct-test-route".to_string());
1089
1090 let handle = tokio::spawn(async move {
1094 consumer.start(ctx).await.unwrap();
1095 });
1096
1097 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1098 assert!(
1099 component
1100 .registry
1101 .lock()
1102 .unwrap_or_else(|e| e.into_inner())
1103 .contains_key("stop-test")
1104 );
1105
1106 let mut stop_consumer = DirectConsumer {
1108 name: "stop-test".to_string(),
1109 registry: Arc::clone(&component.registry),
1110 cancel: Some(token.clone()),
1111 handle: None,
1112 runtime: rt(),
1113 };
1114 stop_consumer.stop().await.unwrap();
1115
1116 let result = tokio::time::timeout(std::time::Duration::from_secs(2), handle).await;
1118 assert!(result.is_ok(), "Consumer loop did not stop within 2s");
1119
1120 assert!(
1122 !component
1123 .registry
1124 .lock()
1125 .unwrap_or_else(|e| e.into_inner())
1126 .contains_key("stop-test")
1127 );
1128 }
1129
1130 #[tokio::test]
1131 async fn test_direct_producer_timeout() {
1132 let component = DirectComponent::new();
1133 let endpoint = component
1134 .create_endpoint("direct:timeout-test", &NoOpComponentContext)
1135 .unwrap();
1136 let mut consumer = endpoint.create_consumer(rt()).unwrap();
1137
1138 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1140 let token = tokio_util::sync::CancellationToken::new();
1141 let ctx = ConsumerContext::new(route_tx, token.clone(), "direct-test-route".to_string());
1142 tokio::spawn(async move {
1143 consumer.start(ctx).await.unwrap();
1144 });
1145
1146 tokio::spawn(async move {
1148 let mut held_reply: Vec<oneshot::Sender<Result<Exchange, CamelError>>> = Vec::new();
1149 while let Some(envelope) = route_rx.recv().await {
1150 held_reply.push(envelope.reply_tx.unwrap());
1151 }
1152 drop(held_reply);
1153 });
1154
1155 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1156
1157 let _ = test_producer_ctx();
1159 let _producer_endpoint = component
1160 .create_endpoint("direct:timeout-test", &NoOpComponentContext)
1161 .unwrap();
1162 let producer = DirectProducer {
1163 name: "timeout-test".to_string(),
1164 registry: Arc::clone(&component.registry),
1165 config: DirectConfig {
1166 name: "timeout-test".to_string(),
1167 timeout_ms: Some(100), block: None,
1169 fail_if_no_consumers: None,
1170 exchange_pattern: None,
1171 },
1172 semaphore: Arc::new(Semaphore::new(1)),
1173 pending_permit: None,
1174 acquire_fut: None,
1175 fail_if_no_consumers: None,
1176 };
1177
1178 let exchange = Exchange::new(Message::new("test"));
1179 let mut svc = producer;
1180 let _ = svc.poll_ready(&mut Context::from_waker(&noop_waker()));
1181 let result = svc.call(exchange).await;
1182 assert!(result.is_err(), "Expected timeout error");
1183 assert!(
1184 result.unwrap_err().to_string().contains("timed out"),
1185 "Expected timeout message"
1186 );
1187
1188 token.cancel();
1189 }
1190
1191 #[tokio::test]
1192 async fn test_send_and_wait_error_increments_errors_metric() {
1193 let errors: Arc<Mutex<Vec<(String, String)>>> = Arc::new(Mutex::new(Vec::new()));
1194 let runtime = Arc::new(RecordingRuntime::new(Arc::clone(&errors)));
1195
1196 let component = DirectComponent::new();
1197 let endpoint = component
1198 .create_endpoint("direct:metrics-error", &NoOpComponentContext)
1199 .unwrap();
1200 let mut consumer = endpoint.create_consumer(runtime).unwrap();
1201
1202 let cancel = tokio_util::sync::CancellationToken::new();
1204 let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1205 let ctx = ConsumerContext::new(route_tx, cancel.clone(), "test-route-id".to_string());
1206
1207 tokio::spawn(async move {
1208 consumer.start(ctx).await.unwrap();
1209 });
1210
1211 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1212
1213 tokio::spawn(async move {
1215 while let Some(envelope) = route_rx.recv().await {
1216 if let Some(tx) = envelope.reply_tx {
1217 let _ = tx.send(Err(CamelError::ProcessorError(
1218 "pipeline failure".to_string(),
1219 )));
1220 }
1221 }
1222 });
1223
1224 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1225
1226 let ctx = test_producer_ctx();
1229 let producer_endpoint = component
1230 .create_endpoint("direct:metrics-error", &NoOpComponentContext)
1231 .unwrap();
1232 let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
1233
1234 let exchange = Exchange::new(Message::new("test"));
1235 let _result = producer.oneshot(exchange).await;
1238
1239 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1240
1241 cancel.cancel();
1242
1243 let recorded = errors.lock().unwrap();
1245 assert_eq!(
1246 recorded.len(),
1247 1,
1248 "expected 1 increment_errors call, got {}: {:?}",
1249 recorded.len(),
1250 *recorded
1251 );
1252 assert_eq!(recorded[0].0, "test-route-id");
1253 assert_eq!(recorded[0].1, "b-prime:direct:send-and-wait");
1254 }
1255
1256 #[test]
1257 fn test_empty_endpoint_name_rejected() {
1258 let result = DirectConfig::from_uri("direct:");
1259 if let Ok(config) = result {
1261 assert!(
1262 validate_name(&config.name).is_err(),
1263 "expected validation error for empty name"
1264 );
1265 }
1266 let component = DirectComponent::new();
1268 let result = component.create_endpoint("direct:", &NoOpComponentContext);
1269 assert!(result.is_err(), "empty endpoint name must be rejected");
1270 }
1271
1272 #[test]
1273 fn test_whitespace_endpoint_name_rejected() {
1274 let result = DirectConfig::from_uri("direct:my endpoint");
1275 if let Ok(config) = result {
1276 assert!(
1277 validate_name(&config.name).is_err(),
1278 "expected validation error for whitespace in name"
1279 );
1280 }
1281 let component = DirectComponent::new();
1282 let result = component.create_endpoint("direct:my endpoint", &NoOpComponentContext);
1283 assert!(result.is_err(), "whitespace endpoint name must be rejected");
1284 }
1285
1286 #[test]
1287 fn test_valid_endpoint_name_accepted() {
1288 let component = DirectComponent::new();
1289 let result = component.create_endpoint("direct:my-endpoint", &NoOpComponentContext);
1290 assert!(result.is_ok(), "valid endpoint name should be accepted");
1291 }
1292}