Skip to main content

camel_component_seda/
lib.rs

1//! In-memory SEDA component for rust-camel — asynchronous staging channel
2//! between routes sharing the same context via bounded queues.
3//!
4//! Main types: `SedaComponent`, `SedaEndpoint`, `SedaConsumer`, `SedaProducer`.
5
6use std::collections::HashMap;
7use std::future::Future;
8use std::pin::Pin;
9use std::sync::atomic::{AtomicU64, Ordering};
10use std::sync::{Arc, Mutex};
11use std::task::{Context, Poll};
12use std::time::Duration;
13
14#[cfg(test)]
15use camel_component_api::test_support::PanicRuntimeObservability;
16#[cfg(test)]
17fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
18    std::sync::Arc::new(PanicRuntimeObservability)
19}
20
21use async_trait::async_trait;
22use tokio::sync::{mpsc, oneshot};
23use tokio::task::JoinHandle;
24use tokio_util::sync::CancellationToken;
25use tower::Service;
26
27use camel_api::BoxProcessorExt;
28use camel_api::component_metadata::{
29    ComponentCapabilities, ComponentMetadata, OptionKind, UriOption,
30};
31use camel_component_api::parse_uri;
32use camel_component_api::{
33    BoxProcessor, CamelError, Component, ComponentContext, ConcurrencyModel, Consumer,
34    ConsumerContext, Endpoint, Exchange, ExchangeEnvelope, ProducerContext,
35};
36use tracing::{info, warn};
37
38// ---------------------------------------------------------------------------
39// Enums
40// ---------------------------------------------------------------------------
41
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43pub enum WaitForTaskToComplete {
44    Never,
45    IfReplyExpected,
46    Always,
47}
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50pub enum ExchangePattern {
51    InOnly,
52    InOut,
53}
54
55// ---------------------------------------------------------------------------
56// SedaConfig
57// ---------------------------------------------------------------------------
58
59/// Configuration parsed from a SEDA URI.
60///
61/// URI format: `seda:name[?options]`
62///
63/// Options are split into two groups:
64/// - **shared**: validated for consistency when multiple endpoints reference
65///   the same endpoint name (`size`, `multiple_consumers`, `exchange_pattern`,
66///   `concurrent_consumers`).
67/// - **producer only**: stored per-endpoint, used only by the producer
68///   (`block_when_full`, `discard_if_no_consumers`, `timeout_ms`,
69///   `wait_for_task_to_complete`).
70#[derive(Debug, Clone)]
71pub struct SedaConfig {
72    pub name: String,
73    pub size: usize,
74    pub concurrent_consumers: usize,
75    pub multiple_consumers: bool,
76    pub block_when_full: bool,
77    pub discard_if_no_consumers: bool,
78    pub timeout_ms: u64,
79    pub wait_for_task_to_complete: WaitForTaskToComplete,
80    pub exchange_pattern: ExchangePattern,
81}
82
83impl SedaConfig {
84    pub fn from_uri(uri: &str) -> Result<Self, CamelError> {
85        let parts = parse_uri(uri)?;
86        if parts.scheme != "seda" {
87            return Err(CamelError::InvalidUri(format!(
88                "invalid scheme '{}', expected 'seda'",
89                parts.scheme
90            )));
91        }
92
93        let name = parts.path;
94        if name.trim().is_empty() {
95            return Err(CamelError::InvalidUri(
96                "seda: endpoint name must not be empty".to_string(),
97            ));
98        }
99        if name.contains(char::is_whitespace) {
100            return Err(CamelError::InvalidUri(
101                "seda: endpoint name must not contain whitespace".to_string(),
102            ));
103        }
104
105        let size: usize = parts
106            .params
107            .get("size")
108            .map(|v| v.parse::<usize>())
109            .transpose()
110            .map_err(|e: std::num::ParseIntError| {
111                CamelError::InvalidUri(format!("invalid size: {e}"))
112            })?
113            .unwrap_or(1000);
114
115        if size == 0 {
116            return Err(CamelError::InvalidUri(
117                "seda: size must be greater than 0".to_string(),
118            ));
119        }
120
121        let concurrent_consumers: usize = parts
122            .params
123            .get("concurrentConsumers")
124            .map(|v| v.parse::<usize>())
125            .transpose()
126            .map_err(|e: std::num::ParseIntError| {
127                CamelError::InvalidUri(format!("invalid concurrentConsumers: {e}"))
128            })?
129            .unwrap_or(1);
130
131        let multiple_consumers = parts
132            .params
133            .get("multipleConsumers")
134            .map(|v| parse_bool("multipleConsumers", v))
135            .transpose()?
136            .unwrap_or(false);
137
138        let block_when_full = parts
139            .params
140            .get("blockWhenFull")
141            .map(|v| parse_bool("blockWhenFull", v))
142            .transpose()?
143            .unwrap_or(false);
144
145        let discard_if_no_consumers = parts
146            .params
147            .get("discardIfNoConsumers")
148            .map(|v| parse_bool("discardIfNoConsumers", v))
149            .transpose()?
150            .unwrap_or(false);
151
152        let timeout_ms: u64 = parts
153            .params
154            .get("timeout")
155            .map(|v| v.parse::<u64>())
156            .transpose()
157            .map_err(|e: std::num::ParseIntError| {
158                CamelError::InvalidUri(format!("invalid timeout: {e}"))
159            })?
160            .unwrap_or(30_000);
161
162        let wait_for_task_to_complete = parts
163            .params
164            .get("waitForTaskToComplete")
165            .map(|v| parse_wait_for_task(v))
166            .transpose()?
167            .unwrap_or(WaitForTaskToComplete::IfReplyExpected);
168
169        let exchange_pattern = parts
170            .params
171            .get("exchangePattern")
172            .map(|v| parse_exchange_pattern(v))
173            .transpose()?
174            .unwrap_or(ExchangePattern::InOnly);
175
176        let concurrent_consumers = if concurrent_consumers == 0 {
177            warn!(name, "concurrentConsumers=0 clamped to 1");
178            1
179        } else {
180            concurrent_consumers
181        };
182
183        Ok(Self {
184            name,
185            size,
186            concurrent_consumers,
187            multiple_consumers,
188            block_when_full,
189            discard_if_no_consumers,
190            timeout_ms,
191            wait_for_task_to_complete,
192            exchange_pattern,
193        })
194    }
195
196    fn is_compatible_with(&self, other: &SedaConfig) -> Result<(), String> {
197        let mut diffs = Vec::new();
198        if self.size != other.size {
199            diffs.push(format!("size: {} vs {}", self.size, other.size));
200        }
201        if self.multiple_consumers != other.multiple_consumers {
202            diffs.push(format!(
203                "multipleConsumers: {} vs {}",
204                self.multiple_consumers, other.multiple_consumers
205            ));
206        }
207        if self.exchange_pattern != other.exchange_pattern {
208            diffs.push(format!(
209                "exchangePattern: {:?} vs {:?}",
210                self.exchange_pattern, other.exchange_pattern
211            ));
212        }
213        if self.concurrent_consumers != other.concurrent_consumers {
214            diffs.push(format!(
215                "concurrentConsumers: {} vs {}",
216                self.concurrent_consumers, other.concurrent_consumers
217            ));
218        }
219        if diffs.is_empty() {
220            Ok(())
221        } else {
222            Err(format!(
223                "endpoint '{}' already exists with different config: {}",
224                self.name,
225                diffs.join(", ")
226            ))
227        }
228    }
229}
230
231fn parse_bool(name: &str, v: &str) -> Result<bool, CamelError> {
232    match v.to_lowercase().as_str() {
233        "true" | "1" | "yes" => Ok(true),
234        "false" | "0" | "no" => Ok(false),
235        _ => Err(CamelError::InvalidUri(format!(
236            "invalid boolean for {name}: '{v}'"
237        ))),
238    }
239}
240
241fn parse_wait_for_task(v: &str) -> Result<WaitForTaskToComplete, CamelError> {
242    match v.to_lowercase().replace('_', "").as_str() {
243        "never" => Ok(WaitForTaskToComplete::Never),
244        "ifreplyexpected" => Ok(WaitForTaskToComplete::IfReplyExpected),
245        "always" => Ok(WaitForTaskToComplete::Always),
246        _ => Err(CamelError::InvalidUri(format!(
247            "invalid waitForTaskToComplete: '{v}' (expected: Never, IfReplyExpected, Always)"
248        ))),
249    }
250}
251
252fn parse_exchange_pattern(v: &str) -> Result<ExchangePattern, CamelError> {
253    match v.to_lowercase().replace('_', "").as_str() {
254        "inonly" => Ok(ExchangePattern::InOnly),
255        "inout" => Ok(ExchangePattern::InOut),
256        _ => Err(CamelError::InvalidUri(format!(
257            "invalid exchangePattern: '{v}' (expected: InOnly, InOut)"
258        ))),
259    }
260}
261
262// ---------------------------------------------------------------------------
263// ConsumerId generator (no uuid dependency needed)
264// ---------------------------------------------------------------------------
265
266static CONSUMER_ID_COUNTER: AtomicU64 = AtomicU64::new(1);
267
268fn next_consumer_id() -> String {
269    format!(
270        "seda-consumer-{}",
271        CONSUMER_ID_COUNTER.fetch_add(1, Ordering::Relaxed)
272    )
273}
274
275// ---------------------------------------------------------------------------
276// SedaMode + SedaEndpointState
277// ---------------------------------------------------------------------------
278
279type ConsumerId = String;
280
281/// Transport mode for a SEDA endpoint.
282///
283/// - `Single`: one bounded mpsc channel, one consumer allowed.
284///   `active` tracks whether a consumer has started (separate from receiver
285///   ownership, which is taken by the forwarder task on start).
286/// - `Fanout`: one bounded mpsc per subscriber, multiple consumers allowed.
287enum SedaMode {
288    Single {
289        tx: mpsc::Sender<ExchangeEnvelope>,
290        rx: Mutex<Option<mpsc::Receiver<ExchangeEnvelope>>>,
291        active: std::sync::atomic::AtomicBool,
292    },
293    Fanout {
294        subscribers: Mutex<HashMap<ConsumerId, mpsc::Sender<ExchangeEnvelope>>>,
295    },
296}
297
298struct SedaEndpointState {
299    config: SedaConfig,
300    mode: SedaMode,
301}
302
303impl SedaEndpointState {
304    fn new(config: &SedaConfig) -> Self {
305        let (tx, rx) = mpsc::channel(config.size);
306        let mode = if config.multiple_consumers {
307            SedaMode::Fanout {
308                subscribers: Mutex::new(HashMap::new()),
309            }
310        } else {
311            SedaMode::Single {
312                tx,
313                rx: Mutex::new(Some(rx)),
314                active: std::sync::atomic::AtomicBool::new(false),
315            }
316        };
317        Self {
318            config: config.clone(),
319            mode,
320        }
321    }
322
323    /// Returns true if at least one consumer has started and not yet stopped.
324    /// For Single mode: checks the `active` flag (not the receiver, which is
325    /// moved into the forwarder task on start).
326    /// For Fanout mode: checks if subscribers map is non-empty.
327    fn has_active_consumers(&self) -> bool {
328        match &self.mode {
329            SedaMode::Single { active, .. } => active.load(Ordering::SeqCst),
330            SedaMode::Fanout { subscribers } => !subscribers
331                .lock()
332                .unwrap_or_else(|e| e.into_inner())
333                .is_empty(),
334        }
335    }
336}
337
338// ---------------------------------------------------------------------------
339// SedaComponent
340// ---------------------------------------------------------------------------
341
342type SedaRegistry = Arc<Mutex<HashMap<String, Arc<SedaEndpointState>>>>;
343
344pub struct SedaComponent {
345    endpoints: SedaRegistry,
346}
347
348impl SedaComponent {
349    pub fn new() -> Self {
350        Self {
351            endpoints: Arc::new(Mutex::new(HashMap::new())),
352        }
353    }
354
355    fn get_or_create_state(
356        &self,
357        config: &SedaConfig,
358    ) -> Result<Arc<SedaEndpointState>, CamelError> {
359        let mut endpoints = self.endpoints.lock().unwrap_or_else(|e| e.into_inner());
360        if let Some(existing) = endpoints.get(&config.name) {
361            existing
362                .config
363                .is_compatible_with(config)
364                .map_err(CamelError::EndpointCreationFailed)?;
365            Ok(Arc::clone(existing))
366        } else {
367            let state = Arc::new(SedaEndpointState::new(config));
368            endpoints.insert(config.name.clone(), Arc::clone(&state));
369            Ok(state)
370        }
371    }
372}
373
374impl Default for SedaComponent {
375    fn default() -> Self {
376        Self::new()
377    }
378}
379
380#[async_trait]
381impl Component for SedaComponent {
382    fn scheme(&self) -> &str {
383        "seda"
384    }
385
386    fn metadata(&self) -> ComponentMetadata {
387        ComponentMetadata {
388            scheme: "seda".to_string(),
389            version: env!("CARGO_PKG_VERSION").to_string(),
390            description: "Asynchronous staged event-driven architecture with bounded queue"
391                .to_string(),
392            uri_syntax: "seda:name?size=1000&concurrentConsumers=1".to_string(),
393            capabilities: ComponentCapabilities {
394                supports_consumer: true,
395                supports_producer: true,
396                ..Default::default()
397            },
398            uri_options: vec![
399                UriOption::new(
400                    "size",
401                    "Bounded queue capacity. Must be > 0",
402                    OptionKind::Int,
403                )
404                .with_default("1000"),
405                UriOption::new(
406                    "concurrentConsumers",
407                    "Consumer concurrency. Clamped to 1 minimum",
408                    OptionKind::Int,
409                )
410                .with_default("1"),
411                UriOption::new(
412                    "multipleConsumers",
413                    "Fanout mode — clone to all subscribers",
414                    OptionKind::Bool,
415                )
416                .with_default("false"),
417                UriOption::new(
418                    "blockWhenFull",
419                    "Block producer when queue full vs fail fast",
420                    OptionKind::Bool,
421                )
422                .with_default("false"),
423                UriOption::new(
424                    "discardIfNoConsumers",
425                    "Silently drop if no consumers vs error",
426                    OptionKind::Bool,
427                )
428                .with_default("false"),
429                UriOption::new(
430                    "timeout",
431                    "Timeout for enqueue and reply wait in milliseconds",
432                    OptionKind::Int,
433                )
434                .with_default("30000"),
435                UriOption::new(
436                    "waitForTaskToComplete",
437                    "When to wait for task completion",
438                    OptionKind::Enum(vec![
439                        "Never".to_string(),
440                        "IfReplyExpected".to_string(),
441                        "Always".to_string(),
442                    ]),
443                )
444                .with_default("IfReplyExpected"),
445                UriOption::new(
446                    "exchangePattern",
447                    "Exchange pattern",
448                    OptionKind::Enum(vec!["InOnly".to_string(), "InOut".to_string()]),
449                )
450                .with_default("InOnly"),
451            ],
452            ..ComponentMetadata::minimal("seda")
453        }
454    }
455
456    fn create_endpoint(
457        &self,
458        uri: &str,
459        _ctx: &dyn ComponentContext,
460    ) -> Result<Box<dyn Endpoint>, CamelError> {
461        let config = SedaConfig::from_uri(uri)?;
462        let state = self.get_or_create_state(&config)?;
463        Ok(Box::new(SedaEndpoint {
464            uri: uri.to_string(),
465            config,
466            state,
467        }))
468    }
469}
470
471// ---------------------------------------------------------------------------
472// SedaEndpoint
473// ---------------------------------------------------------------------------
474
475struct SedaEndpoint {
476    uri: String,
477    config: SedaConfig,
478    state: Arc<SedaEndpointState>,
479}
480
481impl Endpoint for SedaEndpoint {
482    fn uri(&self) -> &str {
483        &self.uri
484    }
485
486    fn create_consumer(
487        &self,
488        rt: Arc<dyn camel_component_api::RuntimeObservability>,
489    ) -> Result<Box<dyn Consumer>, CamelError> {
490        Ok(Box::new(SedaConsumer::new(
491            Arc::clone(&self.state),
492            next_consumer_id(),
493            rt,
494        )))
495    }
496
497    fn create_producer(
498        &self,
499        rt: Arc<dyn camel_component_api::RuntimeObservability>,
500        _ctx: &ProducerContext,
501    ) -> Result<BoxProcessor, CamelError> {
502        let producer = SedaProducer {
503            state: Arc::clone(&self.state),
504            producer_config: ProducerConfig::from(&self.config),
505            runtime: rt,
506        };
507        Ok(BoxProcessor::from_fn(move |ex| {
508            let mut svc = producer.clone();
509            Box::pin(async move { svc.call(ex).await })
510        }))
511    }
512}
513
514/// Per-endpoint producer options. These are NOT shared at the SedaEndpointState
515/// level because two endpoints referencing the same seda name may have different
516/// producer-only options (e.g. different blockWhenFull settings).
517#[derive(Clone)]
518struct ProducerConfig {
519    block_when_full: bool,
520    discard_if_no_consumers: bool,
521    timeout_ms: u64,
522    wait_for_task_to_complete: WaitForTaskToComplete,
523}
524
525impl From<&SedaConfig> for ProducerConfig {
526    fn from(config: &SedaConfig) -> Self {
527        Self {
528            block_when_full: config.block_when_full,
529            discard_if_no_consumers: config.discard_if_no_consumers,
530            timeout_ms: config.timeout_ms,
531            wait_for_task_to_complete: config.wait_for_task_to_complete,
532        }
533    }
534}
535
536// ---------------------------------------------------------------------------
537// SedaConsumer
538// ---------------------------------------------------------------------------
539
540struct SedaConsumer {
541    state: Arc<SedaEndpointState>,
542    consumer_id: ConsumerId,
543    started: bool,
544    cancel_token: CancellationToken,
545    forwarder_handles: Vec<JoinHandle<Result<(), CamelError>>>,
546    /// Phase B will use this for `rt.metrics().increment_errors(...)` and
547    /// `rt.health().force_unhealthy_for_route(...)` calls per ADR-0012.
548    #[allow(dead_code)]
549    runtime: Arc<dyn camel_component_api::RuntimeObservability>,
550}
551
552impl SedaConsumer {
553    fn new(
554        state: Arc<SedaEndpointState>,
555        consumer_id: ConsumerId,
556        runtime: Arc<dyn camel_component_api::RuntimeObservability>,
557    ) -> Self {
558        Self {
559            state,
560            consumer_id,
561            started: false,
562            cancel_token: CancellationToken::new(),
563            forwarder_handles: Vec::new(),
564            runtime,
565        }
566    }
567}
568
569#[async_trait]
570impl Consumer for SedaConsumer {
571    async fn start(&mut self, ctx: ConsumerContext) -> Result<(), CamelError> {
572        if self.started {
573            return Err(CamelError::EndpointCreationFailed(
574                "consumer already started".to_string(),
575            ));
576        }
577
578        match &self.state.mode {
579            SedaMode::Single { rx, active, .. } => {
580                let mut rx_guard = rx.lock().unwrap_or_else(|e| e.into_inner());
581                if rx_guard.is_none() {
582                    return Err(CamelError::EndpointCreationFailed(format!(
583                        "endpoint '{}' already has a registered consumer",
584                        self.state.config.name
585                    )));
586                }
587                active.store(true, Ordering::SeqCst);
588                let receiver = rx_guard.take().ok_or_else(|| {
589                    CamelError::EndpointCreationFailed(format!(
590                        "endpoint '{}' receiver already taken",
591                        self.state.config.name
592                    ))
593                })?;
594                drop(rx_guard);
595
596                let cancel = self.cancel_token.clone();
597                let handle = tokio::spawn(async move {
598                    let mut rx = receiver;
599                    loop {
600                        tokio::select! {
601                            envelope = rx.recv() => {
602                                let Some(envelope) = envelope else { break };
603                                forward_envelope(&ctx, envelope).await;
604                            }
605                            _ = cancel.cancelled() => break,
606                        }
607                    }
608                    Ok(())
609                });
610                self.forwarder_handles.push(handle);
611            }
612            SedaMode::Fanout { subscribers } => {
613                let (tx, rx) = mpsc::channel(self.state.config.size);
614                subscribers
615                    .lock()
616                    .unwrap_or_else(|e| e.into_inner())
617                    .insert(self.consumer_id.clone(), tx);
618
619                let cancel = self.cancel_token.clone();
620                let handle = tokio::spawn(async move {
621                    let mut rx = rx;
622                    loop {
623                        tokio::select! {
624                            envelope = rx.recv() => {
625                                let Some(envelope) = envelope else { break };
626                                forward_envelope(&ctx, envelope).await;
627                            }
628                            _ = cancel.cancelled() => break,
629                        }
630                    }
631                    Ok(())
632                });
633                self.forwarder_handles.push(handle);
634            }
635        }
636
637        self.started = true;
638        info!(
639            name = %self.state.config.name,
640            consumer_id = %self.consumer_id,
641            concurrent = self.state.config.concurrent_consumers,
642            "SEDA consumer started"
643        );
644        Ok(())
645    }
646
647    async fn stop(&mut self) -> Result<(), CamelError> {
648        if !self.started {
649            return Ok(());
650        }
651        self.cancel_token.cancel();
652        for handle in self.forwarder_handles.drain(..) {
653            handle.abort();
654        }
655        match &self.state.mode {
656            SedaMode::Single { active, .. } => {
657                active.store(false, Ordering::SeqCst);
658            }
659            SedaMode::Fanout { subscribers } => {
660                subscribers
661                    .lock()
662                    .unwrap_or_else(|e| e.into_inner())
663                    .remove(&self.consumer_id);
664            }
665        }
666        self.started = false;
667        info!(
668            name = %self.state.config.name,
669            consumer_id = %self.consumer_id,
670            "SEDA consumer stopped"
671        );
672        Ok(())
673    }
674
675    fn concurrency_model(&self) -> ConcurrencyModel {
676        ConcurrencyModel::Concurrent {
677            max: Some(self.state.config.concurrent_consumers),
678        }
679    }
680
681    fn background_task_handle(
682        &mut self,
683    ) -> Option<tokio::task::JoinHandle<Result<(), CamelError>>> {
684        // SEDA may have multiple forwarder handles; return the first one.
685        // The remaining handles are cancelled in stop().
686        self.forwarder_handles.pop()
687    }
688}
689
690/// Forward an envelope from the SEDA queue into the route pipeline.
691///
692/// Key rule: if the envelope carries a `reply_tx`, the forwarder MUST use
693/// `send_and_wait()` to route the pipeline result back to the producer.
694/// This handles both InOut and `waitForTaskToComplete=Always` cases.
695/// If no `reply_tx`, use fire-and-forget `send()`.
696async fn forward_envelope(ctx: &ConsumerContext, envelope: ExchangeEnvelope) {
697    if let Some(reply_tx) = envelope.reply_tx {
698        let result = ctx.send_and_wait(envelope.exchange).await;
699        let _ = reply_tx.send(result);
700    } else {
701        if let Err(e) = ctx.send(envelope.exchange).await {
702            warn!(error = %e, "SEDA consumer send failed");
703        }
704    }
705}
706
707// ---------------------------------------------------------------------------
708// SedaProducer
709// ---------------------------------------------------------------------------
710
711#[derive(Clone)]
712struct SedaProducer {
713    state: Arc<SedaEndpointState>,
714    producer_config: ProducerConfig,
715    /// Phase B will use this for `rt.metrics().increment_errors(...)` and
716    /// `rt.health().force_unhealthy_for_route(...)` calls per ADR-0012.
717    #[allow(dead_code)]
718    runtime: Arc<dyn camel_component_api::RuntimeObservability>,
719}
720
721impl Service<Exchange> for SedaProducer {
722    type Response = Exchange;
723    type Error = CamelError;
724    type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
725
726    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
727        Poll::Ready(Ok(()))
728    }
729
730    fn call(&mut self, exchange: Exchange) -> Self::Future {
731        let state = Arc::clone(&self.state);
732        let producer_config = self.producer_config.clone();
733        let original = exchange.clone();
734        Box::pin(async move {
735            if !state.has_active_consumers() {
736                if producer_config.discard_if_no_consumers {
737                    return Ok(exchange);
738                }
739                return Err(CamelError::EndpointCreationFailed(format!(
740                    "SEDA endpoint '{}' has no active consumers",
741                    state.config.name
742                )));
743            }
744
745            let should_wait = match producer_config.wait_for_task_to_complete {
746                WaitForTaskToComplete::Never => false,
747                WaitForTaskToComplete::Always => true,
748                WaitForTaskToComplete::IfReplyExpected => {
749                    state.config.exchange_pattern == ExchangePattern::InOut
750                }
751            };
752
753            if state.config.multiple_consumers && should_wait {
754                return Err(CamelError::EndpointCreationFailed(
755                    "multipleConsumers=true with waitForTaskToComplete != Never \
756                     is not supported — a single request cannot have N valid \
757                     replies without aggregator semantics"
758                        .to_string(),
759                ));
760            }
761
762            let (reply_tx, reply_rx) = if should_wait {
763                let (tx, rx) = oneshot::channel();
764                (Some(tx), Some(rx))
765            } else {
766                (None, None)
767            };
768
769            let envelope = ExchangeEnvelope { exchange, reply_tx };
770
771            match &state.mode {
772                SedaMode::Single { tx, .. } => {
773                    if producer_config.block_when_full {
774                        let result = tokio::time::timeout(
775                            Duration::from_millis(producer_config.timeout_ms),
776                            tx.send(envelope),
777                        )
778                        .await;
779                        match result {
780                            Ok(Ok(())) => {}
781                            Ok(Err(_)) => return Err(CamelError::ChannelClosed),
782                            Err(_) => {
783                                return Err(CamelError::EndpointCreationFailed(format!(
784                                    "SEDA producer timeout enqueueing on '{}' ({}ms)",
785                                    state.config.name, producer_config.timeout_ms
786                                )));
787                            }
788                        }
789                    } else {
790                        tx.try_send(envelope).map_err(|e| {
791                            if matches!(e, mpsc::error::TrySendError::Full(_)) {
792                                CamelError::EndpointCreationFailed(format!(
793                                    "SEDA queue '{}' is full (size={})",
794                                    state.config.name, state.config.size
795                                ))
796                            } else {
797                                CamelError::ChannelClosed
798                            }
799                        })?;
800                    }
801                }
802                SedaMode::Fanout { subscribers } => {
803                    let sender_list: Vec<mpsc::Sender<ExchangeEnvelope>> = {
804                        let subs_guard = subscribers.lock().unwrap_or_else(|e| e.into_inner());
805                        if subs_guard.is_empty() {
806                            if producer_config.discard_if_no_consumers {
807                                return Ok(original);
808                            }
809                            return Err(CamelError::EndpointCreationFailed(format!(
810                                "SEDA endpoint '{}' has no active subscribers",
811                                state.config.name
812                            )));
813                        }
814                        subs_guard.values().cloned().collect()
815                    };
816
817                    if producer_config.block_when_full {
818                        let mut permits: Vec<mpsc::OwnedPermit<ExchangeEnvelope>> =
819                            Vec::with_capacity(sender_list.len());
820                        for sender in &sender_list {
821                            let result = tokio::time::timeout(
822                                Duration::from_millis(producer_config.timeout_ms),
823                                sender.clone().reserve_owned(),
824                            )
825                            .await;
826                            match result {
827                                Ok(Ok(permit)) => permits.push(permit),
828                                Ok(Err(_)) => return Err(CamelError::ChannelClosed),
829                                Err(_) => {
830                                    return Err(CamelError::EndpointCreationFailed(format!(
831                                        "SEDA fanout timeout on '{}' ({}ms)",
832                                        state.config.name, producer_config.timeout_ms
833                                    )));
834                                }
835                            }
836                        }
837                        for permit in permits {
838                            permit.send(ExchangeEnvelope {
839                                exchange: original.clone(),
840                                reply_tx: None,
841                            });
842                        }
843                    } else {
844                        let mut permits: Vec<mpsc::OwnedPermit<ExchangeEnvelope>> =
845                            Vec::with_capacity(sender_list.len());
846                        for sender in &sender_list {
847                            match sender.clone().try_reserve_owned() {
848                                Ok(permit) => permits.push(permit),
849                                Err(e) => {
850                                    if matches!(e, mpsc::error::TrySendError::Full(_)) {
851                                        return Err(CamelError::EndpointCreationFailed(format!(
852                                            "SEDA queue '{}' subscriber full during fanout (size={})",
853                                            state.config.name, state.config.size
854                                        )));
855                                    } else {
856                                        return Err(CamelError::ChannelClosed);
857                                    }
858                                }
859                            }
860                        }
861                        for permit in permits {
862                            permit.send(ExchangeEnvelope {
863                                exchange: original.clone(),
864                                reply_tx: None,
865                            });
866                        }
867                    }
868                }
869            }
870
871            if !should_wait {
872                return Ok(original);
873            }
874
875            let reply_rx = reply_rx.ok_or(CamelError::ChannelClosed)?;
876            let result =
877                tokio::time::timeout(Duration::from_millis(producer_config.timeout_ms), reply_rx)
878                    .await;
879            match result {
880                Ok(Ok(reply)) => reply,
881                Ok(Err(_)) => Err(CamelError::ChannelClosed),
882                Err(_) => Err(CamelError::EndpointCreationFailed(format!(
883                    "SEDA producer timeout waiting for reply on '{}' ({}ms)",
884                    state.config.name, producer_config.timeout_ms
885                ))),
886            }
887        })
888    }
889}
890
891// ---------------------------------------------------------------------------
892// Tests
893// ---------------------------------------------------------------------------
894
895#[cfg(test)]
896mod config_tests {
897    use super::*;
898
899    #[test]
900    fn test_seda_config_from_uri_minimal() {
901        let config = SedaConfig::from_uri("seda:foo").unwrap();
902        assert_eq!(config.name, "foo");
903        assert_eq!(config.size, 1000);
904        assert_eq!(config.concurrent_consumers, 1);
905        assert!(!config.multiple_consumers);
906        assert!(!config.block_when_full);
907        assert!(!config.discard_if_no_consumers);
908        assert_eq!(config.timeout_ms, 30_000);
909        assert_eq!(
910            config.wait_for_task_to_complete,
911            WaitForTaskToComplete::IfReplyExpected
912        );
913        assert_eq!(config.exchange_pattern, ExchangePattern::InOnly);
914    }
915
916    #[test]
917    fn test_seda_config_from_uri_full() {
918        let config = SedaConfig::from_uri(
919            "seda:bar?size=500&concurrentConsumers=4&multipleConsumers=true\
920             &blockWhenFull=true&discardIfNoConsumers=false&timeout=5000\
921             &waitForTaskToComplete=Never&exchangePattern=InOut",
922        )
923        .unwrap();
924        assert_eq!(config.name, "bar");
925        assert_eq!(config.size, 500);
926        assert_eq!(config.concurrent_consumers, 4);
927        assert!(config.multiple_consumers);
928        assert!(config.block_when_full);
929        assert!(!config.discard_if_no_consumers);
930        assert_eq!(config.timeout_ms, 5000);
931        assert_eq!(
932            config.wait_for_task_to_complete,
933            WaitForTaskToComplete::Never
934        );
935        assert_eq!(config.exchange_pattern, ExchangePattern::InOut);
936    }
937
938    #[test]
939    fn test_seda_config_invalid_scheme() {
940        let err = SedaConfig::from_uri("timer:foo").unwrap_err();
941        assert!(err.to_string().contains("expected 'seda'"));
942    }
943
944    #[test]
945    fn test_seda_config_empty_name() {
946        let err = SedaConfig::from_uri("seda:").unwrap_err();
947        assert!(err.to_string().contains("must not be empty"));
948    }
949
950    #[test]
951    fn test_seda_size_zero() {
952        let err = SedaConfig::from_uri("seda:foo?size=0").unwrap_err();
953        assert!(err.to_string().contains("size must be greater than 0"));
954    }
955
956    #[test]
957    fn test_seda_config_concurrent_consumers_zero_clamped() {
958        let config = SedaConfig::from_uri("seda:foo?concurrentConsumers=0").unwrap();
959        assert_eq!(config.concurrent_consumers, 1);
960    }
961
962    #[test]
963    fn test_seda_config_case_insensitive_enums() {
964        let config =
965            SedaConfig::from_uri("seda:foo?waitForTaskToComplete=never&exchangePattern=inonly")
966                .unwrap();
967        assert_eq!(
968            config.wait_for_task_to_complete,
969            WaitForTaskToComplete::Never
970        );
971        assert_eq!(config.exchange_pattern, ExchangePattern::InOnly);
972    }
973
974    #[test]
975    fn test_seda_config_invalid_enum() {
976        let err = SedaConfig::from_uri("seda:foo?exchangePattern=invalid").unwrap_err();
977        assert!(err.to_string().contains("invalid exchangePattern"));
978    }
979}
980
981#[cfg(test)]
982mod consumer_producer_tests {
983    use super::*;
984    use camel_api::Value;
985    use camel_component_api::Message;
986    use camel_component_api::NoOpComponentContext;
987    use tokio::time::Duration;
988    use tower::ServiceExt;
989
990    fn test_producer_ctx() -> ProducerContext {
991        ProducerContext::default()
992    }
993
994    fn create_component() -> SedaComponent {
995        SedaComponent::new()
996    }
997
998    #[tokio::test]
999    async fn test_seda_single_consumer_producer_roundtrip() {
1000        let comp = create_component();
1001        let ep = comp
1002            .create_endpoint("seda:test1", &NoOpComponentContext)
1003            .unwrap();
1004
1005        let mut consumer = ep.create_consumer(rt()).unwrap();
1006        let (route_tx, mut route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1007        let ctx = ConsumerContext::new(
1008            route_tx,
1009            CancellationToken::new(),
1010            "seda-test-route".to_string(),
1011        );
1012        consumer.start(ctx).await.unwrap();
1013
1014        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1015        let exchange = Exchange::new(Message::new("hello seda"));
1016        let result = producer.oneshot(exchange).await;
1017        assert!(result.is_ok());
1018
1019        let received = tokio::time::timeout(Duration::from_millis(500), route_rx.recv())
1020            .await
1021            .unwrap()
1022            .unwrap();
1023        assert_eq!(received.exchange.input.body.as_text(), Some("hello seda"));
1024
1025        consumer.stop().await.unwrap();
1026    }
1027
1028    #[tokio::test]
1029    async fn test_seda_inout_roundtrip() {
1030        let comp = create_component();
1031        let ep = comp
1032            .create_endpoint("seda:io?exchangePattern=InOut", &NoOpComponentContext)
1033            .unwrap();
1034
1035        let mut consumer = ep.create_consumer(rt()).unwrap();
1036        let (route_tx, _) = mpsc::channel::<ExchangeEnvelope>(16);
1037        let ctx = ConsumerContext::new(
1038            route_tx,
1039            CancellationToken::new(),
1040            "seda-test-route".to_string(),
1041        );
1042        consumer.start(ctx).await.unwrap();
1043
1044        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1045        let exchange = Exchange::new(Message::new("io test"));
1046
1047        let result =
1048            tokio::time::timeout(Duration::from_millis(500), producer.oneshot(exchange)).await;
1049        assert!(result.is_err() || result.unwrap().is_err());
1050
1051        consumer.stop().await.unwrap();
1052    }
1053
1054    #[tokio::test]
1055    async fn test_seda_inonly_fire_and_forget() {
1056        let comp = create_component();
1057        let ep = comp
1058            .create_endpoint("seda:ff", &NoOpComponentContext)
1059            .unwrap();
1060
1061        let mut consumer = ep.create_consumer(rt()).unwrap();
1062        let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1063        let ctx = ConsumerContext::new(
1064            route_tx,
1065            CancellationToken::new(),
1066            "seda-test-route".to_string(),
1067        );
1068        consumer.start(ctx).await.unwrap();
1069
1070        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1071        let exchange = Exchange::new(Message::new("fire and forget"));
1072        let result = producer.oneshot(exchange).await;
1073        assert!(result.is_ok());
1074
1075        consumer.stop().await.unwrap();
1076    }
1077
1078    #[tokio::test]
1079    async fn test_seda_queue_full_fail() {
1080        let comp = create_component();
1081        let ep = comp
1082            .create_endpoint("seda:full?size=2", &NoOpComponentContext)
1083            .unwrap();
1084
1085        let mut consumer = ep.create_consumer(rt()).unwrap();
1086        let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(16);
1087        let ctx = ConsumerContext::new(
1088            route_tx,
1089            CancellationToken::new(),
1090            "seda-test-route".to_string(),
1091        );
1092        consumer.start(ctx).await.unwrap();
1093
1094        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1095        producer
1096            .clone()
1097            .oneshot(Exchange::new(Message::new("1")))
1098            .await
1099            .unwrap();
1100        producer
1101            .clone()
1102            .oneshot(Exchange::new(Message::new("2")))
1103            .await
1104            .unwrap();
1105
1106        let result = producer.oneshot(Exchange::new(Message::new("3"))).await;
1107        assert!(result.is_err());
1108        assert!(result.unwrap_err().to_string().contains("full"));
1109
1110        consumer.stop().await.unwrap();
1111    }
1112
1113    #[tokio::test]
1114    async fn test_seda_block_when_full_with_timeout() {
1115        let comp = create_component();
1116        let ep = comp
1117            .create_endpoint(
1118                "seda:bwf?size=1&blockWhenFull=true&timeout=50",
1119                &NoOpComponentContext,
1120            )
1121            .unwrap();
1122
1123        let mut consumer = ep.create_consumer(rt()).unwrap();
1124        let (route_tx, _route_rx) = mpsc::channel::<ExchangeEnvelope>(1);
1125        route_tx
1126            .send(ExchangeEnvelope {
1127                exchange: Exchange::new(Message::new("dummy")),
1128                reply_tx: None,
1129            })
1130            .await
1131            .unwrap();
1132        let ctx = ConsumerContext::new(
1133            route_tx,
1134            CancellationToken::new(),
1135            "seda-test-route".to_string(),
1136        );
1137        consumer.start(ctx).await.unwrap();
1138
1139        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1140        producer
1141            .clone()
1142            .oneshot(Exchange::new(Message::new("1")))
1143            .await
1144            .unwrap();
1145
1146        producer
1147            .clone()
1148            .oneshot(Exchange::new(Message::new("2")))
1149            .await
1150            .unwrap();
1151
1152        let result = tokio::time::timeout(
1153            Duration::from_millis(200),
1154            producer.oneshot(Exchange::new(Message::new("3"))),
1155        )
1156        .await;
1157        assert!(result.is_ok());
1158        let inner = result.unwrap();
1159        assert!(inner.is_err());
1160        assert!(inner.unwrap_err().to_string().contains("timeout"));
1161
1162        consumer.stop().await.unwrap();
1163    }
1164
1165    #[tokio::test]
1166    async fn test_seda_no_consumers_fail() {
1167        let comp = create_component();
1168        let ep = comp
1169            .create_endpoint("seda:nocons", &NoOpComponentContext)
1170            .unwrap();
1171
1172        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1173        let result = producer.oneshot(Exchange::new(Message::new("test"))).await;
1174        assert!(result.is_err());
1175        assert!(
1176            result
1177                .unwrap_err()
1178                .to_string()
1179                .contains("no active consumers")
1180        );
1181    }
1182
1183    #[tokio::test]
1184    async fn test_seda_no_consumers_discard() {
1185        let comp = create_component();
1186        let ep = comp
1187            .create_endpoint(
1188                "seda:discard?discardIfNoConsumers=true",
1189                &NoOpComponentContext,
1190            )
1191            .unwrap();
1192
1193        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1194        let result = producer.oneshot(Exchange::new(Message::new("test"))).await;
1195        assert!(result.is_ok());
1196    }
1197
1198    #[tokio::test]
1199    async fn test_seda_duplicate_single_consumer() {
1200        let comp = create_component();
1201        let ep = comp
1202            .create_endpoint("seda:dup", &NoOpComponentContext)
1203            .unwrap();
1204
1205        let mut consumer_a = ep.create_consumer(rt()).unwrap();
1206        let (tx_a, _rx_a) = mpsc::channel::<ExchangeEnvelope>(16);
1207        let ctx_a = ConsumerContext::new(
1208            tx_a,
1209            CancellationToken::new(),
1210            "seda-test-route-a".to_string(),
1211        );
1212        consumer_a.start(ctx_a).await.unwrap();
1213
1214        let mut consumer_b = ep.create_consumer(rt()).unwrap();
1215        let (tx_b, _rx_b) = mpsc::channel::<ExchangeEnvelope>(16);
1216        let ctx_b = ConsumerContext::new(
1217            tx_b,
1218            CancellationToken::new(),
1219            "seda-test-route-b".to_string(),
1220        );
1221        let result = consumer_b.start(ctx_b).await;
1222        assert!(result.is_err());
1223        assert!(
1224            result
1225                .unwrap_err()
1226                .to_string()
1227                .contains("already has a registered consumer")
1228        );
1229
1230        consumer_a.stop().await.unwrap();
1231    }
1232
1233    #[tokio::test]
1234    async fn test_seda_fanout_two_consumers() {
1235        let comp = create_component();
1236        let ep = comp
1237            .create_endpoint("seda:fan?multipleConsumers=true", &NoOpComponentContext)
1238            .unwrap();
1239
1240        let mut consumer_a = ep.create_consumer(rt()).unwrap();
1241        let (tx_a, mut rx_a) = mpsc::channel::<ExchangeEnvelope>(16);
1242        let ctx_a = ConsumerContext::new(
1243            tx_a,
1244            CancellationToken::new(),
1245            "seda-test-route-a".to_string(),
1246        );
1247        consumer_a.start(ctx_a).await.unwrap();
1248
1249        let mut consumer_b = ep.create_consumer(rt()).unwrap();
1250        let (tx_b, mut rx_b) = mpsc::channel::<ExchangeEnvelope>(16);
1251        let ctx_b = ConsumerContext::new(
1252            tx_b,
1253            CancellationToken::new(),
1254            "seda-test-route-b".to_string(),
1255        );
1256        consumer_b.start(ctx_b).await.unwrap();
1257
1258        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1259        producer
1260            .oneshot(Exchange::new(Message::new("fanout msg")))
1261            .await
1262            .unwrap();
1263
1264        let recv_a = tokio::time::timeout(Duration::from_millis(500), rx_a.recv())
1265            .await
1266            .unwrap()
1267            .unwrap();
1268        let recv_b = tokio::time::timeout(Duration::from_millis(500), rx_b.recv())
1269            .await
1270            .unwrap()
1271            .unwrap();
1272
1273        assert_eq!(recv_a.exchange.input.body.as_text(), Some("fanout msg"));
1274        assert_eq!(recv_b.exchange.input.body.as_text(), Some("fanout msg"));
1275
1276        consumer_a.stop().await.unwrap();
1277        consumer_b.stop().await.unwrap();
1278    }
1279
1280    #[tokio::test]
1281    async fn test_seda_fanout_inout_rejected() {
1282        let comp = create_component();
1283        let ep = comp
1284            .create_endpoint(
1285                "seda:fanout?multipleConsumers=true&exchangePattern=InOut",
1286                &NoOpComponentContext,
1287            )
1288            .unwrap();
1289
1290        let mut consumer = ep.create_consumer(rt()).unwrap();
1291        let (tx, _rx) = mpsc::channel::<ExchangeEnvelope>(16);
1292        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "seda-test-route".to_string());
1293        consumer.start(ctx).await.unwrap();
1294
1295        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1296        let result = producer.oneshot(Exchange::new(Message::new("test"))).await;
1297        assert!(result.is_err());
1298        assert!(
1299            result
1300                .unwrap_err()
1301                .to_string()
1302                .contains("multipleConsumers")
1303        );
1304
1305        consumer.stop().await.unwrap();
1306    }
1307
1308    #[tokio::test]
1309    async fn test_seda_consumer_stop_unregisters() {
1310        let comp = create_component();
1311        let ep = comp
1312            .create_endpoint("seda:stop", &NoOpComponentContext)
1313            .unwrap();
1314
1315        let mut consumer = ep.create_consumer(rt()).unwrap();
1316        let (tx, _rx) = mpsc::channel::<ExchangeEnvelope>(16);
1317        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "seda-test-route".to_string());
1318        consumer.start(ctx).await.unwrap();
1319
1320        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1321        producer
1322            .clone()
1323            .oneshot(Exchange::new(Message::new("before stop")))
1324            .await
1325            .unwrap();
1326
1327        consumer.stop().await.unwrap();
1328
1329        tokio::time::sleep(Duration::from_millis(50)).await;
1330
1331        let result = producer
1332            .oneshot(Exchange::new(Message::new("after stop")))
1333            .await;
1334        assert!(result.is_err());
1335        assert!(
1336            result
1337                .unwrap_err()
1338                .to_string()
1339                .contains("no active consumers")
1340        );
1341    }
1342
1343    #[test]
1344    fn test_seda_concurrent_consumers_hint() {
1345        let comp = create_component();
1346        let ep = comp
1347            .create_endpoint("seda:conc?concurrentConsumers=4", &NoOpComponentContext)
1348            .unwrap();
1349        let consumer = ep.create_consumer(rt()).unwrap();
1350        assert_eq!(
1351            consumer.concurrency_model(),
1352            ConcurrencyModel::Concurrent { max: Some(4) }
1353        );
1354    }
1355
1356    #[tokio::test]
1357    async fn test_seda_config_mismatch() {
1358        let comp = create_component();
1359        let _ep1 = comp
1360            .create_endpoint("seda:mm?size=100", &NoOpComponentContext)
1361            .unwrap();
1362        let result = comp.create_endpoint("seda:mm?size=200", &NoOpComponentContext);
1363        let err = match result {
1364            Err(e) => e,
1365            Ok(_) => panic!("expected config mismatch error"),
1366        };
1367        assert!(err.to_string().contains("size"));
1368    }
1369
1370    #[tokio::test]
1371    async fn test_seda_wait_always_inonly() {
1372        let comp = create_component();
1373        let ep = comp
1374            .create_endpoint(
1375                "seda:waitalways?waitForTaskToComplete=Always",
1376                &NoOpComponentContext,
1377            )
1378            .unwrap();
1379
1380        let mut consumer = ep.create_consumer(rt()).unwrap();
1381        let (route_tx, _) = mpsc::channel::<ExchangeEnvelope>(16);
1382        let ctx = ConsumerContext::new(
1383            route_tx,
1384            CancellationToken::new(),
1385            "seda-test-route".to_string(),
1386        );
1387        consumer.start(ctx).await.unwrap();
1388
1389        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1390        let result = tokio::time::timeout(
1391            Duration::from_millis(500),
1392            producer.oneshot(Exchange::new(Message::new("always wait"))),
1393        )
1394        .await;
1395        assert!(result.is_err() || result.unwrap().is_err());
1396
1397        consumer.stop().await.unwrap();
1398    }
1399
1400    #[tokio::test]
1401    async fn test_seda_fanout_all_or_nothing() {
1402        let comp = create_component();
1403        let ep = comp
1404            .create_endpoint(
1405                "seda:aon?multipleConsumers=true&size=2",
1406                &NoOpComponentContext,
1407            )
1408            .unwrap();
1409
1410        let mut consumer_a = ep.create_consumer(rt()).unwrap();
1411        let (tx_a, _rx_a) = mpsc::channel::<ExchangeEnvelope>(1);
1412        let ctx_a = ConsumerContext::new(
1413            tx_a,
1414            CancellationToken::new(),
1415            "seda-test-route-a".to_string(),
1416        );
1417        consumer_a.start(ctx_a).await.unwrap();
1418
1419        let mut consumer_b = ep.create_consumer(rt()).unwrap();
1420        let (tx_b, _rx_b) = mpsc::channel::<ExchangeEnvelope>(1);
1421        let ctx_b = ConsumerContext::new(
1422            tx_b,
1423            CancellationToken::new(),
1424            "seda-test-route-b".to_string(),
1425        );
1426        consumer_b.start(ctx_b).await.unwrap();
1427
1428        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1429        producer
1430            .clone()
1431            .oneshot(Exchange::new(Message::new("1")))
1432            .await
1433            .unwrap();
1434        producer
1435            .clone()
1436            .oneshot(Exchange::new(Message::new("2")))
1437            .await
1438            .unwrap();
1439
1440        let result = producer.oneshot(Exchange::new(Message::new("3"))).await;
1441        assert!(result.is_err());
1442        let err_msg = result.unwrap_err().to_string();
1443        assert!(err_msg.contains("full") || err_msg.contains("subscriber"));
1444
1445        consumer_a.stop().await.unwrap();
1446        consumer_b.stop().await.unwrap();
1447    }
1448
1449    #[tokio::test]
1450    async fn test_seda_fanout_block_when_full_rejects_closed_subscriber_without_partial_delivery() {
1451        let comp = create_component();
1452        let ep = comp
1453            .create_endpoint(
1454                "seda:aonblock?multipleConsumers=true&size=2&blockWhenFull=true&timeout=100",
1455                &NoOpComponentContext,
1456            )
1457            .unwrap();
1458
1459        let mut consumer_a = ep.create_consumer(rt()).unwrap();
1460        let (tx_a, mut rx_a) = mpsc::channel::<ExchangeEnvelope>(1);
1461        let ctx_a = ConsumerContext::new(
1462            tx_a,
1463            CancellationToken::new(),
1464            "seda-test-route-a".to_string(),
1465        );
1466        consumer_a.start(ctx_a).await.unwrap();
1467
1468        let state = comp
1469            .endpoints
1470            .lock()
1471            .unwrap_or_else(|e| e.into_inner())
1472            .get("aonblock")
1473            .cloned()
1474            .unwrap();
1475        let (closed_tx, closed_rx) = mpsc::channel::<ExchangeEnvelope>(1);
1476        drop(closed_rx);
1477        match &state.mode {
1478            SedaMode::Fanout { subscribers } => {
1479                subscribers
1480                    .lock()
1481                    .unwrap_or_else(|e| e.into_inner())
1482                    .insert("closed-subscriber".to_string(), closed_tx);
1483            }
1484            SedaMode::Single { .. } => panic!("expected fanout mode"),
1485        }
1486
1487        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1488        let result = producer
1489            .oneshot(Exchange::new(Message::new("partial")))
1490            .await;
1491
1492        assert!(matches!(result, Err(CamelError::ChannelClosed)));
1493        let delivered = tokio::time::timeout(Duration::from_millis(50), rx_a.recv()).await;
1494        assert!(
1495            delivered.is_err(),
1496            "fanout delivered to only one subscriber"
1497        );
1498
1499        consumer_a.stop().await.unwrap();
1500    }
1501
1502    #[tokio::test]
1503    async fn test_seda_discard_if_no_consumers_fanout() {
1504        let comp = create_component();
1505        let ep = comp
1506            .create_endpoint(
1507                "seda:discardfan?multipleConsumers=true&discardIfNoConsumers=true",
1508                &NoOpComponentContext,
1509            )
1510            .unwrap();
1511
1512        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1513        let result = producer
1514            .oneshot(Exchange::new(Message::new("discard")))
1515            .await;
1516        assert!(result.is_ok());
1517    }
1518
1519    #[tokio::test]
1520    async fn test_seda_multiple_producers_single_consumer() {
1521        let comp = create_component();
1522        let ep = comp
1523            .create_endpoint("seda:mpsc", &NoOpComponentContext)
1524            .unwrap();
1525
1526        let mut consumer = ep.create_consumer(rt()).unwrap();
1527        let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(16);
1528        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "seda-test-route".to_string());
1529        consumer.start(ctx).await.unwrap();
1530
1531        let producer_a = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1532        let producer_b = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1533
1534        producer_a
1535            .oneshot(Exchange::new(Message::new("A")))
1536            .await
1537            .unwrap();
1538        producer_b
1539            .oneshot(Exchange::new(Message::new("B")))
1540            .await
1541            .unwrap();
1542
1543        let mut bodies = Vec::new();
1544        for _ in 0..2 {
1545            let received = tokio::time::timeout(Duration::from_millis(500), rx.recv())
1546                .await
1547                .unwrap()
1548                .unwrap();
1549            bodies.push(received.exchange.input.body.as_text().unwrap().to_string());
1550        }
1551        bodies.sort();
1552        assert_eq!(bodies, vec!["A", "B"]);
1553
1554        consumer.stop().await.unwrap();
1555    }
1556
1557    #[tokio::test]
1558    async fn test_seda_inout_timeout_no_reply() {
1559        let comp = create_component();
1560        let ep = comp
1561            .create_endpoint(
1562                "seda:iotimeout?exchangePattern=InOut&timeout=100",
1563                &NoOpComponentContext,
1564            )
1565            .unwrap();
1566
1567        let mut consumer = ep.create_consumer(rt()).unwrap();
1568        let (tx, _rx) = mpsc::channel::<ExchangeEnvelope>(16);
1569        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "seda-test-route".to_string());
1570        consumer.start(ctx).await.unwrap();
1571
1572        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1573        let result = tokio::time::timeout(
1574            Duration::from_millis(500),
1575            producer.oneshot(Exchange::new(Message::new("no reply"))),
1576        )
1577        .await
1578        .unwrap();
1579
1580        assert!(result.is_err());
1581        assert!(result.unwrap_err().to_string().contains("timeout"));
1582
1583        consumer.stop().await.unwrap();
1584    }
1585
1586    #[tokio::test]
1587    async fn test_seda_producer_preserves_headers() {
1588        let comp = create_component();
1589        let ep = comp
1590            .create_endpoint("seda:hdr", &NoOpComponentContext)
1591            .unwrap();
1592
1593        let mut consumer = ep.create_consumer(rt()).unwrap();
1594        let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(16);
1595        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "seda-test-route".to_string());
1596        consumer.start(ctx).await.unwrap();
1597
1598        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1599        let mut msg = Message::new("with headers");
1600        msg.set_header("X-Custom", Value::String("test-value".into()));
1601        msg.set_header("X-Count", Value::Number(42.into()));
1602        producer.oneshot(Exchange::new(msg)).await.unwrap();
1603
1604        let received = tokio::time::timeout(Duration::from_millis(500), rx.recv())
1605            .await
1606            .unwrap()
1607            .unwrap();
1608
1609        assert_eq!(
1610            received.exchange.input.header("X-Custom"),
1611            Some(&Value::String("test-value".into()))
1612        );
1613        assert_eq!(
1614            received.exchange.input.header("X-Count"),
1615            Some(&Value::Number(42.into()))
1616        );
1617
1618        consumer.stop().await.unwrap();
1619    }
1620
1621    #[tokio::test]
1622    async fn test_seda_concurrent_send_receive() {
1623        use std::sync::atomic::AtomicU64;
1624
1625        let comp = create_component();
1626        let ep = comp
1627            .create_endpoint("seda:concsend?size=1000", &NoOpComponentContext)
1628            .unwrap();
1629
1630        let mut consumer = ep.create_consumer(rt()).unwrap();
1631        let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(1000);
1632        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "seda-test-route".to_string());
1633        consumer.start(ctx).await.unwrap();
1634
1635        let counter = Arc::new(AtomicU64::new(0));
1636        let counter_clone = counter.clone();
1637        let recv_handle = tokio::spawn(async move {
1638            while let Some(envelope) = rx.recv().await {
1639                counter_clone.fetch_add(1, Ordering::SeqCst);
1640                let _ = envelope;
1641            }
1642        });
1643
1644        let mut handles = Vec::new();
1645        for i in 0..10u64 {
1646            let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1647            handles.push(tokio::spawn(async move {
1648                for j in 0..10u64 {
1649                    producer
1650                        .clone()
1651                        .oneshot(Exchange::new(Message::new(format!("{}-{}", i, j))))
1652                        .await
1653                        .unwrap();
1654                }
1655            }));
1656        }
1657
1658        for h in handles {
1659            h.await.unwrap();
1660        }
1661
1662        tokio::time::timeout(Duration::from_secs(2), async {
1663            loop {
1664                if counter.load(Ordering::SeqCst) == 100 {
1665                    break;
1666                }
1667                tokio::time::sleep(Duration::from_millis(10)).await;
1668            }
1669        })
1670        .await
1671        .unwrap();
1672
1673        recv_handle.abort();
1674        assert_eq!(counter.load(Ordering::SeqCst), 100);
1675
1676        consumer.stop().await.unwrap();
1677    }
1678
1679    #[tokio::test]
1680    async fn test_seda_size_one_queue() {
1681        let comp = create_component();
1682        let ep = comp
1683            .create_endpoint("seda:sz1?size=1", &NoOpComponentContext)
1684            .unwrap();
1685
1686        let mut consumer = ep.create_consumer(rt()).unwrap();
1687        let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(16);
1688        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "seda-test-route".to_string());
1689        consumer.start(ctx).await.unwrap();
1690
1691        let producer = ep.create_producer(rt(), &test_producer_ctx()).unwrap();
1692        producer
1693            .clone()
1694            .oneshot(Exchange::new(Message::new("1")))
1695            .await
1696            .unwrap();
1697
1698        let result = producer
1699            .clone()
1700            .oneshot(Exchange::new(Message::new("2")))
1701            .await;
1702        assert!(result.is_err());
1703        assert!(result.unwrap_err().to_string().contains("full"));
1704
1705        let _dropped = tokio::time::timeout(Duration::from_millis(500), rx.recv())
1706            .await
1707            .unwrap()
1708            .unwrap();
1709
1710        producer
1711            .oneshot(Exchange::new(Message::new("3")))
1712            .await
1713            .unwrap();
1714
1715        consumer.stop().await.unwrap();
1716    }
1717}