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