Skip to main content

camel_component_direct/
lib.rs

1//! In-memory direct component for rust-camel — synchronous point-to-point
2//! channel between routes sharing the same context with no serialization overhead.
3//!
4//! Main types: `DirectComponent`, `DirectEndpoint`, `DirectConsumer`, `DirectProducer`.
5
6use 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_api::component_metadata::{
19    ComponentCapabilities, ComponentMetadata, OptionKind, UriOption,
20};
21use camel_component_api::parse_uri;
22use camel_component_api::{BoxProcessor, CamelError, Exchange};
23use camel_component_api::{Component, Consumer, ConsumerContext, Endpoint, ProducerContext};
24use tracing::{debug, error, info, warn};
25
26// ---------------------------------------------------------------------------
27// Shared state: maps endpoint names to senders that deliver exchanges to the
28// consumer side.  Each entry holds a sender of `(Exchange, oneshot::Sender)`
29// so the producer can wait for the consumer's pipeline to finish processing
30// and receive the (possibly transformed) exchange back.
31// ---------------------------------------------------------------------------
32
33type DirectSender = mpsc::Sender<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>;
34type DirectRegistry = Arc<Mutex<HashMap<String, DirectSender>>>;
35type AcquirePermitFut =
36    Pin<Box<dyn Future<Output = Result<OwnedSemaphorePermit, AcquireError>> + Send>>;
37
38// ---------------------------------------------------------------------------
39// Validation helpers
40// ---------------------------------------------------------------------------
41
42/// Validate the direct endpoint name (the part after `direct:`).
43fn validate_name(name: &str) -> Result<(), CamelError> {
44    if name.trim().is_empty() {
45        return Err(CamelError::InvalidUri(
46            "direct: endpoint name must not be empty".to_string(),
47        ));
48    }
49    if name.contains(char::is_whitespace) {
50        return Err(CamelError::InvalidUri(
51            "direct: endpoint name must not contain whitespace".to_string(),
52        ));
53    }
54    Ok(())
55}
56
57// ---------------------------------------------------------------------------
58// DirectConfig
59// ---------------------------------------------------------------------------
60
61/// Configuration for Direct endpoints parsed from URIs.
62///
63/// URI format: `direct:name[?timeout_ms=30000]`
64///
65/// Example: `direct:foo` creates an endpoint named "foo"
66#[derive(Debug, Clone)]
67pub struct DirectConfig {
68    /// Endpoint name (path portion).
69    pub name: String,
70    /// Timeout in milliseconds for producer `call()`. Defaults to 30 000 ms.
71    pub timeout_ms: Option<u64>,
72    /// When false, the producer returns immediately if no consumer is registered.
73    /// TODO(DIR-001): implement non-blocking send
74    pub block: Option<bool>,
75    /// When false, skip readiness error if no consumer registered.
76    pub fail_if_no_consumers: Option<bool>,
77    /// TODO(DIR-005): implement exchangePattern override
78    pub exchange_pattern: Option<String>,
79}
80
81impl DirectConfig {
82    pub fn from_uri(uri: &str) -> Result<Self, CamelError> {
83        let parts = parse_uri(uri)?;
84        if parts.scheme != "direct" {
85            return Err(CamelError::InvalidUri(format!(
86                "invalid scheme '{}', expected 'direct'",
87                parts.scheme
88            )));
89        }
90
91        let parse_bool = |name: &str, value: &str| -> Result<bool, CamelError> {
92            match value.to_ascii_lowercase().as_str() {
93                "true" | "1" | "yes" => Ok(true),
94                "false" | "0" | "no" => Ok(false),
95                _ => Err(CamelError::InvalidUri(format!(
96                    "invalid value for {}: invalid boolean value: '{}'",
97                    name, value
98                ))),
99            }
100        };
101
102        let timeout_ms = parts
103            .params
104            .get("timeout_ms")
105            .map(|v| {
106                v.parse::<u64>().map_err(|e| {
107                    CamelError::InvalidUri(format!("invalid value for timeout_ms: {}", e))
108                })
109            })
110            .transpose()?;
111
112        let block = parts
113            .params
114            .get("block")
115            .map(|v| parse_bool("block", v))
116            .transpose()?;
117
118        let fail_if_no_consumers = parts
119            .params
120            .get("fail_if_no_consumers")
121            .or_else(|| parts.params.get("failIfNoConsumers"))
122            .map(|v| parse_bool("fail_if_no_consumers", v))
123            .transpose()?;
124
125        let exchange_pattern = parts
126            .params
127            .get("exchange_pattern")
128            .or_else(|| parts.params.get("exchangePattern"))
129            .cloned();
130
131        Ok(Self {
132            name: parts.path,
133            timeout_ms,
134            block,
135            fail_if_no_consumers,
136            exchange_pattern,
137        })
138    }
139}
140
141// ---------------------------------------------------------------------------
142// DirectComponent
143// ---------------------------------------------------------------------------
144
145/// The Direct component provides in-memory synchronous communication between
146/// routes.
147///
148/// URI format: `direct:name`
149///
150/// A producer sending to `direct:foo` will block until the consumer on
151/// `direct:foo` has finished processing the exchange.
152pub struct DirectComponent {
153    registry: DirectRegistry,
154}
155
156impl DirectComponent {
157    pub fn new() -> Self {
158        Self {
159            registry: Arc::new(Mutex::new(HashMap::new())),
160        }
161    }
162}
163
164impl Default for DirectComponent {
165    fn default() -> Self {
166        Self::new()
167    }
168}
169
170impl Component for DirectComponent {
171    fn scheme(&self) -> &str {
172        "direct"
173    }
174
175    fn metadata(&self) -> ComponentMetadata {
176        ComponentMetadata {
177            scheme: "direct".to_string(),
178            version: env!("CARGO_PKG_VERSION").to_string(),
179            description: "Synchronous in-memory direct invocation between routes".to_string(),
180            uri_syntax: "direct:name".to_string(),
181            capabilities: ComponentCapabilities {
182                supports_consumer: true,
183                supports_producer: true,
184                ..Default::default()
185            },
186            uri_options: vec![
187                UriOption::new(
188                    "timeout_ms",
189                    "Producer call timeout in milliseconds",
190                    OptionKind::Int,
191                )
192                .with_default("30000"),
193                UriOption::new(
194                    "failIfNoConsumers",
195                    "Fail if no consumer registered for the name",
196                    OptionKind::Bool,
197                )
198                .with_default("true"),
199            ],
200            ..ComponentMetadata::minimal("direct")
201        }
202    }
203
204    fn create_endpoint(
205        &self,
206        uri: &str,
207        _ctx: &dyn camel_component_api::ComponentContext,
208    ) -> Result<Box<dyn Endpoint>, CamelError> {
209        let config = DirectConfig::from_uri(uri)?;
210        validate_name(&config.name)?;
211        let name = config.name.clone();
212        debug!(endpoint_name = %name, "direct endpoint created");
213        Ok(Box::new(DirectEndpoint {
214            uri: uri.to_string(),
215            config,
216            registry: Arc::clone(&self.registry),
217        }))
218    }
219}
220
221// ---------------------------------------------------------------------------
222// DirectEndpoint
223// ---------------------------------------------------------------------------
224
225struct DirectEndpoint {
226    uri: String,
227    config: DirectConfig,
228    registry: DirectRegistry,
229}
230
231impl Endpoint for DirectEndpoint {
232    fn uri(&self) -> &str {
233        &self.uri
234    }
235
236    fn create_consumer(
237        &self,
238        rt: Arc<dyn camel_component_api::RuntimeObservability>,
239    ) -> Result<Box<dyn Consumer>, CamelError> {
240        Ok(Box::new(DirectConsumer::new(
241            self.config.name.clone(),
242            Arc::clone(&self.registry),
243            rt,
244        )))
245    }
246
247    fn create_producer(
248        &self,
249        _rt: Arc<dyn camel_component_api::RuntimeObservability>,
250        _ctx: &ProducerContext,
251    ) -> Result<BoxProcessor, CamelError> {
252        Ok(BoxProcessor::new(DirectProducer {
253            name: self.config.name.clone(),
254            registry: Arc::clone(&self.registry),
255            config: self.config.clone(),
256            semaphore: Arc::new(Semaphore::new(1)),
257            pending_permit: None,
258            acquire_fut: None,
259            fail_if_no_consumers: self.config.fail_if_no_consumers,
260        }))
261    }
262}
263
264// ---------------------------------------------------------------------------
265// DirectConsumer
266// ---------------------------------------------------------------------------
267
268/// The Direct consumer registers itself in the shared registry and forwards
269/// incoming exchanges to the route pipeline via `ConsumerContext`.
270struct DirectConsumer {
271    name: String,
272    registry: DirectRegistry,
273    cancel: Option<CancellationToken>,
274    runtime: Arc<dyn camel_component_api::RuntimeObservability>,
275}
276
277impl DirectConsumer {
278    fn new(
279        name: String,
280        registry: DirectRegistry,
281        runtime: Arc<dyn camel_component_api::RuntimeObservability>,
282    ) -> Self {
283        Self {
284            name,
285            registry,
286            cancel: None,
287            runtime,
288        }
289    }
290}
291
292#[async_trait]
293impl Consumer for DirectConsumer {
294    async fn start(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
295        // Create a channel for producers to send exchanges to this consumer.
296        let (tx, mut rx) =
297            mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(32);
298
299        // Register ourselves so producers can find us.
300        {
301            let mut reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
302            if let Some(existing) = reg.get(&self.name)
303                && !existing.is_closed()
304            {
305                return Err(CamelError::EndpointCreationFailed(format!(
306                    "direct endpoint '{}' already has a registered consumer",
307                    self.name
308                )));
309            }
310            reg.insert(self.name.clone(), tx);
311        }
312
313        let name = self.name.clone();
314        let registry = Arc::clone(&self.registry);
315        let cancel = context.cancel_token();
316        let cancel_clone = cancel.clone();
317        let route_id = context.route_id().to_owned();
318        let runtime = Arc::clone(&self.runtime);
319
320        info!(endpoint_name = %self.name, "direct consumer started");
321
322        self.cancel = Some(cancel);
323
324        loop {
325            tokio::select! {
326                _ = cancel_clone.cancelled() => {
327                    debug!(endpoint_name = %name, "direct consumer received cancellation");
328                    break;
329                }
330                msg = rx.recv() => {
331                    match msg {
332                        Some((exchange, reply_tx)) => {
333                            debug!(
334                                endpoint_name = %name,
335                                exchange_id = %exchange.correlation_id,
336                                "direct consumer received exchange"
337                            );
338                            let result = context.send_and_wait(exchange).await;
339                            if let Err(ref err) = result {
340                                // (category b′: send_and_wait returned Err on a normal-data send,
341                                // meaning the route handler did NOT absorb the failure —
342                                // see ADR-0012 "b-bridged discriminator". This emitter is the
343                                // only ERROR signal for the unhandled failure; must stay loud.)
344                                runtime
345                                    .metrics()
346                                    .increment_errors(&route_id, "b-prime:direct:send-and-wait");
347                                // log-policy: outside-contract
348                                error!(
349                                    endpoint_name = %name,
350                                    error = %err,
351                                    "direct consumer pipeline error"
352                                );
353                            }
354                            let _ = reply_tx.send(result);
355                        }
356                        None => break,
357                    }
358                }
359            }
360        }
361
362        // Cleanup: remove from registry on exit
363        {
364            let mut reg = registry.lock().unwrap_or_else(|e| e.into_inner());
365            reg.remove(&name);
366        }
367
368        debug!(endpoint_name = %name, "direct consumer stopped");
369        Ok(())
370    }
371
372    async fn stop(&mut self) -> Result<(), CamelError> {
373        // Cancel the consumer loop if we have a cancellation token.
374        if let Some(cancel) = self.cancel.take() {
375            cancel.cancel();
376        }
377
378        let mut reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
379        reg.remove(&self.name);
380
381        debug!(endpoint_name = %self.name, "direct consumer stopped");
382        Ok(())
383    }
384}
385
386// ---------------------------------------------------------------------------
387// DirectProducer
388// ---------------------------------------------------------------------------
389
390/// The Direct producer sends an exchange to the named direct endpoint and
391/// waits for a reply (synchronous in-memory call).
392struct DirectProducer {
393    name: String,
394    registry: DirectRegistry,
395    config: DirectConfig,
396    semaphore: Arc<Semaphore>,
397    pending_permit: Option<OwnedSemaphorePermit>,
398    acquire_fut: Option<AcquirePermitFut>,
399    fail_if_no_consumers: Option<bool>,
400}
401
402impl Clone for DirectProducer {
403    fn clone(&self) -> Self {
404        Self {
405            name: self.name.clone(),
406            registry: self.registry.clone(),
407            config: self.config.clone(),
408            semaphore: self.semaphore.clone(),
409            pending_permit: None,
410            acquire_fut: None,
411            fail_if_no_consumers: self.fail_if_no_consumers,
412        }
413    }
414}
415
416impl Service<Exchange> for DirectProducer {
417    type Response = Exchange;
418    type Error = CamelError;
419    type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
420
421    fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
422        // If we already hold a permit we are ready.
423        if self.pending_permit.is_some() {
424            return Poll::Ready(Ok(()));
425        }
426
427        // Check that the endpoint is registered.
428        {
429            let reg = self.registry.lock().unwrap_or_else(|e| e.into_inner());
430            match reg.get(&self.name) {
431                None => {
432                    if self.fail_if_no_consumers != Some(false) {
433                        return Poll::Ready(Err(CamelError::EndpointCreationFailed(format!(
434                            "direct endpoint '{}' not registered",
435                            self.name
436                        ))));
437                    }
438                }
439                Some(sender) if sender.is_closed() => {
440                    return Poll::Ready(Err(CamelError::EndpointCreationFailed(format!(
441                        "direct endpoint '{}' channel closed",
442                        self.name
443                    ))));
444                }
445                Some(_) => {}
446            }
447        }
448
449        // Acquire a semaphore permit (bounded concurrency).
450        let fut = self
451            .acquire_fut
452            .get_or_insert_with(|| Box::pin(Arc::clone(&self.semaphore).acquire_owned()));
453        match fut.as_mut().poll(cx) {
454            Poll::Ready(Ok(permit)) => {
455                self.acquire_fut = None;
456                self.pending_permit = Some(permit);
457                Poll::Ready(Ok(()))
458            }
459            Poll::Ready(Err(_)) => Poll::Ready(Err(CamelError::ChannelClosed)),
460            Poll::Pending => Poll::Pending,
461        }
462    }
463
464    fn call(&mut self, exchange: Exchange) -> Self::Future {
465        let _permit = match self.pending_permit.take() {
466            Some(p) => p,
467            None => {
468                return Box::pin(async {
469                    Err(CamelError::ProcessorError(
470                        "call() invoked without poll_ready()".into(),
471                    ))
472                });
473            }
474        };
475
476        let name = self.name.clone();
477        let registry = Arc::clone(&self.registry);
478        let timeout = Duration::from_millis(self.config.timeout_ms.unwrap_or(30_000));
479        let exchange_id = exchange.correlation_id.clone();
480
481        debug!(
482            endpoint_name = %name,
483            exchange_id = %exchange_id,
484            "direct producer call entry"
485        );
486
487        Box::pin(async move {
488            tokio::time::timeout(timeout, async {
489                let sender = {
490                    let reg = registry.lock().unwrap_or_else(|e| e.into_inner());
491                    reg.get(&name)
492                        .ok_or_else(|| {
493                            let err = CamelError::EndpointCreationFailed(format!(
494                                "no consumer registered for direct:{name}"
495                            ));
496                            // log-policy: handler-owned
497                            // (category a: producer send failure inside the route pipeline)
498                            warn!(endpoint_name = %name, error = %err, "direct send failed");
499                            err
500                        })?
501                        .clone()
502                };
503
504                let (reply_tx, reply_rx) = oneshot::channel();
505                sender.send((exchange, reply_tx)).await.map_err(|err| {
506                    // log-policy: handler-owned
507                    // (category a: producer send failure inside the route pipeline)
508                    warn!(endpoint_name = %name, error = %err, "direct send failed");
509                    CamelError::ChannelClosed
510                })?;
511
512                let result = reply_rx.await.map_err(|err| {
513                    // log-policy: handler-owned
514                    // (category a: producer send failure inside the route pipeline)
515                    warn!(endpoint_name = %name, error = %err, "direct send failed");
516                    CamelError::ChannelClosed
517                })?;
518
519                debug!(endpoint_name = %name, "direct message sent");
520                result
521            })
522            .await
523            .map_err(|_| CamelError::ProcessorError(format!("direct:{name} call timed out")))?
524        })
525    }
526}
527
528// ---------------------------------------------------------------------------
529// Tests
530// ---------------------------------------------------------------------------
531
532#[cfg(test)]
533mod tests {
534    use std::sync::Mutex;
535
536    use camel_api::MetricsCollector;
537    use camel_component_api::HealthCheckRegistry;
538    use std::time::Duration;
539    fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
540        // NoOpComponentContext implements RuntimeObservability via blanket
541        // impl and returns a no-op metrics collector — avoids panicking now
542        // that direct consumer calls runtime.metrics() on send_and_wait errors.
543        std::sync::Arc::new(NoOpComponentContext)
544    }
545
546    use super::*;
547    use camel_component_api::ExchangeEnvelope;
548    use camel_component_api::Message;
549    use camel_component_api::NoOpComponentContext;
550    use camel_component_api::RuntimeObservability;
551    use std::task::RawWakerVTable;
552    use tower::ServiceExt;
553
554    // -----------------------------------------------------------------------
555    // Recording metrics collector for testing increment_errors calls
556    // -----------------------------------------------------------------------
557
558    struct RecordingMetrics {
559        errors: Arc<Mutex<Vec<(String, String)>>>,
560    }
561
562    impl MetricsCollector for RecordingMetrics {
563        fn record_exchange_duration(&self, _: &str, _: Duration) {}
564        fn increment_errors(&self, route_id: &str, error_type: &str) {
565            self.errors
566                .lock()
567                .unwrap()
568                .push((route_id.to_string(), error_type.to_string()));
569        }
570        fn increment_exchanges(&self, _: &str) {}
571        fn set_queue_depth(&self, _: &str, _: usize) {}
572        fn record_circuit_breaker_change(&self, _: &str, _: &str, _: &str) {}
573    }
574
575    struct RecordingRuntime {
576        metrics_collector: Arc<RecordingMetrics>,
577    }
578
579    impl RecordingRuntime {
580        fn new(errors: Arc<Mutex<Vec<(String, String)>>>) -> Self {
581            Self {
582                metrics_collector: Arc::new(RecordingMetrics { errors }),
583            }
584        }
585    }
586
587    impl RuntimeObservability for RecordingRuntime {
588        fn metrics(&self) -> Arc<dyn MetricsCollector> {
589            self.metrics_collector.clone() as Arc<dyn MetricsCollector>
590        }
591        fn health(&self) -> Arc<dyn HealthCheckRegistry> {
592            panic!("RecordingRuntime::health not used in this test")
593        }
594    }
595
596    fn noop_waker() -> std::task::Waker {
597        const VTABLE: RawWakerVTable = RawWakerVTable::new(|_| RAW, |_| {}, |_| {}, |_| {});
598        const RAW: std::task::RawWaker = std::task::RawWaker::new(std::ptr::null(), &VTABLE);
599        unsafe { std::task::Waker::from_raw(RAW) }
600    }
601
602    fn test_producer_ctx() -> ProducerContext {
603        ProducerContext::new()
604    }
605
606    #[test]
607    fn test_direct_component_scheme() {
608        let component = DirectComponent::new();
609        assert_eq!(component.scheme(), "direct");
610    }
611
612    #[test]
613    fn test_direct_component_default() {
614        let component = DirectComponent::default();
615        assert_eq!(component.scheme(), "direct");
616    }
617
618    #[test]
619    fn test_direct_config_from_uri() {
620        let config = DirectConfig::from_uri("direct:orders").unwrap();
621        assert_eq!(config.name, "orders");
622    }
623
624    #[test]
625    fn test_direct_endpoint_uri() {
626        let component = DirectComponent::new();
627        let endpoint = component
628            .create_endpoint("direct:uri-check", &NoOpComponentContext)
629            .unwrap();
630        assert_eq!(endpoint.uri(), "direct:uri-check");
631    }
632
633    #[test]
634    fn test_direct_creates_endpoint() {
635        let component = DirectComponent::new();
636        let endpoint = component.create_endpoint("direct:foo", &NoOpComponentContext);
637        assert!(endpoint.is_ok());
638    }
639
640    #[test]
641    fn test_direct_wrong_scheme() {
642        let component = DirectComponent::new();
643        let result = component.create_endpoint("timer:tick", &NoOpComponentContext);
644        assert!(result.is_err());
645    }
646
647    #[test]
648    fn test_direct_endpoint_creates_consumer() {
649        let component = DirectComponent::new();
650        let endpoint = component
651            .create_endpoint("direct:foo", &NoOpComponentContext)
652            .unwrap();
653        assert!(endpoint.create_consumer(rt()).is_ok());
654    }
655
656    #[test]
657    fn test_direct_endpoint_creates_producer() {
658        let ctx = test_producer_ctx();
659        let component = DirectComponent::new();
660        let endpoint = component
661            .create_endpoint("direct:foo", &NoOpComponentContext)
662            .unwrap();
663        assert!(endpoint.create_producer(rt(), &ctx).is_ok());
664    }
665
666    #[test]
667    fn test_direct_empty_name_rejected() {
668        let component = DirectComponent::new();
669        match component.create_endpoint("direct:", &NoOpComponentContext) {
670            Err(e) => assert!(
671                e.to_string().contains("must not be empty"),
672                "unexpected error: {e}"
673            ),
674            Ok(_) => panic!("expected error for empty name"),
675        }
676    }
677
678    #[tokio::test]
679    async fn test_direct_producer_no_consumer_registered() {
680        let ctx = test_producer_ctx();
681        let component = DirectComponent::new();
682        let endpoint = component
683            .create_endpoint("direct:missing", &NoOpComponentContext)
684            .unwrap();
685        let producer = endpoint.create_producer(rt(), &ctx).unwrap();
686
687        let exchange = Exchange::new(Message::new("test"));
688        let result = producer.oneshot(exchange).await;
689        assert!(result.is_err());
690    }
691
692    #[tokio::test]
693    async fn test_direct_duplicate_consumer_returns_error() {
694        let component = DirectComponent::new();
695        let endpoint = component
696            .create_endpoint("direct:dup", &NoOpComponentContext)
697            .unwrap();
698
699        let mut consumer_a = endpoint.create_consumer(rt()).unwrap();
700        let mut consumer_b = endpoint.create_consumer(rt()).unwrap();
701
702        let (route_tx_a, _route_rx_a) = mpsc::channel::<ExchangeEnvelope>(16);
703        let token_a = tokio_util::sync::CancellationToken::new();
704        let ctx_a = ConsumerContext::new(
705            route_tx_a,
706            token_a.clone(),
707            "direct-test-route-a".to_string(),
708        );
709        let handle_a = tokio::spawn(async move {
710            consumer_a.start(ctx_a).await.unwrap();
711        });
712        tokio::time::sleep(std::time::Duration::from_millis(10)).await;
713
714        let (route_tx_b, _route_rx_b) = mpsc::channel::<ExchangeEnvelope>(16);
715        let ctx_b = ConsumerContext::new(
716            route_tx_b,
717            tokio_util::sync::CancellationToken::new(),
718            "direct-test-route-b".to_string(),
719        );
720        let result = consumer_b.start(ctx_b).await;
721
722        assert!(matches!(
723            result,
724            Err(CamelError::EndpointCreationFailed(msg))
725                if msg.contains("already has a registered consumer")
726        ));
727
728        token_a.cancel();
729        handle_a.await.unwrap();
730    }
731
732    #[tokio::test]
733    async fn test_direct_producer_consumer_roundtrip() {
734        let component = DirectComponent::new();
735
736        // Create consumer endpoint and start it
737        let consumer_endpoint = component
738            .create_endpoint("direct:test", &NoOpComponentContext)
739            .unwrap();
740        let mut consumer = consumer_endpoint.create_consumer(rt()).unwrap();
741
742        // The route channel now carries ExchangeEnvelope (request-reply support).
743        let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
744        let ctx = ConsumerContext::new(
745            route_tx,
746            tokio_util::sync::CancellationToken::new(),
747            "direct-test-route".to_string(),
748        );
749
750        // Start the consumer in a background task
751        tokio::spawn(async move {
752            consumer.start(ctx).await.unwrap();
753        });
754
755        // Give the consumer a moment to register
756        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
757
758        // Spawn a pipeline simulator that reads envelopes and replies Ok.
759        tokio::spawn(async move {
760            while let Some(envelope) = route_rx.recv().await {
761                let ExchangeEnvelope { exchange, reply_tx } = envelope;
762                if let Some(tx) = reply_tx {
763                    let _ = tx.send(Ok(exchange));
764                }
765            }
766        });
767
768        // Now send an exchange via the producer
769        let ctx = test_producer_ctx();
770        let producer_endpoint = component
771            .create_endpoint("direct:test", &NoOpComponentContext)
772            .unwrap();
773        let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
774
775        let exchange = Exchange::new(Message::new("hello direct"));
776        let result = producer.oneshot(exchange).await;
777
778        assert!(result.is_ok());
779        let reply = result.unwrap();
780        assert_eq!(reply.input.body.as_text(), Some("hello direct"));
781    }
782
783    #[tokio::test]
784    async fn test_direct_propagates_error_when_no_handler() {
785        let component = DirectComponent::new();
786
787        let consumer_endpoint = component
788            .create_endpoint("direct:err-test", &NoOpComponentContext)
789            .unwrap();
790        let mut consumer = consumer_endpoint.create_consumer(rt()).unwrap();
791
792        let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
793        let ctx = ConsumerContext::new(
794            route_tx,
795            tokio_util::sync::CancellationToken::new(),
796            "direct-test-route".to_string(),
797        );
798
799        tokio::spawn(async move {
800            consumer.start(ctx).await.unwrap();
801        });
802
803        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
804
805        // Pipeline simulator that replies with Err (simulates no error handler).
806        tokio::spawn(async move {
807            while let Some(envelope) = route_rx.recv().await {
808                if let Some(tx) = envelope.reply_tx {
809                    let _ = tx.send(Err(CamelError::ProcessorError("subroute failed".into())));
810                }
811            }
812        });
813
814        let ctx = test_producer_ctx();
815        let producer_endpoint = component
816            .create_endpoint("direct:err-test", &NoOpComponentContext)
817            .unwrap();
818        let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
819
820        let exchange = Exchange::new(Message::new("test"));
821        let result = producer.oneshot(exchange).await;
822        assert!(result.is_err());
823        assert!(matches!(result.unwrap_err(), CamelError::ProcessorError(_)));
824    }
825
826    #[tokio::test]
827    async fn test_direct_consumer_stop_unregisters() {
828        let component = DirectComponent::new();
829        let endpoint = component
830            .create_endpoint("direct:cleanup", &NoOpComponentContext)
831            .unwrap();
832
833        // We need a consumer to register
834        let mut consumer = endpoint.create_consumer(rt()).unwrap();
835
836        let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
837        let ctx = ConsumerContext::new(
838            route_tx,
839            tokio_util::sync::CancellationToken::new(),
840            "direct-test-route".to_string(),
841        );
842
843        // Start consumer in background
844        let handle = tokio::spawn(async move {
845            consumer.start(ctx).await.unwrap();
846        });
847
848        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
849
850        // Verify the name is registered
851        {
852            let reg = component.registry.lock().unwrap_or_else(|e| e.into_inner());
853            assert!(reg.contains_key("cleanup"));
854        }
855
856        // Create a new consumer just to call stop (stop removes from registry)
857        let mut stop_consumer = DirectConsumer {
858            name: "cleanup".to_string(),
859            registry: Arc::clone(&component.registry),
860            cancel: None,
861            runtime: rt(),
862        };
863        stop_consumer.stop().await.unwrap();
864
865        // Verify removed from registry
866        {
867            let reg = component.registry.lock().unwrap_or_else(|e| e.into_inner());
868            assert!(!reg.contains_key("cleanup"));
869        }
870
871        handle.abort();
872    }
873
874    #[tokio::test]
875    async fn test_direct_consumer_respects_cancellation() {
876        use tokio_util::sync::CancellationToken;
877
878        let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
879        let token = CancellationToken::new();
880        let (tx, _rx) = mpsc::channel(16);
881        let ctx = ConsumerContext::new(tx, token.clone(), "direct-test-route".to_string());
882
883        let mut consumer = DirectConsumer {
884            name: "cancel-test".to_string(),
885            registry: registry.clone(),
886            cancel: None,
887            runtime: rt(),
888        };
889
890        let handle = tokio::spawn(async move {
891            consumer.start(ctx).await.unwrap();
892        });
893
894        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
895        assert!(
896            registry
897                .lock()
898                .unwrap_or_else(|e| e.into_inner())
899                .contains_key("cancel-test")
900        );
901
902        token.cancel();
903        // start() now spawns the loop internally and returns immediately,
904        // so the outer handle completes right away. Give the inner task time
905        // to react to the cancellation and clean up the registry.
906        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
907
908        // After cancellation, the consumer should have cleaned up the registry
909        assert!(
910            !registry
911                .lock()
912                .unwrap_or_else(|e| e.into_inner())
913                .contains_key("cancel-test")
914        );
915
916        let _ = handle.await;
917    }
918
919    #[tokio::test]
920    async fn test_direct_consumer_stop_missing_entry_is_ok() {
921        let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
922        let mut consumer = DirectConsumer {
923            name: "never-registered".to_string(),
924            registry,
925            cancel: None,
926            runtime: rt(),
927        };
928        let result = consumer.stop().await;
929        assert!(result.is_ok());
930    }
931
932    #[test]
933    fn test_poll_ready_endpoint_not_registered() {
934        let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
935        let producer = DirectProducer {
936            name: "missing".to_string(),
937            registry,
938            config: DirectConfig {
939                name: "missing".to_string(),
940                timeout_ms: None,
941                block: None,
942                fail_if_no_consumers: None,
943                exchange_pattern: None,
944            },
945            semaphore: Arc::new(Semaphore::new(1)),
946            pending_permit: None,
947            acquire_fut: None,
948            fail_if_no_consumers: None,
949        };
950        let waker = noop_waker();
951        let mut cx = Context::from_waker(&waker);
952        let mut producer = producer;
953        let result = producer.poll_ready(&mut cx);
954        assert!(matches!(
955            result,
956            Poll::Ready(Err(CamelError::EndpointCreationFailed(_)))
957        ));
958    }
959
960    #[test]
961    fn test_poll_ready_endpoint_registered() {
962        let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
963        let (tx, _rx) =
964            mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(1);
965        registry.lock().unwrap().insert("active".to_string(), tx);
966        let producer = DirectProducer {
967            name: "active".to_string(),
968            registry,
969            config: DirectConfig {
970                name: "active".to_string(),
971                timeout_ms: None,
972                block: None,
973                fail_if_no_consumers: None,
974                exchange_pattern: None,
975            },
976            semaphore: Arc::new(Semaphore::new(1)),
977            pending_permit: None,
978            acquire_fut: None,
979            fail_if_no_consumers: None,
980        };
981        let waker = noop_waker();
982        let mut cx = Context::from_waker(&waker);
983        let mut producer = producer;
984        let result = producer.poll_ready(&mut cx);
985        assert!(matches!(result, Poll::Ready(Ok(()))));
986    }
987
988    #[test]
989    fn test_poll_ready_allows_missing_consumer_when_fail_if_no_consumers_false() {
990        let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
991        let producer = DirectProducer {
992            name: "missing-ok".to_string(),
993            registry,
994            config: DirectConfig {
995                name: "missing-ok".to_string(),
996                timeout_ms: None,
997                block: None,
998                fail_if_no_consumers: Some(false),
999                exchange_pattern: None,
1000            },
1001            semaphore: Arc::new(Semaphore::new(1)),
1002            pending_permit: None,
1003            acquire_fut: None,
1004            fail_if_no_consumers: Some(false),
1005        };
1006
1007        let waker = noop_waker();
1008        let mut cx = Context::from_waker(&waker);
1009        let mut producer = producer;
1010        let result = producer.poll_ready(&mut cx);
1011        assert!(matches!(result, Poll::Ready(Ok(()))));
1012    }
1013
1014    #[test]
1015    fn test_poll_ready_channel_closed() {
1016        let registry: DirectRegistry = Arc::new(Mutex::new(HashMap::new()));
1017        let (tx, rx) =
1018            mpsc::channel::<(Exchange, oneshot::Sender<Result<Exchange, CamelError>>)>(1);
1019        drop(rx);
1020        registry.lock().unwrap().insert("closed".to_string(), tx);
1021        let producer = DirectProducer {
1022            name: "closed".to_string(),
1023            registry,
1024            config: DirectConfig {
1025                name: "closed".to_string(),
1026                timeout_ms: None,
1027                block: None,
1028                fail_if_no_consumers: None,
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: None,
1035        };
1036        let waker = noop_waker();
1037        let mut cx = Context::from_waker(&waker);
1038        let mut producer = producer;
1039        let result = producer.poll_ready(&mut cx);
1040        assert!(matches!(
1041            result,
1042            Poll::Ready(Err(CamelError::EndpointCreationFailed(_)))
1043        ));
1044    }
1045
1046    #[tokio::test]
1047    async fn test_direct_stop_cancels_loop() {
1048        use tokio_util::sync::CancellationToken;
1049
1050        let component = DirectComponent::new();
1051        let endpoint = component
1052            .create_endpoint("direct:stop-test", &NoOpComponentContext)
1053            .unwrap();
1054        let mut consumer = endpoint.create_consumer(rt()).unwrap();
1055
1056        let token = CancellationToken::new();
1057        let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1058        let ctx = ConsumerContext::new(route_tx, token.clone(), "direct-test-route".to_string());
1059
1060        // start() runs the consumer loop inline on the managed consumer task.
1061        // We test stop() cancels the loop.
1062        let handle = tokio::spawn(async move {
1063            consumer.start(ctx).await.unwrap();
1064        });
1065
1066        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1067        assert!(
1068            component
1069                .registry
1070                .lock()
1071                .unwrap_or_else(|e| e.into_inner())
1072                .contains_key("stop-test")
1073        );
1074
1075        // Create a new consumer just for stop
1076        let mut stop_consumer = DirectConsumer {
1077            name: "stop-test".to_string(),
1078            registry: Arc::clone(&component.registry),
1079            cancel: Some(token.clone()),
1080            runtime: rt(),
1081        };
1082        stop_consumer.stop().await.unwrap();
1083
1084        // The consumer loop should finish within 2s after stop
1085        let result = tokio::time::timeout(std::time::Duration::from_secs(2), handle).await;
1086        assert!(result.is_ok(), "Consumer loop did not stop within 2s");
1087
1088        // Registry should be cleaned up
1089        assert!(
1090            !component
1091                .registry
1092                .lock()
1093                .unwrap_or_else(|e| e.into_inner())
1094                .contains_key("stop-test")
1095        );
1096    }
1097
1098    #[tokio::test]
1099    async fn test_direct_producer_timeout() {
1100        let component = DirectComponent::new();
1101        let endpoint = component
1102            .create_endpoint("direct:timeout-test", &NoOpComponentContext)
1103            .unwrap();
1104        let mut consumer = endpoint.create_consumer(rt()).unwrap();
1105
1106        // Consumer that never replies (simulates a stuck pipeline)
1107        let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1108        let token = tokio_util::sync::CancellationToken::new();
1109        let ctx = ConsumerContext::new(route_tx, token.clone(), "direct-test-route".to_string());
1110        tokio::spawn(async move {
1111            consumer.start(ctx).await.unwrap();
1112        });
1113
1114        // Drain envelopes but hold onto reply_tx so the producer never gets a reply
1115        tokio::spawn(async move {
1116            let mut held_reply: Vec<oneshot::Sender<Result<Exchange, CamelError>>> = Vec::new();
1117            while let Some(envelope) = route_rx.recv().await {
1118                held_reply.push(envelope.reply_tx.unwrap());
1119            }
1120            drop(held_reply);
1121        });
1122
1123        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1124
1125        // Create producer with a short timeout
1126        let _ = test_producer_ctx();
1127        let _producer_endpoint = component
1128            .create_endpoint("direct:timeout-test", &NoOpComponentContext)
1129            .unwrap();
1130        let producer = DirectProducer {
1131            name: "timeout-test".to_string(),
1132            registry: Arc::clone(&component.registry),
1133            config: DirectConfig {
1134                name: "timeout-test".to_string(),
1135                timeout_ms: Some(100), // 100ms timeout
1136                block: None,
1137                fail_if_no_consumers: None,
1138                exchange_pattern: None,
1139            },
1140            semaphore: Arc::new(Semaphore::new(1)),
1141            pending_permit: None,
1142            acquire_fut: None,
1143            fail_if_no_consumers: None,
1144        };
1145
1146        let exchange = Exchange::new(Message::new("test"));
1147        let mut svc = producer;
1148        let _ = svc.poll_ready(&mut Context::from_waker(&noop_waker()));
1149        let result = svc.call(exchange).await;
1150        assert!(result.is_err(), "Expected timeout error");
1151        assert!(
1152            result.unwrap_err().to_string().contains("timed out"),
1153            "Expected timeout message"
1154        );
1155
1156        token.cancel();
1157    }
1158
1159    #[tokio::test]
1160    async fn test_send_and_wait_error_increments_errors_metric() {
1161        let errors: Arc<Mutex<Vec<(String, String)>>> = Arc::new(Mutex::new(Vec::new()));
1162        let runtime = Arc::new(RecordingRuntime::new(Arc::clone(&errors)));
1163
1164        let component = DirectComponent::new();
1165        let endpoint = component
1166            .create_endpoint("direct:metrics-error", &NoOpComponentContext)
1167            .unwrap();
1168        let mut consumer = endpoint.create_consumer(runtime).unwrap();
1169
1170        // Route that returns Err — simulates unhandled pipeline failure
1171        let cancel = tokio_util::sync::CancellationToken::new();
1172        let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1173        let ctx = ConsumerContext::new(route_tx, cancel.clone(), "test-route-id".to_string());
1174
1175        tokio::spawn(async move {
1176            consumer.start(ctx).await.unwrap();
1177        });
1178
1179        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1180
1181        // Route pipeline: reply with Err for every incoming exchange
1182        tokio::spawn(async move {
1183            while let Some(envelope) = route_rx.recv().await {
1184                if let Some(tx) = envelope.reply_tx {
1185                    let _ = tx.send(Err(CamelError::ProcessorError(
1186                        "pipeline failure".to_string(),
1187                    )));
1188                }
1189            }
1190        });
1191
1192        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1193
1194        // Send an exchange via the producer so the consumer's send_and_wait
1195        // returns Err, triggering the metrics call.
1196        let ctx = test_producer_ctx();
1197        let producer_endpoint = component
1198            .create_endpoint("direct:metrics-error", &NoOpComponentContext)
1199            .unwrap();
1200        let producer = producer_endpoint.create_producer(rt(), &ctx).unwrap();
1201
1202        let exchange = Exchange::new(Message::new("test"));
1203        // The producer oneshot wraps the call — it returns the Err from
1204        // send_and_wait back to us.
1205        let _result = producer.oneshot(exchange).await;
1206
1207        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1208
1209        cancel.cancel();
1210
1211        // Verify MetricsCollector::increment_errors was called
1212        let recorded = errors.lock().unwrap();
1213        assert_eq!(
1214            recorded.len(),
1215            1,
1216            "expected 1 increment_errors call, got {}: {:?}",
1217            recorded.len(),
1218            *recorded
1219        );
1220        assert_eq!(recorded[0].0, "test-route-id");
1221        assert_eq!(recorded[0].1, "b-prime:direct:send-and-wait");
1222    }
1223
1224    #[test]
1225    fn test_empty_endpoint_name_rejected() {
1226        let result = DirectConfig::from_uri("direct:");
1227        // from_uri may parse empty name — validate catches it
1228        if let Ok(config) = result {
1229            assert!(
1230                validate_name(&config.name).is_err(),
1231                "expected validation error for empty name"
1232            );
1233        }
1234        // Also verify via Component (the main entry point)
1235        let component = DirectComponent::new();
1236        let result = component.create_endpoint("direct:", &NoOpComponentContext);
1237        assert!(result.is_err(), "empty endpoint name must be rejected");
1238    }
1239
1240    #[test]
1241    fn test_whitespace_endpoint_name_rejected() {
1242        let result = DirectConfig::from_uri("direct:my endpoint");
1243        if let Ok(config) = result {
1244            assert!(
1245                validate_name(&config.name).is_err(),
1246                "expected validation error for whitespace in name"
1247            );
1248        }
1249        let component = DirectComponent::new();
1250        let result = component.create_endpoint("direct:my endpoint", &NoOpComponentContext);
1251        assert!(result.is_err(), "whitespace endpoint name must be rejected");
1252    }
1253
1254    #[test]
1255    fn test_valid_endpoint_name_accepted() {
1256        let component = DirectComponent::new();
1257        let result = component.create_endpoint("direct:my-endpoint", &NoOpComponentContext);
1258        assert!(result.is_ok(), "valid endpoint name should be accepted");
1259    }
1260}