Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions crates/openshell-core/src/local_api_socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<PathBuf> {
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(
Expand Down
20 changes: 1 addition & 19 deletions crates/openshell-driver-docker/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -815,25 +815,7 @@ impl ComputeDriverService {
/// Return the first responsive local Docker API socket.
#[must_use]
pub fn detect_socket() -> Option<PathBuf> {
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]
Expand Down
196 changes: 185 additions & 11 deletions crates/openshell-driver-vm/src/driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5047,30 +5047,63 @@ fn parse_registry_reference(image_ref: &str) -> Result<Reference, Status> {
}

/// 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<Docker> {
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<PathBuf>,
detect_podman_socket: impl FnOnce() -> Option<PathBuf>,
) -> Option<Docker> {
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<Docker> {
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<PathBuf> {
openshell_core::local_api_socket::detect_docker_socket()
}

fn detect_podman_socket() -> Option<PathBuf> {
Expand Down Expand Up @@ -7431,6 +7464,147 @@ mod tests {
static ENV_LOCK: std::sync::LazyLock<std::sync::Mutex<()>> =
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<AtomicUsize> {
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<Output = ()>) {
// 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();
Expand Down
Loading