Skip to main content

camel_component_container/
lib.rs

1//! Camel Container Component
2//!
3//! This component provides integration with Docker containers, allowing Camel routes
4//! to manage container lifecycle (create, start, stop, remove) and consume container events.
5
6pub mod bundle;
7pub mod health;
8
9pub use bundle::ContainerBundle;
10pub use health::ContainerHealthCheck;
11
12use std::collections::{HashMap, HashSet};
13use std::future::Future;
14use std::pin::Pin;
15use std::sync::{Arc, Mutex};
16use std::task::{Context, Poll};
17
18use async_trait::async_trait;
19use bollard::Docker;
20use bollard::models::{
21    ContainerCreateBody, NetworkConnectRequest, NetworkCreateRequest, NetworkDisconnectRequest,
22};
23use bollard::query_parameters::{
24    CreateContainerOptions, CreateImageOptions, EventsOptions, ListContainersOptions,
25    ListImagesOptions, ListNetworksOptions, LogsOptions, RemoveContainerOptions,
26    StartContainerOptions,
27};
28use bollard::service::{HostConfig, PortBinding};
29use camel_component_api::NetworkRetryPolicy;
30use camel_component_api::parse_uri;
31use camel_component_api::retry_async_cancelable;
32use camel_component_api::{Body, BoxProcessor, CamelError, Exchange, Message};
33use camel_component_api::{
34    Component, Consumer, ConsumerContext, Endpoint, ProducerContext, RuntimeObservability,
35};
36use tower::Service;
37
38/// Global tracker for containers created by this component.
39/// Used for cleanup on shutdown (especially important for hot-reload scenarios).
40static CONTAINER_TRACKER: once_cell::sync::Lazy<Arc<Mutex<HashSet<String>>>> =
41    once_cell::sync::Lazy::new(|| Arc::new(Mutex::new(HashSet::new())));
42
43/// Registers a container ID for tracking (will be cleaned up on shutdown).
44fn track_container(id: String) {
45    if let Ok(mut tracker) = CONTAINER_TRACKER.lock() {
46        tracker.insert(id);
47    }
48}
49
50/// Removes a container ID from tracking (when it's been removed naturally).
51fn untrack_container(id: &str) {
52    if let Ok(mut tracker) = CONTAINER_TRACKER.lock() {
53        tracker.remove(id);
54    }
55}
56
57/// Cleans up all tracked containers. Call this on application shutdown.
58pub async fn cleanup_tracked_containers() {
59    let ids: Vec<String> = {
60        match CONTAINER_TRACKER.lock() {
61            Ok(tracker) => tracker.iter().cloned().collect(),
62            Err(_) => return,
63        }
64    };
65
66    if ids.is_empty() {
67        return;
68    }
69
70    tracing::info!("Cleaning up {} tracked container(s)", ids.len());
71
72    let docker = match Docker::connect_with_local_defaults() {
73        Ok(d) => d,
74        Err(e) => {
75            // log-policy: system-broken
76            tracing::error!("Failed to connect to Docker for cleanup: {}", e);
77            return;
78        }
79    };
80
81    for id in ids {
82        match docker
83            .remove_container(
84                &id,
85                Some(RemoveContainerOptions {
86                    force: true,
87                    ..Default::default()
88                }),
89            )
90            .await
91        {
92            Ok(_) => {
93                tracing::debug!("Cleaned up container {}", id);
94                untrack_container(&id);
95            }
96            Err(e) => {
97                tracing::warn!("Failed to cleanup container {}: {}", id, e);
98            }
99        }
100    }
101}
102
103// Header constants for container operations
104
105/// Timeout (seconds) for connecting to the Docker daemon.
106const DOCKER_CONNECT_TIMEOUT_SECS: u64 = 120;
107
108/// Header key for specifying the container action (e.g., "list", "run", "start", "stop", "remove").
109pub const HEADER_ACTION: &str = "CamelContainerAction";
110
111/// Header key for specifying the container image to use for "run" operations.
112pub const HEADER_IMAGE: &str = "CamelContainerImage";
113
114/// Header key for specifying or receiving the container ID.
115pub const HEADER_CONTAINER_ID: &str = "CamelContainerId";
116
117/// Header key for the log stream type (stdout or stderr).
118pub const HEADER_LOG_STREAM: &str = "CamelContainerLogStream";
119
120/// Header key for the log timestamp.
121pub const HEADER_LOG_TIMESTAMP: &str = "CamelContainerLogTimestamp";
122
123/// Header key for specifying the container name for "run" operations.
124pub const HEADER_CONTAINER_NAME: &str = "CamelContainerName";
125
126/// Header key for the result status of a container operation (e.g., "success").
127pub const HEADER_ACTION_RESULT: &str = "CamelContainerActionResult";
128
129/// Header key for specifying the command to execute in a container.
130pub const HEADER_CMD: &str = "CamelContainerCmd";
131
132/// Header key for specifying the network name for network operations.
133pub const HEADER_NETWORK: &str = "CamelContainerNetwork";
134
135/// Header key for the exit code of an exec operation.
136pub const HEADER_EXIT_CODE: &str = "CamelContainerExitCode";
137
138/// Header key for specifying volume mounts.
139pub const HEADER_VOLUMES: &str = "CamelContainerVolumes";
140
141/// Header key for the exec instance ID.
142pub const HEADER_EXEC_ID: &str = "CamelContainerExecId";
143
144// ---------------------------------------------------------------------------
145// ContainerGlobalConfig
146// ---------------------------------------------------------------------------
147
148/// Per-component reconnect default: unlimited retries (max_attempts=0)
149/// with a fixed 5s delay, preserving the old infinite-reconnect behavior.
150/// Operators can opt into bounded retry via TOML `[reconnect]`.
151fn container_reconnect_default() -> NetworkRetryPolicy {
152    NetworkRetryPolicy {
153        enabled: true,
154        max_attempts: 0, // unlimited
155        initial_delay: std::time::Duration::from_secs(5),
156        multiplier: 1.0, // fixed delay (old behavior was no backoff)
157        max_delay: std::time::Duration::from_secs(5),
158        jitter_factor: 0.0,
159        max_attempts_absolute: None,
160    }
161}
162
163/// Global configuration for Container component.
164/// Supports serde deserialization with defaults and builder methods.
165/// These are the fallback defaults when URI params are not set.
166#[derive(Debug, Clone, PartialEq, serde::Deserialize)]
167#[serde(default)]
168pub struct ContainerGlobalConfig {
169    /// The Docker host URL (default: "unix:///var/run/docker.sock").
170    pub docker_host: String,
171    /// Reconnection policy for events/logs consumers (default: unlimited, 5s fixed delay).
172    #[serde(default = "container_reconnect_default")]
173    pub reconnect: NetworkRetryPolicy,
174}
175
176impl Default for ContainerGlobalConfig {
177    fn default() -> Self {
178        Self {
179            docker_host: "unix:///var/run/docker.sock".to_string(),
180            reconnect: container_reconnect_default(),
181        }
182    }
183}
184
185impl ContainerGlobalConfig {
186    pub fn new() -> Self {
187        Self::default()
188    }
189
190    pub fn with_docker_host(mut self, v: impl Into<String>) -> Self {
191        self.docker_host = v.into();
192        self
193    }
194}
195
196// ---------------------------------------------------------------------------
197// ContainerConfig (endpoint configuration)
198// ---------------------------------------------------------------------------
199
200/// Configuration for the container component endpoint.
201///
202/// This struct holds the parsed URI configuration including the operation type,
203/// optional container image, and Docker host connection details.
204#[derive(Debug, Clone)]
205pub struct ContainerConfig {
206    /// The operation to perform (e.g., "list", "run", "start", "stop", "remove", "events").
207    pub operation: String,
208    /// The container image to use for "run" operations (can be overridden via header).
209    pub image: Option<String>,
210    /// The container name to use for "run" operations (can be overridden via header).
211    pub name: Option<String>,
212    /// The Docker host URL (defaults to "unix:///var/run/docker.sock").
213    pub host: Option<String>,
214    /// Command to run in the container (e.g., "sleep 30").
215    pub cmd: Option<String>,
216    /// Port mappings in format "hostPort:containerPort" (e.g., "8080:80,8443:443").
217    pub ports: Option<String>,
218    /// Environment variables in format "KEY=value,KEY2=value2".
219    pub env: Option<String>,
220    /// Network mode (e.g., "bridge", "host", "none"). Default: "bridge".
221    pub network: Option<String>,
222    /// Container ID or name for logs consumer.
223    pub container_id: Option<String>,
224    /// Follow log output (default: true for consumer).
225    pub follow: bool,
226    /// Include timestamps in logs (default: false).
227    pub timestamps: bool,
228    /// Number of lines to show from the end of logs (default: all).
229    pub tail: Option<String>,
230    /// Automatically pull the image if not present (default: true).
231    pub auto_pull: bool,
232    /// Automatically remove the container when it exits (default: true).
233    pub auto_remove: bool,
234    /// Volume mounts in format "host:container:ro" (e.g., "./html:/usr/share/nginx/html:ro").
235    pub volumes: Option<String>,
236    /// User to run the container or exec command as (e.g., "root").
237    pub user: Option<String>,
238    /// Working directory inside the container.
239    pub workdir: Option<String>,
240    /// Whether to detach from the exec process (default: false).
241    pub detach: bool,
242    /// Network driver for network-create (e.g., "bridge", "overlay").
243    pub driver: Option<String>,
244    /// Whether to force the operation (default: false).
245    pub force: bool,
246    /// Reconnection policy for events/logs consumers (applied from global config).
247    pub reconnect: NetworkRetryPolicy,
248}
249
250impl ContainerConfig {
251    /// Parses a container URI into a `ContainerConfig`.
252    ///
253    /// # Arguments
254    /// * `uri` - The URI to parse (e.g., "container:run?image=alpine")
255    ///
256    /// # Errors
257    /// Returns an error if the URI scheme is not "container".
258    pub fn from_uri(uri: &str) -> Result<Self, CamelError> {
259        let parts = parse_uri(uri)?;
260        if parts.scheme != "container" {
261            return Err(CamelError::InvalidUri(format!(
262                "expected scheme 'container', got '{}'",
263                parts.scheme
264            )));
265        }
266
267        let image = parts.params.get("image").cloned();
268        let name = parts.params.get("name").cloned();
269        let cmd = parts.params.get("cmd").cloned();
270        let ports = parts.params.get("ports").cloned();
271        let env = parts.params.get("env").cloned();
272        let network = parts.params.get("network").cloned();
273        let container_id = parts.params.get("containerId").cloned();
274        let follow = parts
275            .params
276            .get("follow")
277            .map(|v| v.eq_ignore_ascii_case("true"))
278            .unwrap_or(true);
279        let timestamps = parts
280            .params
281            .get("timestamps")
282            .map(|v| v.eq_ignore_ascii_case("true"))
283            .unwrap_or(false);
284        let tail = parts.params.get("tail").cloned();
285        let auto_pull = parts
286            .params
287            .get("autoPull")
288            .map(|v| v.eq_ignore_ascii_case("true"))
289            .unwrap_or(true);
290        let auto_remove = parts
291            .params
292            .get("autoRemove")
293            .map(|v| v.eq_ignore_ascii_case("true"))
294            .unwrap_or(true);
295        // host is only set from URI param; global config defaults are applied later
296        let host = parts.params.get("host").cloned();
297        let volumes = parts.params.get("volumes").cloned();
298        let user = parts.params.get("user").cloned();
299        let workdir = parts.params.get("workdir").cloned();
300        let detach = parts
301            .params
302            .get("detach")
303            .map(|v| v.eq_ignore_ascii_case("true"))
304            .unwrap_or(false);
305        let driver = parts.params.get("driver").cloned();
306        let force = parts
307            .params
308            .get("force")
309            .map(|v| v.eq_ignore_ascii_case("true"))
310            .unwrap_or(false);
311
312        Ok(Self {
313            operation: parts.path,
314            image,
315            name,
316            host,
317            cmd,
318            ports,
319            env,
320            network,
321            container_id,
322            follow,
323            timestamps,
324            tail,
325            auto_pull,
326            auto_remove,
327            volumes,
328            user,
329            workdir,
330            detach,
331            driver,
332            force,
333            reconnect: NetworkRetryPolicy::default(),
334        })
335    }
336
337    /// Apply global config defaults to this endpoint config.
338    /// Only sets values that are currently `None`.
339    fn apply_global_defaults(&mut self, global: &ContainerGlobalConfig) {
340        if self.host.is_none() {
341            self.host = Some(global.docker_host.clone());
342        }
343        self.reconnect = global.reconnect.clone();
344    }
345
346    fn docker_socket_path(&self) -> Result<&str, CamelError> {
347        let host = self.host.as_deref().unwrap_or(if cfg!(windows) {
348            "npipe:////./pipe/docker_engine"
349        } else {
350            "unix:///var/run/docker.sock"
351        });
352
353        if host.starts_with("unix://") || host.starts_with("npipe://") {
354            return Ok(host);
355        }
356
357        if host.contains("://") {
358            return Err(CamelError::ProcessorError(format!(
359                "Unsupported Docker host scheme: {} (only unix:// and npipe:// are supported)",
360                host
361            )));
362        }
363
364        Ok(host)
365    }
366
367    pub fn connect_docker_client(&self) -> Result<Docker, CamelError> {
368        let socket_path = self.docker_socket_path()?;
369        Docker::connect_with_socket(
370            socket_path,
371            DOCKER_CONNECT_TIMEOUT_SECS,
372            bollard::API_DEFAULT_VERSION,
373        )
374        .map_err(|e| {
375            CamelError::ProcessorError(format!("Failed to connect to docker daemon: {}", e))
376        })
377    }
378
379    /// Connects to the Docker daemon using the configured host.
380    ///
381    /// This method establishes a Unix socket connection to Docker and verifies
382    /// the connection by sending a ping request.
383    ///
384    /// # Errors
385    /// Returns an error if the connection fails or the ping request fails.
386    pub async fn connect_docker(&self) -> Result<Docker, CamelError> {
387        let docker = self.connect_docker_client()?;
388        docker
389            .ping()
390            .await
391            .map_err(|e| CamelError::ProcessorError(format!("Docker ping failed: {}", e)))?;
392        Ok(docker)
393    }
394
395    #[allow(clippy::type_complexity)]
396    fn parse_ports(
397        &self,
398    ) -> Result<(Vec<String>, HashMap<String, Option<Vec<PortBinding>>>), CamelError> {
399        let ports_str = match self.ports.as_ref() {
400            Some(s) => s,
401            None => return Ok((Vec::new(), HashMap::new())),
402        };
403
404        let mut exposed_ports: Vec<String> = Vec::new();
405        let mut port_bindings: HashMap<String, Option<Vec<PortBinding>>> = HashMap::new();
406
407        for mapping in ports_str.split(',') {
408            let mapping = mapping.trim();
409            if mapping.is_empty() {
410                continue;
411            }
412
413            let (host_port, container_spec) = mapping.split_once(':').ok_or_else(|| {
414                CamelError::ProcessorError(format!(
415                    "malformed port mapping '{}': expected hostPort:containerPort",
416                    mapping
417                ))
418            })?;
419
420            let (container_port, protocol) = if container_spec.contains('/') {
421                let parts: Vec<&str> = container_spec.split('/').collect();
422                (parts[0], parts[1])
423            } else {
424                (container_spec, "tcp")
425            };
426
427            let container_key = format!("{}/{}", container_port, protocol);
428
429            exposed_ports.push(container_key.clone());
430
431            port_bindings.insert(
432                container_key,
433                Some(vec![PortBinding {
434                    host_ip: None,
435                    host_port: Some(host_port.to_string()),
436                }]),
437            );
438        }
439
440        Ok((exposed_ports, port_bindings))
441    }
442
443    fn parse_env(&self) -> Option<Vec<String>> {
444        let env_str = self.env.as_ref()?;
445
446        let env_vars: Vec<String> = env_str
447            .split(',')
448            .map(|s| s.trim().to_string())
449            .filter(|s| !s.is_empty())
450            .collect();
451
452        if env_vars.is_empty() {
453            None
454        } else {
455            Some(env_vars)
456        }
457    }
458
459    #[cfg(test)]
460    #[allow(clippy::type_complexity)]
461    fn parse_volumes(&self) -> Option<(Vec<String>, Vec<String>)> {
462        self.volumes.as_deref().and_then(parse_volume_str)
463    }
464}
465
466/// Parses a volume specification string into bind mounts and anonymous volumes.
467///
468/// Format: `host:container:ro|rw` for bind mounts, `path` for anonymous volumes,
469/// `path:ro|rw` for anonymous volumes with mode, or `name:container` for named volumes.
470#[allow(clippy::type_complexity)]
471fn parse_volume_str(volumes_str: &str) -> Option<(Vec<String>, Vec<String>)> {
472    let mut binds: Vec<String> = Vec::new();
473    let mut anonymous_volumes: Vec<String> = Vec::new();
474
475    for entry in volumes_str.split(',') {
476        let entry = entry.trim();
477        if entry.is_empty() {
478            continue;
479        }
480
481        let segments: Vec<&str> = entry.split(':').collect();
482
483        match segments.len() {
484            3 => {
485                let source = segments[0];
486                let target = segments[1];
487                let mode = segments[2];
488                if mode != "ro" && mode != "rw" {
489                    continue;
490                }
491                binds.push(format!("{}:{}:{}", source, target, mode));
492            }
493            2 => {
494                let a = segments[0];
495                let b = segments[1];
496                if b == "ro" || b == "rw" {
497                    anonymous_volumes.push(a.to_string());
498                } else {
499                    binds.push(format!("{}:{}", a, b));
500                }
501            }
502            1 => {
503                anonymous_volumes.push(segments[0].to_string());
504            }
505            _ => continue,
506        }
507    }
508
509    if binds.is_empty() && anonymous_volumes.is_empty() {
510        None
511    } else {
512        Some((binds, anonymous_volumes))
513    }
514}
515
516#[derive(Debug, Clone, Copy, PartialEq, Eq)]
517enum ProducerOperation {
518    List,
519    Run,
520    Start,
521    Stop,
522    Remove,
523    Exec,
524    NetworkCreate,
525    NetworkConnect,
526    NetworkDisconnect,
527    NetworkRemove,
528    NetworkList,
529}
530
531fn parse_producer_operation(operation: &str) -> Result<ProducerOperation, CamelError> {
532    match operation {
533        "list" => Ok(ProducerOperation::List),
534        "run" => Ok(ProducerOperation::Run),
535        "start" => Ok(ProducerOperation::Start),
536        "stop" => Ok(ProducerOperation::Stop),
537        "remove" => Ok(ProducerOperation::Remove),
538        "exec" => Ok(ProducerOperation::Exec),
539        "network-create" => Ok(ProducerOperation::NetworkCreate),
540        "network-connect" => Ok(ProducerOperation::NetworkConnect),
541        "network-disconnect" => Ok(ProducerOperation::NetworkDisconnect),
542        "network-remove" => Ok(ProducerOperation::NetworkRemove),
543        "network-list" => Ok(ProducerOperation::NetworkList),
544        _ => Err(CamelError::ProcessorError(format!(
545            "Unknown container operation: {}",
546            operation
547        ))),
548    }
549}
550
551fn resolve_container_name(exchange: &Exchange, config: &ContainerConfig) -> Option<String> {
552    exchange
553        .input
554        .header(HEADER_CONTAINER_NAME)
555        .and_then(|v| v.as_str().map(|s| s.to_string()))
556        .or_else(|| config.name.clone())
557}
558
559async fn image_exists_locally(docker: &Docker, image: &str) -> Result<bool, CamelError> {
560    let images = docker
561        .list_images(None::<ListImagesOptions>)
562        .await
563        .map_err(|e| CamelError::ProcessorError(format!("Failed to list images: {}", e)))?;
564
565    Ok(images.iter().any(|img| {
566        img.repo_tags
567            .iter()
568            .any(|tag| tag == image || tag.starts_with(&format!("{}:", image)))
569    }))
570}
571
572async fn pull_image_with_progress(
573    docker: &Docker,
574    image: &str,
575    timeout_secs: u64,
576) -> Result<(), CamelError> {
577    use futures::StreamExt;
578
579    tracing::info!("Pulling image: {}", image);
580
581    let mut stream = docker.create_image(
582        Some(CreateImageOptions {
583            from_image: Some(image.to_string()),
584            ..Default::default()
585        }),
586        None,
587        None,
588    );
589
590    let start = std::time::Instant::now();
591    let mut last_progress = std::time::Instant::now();
592
593    while let Some(item) = stream.next().await {
594        if start.elapsed().as_secs() > timeout_secs {
595            return Err(CamelError::ProcessorError(format!(
596                "Image pull timeout after {}s. Try manually: docker pull {}",
597                timeout_secs, image
598            )));
599        }
600
601        match item {
602            Ok(update) => {
603                // Log progress every 2 seconds
604                if last_progress.elapsed().as_secs() >= 2 {
605                    if let Some(status) = update.status {
606                        tracing::debug!("Pull progress: {}", status);
607                    }
608                    last_progress = std::time::Instant::now();
609                }
610            }
611            Err(e) => {
612                let err_str = e.to_string().to_lowercase();
613                if err_str.contains("unauthorized") || err_str.contains("401") {
614                    return Err(CamelError::ProcessorError(format!(
615                        "Authentication required for image '{}'. Configure Docker credentials: docker login",
616                        image
617                    )));
618                }
619                if err_str.contains("not found") || err_str.contains("404") {
620                    return Err(CamelError::ProcessorError(format!(
621                        "Image '{}' not found in registry. Check the image name and tag",
622                        image
623                    )));
624                }
625                return Err(CamelError::ProcessorError(format!(
626                    "Failed to pull image '{}': {}",
627                    image, e
628                )));
629            }
630        }
631    }
632
633    tracing::info!("Successfully pulled image: {}", image);
634    Ok(())
635}
636
637async fn ensure_image_available(
638    docker: &Docker,
639    image: &str,
640    auto_pull: bool,
641    timeout_secs: u64,
642) -> Result<(), CamelError> {
643    if image_exists_locally(docker, image).await? {
644        tracing::debug!("Image '{}' already available locally", image);
645        return Ok(());
646    }
647
648    if !auto_pull {
649        return Err(CamelError::ProcessorError(format!(
650            "Image '{}' not found locally. Set autoPull=true to pull automatically, or run: docker pull {}",
651            image, image
652        )));
653    }
654
655    pull_image_with_progress(docker, image, timeout_secs).await
656}
657
658fn format_docker_event(event: &bollard::models::EventMessage) -> String {
659    let action = event.action.as_deref().unwrap_or("unknown");
660    let actor = event.actor.as_ref();
661
662    let container_name = actor
663        .and_then(|a| a.attributes.as_ref())
664        .and_then(|attrs| attrs.get("name"))
665        .map(|s| s.as_str())
666        .unwrap_or("unknown");
667
668    let image = actor
669        .and_then(|a| a.attributes.as_ref())
670        .and_then(|attrs| attrs.get("image"))
671        .map(|s| s.as_str())
672        .unwrap_or("");
673
674    let exit_code = actor
675        .and_then(|a| a.attributes.as_ref())
676        .and_then(|attrs| attrs.get("exitCode"))
677        .map(|s| s.as_str());
678
679    match action {
680        "create" => {
681            if image.is_empty() {
682                format!("[CREATE] Container {}", container_name)
683            } else {
684                format!("[CREATE] Container {} ({})", container_name, image)
685            }
686        }
687        "start" => format!("[START]  Container {}", container_name),
688        "die" => {
689            if let Some(code) = exit_code {
690                format!("[DIE]    Container {} (exit: {})", container_name, code)
691            } else {
692                format!("[DIE]    Container {}", container_name)
693            }
694        }
695        "destroy" => format!("[DESTROY] Container {}", container_name),
696        "stop" => format!("[STOP]   Container {}", container_name),
697        "pause" => format!("[PAUSE]  Container {}", container_name),
698        "unpause" => format!("[UNPAUSE] Container {}", container_name),
699        "restart" => format!("[RESTART] Container {}", container_name),
700        _ => format!("[{}] Container {}", action.to_uppercase(), container_name),
701    }
702}
703
704async fn run_container_with_cleanup<CreateFn, CreateFut, StartFn, StartFut, RemoveFn, RemoveFut>(
705    create: CreateFn,
706    start: StartFn,
707    remove: RemoveFn,
708) -> Result<String, CamelError>
709where
710    CreateFn: FnOnce() -> CreateFut,
711    CreateFut: Future<Output = Result<String, CamelError>>,
712    StartFn: FnOnce(String) -> StartFut,
713    StartFut: Future<Output = Result<(), CamelError>>,
714    RemoveFn: FnOnce(String) -> RemoveFut,
715    RemoveFut: Future<Output = Result<(), CamelError>>,
716{
717    let container_id = create().await?;
718    if let Err(start_err) = start(container_id.clone()).await {
719        if let Err(remove_err) = remove(container_id.clone()).await {
720            return Err(CamelError::ProcessorError(format!(
721                "Failed to start container: {}. Cleanup failed: {}",
722                start_err, remove_err
723            )));
724        }
725        return Err(start_err);
726    }
727
728    Ok(container_id)
729}
730
731async fn handle_list(
732    docker: Docker,
733    _config: ContainerConfig,
734    exchange: &mut Exchange,
735) -> Result<(), CamelError> {
736    let containers = docker
737        .list_containers(None::<ListContainersOptions>)
738        .await
739        .map_err(|e| CamelError::ProcessorError(format!("Failed to list containers: {}", e)))?;
740
741    let json_value = serde_json::to_value(&containers).map_err(|e| {
742        CamelError::ProcessorError(format!("Failed to serialize containers: {}", e))
743    })?;
744
745    exchange.input.body = Body::Json(json_value);
746    exchange.input.set_header(
747        HEADER_ACTION_RESULT,
748        serde_json::Value::String("success".to_string()),
749    );
750    Ok(())
751}
752
753async fn handle_run(
754    docker: Docker,
755    config: ContainerConfig,
756    exchange: &mut Exchange,
757) -> Result<(), CamelError> {
758    let image = exchange
759        .input
760        .header(HEADER_IMAGE)
761        .and_then(|v| v.as_str().map(|s| s.to_string()))
762        .or(config.image.clone())
763        .ok_or_else(|| {
764            CamelError::ProcessorError(
765                "missing image for run operation. Specify in URI (image=alpine) or header (CamelContainerImage)".to_string(),
766            )
767        })?;
768
769    let image = if !image.contains(':') && !image.contains('@') {
770        format!("{}:latest", image)
771    } else {
772        image
773    };
774
775    let pull_timeout = 300;
776    ensure_image_available(&docker, &image, config.auto_pull, pull_timeout)
777        .await
778        .map_err(|e| {
779            CamelError::ProcessorError(format!("Image '{}' not available: {}", image, e))
780        })?;
781
782    let container_name = resolve_container_name(exchange, &config);
783    let container_name_ref = container_name.as_deref().unwrap_or("");
784    let cmd_parts: Option<Vec<String>> = config
785        .cmd
786        .as_ref()
787        .map(|c| c.split_whitespace().map(|s| s.to_string()).collect());
788    let auto_remove = config.auto_remove;
789    let (exposed_ports, port_bindings) = config.parse_ports()?;
790    let env_vars = config.parse_env();
791    let network_mode = config.network.clone();
792
793    let volumes_str = exchange
794        .input
795        .header(HEADER_VOLUMES)
796        .and_then(|v| v.as_str().map(|s| s.to_string()))
797        .or(config.volumes.clone());
798    let (binds, anon_volumes) = volumes_str
799        .as_deref()
800        .and_then(parse_volume_str)
801        .unwrap_or_default();
802
803    let docker_create = docker.clone();
804    let docker_start = docker.clone();
805    let docker_remove = docker.clone();
806
807    let container_id = run_container_with_cleanup(
808        move || async move {
809            let create_options = CreateContainerOptions {
810                name: Some(container_name_ref.to_string()),
811                ..Default::default()
812            };
813            let container_config = ContainerCreateBody {
814                image: Some(image.clone()),
815                cmd: cmd_parts,
816                env: env_vars,
817                exposed_ports: if exposed_ports.is_empty() { None } else { Some(exposed_ports) },
818                volumes: if anon_volumes.is_empty() { None } else { Some(anon_volumes) },
819                host_config: Some(HostConfig {
820                    auto_remove: Some(auto_remove),
821                    port_bindings: if port_bindings.is_empty() { None } else { Some(port_bindings) },
822                    network_mode,
823                    binds: if binds.is_empty() { None } else { Some(binds) },
824                    ..Default::default()
825                }),
826                ..Default::default()
827            };
828
829            let create_response = docker_create
830                .create_container(Some(create_options), container_config)
831                .await
832                .map_err(|e| {
833                    let err_str = e.to_string().to_lowercase();
834                    if err_str.contains("409") || err_str.contains("conflict") {
835                        CamelError::ProcessorError(format!(
836                            "Container name '{}' already exists. Use a unique name or remove the existing container first",
837                            container_name_ref
838                        ))
839                    } else {
840                        CamelError::ProcessorError(format!(
841                            "Failed to create container: {}",
842                            e
843                        ))
844                    }
845                })?;
846
847            Ok(create_response.id)
848        },
849        move |container_id| async move {
850            docker_start
851                .start_container(&container_id, None::<StartContainerOptions>)
852                .await
853                .map_err(|e| {
854                    CamelError::ProcessorError(format!(
855                        "Failed to start container: {}",
856                        e
857                    ))
858                })
859        },
860        move |container_id| async move {
861            docker_remove
862                .remove_container(&container_id, None)
863                .await
864                .map_err(|e| {
865                    CamelError::ProcessorError(format!(
866                        "Failed to remove container after start failure: {}",
867                        e
868                    ))
869                })
870        },
871    )
872    .await?;
873
874    track_container(container_id.clone());
875
876    exchange
877        .input
878        .set_header(HEADER_CONTAINER_ID, serde_json::Value::String(container_id));
879    exchange.input.set_header(
880        HEADER_ACTION_RESULT,
881        serde_json::Value::String("success".to_string()),
882    );
883    Ok(())
884}
885
886async fn handle_lifecycle(
887    docker: Docker,
888    _config: ContainerConfig,
889    exchange: &mut Exchange,
890    operation: ProducerOperation,
891    operation_name: &str,
892) -> Result<(), CamelError> {
893    let container_id = exchange
894        .input
895        .header(HEADER_CONTAINER_ID)
896        .and_then(|v| v.as_str().map(|s| s.to_string()))
897        .ok_or_else(|| {
898            CamelError::ProcessorError(format!(
899                "{} header is required for {} operation",
900                HEADER_CONTAINER_ID, operation_name
901            ))
902        })?;
903
904    match operation {
905        ProducerOperation::Start => {
906            docker
907                .start_container(&container_id, None::<StartContainerOptions>)
908                .await
909                .map_err(|e| {
910                    CamelError::ProcessorError(format!("Failed to start container: {}", e))
911                })?;
912        }
913        ProducerOperation::Stop => {
914            docker
915                .stop_container(&container_id, None)
916                .await
917                .map_err(|e| {
918                    CamelError::ProcessorError(format!("Failed to stop container: {}", e))
919                })?;
920        }
921        ProducerOperation::Remove => {
922            docker
923                .remove_container(&container_id, None)
924                .await
925                .map_err(|e| {
926                    CamelError::ProcessorError(format!("Failed to remove container: {}", e))
927                })?;
928            untrack_container(&container_id);
929        }
930        _ => {}
931    }
932
933    exchange.input.set_header(
934        HEADER_ACTION_RESULT,
935        serde_json::Value::String("success".to_string()),
936    );
937    Ok(())
938}
939
940async fn handle_exec(
941    docker: Docker,
942    config: ContainerConfig,
943    exchange: &mut Exchange,
944) -> Result<(), CamelError> {
945    let container_id = exchange
946        .input
947        .header(HEADER_CONTAINER_ID)
948        .and_then(|v| v.as_str().map(|s| s.to_string()))
949        .or(config.container_id.clone())
950        .ok_or_else(|| {
951            CamelError::ProcessorError(format!(
952                "{} header or containerId param is required for exec operation",
953                HEADER_CONTAINER_ID
954            ))
955        })?;
956
957    let cmd = exchange
958        .input
959        .header(HEADER_CMD)
960        .and_then(|v| v.as_str().map(|s| s.to_string()))
961        .or(config.cmd.clone())
962        .ok_or_else(|| {
963            CamelError::ProcessorError(
964                "CamelContainerCmd header or cmd param is required for exec operation".to_string(),
965            )
966        })?;
967
968    let cmd_parts: Vec<String> = cmd.split_whitespace().map(|s| s.to_string()).collect();
969    let env_vars = config.parse_env();
970
971    let exec_config = bollard::exec::CreateExecOptions {
972        cmd: Some(cmd_parts),
973        env: env_vars,
974        user: config.user.clone(),
975        working_dir: config.workdir.clone(),
976        attach_stdout: Some(true),
977        attach_stderr: Some(true),
978        ..Default::default()
979    };
980
981    let create_result = docker
982        .create_exec(&container_id, exec_config)
983        .await
984        .map_err(|e| {
985            let err_str = e.to_string().to_lowercase();
986            if err_str.contains("404") || err_str.contains("no such") {
987                CamelError::ProcessorError(format!(
988                    "Container '{}' not found for exec",
989                    container_id
990                ))
991            } else {
992                CamelError::ProcessorError(format!("Failed to create exec: {}", e))
993            }
994        })?;
995
996    let exec_id = create_result.id;
997
998    if config.detach {
999        docker
1000            .start_exec(
1001                &exec_id,
1002                Some(bollard::exec::StartExecOptions {
1003                    detach: true,
1004                    ..Default::default()
1005                }),
1006            )
1007            .await
1008            .map_err(|e| {
1009                CamelError::ProcessorError(format!("Failed to start exec (detached): {}", e))
1010            })?;
1011
1012        exchange
1013            .input
1014            .set_header(HEADER_EXEC_ID, serde_json::Value::String(exec_id));
1015        exchange
1016            .input
1017            .set_header(HEADER_CONTAINER_ID, serde_json::Value::String(container_id));
1018    } else {
1019        let start_result = docker
1020            .start_exec(&exec_id, None)
1021            .await
1022            .map_err(|e| CamelError::ProcessorError(format!("Failed to start exec: {}", e)))?;
1023
1024        let mut output = String::new();
1025
1026        match start_result {
1027            bollard::exec::StartExecResults::Attached {
1028                output: mut stream, ..
1029            } => {
1030                use futures::StreamExt;
1031                while let Some(msg) = stream.next().await {
1032                    match msg {
1033                        Ok(bollard::container::LogOutput::StdOut { message }) => {
1034                            output.push_str(&String::from_utf8_lossy(&message));
1035                        }
1036                        Ok(bollard::container::LogOutput::StdErr { message }) => {
1037                            output.push_str(&String::from_utf8_lossy(&message));
1038                        }
1039                        Ok(_) => {}
1040                        Err(e) => {
1041                            output.push_str(&format!("[error reading stream: {}]", e));
1042                        }
1043                    }
1044                }
1045            }
1046            bollard::exec::StartExecResults::Detached => {}
1047        }
1048
1049        let inspect = docker
1050            .inspect_exec(&exec_id)
1051            .await
1052            .map_err(|e| CamelError::ProcessorError(format!("Failed to inspect exec: {}", e)))?;
1053
1054        let exit_code: i64 = inspect.exit_code.ok_or_else(|| {
1055            CamelError::ProcessorError("container exec returned no exit code".into())
1056        })?;
1057
1058        let output = output.trim_end().to_string();
1059        exchange.input.body = Body::Text(output);
1060        exchange.input.set_header(
1061            HEADER_EXIT_CODE,
1062            serde_json::Value::Number(exit_code.into()),
1063        );
1064        exchange
1065            .input
1066            .set_header(HEADER_CONTAINER_ID, serde_json::Value::String(container_id));
1067    }
1068
1069    exchange.input.set_header(
1070        HEADER_ACTION_RESULT,
1071        serde_json::Value::String("success".to_string()),
1072    );
1073    Ok(())
1074}
1075
1076async fn handle_network_create(
1077    docker: Docker,
1078    config: ContainerConfig,
1079    exchange: &mut Exchange,
1080) -> Result<(), CamelError> {
1081    let network_name = exchange
1082        .input
1083        .header(HEADER_CONTAINER_NAME)
1084        .and_then(|v| v.as_str().map(|s| s.to_string()))
1085        .or(config.name.clone())
1086        .ok_or_else(|| {
1087            CamelError::ProcessorError(
1088                "CamelContainerName header or name param is required for network-create"
1089                    .to_string(),
1090            )
1091        })?;
1092
1093    let driver = config.driver.as_deref().unwrap_or("bridge");
1094
1095    let options = NetworkCreateRequest {
1096        name: network_name.clone(),
1097        driver: Some(driver.to_string()),
1098        ..Default::default()
1099    };
1100
1101    let result = docker.create_network(options).await.map_err(|e| {
1102        let err_str = e.to_string().to_lowercase();
1103        if err_str.contains("409") || err_str.contains("already exists") {
1104            CamelError::ProcessorError(format!("Network '{}' already exists", network_name))
1105        } else {
1106            CamelError::ProcessorError(format!("Failed to create network: {}", e))
1107        }
1108    })?;
1109
1110    let network_id = result.id.clone();
1111    let json_value = serde_json::to_value(&result).map_err(|e| {
1112        CamelError::ProcessorError(format!("Failed to serialize network response: {}", e))
1113    })?;
1114
1115    exchange.input.body = Body::Json(json_value);
1116    exchange
1117        .input
1118        .set_header(HEADER_NETWORK, serde_json::Value::String(network_id));
1119    exchange.input.set_header(
1120        HEADER_ACTION_RESULT,
1121        serde_json::Value::String("success".to_string()),
1122    );
1123    Ok(())
1124}
1125
1126async fn handle_network_connect(
1127    docker: Docker,
1128    config: ContainerConfig,
1129    exchange: &mut Exchange,
1130) -> Result<(), CamelError> {
1131    let network = exchange
1132        .input
1133        .header(HEADER_NETWORK)
1134        .and_then(|v| v.as_str().map(|s| s.to_string()))
1135        .or(config.network.clone())
1136        .ok_or_else(|| {
1137            CamelError::ProcessorError(
1138                "CamelContainerNetwork header or network param is required for network-connect"
1139                    .to_string(),
1140            )
1141        })?;
1142
1143    let container = exchange
1144        .input
1145        .header(HEADER_CONTAINER_ID)
1146        .and_then(|v| v.as_str().map(|s| s.to_string()))
1147        .or(config.container_id.clone())
1148        .ok_or_else(|| {
1149            CamelError::ProcessorError(
1150                "CamelContainerId header or container param is required for network-connect"
1151                    .to_string(),
1152            )
1153        })?;
1154
1155    docker
1156        .connect_network(
1157            &network,
1158            NetworkConnectRequest {
1159                container,
1160                ..Default::default()
1161            },
1162        )
1163        .await
1164        .map_err(|e| {
1165            let err_str = e.to_string().to_lowercase();
1166            if err_str.contains("404") || err_str.contains("not found") {
1167                CamelError::ProcessorError(format!("Network '{}' or container not found", network))
1168            } else {
1169                CamelError::ProcessorError(format!("Failed to connect to network: {}", e))
1170            }
1171        })?;
1172
1173    exchange.input.set_header(
1174        HEADER_ACTION_RESULT,
1175        serde_json::Value::String("success".to_string()),
1176    );
1177    Ok(())
1178}
1179
1180async fn handle_network_disconnect(
1181    docker: Docker,
1182    config: ContainerConfig,
1183    exchange: &mut Exchange,
1184) -> Result<(), CamelError> {
1185    let network = exchange
1186        .input
1187        .header(HEADER_NETWORK)
1188        .and_then(|v| v.as_str().map(|s| s.to_string()))
1189        .or(config.network.clone())
1190        .ok_or_else(|| {
1191            CamelError::ProcessorError(
1192                "CamelContainerNetwork header or network param is required for network-disconnect"
1193                    .to_string(),
1194            )
1195        })?;
1196
1197    let container = exchange
1198        .input
1199        .header(HEADER_CONTAINER_ID)
1200        .and_then(|v| v.as_str().map(|s| s.to_string()))
1201        .or(config.container_id.clone())
1202        .ok_or_else(|| {
1203            CamelError::ProcessorError(
1204                "CamelContainerId header or container param is required for network-disconnect"
1205                    .to_string(),
1206            )
1207        })?;
1208
1209    docker
1210        .disconnect_network(
1211            &network,
1212            NetworkDisconnectRequest {
1213                container,
1214                force: Some(config.force),
1215            },
1216        )
1217        .await
1218        .map_err(|e| {
1219            let err_str = e.to_string().to_lowercase();
1220            if err_str.contains("404") || err_str.contains("not found") {
1221                CamelError::ProcessorError(format!("Network '{}' or container not found", network))
1222            } else {
1223                CamelError::ProcessorError(format!("Failed to disconnect from network: {}", e))
1224            }
1225        })?;
1226
1227    exchange.input.set_header(
1228        HEADER_ACTION_RESULT,
1229        serde_json::Value::String("success".to_string()),
1230    );
1231    Ok(())
1232}
1233
1234async fn handle_network_remove(
1235    docker: Docker,
1236    config: ContainerConfig,
1237    exchange: &mut Exchange,
1238) -> Result<(), CamelError> {
1239    let network = exchange
1240        .input
1241        .header(HEADER_NETWORK)
1242        .and_then(|v| v.as_str().map(|s| s.to_string()))
1243        .or(config.network.clone())
1244        .ok_or_else(|| {
1245            CamelError::ProcessorError(
1246                "CamelContainerNetwork header or network param is required for network-remove"
1247                    .to_string(),
1248            )
1249        })?;
1250
1251    docker.remove_network(&network).await.map_err(|e| {
1252        let err_str = e.to_string().to_lowercase();
1253        if err_str.contains("404") || err_str.contains("not found") {
1254            CamelError::ProcessorError(format!("Network '{}' not found", network))
1255        } else if err_str.contains("409") || err_str.contains("in use") {
1256            CamelError::ProcessorError(format!(
1257                "Network '{}' is in use and cannot be removed",
1258                network
1259            ))
1260        } else {
1261            CamelError::ProcessorError(format!("Failed to remove network: {}", e))
1262        }
1263    })?;
1264
1265    exchange.input.set_header(
1266        HEADER_ACTION_RESULT,
1267        serde_json::Value::String("success".to_string()),
1268    );
1269    Ok(())
1270}
1271
1272async fn handle_network_list(
1273    docker: Docker,
1274    _config: ContainerConfig,
1275    exchange: &mut Exchange,
1276) -> Result<(), CamelError> {
1277    let networks = docker
1278        .list_networks(None::<ListNetworksOptions>)
1279        .await
1280        .map_err(|e| CamelError::ProcessorError(format!("Failed to list networks: {}", e)))?;
1281
1282    let json_value = serde_json::to_value(&networks)
1283        .map_err(|e| CamelError::ProcessorError(format!("Failed to serialize networks: {}", e)))?;
1284
1285    exchange.input.body = Body::Json(json_value);
1286    exchange.input.set_header(
1287        HEADER_ACTION_RESULT,
1288        serde_json::Value::String("success".to_string()),
1289    );
1290    Ok(())
1291}
1292
1293/// Producer for executing container operations.
1294///
1295/// This producer handles synchronous container operations like listing,
1296/// creating, starting, stopping, and removing containers.
1297#[derive(Clone)]
1298pub struct ContainerProducer {
1299    config: ContainerConfig,
1300    docker: Docker,
1301}
1302
1303impl Service<Exchange> for ContainerProducer {
1304    type Response = Exchange;
1305    type Error = CamelError;
1306    type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
1307
1308    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
1309        Poll::Ready(Ok(()))
1310    }
1311
1312    fn call(&mut self, mut exchange: Exchange) -> Self::Future {
1313        let config = self.config.clone();
1314        let docker = self.docker.clone();
1315        Box::pin(async move {
1316            let operation_name = exchange
1317                .input
1318                .header(HEADER_ACTION)
1319                .and_then(|v| v.as_str().map(|s| s.to_string()))
1320                .unwrap_or_else(|| config.operation.clone());
1321
1322            let operation = parse_producer_operation(&operation_name)?;
1323
1324            match operation {
1325                ProducerOperation::List => {
1326                    handle_list(docker, config, &mut exchange).await?;
1327                }
1328                ProducerOperation::Run => {
1329                    handle_run(docker, config, &mut exchange).await?;
1330                }
1331                ProducerOperation::Start => {
1332                    handle_lifecycle(docker, config, &mut exchange, operation, &operation_name)
1333                        .await?;
1334                }
1335                ProducerOperation::Stop => {
1336                    handle_lifecycle(docker, config, &mut exchange, operation, &operation_name)
1337                        .await?;
1338                }
1339                ProducerOperation::Remove => {
1340                    handle_lifecycle(docker, config, &mut exchange, operation, &operation_name)
1341                        .await?;
1342                }
1343                ProducerOperation::Exec => {
1344                    handle_exec(docker, config, &mut exchange).await?;
1345                }
1346                ProducerOperation::NetworkCreate => {
1347                    handle_network_create(docker, config, &mut exchange).await?;
1348                }
1349                ProducerOperation::NetworkConnect => {
1350                    handle_network_connect(docker, config, &mut exchange).await?;
1351                }
1352                ProducerOperation::NetworkDisconnect => {
1353                    handle_network_disconnect(docker, config, &mut exchange).await?;
1354                }
1355                ProducerOperation::NetworkRemove => {
1356                    handle_network_remove(docker, config, &mut exchange).await?;
1357                }
1358                ProducerOperation::NetworkList => {
1359                    handle_network_list(docker, config, &mut exchange).await?;
1360                }
1361            }
1362
1363            Ok(exchange)
1364        })
1365    }
1366}
1367
1368/// Consumer for receiving Docker container events or logs.
1369///
1370/// This consumer subscribes to Docker events or container logs and forwards them
1371/// to the route as exchanges. It implements automatic reconnection on connection failures.
1372pub struct ContainerConsumer {
1373    config: ContainerConfig,
1374    /// ADR-0012 observability handle: `rt.metrics().increment_errors(...)` and
1375    /// `rt.health().force_unhealthy_for_route(...)` calls.
1376    runtime: Arc<dyn RuntimeObservability>,
1377}
1378
1379impl ContainerConsumer {
1380    pub fn new(config: ContainerConfig, runtime: Arc<dyn RuntimeObservability>) -> Self {
1381        Self { config, runtime }
1382    }
1383}
1384
1385#[async_trait]
1386impl Consumer for ContainerConsumer {
1387    async fn start(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
1388        match self.config.operation.as_str() {
1389            "events" => self.start_events_consumer(context).await,
1390            "logs" => self.start_logs_consumer(context).await,
1391            _ => Err(CamelError::EndpointCreationFailed(format!(
1392                "Consumer only supports 'events' or 'logs' operations, got '{}'",
1393                self.config.operation
1394            ))),
1395        }
1396    }
1397
1398    async fn stop(&mut self) -> Result<(), CamelError> {
1399        Ok(())
1400    }
1401
1402    fn concurrency_model(&self) -> camel_component_api::ConcurrencyModel {
1403        camel_component_api::ConcurrencyModel::Concurrent { max: None }
1404    }
1405}
1406
1407impl ContainerConsumer {
1408    async fn start_events_consumer(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
1409        use futures::StreamExt;
1410
1411        let cancel = context.cancel_token();
1412
1413        // Outer reconnect loop: when the inner event stream breaks, this loop
1414        // reconnects. The connect_docker() retry is now backed by
1415        // retry_async_cancelable (migrated from manual loop in rc-k9c).
1416        loop {
1417            let docker = match retry_async_cancelable(
1418                &self.config.reconnect,
1419                Some("container-events"),
1420                || async { self.config.connect_docker().await },
1421                |_| true,
1422                &cancel,
1423            )
1424            .await
1425            {
1426                Ok(d) => d,
1427                Err(_) if context.is_cancelled() => {
1428                    tracing::info!("Container events consumer shutting down");
1429                    return Ok(());
1430                }
1431                Err(e) => {
1432                    self.runtime
1433                        .metrics()
1434                        .increment_errors(context.route_id(), "e:container:events-connect");
1435                    // log-policy: outside-contract
1436                    tracing::error!(error = %e, "Container events consumer exhausted reconnect attempts");
1437                    return Err(e);
1438                }
1439            };
1440
1441            let mut event_stream = docker.events(None::<EventsOptions>);
1442
1443            loop {
1444                tokio::select! {
1445                    _ = context.cancelled() => {
1446                        tracing::info!("Container events consumer shutting down");
1447                        return Ok(());
1448                    }
1449
1450                    msg = event_stream.next() => {
1451                        match msg {
1452                            Some(Ok(event)) => {
1453                                let formatted = format_docker_event(&event);
1454                                let message = Message::new(Body::Text(formatted));
1455                                let exchange = Exchange::new(message);
1456
1457                                if let Err(e) = context.send(exchange).await {
1458                                    // log-policy: system-broken
1459                                    tracing::error!("Failed to send exchange: {:?}", e);
1460                                    break;
1461                                }
1462                            }
1463                            Some(Err(e)) => {
1464                                self.runtime.metrics().increment_errors(context.route_id(), "e:container:events-stream");
1465                                // log-policy: outside-contract
1466                                tracing::error!("Docker event stream error: {}. Reconnecting...", e);
1467                                break;
1468                            }
1469                            None => {
1470                                tracing::info!("Docker event stream ended. Reconnecting...");
1471                                break;
1472                            }
1473                        }
1474                    }
1475                }
1476            }
1477
1478            tokio::select! {
1479                _ = context.cancelled() => {
1480                    tracing::info!("Container events consumer shutting down");
1481                    return Ok(());
1482                }
1483                _ = tokio::time::sleep(std::time::Duration::from_secs(1)) => {}
1484            }
1485        }
1486    }
1487
1488    async fn start_logs_consumer(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
1489        use futures::StreamExt;
1490
1491        let container_id = self.config.container_id.clone().ok_or_else(|| {
1492            CamelError::EndpointCreationFailed(
1493                "containerId is required for logs consumer. Use container:logs?containerId=xxx"
1494                    .to_string(),
1495            )
1496        })?;
1497
1498        let cancel = context.cancel_token();
1499
1500        // Outer reconnect loop: when the inner log stream breaks, this loop
1501        // reconnects. The connect_docker() retry is now backed by
1502        // retry_async_cancelable (migrated from manual loop in rc-k9c).
1503        loop {
1504            let docker = match retry_async_cancelable(
1505                &self.config.reconnect,
1506                Some("container-logs"),
1507                || async { self.config.connect_docker().await },
1508                |_| true,
1509                &cancel,
1510            )
1511            .await
1512            {
1513                Ok(d) => d,
1514                Err(_) if context.is_cancelled() => {
1515                    tracing::info!("Container logs consumer shutting down");
1516                    return Ok(());
1517                }
1518                Err(e) => {
1519                    self.runtime
1520                        .metrics()
1521                        .increment_errors(context.route_id(), "e:container:logs-connect");
1522                    // log-policy: outside-contract
1523                    tracing::error!(error = %e, "Container logs consumer exhausted reconnect attempts");
1524                    return Err(e);
1525                }
1526            };
1527
1528            let tail = self
1529                .config
1530                .tail
1531                .clone()
1532                .unwrap_or_else(|| "all".to_string());
1533
1534            let options = LogsOptions {
1535                follow: self.config.follow,
1536                stdout: true,
1537                stderr: true,
1538                timestamps: self.config.timestamps,
1539                tail,
1540                ..Default::default()
1541            };
1542
1543            let mut log_stream = docker.logs(&container_id, Some(options));
1544            let container_id_header = container_id.clone();
1545
1546            loop {
1547                tokio::select! {
1548                    _ = context.cancelled() => {
1549                        tracing::info!("Container logs consumer shutting down");
1550                        return Ok(());
1551                    }
1552
1553                    msg = log_stream.next() => {
1554                        match msg {
1555                            Some(Ok(log_output)) => {
1556                                let (stream_type, content) = match log_output {
1557                                    bollard::container::LogOutput::StdOut { message } => {
1558                                        ("stdout", String::from_utf8_lossy(&message).into_owned())
1559                                    }
1560                                    bollard::container::LogOutput::StdErr { message } => {
1561                                        ("stderr", String::from_utf8_lossy(&message).into_owned())
1562                                    }
1563                                    bollard::container::LogOutput::Console { message } => {
1564                                        ("console", String::from_utf8_lossy(&message).into_owned())
1565                                    }
1566                                    bollard::container::LogOutput::StdIn { message } => {
1567                                        ("stdin", String::from_utf8_lossy(&message).into_owned())
1568                                    }
1569                                };
1570
1571                                let content = content.trim_end();
1572                                if content.is_empty() {
1573                                    continue;
1574                                }
1575
1576                                let mut message = Message::new(Body::Text(content.to_string()));
1577                                message.set_header(
1578                                    HEADER_CONTAINER_ID,
1579                                    serde_json::Value::String(container_id_header.clone()),
1580                                );
1581                                message.set_header(
1582                                    HEADER_LOG_STREAM,
1583                                    serde_json::Value::String(stream_type.to_string()),
1584                                );
1585
1586                                if self.config.timestamps
1587                                    && let Some(ts) = extract_timestamp(content) {
1588                                        message.set_header(
1589                                            HEADER_LOG_TIMESTAMP,
1590                                            serde_json::Value::String(ts),
1591                                        );
1592                                    }
1593
1594                                let exchange = Exchange::new(message);
1595
1596                                if let Err(e) = context.send(exchange).await {
1597                                    // log-policy: system-broken
1598                                    tracing::error!("Failed to send log exchange: {:?}", e);
1599                                    break;
1600                                }
1601                            }
1602                            Some(Err(e)) => {
1603                                self.runtime.metrics().increment_errors(context.route_id(), "e:container:logs-stream");
1604                                // log-policy: outside-contract
1605                                tracing::error!("Docker log stream error: {}. Reconnecting...", e);
1606                                break;
1607                            }
1608                            None => {
1609                                if self.config.follow {
1610                                    tracing::info!("Docker log stream ended. Reconnecting...");
1611                                    break;
1612                                } else {
1613                                    tracing::info!("Container logs consumer finished (follow=false)");
1614                                    return Ok(());
1615                                }
1616                            }
1617                        }
1618                    }
1619                }
1620            }
1621
1622            tokio::select! {
1623                _ = context.cancelled() => {
1624                    tracing::info!("Container logs consumer shutting down");
1625                    return Ok(());
1626                }
1627                _ = tokio::time::sleep(std::time::Duration::from_secs(1)) => {}
1628            }
1629        }
1630    }
1631}
1632
1633fn extract_timestamp(log_line: &str) -> Option<String> {
1634    let parts: Vec<&str> = log_line.splitn(2, ' ').collect();
1635    if parts.len() > 1 && parts[0].contains('T') {
1636        Some(parts[0].to_string())
1637    } else {
1638        None
1639    }
1640}
1641
1642/// Component for creating container endpoints.
1643///
1644/// This component handles URIs with the "container" scheme and creates
1645/// appropriate producer and consumer endpoints for Docker operations.
1646///
1647/// Containers created via `run` operation are tracked globally and can be
1648/// cleaned up on shutdown by calling `cleanup_tracked_containers()`.
1649pub struct ContainerComponent {
1650    config: Option<ContainerGlobalConfig>,
1651}
1652
1653impl ContainerComponent {
1654    /// Creates a new container component instance without global config.
1655    pub fn new() -> Self {
1656        Self { config: None }
1657    }
1658
1659    /// Creates a container component with the given global config.
1660    pub fn with_config(config: ContainerGlobalConfig) -> Self {
1661        Self {
1662            config: Some(config),
1663        }
1664    }
1665
1666    /// Creates a container component with optional global config.
1667    pub fn with_optional_config(config: Option<ContainerGlobalConfig>) -> Self {
1668        Self { config }
1669    }
1670}
1671
1672impl Default for ContainerComponent {
1673    fn default() -> Self {
1674        Self::new()
1675    }
1676}
1677
1678impl Component for ContainerComponent {
1679    fn scheme(&self) -> &str {
1680        "container"
1681    }
1682
1683    fn create_endpoint(
1684        &self,
1685        uri: &str,
1686        ctx: &dyn camel_component_api::ComponentContext,
1687    ) -> Result<Box<dyn Endpoint>, CamelError> {
1688        let mut config = ContainerConfig::from_uri(uri)?;
1689        // Apply global defaults if present and URI didn't set them
1690        if let Some(ref global) = self.config {
1691            config.apply_global_defaults(global);
1692        }
1693        let health_check = ContainerHealthCheck::new(&config);
1694        ctx.register_current_route_health_check(Arc::new(health_check));
1695        Ok(Box::new(ContainerEndpoint {
1696            uri: uri.to_string(),
1697            config,
1698        }))
1699    }
1700}
1701
1702/// Endpoint for container operations.
1703///
1704/// This endpoint creates producers for executing container operations
1705/// and consumers for receiving container events.
1706// TODO(CON-003): Forward container health status (inspect healthcheck / Health field) to
1707// Camel's health subsystem so the route can react to unhealthy containers.
1708pub struct ContainerEndpoint {
1709    uri: String,
1710    config: ContainerConfig,
1711}
1712
1713impl ContainerEndpoint {
1714    /// Returns the Docker host configured for this endpoint.
1715    /// Returns `None` if not set (for testing purposes).
1716    pub fn docker_host(&self) -> Option<&str> {
1717        self.config.host.as_deref()
1718    }
1719}
1720
1721impl Endpoint for ContainerEndpoint {
1722    fn uri(&self) -> &str {
1723        &self.uri
1724    }
1725
1726    fn create_consumer(
1727        &self,
1728        rt: Arc<dyn camel_component_api::RuntimeObservability>,
1729    ) -> Result<Box<dyn Consumer>, CamelError> {
1730        Ok(Box::new(ContainerConsumer::new(self.config.clone(), rt)))
1731    }
1732
1733    fn create_producer(
1734        &self,
1735        _rt: Arc<dyn camel_component_api::RuntimeObservability>,
1736        _ctx: &ProducerContext,
1737    ) -> Result<BoxProcessor, CamelError> {
1738        let docker = self.config.connect_docker_client()?;
1739        Ok(BoxProcessor::new(ContainerProducer {
1740            config: self.config.clone(),
1741            docker,
1742        }))
1743    }
1744}
1745
1746#[cfg(test)]
1747mod tests {
1748    use camel_component_api::test_support::PanicRuntimeObservability;
1749    fn test_rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
1750        std::sync::Arc::new(PanicRuntimeObservability)
1751    }
1752    fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
1753        std::sync::Arc::new(PanicRuntimeObservability)
1754    }
1755
1756    use super::*;
1757    use camel_api::MetricsCollector;
1758    use camel_component_api::HealthCheckRegistry;
1759    use camel_component_api::NoOpComponentContext;
1760
1761    #[test]
1762    fn test_container_config() {
1763        let config = ContainerConfig::from_uri("container:run?image=alpine").unwrap();
1764        assert_eq!(config.operation, "run");
1765        assert_eq!(config.image.as_deref(), Some("alpine"));
1766        // host is None by default; global config applies it later
1767        assert!(config.host.is_none());
1768    }
1769
1770    #[test]
1771    fn test_global_config_applied_to_endpoint() {
1772        // When global config is set and URI doesn't specify host,
1773        // apply_global_defaults should set host from global config.
1774        let global =
1775            ContainerGlobalConfig::default().with_docker_host("unix:///custom/docker.sock");
1776        let mut config = ContainerConfig::from_uri("container:run?image=alpine").unwrap();
1777        assert!(
1778            config.host.is_none(),
1779            "URI without ?host= should leave host as None"
1780        );
1781        config.apply_global_defaults(&global);
1782        assert_eq!(
1783            config.host.as_deref(),
1784            Some("unix:///custom/docker.sock"),
1785            "global docker_host must be applied when URI did not set host"
1786        );
1787    }
1788
1789    #[test]
1790    fn test_uri_param_wins_over_global_config() {
1791        // When URI explicitly sets host param, apply_global_defaults must NOT override it.
1792        let global =
1793            ContainerGlobalConfig::default().with_docker_host("unix:///custom/docker.sock");
1794        let mut config =
1795            ContainerConfig::from_uri("container:run?image=alpine&host=unix:///override.sock")
1796                .unwrap();
1797        assert_eq!(
1798            config.host.as_deref(),
1799            Some("unix:///override.sock"),
1800            "URI-set host should be parsed correctly"
1801        );
1802        config.apply_global_defaults(&global);
1803        assert_eq!(
1804            config.host.as_deref(),
1805            Some("unix:///override.sock"),
1806            "global config must NOT override a host already set by URI"
1807        );
1808    }
1809
1810    #[test]
1811    fn test_container_config_parses_name() {
1812        let config = ContainerConfig::from_uri("container:run?name=my-container").unwrap();
1813        assert_eq!(config.name.as_deref(), Some("my-container"));
1814    }
1815
1816    #[test]
1817    fn test_parse_producer_operation_known() {
1818        assert_eq!(
1819            parse_producer_operation("list").unwrap(),
1820            ProducerOperation::List
1821        );
1822        assert_eq!(
1823            parse_producer_operation("run").unwrap(),
1824            ProducerOperation::Run
1825        );
1826        assert_eq!(
1827            parse_producer_operation("start").unwrap(),
1828            ProducerOperation::Start
1829        );
1830        assert_eq!(
1831            parse_producer_operation("stop").unwrap(),
1832            ProducerOperation::Stop
1833        );
1834        assert_eq!(
1835            parse_producer_operation("remove").unwrap(),
1836            ProducerOperation::Remove
1837        );
1838    }
1839
1840    #[test]
1841    fn test_parse_producer_operation_unknown() {
1842        let err = parse_producer_operation("destruir_mundo").unwrap_err();
1843        match err {
1844            CamelError::ProcessorError(msg) => {
1845                assert!(
1846                    msg.contains("Unknown container operation"),
1847                    "Unexpected error message: {}",
1848                    msg
1849                );
1850            }
1851            _ => panic!("Expected ProcessorError for unknown operation"),
1852        }
1853    }
1854
1855    #[test]
1856    fn test_parse_producer_operation_new_variants() {
1857        assert_eq!(
1858            parse_producer_operation("exec").unwrap(),
1859            ProducerOperation::Exec
1860        );
1861        assert_eq!(
1862            parse_producer_operation("network-create").unwrap(),
1863            ProducerOperation::NetworkCreate
1864        );
1865        assert_eq!(
1866            parse_producer_operation("network-connect").unwrap(),
1867            ProducerOperation::NetworkConnect
1868        );
1869        assert_eq!(
1870            parse_producer_operation("network-disconnect").unwrap(),
1871            ProducerOperation::NetworkDisconnect
1872        );
1873        assert_eq!(
1874            parse_producer_operation("network-remove").unwrap(),
1875            ProducerOperation::NetworkRemove
1876        );
1877        assert_eq!(
1878            parse_producer_operation("network-list").unwrap(),
1879            ProducerOperation::NetworkList
1880        );
1881    }
1882
1883    #[test]
1884    fn test_resolve_container_name_header_overrides_config() {
1885        let config = ContainerConfig::from_uri("container:run?name=config-name").unwrap();
1886        let mut exchange = Exchange::new(Message::new(""));
1887        exchange.input.set_header(
1888            HEADER_CONTAINER_NAME,
1889            serde_json::Value::String("header-name".to_string()),
1890        );
1891
1892        let resolved = resolve_container_name(&exchange, &config);
1893        assert_eq!(resolved.as_deref(), Some("header-name"));
1894    }
1895
1896    #[test]
1897    fn test_container_config_rejects_tcp_host() {
1898        let config = ContainerConfig::from_uri("container:list?host=tcp://localhost:2375").unwrap();
1899        let err = config.connect_docker_client().unwrap_err();
1900        match err {
1901            CamelError::ProcessorError(msg) => {
1902                assert!(
1903                    msg.to_lowercase().contains("tcp"),
1904                    "Expected TCP scheme error, got: {}",
1905                    msg
1906                );
1907            }
1908            _ => panic!("Expected ProcessorError for unsupported tcp host"),
1909        }
1910    }
1911
1912    #[tokio::test]
1913    async fn test_run_container_with_cleanup_removes_on_start_failure() {
1914        let remove_called = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
1915        let remove_called_clone = remove_called.clone();
1916
1917        let result = run_container_with_cleanup(
1918            || async { Ok("container-123".to_string()) },
1919            |_id| async move {
1920                Err(CamelError::ProcessorError(
1921                    "Failed to start container".to_string(),
1922                ))
1923            },
1924            move |_id| {
1925                let remove_called_inner = remove_called_clone.clone();
1926                async move {
1927                    remove_called_inner.store(true, std::sync::atomic::Ordering::SeqCst);
1928                    Ok(())
1929                }
1930            },
1931        )
1932        .await;
1933
1934        assert!(result.is_err(), "Expected start failure to bubble up");
1935        assert!(
1936            remove_called.load(std::sync::atomic::Ordering::SeqCst),
1937            "Expected cleanup to remove container"
1938        );
1939    }
1940
1941    #[test]
1942    fn test_container_component_creates_endpoint() {
1943        let component = ContainerComponent::new();
1944        assert_eq!(component.scheme(), "container");
1945        let ctx = NoOpComponentContext;
1946        let endpoint = component
1947            .create_endpoint("container:run?image=alpine", &ctx)
1948            .unwrap();
1949        assert_eq!(endpoint.uri(), "container:run?image=alpine");
1950    }
1951
1952    #[test]
1953    fn test_container_config_parses_ports() {
1954        let config =
1955            ContainerConfig::from_uri("container:run?image=nginx&ports=8080:80,8443:443").unwrap();
1956        assert_eq!(config.ports.as_deref(), Some("8080:80,8443:443"));
1957    }
1958
1959    #[test]
1960    fn test_container_config_parses_env() {
1961        let config =
1962            ContainerConfig::from_uri("container:run?image=nginx&env=FOO=bar,BAZ=qux").unwrap();
1963        assert_eq!(config.env.as_deref(), Some("FOO=bar,BAZ=qux"));
1964    }
1965
1966    #[test]
1967    fn test_container_config_parses_logs_options() {
1968        let config = ContainerConfig::from_uri(
1969            "container:logs?containerId=my-app&follow=true&timestamps=true&tail=100",
1970        )
1971        .unwrap();
1972        assert_eq!(config.operation, "logs");
1973        assert_eq!(config.container_id.as_deref(), Some("my-app"));
1974        assert!(config.follow);
1975        assert!(config.timestamps);
1976        assert_eq!(config.tail.as_deref(), Some("100"));
1977    }
1978
1979    #[test]
1980    fn test_container_config_logs_defaults() {
1981        let config = ContainerConfig::from_uri("container:logs?containerId=test").unwrap();
1982        assert!(config.follow); // default: true
1983        assert!(!config.timestamps); // default: false
1984        assert!(config.tail.is_none()); // default: None (all)
1985    }
1986
1987    #[test]
1988    fn test_parse_ports_single() {
1989        let config = ContainerConfig::from_uri("container:run?image=nginx&ports=8080:80").unwrap();
1990        let (exposed, bindings) = config.parse_ports().unwrap();
1991
1992        assert!(exposed.contains(&"80/tcp".to_string()));
1993        assert!(bindings.contains_key("80/tcp"));
1994
1995        let binding = bindings.get("80/tcp").unwrap().as_ref().unwrap();
1996        assert_eq!(binding.len(), 1);
1997        assert_eq!(binding[0].host_port, Some("8080".to_string()));
1998    }
1999
2000    #[test]
2001    fn test_parse_ports_multiple() {
2002        let config =
2003            ContainerConfig::from_uri("container:run?image=nginx&ports=8080:80,8443:443").unwrap();
2004        let (exposed, bindings) = config.parse_ports().unwrap();
2005
2006        assert!(exposed.contains(&"80/tcp".to_string()));
2007        assert!(exposed.contains(&"443/tcp".to_string()));
2008        assert_eq!(bindings.len(), 2);
2009    }
2010
2011    #[test]
2012    fn test_parse_ports_with_protocol() {
2013        let config =
2014            ContainerConfig::from_uri("container:run?image=nginx&ports=8080:80/tcp,5353:53/udp")
2015                .unwrap();
2016        let (exposed, _bindings) = config.parse_ports().unwrap();
2017
2018        assert!(exposed.contains(&"80/tcp".to_string()));
2019        assert!(exposed.contains(&"53/udp".to_string()));
2020    }
2021
2022    #[test]
2023    fn test_parse_ports_none() {
2024        let config = ContainerConfig::from_uri("container:run?image=nginx").unwrap();
2025        let (exposed, bindings) = config.parse_ports().unwrap();
2026        assert!(exposed.is_empty());
2027        assert!(bindings.is_empty());
2028    }
2029
2030    #[test]
2031    fn test_parse_env_single() {
2032        let config = ContainerConfig::from_uri("container:run?image=nginx&env=FOO=bar").unwrap();
2033        let env = config.parse_env().unwrap();
2034
2035        assert_eq!(env.len(), 1);
2036        assert_eq!(env[0], "FOO=bar");
2037    }
2038
2039    #[test]
2040    fn test_parse_env_multiple() {
2041        let config =
2042            ContainerConfig::from_uri("container:run?image=nginx&env=FOO=bar,BAZ=qux,NUM=123")
2043                .unwrap();
2044        let env = config.parse_env().unwrap();
2045
2046        assert_eq!(env.len(), 3);
2047        assert!(env.contains(&"FOO=bar".to_string()));
2048        assert!(env.contains(&"BAZ=qux".to_string()));
2049        assert!(env.contains(&"NUM=123".to_string()));
2050    }
2051
2052    #[test]
2053    fn test_parse_env_none() {
2054        let config = ContainerConfig::from_uri("container:run?image=nginx").unwrap();
2055        assert!(config.parse_env().is_none());
2056    }
2057
2058    use camel_component_api::Message;
2059    use std::sync::Arc;
2060
2061    #[tokio::test]
2062    async fn test_container_producer_connection_error_on_invalid_host() {
2063        // Test that an invalid host (nonexistent socket) results in a connection error
2064        let component = ContainerComponent::new();
2065        let ctx = NoOpComponentContext;
2066        let endpoint = component
2067            .create_endpoint("container:list?host=unix:///nonexistent/docker.sock", &ctx)
2068            .unwrap();
2069
2070        let ctx = ProducerContext::new();
2071        let result = endpoint.create_producer(rt(), &ctx);
2072
2073        // The producer should return an error because it cannot connect to the invalid socket
2074        assert!(
2075            result.is_err(),
2076            "Expected error when connecting to invalid host"
2077        );
2078        let err = result.unwrap_err();
2079        match &err {
2080            CamelError::ProcessorError(msg) => {
2081                assert!(
2082                    msg.to_lowercase().contains("connection")
2083                        || msg.to_lowercase().contains("connect")
2084                        || msg.to_lowercase().contains("socket")
2085                        || msg.contains("docker"),
2086                    "Error message should indicate connection failure, got: {}",
2087                    msg
2088                );
2089            }
2090            _ => panic!("Expected ProcessorError, got: {:?}", err),
2091        }
2092    }
2093
2094    /// Test that consumer returns an error for unsupported operations.
2095    #[tokio::test]
2096    async fn test_container_consumer_unsupported_operation() {
2097        use tokio::sync::mpsc;
2098
2099        let component = ContainerComponent::new();
2100        let ctx = NoOpComponentContext;
2101        let endpoint = component.create_endpoint("container:run", &ctx).unwrap();
2102        let mut consumer = endpoint.create_consumer(rt()).unwrap();
2103
2104        // Create a minimal ConsumerContext
2105        let (tx, _rx) = mpsc::channel(16);
2106        let cancel_token = tokio_util::sync::CancellationToken::new();
2107        let context = ConsumerContext::new(tx, cancel_token, "container-test-route".to_string());
2108
2109        let result = consumer.start(context).await;
2110
2111        // Should return error because "run" is not a supported consumer operation
2112        assert!(
2113            result.is_err(),
2114            "Expected error for unsupported consumer operation"
2115        );
2116        let err = result.unwrap_err();
2117        match &err {
2118            CamelError::EndpointCreationFailed(msg) => {
2119                assert!(
2120                    msg.contains("Consumer only supports 'events' or 'logs'"),
2121                    "Error message should mention events or logs support, got: {}",
2122                    msg
2123                );
2124            }
2125            _ => panic!("Expected EndpointCreationFailed error, got: {:?}", err),
2126        }
2127    }
2128
2129    #[test]
2130    fn test_container_consumer_concurrency_model_is_concurrent() {
2131        let consumer = ContainerConsumer {
2132            config: ContainerConfig::from_uri("container:events").unwrap(),
2133            runtime: test_rt(),
2134        };
2135
2136        assert_eq!(
2137            consumer.concurrency_model(),
2138            camel_component_api::ConcurrencyModel::Concurrent { max: None }
2139        );
2140    }
2141
2142    #[test]
2143    fn test_container_config_parses_volumes() {
2144        let config = ContainerConfig::from_uri(
2145            "container:run?image=nginx&volumes=./html:/usr/share/nginx/html:ro",
2146        )
2147        .unwrap();
2148        assert_eq!(
2149            config.volumes.as_deref(),
2150            Some("./html:/usr/share/nginx/html:ro")
2151        );
2152    }
2153
2154    #[test]
2155    fn test_container_config_parses_exec_params() {
2156        let config = ContainerConfig::from_uri(
2157            "container:exec?containerId=my-app&cmd=ls /app&user=root&workdir=/tmp&detach=true",
2158        )
2159        .unwrap();
2160        assert_eq!(config.operation, "exec");
2161        assert_eq!(config.container_id.as_deref(), Some("my-app"));
2162        assert_eq!(config.cmd.as_deref(), Some("ls /app"));
2163        assert_eq!(config.user.as_deref(), Some("root"));
2164        assert_eq!(config.workdir.as_deref(), Some("/tmp"));
2165        assert!(config.detach);
2166    }
2167
2168    #[test]
2169    fn test_container_config_parses_network_create_params() {
2170        let config =
2171            ContainerConfig::from_uri("container:network-create?name=my-net&driver=bridge")
2172                .unwrap();
2173        assert_eq!(config.operation, "network-create");
2174        assert_eq!(config.name.as_deref(), Some("my-net"));
2175        assert_eq!(config.driver.as_deref(), Some("bridge"));
2176    }
2177
2178    #[test]
2179    fn test_container_config_defaults_new_fields() {
2180        let config = ContainerConfig::from_uri("container:list").unwrap();
2181        assert!(config.volumes.is_none());
2182        assert!(config.user.is_none());
2183        assert!(config.workdir.is_none());
2184        assert!(!config.detach);
2185        assert!(config.driver.is_none());
2186        assert!(!config.force);
2187    }
2188
2189    #[test]
2190    fn test_parse_volumes_bind_mount() {
2191        let config = ContainerConfig::from_uri(
2192            "container:run?image=nginx&volumes=./html:/usr/share/nginx/html:ro",
2193        )
2194        .unwrap();
2195        let (binds, anon) = config.parse_volumes().unwrap();
2196        assert_eq!(binds, vec!["./html:/usr/share/nginx/html:ro"]);
2197        assert!(anon.is_empty());
2198    }
2199
2200    #[test]
2201    fn test_parse_volumes_named_volume() {
2202        let config =
2203            ContainerConfig::from_uri("container:run?image=postgres&volumes=data:/var/lib/data")
2204                .unwrap();
2205        let (binds, anon) = config.parse_volumes().unwrap();
2206        assert_eq!(binds, vec!["data:/var/lib/data"]);
2207        assert!(anon.is_empty());
2208    }
2209
2210    #[test]
2211    fn test_parse_volumes_anonymous() {
2212        let config =
2213            ContainerConfig::from_uri("container:run?image=alpine&volumes=/tmp/app-data").unwrap();
2214        let (binds, anon) = config.parse_volumes().unwrap();
2215        assert!(binds.is_empty());
2216        assert!(anon.contains(&"/tmp/app-data".to_string()));
2217    }
2218
2219    #[test]
2220    fn test_parse_volumes_anonymous_with_mode() {
2221        let config =
2222            ContainerConfig::from_uri("container:run?image=alpine&volumes=/tmp/app-data:ro")
2223                .unwrap();
2224        let (binds, anon) = config.parse_volumes().unwrap();
2225        assert!(binds.is_empty());
2226        assert!(anon.contains(&"/tmp/app-data".to_string()));
2227    }
2228
2229    #[test]
2230    fn test_parse_volumes_multiple() {
2231        let config = ContainerConfig::from_uri(
2232            "container:run?image=nginx&volumes=./html:/usr/share/nginx/html:ro,data:/var/log/app",
2233        )
2234        .unwrap();
2235        let (binds, anon) = config.parse_volumes().unwrap();
2236        assert_eq!(binds.len(), 2);
2237        assert!(binds.contains(&"./html:/usr/share/nginx/html:ro".to_string()));
2238        assert!(binds.contains(&"data:/var/log/app".to_string()));
2239        assert!(anon.is_empty());
2240    }
2241
2242    #[test]
2243    fn test_parse_volumes_mixed() {
2244        let config = ContainerConfig::from_uri(
2245            "container:run?image=nginx&volumes=./html:/usr/share/nginx/html:ro,/tmp/cache",
2246        )
2247        .unwrap();
2248        let (binds, anon) = config.parse_volumes().unwrap();
2249        assert_eq!(binds.len(), 1);
2250        assert!(anon.contains(&"/tmp/cache".to_string()));
2251    }
2252
2253    #[test]
2254    fn test_parse_volumes_none() {
2255        let config = ContainerConfig::from_uri("container:run?image=nginx").unwrap();
2256        assert!(config.parse_volumes().is_none());
2257    }
2258
2259    #[test]
2260    fn test_parse_volumes_empty_entry_skipped() {
2261        let config = ContainerConfig::from_uri("container:run?image=nginx&volumes=,,").unwrap();
2262        assert!(config.parse_volumes().is_none());
2263    }
2264
2265    #[test]
2266    fn test_parse_volumes_rw_mode() {
2267        let config =
2268            ContainerConfig::from_uri("container:run?image=nginx&volumes=./data:/app/data:rw")
2269                .unwrap();
2270        let (binds, _) = config.parse_volumes().unwrap();
2271        assert_eq!(binds, vec!["./data:/app/data:rw"]);
2272    }
2273
2274    #[test]
2275    fn test_container_config_from_uri_parses_false_flags() {
2276        let config = ContainerConfig::from_uri(
2277            "container:logs?containerId=a&follow=false&timestamps=FALSE&autoPull=false&autoRemove=False&detach=TRUE&force=true",
2278        )
2279        .unwrap();
2280        assert!(!config.follow);
2281        assert!(!config.timestamps);
2282        assert!(!config.auto_pull);
2283        assert!(!config.auto_remove);
2284        assert!(config.detach);
2285        assert!(config.force);
2286    }
2287
2288    #[test]
2289    fn test_docker_socket_path_validation() {
2290        let unix_cfg =
2291            ContainerConfig::from_uri("container:list?host=unix:///tmp/docker.sock").unwrap();
2292        assert_eq!(
2293            unix_cfg.docker_socket_path().unwrap(),
2294            "unix:///tmp/docker.sock"
2295        );
2296
2297        let npipe_cfg =
2298            ContainerConfig::from_uri("container:list?host=npipe:////./pipe/docker_engine")
2299                .unwrap();
2300        assert_eq!(
2301            npipe_cfg.docker_socket_path().unwrap(),
2302            "npipe:////./pipe/docker_engine"
2303        );
2304
2305        let plain_cfg =
2306            ContainerConfig::from_uri("container:list?host=/var/run/docker.sock").unwrap();
2307        assert_eq!(
2308            plain_cfg.docker_socket_path().unwrap(),
2309            "/var/run/docker.sock"
2310        );
2311
2312        let bad_cfg =
2313            ContainerConfig::from_uri("container:list?host=http://localhost:2375").unwrap();
2314        assert!(bad_cfg.docker_socket_path().is_err());
2315    }
2316
2317    #[test]
2318    fn test_parse_ports_invalid_and_whitespace_entries() {
2319        // "8080" without colon is malformed — must fail-fast
2320        let cfg = ContainerConfig::from_uri("container:run?ports= , ,8080").unwrap();
2321        let err = cfg.parse_ports().unwrap_err();
2322        assert!(
2323            err.to_string().contains("malformed port mapping"),
2324            "expected malformed port error, got: {}",
2325            err
2326        );
2327
2328        let cfg =
2329            ContainerConfig::from_uri("container:run?ports= 8080:80 ,  5353:53/udp ").unwrap();
2330        let (exposed, bindings) = cfg.parse_ports().unwrap();
2331        assert!(exposed.contains(&"80/tcp".to_string()));
2332        assert!(exposed.contains(&"53/udp".to_string()));
2333        assert_eq!(bindings.len(), 2);
2334    }
2335
2336    #[test]
2337    fn test_parse_ports_fail_fast_on_malformed() {
2338        // First entry valid, second malformed — must fail on the malformed one
2339        let cfg = ContainerConfig::from_uri("container:run?ports=8080:80,badentry").unwrap();
2340        let err = cfg.parse_ports().unwrap_err();
2341        assert!(
2342            err.to_string()
2343                .contains("malformed port mapping 'badentry'"),
2344            "expected fail-fast on 'badentry', got: {}",
2345            err
2346        );
2347    }
2348
2349    #[test]
2350    fn test_parse_env_trims_and_filters_empty_items() {
2351        let cfg = ContainerConfig::from_uri("container:run?env= FOO=bar , ,BAZ=qux, ").unwrap();
2352        let env = cfg.parse_env().unwrap();
2353        assert_eq!(env, vec!["FOO=bar".to_string(), "BAZ=qux".to_string()]);
2354
2355        let cfg = ContainerConfig::from_uri("container:run?env= , , ").unwrap();
2356        assert!(cfg.parse_env().is_none());
2357    }
2358
2359    #[test]
2360    fn test_parse_volume_str_rejects_invalid_mode_and_accepts_mixed() {
2361        assert!(parse_volume_str("/host:/ctr:badmode").is_none());
2362        let (binds, anon) = parse_volume_str("a:/b:rw,/tmp/cache,/tmp/logs:ro").unwrap();
2363        assert!(binds.contains(&"a:/b:rw".to_string()));
2364        assert!(anon.contains(&"/tmp/cache".to_string()));
2365        assert!(anon.contains(&"/tmp/logs".to_string()));
2366    }
2367
2368    #[test]
2369    fn test_format_docker_event_variants_and_timestamp_extraction() {
2370        let mut attrs = std::collections::HashMap::new();
2371        attrs.insert("name".to_string(), "demo".to_string());
2372        attrs.insert("image".to_string(), "alpine:latest".to_string());
2373        attrs.insert("exitCode".to_string(), "137".to_string());
2374        let actor = bollard::models::EventActor {
2375            id: None,
2376            attributes: Some(attrs),
2377        };
2378
2379        let create_event = bollard::models::EventMessage {
2380            action: Some("create".to_string()),
2381            actor: Some(actor.clone()),
2382            ..Default::default()
2383        };
2384        assert_eq!(
2385            format_docker_event(&create_event),
2386            "[CREATE] Container demo (alpine:latest)"
2387        );
2388
2389        let die_event = bollard::models::EventMessage {
2390            action: Some("die".to_string()),
2391            actor: Some(actor),
2392            ..Default::default()
2393        };
2394        assert_eq!(
2395            format_docker_event(&die_event),
2396            "[DIE]    Container demo (exit: 137)"
2397        );
2398
2399        let other_event = bollard::models::EventMessage {
2400            action: Some("oom".to_string()),
2401            actor: None,
2402            ..Default::default()
2403        };
2404        assert_eq!(format_docker_event(&other_event), "[OOM] Container unknown");
2405
2406        assert_eq!(
2407            extract_timestamp("2024-01-01T00:00:00Z hello"),
2408            Some("2024-01-01T00:00:00Z".to_string())
2409        );
2410        assert_eq!(extract_timestamp("hello world"), None);
2411    }
2412
2413    #[tokio::test]
2414    async fn test_run_container_with_cleanup_error_paths() {
2415        let create_fail = run_container_with_cleanup(
2416            || async { Err(CamelError::ProcessorError("create-fail".to_string())) },
2417            |_id| async move { Ok(()) },
2418            |_id| async move { Ok(()) },
2419        )
2420        .await;
2421        assert!(
2422            matches!(create_fail, Err(CamelError::ProcessorError(msg)) if msg == "create-fail")
2423        );
2424
2425        let cleanup_fail = run_container_with_cleanup(
2426            || async { Ok("cid-1".to_string()) },
2427            |_id| async move { Err(CamelError::ProcessorError("start-fail".to_string())) },
2428            |_id| async move { Err(CamelError::ProcessorError("remove-fail".to_string())) },
2429        )
2430        .await;
2431        match cleanup_fail {
2432            Err(CamelError::ProcessorError(msg)) => {
2433                assert!(msg.contains("Failed to start container"));
2434                assert!(msg.contains("Cleanup failed"));
2435            }
2436            other => panic!("unexpected result: {:?}", other),
2437        }
2438    }
2439
2440    #[tokio::test]
2441    async fn test_logs_consumer_requires_container_id() {
2442        use tokio::sync::mpsc;
2443
2444        let mut consumer = ContainerConsumer {
2445            config: ContainerConfig::from_uri("container:logs").unwrap(),
2446            runtime: test_rt(),
2447        };
2448        let (tx, _rx) = mpsc::channel(4);
2449        let context = ConsumerContext::new(
2450            tx,
2451            tokio_util::sync::CancellationToken::new(),
2452            "container-test-route".to_string(),
2453        );
2454
2455        let err = consumer.start(context).await.unwrap_err();
2456        match err {
2457            CamelError::EndpointCreationFailed(msg) => {
2458                assert!(msg.contains("containerId is required for logs consumer"));
2459            }
2460            other => panic!("unexpected error: {:?}", other),
2461        }
2462    }
2463
2464    #[tokio::test]
2465    async fn test_events_consumer_stops_immediately_when_cancelled() {
2466        use tokio::sync::mpsc;
2467
2468        let mut consumer = ContainerConsumer {
2469            config: ContainerConfig::from_uri("container:events").unwrap(),
2470            runtime: test_rt(),
2471        };
2472
2473        let (tx, _rx) = mpsc::channel(4);
2474        let cancel = tokio_util::sync::CancellationToken::new();
2475        cancel.cancel();
2476        let context = ConsumerContext::new(tx, cancel, "container-test-route".to_string());
2477
2478        let result = consumer.start(context).await;
2479        assert!(result.is_ok());
2480    }
2481
2482    #[test]
2483    fn test_global_config_constructors_and_endpoint_docker_host() {
2484        let global = ContainerGlobalConfig::new().with_docker_host("unix:///tmp/docker.sock");
2485        let mut cfg = ContainerConfig::from_uri("container:list").unwrap();
2486        cfg.apply_global_defaults(&global);
2487
2488        let endpoint = ContainerEndpoint {
2489            uri: "container:list".to_string(),
2490            config: cfg,
2491        };
2492
2493        assert_eq!(endpoint.docker_host(), Some("unix:///tmp/docker.sock"));
2494
2495        let component = ContainerComponent::with_config(global);
2496        assert_eq!(component.scheme(), "container");
2497    }
2498
2499    #[test]
2500    fn container_global_config_has_reconnect_policy() {
2501        let cfg = ContainerGlobalConfig::default();
2502        assert_eq!(cfg.reconnect.max_attempts, 0); // unlimited
2503        assert!(cfg.reconnect.enabled);
2504    }
2505
2506    /// Regression: max_attempts=N → exactly N invocations (caught OpenSearch off-by-one 1f5c4c2a).
2507    /// Replicates the exact retry loop from ContainerConsumer::{start_events_consumer,start_logs_consumer}
2508    /// (lib.rs:~1408-1431, ~1503-1525):
2509    ///   attempt starts at 0, incremented on error, !should_retry(attempt), delay_for(attempt-1)
2510    #[tokio::test]
2511    async fn retry_loop_invokes_operation_exactly_max_attempts_times() {
2512        use std::sync::Arc;
2513        use std::sync::atomic::{AtomicU32, Ordering};
2514        use std::time::Duration;
2515
2516        let policy = NetworkRetryPolicy {
2517            max_attempts: 3,
2518            initial_delay: Duration::from_millis(1),
2519            max_delay: Duration::from_millis(1),
2520            multiplier: 1.0,
2521            ..NetworkRetryPolicy::default()
2522        };
2523
2524        let calls = Arc::new(AtomicU32::new(0));
2525        let calls_clone = Arc::clone(&calls);
2526        let mut attempt: u32 = 0;
2527
2528        loop {
2529            calls_clone.fetch_add(1, Ordering::SeqCst);
2530            let result: Result<(), ()> = Err(());
2531            match result {
2532                Ok(_) => {
2533                    break;
2534                }
2535                Err(_) => {
2536                    attempt += 1;
2537                    if !policy.should_retry(attempt) {
2538                        break;
2539                    }
2540                    let delay = policy.delay_for(attempt - 1);
2541                    tokio::time::sleep(delay).await;
2542                    continue;
2543                }
2544            }
2545        }
2546
2547        assert_eq!(
2548            calls.load(Ordering::SeqCst),
2549            3,
2550            "max_attempts=3 must yield exactly 3 invocations"
2551        );
2552    }
2553
2554    // -----------------------------------------------------------------------
2555    // ADR-0012 (e) metric wiring regression tests
2556    // -----------------------------------------------------------------------
2557
2558    /// Regression: events-connect error path calls increment_errors with
2559    /// correct route_id and label. Uses an unsupported tcp:// host to trigger
2560    /// the error path WITHOUT needing a real Docker daemon (docker_socket_path
2561    /// returns Err on non-unix/npipe schemes).
2562    #[tokio::test]
2563    async fn events_connect_error_increments_metrics() {
2564        use std::sync::Mutex;
2565        use std::time::Duration;
2566
2567        struct RecordingMetrics(Mutex<Vec<(String, String)>>);
2568
2569        impl MetricsCollector for RecordingMetrics {
2570            fn record_exchange_duration(&self, _: &str, _: Duration) {}
2571            fn increment_errors(&self, route_id: &str, error_type: &str) {
2572                self.0
2573                    .lock()
2574                    .unwrap()
2575                    .push((route_id.to_string(), error_type.to_string()));
2576            }
2577            fn increment_exchanges(&self, _: &str) {}
2578            fn set_queue_depth(&self, _: &str, _: usize) {}
2579            fn record_circuit_breaker_change(&self, _: &str, _: &str, _: &str) {}
2580        }
2581
2582        struct RecordingRuntime {
2583            metrics: Arc<RecordingMetrics>,
2584        }
2585
2586        impl HealthCheckRegistry for RecordingRuntime {
2587            fn force_unhealthy_for_route(&self, _: &str, _: &str, _: &str) {}
2588        }
2589
2590        impl RuntimeObservability for RecordingRuntime {
2591            fn metrics(&self) -> Arc<dyn MetricsCollector> {
2592                self.metrics.clone()
2593            }
2594            fn health(&self) -> Arc<dyn HealthCheckRegistry> {
2595                Arc::new(camel_component_api::NoOpHealthCheckRegistry)
2596            }
2597        }
2598
2599        let recording = Arc::new(RecordingMetrics(Mutex::new(Vec::new())));
2600        let rt: Arc<dyn RuntimeObservability> = Arc::new(RecordingRuntime {
2601            metrics: recording.clone(),
2602        });
2603
2604        // Use unsupported tcp:// host so docker_socket_path() fails fast
2605        // without any real I/O. Disable retry so the first failure propagates
2606        // immediately to the Err(e) arm.
2607        let mut config = ContainerConfig::from_uri("container:events").unwrap();
2608        config.host = Some("tcp://192.0.2.1:2375".to_string());
2609        config.reconnect = NetworkRetryPolicy::disabled();
2610
2611        let mut consumer = ContainerConsumer {
2612            config,
2613            runtime: rt,
2614        };
2615
2616        let (tx, _rx) = tokio::sync::mpsc::channel(4);
2617        let cancel = tokio_util::sync::CancellationToken::new();
2618        let context = ConsumerContext::new(tx, cancel, "events-test-route".to_string());
2619
2620        let result = consumer.start(context).await;
2621        assert!(result.is_err(), "expected Docker connection error");
2622
2623        let errors = recording.0.lock().unwrap();
2624        assert_eq!(
2625            errors.len(),
2626            1,
2627            "expected exactly one increment_errors call"
2628        );
2629        assert_eq!(errors[0].0, "events-test-route");
2630        assert_eq!(errors[0].1, "e:container:events-connect");
2631    }
2632}