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 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
743// ── JmsEndpoint ──────────────────────────────────────────────────────────────
744
745struct 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    /// Phase B will use this for `rt.metrics().increment_errors(...)` and
794    /// `rt.health().force_unhealthy_for_route(...)` calls per ADR-0012.
795    #[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        // Check existing slot state if one is already present.
806        // If no slot exists yet, return Ready — call() will handle async bridge start.
807        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                    // Register a waker so the executor is notified when bridge state
812                    // changes. Without this, Poll::Pending would stall callers that use
813                    // strict Tower semantics (poll_ready loop before call()).
814                    // Guard with try_current: fall back to wake_by_ref when no Tokio
815                    // runtime is active (e.g. unit tests).
816                    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                                // Do NOT automatically resend — the first send may have reached
877                                // the broker even though the ack failed. Resending non-idempotent
878                                // writes causes duplicates. Return the original error so the caller
879                                // can decide whether to retry.
880                                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
911// ── Helpers ──────────────────────────────────────────────────────────────────
912
913/// Redact userinfo (username:password@) from a broker URL for safe logging.
914/// Handles URLs like `tcp://user:pass@host:61616` → `tcp://***@host:61616`.
915fn redact_url(url: &str) -> String {
916    // Find the scheme separator (://)
917    if let Some(pos) = url.find("://") {
918        let scheme = &url[..pos + 3]; // includes "://"
919        let rest = &url[pos + 3..];
920        // Find @ in the remainder — everything before @ is userinfo
921        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    // Typed variant matching: only ProcessorError messages that start with
930    // the well-known transport prefix are classified as transport errors.
931    // This rejects Config errors, business errors, and other CamelError variants
932    // without relying on the Display wrapper formatting.
933    match err {
934        CamelError::ProcessorError(msg) => msg.starts_with(BRIDGE_TRANSPORT_ERROR_PREFIX),
935        _ => false,
936    }
937}
938
939// ── Unit tests ───────────────────────────────────────────────────────────────
940
941#[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        // Serialize against other tests mutating CAMEL_JMS_BRIDGE_BINARY_PATH (rc-alwn).
1124        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                    // SAFETY: restoring process env in test scope.
1134                    unsafe { std::env::set_var(self.key, v) };
1135                } else {
1136                    // SAFETY: restoring process env in test scope.
1137                    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        // SAFETY: test-scoped env mutation.
1148        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        // Serialize against other tests mutating CAMEL_JMS_BRIDGE_BINARY_PATH (rc-alwn).
1192        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                    // SAFETY: restoring process env in test scope.
1202                    unsafe { std::env::set_var(self.key, v) };
1203                } else {
1204                    // SAFETY: restoring process env in test scope.
1205                    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        // SAFETY: test-scoped env mutation.
1216        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    /// A send transport error should trigger a channel refresh attempt first.
1254    /// If refresh cannot be performed (e.g. no running bridge process metadata),
1255    /// the producer requests a bridge restart as fallback.
1256    #[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        // Build a lazy channel to a port where nothing is listening.
1263        // connect_lazy() succeeds immediately; the error manifests on the actual RPC call.
1264        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        // Manually insert a slot with the dead-channel in Ready state.
1279        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        // The send will fail because the channel points to a dead port.
1306        let result = producer.call(exchange).await;
1307        assert!(result.is_err(), "expected send to fail");
1308
1309        // Refresh cannot run in this setup (slot has no BridgeProcess), so the
1310        // fallback path requests a restart.
1311        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    // ── JMS-007: Transport error classification ──────────────────────────────
1320
1321    #[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        // Verify the constant prefix matches what producer.rs and consumer.rs emit.
1366        // If this test fails, the constant has drifted from the error format strings.
1367        assert!(
1368            BRIDGE_TRANSPORT_ERROR_PREFIX.starts_with("JMS gRPC "),
1369            "prefix must start with 'JMS gRPC '"
1370        );
1371    }
1372
1373    // ── JMS-006: max_bridges enforcement ─────────────────────────────────────
1374
1375    #[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        // Manually insert two slots to simulate existing bridges.
1417        // max_bridges counts ALL slots in the map (not just active states).
1418        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        // Attempting to create a third slot should fail.
1434        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        // Insert one slot in Degraded state (not counted as active).
1465        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        // b1 is Degraded (not active), so creating b1's slot returns existing.
1479        // The max_bridges check only applies to new slots.
1480        let result = pool.get_or_create_slot("b1").await;
1481        assert!(result.is_ok(), "existing slot must be returned");
1482    }
1483
1484    // ── JMS-003: poll_ready reflects bridge state ────────────────────────────
1485
1486    #[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        // No slot exists yet — poll_ready should return Ready so call() can start the bridge.
1694        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    // ── JMS-001: Health monitor lifecycle ────────────────────────────────────
1703
1704    #[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        // Create a slot manually with a spawned health monitor task.
1726        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        // Spawn a simple monitor that exits on Stopped.
1733        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        // Shutdown should complete without hanging — the monitor exits on Stopped.
1757        let result = pool.shutdown().await;
1758        // May report errors from bridge process (none in this test), but must not hang.
1759        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    // ── JMS-008: URL redaction for safe logging ──────────────────────────────
1808
1809    #[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    // ── JMS-009: max_bridges race condition under concurrency ────────────────
1839
1840    #[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    // ── JMS-010: Transport error does NOT auto-resend ────────────────────────
1919
1920    #[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    // ── D-L12: Drop impl cleans up slots without explicit shutdown ───────────
1985
1986    #[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        // Create a slot with a health monitor that exits on BridgeState::Stopped.
2009        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 the pool WITHOUT calling shutdown() — the Drop impl must fire.
2043        drop(pool);
2044
2045        // Give the spawned cleanup task time to send Stopped and await the monitor.
2046        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}