1use std::collections::HashMap;
2use std::future::Future;
3use std::path::PathBuf;
4use std::pin::Pin;
5use std::sync::Arc;
6use std::sync::atomic::{AtomicBool, Ordering};
7use std::task::{Context, Poll};
8use std::time::{Duration, Instant};
9
10use camel_bridge::{
11 download::ensure_binary,
12 health::wait_for_health,
13 process::{BridgeProcess, BridgeProcessConfig, BrokerType},
14};
15use camel_component_api::{
16 BoxProcessor, CamelError, Component, Consumer, Endpoint, Exchange, NetworkRetryPolicy,
17 ProducerContext,
18};
19use dashmap::DashMap;
20use tokio::sync::{Mutex, watch};
21use tonic::transport::Channel;
22use tower::Service;
23use tracing::{info, warn};
24
25use crate::config::{BrokerConfig, JmsEndpointConfig, JmsPoolConfig};
26use crate::consumer::JmsConsumer;
27use crate::health::JmsHealthCheck;
28use crate::producer::JmsProducer;
29use crate::proto::{HealthRequest, bridge_service_client::BridgeServiceClient};
30
31pub const BRIDGE_TRANSPORT_ERROR_PREFIX: &str = "JMS gRPC ";
39const MAX_RESTART_ATTEMPTS: u32 = 10;
40
41#[derive(Debug, Clone)]
44pub enum BridgeState {
45 Starting,
46 Ready { channel: Channel },
47 Degraded(String),
48 Restarting { attempt: u32, next_at: Instant },
49 Stopped,
50}
51
52pub struct BridgeSlot {
55 pub name: String,
56 pub broker_url: String,
57 pub broker_type: BrokerType,
58 pub credentials: Option<(String, String)>,
59 pub state_rx: watch::Receiver<BridgeState>,
60 pub(crate) state_tx: watch::Sender<BridgeState>,
61 pub process: Arc<tokio::sync::Mutex<Option<BridgeProcess>>>,
63 pub(crate) health_monitor_handle: Arc<tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>>,
66}
67
68impl std::fmt::Debug for BridgeSlot {
69 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
70 f.debug_struct("BridgeSlot")
71 .field("name", &self.name)
72 .field("broker_url", &self.broker_url)
73 .field("broker_type", &self.broker_type)
74 .finish()
75 }
76}
77
78pub struct JmsBridgePool {
81 pub(crate) slots: DashMap<String, Arc<BridgeSlot>>,
82 pub(crate) config: HashMap<String, BrokerConfig>,
83 pub(crate) bridge_start_timeout_ms: u64,
84 pub(crate) reconnect: NetworkRetryPolicy,
85 pub(crate) health_check_interval_ms: u64,
86 pub(crate) bridge_version: String,
87 pub(crate) bridge_cache_dir: PathBuf,
88 pub(crate) max_bridges: usize,
90 bridge_create_lock: Mutex<()>,
92 pub(crate) shutting_down: Arc<AtomicBool>,
93}
94
95impl JmsBridgePool {
96 pub fn from_config(pool_config: JmsPoolConfig) -> Result<Self, CamelError> {
97 pool_config.validate()?;
98 let mut reconnect = pool_config.reconnect;
101 if pool_config.broker_reconnect_interval_ms
102 != crate::config::default_broker_reconnect_interval_ms()
103 {
104 reconnect.initial_delay =
105 Duration::from_millis(pool_config.broker_reconnect_interval_ms);
106 }
107 Ok(Self {
108 slots: DashMap::new(),
109 config: pool_config.brokers,
110 bridge_start_timeout_ms: pool_config.bridge_start_timeout_ms,
111 reconnect,
112 health_check_interval_ms: pool_config.health_check_interval_ms,
113 bridge_version: crate::BRIDGE_VERSION.to_string(),
114 bridge_cache_dir: pool_config.bridge_cache_dir,
115 max_bridges: pool_config.max_bridges,
116 bridge_create_lock: Mutex::new(()),
117 shutting_down: Arc::new(AtomicBool::new(false)),
118 })
119 }
120
121 pub fn resolve_broker_name(&self, name: Option<&str>) -> Result<String, CamelError> {
128 match name {
129 Some(n) => {
130 if self.config.contains_key(n) {
131 Ok(n.to_string())
132 } else {
133 Err(CamelError::ProcessorError(format!(
134 "Unknown JMS broker '{n}' — declare it in [components.jms.brokers] in Camel.toml",
135 )))
136 }
137 }
138 None => match self.config.len() {
139 0 => Err(CamelError::ProcessorError(
140 "No JMS brokers configured — declare at least one in [components.jms.brokers] in Camel.toml".to_string(),
141 )),
142 1 => Ok(self.config.keys().next().unwrap().clone()), _ => Err(CamelError::ProcessorError(format!(
144 "Multiple JMS brokers configured ({}); specify one with ?broker=<name> in the URI",
145 self.config.keys().cloned().collect::<Vec<_>>().join(", ")
146 ))),
147 },
148 }
149 }
150
151 pub fn resolve_broker_type(&self, scheme: &str, broker_name: &str) -> BrokerType {
153 let config_type = self
154 .config
155 .get(broker_name)
156 .map(|c| c.broker_type.clone())
157 .unwrap_or(BrokerType::Generic);
158
159 match scheme {
160 "activemq" => {
161 if config_type != BrokerType::ActiveMq && config_type != BrokerType::Generic {
162 warn!(
163 "Scheme 'activemq' overrides configured broker_type '{:?}' for broker '{}'",
164 config_type, broker_name
165 );
166 }
167 BrokerType::ActiveMq
168 }
169 "artemis" => {
170 if config_type != BrokerType::Artemis && config_type != BrokerType::Generic {
171 warn!(
172 "Scheme 'artemis' overrides configured broker_type '{:?}' for broker '{}'",
173 config_type, broker_name
174 );
175 }
176 BrokerType::Artemis
177 }
178 _ => config_type,
179 }
180 }
181
182 pub async fn get_or_create_slot(
185 &self,
186 broker_name: &str,
187 ) -> Result<Arc<BridgeSlot>, CamelError> {
188 if let Some(slot) = self.slots.get(broker_name) {
189 return Ok(Arc::clone(&*slot));
190 }
191
192 let _guard = self.bridge_create_lock.lock().await;
195
196 if let Some(slot) = self.slots.get(broker_name) {
198 return Ok(Arc::clone(&*slot));
199 }
200
201 let total_count = self.slots.len();
208 if total_count >= self.max_bridges {
209 return Err(CamelError::Config(format!(
210 "JMS bridge limit reached: {total_count} bridge(s) >= max_bridges ({})",
211 self.max_bridges
212 )));
213 }
214
215 let broker_config = self.config.get(broker_name).ok_or_else(|| {
216 CamelError::ProcessorError(format!("Unknown JMS broker '{}'", broker_name))
217 })?;
218
219 let broker_url = broker_config.broker_url.clone();
221 let broker_type = broker_config.broker_type.clone();
222 let credentials = match (&broker_config.username, &broker_config.password) {
223 (Some(u), Some(p)) => Some((u.clone(), p.clone())),
224 _ => None,
225 };
226
227 let slot = match self.slots.entry(broker_name.to_string()) {
228 dashmap::Entry::Occupied(existing) => {
229 return Ok(Arc::clone(existing.get()));
230 }
231 dashmap::Entry::Vacant(entry) => {
232 let (state_tx, state_rx) = watch::channel(BridgeState::Starting);
233 let slot = Arc::new(BridgeSlot {
234 name: broker_name.to_string(),
235 broker_url: broker_url.clone(),
236 broker_type: broker_type.clone(),
237 credentials: credentials.clone(),
238 state_rx,
239 state_tx,
240 process: Arc::new(tokio::sync::Mutex::new(None)),
241 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
242 });
243 entry.insert(Arc::clone(&slot));
244 slot
245 }
246 };
247
248 let start_result = Self::start_bridge_process(
249 &self.bridge_version,
250 &self.bridge_cache_dir,
251 self.bridge_start_timeout_ms,
252 &broker_url,
253 &broker_type,
254 &credentials,
255 )
256 .await;
257
258 match start_result {
259 Ok((process, channel)) => {
260 {
261 let mut guard = slot.process.lock().await;
262 *guard = Some(process);
263 }
264 let _ = slot.state_tx.send(BridgeState::Ready { channel });
265 }
266 Err(e) => {
267 let _ = slot
268 .state_tx
269 .send(BridgeState::Degraded(format!("Initial start failed: {e}")));
270 }
271 }
272
273 self.spawn_health_monitor(Arc::clone(&slot)).await;
274
275 Ok(slot)
276 }
277
278 pub fn restart_slot(&self, broker_name: &str) {
280 if let Some(slot) = self.slots.get(broker_name) {
281 let _ = slot.state_tx.send(BridgeState::Restarting {
282 attempt: 0,
283 next_at: Instant::now(),
284 });
285 }
286 }
287
288 pub async fn refresh_slot_channel(&self, broker_name: &str) -> Result<(), CamelError> {
293 let slot = self
294 .slots
295 .get(broker_name)
296 .map(|s| Arc::clone(&*s))
297 .ok_or_else(|| {
298 CamelError::ProcessorError(format!("Unknown JMS broker '{}'", broker_name))
299 })?;
300
301 let channel = {
306 let guard = slot.process.lock().await;
307 let process = guard.as_ref().ok_or_else(|| {
308 CamelError::ProcessorError(format!(
309 "JMS broker '{}' has no running bridge process",
310 broker_name
311 ))
312 })?;
313 process.connect().await.map_err(|e| {
314 CamelError::ProcessorError(format!(
315 "JMS broker '{}' channel refresh failed: {}",
316 broker_name, e
317 ))
318 })?
319 };
320
321 let _ = slot.state_tx.send(BridgeState::Ready { channel });
322 Ok(())
323 }
324
325 pub fn begin_shutdown(&self) {
326 self.shutting_down.store(true, Ordering::SeqCst);
327 }
328
329 pub async fn shutdown(&self) -> Result<(), CamelError> {
331 self.begin_shutdown();
332 let names: Vec<String> = self.slots.iter().map(|e| e.key().clone()).collect();
333 let mut errors: Vec<String> = Vec::new();
334
335 for name in names {
336 if let Some((_, slot)) = self.slots.remove(&name) {
337 let _ = slot.state_tx.send(BridgeState::Stopped);
339
340 let process = {
342 let mut guard = slot.process.lock().await;
343 guard.take()
344 };
345 if let Some(p) = process
346 && let Err(e) = p.stop().await
347 {
348 errors.push(format!("broker '{}': process stop failed: {e}", slot.name));
349 }
350
351 let monitor_handle = {
353 let mut guard = slot.health_monitor_handle.lock().await;
354 guard.take()
355 };
356 if let Some(mut h) = monitor_handle
357 && tokio::time::timeout(Duration::from_secs(5), &mut h)
358 .await
359 .is_err()
360 {
361 h.abort();
362 let _ = h.await;
363 warn!(
364 "health monitor for '{}' did not stop in 5s; aborted",
365 slot.name
366 );
367 }
368 }
369 }
370
371 if errors.is_empty() {
372 Ok(())
373 } else {
374 Err(CamelError::ProcessorError(format!(
375 "JMS pool shutdown completed with {} error(s): {}",
376 errors.len(),
377 errors.join("; ")
378 )))
379 }
380 }
381
382 async fn spawn_health_monitor(&self, slot: Arc<BridgeSlot>) {
383 let health_interval = self.health_check_interval_ms;
384 let bridge_version = self.bridge_version.clone();
385 let bridge_cache_dir = self.bridge_cache_dir.clone();
386 let start_timeout_ms = self.bridge_start_timeout_ms;
387 let handle_ref = Arc::clone(&slot.health_monitor_handle);
388 let shutting_down = Arc::clone(&self.shutting_down);
389 let broker_name = slot.name.clone();
390
391 let handle = tokio::spawn(async move {
392 loop {
393 let state = slot.state_rx.borrow().clone();
394 match state {
395 BridgeState::Stopped => {
396 info!("Health monitor for '{}' exiting (Stopped)", slot.name);
397 break;
398 }
399 BridgeState::Ready { ref channel } => {
400 tokio::time::sleep(Duration::from_millis(health_interval)).await;
401 let mut client = BridgeServiceClient::new(channel.clone());
402 let health_timeout = Duration::from_secs(3);
403 match tokio::time::timeout(health_timeout, client.health(HealthRequest {}))
404 .await
405 {
406 Ok(Ok(_)) => {}
407 Ok(Err(e)) => {
408 warn!(
409 "Health check failed for broker '{}': {e}. Marking Degraded.",
410 slot.name
411 );
412 let _ = slot.state_tx.send(BridgeState::Degraded(e.to_string()));
413 }
414 Err(_) => {
415 let msg = format!(
416 "health RPC timed out after {}ms",
417 health_timeout.as_millis()
418 );
419 warn!(
420 "Health check timed out for broker '{}': {}. Marking Degraded.",
421 slot.name, msg
422 );
423 let _ = slot.state_tx.send(BridgeState::Degraded(msg));
424 }
425 }
426 }
427 BridgeState::Degraded(_) | BridgeState::Starting => {
428 if matches!(*slot.state_rx.borrow(), BridgeState::Stopped) {
429 break;
430 }
431 if shutting_down.load(Ordering::SeqCst) {
432 tracing::info!(
433 "Pool shutting down — not restarting bridge for broker '{}'",
434 broker_name
435 );
436 break;
437 }
438 let _ = slot.state_tx.send(BridgeState::Restarting {
439 attempt: 0,
440 next_at: Instant::now(),
441 });
442 }
443 BridgeState::Restarting { attempt, next_at } => {
444 if shutting_down.load(Ordering::SeqCst) {
445 tracing::info!(
446 "Pool shutting down — aborting restart for broker '{}'",
447 broker_name
448 );
449 break;
450 }
451
452 let now = Instant::now();
453 if now < next_at {
454 tokio::time::sleep(next_at - now).await;
455 }
456
457 info!(
458 "Restarting bridge for broker '{}' (attempt {})",
459 slot.name,
460 attempt + 1
461 );
462
463 if attempt >= MAX_RESTART_ATTEMPTS {
464 tracing::error!(
466 "Max restart attempts ({}) reached for broker '{}' — staying degraded",
467 attempt,
468 broker_name
469 );
470 let _ = slot.state_tx.send(BridgeState::Degraded(format!(
471 "max restart attempts ({}) exceeded",
472 attempt
473 )));
474 break;
475 }
476
477 let old_process = {
478 let mut guard = slot.process.lock().await;
479 guard.take()
480 };
481 if let Some(p) = old_process {
482 let _ = p.stop().await;
483 }
484
485 let start_result = Self::start_bridge_process(
486 &bridge_version,
487 &bridge_cache_dir,
488 start_timeout_ms,
489 &slot.broker_url,
490 &slot.broker_type,
491 &slot.credentials,
492 )
493 .await;
494
495 match start_result {
496 Ok((process, channel)) => {
497 if matches!(*slot.state_rx.borrow(), BridgeState::Stopped) {
500 let _ = process.stop().await;
501 break;
502 }
503 {
504 let mut guard = slot.process.lock().await;
505 *guard = Some(process);
506 }
507 let _ = slot.state_tx.send(BridgeState::Ready { channel });
508 info!("Broker '{}' bridge restarted successfully", slot.name);
509 }
510 Err(e) => {
511 if matches!(*slot.state_rx.borrow(), BridgeState::Stopped) {
513 break;
514 }
515 let delay_secs = std::cmp::min(5 * 2u64.pow(attempt), 120);
516 let next = Instant::now() + Duration::from_secs(delay_secs);
517 warn!(
518 "Failed to restart bridge for '{}' (attempt {}): {e}. Retry in {delay_secs}s",
519 slot.name,
520 attempt + 1
521 );
522 let _ = slot.state_tx.send(BridgeState::Restarting {
523 attempt: attempt + 1,
524 next_at: next,
525 });
526 }
527 }
528 }
529 }
530 }
531 });
532
533 let mut guard = handle_ref.lock().await;
535 *guard = Some(handle);
536 }
537
538 async fn start_bridge_process(
539 bridge_version: &str,
540 bridge_cache_dir: &std::path::Path,
541 start_timeout_ms: u64,
542 broker_url: &str,
543 broker_type: &BrokerType,
544 credentials: &Option<(String, String)>,
545 ) -> Result<(BridgeProcess, Channel), CamelError> {
546 info!(
547 "Starting JMS bridge process for {}...",
548 redact_url(broker_url)
549 );
550 let binary_path = ensure_binary(bridge_version, bridge_cache_dir)
551 .await
552 .map_err(|e| {
553 CamelError::ProcessorError(format!("JMS bridge binary unavailable: {e}"))
554 })?;
555
556 let process_config = BridgeProcessConfig::jms(
557 binary_path,
558 broker_url.to_string(),
559 broker_type.clone(),
560 credentials.as_ref().map(|(u, _)| u.clone()),
561 credentials
562 .as_ref()
563 .map(|(_, p)| camel_bridge::process::Redacted::new(p.clone())),
564 start_timeout_ms,
565 );
566
567 let total_timeout = Duration::from_millis(start_timeout_ms);
568 let result = tokio::time::timeout(total_timeout, async {
569 let (process, channel) = BridgeProcess::start_and_connect(&process_config)
570 .await
571 .map_err(|e| CamelError::ProcessorError(format!("JMS bridge start failed: {e}")))?;
572
573 wait_for_health(&channel, Duration::from_secs(10), |ch| {
574 let mut client = BridgeServiceClient::new(ch);
575 async move {
576 let resp = client.health(HealthRequest {}).await?;
577 Ok(resp.into_inner().healthy)
578 }
579 })
580 .await
581 .map_err(|e| {
582 CamelError::ProcessorError(format!("JMS bridge health check failed: {e}"))
583 })?;
584
585 Ok::<(BridgeProcess, Channel), CamelError>((process, channel))
586 })
587 .await
588 .map_err(|_| {
589 CamelError::ProcessorError(format!(
590 "JMS bridge start timed out after {}ms",
591 start_timeout_ms
592 ))
593 })??;
594
595 Ok(result)
596 }
597}
598
599impl Drop for JmsBridgePool {
602 fn drop(&mut self) {
603 self.shutting_down.store(true, Ordering::SeqCst);
604
605 let slots: Vec<(String, Arc<BridgeSlot>)> = self
607 .slots
608 .iter()
609 .map(|e| (e.key().clone(), e.value().clone()))
610 .collect();
611
612 if slots.is_empty() {
613 return;
614 }
615
616 match tokio::runtime::Handle::try_current() {
617 Ok(handle) => {
618 drop(handle.spawn(async move {
621 for (_name, slot) in slots {
622 let _ = slot.state_tx.send(BridgeState::Stopped);
624
625 let process = {
627 let mut guard = slot.process.lock().await;
628 guard.take()
629 };
630 if let Some(p) = process {
631 let _ = p.stop().await;
632 }
633
634 let monitor_handle = {
636 let mut guard = slot.health_monitor_handle.lock().await;
637 guard.take()
638 };
639 if let Some(mut h) = monitor_handle
640 && tokio::time::timeout(Duration::from_secs(5), &mut h)
641 .await
642 .is_err()
643 {
644 h.abort();
645 let _ = h.await;
646 warn!(
647 "health monitor for '{}' did not stop in 5s; aborted",
648 slot.name
649 );
650 }
651 }
652 }));
653 }
654 Err(_) => {
655 warn!("JmsBridgePool dropped outside tokio runtime; bridges not cleaned up");
656 }
657 }
658 }
659}
660
661#[derive(Clone)]
664pub struct JmsComponent {
665 scheme: String,
666 pool: Arc<JmsBridgePool>,
667}
668
669impl JmsComponent {
670 pub fn with_scheme(scheme: impl Into<String>, pool: Arc<JmsBridgePool>) -> Self {
671 Self {
672 scheme: scheme.into(),
673 pool,
674 }
675 }
676
677 pub fn scheme(&self) -> &str {
678 &self.scheme
679 }
680
681 #[cfg(test)]
683 pub async fn send_for_test(
684 &self,
685 destination: &str,
686 body: &[u8],
687 content_type: &str,
688 ) -> Result<String, CamelError> {
689 let broker_name = self.pool.resolve_broker_name(None)?;
690 let slot = self.pool.get_or_create_slot(&broker_name).await?;
691 let channel = match &*slot.state_rx.borrow() {
692 BridgeState::Ready { channel } => channel.clone(),
693 other => {
694 return Err(CamelError::ProcessorError(format!(
695 "Bridge not ready: {:?}",
696 other
697 )));
698 }
699 };
700 let mut client = BridgeServiceClient::new(channel);
701 let r = client
702 .send(crate::proto::SendRequest {
703 destination: destination.to_string(),
704 body: body.to_vec(),
705 headers: Default::default(),
706 content_type: content_type.to_string(),
707 })
708 .await
709 .map_err(|e| CamelError::ProcessorError(format!("test send error: {e}")))?;
710 Ok(r.into_inner().message_id)
711 }
712}
713
714impl Component for JmsComponent {
715 fn scheme(&self) -> &str {
716 &self.scheme
717 }
718
719 fn metadata(&self) -> camel_component_api::ComponentMetadata {
720 crate::metadata::JmsMetadataDescriptor::metadata()
721 }
722
723 fn create_endpoint(
724 &self,
725 uri: &str,
726 ctx: &dyn camel_component_api::ComponentContext,
727 ) -> Result<Box<dyn Endpoint>, CamelError> {
728 let endpoint_config = JmsEndpointConfig::from_uri(uri)?;
729 let broker_name = self
730 .pool
731 .resolve_broker_name(endpoint_config.broker_name.as_deref())?;
732 let resolved_broker_type = self.pool.resolve_broker_type(&self.scheme, &broker_name);
733
734 let health_check = JmsHealthCheck::new(Arc::clone(&self.pool), broker_name.clone());
735 ctx.register_current_route_health_check(Arc::new(health_check));
736
737 Ok(Box::new(JmsEndpoint {
738 pool: Arc::clone(&self.pool),
739 uri: uri.to_string(),
740 broker_name,
741 resolved_broker_type,
742 endpoint_config,
743 }))
744 }
745}
746
747struct JmsEndpoint {
750 pool: Arc<JmsBridgePool>,
751 uri: String,
752 broker_name: String,
753 resolved_broker_type: BrokerType,
754 endpoint_config: JmsEndpointConfig,
755}
756
757impl Endpoint for JmsEndpoint {
758 fn uri(&self) -> &str {
759 &self.uri
760 }
761
762 fn create_producer(
763 &self,
764 rt: Arc<dyn camel_component_api::RuntimeObservability>,
765 _ctx: &ProducerContext,
766 ) -> Result<BoxProcessor, CamelError> {
767 Ok(BoxProcessor::new(LazyJmsProducer {
768 pool: Arc::clone(&self.pool),
769 broker_name: self.broker_name.clone(),
770 endpoint_config: self.endpoint_config.clone(),
771 resolved_broker_type: self.resolved_broker_type.clone(),
772 runtime: rt,
773 }))
774 }
775
776 fn create_consumer(
777 &self,
778 rt: Arc<dyn camel_component_api::RuntimeObservability>,
779 ) -> Result<Box<dyn Consumer>, CamelError> {
780 Ok(Box::new(JmsConsumer::new(
781 Arc::clone(&self.pool),
782 self.broker_name.clone(),
783 self.endpoint_config.clone(),
784 self.pool.reconnect.clone(),
785 rt,
786 )))
787 }
788}
789
790#[derive(Clone)]
791struct LazyJmsProducer {
792 pool: Arc<JmsBridgePool>,
793 broker_name: String,
794 endpoint_config: JmsEndpointConfig,
795 #[allow(dead_code)]
796 resolved_broker_type: BrokerType,
797 #[allow(dead_code)]
800 runtime: Arc<dyn camel_component_api::RuntimeObservability>,
801}
802
803impl Service<Exchange> for LazyJmsProducer {
804 type Response = Exchange;
805 type Error = CamelError;
806 type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
807
808 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
809 if let Some(slot) = self.pool.slots.get(&self.broker_name) {
812 match &*slot.state_rx.borrow() {
813 BridgeState::Ready { .. } => return Poll::Ready(Ok(())),
814 BridgeState::Starting | BridgeState::Restarting { .. } => {
815 let waker = cx.waker().clone();
821 let mut rx = slot.state_rx.clone();
822 if let Ok(handle) = tokio::runtime::Handle::try_current() {
823 handle.spawn(async move {
824 let _ = rx.changed().await;
825 waker.wake();
826 });
827 } else {
828 waker.wake_by_ref();
829 }
830 return Poll::Pending;
831 }
832 BridgeState::Degraded(reason) => {
833 return Poll::Ready(Err(CamelError::ProcessorError(format!(
834 "JMS broker '{}' is degraded: {}",
835 self.broker_name, reason
836 ))));
837 }
838 BridgeState::Stopped => {
839 return Poll::Ready(Err(CamelError::ConsumerStopping));
840 }
841 }
842 }
843 Poll::Ready(Ok(()))
844 }
845
846 fn call(&mut self, exchange: Exchange) -> Self::Future {
847 let pool = Arc::clone(&self.pool);
848 let broker_name = self.broker_name.clone();
849 let endpoint_config = self.endpoint_config.clone();
850
851 Box::pin(async move {
852 let slot = pool.get_or_create_slot(&broker_name).await?;
853 let mut rx = slot.state_rx.clone();
854
855 loop {
856 let state = rx.borrow().clone();
857 match state {
858 BridgeState::Ready { channel } => {
859 let mut producer = JmsProducer::new(channel, endpoint_config.clone());
860 match producer.call(exchange).await {
861 Ok(done) => return Ok(done),
862 Err(first_err) if is_bridge_transport_error(&first_err) => {
863 warn!(
864 broker = %broker_name,
865 error = %first_err,
866 "JMS send transport error; refreshing channel (no automatic resend)"
867 );
868
869 if let Err(refresh_err) =
870 pool.refresh_slot_channel(&broker_name).await
871 {
872 warn!(
873 broker = %broker_name,
874 error = %refresh_err,
875 "JMS channel refresh failed; requesting bridge restart"
876 );
877 pool.restart_slot(&broker_name);
878 }
879
880 return Err(first_err);
885 }
886 Err(other_err) => return Err(other_err),
887 }
888 }
889 BridgeState::Degraded(reason) => {
890 return Err(CamelError::ProcessorError(format!(
891 "JMS broker '{}' is degraded: {}",
892 broker_name, reason
893 )));
894 }
895 BridgeState::Stopped => {
896 return Err(CamelError::ProcessorError(format!(
897 "JMS broker '{}' is stopped",
898 broker_name
899 )));
900 }
901 BridgeState::Starting | BridgeState::Restarting { .. } => {
902 if rx.changed().await.is_err() {
903 return Err(CamelError::ProcessorError(format!(
904 "JMS broker '{}' state channel closed",
905 broker_name
906 )));
907 }
908 }
909 }
910 }
911 })
912 }
913}
914
915fn redact_url(url: &str) -> String {
920 if let Some(pos) = url.find("://") {
922 let scheme = &url[..pos + 3]; let rest = &url[pos + 3..];
924 if let Some(at_pos) = rest.find('@') {
926 return format!("{}***@{}", scheme, &rest[at_pos + 1..]);
927 }
928 }
929 url.to_string()
930}
931
932pub fn is_bridge_transport_error(err: &CamelError) -> bool {
933 match err {
938 CamelError::ProcessorError(msg) => msg.starts_with(BRIDGE_TRANSPORT_ERROR_PREFIX),
939 _ => false,
940 }
941}
942
943#[cfg(test)]
946mod tests {
947 use camel_component_api::test_support::PanicRuntimeObservability;
948 fn test_rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
949 std::sync::Arc::new(PanicRuntimeObservability)
950 }
951 fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
952 std::sync::Arc::new(PanicRuntimeObservability)
953 }
954
955 use super::*;
956 use crate::config::{BrokerConfig, JmsPoolConfig};
957 use std::collections::HashMap;
958
959 #[test]
960 fn from_config_accepts_empty_brokers() {
961 let pool_config = JmsPoolConfig::default();
962 let result = JmsBridgePool::from_config(pool_config);
963 assert!(result.is_ok());
964 }
965
966 #[test]
967 fn resolve_broker_name_with_explicit_name() {
968 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
969 "tcp://localhost:61616",
970 BrokerType::ActiveMq,
971 ))
972 .unwrap();
973 assert_eq!(
974 pool.resolve_broker_name(Some("default")).unwrap(),
975 "default"
976 );
977 }
978
979 #[test]
980 fn resolve_broker_name_default() {
981 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
982 "tcp://localhost:61616",
983 BrokerType::ActiveMq,
984 ))
985 .unwrap();
986 assert_eq!(pool.resolve_broker_name(None).unwrap(), "default");
987 }
988
989 #[test]
990 fn resolve_broker_name_unknown_returns_error() {
991 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
992 "tcp://localhost:61616",
993 BrokerType::ActiveMq,
994 ))
995 .unwrap();
996 let err = pool.resolve_broker_name(Some("unknown")).unwrap_err();
997 assert!(
998 err.to_string().contains("Unknown JMS broker 'unknown'"),
999 "got: {}",
1000 err
1001 );
1002 }
1003
1004 #[test]
1005 fn resolve_broker_type_scheme_overrides() {
1006 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1007 "tcp://localhost:61616",
1008 BrokerType::Generic,
1009 ))
1010 .unwrap();
1011 assert_eq!(
1012 pool.resolve_broker_type("activemq", "default"),
1013 BrokerType::ActiveMq
1014 );
1015 assert_eq!(
1016 pool.resolve_broker_type("artemis", "default"),
1017 BrokerType::Artemis
1018 );
1019 assert_eq!(
1020 pool.resolve_broker_type("jms", "default"),
1021 BrokerType::Generic
1022 );
1023 }
1024
1025 #[test]
1026 fn resolve_broker_type_activemq_scheme_overrides_artemis_config() {
1027 let pool = JmsBridgePool::from_config(JmsPoolConfig {
1028 brokers: HashMap::from([(
1029 "main".to_string(),
1030 BrokerConfig {
1031 broker_url: "tcp://localhost:61616".to_string(),
1032 broker_type: BrokerType::Artemis,
1033 username: None,
1034 password: None,
1035 },
1036 )]),
1037 ..JmsPoolConfig::default()
1038 })
1039 .unwrap();
1040 assert_eq!(
1041 pool.resolve_broker_type("activemq", "main"),
1042 BrokerType::ActiveMq
1043 );
1044 assert_eq!(pool.resolve_broker_type("jms", "main"), BrokerType::Artemis);
1045 }
1046
1047 #[test]
1048 fn create_endpoint_resolves_broker() {
1049 let pool = Arc::new(
1050 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1051 "tcp://localhost:61616",
1052 BrokerType::ActiveMq,
1053 ))
1054 .unwrap(),
1055 );
1056 let component = JmsComponent::with_scheme("jms", pool);
1057 let endpoint = component.create_endpoint(
1058 "jms:queue:orders",
1059 &camel_component_api::NoOpComponentContext,
1060 );
1061 assert!(endpoint.is_ok(), "got: {:?}", endpoint.err());
1062 }
1063
1064 #[test]
1065 fn create_endpoint_rejects_wrong_scheme() {
1066 let pool = Arc::new(
1067 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1068 "tcp://localhost:61616",
1069 BrokerType::ActiveMq,
1070 ))
1071 .unwrap(),
1072 );
1073 let component = JmsComponent::with_scheme("jms", pool);
1074 let err = component
1075 .create_endpoint("kafka:orders", &camel_component_api::NoOpComponentContext)
1076 .err()
1077 .unwrap();
1078 assert!(
1079 err.to_string()
1080 .contains("expected scheme 'jms', 'activemq', or 'artemis'"),
1081 "got: {}",
1082 err
1083 );
1084 }
1085
1086 #[test]
1087 fn create_endpoint_with_explicit_broker_param() {
1088 let pool = Arc::new(
1089 JmsBridgePool::from_config(JmsPoolConfig {
1090 brokers: HashMap::from([
1091 (
1092 "primary".to_string(),
1093 BrokerConfig {
1094 broker_url: "tcp://primary:61616".to_string(),
1095 broker_type: BrokerType::ActiveMq,
1096 username: None,
1097 password: None,
1098 },
1099 ),
1100 (
1101 "secondary".to_string(),
1102 BrokerConfig {
1103 broker_url: "tcp://secondary:61616".to_string(),
1104 broker_type: BrokerType::Artemis,
1105 username: None,
1106 password: None,
1107 },
1108 ),
1109 ]),
1110 ..JmsPoolConfig::default()
1111 })
1112 .unwrap(),
1113 );
1114 let component = JmsComponent::with_scheme("jms", Arc::clone(&pool));
1115 let endpoint = component.create_endpoint(
1116 "jms:queue:orders?broker=secondary",
1117 &camel_component_api::NoOpComponentContext,
1118 );
1119 assert!(endpoint.is_ok(), "got: {:?}", endpoint.err());
1120 }
1121
1122 #[tokio::test]
1123 #[allow(clippy::await_holding_lock)]
1124 async fn concurrent_get_or_create_slot_no_deadlock() {
1125 use tokio::time::timeout;
1126
1127 let _env_lock_guard = crate::BRIDGE_ENV_LOCK.lock().unwrap();
1129
1130 struct EnvGuard {
1131 key: &'static str,
1132 prev: Option<std::ffi::OsString>,
1133 }
1134 impl Drop for EnvGuard {
1135 fn drop(&mut self) {
1136 if let Some(v) = &self.prev {
1137 unsafe { std::env::set_var(self.key, v) };
1139 } else {
1140 unsafe { std::env::remove_var(self.key) };
1142 }
1143 }
1144 }
1145
1146 let env_key = "CAMEL_JMS_BRIDGE_BINARY_PATH";
1147 let _guard = EnvGuard {
1148 key: env_key,
1149 prev: std::env::var_os(env_key),
1150 };
1151 unsafe { std::env::set_var(env_key, "/bin/false") };
1153
1154 let pool = Arc::new(
1155 JmsBridgePool::from_config(JmsPoolConfig {
1156 brokers: HashMap::from([(
1157 "test".to_string(),
1158 BrokerConfig {
1159 broker_url: "tcp://localhost:61616".to_string(),
1160 broker_type: BrokerType::ActiveMq,
1161 username: None,
1162 password: None,
1163 },
1164 )]),
1165 bridge_start_timeout_ms: 100,
1166 ..JmsPoolConfig::default()
1167 })
1168 .unwrap(),
1169 );
1170
1171 let handles: Vec<_> = (0..5)
1172 .map(|_| {
1173 let pool = Arc::clone(&pool);
1174 tokio::spawn(async move {
1175 let _ = pool.get_or_create_slot("test").await;
1176 })
1177 })
1178 .collect();
1179
1180 let result = timeout(Duration::from_secs(5), async {
1181 for h in handles {
1182 let _ = h.await;
1183 }
1184 })
1185 .await;
1186
1187 assert!(result.is_ok(), "Concurrent get_or_create_slot deadlocked!");
1188 }
1189
1190 #[tokio::test]
1191 #[allow(clippy::await_holding_lock)]
1192 async fn lazy_producer_reports_degraded_when_bridge_start_fails() {
1193 use tower::Service;
1194
1195 let _env_lock_guard = crate::BRIDGE_ENV_LOCK.lock().unwrap();
1197
1198 struct EnvGuard {
1199 key: &'static str,
1200 prev: Option<std::ffi::OsString>,
1201 }
1202 impl Drop for EnvGuard {
1203 fn drop(&mut self) {
1204 if let Some(v) = &self.prev {
1205 unsafe { std::env::set_var(self.key, v) };
1207 } else {
1208 unsafe { std::env::remove_var(self.key) };
1210 }
1211 }
1212 }
1213
1214 let env_key = "CAMEL_JMS_BRIDGE_BINARY_PATH";
1215 let _guard = EnvGuard {
1216 key: env_key,
1217 prev: std::env::var_os(env_key),
1218 };
1219 unsafe { std::env::set_var(env_key, "/bin/false") };
1221
1222 let pool = Arc::new(
1223 JmsBridgePool::from_config(JmsPoolConfig {
1224 brokers: HashMap::from([(
1225 "default".to_string(),
1226 BrokerConfig {
1227 broker_url: "tcp://localhost:61616".to_string(),
1228 broker_type: BrokerType::ActiveMq,
1229 username: None,
1230 password: None,
1231 },
1232 )]),
1233 bridge_start_timeout_ms: 100,
1234 ..JmsPoolConfig::default()
1235 })
1236 .unwrap(),
1237 );
1238
1239 let component = JmsComponent::with_scheme("jms", pool);
1240 let endpoint = component
1241 .create_endpoint(
1242 "jms:queue:orders",
1243 &camel_component_api::NoOpComponentContext,
1244 )
1245 .unwrap();
1246 let mut producer = endpoint
1247 .create_producer(rt(), &camel_component_api::ProducerContext::default())
1248 .unwrap();
1249
1250 let mut exchange = Exchange::default();
1251 exchange.input.body = camel_component_api::Body::Text("hello".to_string());
1252
1253 let err = producer.call(exchange).await.unwrap_err();
1254 assert!(err.to_string().contains("is degraded"), "got: {}", err);
1255 }
1256
1257 #[tokio::test]
1261 async fn lazy_producer_requests_restart_when_refresh_unavailable() {
1262 use tokio::sync::watch;
1263 use tonic::transport::Endpoint as TonicEndpoint;
1264 use tower::Service;
1265
1266 let dead_channel = TonicEndpoint::from_static("http://127.0.0.1:1").connect_lazy();
1269
1270 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1271 channel: dead_channel.clone(),
1272 });
1273
1274 let pool = Arc::new(
1275 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1276 "tcp://localhost:61616",
1277 BrokerType::ActiveMq,
1278 ))
1279 .unwrap(),
1280 );
1281
1282 let slot = Arc::new(BridgeSlot {
1284 name: "default".to_string(),
1285 broker_url: "tcp://localhost:61616".to_string(),
1286 broker_type: BrokerType::ActiveMq,
1287 credentials: None,
1288 state_rx: state_rx.clone(),
1289 state_tx: state_tx.clone(),
1290 process: Arc::new(tokio::sync::Mutex::new(None)),
1291 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1292 });
1293 pool.slots.insert("default".to_string(), Arc::clone(&slot));
1294
1295 let endpoint_config =
1296 crate::config::JmsEndpointConfig::from_uri("jms:queue:test-retry").unwrap();
1297
1298 let mut producer = LazyJmsProducer {
1299 pool: Arc::clone(&pool),
1300 broker_name: "default".to_string(),
1301 endpoint_config,
1302 resolved_broker_type: BrokerType::ActiveMq,
1303 runtime: test_rt(),
1304 };
1305
1306 let mut exchange = Exchange::default();
1307 exchange.input.body = camel_component_api::Body::Text("hello".to_string());
1308
1309 let result = producer.call(exchange).await;
1311 assert!(result.is_err(), "expected send to fail");
1312
1313 let state_after = state_rx.borrow().clone();
1316 assert!(
1317 matches!(state_after, BridgeState::Restarting { .. }),
1318 "slot must enter Restarting when refresh is unavailable; got: {:?}",
1319 state_after
1320 );
1321 }
1322
1323 #[test]
1326 fn transport_error_detects_send_error() {
1327 let err = CamelError::ProcessorError(format!(
1328 "{}send error: connection refused",
1329 BRIDGE_TRANSPORT_ERROR_PREFIX
1330 ));
1331 assert!(
1332 is_bridge_transport_error(&err),
1333 "send error must be classified as transport"
1334 );
1335 }
1336
1337 #[test]
1338 fn transport_error_detects_subscribe_error() {
1339 let err = CamelError::ProcessorError(format!(
1340 "{}subscribe error: stream reset",
1341 BRIDGE_TRANSPORT_ERROR_PREFIX
1342 ));
1343 assert!(
1344 is_bridge_transport_error(&err),
1345 "subscribe error must be classified as transport"
1346 );
1347 }
1348
1349 #[test]
1350 fn transport_error_rejects_business_errors() {
1351 let err = CamelError::ProcessorError("JMS broker 'main' is degraded: timeout".to_string());
1352 assert!(
1353 !is_bridge_transport_error(&err),
1354 "degraded state error must NOT be transport"
1355 );
1356 }
1357
1358 #[test]
1359 fn transport_error_rejects_config_errors() {
1360 let err = CamelError::Config("bridge_start_timeout_ms must be > 0".to_string());
1361 assert!(
1362 !is_bridge_transport_error(&err),
1363 "config error must NOT be transport"
1364 );
1365 }
1366
1367 #[test]
1368 fn transport_error_prefix_is_used_by_producer_and_consumer() {
1369 assert!(
1372 BRIDGE_TRANSPORT_ERROR_PREFIX.starts_with("JMS gRPC "),
1373 "prefix must start with 'JMS gRPC '"
1374 );
1375 }
1376
1377 #[tokio::test]
1380 async fn pool_enforces_max_bridges_limit() {
1381 use tokio::sync::watch;
1382
1383 let pool = Arc::new(
1384 JmsBridgePool::from_config(JmsPoolConfig {
1385 brokers: HashMap::from([
1386 (
1387 "b1".to_string(),
1388 BrokerConfig {
1389 broker_url: "tcp://b1:61616".to_string(),
1390 broker_type: BrokerType::ActiveMq,
1391 username: None,
1392 password: None,
1393 },
1394 ),
1395 (
1396 "b2".to_string(),
1397 BrokerConfig {
1398 broker_url: "tcp://b2:61616".to_string(),
1399 broker_type: BrokerType::ActiveMq,
1400 username: None,
1401 password: None,
1402 },
1403 ),
1404 (
1405 "b3".to_string(),
1406 BrokerConfig {
1407 broker_url: "tcp://b3:61616".to_string(),
1408 broker_type: BrokerType::ActiveMq,
1409 username: None,
1410 password: None,
1411 },
1412 ),
1413 ]),
1414 max_bridges: 2,
1415 ..JmsPoolConfig::default()
1416 })
1417 .unwrap(),
1418 );
1419
1420 for name in &["b1", "b2"] {
1423 let (state_tx, state_rx) = watch::channel(BridgeState::Starting);
1424 let slot = Arc::new(BridgeSlot {
1425 name: name.to_string(),
1426 broker_url: format!("tcp://{name}:61616"),
1427 broker_type: BrokerType::ActiveMq,
1428 credentials: None,
1429 state_rx,
1430 state_tx,
1431 process: Arc::new(tokio::sync::Mutex::new(None)),
1432 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1433 });
1434 pool.slots.insert(name.to_string(), slot);
1435 }
1436
1437 let err = pool.get_or_create_slot("b3").await.unwrap_err();
1439 assert!(
1440 err.to_string().contains("max_bridges"),
1441 "expected max_bridges error, got: {}",
1442 err
1443 );
1444 }
1445
1446 #[tokio::test]
1447 async fn pool_allows_slot_when_below_max_bridges() {
1448 use tokio::sync::watch;
1449
1450 let pool = Arc::new(
1451 JmsBridgePool::from_config(JmsPoolConfig {
1452 brokers: HashMap::from([(
1453 "b1".to_string(),
1454 BrokerConfig {
1455 broker_url: "tcp://b1:61616".to_string(),
1456 broker_type: BrokerType::ActiveMq,
1457 username: None,
1458 password: None,
1459 },
1460 )]),
1461 max_bridges: 2,
1462 bridge_start_timeout_ms: 100,
1463 ..JmsPoolConfig::default()
1464 })
1465 .unwrap(),
1466 );
1467
1468 let (state_tx, state_rx) = watch::channel(BridgeState::Degraded("test".to_string()));
1470 let slot = Arc::new(BridgeSlot {
1471 name: "b1".to_string(),
1472 broker_url: "tcp://b1:61616".to_string(),
1473 broker_type: BrokerType::ActiveMq,
1474 credentials: None,
1475 state_rx,
1476 state_tx,
1477 process: Arc::new(tokio::sync::Mutex::new(None)),
1478 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1479 });
1480 pool.slots.insert("b1".to_string(), slot);
1481
1482 let result = pool.get_or_create_slot("b1").await;
1485 assert!(result.is_ok(), "existing slot must be returned");
1486 }
1487
1488 #[tokio::test]
1491 async fn poll_ready_returns_pending_when_starting() {
1492 use tokio::sync::watch;
1493 use tower::Service;
1494
1495 let pool = Arc::new(
1496 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1497 "tcp://localhost:61616",
1498 BrokerType::ActiveMq,
1499 ))
1500 .unwrap(),
1501 );
1502
1503 let (state_tx, state_rx) = watch::channel(BridgeState::Starting);
1504 let slot = Arc::new(BridgeSlot {
1505 name: "default".to_string(),
1506 broker_url: "tcp://localhost:61616".to_string(),
1507 broker_type: BrokerType::ActiveMq,
1508 credentials: None,
1509 state_rx,
1510 state_tx,
1511 process: Arc::new(tokio::sync::Mutex::new(None)),
1512 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1513 });
1514 pool.slots.insert("default".to_string(), slot);
1515
1516 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1517 let mut producer = LazyJmsProducer {
1518 pool: Arc::clone(&pool),
1519 broker_name: "default".to_string(),
1520 endpoint_config,
1521 resolved_broker_type: BrokerType::ActiveMq,
1522 runtime: test_rt(),
1523 };
1524
1525 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1526 assert!(
1527 matches!(result, Poll::Pending),
1528 "poll_ready must be Pending when Starting; got: {:?}",
1529 result
1530 );
1531 }
1532
1533 #[tokio::test]
1534 async fn poll_ready_returns_error_when_degraded() {
1535 use tokio::sync::watch;
1536 use tower::Service;
1537
1538 let pool = Arc::new(
1539 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1540 "tcp://localhost:61616",
1541 BrokerType::ActiveMq,
1542 ))
1543 .unwrap(),
1544 );
1545
1546 let (state_tx, state_rx) =
1547 watch::channel(BridgeState::Degraded("health check failed".to_string()));
1548 let slot = Arc::new(BridgeSlot {
1549 name: "default".to_string(),
1550 broker_url: "tcp://localhost:61616".to_string(),
1551 broker_type: BrokerType::ActiveMq,
1552 credentials: None,
1553 state_rx,
1554 state_tx,
1555 process: Arc::new(tokio::sync::Mutex::new(None)),
1556 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1557 });
1558 pool.slots.insert("default".to_string(), slot);
1559
1560 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1561 let mut producer = LazyJmsProducer {
1562 pool: Arc::clone(&pool),
1563 broker_name: "default".to_string(),
1564 endpoint_config,
1565 resolved_broker_type: BrokerType::ActiveMq,
1566 runtime: test_rt(),
1567 };
1568
1569 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1570 assert!(
1571 matches!(result, Poll::Ready(Err(_))),
1572 "poll_ready must be Err when Degraded; got: {:?}",
1573 result
1574 );
1575 let err_msg = match result {
1576 Poll::Ready(Err(e)) => e.to_string(),
1577 _ => unreachable!(),
1578 };
1579 assert!(
1580 err_msg.contains("degraded"),
1581 "error must mention degraded: {}",
1582 err_msg
1583 );
1584 }
1585
1586 #[tokio::test]
1587 async fn lazy_producer_poll_ready_returns_consumer_stopping_on_stopped() {
1588 use tokio::sync::watch;
1589 use tower::Service;
1590
1591 let pool = Arc::new(
1592 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1593 "tcp://localhost:61616",
1594 BrokerType::ActiveMq,
1595 ))
1596 .unwrap(),
1597 );
1598
1599 let (state_tx, state_rx) = watch::channel(BridgeState::Stopped);
1600 let slot = Arc::new(BridgeSlot {
1601 name: "default".to_string(),
1602 broker_url: "tcp://localhost:61616".to_string(),
1603 broker_type: BrokerType::ActiveMq,
1604 credentials: None,
1605 state_rx,
1606 state_tx,
1607 process: Arc::new(tokio::sync::Mutex::new(None)),
1608 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1609 });
1610 pool.slots.insert("default".to_string(), slot);
1611
1612 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1613 let mut producer = LazyJmsProducer {
1614 pool: Arc::clone(&pool),
1615 broker_name: "default".to_string(),
1616 endpoint_config,
1617 resolved_broker_type: BrokerType::ActiveMq,
1618 runtime: test_rt(),
1619 };
1620
1621 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1622 assert!(
1623 matches!(result, Poll::Ready(Err(CamelError::ConsumerStopping))),
1624 "poll_ready must return ConsumerStopping when Stopped; got: {:?}",
1625 result
1626 );
1627 }
1628
1629 #[tokio::test]
1630 async fn poll_ready_returns_ready_when_slot_ready() {
1631 use tokio::sync::watch;
1632 use tonic::transport::Endpoint as TonicEndpoint;
1633 use tower::Service;
1634
1635 let pool = Arc::new(
1636 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1637 "tcp://localhost:61616",
1638 BrokerType::ActiveMq,
1639 ))
1640 .unwrap(),
1641 );
1642
1643 let lazy_channel = TonicEndpoint::from_static("http://127.0.0.1:1").connect_lazy();
1644 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1645 channel: lazy_channel,
1646 });
1647 let slot = Arc::new(BridgeSlot {
1648 name: "default".to_string(),
1649 broker_url: "tcp://localhost:61616".to_string(),
1650 broker_type: BrokerType::ActiveMq,
1651 credentials: None,
1652 state_rx,
1653 state_tx,
1654 process: Arc::new(tokio::sync::Mutex::new(None)),
1655 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1656 });
1657 pool.slots.insert("default".to_string(), slot);
1658
1659 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1660 let mut producer = LazyJmsProducer {
1661 pool: Arc::clone(&pool),
1662 broker_name: "default".to_string(),
1663 endpoint_config,
1664 resolved_broker_type: BrokerType::ActiveMq,
1665 runtime: test_rt(),
1666 };
1667
1668 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1669 assert!(
1670 matches!(result, Poll::Ready(Ok(()))),
1671 "poll_ready must be Ready(Ok) when bridge is Ready; got: {:?}",
1672 result
1673 );
1674 }
1675
1676 #[tokio::test]
1677 async fn poll_ready_returns_ready_when_no_slot_exists() {
1678 use tower::Service;
1679
1680 let pool = Arc::new(
1681 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1682 "tcp://localhost:61616",
1683 BrokerType::ActiveMq,
1684 ))
1685 .unwrap(),
1686 );
1687
1688 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1689 let mut producer = LazyJmsProducer {
1690 pool: Arc::clone(&pool),
1691 broker_name: "default".to_string(),
1692 endpoint_config,
1693 resolved_broker_type: BrokerType::ActiveMq,
1694 runtime: test_rt(),
1695 };
1696
1697 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1699 assert!(
1700 matches!(result, Poll::Ready(Ok(()))),
1701 "poll_ready must be Ready(Ok) when no slot exists; got: {:?}",
1702 result
1703 );
1704 }
1705
1706 #[tokio::test]
1709 async fn pool_shutdown_awaits_health_monitor() {
1710 use tokio::sync::watch;
1711
1712 let pool = Arc::new(
1713 JmsBridgePool::from_config(JmsPoolConfig {
1714 brokers: HashMap::from([(
1715 "default".to_string(),
1716 BrokerConfig {
1717 broker_url: "tcp://localhost:61616".to_string(),
1718 broker_type: BrokerType::ActiveMq,
1719 username: None,
1720 password: None,
1721 },
1722 )]),
1723 health_check_interval_ms: 100,
1724 ..JmsPoolConfig::default()
1725 })
1726 .unwrap(),
1727 );
1728
1729 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1731 channel: tonic::transport::Endpoint::from_static("http://127.0.0.1:1").connect_lazy(),
1732 });
1733 let monitor_handle_ref: Arc<tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>> =
1734 Arc::new(tokio::sync::Mutex::new(None));
1735
1736 let state_rx_clone = state_rx.clone();
1738 let handle = tokio::spawn(async move {
1739 loop {
1740 if matches!(*state_rx_clone.borrow(), BridgeState::Stopped) {
1741 break;
1742 }
1743 tokio::time::sleep(Duration::from_millis(50)).await;
1744 }
1745 });
1746 *monitor_handle_ref.lock().await = Some(handle);
1747
1748 let slot = Arc::new(BridgeSlot {
1749 name: "default".to_string(),
1750 broker_url: "tcp://localhost:61616".to_string(),
1751 broker_type: BrokerType::ActiveMq,
1752 credentials: None,
1753 state_rx,
1754 state_tx,
1755 process: Arc::new(tokio::sync::Mutex::new(None)),
1756 health_monitor_handle: monitor_handle_ref,
1757 });
1758 pool.slots.insert("default".to_string(), slot);
1759
1760 let result = pool.shutdown().await;
1762 let _ = result;
1764 }
1765
1766 #[tokio::test]
1767 async fn health_monitor_handle_stored_after_spawn() {
1768 use tokio::sync::watch;
1769
1770 let pool = Arc::new(
1771 JmsBridgePool::from_config(JmsPoolConfig {
1772 brokers: HashMap::from([(
1773 "default".to_string(),
1774 BrokerConfig {
1775 broker_url: "tcp://localhost:61616".to_string(),
1776 broker_type: BrokerType::ActiveMq,
1777 username: None,
1778 password: None,
1779 },
1780 )]),
1781 health_check_interval_ms: 100,
1782 bridge_start_timeout_ms: 100,
1783 ..JmsPoolConfig::default()
1784 })
1785 .unwrap(),
1786 );
1787
1788 let (state_tx, state_rx) = watch::channel(BridgeState::Stopped);
1789 let slot = Arc::new(BridgeSlot {
1790 name: "default".to_string(),
1791 broker_url: "tcp://localhost:61616".to_string(),
1792 broker_type: BrokerType::ActiveMq,
1793 credentials: None,
1794 state_rx,
1795 state_tx,
1796 process: Arc::new(tokio::sync::Mutex::new(None)),
1797 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1798 });
1799 pool.slots.insert("default".to_string(), Arc::clone(&slot));
1800
1801 pool.spawn_health_monitor(Arc::clone(&slot)).await;
1802
1803 tokio::time::sleep(Duration::from_millis(50)).await;
1804 let guard = slot.health_monitor_handle.lock().await;
1805 assert!(
1806 guard.is_some(),
1807 "health monitor handle must be stored after spawn_health_monitor"
1808 );
1809 }
1810
1811 #[test]
1814 fn redact_url_strips_userinfo_with_password() {
1815 assert_eq!(
1816 redact_url("tcp://admin:s3cret@broker:61616"),
1817 "tcp://***@broker:61616"
1818 );
1819 }
1820
1821 #[test]
1822 fn redact_url_strips_userinfo_without_password() {
1823 assert_eq!(
1824 redact_url("tcp://admin@broker:61616"),
1825 "tcp://***@broker:61616"
1826 );
1827 }
1828
1829 #[test]
1830 fn redact_url_passes_clean_url_unchanged() {
1831 assert_eq!(redact_url("tcp://localhost:61616"), "tcp://localhost:61616");
1832 }
1833
1834 #[test]
1835 fn redact_url_handles_ssl_scheme() {
1836 assert_eq!(
1837 redact_url("ssl://user:pass@secure-broker:61617"),
1838 "ssl://***@secure-broker:61617"
1839 );
1840 }
1841
1842 #[tokio::test]
1845 async fn concurrent_slot_creation_respects_max_bridges() {
1846 let pool = Arc::new(
1847 JmsBridgePool::from_config(JmsPoolConfig {
1848 brokers: HashMap::from([
1849 (
1850 "b1".to_string(),
1851 BrokerConfig {
1852 broker_url: "tcp://b1:61616".to_string(),
1853 broker_type: BrokerType::ActiveMq,
1854 username: None,
1855 password: None,
1856 },
1857 ),
1858 (
1859 "b2".to_string(),
1860 BrokerConfig {
1861 broker_url: "tcp://b2:61616".to_string(),
1862 broker_type: BrokerType::ActiveMq,
1863 username: None,
1864 password: None,
1865 },
1866 ),
1867 (
1868 "b3".to_string(),
1869 BrokerConfig {
1870 broker_url: "tcp://b3:61616".to_string(),
1871 broker_type: BrokerType::ActiveMq,
1872 username: None,
1873 password: None,
1874 },
1875 ),
1876 ]),
1877 max_bridges: 2,
1878 bridge_start_timeout_ms: 100,
1879 ..JmsPoolConfig::default()
1880 })
1881 .unwrap(),
1882 );
1883
1884 let (state_tx, state_rx) = watch::channel(BridgeState::Starting);
1885 for name in &["b1", "b2"] {
1886 let slot = Arc::new(BridgeSlot {
1887 name: name.to_string(),
1888 broker_url: format!("tcp://{name}:61616"),
1889 broker_type: BrokerType::ActiveMq,
1890 credentials: None,
1891 state_rx: state_rx.clone(),
1892 state_tx: state_tx.clone(),
1893 process: Arc::new(tokio::sync::Mutex::new(None)),
1894 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1895 });
1896 pool.slots.insert(name.to_string(), slot);
1897 }
1898
1899 assert_eq!(pool.slots.len(), 2);
1900
1901 let guard = pool.bridge_create_lock.lock().await;
1902 let total_count = pool.slots.len();
1903 assert!(total_count >= pool.max_bridges);
1904 let result = if total_count >= pool.max_bridges {
1905 Err(CamelError::Config(format!(
1906 "JMS bridge limit reached: {total_count} bridge(s) >= max_bridges ({})",
1907 pool.max_bridges
1908 )))
1909 } else {
1910 Ok(())
1911 };
1912 drop(guard);
1913
1914 assert!(result.is_err(), "3rd broker should be rejected");
1915 let err_msg = result.unwrap_err().to_string();
1916 assert!(
1917 err_msg.contains("max_bridges"),
1918 "error must mention max_bridges, got: {err_msg}"
1919 );
1920 }
1921
1922 #[tokio::test]
1925 async fn transport_error_refreshes_channel_but_does_not_resend() {
1926 use tokio::sync::watch;
1927 use tonic::transport::Endpoint as TonicEndpoint;
1928 use tower::Service;
1929
1930 let dead_channel = TonicEndpoint::from_static("http://127.0.0.1:1").connect_lazy();
1931
1932 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1933 channel: dead_channel.clone(),
1934 });
1935
1936 let pool = Arc::new(
1937 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1938 "tcp://localhost:61616",
1939 BrokerType::ActiveMq,
1940 ))
1941 .unwrap(),
1942 );
1943
1944 let slot = Arc::new(BridgeSlot {
1945 name: "default".to_string(),
1946 broker_url: "tcp://localhost:61616".to_string(),
1947 broker_type: BrokerType::ActiveMq,
1948 credentials: None,
1949 state_rx: state_rx.clone(),
1950 state_tx: state_tx.clone(),
1951 process: Arc::new(tokio::sync::Mutex::new(None)),
1952 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1953 });
1954 pool.slots.insert("default".to_string(), Arc::clone(&slot));
1955
1956 let endpoint_config =
1957 crate::config::JmsEndpointConfig::from_uri("jms:queue:test-no-resend").unwrap();
1958
1959 let mut producer = LazyJmsProducer {
1960 pool: Arc::clone(&pool),
1961 broker_name: "default".to_string(),
1962 endpoint_config,
1963 resolved_broker_type: BrokerType::ActiveMq,
1964 runtime: test_rt(),
1965 };
1966
1967 let mut exchange = Exchange::default();
1968 exchange.input.body = camel_component_api::Body::Text("hello".to_string());
1969
1970 let result = producer.call(exchange).await;
1971 assert!(result.is_err(), "expected send to fail");
1972
1973 let state_after = state_rx.borrow().clone();
1974 assert!(
1975 matches!(state_after, BridgeState::Restarting { .. }),
1976 "slot must enter Restarting; got: {:?}",
1977 state_after
1978 );
1979
1980 let err_msg = result.unwrap_err().to_string();
1981 assert!(
1982 err_msg.contains(BRIDGE_TRANSPORT_ERROR_PREFIX),
1983 "error must be original transport error, got: {}",
1984 err_msg
1985 );
1986 }
1987
1988 #[tokio::test]
1991 async fn test_jms_bridge_pool_drop_cleans_up_slots() {
1992 use std::sync::atomic::{AtomicBool, Ordering};
1993 use tokio::sync::watch;
1994
1995 let pool = Arc::new(
1996 JmsBridgePool::from_config(JmsPoolConfig {
1997 brokers: HashMap::from([(
1998 "default".to_string(),
1999 BrokerConfig {
2000 broker_url: "tcp://localhost:61616".to_string(),
2001 broker_type: BrokerType::ActiveMq,
2002 username: None,
2003 password: None,
2004 },
2005 )]),
2006 health_check_interval_ms: 50,
2007 ..JmsPoolConfig::default()
2008 })
2009 .unwrap(),
2010 );
2011
2012 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
2014 channel: tonic::transport::Endpoint::from_static("http://127.0.0.1:1").connect_lazy(),
2015 });
2016
2017 let monitor_exited = Arc::new(AtomicBool::new(false));
2018 let exited = monitor_exited.clone();
2019 let rx = state_rx.clone();
2020 let handle = tokio::spawn(async move {
2021 loop {
2022 if matches!(*rx.borrow(), BridgeState::Stopped) {
2023 break;
2024 }
2025 tokio::time::sleep(Duration::from_millis(10)).await;
2026 }
2027 exited.store(true, Ordering::SeqCst);
2028 });
2029
2030 let monitor_handle_ref: Arc<tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>> =
2031 Arc::new(tokio::sync::Mutex::new(None));
2032 *monitor_handle_ref.lock().await = Some(handle);
2033
2034 let slot = Arc::new(BridgeSlot {
2035 name: "default".to_string(),
2036 broker_url: "tcp://localhost:61616".to_string(),
2037 broker_type: BrokerType::ActiveMq,
2038 credentials: None,
2039 state_rx,
2040 state_tx,
2041 process: Arc::new(tokio::sync::Mutex::new(None)),
2042 health_monitor_handle: monitor_handle_ref,
2043 });
2044 pool.slots.insert("default".to_string(), slot);
2045
2046 drop(pool);
2048
2049 tokio::time::sleep(Duration::from_millis(200)).await;
2051
2052 assert!(
2053 monitor_exited.load(Ordering::SeqCst),
2054 "health monitor should have exited after pool drop (Stopped signal sent by Drop)"
2055 );
2056 }
2057}