use super::{Executor, GpuRequest, OwnedRunner, ProvisionError, RunnerSpec, RunnerState};
use async_trait::async_trait;
use std::sync::Arc;
use std::time::Duration;
fn is_safe_docker_identifier(s: &str) -> bool {
if s.is_empty() || s.len() > 255 {
return false;
}
if s.starts_with('-') {
return false;
}
s.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'.' | b'_' | b'-' | b'/' | b':' | b'@'))
}
pub(crate) fn is_docker_not_found(stderr: &str) -> bool {
let lower = stderr.to_ascii_lowercase();
lower.contains("no such object") || lower.contains("no such container")
}
fn runner_listener_running(bin: &str, container: &str) -> bool {
std::process::Command::new(bin)
.args(["exec", container, "pgrep", "-x", "Runner.Listener"])
.output()
.map(|o| o.status.success())
.unwrap_or(false)
}
pub(super) fn map_container_status(status: &str) -> RunnerState {
match status {
"running" => RunnerState::Healthy,
"created" | "restarting" => RunnerState::Starting,
"exited" | "dead" => RunnerState::Terminated {
exit_code: None,
last_logs: String::new(),
},
_ => RunnerState::Starting, }
}
pub struct DockerExecutor {
client: Arc<crate::docker::client::DockerClient>,
}
impl DockerExecutor {
pub fn new() -> Self {
Self {
client: Arc::new(crate::docker::client::DockerClient::new()),
}
}
pub fn client_ping(&self) -> Result<String, crate::docker::errors::DockerError> {
self.client.ping()
}
}
#[async_trait]
impl Executor for DockerExecutor {
fn settle_timeout(&self) -> Duration {
Duration::from_secs(120)
}
fn validate(&self, spec: &RunnerSpec) -> Result<(), ProvisionError> {
if !is_safe_docker_identifier(&spec.name) {
return Err(ProvisionError::Incompatible(format!(
"unsafe docker runner name '{}'",
spec.name
)));
}
if !is_safe_docker_identifier(&spec.image) {
return Err(ProvisionError::Incompatible(format!(
"unsafe docker image '{}'",
spec.image
)));
}
Ok(())
}
async fn inspect(&self, name: &str) -> Result<RunnerState, ProvisionError> {
let bin = "docker";
let out = std::process::Command::new(bin)
.args([
"inspect",
"--format",
"{{.State.Status}} {{.State.ExitCode}}",
"--",
name,
])
.output()
.map_err(|e| ProvisionError::transient(format!("docker inspect: {e}")))?;
if !out.status.success() {
let stderr = String::from_utf8_lossy(&out.stderr);
if is_docker_not_found(&stderr) {
return Ok(RunnerState::Absent);
}
return Err(ProvisionError::transient(format!(
"docker inspect failed: {}",
stderr
)));
}
let raw = String::from_utf8_lossy(&out.stdout);
let raw = raw.trim();
let mut parts = raw.split_whitespace();
let status = parts.next().unwrap_or("");
let exit_code: Option<i32> = parts.next().and_then(|s| s.parse().ok());
let mut state = map_container_status(status);
if let RunnerState::Terminated { .. } = state {
let logs = std::process::Command::new(bin)
.args(["logs", "--tail", "30", "--", name])
.output()
.map(|o| {
let s = String::from_utf8_lossy(&o.stdout);
let e = String::from_utf8_lossy(&o.stderr);
format!("{e}\n{s}")
})
.unwrap_or_default();
state = RunnerState::Terminated {
exit_code,
last_logs: logs,
};
}
if state == RunnerState::Healthy && !runner_listener_running(bin, name) {
state = RunnerState::Starting;
}
Ok(state)
}
async fn spawn(&self, spec: &RunnerSpec) -> Result<(), ProvisionError> {
use crate::docker::models::{ContainerCommand, GpuSelection, RunnerContainerSpec};
let gpus = match &spec.gpu {
GpuRequest::None => GpuSelection::None,
GpuRequest::All => GpuSelection::All,
GpuRequest::Count(n) => GpuSelection::Count(*n),
};
let container_spec = RunnerContainerSpec {
name: spec.name.clone(),
image: spec.image.clone(),
gpus,
cpus: Some(spec.cpu),
memory_gb: Some(spec.memory_gb),
env: vec![("CIRUN_RUNNER_NAME".into(), spec.name.clone())],
command: ContainerCommand::Script(spec.provision_script.clone()),
privileged: spec.docker_privileged,
mount_docker_socket: spec.docker_mount_socket,
};
self.client
.run_runner(&container_spec)
.map(|_| ())
.map_err(|e| ProvisionError::transient(format!("docker run failed: {e}")))
}
async fn kill(&self, name: &str) -> Result<(), ProvisionError> {
self.client
.stop_and_remove(name)
.map_err(|e| ProvisionError::transient(format!("docker rm: {e}")))
}
async fn list_owned(&self) -> Result<Vec<OwnedRunner>, ProvisionError> {
let infos = self
.client
.list_runner_containers("cirun.runner=true")
.map_err(|e| ProvisionError::transient(format!("docker ps: {e}")))?;
Ok(infos
.into_iter()
.map(|i| OwnedRunner {
name: i.name,
state: map_container_status(&i.state),
})
.collect())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn running_maps_to_healthy() {
assert_eq!(map_container_status("running"), RunnerState::Healthy);
}
#[test]
fn created_maps_to_starting() {
assert_eq!(map_container_status("created"), RunnerState::Starting);
}
#[test]
fn restarting_maps_to_starting() {
assert_eq!(map_container_status("restarting"), RunnerState::Starting);
}
#[test]
fn exited_maps_to_terminated() {
assert!(matches!(
map_container_status("exited"),
RunnerState::Terminated { .. }
));
}
#[test]
fn dead_maps_to_terminated() {
assert!(matches!(
map_container_status("dead"),
RunnerState::Terminated { .. }
));
}
#[test]
fn safe_identifier_accepts_normal_images_and_names() {
for s in [
"ubuntu:24.04",
"ghcr.io/aktech/runner:latest",
"cirun-gpu-runner:latest",
"image@sha256:abcdef0123",
"cirun-aktech--repo-1a5bc0163a",
] {
assert!(is_safe_docker_identifier(s), "rejected: {s}");
}
}
#[test]
fn safe_identifier_rejects_flag_shapes() {
for s in [
"--privileged",
"-v=/:/host",
"--security-opt=apparmor=unconfined",
"-",
] {
assert!(!is_safe_docker_identifier(s), "accepted leading-dash: {s}");
}
}
#[test]
fn safe_identifier_rejects_shell_metachars_and_whitespace() {
for s in [
"ubuntu;rm -rf /",
"ubuntu image",
"ubuntu`whoami`",
"ubuntu|sh",
] {
assert!(!is_safe_docker_identifier(s), "accepted: {s:?}");
}
}
fn rspec(name: &str, image: &str) -> RunnerSpec {
RunnerSpec {
name: name.into(),
provision_script: String::new(),
image: image.into(),
cpu: 2,
memory_gb: 4,
disk_gb: 20,
gpu: GpuRequest::None,
docker_privileged: false,
docker_mount_socket: false,
login: crate::executor::RunnerLogin {
username: "u".into(),
password: "p".into(),
},
}
}
#[test]
fn validate_rejects_flag_shaped_image() {
let exec = DockerExecutor::new();
let err = exec
.validate(&rspec("ok", "--privileged"))
.expect_err("must reject flag-shaped image");
assert!(matches!(err, ProvisionError::Incompatible(_)));
}
#[test]
fn validate_rejects_flag_shaped_name() {
let exec = DockerExecutor::new();
let err = exec
.validate(&rspec("--rm", "ubuntu:24.04"))
.expect_err("must reject flag-shaped name");
assert!(matches!(err, ProvisionError::Incompatible(_)));
}
#[test]
fn validate_accepts_normal_spec() {
let exec = DockerExecutor::new();
assert!(exec.validate(&rspec("cirun-r1", "ubuntu:24.04")).is_ok());
}
#[test]
fn not_found_classifier_matches_macos_docker_desktop_form() {
assert!(is_docker_not_found(
"Error: No such object: cirun-aktech--demo-25b350fc9b\n"
));
}
#[test]
fn not_found_classifier_matches_linux_docker_ce_form() {
assert!(is_docker_not_found(
"\nerror: no such object: cirun-aktech--demo-25b350fc9b\n"
));
}
#[test]
fn not_found_classifier_matches_no_such_container_phrasing() {
assert!(is_docker_not_found("Error: No such container: foo"));
assert!(is_docker_not_found("error: no such container: foo"));
}
#[test]
fn not_found_classifier_rejects_unrelated_errors() {
assert!(!is_docker_not_found(
"Cannot connect to the Docker daemon at unix:///var/run/docker.sock"
));
assert!(!is_docker_not_found(
"permission denied while trying to connect"
));
}
}