1use std::{
2 collections::{hash_map::DefaultHasher, HashMap},
3 fmt,
4 future::Future,
5 hash::{Hash as _, Hasher as _},
6 path::{Path, PathBuf},
7 time::{Duration, Instant},
8};
9
10use bollard::{
11 container::{LogOutput, NetworkingConfig as ContainerNetworkingConfig},
12 errors::Error,
13 exec::{CreateExecOptions, StartExecResults},
14 models::{
15 ContainerCreateBody, ContainerStateStatusEnum, EndpointSettings, HealthConfig, HealthStatusEnum, HostConfig,
16 HostConfigCgroupnsModeEnum, Ipam, NetworkConnectRequest, NetworkCreateRequest, VolumeCreateRequest,
17 },
18 query_parameters::{CreateContainerOptionsBuilder, CreateImageOptions, ListContainersOptionsBuilder, LogsOptions},
19 Docker,
20};
21use futures::{StreamExt as _, TryStreamExt as _};
22use saluki_error::{generic_error, ErrorContext as _, GenericError};
23use tokio::{
24 io::{AsyncWriteExt as _, BufWriter},
25 time::sleep,
26};
27use tracing::{debug, error, trace, warn};
28
29use crate::config::{DatadogIntakeConfig, MillstoneConfig, TargetConfig};
30
31const MILLSTONE_CONFIG_PATH_INTERNAL: &str = "/etc/millstone/config.toml";
32
33pub const DEFAULT_ALPINE_IMAGE: &str = "alpine:latest";
35const DATADOG_INTAKE_HEALTHCHECK_INTERVAL: Duration = Duration::from_secs(1);
36const DATADOG_INTAKE_HEALTHCHECK_TIMEOUT: Duration = Duration::from_secs(1);
37const DATADOG_INTAKE_HEALTHCHECK_RETRIES: i64 = 30;
38const DATADOG_INTAKE_HEALTHCHECK_START_PERIOD: Duration = Duration::from_secs(1);
39const DATADOG_INTAKE_HEALTHCHECK_START_INTERVAL: Duration = Duration::from_secs(1);
40const DATADOG_INTAKE_HEALTHCHECK_COMMAND: &str = concat!(
41 "exec 3<>/dev/tcp/127.0.0.1/2049 && ",
42 "printf 'GET /ready HTTP/1.1\\r\\nHost: localhost\\r\\nConnection: close\\r\\n\\r\\n' >&3 && ",
43 "grep -q '200 OK' <&3"
44);
45
46const TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS: u32 = 3;
51const TRANSIENT_ERROR_RETRY_BASE_BACKOFF: Duration = Duration::from_millis(250);
52const TRANSIENT_ERROR_RETRY_MAX_JITTER: Duration = Duration::from_millis(500);
53
54pub enum ExitStatus {
55 Success,
56 Failed { code: i64, error: String },
57}
58
59impl fmt::Display for ExitStatus {
60 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
61 match self {
62 ExitStatus::Success => write!(f, "success (0)"),
63 ExitStatus::Failed { code, error } => write!(f, "failed (exit code: {}, error: {})", code, error),
64 }
65 }
66}
67
68#[derive(Clone, Copy, Debug, Eq, PartialEq)]
70pub enum ContainerOs {
71 Linux,
73 Windows,
75}
76
77#[derive(Clone)]
81pub struct DriverConfig {
82 driver_id: &'static str,
83 image: String,
84 entrypoint: Option<Vec<String>>,
85 command: Option<Vec<String>>,
86 env: Vec<String>,
87 binds: Vec<String>,
88 healthcheck: Option<HealthConfig>,
89 exposed_ports: Vec<(&'static str, u16)>,
90 container_os: ContainerOs,
91 host_cgroup_namespace: bool,
92 additional_volume_mounts: Vec<String>,
99
100 network_aliases: Vec<String>,
107
108 additional_networks: Vec<String>,
115
116 alpine_image: String,
120}
121
122impl DriverConfig {
123 pub async fn millstone(config: MillstoneConfig) -> Result<Self, GenericError> {
124 match tokio::fs::metadata(&config.config_path).await {
126 Ok(metadata) if metadata.is_file() => {}
127 Ok(_) => {
128 return Err(generic_error!(
129 "Specified millstone configuration path ({}) does not point to a file.",
130 config.config_path.display()
131 ))
132 }
133 Err(e) => {
134 return Err(generic_error!(
135 "Failed to ensure specified millstone configuration ({}) exists locally: {}",
136 config.config_path.display(),
137 e
138 ))
139 }
140 }
141
142 let millstone_binary_path = config
143 .binary_path
144 .unwrap_or_else(|| "/usr/local/bin/millstone".to_string());
145 let entrypoint = vec![millstone_binary_path, MILLSTONE_CONFIG_PATH_INTERNAL.to_string()];
146
147 let driver_config = Self::from_image("millstone", config.image)
148 .with_entrypoint(entrypoint)
149 .with_bind_mount(config.config_path, MILLSTONE_CONFIG_PATH_INTERNAL);
150
151 Ok(driver_config)
152 }
153
154 pub async fn datadog_intake(config: DatadogIntakeConfig) -> Result<Self, GenericError> {
155 let datadog_intake_binary_path = config
156 .binary_path
157 .unwrap_or_else(|| "/usr/local/bin/datadog-intake".to_string());
158 let entrypoint = vec![datadog_intake_binary_path];
159
160 let driver_config = DriverConfig::from_image("datadog-intake", config.image)
161 .with_entrypoint(entrypoint)
162 .with_healthcheck(
163 vec![
164 "/bin/bash".to_string(),
165 "-c".to_string(),
166 DATADOG_INTAKE_HEALTHCHECK_COMMAND.to_string(),
167 ],
168 DATADOG_INTAKE_HEALTHCHECK_INTERVAL,
169 DATADOG_INTAKE_HEALTHCHECK_TIMEOUT,
170 DATADOG_INTAKE_HEALTHCHECK_RETRIES,
171 DATADOG_INTAKE_HEALTHCHECK_START_PERIOD,
172 DATADOG_INTAKE_HEALTHCHECK_START_INTERVAL,
173 )
174 .with_exposed_port("tcp", 2049);
177
178 Ok(driver_config)
179 }
180
181 pub async fn target(target_id: &'static str, config: TargetConfig) -> Result<Self, GenericError> {
182 let driver_config = DriverConfig::from_image(target_id, config.image)
183 .with_entrypoint(config.entrypoint)
184 .with_command(config.command)
185 .with_env_vars(config.additional_env_vars)
186 .with_container_os(config.container_os)
187 .with_host_cgroup_namespace(config.host_cgroup_namespace);
188
189 Ok(driver_config)
190 }
191
192 pub fn from_image(driver_id: &'static str, image: String) -> Self {
194 Self {
195 driver_id,
196 image,
197 entrypoint: None,
198 command: None,
199 env: vec![],
200 binds: vec![],
201 healthcheck: None,
202 exposed_ports: vec![],
203 container_os: ContainerOs::Linux,
204 host_cgroup_namespace: false,
205 additional_volume_mounts: vec![],
206 network_aliases: vec![],
207 additional_networks: vec![],
208 alpine_image: DEFAULT_ALPINE_IMAGE.to_string(),
209 }
210 }
211
212 pub fn with_alpine_image(mut self, alpine_image: impl Into<String>) -> Self {
214 self.alpine_image = alpine_image.into();
215 self
216 }
217
218 pub fn with_entrypoint(mut self, entrypoint: Vec<String>) -> Self {
222 if !entrypoint.is_empty() {
223 self.entrypoint = Some(entrypoint);
224 }
225 self
226 }
227
228 pub fn with_command(mut self, command: Vec<String>) -> Self {
232 if !command.is_empty() {
233 self.command = Some(command);
234 }
235 self
236 }
237
238 pub fn with_env_var<K, V>(mut self, key: K, value: V) -> Self
240 where
241 K: AsRef<str>,
242 V: AsRef<str>,
243 {
244 self.env.push(format!("{}={}", key.as_ref(), value.as_ref()));
245 self
246 }
247
248 pub fn with_env_vars(mut self, env: Vec<String>) -> Self {
250 self.env.extend(env);
251 self
252 }
253
254 pub fn with_bind_mount<HP, CP>(mut self, host_path: HP, container_path: CP) -> Self
259 where
260 HP: AsRef<Path>,
261 CP: AsRef<Path>,
262 {
263 let bind_mount = format!("{}:{}", host_path.as_ref().display(), container_path.as_ref().display());
264 self.binds.push(bind_mount);
265 self
266 }
267
268 pub fn with_readonly_bind_mount<HP, CP>(mut self, host_path: HP, container_path: CP) -> Self
272 where
273 HP: AsRef<Path>,
274 CP: AsRef<Path>,
275 {
276 let bind_mount = format!(
277 "{}:{}:ro",
278 host_path.as_ref().display(),
279 container_path.as_ref().display()
280 );
281 self.binds.push(bind_mount);
282 self
283 }
284
285 pub fn with_healthcheck(
287 mut self, mut test_command: Vec<String>, interval: Duration, timeout: Duration, retries: i64,
288 start_period: Duration, start_interval: Duration,
289 ) -> Self {
290 test_command.insert(0, "CMD".to_string());
293
294 self.healthcheck = Some(HealthConfig {
295 test: Some(test_command),
296 interval: Some(interval.as_nanos() as i64),
297 timeout: Some(timeout.as_nanos() as i64),
298 retries: Some(retries),
299 start_period: Some(start_period.as_nanos() as i64),
300 start_interval: Some(start_interval.as_nanos() as i64),
301 });
302 self
303 }
304
305 pub fn with_network_alias(mut self, alias: impl Into<String>) -> Self {
310 self.network_aliases.push(alias.into());
311 self
312 }
313
314 pub fn with_network(mut self, network: impl Into<String>) -> Self {
320 self.additional_networks.push(network.into());
321 self
322 }
323
324 pub fn with_volume_mount(mut self, volume_name: impl Into<String>, container_path: impl AsRef<Path>) -> Self {
332 self.additional_volume_mounts
333 .push(format!("{}:{}", volume_name.into(), container_path.as_ref().display()));
334 self
335 }
336
337 pub fn with_exposed_port(mut self, protocol: &'static str, internal_port: u16) -> Self {
343 self.exposed_ports.push((protocol, internal_port));
344 self
345 }
346
347 pub fn with_container_os(mut self, container_os: ContainerOs) -> Self {
354 self.container_os = container_os;
355 self
356 }
357
358 pub fn with_host_cgroup_namespace(mut self, host_cgroup_namespace: bool) -> Self {
360 self.host_cgroup_namespace = host_cgroup_namespace;
361 self
362 }
363
364 fn needs_shared_volume_permission_fixup(&self) -> bool {
373 self.container_os == ContainerOs::Linux
374 }
375
376 fn network_driver(&self) -> &'static str {
381 match self.container_os {
382 ContainerOs::Linux => "bridge",
383 ContainerOs::Windows => "nat",
384 }
385 }
386
387 fn port_publishing_options(&self) -> (Option<bool>, Option<Vec<String>>) {
388 if self.exposed_ports.is_empty() {
389 return (None, None);
390 }
391
392 let exposed_ports = self
393 .exposed_ports
394 .iter()
395 .map(|(protocol, internal_port)| format!("{}/{}", internal_port, protocol))
396 .collect();
397 let publish_all_ports = match self.container_os {
398 ContainerOs::Linux => Some(true),
399 ContainerOs::Windows => None,
400 };
401
402 (publish_all_ports, Some(exposed_ports))
403 }
404
405 fn container_binds_from(&self, isolation_group_name: &str, mut binds: Vec<String>) -> Vec<String> {
414 match self.container_os {
415 ContainerOs::Linux => {
416 binds.push(format!("{}:/airlock:z", isolation_group_name));
417 binds.push("/proc:/host/proc:ro".to_string());
418 binds.push("/sys/fs/cgroup:/host/sys/fs/cgroup:ro".to_string());
419 binds.push("/var/run/docker.sock:/var/run/docker.sock:ro".to_string());
420 }
421 ContainerOs::Windows => {
422 binds.push(format!("{}:C:\\airlock", isolation_group_name));
423 }
424 }
425
426 binds.extend(self.additional_volume_mounts.clone());
427 binds
428 }
429}
430
431#[derive(Debug, Default)]
433pub struct DriverDetails {
434 container_name: String,
435 container_ip: Option<String>,
436 port_mappings: Option<HashMap<String, u16>>,
437}
438
439fn insert_port_mapping_if_parseable(
445 port_mappings: &mut HashMap<String, u16>, internal_port: impl Into<String>, host_port: Option<&str>,
446) {
447 if let Some(host_port) = host_port.and_then(|value| value.parse::<u16>().ok()) {
448 port_mappings.insert(internal_port.into(), host_port);
449 }
450}
451
452impl DriverDetails {
453 pub fn container_name(&self) -> &str {
455 &self.container_name
456 }
457
458 pub fn container_ip(&self) -> Option<&str> {
460 self.container_ip.as_deref()
461 }
462
463 pub fn try_get_exposed_port(&self, protocol: &str, internal_port: u16) -> Option<u16> {
469 self.port_mappings
470 .as_ref()
471 .and_then(|port_mappings| port_mappings.get(&format!("{}/{}", internal_port, protocol)).copied())
472 }
473}
474
475pub struct Driver {
477 isolation_group_id: String,
478 isolation_group_name: String,
479 container_name: String,
480 config: DriverConfig,
481 docker: Docker,
482 log_dir: Option<PathBuf>,
483}
484
485impl Driver {
486 pub fn from_config(isolation_group_id: String, config: DriverConfig) -> Result<Self, GenericError> {
506 let docker = crate::docker::connect()?;
507
508 Ok(Self {
509 isolation_group_name: format!("airlock-{}", isolation_group_id),
510 container_name: format!("airlock-{}-{}", isolation_group_id, config.driver_id),
511 isolation_group_id,
512 config,
513 docker,
514 log_dir: None,
515 })
516 }
517
518 pub fn with_logging(mut self, log_dir: PathBuf) -> Self {
524 self.log_dir = Some(log_dir);
525 self
526 }
527
528 pub fn driver_id(&self) -> &'static str {
532 self.config.driver_id
533 }
534
535 pub async fn clean_related_resources(isolation_group_id: String) -> Result<(), GenericError> {
544 let docker = crate::docker::connect()?;
545
546 let isolation_group_name = format!("airlock-{}", isolation_group_id);
547 let isolation_group_label = format!("airlock-isolation-group={}", isolation_group_id);
548
549 let list_filters: HashMap<&str, Vec<&str>> =
551 [("label", vec!["created_by=airlock", isolation_group_label.as_str()])]
552 .into_iter()
553 .collect();
554 let list_options = Some(
555 ListContainersOptionsBuilder::default()
556 .all(true)
557 .filters(&list_filters)
558 .build(),
559 );
560 let containers = docker.list_containers(list_options).await.with_error_context(|| {
561 format!(
562 "Failed to list containers attached to isolation group '{}'.",
563 isolation_group_id
564 )
565 })?;
566
567 for container in containers {
568 let container_name = match container.id {
569 Some(id) => id,
570 None => {
571 debug!("Listed container had no ID. Skipping removal.");
572 continue;
573 }
574 };
575
576 if let Err(e) = docker.stop_container(container_name.as_str(), None).await {
577 error!(error = %e, "Failed to stop container '{}'.", container_name);
578 continue;
579 } else {
580 debug!("Stopped container '{}'.", container_name);
581 }
582
583 if let Err(e) = docker.remove_container(container_name.as_str(), None).await {
584 error!(error = %e, "Failed to remove container '{}'.", container_name);
585 continue;
586 } else {
587 debug!("Removed container '{}'.", container_name);
588 }
589 }
590
591 if let Err(e) = docker
593 .remove_volume(
594 isolation_group_name.as_str(),
595 None::<bollard::query_parameters::RemoveVolumeOptions>,
596 )
597 .await
598 {
599 error!(error = %e, "Failed to remove shared volume '{}'.", isolation_group_name);
600 } else {
601 debug!("Removed shared volume '{}'.", isolation_group_name);
602 }
603
604 if let Err(e) = docker.remove_network(isolation_group_name.as_str()).await {
606 error!(error = %e, "Failed to remove shared network '{}'.", isolation_group_name);
607 } else {
608 debug!("Removed shared network '{}'.", isolation_group_name);
609 }
610
611 Ok(())
612 }
613
614 async fn create_network_if_missing(&self) -> Result<(), GenericError> {
615 with_transient_error_retry("network creation", &self.isolation_group_id, move || {
616 self.create_network_if_missing_inner()
617 })
618 .await
619 }
620
621 async fn create_network_if_missing_inner(&self) -> Result<(), GenericError> {
622 let networks = self.docker.list_networks(None).await?;
624 if networks
625 .iter()
626 .any(|network| network.name.as_deref() == Some(self.isolation_group_name.as_str()))
627 {
628 debug!("Network '{}' already exists.", self.isolation_group_name);
629 return Ok(());
630 }
631
632 debug!(
633 driver_id = self.config.driver_id,
634 isolation_group = self.isolation_group_id,
635 "Network '{}' does not exist. Creating...",
636 self.isolation_group_name
637 );
638
639 let network_options = NetworkCreateRequest {
641 name: self.isolation_group_name.clone(),
642 driver: Some(self.config.network_driver().to_string()),
643 ipam: Some(Ipam::default()),
644 enable_ipv6: Some(false),
645 labels: Some(get_default_airlock_labels(self.isolation_group_id.as_str())),
646 ..Default::default()
647 };
648 let response = self.docker.create_network(network_options).await?;
649 debug!(
650 driver_id = self.config.driver_id,
651 isolation_group = self.isolation_group_id,
652 "Created network '{}' (ID: {:?}).",
653 self.isolation_group_name,
654 response.id
655 );
656
657 Ok(())
658 }
659
660 async fn create_image_if_missing_inner(&self, image: &str) -> Result<(), GenericError> {
661 let image_options = CreateImageOptions {
662 from_image: Some(image.to_string()),
663 ..Default::default()
664 };
665
666 let mut create_stream = self.docker.create_image(Some(image_options), None, None);
667 while let Some(info) = create_stream.next().await {
668 trace!(
669 driver_id = self.config.driver_id,
670 isolation_group = self.isolation_group_id,
671 image,
672 "Received image pull update: {:?}",
673 info
674 );
675 }
676
677 Ok(())
678 }
679
680 async fn create_image_if_missing(&self) -> Result<(), GenericError> {
681 debug!(
682 driver_id = self.config.driver_id,
683 isolation_group = self.isolation_group_id,
684 "Pulling image '{}'...",
685 self.config.image
686 );
687
688 self.create_image_if_missing_inner(self.config.image.as_str()).await?;
689
690 debug!(
691 driver_id = self.config.driver_id,
692 isolation_group = self.isolation_group_id,
693 "Pulled image '{}'.",
694 self.config.image
695 );
696
697 Ok(())
698 }
699
700 async fn create_volume_if_missing(&self) -> Result<(), GenericError> {
701 let volumes = self
703 .docker
704 .list_volumes(None::<bollard::query_parameters::ListVolumesOptions>)
705 .await?;
706 if volumes
707 .volumes
708 .iter()
709 .flatten()
710 .any(|volume| volume.name == self.isolation_group_name.as_str())
711 {
712 debug!("Shared volume '{}' already exists.", self.isolation_group_name);
713 return Ok(());
714 }
715
716 debug!(
717 driver_id = self.config.driver_id,
718 isolation_group = self.isolation_group_id,
719 "Shared volume '{}' does not exist. Creating...",
720 self.isolation_group_name
721 );
722
723 let volume_options = VolumeCreateRequest {
724 name: Some(self.isolation_group_name.clone()),
725 driver: Some("local".to_string()),
726 labels: Some(get_default_airlock_labels(self.isolation_group_id.as_str())),
727 ..Default::default()
728 };
729 self.docker.create_volume(volume_options).await?;
730
731 debug!(
732 driver_id = self.config.driver_id,
733 isolation_group = self.isolation_group_id,
734 "Created shared volume '{}'.",
735 self.isolation_group_name
736 );
737
738 Ok(())
739 }
740
741 async fn adjust_shared_volume_permissions(&self) -> Result<(), GenericError> {
742 debug!(
743 driver_id = self.config.driver_id,
744 isolation_group = self.isolation_group_id,
745 "Adjusting permissions on shared volume '{}'...",
746 self.container_name
747 );
748
749 let image = self.config.alpine_image.clone();
751 self.create_image_if_missing_inner(&image).await?;
752
753 let container_name = format!("airlock-{}-volume-fix-up", self.isolation_group_id);
754 let entrypoint = vec![
755 "chmod".to_string(),
756 "-R".to_string(),
757 "777".to_string(),
758 "/airlock".to_string(),
759 ];
760 let _ = self
761 .create_container_inner(container_name.clone(), image, Some(entrypoint), None, vec![], None)
762 .await?;
763
764 self.start_container_inner(&container_name).await?;
765 self.wait_for_container_exit_inner(&container_name).await?;
766 self.cleanup_inner(&container_name).await?;
767
768 Ok(())
769 }
770
771 async fn create_container_inner(
772 &self, container_name: String, image: String, entrypoint: Option<Vec<String>>, cmd: Option<Vec<String>>,
773 binds: Vec<String>, env: Option<Vec<String>>,
774 ) -> Result<String, GenericError> {
775 let binds = self.config.container_binds_from(&self.isolation_group_name, binds);
776
777 let networking_config = if !self.config.network_aliases.is_empty() {
779 let mut endpoints = HashMap::new();
780 endpoints.insert(
781 self.isolation_group_name.clone(),
782 EndpointSettings {
783 aliases: Some(self.config.network_aliases.clone()),
784 ..Default::default()
785 },
786 );
787 Some(
788 ContainerNetworkingConfig {
789 endpoints_config: endpoints,
790 }
791 .into(),
792 )
793 } else {
794 None
795 };
796
797 let (publish_all_ports, exposed_ports) = self.config.port_publishing_options();
798
799 let pid_mode = match self.config.container_os {
804 ContainerOs::Linux => Some("host".to_string()),
805 ContainerOs::Windows => None,
806 };
807 let cgroupns_mode = (self.config.container_os == ContainerOs::Linux && self.config.host_cgroup_namespace)
808 .then_some(HostConfigCgroupnsModeEnum::HOST);
809
810 let container_config = ContainerCreateBody {
811 hostname: Some(self.config.driver_id.to_string()),
812 env,
813 image: Some(image),
814 entrypoint,
815 cmd,
816 host_config: Some(HostConfig {
817 binds: Some(binds),
818 network_mode: Some(self.isolation_group_name.clone()),
819 publish_all_ports,
820 pid_mode,
821 cgroupns_mode,
822 ..Default::default()
823 }),
824 healthcheck: self.config.healthcheck.clone(),
825 exposed_ports,
826 labels: Some(get_default_airlock_labels(self.isolation_group_id.as_str())),
827 networking_config,
828 ..Default::default()
829 };
830
831 let create_options = CreateContainerOptionsBuilder::default().name(&container_name).build();
832
833 let response = self
834 .docker
835 .create_container(Some(create_options), container_config)
836 .await?;
837
838 Ok(response.id)
839 }
840
841 async fn create_container(&self) -> Result<(), GenericError> {
842 debug!(
843 driver_id = self.config.driver_id,
844 isolation_group = self.isolation_group_id,
845 "Creating container '{}'...",
846 self.container_name
847 );
848
849 let container_id = self
850 .create_container_inner(
851 self.container_name.clone(),
852 self.config.image.clone(),
853 self.config.entrypoint.clone(),
854 self.config.command.clone(),
855 self.config.binds.clone(),
856 Some(self.config.env.clone()),
857 )
858 .await?;
859
860 debug!(
861 driver_id = self.config.driver_id,
862 isolation_group = self.isolation_group_id,
863 "Created container '{}' (ID: {}).",
864 self.container_name,
865 container_id
866 );
867
868 Ok(())
869 }
870
871 async fn start_container_inner(&self, container_name: &str) -> Result<DriverDetails, GenericError> {
872 with_transient_error_retry("container start", &self.isolation_group_id, move || async move {
874 self.docker
875 .start_container(container_name, None)
876 .await
877 .map_err(GenericError::from)
878 })
879 .await?;
880
881 let mut details = DriverDetails {
882 container_name: container_name.to_string(),
883 ..Default::default()
884 };
885
886 let response = self.docker.inspect_container(container_name, None).await?;
887 if let Some(network_settings) = response.network_settings {
888 if let Some(networks) = network_settings.networks.as_ref() {
892 details.container_ip = networks
893 .get(&self.isolation_group_name)
894 .and_then(|settings| settings.ip_address.clone())
895 .filter(|address| !address.is_empty());
896 }
897
898 if let Some(ports) = network_settings.ports {
899 let port_mappings = details.port_mappings.get_or_insert_with(HashMap::new);
900 for (internal_port, bindings) in ports {
901 if let Some(bindings) = bindings {
902 for binding in bindings {
903 insert_port_mapping_if_parseable(
904 port_mappings,
905 internal_port.clone(),
906 binding.host_port.as_deref(),
907 );
908 if port_mappings.contains_key(&internal_port) {
909 break;
910 }
911 }
912 }
913 }
914 }
915 }
916
917 Ok(details)
918 }
919
920 async fn start_container(&self) -> Result<DriverDetails, GenericError> {
921 debug!(
922 driver_id = self.config.driver_id,
923 isolation_group = self.isolation_group_id,
924 "Starting container '{}'...",
925 self.container_name
926 );
927
928 let details = self.start_container_inner(&self.container_name).await?;
929
930 if let Some(log_dir) = self.log_dir.clone() {
931 debug!(
932 "Capturing logs for container '{}' to {}...",
933 self.container_name,
934 log_dir.display()
935 );
936
937 self.capture_container_logs(log_dir, self.config.driver_id, &self.container_name)
938 .await?;
939 }
940
941 debug!(
942 driver_id = self.config.driver_id,
943 isolation_group = self.isolation_group_id,
944 "Started container '{}'.",
945 self.container_name
946 );
947
948 Ok(details)
949 }
950
951 async fn connect_to_additional_networks(&self) -> Result<(), GenericError> {
956 for network in &self.config.additional_networks {
957 self.docker
958 .connect_network(
959 network,
960 NetworkConnectRequest {
961 container: self.container_name.clone(),
962 endpoint_config: None,
963 },
964 )
965 .await
966 .map_err(|e| {
967 generic_error!(
968 "Failed to connect container '{}' to network '{}': {}",
969 self.container_name,
970 network,
971 e
972 )
973 })?;
974 }
975 Ok(())
976 }
977
978 pub async fn start(&mut self) -> Result<DriverDetails, GenericError> {
985 self.create_network_if_missing().await?;
986 self.create_image_if_missing().await?;
987 self.create_volume_if_missing().await?;
988 if self.config.needs_shared_volume_permission_fixup() {
989 self.adjust_shared_volume_permissions().await?;
990 }
991
992 self.create_container().await?;
993 self.connect_to_additional_networks().await?;
994 self.start_container().await
995 }
996
997 pub async fn wait_for_container_healthy(&mut self) -> Result<(), GenericError> {
1005 loop {
1006 let response = self.docker.inspect_container(&self.container_name, None).await?;
1008 let state = response
1009 .state
1010 .ok_or_else(|| generic_error!("Container state should be present."))?;
1011
1012 let status = state
1014 .status
1015 .ok_or_else(|| generic_error!("Container status should be present."))?;
1016 if status != ContainerStateStatusEnum::RUNNING {
1017 return Err(generic_error!(
1018 "Container exited unexpectedly (driver_id: {}, container: {}). Check logs in the test run directory.",
1019 self.config.driver_id,
1020 self.container_name
1021 ));
1022 }
1023
1024 if let Some(health_status) = state.health.and_then(|h| h.status) {
1025 match health_status {
1026 HealthStatusEnum::EMPTY | HealthStatusEnum::NONE | HealthStatusEnum::HEALTHY => {
1028 debug!(
1029 driver_id = self.config.driver_id,
1030 "Container '{}' healthy or no healthcheck defined. Proceeding.", &self.container_name
1031 );
1032 return Ok(());
1033 }
1034
1035 HealthStatusEnum::STARTING => {
1037 debug!(
1038 driver_id = self.config.driver_id,
1039 "Container '{}' not yet healthy. Waiting...", &self.container_name
1040 );
1041 }
1042
1043 HealthStatusEnum::UNHEALTHY => {
1044 return Err(generic_error!(
1045 "Container became unhealthy (driver_id: {}, container: {}). Check logs in the test run directory.",
1046 self.config.driver_id,
1047 self.container_name
1048 ));
1049 }
1050 }
1051 } else {
1052 debug!(
1053 driver_id = self.config.driver_id,
1054 "Container '{}' has no healthcheck defined. Proceeding.", &self.container_name
1055 );
1056 return Ok(());
1057 }
1058
1059 sleep(Duration::from_secs(1)).await;
1061 }
1062 }
1063
1064 async fn wait_for_container_exit_inner(&self, container_name: &str) -> Result<ExitStatus, GenericError> {
1065 let mut wait_stream = self.docker.wait_container(container_name, None);
1066 match wait_stream.next().await {
1067 Some(result) => match result {
1068 Ok(response) => {
1069 assert_eq!(response.error, None);
1077 assert_eq!(response.status_code, 0);
1078
1079 Ok(ExitStatus::Success)
1080 }
1081
1082 Err(Error::DockerContainerWaitError { error, code }) => {
1083 let error = if error.is_empty() {
1084 String::from("<no error message provided>")
1085 } else {
1086 error
1087 };
1088 Ok(ExitStatus::Failed { code, error })
1089 }
1090
1091 Err(e) => Err(generic_error!("Failed to wait for container to finish: {:?}", e)),
1092 },
1093 None => unreachable!("Docker wait stream ended unexpectedly."),
1094 }
1095 }
1096
1097 pub async fn wait_for_container_exit(&self) -> Result<ExitStatus, GenericError> {
1106 debug!(
1107 driver_id = self.config.driver_id,
1108 isolation_group = self.isolation_group_id,
1109 "Waiting for container '{}' to finish...",
1110 &self.container_name
1111 );
1112
1113 let exit_status = self.wait_for_container_exit_inner(&self.container_name).await?;
1114
1115 debug!(
1116 driver_id = self.config.driver_id,
1117 isolation_group = self.isolation_group_id,
1118 "Container '{}' finished successfully.",
1119 &self.container_name
1120 );
1121
1122 Ok(exit_status)
1123 }
1124
1125 pub async fn exec_in_container(&self, cmd: Vec<String>) -> Result<String, GenericError> {
1134 let exec_opts = CreateExecOptions {
1135 attach_stdout: Some(true),
1136 attach_stderr: Some(false),
1137 cmd: Some(cmd.clone()),
1138 ..Default::default()
1139 };
1140
1141 let exec = self
1142 .docker
1143 .create_exec(&self.container_name, exec_opts)
1144 .await
1145 .with_error_context(|| format!("Failed to create exec instance for container {}.", self.container_name))?;
1146
1147 let exec_id = exec.id.clone();
1148
1149 let output = self
1150 .docker
1151 .start_exec(&exec.id, None)
1152 .await
1153 .with_error_context(|| format!("Failed to start exec for container {}.", self.container_name))?;
1154
1155 let mut stdout = String::new();
1156 if let StartExecResults::Attached { mut output, .. } = output {
1157 while let Some(chunk) = output.try_next().await? {
1158 if let LogOutput::StdOut { message } = chunk {
1159 stdout.push_str(&String::from_utf8_lossy(&message));
1160 }
1161 }
1162 }
1163
1164 let inspect = self
1166 .docker
1167 .inspect_exec(&exec_id)
1168 .await
1169 .error_context("Failed to inspect exec result.")?;
1170
1171 if let Some(code) = inspect.exit_code {
1172 if code != 0 {
1173 return Err(generic_error!(
1174 "Command {:?} exited with code {} in container {}.",
1175 cmd,
1176 code,
1177 self.container_name
1178 ));
1179 }
1180 }
1181
1182 Ok(stdout)
1183 }
1184
1185 async fn cleanup_inner(&self, container_name: &str) -> Result<(), GenericError> {
1186 self.docker.stop_container(container_name, None).await?;
1187 self.docker.remove_container(container_name, None).await?;
1188
1189 Ok(())
1190 }
1191
1192 pub async fn cleanup(self) -> Result<(), GenericError> {
1198 debug!(
1199 driver_id = self.config.driver_id,
1200 isolation_group = self.isolation_group_id,
1201 "Cleaning up container '{}'...",
1202 self.container_name
1203 );
1204
1205 let start = Instant::now();
1206
1207 self.cleanup_inner(&self.container_name).await?;
1208
1209 debug!(
1210 driver_id = self.config.driver_id,
1211 isolation_group = self.isolation_group_id,
1212 "Container '{}' removed after {:?}.",
1213 self.container_name,
1214 start.elapsed()
1215 );
1216
1217 Ok(())
1218 }
1219
1220 async fn capture_container_logs(
1221 &self, container_log_dir: PathBuf, log_name: &str, container_name: &str,
1222 ) -> Result<(), GenericError> {
1223 tokio::fs::create_dir_all(&container_log_dir)
1226 .await
1227 .error_context("Failed to create logs directory. Possible permissions issue.")?;
1228
1229 let stdout_log_path = container_log_dir.join(format!("{}.stdout.log", log_name));
1230 let stderr_log_path = container_log_dir.join(format!("{}.stderr.log", log_name));
1231
1232 let mut stdout_file = tokio::fs::File::create(&stdout_log_path)
1233 .await
1234 .map(BufWriter::new)
1235 .error_context("Failed to create standard output log file. Possible permissions issue.")?;
1236 let mut stderr_file = tokio::fs::File::create(&stderr_log_path)
1237 .await
1238 .map(BufWriter::new)
1239 .error_context("Failed to create standard error log file. Possible permissions issue.")?;
1240
1241 let logs_config = LogsOptions {
1243 follow: true,
1244 stdout: true,
1245 stderr: true,
1246 ..Default::default()
1247 };
1248 let mut log_stream = self.docker.logs(container_name, Some(logs_config));
1249
1250 tokio::spawn(async move {
1251 while let Some(log_result) = log_stream.next().await {
1252 match log_result {
1253 Ok(log) => match log {
1254 LogOutput::StdErr { message } => {
1255 if let Err(e) = stderr_file.write_all(&strip_ansi_codes(&message)).await {
1256 error!(error = %e, "Failed to write log line to standard error log file.");
1257 break;
1258 }
1259 if let Err(e) = stderr_file.flush().await {
1260 error!(error = %e, "Failed to flush standard error log file.");
1261 break;
1262 }
1263 }
1264 LogOutput::StdOut { message } => {
1265 if let Err(e) = stdout_file.write_all(&strip_ansi_codes(&message)).await {
1266 error!(error = %e, "Failed to write log line to standard output log file.");
1267 break;
1268 }
1269 if let Err(e) = stdout_file.flush().await {
1270 error!(error = %e, "Failed to flush standard output log file.");
1271 break;
1272 }
1273 }
1274 LogOutput::StdIn { .. } | LogOutput::Console { .. } => {}
1275 },
1276 Err(e) => {
1277 error!(error = %e, "Failed to read log line from container.");
1278 break;
1279 }
1280 }
1281 }
1282
1283 if let Err(e) = stdout_file.get_mut().sync_all().await {
1285 error!(error = %e, "Failed to fsync standard output log file.");
1286 }
1287
1288 if let Err(e) = stderr_file.get_mut().sync_all().await {
1289 error!(error = %e, "Failed to fsync standard error log file.");
1290 }
1291 });
1292
1293 Ok(())
1294 }
1295}
1296
1297fn strip_ansi_codes(input: &[u8]) -> Vec<u8> {
1299 let mut out = Vec::with_capacity(input.len());
1300 let mut i = 0;
1301 while i < input.len() {
1302 if input[i] == 0x1b && input.get(i + 1) == Some(&b'[') {
1303 i += 2;
1304 while i < input.len() && !input[i].is_ascii_alphabetic() {
1305 i += 1;
1306 }
1307 i += 1;
1308 } else {
1309 out.push(input[i]);
1310 i += 1;
1311 }
1312 }
1313 out
1314}
1315
1316fn is_transient_error(error: &GenericError) -> bool {
1318 is_transient_hns_error(error)
1319}
1320
1321fn is_transient_hns_error(error: &GenericError) -> bool {
1333 error.downcast_ref::<Error>().is_some_and(|error| match error {
1334 Error::DockerResponseServerError { message, .. } => message.contains("hnsCall") && message.contains("0x1392"),
1335 _ => false,
1336 })
1337}
1338
1339fn transient_error_retry_backoff(attempt: u32, seed: &str) -> Duration {
1345 let mut hasher = DefaultHasher::new();
1346 seed.hash(&mut hasher);
1347 attempt.hash(&mut hasher);
1348
1349 let jitter_ms = hasher.finish() % (TRANSIENT_ERROR_RETRY_MAX_JITTER.as_millis() as u64 + 1);
1350
1351 TRANSIENT_ERROR_RETRY_BASE_BACKOFF * attempt + Duration::from_millis(jitter_ms)
1352}
1353
1354async fn with_transient_error_retry<F, Fut, T>(
1366 operation: &str, seed: &str, mut run_attempt: F,
1367) -> Result<T, GenericError>
1368where
1369 F: FnMut() -> Fut,
1370 Fut: Future<Output = Result<T, GenericError>>,
1371{
1372 let mut attempt = 1;
1373 loop {
1374 match run_attempt().await {
1375 Ok(value) => return Ok(value),
1376 Err(e) => {
1377 if attempt >= TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS || !is_transient_error(&e) {
1378 return Err(e);
1379 }
1380
1381 let backoff = transient_error_retry_backoff(attempt, seed);
1382 warn!(
1383 isolation_group = seed,
1384 error = %e,
1385 "Transient error during {}. Retrying in {:?} (attempt {} of {}).",
1386 operation,
1387 backoff,
1388 attempt + 1,
1389 TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS
1390 );
1391
1392 sleep(backoff).await;
1393 attempt += 1;
1394 }
1395 }
1396 }
1397}
1398
1399fn get_default_airlock_labels(isolation_group_id: &str) -> HashMap<String, String> {
1400 let mut labels = HashMap::new();
1401 labels.insert("created_by".to_string(), "airlock".to_string());
1402 labels.insert("airlock-isolation-group".to_string(), isolation_group_id.to_string());
1403 labels
1404}
1405
1406#[cfg(test)]
1407mod tests {
1408 use super::*;
1409
1410 fn docker_server_error(status_code: u16, message: &str) -> GenericError {
1411 Error::DockerResponseServerError {
1412 status_code,
1413 message: message.to_string(),
1414 }
1415 .into()
1416 }
1417
1418 #[test]
1419 fn hns_object_already_exists_is_treated_as_transient() {
1420 let error = docker_server_error(
1422 500,
1423 "failed during hnsCallRawResponse: hnsCall failed in Win32: The object already exists. (0x1392)",
1424 );
1425
1426 assert!(is_transient_hns_error(&error));
1427 }
1428
1429 #[test]
1430 fn transient_hns_error_is_still_detected_through_added_context() {
1431 let error = docker_server_error(
1432 500,
1433 "failed during hnsCallRawResponse: hnsCall failed in Win32: The object already exists. (0x1392)",
1434 )
1435 .context("Failed to start container 'airlock-q9183x1a-target'");
1436
1437 assert!(is_transient_hns_error(&error));
1438 }
1439
1440 #[test]
1441 fn other_failures_are_not_treated_as_transient_hns_errors() {
1442 assert!(!is_transient_hns_error(&docker_server_error(
1444 500,
1445 "failed during hnsCallRawResponse: hnsCall failed in Win32: The system cannot find the file specified. (0x2)",
1446 )));
1447
1448 assert!(!is_transient_hns_error(&docker_server_error(
1450 409,
1451 "Conflict. The container name \"/airlock-q9183x1a-target\" is already in use",
1452 )));
1453
1454 assert!(!is_transient_hns_error(&generic_error!("something else went wrong")));
1455 }
1456
1457 #[test]
1458 fn transient_error_retry_backoff_grows_per_attempt_and_stays_within_bounds() {
1459 for attempt in 1..TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS {
1460 let backoff = transient_error_retry_backoff(attempt, "q9183x1a");
1461
1462 assert!(backoff >= TRANSIENT_ERROR_RETRY_BASE_BACKOFF * attempt);
1463 assert!(backoff <= TRANSIENT_ERROR_RETRY_BASE_BACKOFF * attempt + TRANSIENT_ERROR_RETRY_MAX_JITTER);
1464 }
1465 }
1466
1467 #[test]
1468 fn transient_error_retry_backoff_differs_between_isolation_groups() {
1469 assert_ne!(
1472 transient_error_retry_backoff(1, "q9183x1a"),
1473 transient_error_retry_backoff(1, "fajekrvm")
1474 );
1475 }
1476
1477 #[tokio::test]
1478 async fn with_transient_error_retry_retries_transient_failures_until_one_succeeds() {
1479 let attempts = std::cell::Cell::new(0);
1480
1481 let result = with_transient_error_retry("test operation", "q9183x1a", || async {
1482 attempts.set(attempts.get() + 1);
1483 if attempts.get() < TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS {
1484 Err(docker_server_error(
1485 500,
1486 "failed during hnsCallRawResponse: hnsCall failed in Win32: The object already exists. (0x1392)",
1487 ))
1488 } else {
1489 Ok(attempts.get())
1490 }
1491 })
1492 .await;
1493
1494 assert_eq!(result.unwrap(), TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS);
1495 assert_eq!(attempts.get(), TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS);
1496 }
1497
1498 #[tokio::test]
1499 async fn with_transient_error_retryy_gives_up_after_the_attempt_budget() {
1500 let attempts = std::cell::Cell::new(0);
1501
1502 let result: Result<(), GenericError> = with_transient_error_retry("test operation", "q9183x1a", || async {
1503 attempts.set(attempts.get() + 1);
1504 Err(docker_server_error(
1505 500,
1506 "failed during hnsCallRawResponse: hnsCall failed in Win32: The object already exists. (0x1392)",
1507 ))
1508 })
1509 .await;
1510
1511 assert!(result.is_err());
1512 assert_eq!(attempts.get(), TRANSIENT_ERROR_RETRY_MAX_ATTEMPTS);
1513 }
1514
1515 #[tokio::test]
1516 async fn with_transient_error_retry_does_not_retry_unrelated_failures() {
1517 let attempts = std::cell::Cell::new(0);
1518
1519 let result: Result<(), GenericError> = with_transient_error_retry("test operation", "q9183x1a", || async {
1520 attempts.set(attempts.get() + 1);
1521 Err(generic_error!("something else went wrong"))
1522 })
1523 .await;
1524
1525 assert!(result.is_err());
1526 assert_eq!(attempts.get(), 1);
1527 }
1528
1529 #[test]
1530 fn alpine_image_defaults_to_docker_hub_and_can_be_overridden() {
1531 let config = DriverConfig::from_image("target", "example:latest".to_string());
1532 assert_eq!(config.alpine_image, "alpine:latest");
1533
1534 let config = config.with_alpine_image("registry.example/alpine:3.20");
1535 assert_eq!(config.alpine_image, "registry.example/alpine:3.20");
1536 }
1537
1538 #[test]
1539 fn default_linux_container_binds_include_airlock_and_linux_host_resources() {
1540 let config = DriverConfig::from_image("target", "example:latest".to_string());
1541
1542 let binds = config.container_binds_from("airlock-test", config.binds.clone());
1543
1544 assert!(binds.contains(&"airlock-test:/airlock:z".to_string()));
1545 assert!(binds.contains(&"/proc:/host/proc:ro".to_string()));
1546 assert!(binds.contains(&"/sys/fs/cgroup:/host/sys/fs/cgroup:ro".to_string()));
1547 assert!(binds.contains(&"/var/run/docker.sock:/var/run/docker.sock:ro".to_string()));
1548 }
1549
1550 #[test]
1551 fn windows_container_binds_use_windows_airlock_and_skip_linux_host_resources() {
1552 let config =
1553 DriverConfig::from_image("target", "example:latest".to_string()).with_container_os(ContainerOs::Windows);
1554
1555 let binds = config.container_binds_from("airlock-test", config.binds.clone());
1556
1557 assert!(binds.contains(&"airlock-test:C:\\airlock".to_string()));
1558 assert!(!binds.iter().any(|bind| bind.contains("/proc")));
1559 assert!(!binds.iter().any(|bind| bind.contains("/sys/fs/cgroup")));
1560 assert!(!binds.iter().any(|bind| bind.contains("/var/run/docker.sock")));
1561 assert!(!binds.iter().any(|bind| bind.ends_with(":z")));
1562 }
1563
1564 #[tokio::test]
1565 async fn target_config_preserves_windows_container_os() {
1566 let target = TargetConfig {
1567 image: "example:latest".to_string(),
1568 entrypoint: vec![],
1569 command: vec![],
1570 additional_env_vars: vec![],
1571 container_os: ContainerOs::Windows,
1572 host_cgroup_namespace: false,
1573 };
1574
1575 let config = DriverConfig::target("target", target).await.unwrap();
1576
1577 assert_eq!(config.container_os, ContainerOs::Windows);
1578 assert!(!config.host_cgroup_namespace);
1579 }
1580
1581 #[tokio::test]
1582 async fn target_config_preserves_host_cgroup_namespace() {
1583 let target = TargetConfig {
1584 image: "example:latest".to_string(),
1585 entrypoint: vec![],
1586 command: vec![],
1587 additional_env_vars: vec![],
1588 container_os: ContainerOs::Linux,
1589 host_cgroup_namespace: true,
1590 };
1591
1592 let config = DriverConfig::target("target", target).await.unwrap();
1593
1594 assert!(config.host_cgroup_namespace);
1595 }
1596
1597 #[test]
1598 fn port_mapping_inserts_parseable_host_port() {
1599 let mut mappings = HashMap::new();
1600
1601 insert_port_mapping_if_parseable(&mut mappings, "55100/tcp", Some("49152"));
1602
1603 assert_eq!(mappings.get("55100/tcp"), Some(&49152));
1604 }
1605
1606 #[test]
1607 fn port_mapping_ignores_invalid_host_port() {
1608 let mut mappings = HashMap::new();
1609
1610 insert_port_mapping_if_parseable(&mut mappings, "55100/tcp", Some("not-a-port"));
1611
1612 assert!(!mappings.contains_key("55100/tcp"));
1613 }
1614
1615 #[test]
1616 fn windows_container_skips_shared_volume_permission_fixup() {
1617 let config =
1618 DriverConfig::from_image("target", "example:latest".to_string()).with_container_os(ContainerOs::Windows);
1619
1620 assert!(!config.needs_shared_volume_permission_fixup());
1621 }
1622
1623 #[test]
1624 fn windows_container_uses_nat_network_driver() {
1625 let config =
1626 DriverConfig::from_image("target", "example:latest".to_string()).with_container_os(ContainerOs::Windows);
1627
1628 assert_eq!(config.network_driver(), "nat");
1629 }
1630
1631 #[test]
1632 fn windows_container_exposes_ports_without_publishing_to_host() {
1633 let config = DriverConfig::from_image("target", "example:latest".to_string())
1634 .with_container_os(ContainerOs::Windows)
1635 .with_exposed_port("udp", 58125);
1636
1637 let (publish_all_ports, exposed_ports) = config.port_publishing_options();
1638
1639 assert_eq!(publish_all_ports, None);
1640 assert_eq!(exposed_ports, Some(vec!["58125/udp".to_string()]));
1641 }
1642
1643 #[test]
1644 fn linux_container_exposes_ports_and_publishes_to_host() {
1645 let config = DriverConfig::from_image("target", "example:latest".to_string()).with_exposed_port("tcp", 55100);
1646
1647 let (publish_all_ports, exposed_ports) = config.port_publishing_options();
1648
1649 assert_eq!(publish_all_ports, Some(true));
1650 assert_eq!(exposed_ports, Some(vec!["55100/tcp".to_string()]));
1651 }
1652
1653 #[test]
1654 fn linux_container_uses_bridge_network_driver() {
1655 let config = DriverConfig::from_image("target", "example:latest".to_string());
1656
1657 assert_eq!(config.network_driver(), "bridge");
1658 }
1659}