1use 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#[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#[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
262static 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
275type ConsumerId = String;
280
281enum 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 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
338type 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
471struct 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#[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
536struct SedaConsumer {
541 state: Arc<SedaEndpointState>,
542 consumer_id: ConsumerId,
543 started: bool,
544 cancel_token: CancellationToken,
545 forwarder_handles: Vec<JoinHandle<Result<(), CamelError>>>,
546 #[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 self.forwarder_handles.pop()
687 }
688}
689
690async 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#[derive(Clone)]
712struct SedaProducer {
713 state: Arc<SedaEndpointState>,
714 producer_config: ProducerConfig,
715 #[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#[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}