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 create_endpoint(
720 &self,
721 uri: &str,
722 ctx: &dyn camel_component_api::ComponentContext,
723 ) -> Result<Box<dyn Endpoint>, CamelError> {
724 let endpoint_config = JmsEndpointConfig::from_uri(uri)?;
725 let broker_name = self
726 .pool
727 .resolve_broker_name(endpoint_config.broker_name.as_deref())?;
728 let resolved_broker_type = self.pool.resolve_broker_type(&self.scheme, &broker_name);
729
730 let health_check = JmsHealthCheck::new(Arc::clone(&self.pool), broker_name.clone());
731 ctx.register_current_route_health_check(Arc::new(health_check));
732
733 Ok(Box::new(JmsEndpoint {
734 pool: Arc::clone(&self.pool),
735 uri: uri.to_string(),
736 broker_name,
737 resolved_broker_type,
738 endpoint_config,
739 }))
740 }
741}
742
743struct JmsEndpoint {
746 pool: Arc<JmsBridgePool>,
747 uri: String,
748 broker_name: String,
749 resolved_broker_type: BrokerType,
750 endpoint_config: JmsEndpointConfig,
751}
752
753impl Endpoint for JmsEndpoint {
754 fn uri(&self) -> &str {
755 &self.uri
756 }
757
758 fn create_producer(
759 &self,
760 rt: Arc<dyn camel_component_api::RuntimeObservability>,
761 _ctx: &ProducerContext,
762 ) -> Result<BoxProcessor, CamelError> {
763 Ok(BoxProcessor::new(LazyJmsProducer {
764 pool: Arc::clone(&self.pool),
765 broker_name: self.broker_name.clone(),
766 endpoint_config: self.endpoint_config.clone(),
767 resolved_broker_type: self.resolved_broker_type.clone(),
768 runtime: rt,
769 }))
770 }
771
772 fn create_consumer(
773 &self,
774 rt: Arc<dyn camel_component_api::RuntimeObservability>,
775 ) -> Result<Box<dyn Consumer>, CamelError> {
776 Ok(Box::new(JmsConsumer::new(
777 Arc::clone(&self.pool),
778 self.broker_name.clone(),
779 self.endpoint_config.clone(),
780 self.pool.reconnect.clone(),
781 rt,
782 )))
783 }
784}
785
786#[derive(Clone)]
787struct LazyJmsProducer {
788 pool: Arc<JmsBridgePool>,
789 broker_name: String,
790 endpoint_config: JmsEndpointConfig,
791 #[allow(dead_code)]
792 resolved_broker_type: BrokerType,
793 #[allow(dead_code)]
796 runtime: Arc<dyn camel_component_api::RuntimeObservability>,
797}
798
799impl Service<Exchange> for LazyJmsProducer {
800 type Response = Exchange;
801 type Error = CamelError;
802 type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
803
804 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
805 if let Some(slot) = self.pool.slots.get(&self.broker_name) {
808 match &*slot.state_rx.borrow() {
809 BridgeState::Ready { .. } => return Poll::Ready(Ok(())),
810 BridgeState::Starting | BridgeState::Restarting { .. } => {
811 let waker = cx.waker().clone();
817 let mut rx = slot.state_rx.clone();
818 if let Ok(handle) = tokio::runtime::Handle::try_current() {
819 handle.spawn(async move {
820 let _ = rx.changed().await;
821 waker.wake();
822 });
823 } else {
824 waker.wake_by_ref();
825 }
826 return Poll::Pending;
827 }
828 BridgeState::Degraded(reason) => {
829 return Poll::Ready(Err(CamelError::ProcessorError(format!(
830 "JMS broker '{}' is degraded: {}",
831 self.broker_name, reason
832 ))));
833 }
834 BridgeState::Stopped => {
835 return Poll::Ready(Err(CamelError::ConsumerStopping));
836 }
837 }
838 }
839 Poll::Ready(Ok(()))
840 }
841
842 fn call(&mut self, exchange: Exchange) -> Self::Future {
843 let pool = Arc::clone(&self.pool);
844 let broker_name = self.broker_name.clone();
845 let endpoint_config = self.endpoint_config.clone();
846
847 Box::pin(async move {
848 let slot = pool.get_or_create_slot(&broker_name).await?;
849 let mut rx = slot.state_rx.clone();
850
851 loop {
852 let state = rx.borrow().clone();
853 match state {
854 BridgeState::Ready { channel } => {
855 let mut producer = JmsProducer::new(channel, endpoint_config.clone());
856 match producer.call(exchange).await {
857 Ok(done) => return Ok(done),
858 Err(first_err) if is_bridge_transport_error(&first_err) => {
859 warn!(
860 broker = %broker_name,
861 error = %first_err,
862 "JMS send transport error; refreshing channel (no automatic resend)"
863 );
864
865 if let Err(refresh_err) =
866 pool.refresh_slot_channel(&broker_name).await
867 {
868 warn!(
869 broker = %broker_name,
870 error = %refresh_err,
871 "JMS channel refresh failed; requesting bridge restart"
872 );
873 pool.restart_slot(&broker_name);
874 }
875
876 return Err(first_err);
881 }
882 Err(other_err) => return Err(other_err),
883 }
884 }
885 BridgeState::Degraded(reason) => {
886 return Err(CamelError::ProcessorError(format!(
887 "JMS broker '{}' is degraded: {}",
888 broker_name, reason
889 )));
890 }
891 BridgeState::Stopped => {
892 return Err(CamelError::ProcessorError(format!(
893 "JMS broker '{}' is stopped",
894 broker_name
895 )));
896 }
897 BridgeState::Starting | BridgeState::Restarting { .. } => {
898 if rx.changed().await.is_err() {
899 return Err(CamelError::ProcessorError(format!(
900 "JMS broker '{}' state channel closed",
901 broker_name
902 )));
903 }
904 }
905 }
906 }
907 })
908 }
909}
910
911fn redact_url(url: &str) -> String {
916 if let Some(pos) = url.find("://") {
918 let scheme = &url[..pos + 3]; let rest = &url[pos + 3..];
920 if let Some(at_pos) = rest.find('@') {
922 return format!("{}***@{}", scheme, &rest[at_pos + 1..]);
923 }
924 }
925 url.to_string()
926}
927
928pub fn is_bridge_transport_error(err: &CamelError) -> bool {
929 match err {
934 CamelError::ProcessorError(msg) => msg.starts_with(BRIDGE_TRANSPORT_ERROR_PREFIX),
935 _ => false,
936 }
937}
938
939#[cfg(test)]
942mod tests {
943 use camel_component_api::test_support::PanicRuntimeObservability;
944 fn test_rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
945 std::sync::Arc::new(PanicRuntimeObservability)
946 }
947 fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
948 std::sync::Arc::new(PanicRuntimeObservability)
949 }
950
951 use super::*;
952 use crate::config::{BrokerConfig, JmsPoolConfig};
953 use std::collections::HashMap;
954
955 #[test]
956 fn from_config_accepts_empty_brokers() {
957 let pool_config = JmsPoolConfig::default();
958 let result = JmsBridgePool::from_config(pool_config);
959 assert!(result.is_ok());
960 }
961
962 #[test]
963 fn resolve_broker_name_with_explicit_name() {
964 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
965 "tcp://localhost:61616",
966 BrokerType::ActiveMq,
967 ))
968 .unwrap();
969 assert_eq!(
970 pool.resolve_broker_name(Some("default")).unwrap(),
971 "default"
972 );
973 }
974
975 #[test]
976 fn resolve_broker_name_default() {
977 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
978 "tcp://localhost:61616",
979 BrokerType::ActiveMq,
980 ))
981 .unwrap();
982 assert_eq!(pool.resolve_broker_name(None).unwrap(), "default");
983 }
984
985 #[test]
986 fn resolve_broker_name_unknown_returns_error() {
987 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
988 "tcp://localhost:61616",
989 BrokerType::ActiveMq,
990 ))
991 .unwrap();
992 let err = pool.resolve_broker_name(Some("unknown")).unwrap_err();
993 assert!(
994 err.to_string().contains("Unknown JMS broker 'unknown'"),
995 "got: {}",
996 err
997 );
998 }
999
1000 #[test]
1001 fn resolve_broker_type_scheme_overrides() {
1002 let pool = JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1003 "tcp://localhost:61616",
1004 BrokerType::Generic,
1005 ))
1006 .unwrap();
1007 assert_eq!(
1008 pool.resolve_broker_type("activemq", "default"),
1009 BrokerType::ActiveMq
1010 );
1011 assert_eq!(
1012 pool.resolve_broker_type("artemis", "default"),
1013 BrokerType::Artemis
1014 );
1015 assert_eq!(
1016 pool.resolve_broker_type("jms", "default"),
1017 BrokerType::Generic
1018 );
1019 }
1020
1021 #[test]
1022 fn resolve_broker_type_activemq_scheme_overrides_artemis_config() {
1023 let pool = JmsBridgePool::from_config(JmsPoolConfig {
1024 brokers: HashMap::from([(
1025 "main".to_string(),
1026 BrokerConfig {
1027 broker_url: "tcp://localhost:61616".to_string(),
1028 broker_type: BrokerType::Artemis,
1029 username: None,
1030 password: None,
1031 },
1032 )]),
1033 ..JmsPoolConfig::default()
1034 })
1035 .unwrap();
1036 assert_eq!(
1037 pool.resolve_broker_type("activemq", "main"),
1038 BrokerType::ActiveMq
1039 );
1040 assert_eq!(pool.resolve_broker_type("jms", "main"), BrokerType::Artemis);
1041 }
1042
1043 #[test]
1044 fn create_endpoint_resolves_broker() {
1045 let pool = Arc::new(
1046 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1047 "tcp://localhost:61616",
1048 BrokerType::ActiveMq,
1049 ))
1050 .unwrap(),
1051 );
1052 let component = JmsComponent::with_scheme("jms", pool);
1053 let endpoint = component.create_endpoint(
1054 "jms:queue:orders",
1055 &camel_component_api::NoOpComponentContext,
1056 );
1057 assert!(endpoint.is_ok(), "got: {:?}", endpoint.err());
1058 }
1059
1060 #[test]
1061 fn create_endpoint_rejects_wrong_scheme() {
1062 let pool = Arc::new(
1063 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1064 "tcp://localhost:61616",
1065 BrokerType::ActiveMq,
1066 ))
1067 .unwrap(),
1068 );
1069 let component = JmsComponent::with_scheme("jms", pool);
1070 let err = component
1071 .create_endpoint("kafka:orders", &camel_component_api::NoOpComponentContext)
1072 .err()
1073 .unwrap();
1074 assert!(
1075 err.to_string()
1076 .contains("expected scheme 'jms', 'activemq', or 'artemis'"),
1077 "got: {}",
1078 err
1079 );
1080 }
1081
1082 #[test]
1083 fn create_endpoint_with_explicit_broker_param() {
1084 let pool = Arc::new(
1085 JmsBridgePool::from_config(JmsPoolConfig {
1086 brokers: HashMap::from([
1087 (
1088 "primary".to_string(),
1089 BrokerConfig {
1090 broker_url: "tcp://primary:61616".to_string(),
1091 broker_type: BrokerType::ActiveMq,
1092 username: None,
1093 password: None,
1094 },
1095 ),
1096 (
1097 "secondary".to_string(),
1098 BrokerConfig {
1099 broker_url: "tcp://secondary:61616".to_string(),
1100 broker_type: BrokerType::Artemis,
1101 username: None,
1102 password: None,
1103 },
1104 ),
1105 ]),
1106 ..JmsPoolConfig::default()
1107 })
1108 .unwrap(),
1109 );
1110 let component = JmsComponent::with_scheme("jms", Arc::clone(&pool));
1111 let endpoint = component.create_endpoint(
1112 "jms:queue:orders?broker=secondary",
1113 &camel_component_api::NoOpComponentContext,
1114 );
1115 assert!(endpoint.is_ok(), "got: {:?}", endpoint.err());
1116 }
1117
1118 #[tokio::test]
1119 #[allow(clippy::await_holding_lock)]
1120 async fn concurrent_get_or_create_slot_no_deadlock() {
1121 use tokio::time::timeout;
1122
1123 let _env_lock_guard = crate::BRIDGE_ENV_LOCK.lock().unwrap();
1125
1126 struct EnvGuard {
1127 key: &'static str,
1128 prev: Option<std::ffi::OsString>,
1129 }
1130 impl Drop for EnvGuard {
1131 fn drop(&mut self) {
1132 if let Some(v) = &self.prev {
1133 unsafe { std::env::set_var(self.key, v) };
1135 } else {
1136 unsafe { std::env::remove_var(self.key) };
1138 }
1139 }
1140 }
1141
1142 let env_key = "CAMEL_JMS_BRIDGE_BINARY_PATH";
1143 let _guard = EnvGuard {
1144 key: env_key,
1145 prev: std::env::var_os(env_key),
1146 };
1147 unsafe { std::env::set_var(env_key, "/bin/false") };
1149
1150 let pool = Arc::new(
1151 JmsBridgePool::from_config(JmsPoolConfig {
1152 brokers: HashMap::from([(
1153 "test".to_string(),
1154 BrokerConfig {
1155 broker_url: "tcp://localhost:61616".to_string(),
1156 broker_type: BrokerType::ActiveMq,
1157 username: None,
1158 password: None,
1159 },
1160 )]),
1161 bridge_start_timeout_ms: 100,
1162 ..JmsPoolConfig::default()
1163 })
1164 .unwrap(),
1165 );
1166
1167 let handles: Vec<_> = (0..5)
1168 .map(|_| {
1169 let pool = Arc::clone(&pool);
1170 tokio::spawn(async move {
1171 let _ = pool.get_or_create_slot("test").await;
1172 })
1173 })
1174 .collect();
1175
1176 let result = timeout(Duration::from_secs(5), async {
1177 for h in handles {
1178 let _ = h.await;
1179 }
1180 })
1181 .await;
1182
1183 assert!(result.is_ok(), "Concurrent get_or_create_slot deadlocked!");
1184 }
1185
1186 #[tokio::test]
1187 #[allow(clippy::await_holding_lock)]
1188 async fn lazy_producer_reports_degraded_when_bridge_start_fails() {
1189 use tower::Service;
1190
1191 let _env_lock_guard = crate::BRIDGE_ENV_LOCK.lock().unwrap();
1193
1194 struct EnvGuard {
1195 key: &'static str,
1196 prev: Option<std::ffi::OsString>,
1197 }
1198 impl Drop for EnvGuard {
1199 fn drop(&mut self) {
1200 if let Some(v) = &self.prev {
1201 unsafe { std::env::set_var(self.key, v) };
1203 } else {
1204 unsafe { std::env::remove_var(self.key) };
1206 }
1207 }
1208 }
1209
1210 let env_key = "CAMEL_JMS_BRIDGE_BINARY_PATH";
1211 let _guard = EnvGuard {
1212 key: env_key,
1213 prev: std::env::var_os(env_key),
1214 };
1215 unsafe { std::env::set_var(env_key, "/bin/false") };
1217
1218 let pool = Arc::new(
1219 JmsBridgePool::from_config(JmsPoolConfig {
1220 brokers: HashMap::from([(
1221 "default".to_string(),
1222 BrokerConfig {
1223 broker_url: "tcp://localhost:61616".to_string(),
1224 broker_type: BrokerType::ActiveMq,
1225 username: None,
1226 password: None,
1227 },
1228 )]),
1229 bridge_start_timeout_ms: 100,
1230 ..JmsPoolConfig::default()
1231 })
1232 .unwrap(),
1233 );
1234
1235 let component = JmsComponent::with_scheme("jms", pool);
1236 let endpoint = component
1237 .create_endpoint(
1238 "jms:queue:orders",
1239 &camel_component_api::NoOpComponentContext,
1240 )
1241 .unwrap();
1242 let mut producer = endpoint
1243 .create_producer(rt(), &camel_component_api::ProducerContext::default())
1244 .unwrap();
1245
1246 let mut exchange = Exchange::default();
1247 exchange.input.body = camel_component_api::Body::Text("hello".to_string());
1248
1249 let err = producer.call(exchange).await.unwrap_err();
1250 assert!(err.to_string().contains("is degraded"), "got: {}", err);
1251 }
1252
1253 #[tokio::test]
1257 async fn lazy_producer_requests_restart_when_refresh_unavailable() {
1258 use tokio::sync::watch;
1259 use tonic::transport::Endpoint as TonicEndpoint;
1260 use tower::Service;
1261
1262 let dead_channel = TonicEndpoint::from_static("http://127.0.0.1:1").connect_lazy();
1265
1266 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1267 channel: dead_channel.clone(),
1268 });
1269
1270 let pool = Arc::new(
1271 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1272 "tcp://localhost:61616",
1273 BrokerType::ActiveMq,
1274 ))
1275 .unwrap(),
1276 );
1277
1278 let slot = Arc::new(BridgeSlot {
1280 name: "default".to_string(),
1281 broker_url: "tcp://localhost:61616".to_string(),
1282 broker_type: BrokerType::ActiveMq,
1283 credentials: None,
1284 state_rx: state_rx.clone(),
1285 state_tx: state_tx.clone(),
1286 process: Arc::new(tokio::sync::Mutex::new(None)),
1287 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1288 });
1289 pool.slots.insert("default".to_string(), Arc::clone(&slot));
1290
1291 let endpoint_config =
1292 crate::config::JmsEndpointConfig::from_uri("jms:queue:test-retry").unwrap();
1293
1294 let mut producer = LazyJmsProducer {
1295 pool: Arc::clone(&pool),
1296 broker_name: "default".to_string(),
1297 endpoint_config,
1298 resolved_broker_type: BrokerType::ActiveMq,
1299 runtime: test_rt(),
1300 };
1301
1302 let mut exchange = Exchange::default();
1303 exchange.input.body = camel_component_api::Body::Text("hello".to_string());
1304
1305 let result = producer.call(exchange).await;
1307 assert!(result.is_err(), "expected send to fail");
1308
1309 let state_after = state_rx.borrow().clone();
1312 assert!(
1313 matches!(state_after, BridgeState::Restarting { .. }),
1314 "slot must enter Restarting when refresh is unavailable; got: {:?}",
1315 state_after
1316 );
1317 }
1318
1319 #[test]
1322 fn transport_error_detects_send_error() {
1323 let err = CamelError::ProcessorError(format!(
1324 "{}send error: connection refused",
1325 BRIDGE_TRANSPORT_ERROR_PREFIX
1326 ));
1327 assert!(
1328 is_bridge_transport_error(&err),
1329 "send error must be classified as transport"
1330 );
1331 }
1332
1333 #[test]
1334 fn transport_error_detects_subscribe_error() {
1335 let err = CamelError::ProcessorError(format!(
1336 "{}subscribe error: stream reset",
1337 BRIDGE_TRANSPORT_ERROR_PREFIX
1338 ));
1339 assert!(
1340 is_bridge_transport_error(&err),
1341 "subscribe error must be classified as transport"
1342 );
1343 }
1344
1345 #[test]
1346 fn transport_error_rejects_business_errors() {
1347 let err = CamelError::ProcessorError("JMS broker 'main' is degraded: timeout".to_string());
1348 assert!(
1349 !is_bridge_transport_error(&err),
1350 "degraded state error must NOT be transport"
1351 );
1352 }
1353
1354 #[test]
1355 fn transport_error_rejects_config_errors() {
1356 let err = CamelError::Config("bridge_start_timeout_ms must be > 0".to_string());
1357 assert!(
1358 !is_bridge_transport_error(&err),
1359 "config error must NOT be transport"
1360 );
1361 }
1362
1363 #[test]
1364 fn transport_error_prefix_is_used_by_producer_and_consumer() {
1365 assert!(
1368 BRIDGE_TRANSPORT_ERROR_PREFIX.starts_with("JMS gRPC "),
1369 "prefix must start with 'JMS gRPC '"
1370 );
1371 }
1372
1373 #[tokio::test]
1376 async fn pool_enforces_max_bridges_limit() {
1377 use tokio::sync::watch;
1378
1379 let pool = Arc::new(
1380 JmsBridgePool::from_config(JmsPoolConfig {
1381 brokers: HashMap::from([
1382 (
1383 "b1".to_string(),
1384 BrokerConfig {
1385 broker_url: "tcp://b1:61616".to_string(),
1386 broker_type: BrokerType::ActiveMq,
1387 username: None,
1388 password: None,
1389 },
1390 ),
1391 (
1392 "b2".to_string(),
1393 BrokerConfig {
1394 broker_url: "tcp://b2:61616".to_string(),
1395 broker_type: BrokerType::ActiveMq,
1396 username: None,
1397 password: None,
1398 },
1399 ),
1400 (
1401 "b3".to_string(),
1402 BrokerConfig {
1403 broker_url: "tcp://b3:61616".to_string(),
1404 broker_type: BrokerType::ActiveMq,
1405 username: None,
1406 password: None,
1407 },
1408 ),
1409 ]),
1410 max_bridges: 2,
1411 ..JmsPoolConfig::default()
1412 })
1413 .unwrap(),
1414 );
1415
1416 for name in &["b1", "b2"] {
1419 let (state_tx, state_rx) = watch::channel(BridgeState::Starting);
1420 let slot = Arc::new(BridgeSlot {
1421 name: name.to_string(),
1422 broker_url: format!("tcp://{name}:61616"),
1423 broker_type: BrokerType::ActiveMq,
1424 credentials: None,
1425 state_rx,
1426 state_tx,
1427 process: Arc::new(tokio::sync::Mutex::new(None)),
1428 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1429 });
1430 pool.slots.insert(name.to_string(), slot);
1431 }
1432
1433 let err = pool.get_or_create_slot("b3").await.unwrap_err();
1435 assert!(
1436 err.to_string().contains("max_bridges"),
1437 "expected max_bridges error, got: {}",
1438 err
1439 );
1440 }
1441
1442 #[tokio::test]
1443 async fn pool_allows_slot_when_below_max_bridges() {
1444 use tokio::sync::watch;
1445
1446 let pool = Arc::new(
1447 JmsBridgePool::from_config(JmsPoolConfig {
1448 brokers: HashMap::from([(
1449 "b1".to_string(),
1450 BrokerConfig {
1451 broker_url: "tcp://b1:61616".to_string(),
1452 broker_type: BrokerType::ActiveMq,
1453 username: None,
1454 password: None,
1455 },
1456 )]),
1457 max_bridges: 2,
1458 bridge_start_timeout_ms: 100,
1459 ..JmsPoolConfig::default()
1460 })
1461 .unwrap(),
1462 );
1463
1464 let (state_tx, state_rx) = watch::channel(BridgeState::Degraded("test".to_string()));
1466 let slot = Arc::new(BridgeSlot {
1467 name: "b1".to_string(),
1468 broker_url: "tcp://b1:61616".to_string(),
1469 broker_type: BrokerType::ActiveMq,
1470 credentials: None,
1471 state_rx,
1472 state_tx,
1473 process: Arc::new(tokio::sync::Mutex::new(None)),
1474 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1475 });
1476 pool.slots.insert("b1".to_string(), slot);
1477
1478 let result = pool.get_or_create_slot("b1").await;
1481 assert!(result.is_ok(), "existing slot must be returned");
1482 }
1483
1484 #[tokio::test]
1487 async fn poll_ready_returns_pending_when_starting() {
1488 use tokio::sync::watch;
1489 use tower::Service;
1490
1491 let pool = Arc::new(
1492 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1493 "tcp://localhost:61616",
1494 BrokerType::ActiveMq,
1495 ))
1496 .unwrap(),
1497 );
1498
1499 let (state_tx, state_rx) = watch::channel(BridgeState::Starting);
1500 let slot = Arc::new(BridgeSlot {
1501 name: "default".to_string(),
1502 broker_url: "tcp://localhost:61616".to_string(),
1503 broker_type: BrokerType::ActiveMq,
1504 credentials: None,
1505 state_rx,
1506 state_tx,
1507 process: Arc::new(tokio::sync::Mutex::new(None)),
1508 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1509 });
1510 pool.slots.insert("default".to_string(), slot);
1511
1512 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1513 let mut producer = LazyJmsProducer {
1514 pool: Arc::clone(&pool),
1515 broker_name: "default".to_string(),
1516 endpoint_config,
1517 resolved_broker_type: BrokerType::ActiveMq,
1518 runtime: test_rt(),
1519 };
1520
1521 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1522 assert!(
1523 matches!(result, Poll::Pending),
1524 "poll_ready must be Pending when Starting; got: {:?}",
1525 result
1526 );
1527 }
1528
1529 #[tokio::test]
1530 async fn poll_ready_returns_error_when_degraded() {
1531 use tokio::sync::watch;
1532 use tower::Service;
1533
1534 let pool = Arc::new(
1535 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1536 "tcp://localhost:61616",
1537 BrokerType::ActiveMq,
1538 ))
1539 .unwrap(),
1540 );
1541
1542 let (state_tx, state_rx) =
1543 watch::channel(BridgeState::Degraded("health check failed".to_string()));
1544 let slot = Arc::new(BridgeSlot {
1545 name: "default".to_string(),
1546 broker_url: "tcp://localhost:61616".to_string(),
1547 broker_type: BrokerType::ActiveMq,
1548 credentials: None,
1549 state_rx,
1550 state_tx,
1551 process: Arc::new(tokio::sync::Mutex::new(None)),
1552 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1553 });
1554 pool.slots.insert("default".to_string(), slot);
1555
1556 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1557 let mut producer = LazyJmsProducer {
1558 pool: Arc::clone(&pool),
1559 broker_name: "default".to_string(),
1560 endpoint_config,
1561 resolved_broker_type: BrokerType::ActiveMq,
1562 runtime: test_rt(),
1563 };
1564
1565 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1566 assert!(
1567 matches!(result, Poll::Ready(Err(_))),
1568 "poll_ready must be Err when Degraded; got: {:?}",
1569 result
1570 );
1571 let err_msg = match result {
1572 Poll::Ready(Err(e)) => e.to_string(),
1573 _ => unreachable!(),
1574 };
1575 assert!(
1576 err_msg.contains("degraded"),
1577 "error must mention degraded: {}",
1578 err_msg
1579 );
1580 }
1581
1582 #[tokio::test]
1583 async fn lazy_producer_poll_ready_returns_consumer_stopping_on_stopped() {
1584 use tokio::sync::watch;
1585 use tower::Service;
1586
1587 let pool = Arc::new(
1588 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1589 "tcp://localhost:61616",
1590 BrokerType::ActiveMq,
1591 ))
1592 .unwrap(),
1593 );
1594
1595 let (state_tx, state_rx) = watch::channel(BridgeState::Stopped);
1596 let slot = Arc::new(BridgeSlot {
1597 name: "default".to_string(),
1598 broker_url: "tcp://localhost:61616".to_string(),
1599 broker_type: BrokerType::ActiveMq,
1600 credentials: None,
1601 state_rx,
1602 state_tx,
1603 process: Arc::new(tokio::sync::Mutex::new(None)),
1604 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1605 });
1606 pool.slots.insert("default".to_string(), slot);
1607
1608 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1609 let mut producer = LazyJmsProducer {
1610 pool: Arc::clone(&pool),
1611 broker_name: "default".to_string(),
1612 endpoint_config,
1613 resolved_broker_type: BrokerType::ActiveMq,
1614 runtime: test_rt(),
1615 };
1616
1617 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1618 assert!(
1619 matches!(result, Poll::Ready(Err(CamelError::ConsumerStopping))),
1620 "poll_ready must return ConsumerStopping when Stopped; got: {:?}",
1621 result
1622 );
1623 }
1624
1625 #[tokio::test]
1626 async fn poll_ready_returns_ready_when_slot_ready() {
1627 use tokio::sync::watch;
1628 use tonic::transport::Endpoint as TonicEndpoint;
1629 use tower::Service;
1630
1631 let pool = Arc::new(
1632 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1633 "tcp://localhost:61616",
1634 BrokerType::ActiveMq,
1635 ))
1636 .unwrap(),
1637 );
1638
1639 let lazy_channel = TonicEndpoint::from_static("http://127.0.0.1:1").connect_lazy();
1640 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1641 channel: lazy_channel,
1642 });
1643 let slot = Arc::new(BridgeSlot {
1644 name: "default".to_string(),
1645 broker_url: "tcp://localhost:61616".to_string(),
1646 broker_type: BrokerType::ActiveMq,
1647 credentials: None,
1648 state_rx,
1649 state_tx,
1650 process: Arc::new(tokio::sync::Mutex::new(None)),
1651 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1652 });
1653 pool.slots.insert("default".to_string(), slot);
1654
1655 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1656 let mut producer = LazyJmsProducer {
1657 pool: Arc::clone(&pool),
1658 broker_name: "default".to_string(),
1659 endpoint_config,
1660 resolved_broker_type: BrokerType::ActiveMq,
1661 runtime: test_rt(),
1662 };
1663
1664 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1665 assert!(
1666 matches!(result, Poll::Ready(Ok(()))),
1667 "poll_ready must be Ready(Ok) when bridge is Ready; got: {:?}",
1668 result
1669 );
1670 }
1671
1672 #[tokio::test]
1673 async fn poll_ready_returns_ready_when_no_slot_exists() {
1674 use tower::Service;
1675
1676 let pool = Arc::new(
1677 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1678 "tcp://localhost:61616",
1679 BrokerType::ActiveMq,
1680 ))
1681 .unwrap(),
1682 );
1683
1684 let endpoint_config = crate::config::JmsEndpointConfig::from_uri("jms:queue:test").unwrap();
1685 let mut producer = LazyJmsProducer {
1686 pool: Arc::clone(&pool),
1687 broker_name: "default".to_string(),
1688 endpoint_config,
1689 resolved_broker_type: BrokerType::ActiveMq,
1690 runtime: test_rt(),
1691 };
1692
1693 let result = producer.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
1695 assert!(
1696 matches!(result, Poll::Ready(Ok(()))),
1697 "poll_ready must be Ready(Ok) when no slot exists; got: {:?}",
1698 result
1699 );
1700 }
1701
1702 #[tokio::test]
1705 async fn pool_shutdown_awaits_health_monitor() {
1706 use tokio::sync::watch;
1707
1708 let pool = Arc::new(
1709 JmsBridgePool::from_config(JmsPoolConfig {
1710 brokers: HashMap::from([(
1711 "default".to_string(),
1712 BrokerConfig {
1713 broker_url: "tcp://localhost:61616".to_string(),
1714 broker_type: BrokerType::ActiveMq,
1715 username: None,
1716 password: None,
1717 },
1718 )]),
1719 health_check_interval_ms: 100,
1720 ..JmsPoolConfig::default()
1721 })
1722 .unwrap(),
1723 );
1724
1725 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1727 channel: tonic::transport::Endpoint::from_static("http://127.0.0.1:1").connect_lazy(),
1728 });
1729 let monitor_handle_ref: Arc<tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>> =
1730 Arc::new(tokio::sync::Mutex::new(None));
1731
1732 let state_rx_clone = state_rx.clone();
1734 let handle = tokio::spawn(async move {
1735 loop {
1736 if matches!(*state_rx_clone.borrow(), BridgeState::Stopped) {
1737 break;
1738 }
1739 tokio::time::sleep(Duration::from_millis(50)).await;
1740 }
1741 });
1742 *monitor_handle_ref.lock().await = Some(handle);
1743
1744 let slot = Arc::new(BridgeSlot {
1745 name: "default".to_string(),
1746 broker_url: "tcp://localhost:61616".to_string(),
1747 broker_type: BrokerType::ActiveMq,
1748 credentials: None,
1749 state_rx,
1750 state_tx,
1751 process: Arc::new(tokio::sync::Mutex::new(None)),
1752 health_monitor_handle: monitor_handle_ref,
1753 });
1754 pool.slots.insert("default".to_string(), slot);
1755
1756 let result = pool.shutdown().await;
1758 let _ = result;
1760 }
1761
1762 #[tokio::test]
1763 async fn health_monitor_handle_stored_after_spawn() {
1764 use tokio::sync::watch;
1765
1766 let pool = Arc::new(
1767 JmsBridgePool::from_config(JmsPoolConfig {
1768 brokers: HashMap::from([(
1769 "default".to_string(),
1770 BrokerConfig {
1771 broker_url: "tcp://localhost:61616".to_string(),
1772 broker_type: BrokerType::ActiveMq,
1773 username: None,
1774 password: None,
1775 },
1776 )]),
1777 health_check_interval_ms: 100,
1778 bridge_start_timeout_ms: 100,
1779 ..JmsPoolConfig::default()
1780 })
1781 .unwrap(),
1782 );
1783
1784 let (state_tx, state_rx) = watch::channel(BridgeState::Stopped);
1785 let slot = Arc::new(BridgeSlot {
1786 name: "default".to_string(),
1787 broker_url: "tcp://localhost:61616".to_string(),
1788 broker_type: BrokerType::ActiveMq,
1789 credentials: None,
1790 state_rx,
1791 state_tx,
1792 process: Arc::new(tokio::sync::Mutex::new(None)),
1793 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1794 });
1795 pool.slots.insert("default".to_string(), Arc::clone(&slot));
1796
1797 pool.spawn_health_monitor(Arc::clone(&slot)).await;
1798
1799 tokio::time::sleep(Duration::from_millis(50)).await;
1800 let guard = slot.health_monitor_handle.lock().await;
1801 assert!(
1802 guard.is_some(),
1803 "health monitor handle must be stored after spawn_health_monitor"
1804 );
1805 }
1806
1807 #[test]
1810 fn redact_url_strips_userinfo_with_password() {
1811 assert_eq!(
1812 redact_url("tcp://admin:s3cret@broker:61616"),
1813 "tcp://***@broker:61616"
1814 );
1815 }
1816
1817 #[test]
1818 fn redact_url_strips_userinfo_without_password() {
1819 assert_eq!(
1820 redact_url("tcp://admin@broker:61616"),
1821 "tcp://***@broker:61616"
1822 );
1823 }
1824
1825 #[test]
1826 fn redact_url_passes_clean_url_unchanged() {
1827 assert_eq!(redact_url("tcp://localhost:61616"), "tcp://localhost:61616");
1828 }
1829
1830 #[test]
1831 fn redact_url_handles_ssl_scheme() {
1832 assert_eq!(
1833 redact_url("ssl://user:pass@secure-broker:61617"),
1834 "ssl://***@secure-broker:61617"
1835 );
1836 }
1837
1838 #[tokio::test]
1841 async fn concurrent_slot_creation_respects_max_bridges() {
1842 let pool = Arc::new(
1843 JmsBridgePool::from_config(JmsPoolConfig {
1844 brokers: HashMap::from([
1845 (
1846 "b1".to_string(),
1847 BrokerConfig {
1848 broker_url: "tcp://b1:61616".to_string(),
1849 broker_type: BrokerType::ActiveMq,
1850 username: None,
1851 password: None,
1852 },
1853 ),
1854 (
1855 "b2".to_string(),
1856 BrokerConfig {
1857 broker_url: "tcp://b2:61616".to_string(),
1858 broker_type: BrokerType::ActiveMq,
1859 username: None,
1860 password: None,
1861 },
1862 ),
1863 (
1864 "b3".to_string(),
1865 BrokerConfig {
1866 broker_url: "tcp://b3:61616".to_string(),
1867 broker_type: BrokerType::ActiveMq,
1868 username: None,
1869 password: None,
1870 },
1871 ),
1872 ]),
1873 max_bridges: 2,
1874 bridge_start_timeout_ms: 100,
1875 ..JmsPoolConfig::default()
1876 })
1877 .unwrap(),
1878 );
1879
1880 let (state_tx, state_rx) = watch::channel(BridgeState::Starting);
1881 for name in &["b1", "b2"] {
1882 let slot = Arc::new(BridgeSlot {
1883 name: name.to_string(),
1884 broker_url: format!("tcp://{name}:61616"),
1885 broker_type: BrokerType::ActiveMq,
1886 credentials: None,
1887 state_rx: state_rx.clone(),
1888 state_tx: state_tx.clone(),
1889 process: Arc::new(tokio::sync::Mutex::new(None)),
1890 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1891 });
1892 pool.slots.insert(name.to_string(), slot);
1893 }
1894
1895 assert_eq!(pool.slots.len(), 2);
1896
1897 let guard = pool.bridge_create_lock.lock().await;
1898 let total_count = pool.slots.len();
1899 assert!(total_count >= pool.max_bridges);
1900 let result = if total_count >= pool.max_bridges {
1901 Err(CamelError::Config(format!(
1902 "JMS bridge limit reached: {total_count} bridge(s) >= max_bridges ({})",
1903 pool.max_bridges
1904 )))
1905 } else {
1906 Ok(())
1907 };
1908 drop(guard);
1909
1910 assert!(result.is_err(), "3rd broker should be rejected");
1911 let err_msg = result.unwrap_err().to_string();
1912 assert!(
1913 err_msg.contains("max_bridges"),
1914 "error must mention max_bridges, got: {err_msg}"
1915 );
1916 }
1917
1918 #[tokio::test]
1921 async fn transport_error_refreshes_channel_but_does_not_resend() {
1922 use tokio::sync::watch;
1923 use tonic::transport::Endpoint as TonicEndpoint;
1924 use tower::Service;
1925
1926 let dead_channel = TonicEndpoint::from_static("http://127.0.0.1:1").connect_lazy();
1927
1928 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
1929 channel: dead_channel.clone(),
1930 });
1931
1932 let pool = Arc::new(
1933 JmsBridgePool::from_config(JmsPoolConfig::single_broker(
1934 "tcp://localhost:61616",
1935 BrokerType::ActiveMq,
1936 ))
1937 .unwrap(),
1938 );
1939
1940 let slot = Arc::new(BridgeSlot {
1941 name: "default".to_string(),
1942 broker_url: "tcp://localhost:61616".to_string(),
1943 broker_type: BrokerType::ActiveMq,
1944 credentials: None,
1945 state_rx: state_rx.clone(),
1946 state_tx: state_tx.clone(),
1947 process: Arc::new(tokio::sync::Mutex::new(None)),
1948 health_monitor_handle: Arc::new(tokio::sync::Mutex::new(None)),
1949 });
1950 pool.slots.insert("default".to_string(), Arc::clone(&slot));
1951
1952 let endpoint_config =
1953 crate::config::JmsEndpointConfig::from_uri("jms:queue:test-no-resend").unwrap();
1954
1955 let mut producer = LazyJmsProducer {
1956 pool: Arc::clone(&pool),
1957 broker_name: "default".to_string(),
1958 endpoint_config,
1959 resolved_broker_type: BrokerType::ActiveMq,
1960 runtime: test_rt(),
1961 };
1962
1963 let mut exchange = Exchange::default();
1964 exchange.input.body = camel_component_api::Body::Text("hello".to_string());
1965
1966 let result = producer.call(exchange).await;
1967 assert!(result.is_err(), "expected send to fail");
1968
1969 let state_after = state_rx.borrow().clone();
1970 assert!(
1971 matches!(state_after, BridgeState::Restarting { .. }),
1972 "slot must enter Restarting; got: {:?}",
1973 state_after
1974 );
1975
1976 let err_msg = result.unwrap_err().to_string();
1977 assert!(
1978 err_msg.contains(BRIDGE_TRANSPORT_ERROR_PREFIX),
1979 "error must be original transport error, got: {}",
1980 err_msg
1981 );
1982 }
1983
1984 #[tokio::test]
1987 async fn test_jms_bridge_pool_drop_cleans_up_slots() {
1988 use std::sync::atomic::{AtomicBool, Ordering};
1989 use tokio::sync::watch;
1990
1991 let pool = Arc::new(
1992 JmsBridgePool::from_config(JmsPoolConfig {
1993 brokers: HashMap::from([(
1994 "default".to_string(),
1995 BrokerConfig {
1996 broker_url: "tcp://localhost:61616".to_string(),
1997 broker_type: BrokerType::ActiveMq,
1998 username: None,
1999 password: None,
2000 },
2001 )]),
2002 health_check_interval_ms: 50,
2003 ..JmsPoolConfig::default()
2004 })
2005 .unwrap(),
2006 );
2007
2008 let (state_tx, state_rx) = watch::channel(BridgeState::Ready {
2010 channel: tonic::transport::Endpoint::from_static("http://127.0.0.1:1").connect_lazy(),
2011 });
2012
2013 let monitor_exited = Arc::new(AtomicBool::new(false));
2014 let exited = monitor_exited.clone();
2015 let rx = state_rx.clone();
2016 let handle = tokio::spawn(async move {
2017 loop {
2018 if matches!(*rx.borrow(), BridgeState::Stopped) {
2019 break;
2020 }
2021 tokio::time::sleep(Duration::from_millis(10)).await;
2022 }
2023 exited.store(true, Ordering::SeqCst);
2024 });
2025
2026 let monitor_handle_ref: Arc<tokio::sync::Mutex<Option<tokio::task::JoinHandle<()>>>> =
2027 Arc::new(tokio::sync::Mutex::new(None));
2028 *monitor_handle_ref.lock().await = Some(handle);
2029
2030 let slot = Arc::new(BridgeSlot {
2031 name: "default".to_string(),
2032 broker_url: "tcp://localhost:61616".to_string(),
2033 broker_type: BrokerType::ActiveMq,
2034 credentials: None,
2035 state_rx,
2036 state_tx,
2037 process: Arc::new(tokio::sync::Mutex::new(None)),
2038 health_monitor_handle: monitor_handle_ref,
2039 });
2040 pool.slots.insert("default".to_string(), slot);
2041
2042 drop(pool);
2044
2045 tokio::time::sleep(Duration::from_millis(200)).await;
2047
2048 assert!(
2049 monitor_exited.load(Ordering::SeqCst),
2050 "health monitor should have exited after pool drop (Stopped signal sent by Drop)"
2051 );
2052 }
2053}