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