airlock/
driver.rs

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
33/// Default image for the shared-volume permission fix-up container.
34pub 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
46/// Configuration for transient error retries.
47///
48/// Transient error retries are a class of retries that are applied to a select few operations that
49/// have known error cases where retrying is safe and likely to fix the issue.
50const 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/// Container operating system for a driver target.
69#[derive(Clone, Copy, Debug, Eq, PartialEq)]
70pub enum ContainerOs {
71    /// Linux container defaults.
72    Linux,
73    /// Windows container defaults.
74    Windows,
75}
76
77/// Driver configuration.
78///
79/// This is the basic set of configuration options needed to spawn the container for a given driver.
80#[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 named Docker volume mounts, in `volume_name:/container/path` format.
93    ///
94    /// Unlike bind mounts specified via [`with_bind_mount`][Self::with_bind_mount], these reference
95    /// existing named Docker volumes rather than host filesystem paths. Used to mount volumes that
96    /// belong to other isolation groups (for example, a shared millstone mounting both the baseline
97    /// and comparison agent volumes).
98    additional_volume_mounts: Vec<String>,
99
100    /// DNS aliases for this container on its primary network.
101    ///
102    /// Set via `NetworkingConfig.EndpointsConfig` at container creation time. Other containers on
103    /// the same network can reach this container using any of these aliases in addition to its
104    /// hostname. Used to give agent containers unambiguous names (for example, `"baseline"`, `"comparison"`)
105    /// that the shared millstone can use to address each one independently.
106    network_aliases: Vec<String>,
107
108    /// Additional Docker networks to connect this container to after creation.
109    ///
110    /// The primary network is set via `HostConfig.NetworkMode`. Each network listed here is joined
111    /// via a separate `docker network connect` call after the container is created but before it's
112    /// started. Used to connect the shared millstone container to both agent networks so it can
113    /// reach `baseline` and `comparison` by hostname.
114    additional_networks: Vec<String>,
115
116    /// Image used for the shared-volume permission fix-up container.
117    ///
118    /// Defaults to [`DEFAULT_ALPINE_IMAGE`].
119    alpine_image: String,
120}
121
122impl DriverConfig {
123    pub async fn millstone(config: MillstoneConfig) -> Result<Self, GenericError> {
124        // Ensure the given configuration file path actually exists.
125        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            // Map our intake port to an ephemeral port on the host side, which we'll query once the container has been
175            // started so that we can connect to it.
176            .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    /// Creates a new `DriverConfig` from the given driver identifier and container image reference.
193    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    /// Sets the image used for the shared-volume permission fix-up container.
213    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    /// Sets the entrypoint for the container.
219    ///
220    /// If `entrypoint` is empty, the default entrypoint will be used.
221    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    /// Sets the command for the container.
229    ///
230    /// If `command` is empty, the default command will be used.
231    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    /// Adds an environment variable to the container.
239    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    /// Adds environment variables to the container.
249    pub fn with_env_vars(mut self, env: Vec<String>) -> Self {
250        self.env.extend(env);
251        self
252    }
253
254    /// Adds a bind mount to the container.
255    ///
256    /// `host_path` represents the path on the host to mount, while `container_path` represents the path on the
257    /// container side to mount it to. Bind mounts can be either files or directories.
258    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    /// Adds a read-only bind mount to the container.
269    ///
270    /// Same as [`with_bind_mount`][Self::with_bind_mount] but the container can't modify the mounted path.
271    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    /// Sets the healthcheck for the container.
286    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        // We manually insert "CMD" as the first value in the command array, so that it doesn't have to be done by the
291        // caller, since it's some goofy ass syntax to have to know about.
292        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    /// Adds a DNS alias for this container on its primary network.
306    ///
307    /// Other containers on the same network can resolve this container by `alias` in addition to
308    /// its hostname. Call this before the container is started.
309    pub fn with_network_alias(mut self, alias: impl Into<String>) -> Self {
310        self.network_aliases.push(alias.into());
311        self
312    }
313
314    /// Connects this container to an additional Docker network after creation.
315    ///
316    /// The primary network is always the container's isolation group network. Each network added
317    /// here is joined via `docker network connect` after the container is created but before it
318    /// is started, so the container is reachable on all listed networks from the moment it runs.
319    pub fn with_network(mut self, network: impl Into<String>) -> Self {
320        self.additional_networks.push(network.into());
321        self
322    }
323
324    /// Mounts a named Docker volume into the container at the given path.
325    ///
326    /// Unlike [`with_bind_mount`][Self::with_bind_mount], this references a named Docker volume
327    /// rather than a host filesystem path. The volume must already exist when the container starts.
328    /// This is useful for mounting volumes that belong to other isolation groups: for example,
329    /// a shared millstone container that needs to reach the DogStatsD sockets of both the baseline
330    /// and comparison agent containers.
331    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    /// Adds an exposed port to the container.
338    ///
339    /// The `protocol` should be either `tcp` or `udp`. Linux containers publish exposed ports to ephemeral host ports,
340    /// which are returned in [`DriverDetails`] after starting the driver. Windows containers keep exposed ports internal
341    /// to the container network because Panoramic probes them from inside the container or via the container IP.
342    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    /// Sets the operating system this container will run as.
348    ///
349    /// The OS choice drives several non-portable defaults (network driver, default binds, host
350    /// resources to share, container path conventions) that the rest of the driver applies
351    /// automatically through the helpers below. Callers should set this before any binds or
352    /// health checks are added so OS-specific defaults are appended consistently.
353    pub fn with_container_os(mut self, container_os: ContainerOs) -> Self {
354        self.container_os = container_os;
355        self
356    }
357
358    /// Configures whether a Linux target joins the Docker host's cgroup namespace.
359    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    /// Whether the shared `/airlock` volume needs a one-shot world-writable chmod fix-up.
365    ///
366    /// Linux Docker volumes default to root-owned with restrictive permissions, so containers
367    /// running as non-root users (the Datadog Agent image, in particular) cannot write to
368    /// `/airlock` without an out-of-band chmod. We do that fix-up by spawning a short-lived
369    /// Alpine container that owns the volume mount and runs `chmod -R 777 /airlock`. Windows
370    /// containers do not have the same UID/permission model and the fix-up is unnecessary
371    /// (and unsupported, since Alpine is a Linux image).
372    fn needs_shared_volume_permission_fixup(&self) -> bool {
373        self.container_os == ContainerOs::Linux
374    }
375
376    /// Docker network driver to use for the isolation group network on this container's OS.
377    ///
378    /// Linux containers use the `bridge` driver; Windows containers use `nat` (the only
379    /// driver that supports container-to-container traffic on a single Windows host).
380    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    /// Returns the full set of bind mounts to apply to this container, including OS-specific
406    /// defaults and any additional named volume mounts.
407    ///
408    /// Linux containers receive the shared `/airlock` volume plus read-only mounts of host
409    /// paths needed for origin detection (`/proc`, `/sys/fs/cgroup`, the Docker socket).
410    /// Windows containers receive only the shared `C:\airlock` volume; the host-resource
411    /// mounts have no Windows-container equivalent and the `:z` shared-relabel mount option is
412    /// Linux-specific.
413    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/// Detailed information about the spawned container.
432#[derive(Debug, Default)]
433pub struct DriverDetails {
434    container_name: String,
435    container_ip: Option<String>,
436    port_mappings: Option<HashMap<String, u16>>,
437}
438
439/// Inserts an `internal_port` -> host port mapping when `host_port` parses as a valid `u16`.
440///
441/// Docker reports each binding's host port as a string, and we treat values that don't parse as
442/// "no mapping available" rather than failing the whole inspect call. `internal_port` is the
443/// existing key (already including the protocol suffix, for example `"58125/udp"`).
444fn 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    /// Returns the name of the container.
454    pub fn container_name(&self) -> &str {
455        &self.container_name
456    }
457
458    /// Returns the container IP address on its primary Docker network, if known.
459    pub fn container_ip(&self) -> Option<&str> {
460        self.container_ip.as_deref()
461    }
462
463    /// Attempts to look up a mapped ephemeral port for the given exposed port.
464    ///
465    /// The same `protocol` and internal port values used to expose the port must be used here. If the given
466    /// protocol/port combination wasn't exposed, `None` is returned. Otherwise, the mapped ephemeral port is returned.
467    /// This port is exposed on `0.0.0.0` on the host side.
468    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
475/// Container driver.
476pub 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    /// Creates a new `Driver` from the given isolation group ID and configuration.
487    ///
488    /// # Isolation group
489    ///
490    /// The isolation group ID serves as a unique identifier to be used for both the name of the container as well as
491    /// the shared resources that are created and attached to the container. If two drivers share the same isolation
492    /// group ID, the containers they spawn will be located in the same network namespace, have access to the same
493    /// shared Airlock volume, etc.
494    ///
495    /// # Shared volume
496    ///
497    /// The container will have a volume bind-mounted at `/airlock` that's shared between all containers in the same
498    /// isolation group. This volume is mounted as world writeable (777) so all containers can freely read and write to
499    /// it. This makes it easier for containers to share data between one another, but also means that care should be
500    /// taken to avoid conflicts between trying to write to the same file, etc.
501    ///
502    /// # Errors
503    ///
504    /// If the Docker client can't be created/configured, an error will be returned.
505    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    /// Configures the driver to capture container logs.
519    ///
520    /// The logs will be stored in the given directory, under a subdirectory named after the isolation group ID. Each
521    /// container will get a log for standard output and standard error, following the pattern of `<container
522    /// name>.[stdout|stderr].log`.
523    pub fn with_logging(mut self, log_dir: PathBuf) -> Self {
524        self.log_dir = Some(log_dir);
525        self
526    }
527
528    /// Returns the string identifier of the driver.
529    ///
530    /// This is generally a shorthand of the application/service, such as `dogstatsd` or `millstone`.
531    pub fn driver_id(&self) -> &'static str {
532        self.config.driver_id
533    }
534
535    /// Clean up any containers, networks, and volumes related to the given isolation group ID.
536    ///
537    /// This is a free function to facilitate cleaning up resources after a number of drivers are run.
538    ///
539    /// # Errors
540    ///
541    /// If the Docker client can't be created/configured, or there is an error when finding or removing any of the
542    /// related resources, an error will be returned.
543    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        // Remove any containers related to the isolation group. We do so forcefully.
550        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        // Remove the shared volume.
592        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        // Remove the network.
605        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        // See if the network already exists or not.
623        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        // Create the network since it doesn't yet exist.
640        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        // Check to see if the shared volume already exists.
702        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        // We spin up a minimal Alpine container, chmod the directory bind-mounted to the shared volume, and that's it.
750        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        // Set up NetworkingConfig to apply aliases on the primary network, if any are configured.
778        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        // Linux test containers run with `pid_mode=host` so origin-detection logic in ADP and
800        // the Core Agent can see processes on the runner. Windows containers do not support
801        // host PID mode, so we leave it unset and accept that Windows-runtime tests don't
802        // exercise the host-pid origin-detection path.
803        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        // Retrying is safe: the daemon treats starting an already-running container as a no-op rather than an error.
873        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            // Look up the IP only on the primary isolation-group network. Falling back to
889            // "any other network's IP" would be non-deterministic and effectively wrong for
890            // assertion targeting.
891            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    /// Connects this container to each network listed in `additional_networks`.
952    ///
953    /// Called after container creation but before start, so the container is already reachable
954    /// on all configured networks from the moment it begins running.
955    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    /// Starts the container, creating any necessary resources.
979    ///
980    /// # Errors
981    ///
982    /// If there is an error while creating the network or shared volume, while pulling the container image, or while
983    /// creating or starting the container, it will be returned.
984    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    /// Waits until the container is marked as healthy.
998    ///
999    /// If the container has no health checks defined, this returns early and does no waiting.
1000    ///
1001    /// # Errors
1002    ///
1003    /// If there is an error while inspecting the container, it will be returned.
1004    pub async fn wait_for_container_healthy(&mut self) -> Result<(), GenericError> {
1005        loop {
1006            // Inspect the container, and see if it even has any health checks defined. If not, then we can return early.
1007            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            // Make sure the container is actually running.
1013            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                    // No healthcheck defined, or healthy, so we're good to go.
1027                    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                    // Not healthy yet, so we'll keep waiting.
1036                    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            // Wait for a second and then check again.
1060            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                    // When the exit code is non-zero, `bollard` transforms the normal `ContainerWaitResponse` into
1070                    // `Error::DockerContainerWaitError`, which is why we have these asserts here to catch any scenario
1071                    // where there's _somehow_ an error condition being indicated without it having been transformed into
1072                    // `Error::DockerContainerWaitError`.
1073                    //
1074                    // Essentially, getting to this point should imply successfully exiting, but the API isn't very
1075                    // ergonomic in that regard, so we're just making sure.
1076                    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    /// Waits for the container to finish successfully.
1098    ///
1099    /// The container's exit status is returned, indicating success (exit code 0) or failure (exit code != 0), including
1100    /// any error message related to the failure.
1101    ///
1102    /// # Errors
1103    ///
1104    /// If an error is encountered while waiting for the container to exit, it will be returned.
1105    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    /// Executes a command inside the running container and returns its stdout.
1126    ///
1127    /// The command runs as root with no TTY. Stderr is discarded, and only stdout is returned. If the command exits with a
1128    /// nonzero status, an error is returned.
1129    ///
1130    /// # Errors
1131    ///
1132    /// If the exec creation, start, output collection, or command exit code indicates failure, an error is returned.
1133    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        // Check the command's exit code.
1165        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    /// Cleans up the container, stopping and removing it from the system.
1193    ///
1194    /// # Errors
1195    ///
1196    /// If there is an error while stopping or removing the container, it will be returned.
1197    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        // Make sure the directories exist first and prepare the files, just to get any permissions issues out of the
1224        // way up front before we spawn our background task.
1225        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        // Spawn a background task to capture the logs.
1242        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            // One final fsync to ensure the logs are fully written to disk.
1284            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
1297/// Removes ANSI escape sequences (`ESC[...letter`) from a byte slice.
1298fn 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
1316/// Returns `true` if `error` is a transient error that can be retried.
1317fn is_transient_error(error: &GenericError) -> bool {
1318    is_transient_hns_error(error)
1319}
1320
1321/// Returns `true` if `error` is a transient Windows HNS collision.
1322///
1323/// Windows containers are networked through the Host Network Service, and HNS serializes poorly: when several
1324/// containers create a NAT network or attach an endpoint at the same moment, the losing caller gets a 500 from the
1325/// daemon carrying `hnsCall failed in Win32: The object already exists. (0x1392)`. The object in question is internal
1326/// to HNS -- airlock names every network after a unique isolation group, so this is never a name conflict of ours --
1327/// and the conflict clears once the winning call completes.
1328///
1329/// The predicate matches on the message rather than on [`ContainerOs`] because `hnsCall` has no analogue on Linux,
1330/// which makes it inert there. It deliberately requires the `0x1392` code as well as `hnsCall`: other HNS failures
1331/// indicate real breakage and should fail fast rather than be papered over by a retry.
1332fn 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
1339/// Returns how long to wait before the given transient error retry attempt.
1340///
1341/// `seed` keeps concurrent callers out of step with each other. Isolation group IDs are unique per
1342/// test, so two drivers that collided at the same instant draw different jitter and don't line
1343/// back up on the retry.
1344fn 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
1354/// Runs `run_attempt`, retrying it when it fails with a known transient error.
1355///
1356/// Any other error, and the error from the final attempt, is returned to the caller untouched. "Known" transient errors
1357/// are hard-coded/categorized by hand, so this function mostly represents known areas where we've seen transient issues
1358/// that can generally be solved by simply retrying.
1359///
1360/// `run_attempt` must be safe to re-run: every call site either re-checks for the resource it is creating, or targets a
1361/// Docker endpoint that is idempotent.
1362///
1363/// Each retry is logged at warning level so that runs which only passed because of a retry stay visible in CI output,
1364/// rather than the mitigation quietly hiding how often the collision fires.
1365async 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        // Verbatim message from a failed `test-integration-windows-amd64` run.
1421        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        // A different HNS failure mode: real breakage, not a collision worth papering over.
1443        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        // A name conflict of ours, which a retry would never clear.
1449        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        // Racers that collided at the same instant must not retry on the same schedule, or they
1470        // just collide again. Isolation group IDs are taken from real failing runs.
1471        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}