1pub 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
38static CONTAINER_TRACKER: once_cell::sync::Lazy<Arc<Mutex<HashSet<String>>>> =
41 once_cell::sync::Lazy::new(|| Arc::new(Mutex::new(HashSet::new())));
42
43fn track_container(id: String) {
45 if let Ok(mut tracker) = CONTAINER_TRACKER.lock() {
46 tracker.insert(id);
47 }
48}
49
50fn untrack_container(id: &str) {
52 if let Ok(mut tracker) = CONTAINER_TRACKER.lock() {
53 tracker.remove(id);
54 }
55}
56
57pub 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 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
103const DOCKER_CONNECT_TIMEOUT_SECS: u64 = 120;
107
108pub const HEADER_ACTION: &str = "CamelContainerAction";
110
111pub const HEADER_IMAGE: &str = "CamelContainerImage";
113
114pub const HEADER_CONTAINER_ID: &str = "CamelContainerId";
116
117pub const HEADER_LOG_STREAM: &str = "CamelContainerLogStream";
119
120pub const HEADER_LOG_TIMESTAMP: &str = "CamelContainerLogTimestamp";
122
123pub const HEADER_CONTAINER_NAME: &str = "CamelContainerName";
125
126pub const HEADER_ACTION_RESULT: &str = "CamelContainerActionResult";
128
129pub const HEADER_CMD: &str = "CamelContainerCmd";
131
132pub const HEADER_NETWORK: &str = "CamelContainerNetwork";
134
135pub const HEADER_EXIT_CODE: &str = "CamelContainerExitCode";
137
138pub const HEADER_VOLUMES: &str = "CamelContainerVolumes";
140
141pub const HEADER_EXEC_ID: &str = "CamelContainerExecId";
143
144fn container_reconnect_default() -> NetworkRetryPolicy {
152 NetworkRetryPolicy {
153 enabled: true,
154 max_attempts: 0, initial_delay: std::time::Duration::from_secs(5),
156 multiplier: 1.0, max_delay: std::time::Duration::from_secs(5),
158 jitter_factor: 0.0,
159 max_attempts_absolute: None,
160 }
161}
162
163#[derive(Debug, Clone, PartialEq, serde::Deserialize)]
167#[serde(default)]
168pub struct ContainerGlobalConfig {
169 pub docker_host: String,
171 #[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#[derive(Debug, Clone)]
205pub struct ContainerConfig {
206 pub operation: String,
208 pub image: Option<String>,
210 pub name: Option<String>,
212 pub host: Option<String>,
214 pub cmd: Option<String>,
216 pub ports: Option<String>,
218 pub env: Option<String>,
220 pub network: Option<String>,
222 pub container_id: Option<String>,
224 pub follow: bool,
226 pub timestamps: bool,
228 pub tail: Option<String>,
230 pub auto_pull: bool,
232 pub auto_remove: bool,
234 pub volumes: Option<String>,
236 pub user: Option<String>,
238 pub workdir: Option<String>,
240 pub detach: bool,
242 pub driver: Option<String>,
244 pub force: bool,
246 pub reconnect: NetworkRetryPolicy,
248}
249
250impl ContainerConfig {
251 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 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 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 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#[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 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#[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
1368pub struct ContainerConsumer {
1373 config: ContainerConfig,
1374 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 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 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 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 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 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 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 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 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
1642pub struct ContainerComponent {
1650 config: Option<ContainerGlobalConfig>,
1651}
1652
1653impl ContainerComponent {
1654 pub fn new() -> Self {
1656 Self { config: None }
1657 }
1658
1659 pub fn with_config(config: ContainerGlobalConfig) -> Self {
1661 Self {
1662 config: Some(config),
1663 }
1664 }
1665
1666 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 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
1702pub struct ContainerEndpoint {
1709 uri: String,
1710 config: ContainerConfig,
1711}
1712
1713impl ContainerEndpoint {
1714 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 assert!(config.host.is_none());
1768 }
1769
1770 #[test]
1771 fn test_global_config_applied_to_endpoint() {
1772 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 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×tamps=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); assert!(!config.timestamps); assert!(config.tail.is_none()); }
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 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 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 #[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 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 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×tamps=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 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 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); assert!(cfg.reconnect.enabled);
2504 }
2505
2506 #[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 #[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 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}