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