diff --git a/crates/openshell-core/src/local_api_socket.rs b/crates/openshell-core/src/local_api_socket.rs index 6804ae513a..7dbf2237fb 100644 --- a/crates/openshell-core/src/local_api_socket.rs +++ b/crates/openshell-core/src/local_api_socket.rs @@ -5,6 +5,30 @@ use std::path::{Path, PathBuf}; +/// Return the first responsive local Docker API socket. +#[must_use] +pub fn detect_docker_socket() -> Option { + let mut candidates = Vec::new(); + if let Ok(host) = std::env::var("DOCKER_HOST") + && let Some(path) = host.trim().strip_prefix("unix://") + && !path.is_empty() + { + candidates.push(PathBuf::from(path)); + } + candidates.push(PathBuf::from("/var/run/docker.sock")); + if let Some(home) = std::env::var_os("HOME") { + candidates.push(PathBuf::from(home).join(".docker/run/docker.sock")); + } + if let Some(runtime_dir) = std::env::var_os("XDG_RUNTIME_DIR") { + candidates.push(PathBuf::from(runtime_dir).join("docker.sock")); + } + first_responsive_socket(&candidates, |response| { + http_response_is_success(response) + && contains_ascii(response, b"Api-Version:") + && !contains_ascii(response, b"Libpod-Api-Version:") + }) +} + /// Return the first candidate whose HTTP ping response is accepted. #[must_use] pub fn first_responsive_socket( diff --git a/crates/openshell-driver-docker/src/lib.rs b/crates/openshell-driver-docker/src/lib.rs index dad13d74a4..ae2e9c9c63 100644 --- a/crates/openshell-driver-docker/src/lib.rs +++ b/crates/openshell-driver-docker/src/lib.rs @@ -815,25 +815,7 @@ impl ComputeDriverService { /// Return the first responsive local Docker API socket. #[must_use] pub fn detect_socket() -> Option { - let mut candidates = Vec::new(); - if let Ok(host) = std::env::var("DOCKER_HOST") - && let Some(path) = host.trim().strip_prefix("unix://") - && !path.is_empty() - { - candidates.push(PathBuf::from(path)); - } - candidates.push(PathBuf::from("/var/run/docker.sock")); - if let Some(home) = std::env::var_os("HOME") { - candidates.push(PathBuf::from(home).join(".docker/run/docker.sock")); - } - if let Some(runtime_dir) = std::env::var_os("XDG_RUNTIME_DIR") { - candidates.push(PathBuf::from(runtime_dir).join("docker.sock")); - } - openshell_core::local_api_socket::first_responsive_socket(&candidates, |response| { - openshell_core::local_api_socket::http_response_is_success(response) - && openshell_core::local_api_socket::contains_ascii(response, b"Api-Version:") - && !openshell_core::local_api_socket::contains_ascii(response, b"Libpod-Api-Version:") - }) + openshell_core::local_api_socket::detect_docker_socket() } #[must_use] diff --git a/crates/openshell-driver-vm/src/driver.rs b/crates/openshell-driver-vm/src/driver.rs index 799cb0b6e8..a3a052fb33 100644 --- a/crates/openshell-driver-vm/src/driver.rs +++ b/crates/openshell-driver-vm/src/driver.rs @@ -5047,30 +5047,63 @@ fn parse_registry_reference(image_ref: &str) -> Result { } /// Try to connect to a local container engine (Docker or Podman). -/// -/// Tries Docker first (`connect_with_local_defaults`, which respects -/// `DOCKER_HOST`). If Docker is unavailable, falls back to the Podman -/// socket, which exposes a Docker-compatible API. async fn connect_local_container_engine() -> Option { + connect_container_engine(detect_docker_socket, detect_podman_socket).await +} + +/// Connect to the first reachable engine, trying in order: +/// +/// 1. `connect_with_local_defaults`, which honours any `DOCKER_HOST` scheme +/// (including `tcp://` and `ssh://`) and otherwise `/var/run/docker.sock`. +/// 2. A Docker socket found by the shared core discovery. Docker Desktop +/// for macOS may expose only a per-user socket, which step 1 never probes. +/// A failed `DOCKER_HOST` does not stop the search: an unreachable engine +/// is replaced by a live local one, as it already was by Podman. +/// 3. The Podman socket, which exposes a Docker-compatible API. +/// +/// Returning `None` is not an error: the caller falls back to pulling the +/// image from its registry. The detectors are parameters so the ordering can +/// be tested without depending on which engines the host runs. +async fn connect_container_engine( + detect_docker_socket: impl FnOnce() -> Option, + detect_podman_socket: impl FnOnce() -> Option, +) -> Option { if let Ok(docker) = Docker::connect_with_local_defaults() && docker.ping().await.is_ok() { return Some(docker); } - let podman_socket = detect_podman_socket()?; - if let Ok(docker) = - Docker::connect_with_unix(podman_socket.to_str()?, 120, bollard::API_DEFAULT_VERSION) - && docker.ping().await.is_ok() + if let Some(socket) = detect_docker_socket() + && let Some(docker) = connect_unix_engine(&socket).await { info!( - socket = %podman_socket.display(), - "vm driver: connected to Podman (Docker-compatible API)" + socket = %socket.display(), + "vm driver: connected to Docker through a discovered socket" ); return Some(docker); } - None + let socket = detect_podman_socket()?; + let docker = connect_unix_engine(&socket).await?; + info!( + socket = %socket.display(), + "vm driver: connected to Podman (Docker-compatible API)" + ); + Some(docker) +} + +/// Connect to a Docker-compatible API on a Unix socket, or `None` if it does +/// not answer a ping. +async fn connect_unix_engine(socket: &Path) -> Option { + let docker = + Docker::connect_with_unix(socket.to_str()?, 120, bollard::API_DEFAULT_VERSION).ok()?; + docker.ping().await.ok()?; + Some(docker) +} + +fn detect_docker_socket() -> Option { + openshell_core::local_api_socket::detect_docker_socket() } fn detect_podman_socket() -> Option { @@ -7431,6 +7464,147 @@ mod tests { static ENV_LOCK: std::sync::LazyLock> = std::sync::LazyLock::new(|| std::sync::Mutex::new(())); + /// Serve the Docker `/_ping` endpoint on `socket` and count the requests. + fn spawn_fake_engine(socket: &Path) -> Arc { + let listener = tokio::net::UnixListener::bind(socket).expect("bind fake engine socket"); + let requests = Arc::new(AtomicUsize::new(0)); + let served = Arc::clone(&requests); + tokio::spawn(async move { + while let Ok((mut stream, _)) = listener.accept().await { + let served = Arc::clone(&served); + tokio::spawn(async move { + use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _}; + let mut request = [0_u8; 1024]; + let _ = stream.read(&mut request).await; + served.fetch_add(1, Ordering::SeqCst); + let _ = stream + .write_all( + b"HTTP/1.1 200 OK\r\nApi-Version: 1.43\r\n\ + Content-Type: text/plain\r\nContent-Length: 2\r\n\ + Connection: close\r\n\r\nOK", + ) + .await; + }); + } + }); + requests + } + + /// Run `body` with `DOCKER_HOST` naming a socket that does not exist, so + /// the default connection step fails on every host, including one that + /// really runs Docker. + fn run_with_unreachable_default_engine(dir: &Path, body: impl Future) { + // A failed assertion in one of these tests must not poison the lock and + // hide the result of the others. + let _guard = ENV_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let missing = format!("unix://{}", dir.join("absent.sock").display()); + temp_env::with_var("DOCKER_HOST", Some(missing), || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(body); + }); + } + + #[test] + fn container_engine_uses_docker_socket_found_by_discovery() { + let dir = tempfile::tempdir().unwrap(); + let socket = dir.path().join("docker.sock"); + run_with_unreachable_default_engine(dir.path(), async { + let requests = spawn_fake_engine(&socket); + + let engine = connect_container_engine(|| Some(socket.clone()), || None).await; + + assert!( + engine.is_some(), + "a Docker socket outside DOCKER_HOST and the default path must be used" + ); + assert!(requests.load(Ordering::SeqCst) >= 1); + }); + } + + #[test] + fn container_engine_prefers_discovered_docker_over_podman() { + let dir = tempfile::tempdir().unwrap(); + let docker_socket = dir.path().join("docker.sock"); + let podman_socket = dir.path().join("podman.sock"); + run_with_unreachable_default_engine(dir.path(), async { + let docker_requests = spawn_fake_engine(&docker_socket); + let podman_requests = spawn_fake_engine(&podman_socket); + + let engine = connect_container_engine( + || Some(docker_socket.clone()), + || Some(podman_socket.clone()), + ) + .await; + + assert!(engine.is_some()); + assert!(docker_requests.load(Ordering::SeqCst) >= 1); + assert_eq!(podman_requests.load(Ordering::SeqCst), 0); + }); + } + + #[test] + fn container_engine_falls_back_to_podman_without_a_docker_socket() { + let dir = tempfile::tempdir().unwrap(); + let podman_socket = dir.path().join("podman.sock"); + run_with_unreachable_default_engine(dir.path(), async { + let podman_requests = spawn_fake_engine(&podman_socket); + + let engine = connect_container_engine(|| None, || Some(podman_socket.clone())).await; + + assert!(engine.is_some()); + assert!(podman_requests.load(Ordering::SeqCst) >= 1); + }); + } + + #[test] + fn container_engine_default_detector_uses_shared_docker_discovery() { + let dir = tempfile::tempdir().unwrap(); + let socket = dir.path().join("real.sock"); + let _guard = ENV_LOCK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + // Discovery honours a `unix://` DOCKER_HOST first, so this finds the + // fake engine on every host, whatever sockets it really has. + temp_env::with_var( + "DOCKER_HOST", + Some(format!("unix://{}", socket.display())), + || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap() + .block_on(async { + let _requests = spawn_fake_engine(&socket); + // The probe is blocking I/O; run it off the runtime + // thread that serves the fake engine. + let found = tokio::task::spawn_blocking(detect_docker_socket) + .await + .unwrap(); + assert_eq!(found, Some(socket.clone())); + }); + }, + ); + } + + #[test] + fn container_engine_is_unavailable_when_no_discovered_socket_answers() { + let dir = tempfile::tempdir().unwrap(); + let stale_socket = dir.path().join("stale.sock"); + // Leave the socket file behind with nothing listening on it. + drop(std::os::unix::net::UnixListener::bind(&stale_socket).unwrap()); + assert!(stale_socket.exists()); + run_with_unreachable_default_engine(dir.path(), async { + let engine = connect_container_engine(|| Some(stale_socket.clone()), || None).await; + + assert!(engine.is_none()); + }); + } + #[test] fn vm_console_diagnostic_is_bounded_to_the_tail() { let directory = tempfile::tempdir().unwrap();