Skip to main content

camel_component_jms/
component.rs

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
31// ── Transport error classification ───────────────────────────────────────────
32
33/// Shared constant prefix for all bridge transport errors.
34///
35/// Both error-producing sites (producer.rs, consumer.rs) and the detection
36/// helper (`is_bridge_transport_error`) reference this constant so that a
37/// format-string drift cannot silently break retry logic.
38pub const BRIDGE_TRANSPORT_ERROR_PREFIX: &str = "JMS gRPC ";
39const MAX_RESTART_ATTEMPTS: u32 = 10;
40
41// ── BridgeState ──────────────────────────────────────────────────────────────
42
43#[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
52// ── BridgeSlot ───────────────────────────────────────────────────────────────
53
54pub 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    /// BridgeProcess::stop(mut self) takes ownership — Mutex<Option<>> is required.
62    pub process: Arc<tokio::sync::Mutex<Option<BridgeProcess>>>,
63    /// JoinHandle of the health monitor task for this slot.
64    /// Stored so that shutdown can await the monitor and observe panics.
65    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
78// ── JmsBridgePool ────────────────────────────────────────────────────────────
79
80pub 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    /// Maximum number of concurrently active (Starting/Ready) bridges.
89    pub(crate) max_bridges: usize,
90    /// Serializes bridge admission (check + insert) to prevent race on max_bridges.
91    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        // Backward compat: broker_reconnect_interval_ms overrides
99        // reconnect.initial_delay when explicitly set (non-default).
100        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    /// Resolve broker name from the URI `broker=` param.
122    ///
123    /// - If `Some(name)` → validate it exists in config and return it.
124    /// - If `None` and exactly one broker is configured → use it implicitly.
125    /// - If `None` and multiple brokers are configured → error asking for `?broker=`.
126    /// - If `None` and no brokers are configured → error asking to declare brokers.
127    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()), // allow-unwrap
143                _ => 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    /// Resolve broker type: activemq/artemis schemes hard-override config type; jms uses config.
152    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    /// Get or create a BridgeSlot for the given broker name.
183    /// If the slot doesn't exist, starts the bridge process and spawns the health monitor.
184    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        // Serialize admission: check + insert in single critical section to prevent
193        // concurrent creators from both seeing count below limit and exceeding max_bridges.
194        let _guard = self.bridge_create_lock.lock().await;
195
196        // Re-check after acquiring lock (another caller may have inserted while we waited).
197        if let Some(slot) = self.slots.get(broker_name) {
198            return Ok(Arc::clone(&*slot));
199        }
200
201        // Enforce max_bridges: count ALL slots in the map under the admission lock.
202        // We count total slots (not just Starting/Ready) because any inserted slot
203        // represents an allocated bridge — even Degraded/Restarting slots hold resources
204        // and the bridge process may still be running. Counting only active states would
205        // allow a race: slot A transitions Starting→Degraded between two callers' checks,
206        // letting both pass the limit.
207        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        // Clone all required broker data before touching DashMap::entry().
220        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    /// Signal a slot to restart (called by producers on transport errors).
279    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    /// Recreate the tonic channel for an existing running bridge process.
289    ///
290    /// Useful when a channel becomes stale after transport-level failures while
291    /// the underlying bridge process is still alive.
292    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        // Lock is held during connect() (TLS handshake + retries) — widened from
302        // the old grpc_port() read. Bounded: localhost TLS (~10ms typical), per-slot
303        // lock (only blocks same-broker operations). Acceptable trade-off vs exposing
304        // pub(crate) TLS material to component crates.
305        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    /// Shutdown all slots: stop all bridge processes and await health monitors.
330    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                // Signal the health monitor to stop.
338                let _ = slot.state_tx.send(BridgeState::Stopped);
339
340                // Stop the bridge process.
341                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                // Await the health monitor task with timeout; abort if it doesn't stop.
352                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                            // log-policy: system-broken
465                            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                                // Guard: don't resurrect a stopped slot (shutdown may have run
498                                // while this async bridge start was in-flight).
499                                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                                // Guard: don't schedule retries after shutdown.
512                                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        // Store the handle so shutdown can await the monitor.
534        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
599// ── Drop impl: cleanup on pool drop without explicit shutdown ─────────────────
600
601impl Drop for JmsBridgePool {
602    fn drop(&mut self) {
603        self.shutting_down.store(true, Ordering::SeqCst);
604
605        // Clone bridge slots before we lose access to self.slots.
606        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                // Spawn the per-bridge cleanup (same sequence as shutdown()).
619                // Dropping the JoinHandle detaches the task — it runs to completion.
620                drop(handle.spawn(async move {
621                    for (_name, slot) in slots {
622                        // Signal the health monitor to stop.
623                        let _ = slot.state_tx.send(BridgeState::Stopped);
624
625                        // Stop the bridge process.
626                        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                        // Await the health monitor task with timeout; abort if needed.
635                        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// ── JmsComponent ─────────────────────────────────────────────────────────────
662
663#[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    /// Test helper: send a message directly without going through a route.
682    #[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
747// ── JmsEndpoint ──────────────────────────────────────────────────────────────
748
749struct 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    /// Phase B will use this for `rt.metrics().increment_errors(...)` and
798    /// `rt.health().force_unhealthy_for_route(...)` calls per ADR-0012.
799    #[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        // Check existing slot state if one is already present.
810        // If no slot exists yet, return Ready — call() will handle async bridge start.
811        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                    // Register a waker so the executor is notified when bridge state
816                    // changes. Without this, Poll::Pending would stall callers that use
817                    // strict Tower semantics (poll_ready loop before call()).
818                    // Guard with try_current: fall back to wake_by_ref when no Tokio
819                    // runtime is active (e.g. unit tests).
820                    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                                // Do NOT automatically resend — the first send may have reached
881                                // the broker even though the ack failed. Resending non-idempotent
882                                // writes causes duplicates. Return the original error so the caller
883                                // can decide whether to retry.
884                                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
915// ── Helpers ──────────────────────────────────────────────────────────────────
916
917/// Redact userinfo (username:password@) from a broker URL for safe logging.
918/// Handles URLs like `tcp://user:pass@host:61616` → `tcp://***@host:61616`.
919fn redact_url(url: &str) -> String {
920    // Find the scheme separator (://)
921    if let Some(pos) = url.find("://") {
922        let scheme = &url[..pos + 3]; // includes "://"
923        let rest = &url[pos + 3..];
924        // Find @ in the remainder — everything before @ is userinfo
925        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    // Typed variant matching: only ProcessorError messages that start with
934    // the well-known transport prefix are classified as transport errors.
935    // This rejects Config errors, business errors, and other CamelError variants
936    // without relying on the Display wrapper formatting.
937    match err {
938        CamelError::ProcessorError(msg) => msg.starts_with(BRIDGE_TRANSPORT_ERROR_PREFIX),
939        _ => false,
940    }
941}
942
943// ── Unit tests ───────────────────────────────────────────────────────────────
944
945#[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        // Serialize against other tests mutating CAMEL_JMS_BRIDGE_BINARY_PATH (rc-alwn).
1128        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                    // SAFETY: restoring process env in test scope.
1138                    unsafe { std::env::set_var(self.key, v) };
1139                } else {
1140                    // SAFETY: restoring process env in test scope.
1141                    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        // SAFETY: test-scoped env mutation.
1152        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        // Serialize against other tests mutating CAMEL_JMS_BRIDGE_BINARY_PATH (rc-alwn).
1196        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                    // SAFETY: restoring process env in test scope.
1206                    unsafe { std::env::set_var(self.key, v) };
1207                } else {
1208                    // SAFETY: restoring process env in test scope.
1209                    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        // SAFETY: test-scoped env mutation.
1220        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    /// A send transport error should trigger a channel refresh attempt first.
1258    /// If refresh cannot be performed (e.g. no running bridge process metadata),
1259    /// the producer requests a bridge restart as fallback.
1260    #[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        // Build a lazy channel to a port where nothing is listening.
1267        // connect_lazy() succeeds immediately; the error manifests on the actual RPC call.
1268        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        // Manually insert a slot with the dead-channel in Ready state.
1283        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        // The send will fail because the channel points to a dead port.
1310        let result = producer.call(exchange).await;
1311        assert!(result.is_err(), "expected send to fail");
1312
1313        // Refresh cannot run in this setup (slot has no BridgeProcess), so the
1314        // fallback path requests a restart.
1315        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    // ── JMS-007: Transport error classification ──────────────────────────────
1324
1325    #[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        // Verify the constant prefix matches what producer.rs and consumer.rs emit.
1370        // If this test fails, the constant has drifted from the error format strings.
1371        assert!(
1372            BRIDGE_TRANSPORT_ERROR_PREFIX.starts_with("JMS gRPC "),
1373            "prefix must start with 'JMS gRPC '"
1374        );
1375    }
1376
1377    // ── JMS-006: max_bridges enforcement ─────────────────────────────────────
1378
1379    #[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        // Manually insert two slots to simulate existing bridges.
1421        // max_bridges counts ALL slots in the map (not just active states).
1422        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        // Attempting to create a third slot should fail.
1438        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        // Insert one slot in Degraded state (not counted as active).
1469        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        // b1 is Degraded (not active), so creating b1's slot returns existing.
1483        // The max_bridges check only applies to new slots.
1484        let result = pool.get_or_create_slot("b1").await;
1485        assert!(result.is_ok(), "existing slot must be returned");
1486    }
1487
1488    // ── JMS-003: poll_ready reflects bridge state ────────────────────────────
1489
1490    #[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        // No slot exists yet — poll_ready should return Ready so call() can start the bridge.
1698        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    // ── JMS-001: Health monitor lifecycle ────────────────────────────────────
1707
1708    #[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        // Create a slot manually with a spawned health monitor task.
1730        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        // Spawn a simple monitor that exits on Stopped.
1737        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        // Shutdown should complete without hanging — the monitor exits on Stopped.
1761        let result = pool.shutdown().await;
1762        // May report errors from bridge process (none in this test), but must not hang.
1763        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    // ── JMS-008: URL redaction for safe logging ──────────────────────────────
1812
1813    #[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    // ── JMS-009: max_bridges race condition under concurrency ────────────────
1843
1844    #[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    // ── JMS-010: Transport error does NOT auto-resend ────────────────────────
1923
1924    #[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    // ── D-L12: Drop impl cleans up slots without explicit shutdown ───────────
1989
1990    #[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        // Create a slot with a health monitor that exits on BridgeState::Stopped.
2013        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 the pool WITHOUT calling shutdown() — the Drop impl must fire.
2047        drop(pool);
2048
2049        // Give the spawned cleanup task time to send Stopped and await the monitor.
2050        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}