use crate::{Boot, BootSpec, Error, ImageSource, Lifecycle, Machine, PowerState, Result};
use std::time::{Duration, Instant};
#[cfg(feature = "backend-oci")]
use std::sync::{Arc, Mutex};
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum Readiness {
#[default]
Running,
LogMatch(String),
}
#[cfg_attr(not(feature = "backend-oci"), allow(dead_code))]
fn poll_until_ready(
label: &str,
timeout: Duration,
interval: Duration,
mut check: impl FnMut() -> Result<bool>,
) -> Result<()> {
let deadline = Instant::now() + timeout;
loop {
if check()? {
return Ok(());
}
let now = Instant::now();
if now >= deadline {
return Err(Error::Backend(format!(
"container `{label}` not ready within {timeout:?}"
)));
}
let remaining = deadline.saturating_duration_since(now);
std::thread::sleep(interval.min(remaining));
}
}
#[derive(Default, Clone)]
pub struct ContainerBoot {
#[cfg(feature = "backend-oci")]
engine: Arc<Mutex<Option<Arc<Engine>>>>,
}
impl std::fmt::Debug for ContainerBoot {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ContainerBoot").finish_non_exhaustive()
}
}
impl ContainerBoot {
pub fn new() -> Self {
Self::default()
}
pub fn image_ref<'a>(&self, spec: &'a BootSpec) -> Result<&'a str> {
match &spec.image {
ImageSource::OciImage(r) => Ok(r.as_str()),
other => Err(Error::Spec(format!(
"container backend needs an OCI image, got {other:?}"
))),
}
}
#[cfg_attr(not(feature = "backend-oci"), allow(dead_code))]
fn container_name(spec: &BootSpec) -> String {
format!("draupnir-{}", spec.name)
}
pub fn wait_ready(
&self,
machine: &Machine,
readiness: &Readiness,
timeout: Duration,
) -> Result<()> {
#[cfg(feature = "backend-oci")]
{
let engine = self.engine()?;
let name = machine.id.clone();
poll_until_ready(&name, timeout, Duration::from_millis(200), || {
match readiness {
Readiness::Running => Ok(matches!(engine.power_state(&name), PowerState::On)),
Readiness::LogMatch(needle) => engine.log_contains(&name, needle),
}
})
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = (machine, readiness, timeout);
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature".into(),
))
}
}
#[cfg(feature = "backend-oci")]
fn engine(&self) -> Result<Arc<Engine>> {
let mut guard = self.engine.lock().unwrap();
if let Some(e) = guard.as_ref() {
return Ok(Arc::clone(e));
}
let e = Arc::new(Engine::connect()?);
*guard = Some(Arc::clone(&e));
Ok(e)
}
}
impl Boot for ContainerBoot {
fn boot(&self, spec: &BootSpec) -> Result<Machine> {
spec.validate()?;
let image = self.image_ref(spec)?;
#[cfg(feature = "backend-oci")]
{
let name = Self::container_name(spec);
let env: Vec<String> = spec.env.iter().map(|(k, v)| format!("{k}={v}")).collect();
self.engine()?.create_and_start(image, &name, &env, &spec.cmd, &spec.ports)?;
Ok(Machine::started(name, spec))
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = image;
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
}
impl Lifecycle for ContainerBoot {
fn power_on(&self, machine: &Machine) -> Result<()> {
#[cfg(feature = "backend-oci")]
{
self.engine()?.start(&machine.id)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
Err(Error::Unsupported("container backend needs the `backend-oci` feature".into()))
}
}
fn power_off(&self, machine: &Machine) -> Result<()> {
#[cfg(feature = "backend-oci")]
{
self.engine()?.stop(&machine.id);
Ok(())
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
Err(Error::Unsupported("container backend needs the `backend-oci` feature".into()))
}
}
fn status(&self, machine: &Machine) -> Result<PowerState> {
#[cfg(feature = "backend-oci")]
{
Ok(self.engine()?.power_state(&machine.id))
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
Err(Error::Unsupported("container backend needs the `backend-oci` feature".into()))
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ContainerState {
Running,
Exited(i64),
Gone,
}
pub trait ContainerControl {
fn container_state(&self, machine: &Machine) -> ContainerState;
fn drain_logs(&self, machine: &Machine) -> Vec<String>;
fn drain_logs_split(&self, machine: &Machine) -> (Vec<String>, Vec<String>) {
(self.drain_logs(machine), Vec::new())
}
fn stop(&self, machine: &Machine);
}
impl ContainerControl for ContainerBoot {
fn container_state(&self, machine: &Machine) -> ContainerState {
#[cfg(feature = "backend-oci")]
{
match self.engine() {
Ok(e) => e.container_state(&machine.id),
Err(_) => ContainerState::Gone,
}
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
ContainerState::Gone
}
}
fn drain_logs(&self, machine: &Machine) -> Vec<String> {
#[cfg(feature = "backend-oci")]
{
match self.engine() {
Ok(e) => e.drain_logs(&machine.id),
Err(_) => Vec::new(),
}
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
Vec::new()
}
}
fn drain_logs_split(&self, machine: &Machine) -> (Vec<String>, Vec<String>) {
#[cfg(feature = "backend-oci")]
{
match self.engine() {
Ok(e) => e.drain_logs_split(&machine.id),
Err(_) => (Vec::new(), Vec::new()),
}
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
(Vec::new(), Vec::new())
}
}
fn stop(&self, machine: &Machine) {
#[cfg(feature = "backend-oci")]
{
if let Ok(e) = self.engine() {
e.stop(&machine.id);
}
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct RunOutcome {
pub exit_code: Option<i64>,
pub stdout: Vec<String>,
pub stderr: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RunOptions {
pub timeout: Option<Duration>,
pub poll_interval: Duration,
}
impl Default for RunOptions {
fn default() -> Self {
Self { timeout: None, poll_interval: Duration::from_millis(200) }
}
}
impl RunOptions {
pub fn poll_every(poll_interval: Duration) -> Self {
Self { timeout: None, poll_interval }
}
pub fn bounded(timeout: Duration, poll_interval: Duration) -> Self {
Self { timeout: Some(timeout), poll_interval }
}
}
pub fn run_to_completion<B>(backend: &B, spec: &BootSpec, opts: &RunOptions) -> Result<RunOutcome>
where
B: Boot + ContainerControl,
{
let machine = backend.boot(spec)?;
let mut stdout: Vec<String> = Vec::new();
let mut stderr: Vec<String> = Vec::new();
let deadline = opts.timeout.map(|t| Instant::now() + t);
let exit_code = loop {
let (mut out, mut err) = backend.drain_logs_split(&machine);
stdout.append(&mut out);
stderr.append(&mut err);
match backend.container_state(&machine) {
ContainerState::Exited(code) => break Some(code),
ContainerState::Gone => break None,
ContainerState::Running => {}
}
if let Some(dl) = deadline {
if Instant::now() >= dl {
backend.stop(&machine);
return Err(Error::Backend(format!(
"container `{}` did not run to completion within {:?}",
machine.id,
opts.timeout.unwrap()
)));
}
}
std::thread::sleep(opts.poll_interval);
};
let (mut out, mut err) = backend.drain_logs_split(&machine);
stdout.append(&mut out);
stderr.append(&mut err);
backend.stop(&machine);
Ok(RunOutcome { exit_code, stdout, stderr })
}
impl ContainerBoot {
pub fn run_to_completion(&self, spec: &BootSpec, opts: &RunOptions) -> Result<RunOutcome> {
#[cfg(feature = "backend-oci")]
{
run_to_completion(self, spec, opts)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = (spec, opts);
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
}
#[cfg(feature = "backend-oci")]
type LogBuf = Arc<Mutex<Vec<(bool, String)>>>;
#[cfg(feature = "backend-oci")]
struct Engine {
docker: bollard::Docker,
rt: tokio::runtime::Runtime,
logs: Mutex<std::collections::HashMap<String, LogBuf>>,
}
#[cfg(feature = "backend-oci")]
impl Engine {
fn socket_url() -> String {
if let Ok(h) = std::env::var("DOCKER_HOST") {
return h;
}
let xdg = std::env::var("XDG_RUNTIME_DIR").unwrap_or_else(|_| "/run/user/1000".into());
format!("unix://{xdg}/podman/podman.sock")
}
fn connect() -> Result<Self> {
let url = Self::socket_url();
let path = url.strip_prefix("unix://").unwrap_or(&url);
if url.starts_with("unix://") && !std::path::Path::new(path).exists() {
return Err(Error::Backend(format!(
"podman/Docker API socket not found at {path} (enable with \
`systemctl --user enable --now podman.socket`, or point DOCKER_HOST at a running socket)"
)));
}
let docker = bollard::Docker::connect_with_unix(&url, 120, bollard::API_DEFAULT_VERSION)
.map_err(|e| Error::Backend(format!("connect container socket {url}: {e}")))?;
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(2)
.enable_all()
.build()
.map_err(|e| Error::Backend(format!("build tokio runtime for bollard: {e}")))?;
Ok(Engine { docker, rt, logs: Mutex::new(std::collections::HashMap::new()) })
}
pub fn create_body(image: &str, env: &[String], cmd: &[String], ports: &[u16]) -> bollard::models::ContainerCreateBody {
use bollard::models::{ContainerCreateBody, HostConfig, PortBinding};
use std::collections::HashMap;
let mut exposed: Vec<String> = Vec::new();
let mut bindings: HashMap<String, Option<Vec<PortBinding>>> = HashMap::new();
for p in ports {
let key = format!("{p}/tcp");
exposed.push(key.clone());
bindings.insert(
key,
Some(vec![PortBinding {
host_ip: Some("0.0.0.0".to_string()),
host_port: Some(p.to_string()),
}]),
);
}
let host_config = if bindings.is_empty() {
None
} else {
Some(HostConfig { port_bindings: Some(bindings), ..Default::default() })
};
ContainerCreateBody {
image: Some(image.to_string()),
cmd: if cmd.is_empty() { None } else { Some(cmd.to_vec()) },
env: if env.is_empty() { None } else { Some(env.to_vec()) },
exposed_ports: if exposed.is_empty() { None } else { Some(exposed) },
host_config,
..Default::default()
}
}
fn create_and_start(&self, image: &str, name: &str, env: &[String], cmd: &[String], ports: &[u16]) -> Result<()> {
use bollard::query_parameters::{
CreateContainerOptionsBuilder, CreateImageOptionsBuilder, RemoveContainerOptionsBuilder,
StartContainerOptions,
};
use futures::StreamExt;
let docker = &self.docker;
self.rt.block_on(async {
let _ = docker
.remove_container(name, Some(RemoveContainerOptionsBuilder::new().force(true).build()))
.await;
if docker.inspect_image(image).await.is_err() {
let (repo, tag) = image.rsplit_once(':').unwrap_or((image, "latest"));
let opts = CreateImageOptionsBuilder::new().from_image(repo).tag(tag).build();
let mut pull = docker.create_image(Some(opts), None, None);
while let Some(item) = pull.next().await {
item.map_err(|e| Error::Backend(format!("pull image {image}: {e}")))?;
}
}
let body = Self::create_body(image, env, cmd, ports);
docker
.create_container(Some(CreateContainerOptionsBuilder::new().name(name).build()), body)
.await
.map_err(|e| Error::Backend(format!("create container {name}: {e}")))?;
docker
.start_container(name, None::<StartContainerOptions>)
.await
.map_err(|e| Error::Backend(format!("start container {name}: {e}")))?;
Ok::<(), Error>(())
})?;
let buf: LogBuf = Arc::new(Mutex::new(Vec::new()));
self.logs.lock().unwrap().insert(name.to_string(), Arc::clone(&buf));
let docker = self.docker.clone();
let name_owned = name.to_string();
self.rt.spawn(async move {
use bollard::container::LogOutput;
use bollard::query_parameters::LogsOptionsBuilder;
let mut stream = docker.logs(
&name_owned,
Some(LogsOptionsBuilder::new().follow(true).stdout(true).stderr(true).build()),
);
while let Some(item) = stream.next().await {
match item {
Ok(out) => {
let is_err = matches!(out, LogOutput::StdErr { .. });
let line = LogOutput::to_string(&out);
let line = line.trim_end_matches(['\n', '\r']).to_string();
if !line.is_empty() {
buf.lock().unwrap().push((is_err, line));
}
}
Err(_) => break,
}
}
});
Ok(())
}
fn container_state(&self, name: &str) -> ContainerState {
use bollard::query_parameters::InspectContainerOptions;
self.rt.block_on(async {
match self.docker.inspect_container(name, None::<InspectContainerOptions>).await {
Ok(info) => {
let state = info.state;
let running = state.as_ref().and_then(|s| s.running).unwrap_or(false);
if running {
ContainerState::Running
} else {
ContainerState::Exited(state.and_then(|s| s.exit_code).unwrap_or(0))
}
}
Err(_) => ContainerState::Gone,
}
})
}
fn drain_logs(&self, name: &str) -> Vec<String> {
match self.logs.lock().unwrap().get(name) {
Some(buf) => std::mem::take(&mut *buf.lock().unwrap())
.into_iter()
.map(|(_, line)| line)
.collect(),
None => Vec::new(),
}
}
fn drain_logs_split(&self, name: &str) -> (Vec<String>, Vec<String>) {
match self.logs.lock().unwrap().get(name) {
Some(buf) => {
let mut out = Vec::new();
let mut err = Vec::new();
for (is_err, line) in std::mem::take(&mut *buf.lock().unwrap()) {
if is_err {
err.push(line);
} else {
out.push(line);
}
}
(out, err)
}
None => (Vec::new(), Vec::new()),
}
}
fn start(&self, name: &str) -> Result<()> {
use bollard::query_parameters::StartContainerOptions;
self.rt.block_on(async {
self.docker
.start_container(name, None::<StartContainerOptions>)
.await
.map_err(|e| Error::Backend(format!("start container {name}: {e}")))
})
}
fn stop(&self, name: &str) {
use bollard::query_parameters::{RemoveContainerOptionsBuilder, StopContainerOptions};
self.rt.block_on(async {
let _ = self
.docker
.stop_container(name, None::<StopContainerOptions>)
.await;
let _ = self
.docker
.remove_container(name, Some(RemoveContainerOptionsBuilder::new().force(true).build()))
.await;
});
self.logs.lock().unwrap().remove(name);
}
fn power_state(&self, name: &str) -> PowerState {
use bollard::query_parameters::InspectContainerOptions;
self.rt.block_on(async {
match self
.docker
.inspect_container(name, None::<InspectContainerOptions>)
.await
{
Ok(info) => {
let running = info.state.as_ref().and_then(|s| s.running).unwrap_or(false);
if running {
PowerState::On
} else {
PowerState::Off
}
}
Err(_) => PowerState::Off,
}
})
}
fn log_contains(&self, name: &str, needle: &str) -> Result<bool> {
use bollard::query_parameters::LogsOptionsBuilder;
use futures::StreamExt;
self.rt.block_on(async {
let opts = LogsOptionsBuilder::new().stdout(true).stderr(true).build();
let mut stream = self.docker.logs(name, Some(opts));
let mut buf = String::new();
while let Some(item) = stream.next().await {
match item {
Ok(chunk) => buf.push_str(&String::from_utf8_lossy(&chunk.into_bytes())),
Err(e) => {
return Err(Error::Backend(format!("read logs for {name}: {e}")))
}
}
}
Ok(buf.contains(needle))
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn image_ref_extracts_the_oci_reference() {
let spec = BootSpec::container("cache", "docker.io/library/redis:7");
assert_eq!(ContainerBoot::new().image_ref(&spec).unwrap(), "docker.io/library/redis:7");
}
#[test]
fn image_ref_rejects_a_non_oci_image() {
let mut spec = BootSpec::container("bad", "redis:7");
spec.image = ImageSource::Iso("/boot.iso".into());
assert!(matches!(ContainerBoot::new().image_ref(&spec), Err(Error::Spec(_))));
}
#[test]
fn container_name_is_derived_from_the_spec_name() {
let spec = BootSpec::container("cache", "redis:7");
assert_eq!(ContainerBoot::container_name(&spec), "draupnir-cache");
}
#[test]
fn wait_ready_returns_a_clear_timeout_error_when_never_ready() {
let err = poll_until_ready(
"cache",
Duration::from_millis(40),
Duration::from_millis(5),
|| Ok(false),
)
.unwrap_err();
match err {
Error::Backend(m) => {
assert!(m.contains("cache"), "names the instance: {m}");
assert!(m.contains("not ready"), "says it wasn't ready: {m}");
}
other => panic!("expected Error::Backend, got {other:?}"),
}
}
#[test]
fn wait_ready_returns_ok_as_soon_as_the_probe_reports_ready() {
let mut n = 0;
let r = poll_until_ready("cache", Duration::from_secs(5), Duration::from_millis(1), || {
n += 1;
Ok(n >= 3)
});
assert!(r.is_ok());
assert_eq!(n, 3);
}
#[test]
fn wait_ready_fails_fast_on_a_backend_error() {
let r = poll_until_ready("cache", Duration::from_secs(5), Duration::from_millis(1), || {
Err(Error::Backend("socket vanished".into()))
});
assert!(matches!(r, Err(Error::Backend(m)) if m.contains("socket vanished")));
}
#[test]
fn readiness_defaults_to_running() {
assert_eq!(Readiness::default(), Readiness::Running);
}
#[test]
fn container_boot_carries_cmd_and_ports_through_the_spec() {
let spec = BootSpec::container("web", "docker.io/library/nginx:alpine")
.with_cmd(["nginx", "-g", "daemon off;"])
.with_port(8080)
.with_port(8443)
.with_env("TZ", "UTC");
assert_eq!(spec.cmd, vec!["nginx", "-g", "daemon off;"]);
assert_eq!(spec.ports, vec![8080, 8443]);
assert_eq!(spec.env.get("TZ").map(String::as_str), Some("UTC"));
spec.validate().unwrap();
}
#[test]
fn container_control_is_an_honest_noop_without_the_engine() {
let boot = ContainerBoot::new();
let m = Machine::started("draupnir-x", &BootSpec::container("x", "redis:7"));
assert_eq!(boot.container_state(&m), ContainerState::Gone);
assert!(boot.drain_logs(&m).is_empty());
boot.stop(&m); }
#[cfg(feature = "backend-oci")]
#[test]
fn create_body_carries_env_cmd_and_port_bindings() {
let body = Engine::create_body(
"app:1",
&["A=1".to_string(), "B=2".to_string()],
&["/bin/app".to_string(), "--serve".to_string()],
&[8080],
);
assert_eq!(body.image.as_deref(), Some("app:1"));
assert_eq!(body.cmd, Some(vec!["/bin/app".to_string(), "--serve".to_string()]));
assert_eq!(body.env, Some(vec!["A=1".to_string(), "B=2".to_string()]));
let exposed = body.exposed_ports.expect("exposed ports set");
assert!(exposed.iter().any(|s| s == "8080/tcp"));
let hc = body.host_config.expect("host config");
let b = hc.port_bindings.expect("bindings").get("8080/tcp").and_then(|v| v.clone()).expect("8080");
assert_eq!(b[0].host_port.as_deref(), Some("8080"));
}
use std::cell::RefCell;
#[derive(Default)]
struct ScriptedBackend {
booted: RefCell<Vec<String>>,
stopped: RefCell<Vec<String>>,
states: RefCell<std::collections::VecDeque<ContainerState>>,
log_batches: RefCell<std::collections::VecDeque<(Vec<String>, Vec<String>)>>,
}
impl ScriptedBackend {
fn with_states(states: Vec<ContainerState>) -> Self {
Self { states: RefCell::new(states.into()), ..Default::default() }
}
fn push_logs(&self, out: &[&str], err: &[&str]) {
self.log_batches.borrow_mut().push_back((
out.iter().map(|s| s.to_string()).collect(),
err.iter().map(|s| s.to_string()).collect(),
));
}
}
impl Boot for ScriptedBackend {
fn boot(&self, spec: &BootSpec) -> Result<Machine> {
self.booted.borrow_mut().push(spec.name.clone());
Ok(Machine::started(format!("draupnir-{}", spec.name), spec))
}
}
impl ContainerControl for ScriptedBackend {
fn container_state(&self, _m: &Machine) -> ContainerState {
let mut q = self.states.borrow_mut();
if q.len() > 1 {
q.pop_front().unwrap()
} else {
q.front().cloned().unwrap_or(ContainerState::Gone)
}
}
fn drain_logs(&self, _m: &Machine) -> Vec<String> {
let (mut out, mut err) = self.log_batches.borrow_mut().pop_front().unwrap_or_default();
out.append(&mut err);
out
}
fn drain_logs_split(&self, _m: &Machine) -> (Vec<String>, Vec<String>) {
self.log_batches.borrow_mut().pop_front().unwrap_or_default()
}
fn stop(&self, m: &Machine) {
self.stopped.borrow_mut().push(m.id.clone());
}
}
fn fast_opts() -> RunOptions {
RunOptions::poll_every(Duration::from_millis(1))
}
#[test]
fn run_to_completion_starts_waits_collects_split_logs_and_exit_code() {
let backend = ScriptedBackend::with_states(vec![
ContainerState::Running,
ContainerState::Running,
ContainerState::Exited(0),
]);
backend.push_logs(&["booting"], &[]); backend.push_logs(&["serving"], &["a warning"]); backend.push_logs(&["bye"], &[]);
let spec = BootSpec::container("job", "docker.io/library/busybox:latest");
let out = run_to_completion(&backend, &spec, &fast_opts()).unwrap();
assert_eq!(out.exit_code, Some(0));
assert_eq!(out.stdout, vec!["booting", "serving", "bye"]);
assert_eq!(out.stderr, vec!["a warning"]);
assert_eq!(backend.booted.borrow().as_slice(), &["job".to_string()]);
assert_eq!(backend.stopped.borrow().as_slice(), &["draupnir-job".to_string()]);
}
#[test]
fn run_to_completion_surfaces_a_nonzero_exit_as_ok_not_err() {
let backend = ScriptedBackend::with_states(vec![ContainerState::Exited(137)]);
backend.push_logs(&[], &["oom-killed"]);
let spec = BootSpec::container("crash", "img:1");
let out = run_to_completion(&backend, &spec, &fast_opts()).unwrap();
assert_eq!(out.exit_code, Some(137));
assert_eq!(out.stderr, vec!["oom-killed"]);
}
#[test]
fn run_to_completion_reports_none_when_the_container_is_gone() {
let backend = ScriptedBackend::with_states(vec![ContainerState::Gone]);
let spec = BootSpec::container("vanished", "img:1");
let out = run_to_completion(&backend, &spec, &fast_opts()).unwrap();
assert_eq!(out.exit_code, None);
assert_eq!(backend.stopped.borrow().len(), 1, "still removed on the way out");
}
#[test]
fn run_to_completion_times_out_with_a_clear_error_and_stops_the_container() {
let backend = ScriptedBackend::with_states(vec![ContainerState::Running]);
let spec = BootSpec::container("hang", "img:1");
let opts = RunOptions::bounded(Duration::from_millis(20), Duration::from_millis(2));
let err = run_to_completion(&backend, &spec, &opts).unwrap_err();
match err {
Error::Backend(m) => {
assert!(m.contains("draupnir-hang"), "names the instance: {m}");
assert!(m.contains("run to completion"), "says what timed out: {m}");
}
other => panic!("expected Error::Backend, got {other:?}"),
}
assert_eq!(backend.stopped.borrow().len(), 1, "container stopped on timeout");
}
#[test]
fn run_to_completion_propagates_a_boot_failure_without_polling() {
struct FailBoot;
impl Boot for FailBoot {
fn boot(&self, _spec: &BootSpec) -> Result<Machine> {
Err(Error::Backend("no socket".into()))
}
}
impl ContainerControl for FailBoot {
fn container_state(&self, _m: &Machine) -> ContainerState {
panic!("must not poll after a boot failure")
}
fn drain_logs(&self, _m: &Machine) -> Vec<String> {
Vec::new()
}
fn stop(&self, _m: &Machine) {
panic!("must not stop after a boot failure")
}
}
let spec = BootSpec::container("x", "img:1");
let err = run_to_completion(&FailBoot, &spec, &fast_opts()).unwrap_err();
assert!(matches!(err, Error::Backend(m) if m.contains("no socket")));
}
#[test]
fn default_drain_logs_split_routes_combined_logs_to_stdout() {
struct CombinedOnly;
impl ContainerControl for CombinedOnly {
fn container_state(&self, _m: &Machine) -> ContainerState {
ContainerState::Gone
}
fn drain_logs(&self, _m: &Machine) -> Vec<String> {
vec!["one".into(), "two".into()]
}
fn stop(&self, _m: &Machine) {}
}
let m = Machine::started("draupnir-x", &BootSpec::container("x", "img:1"));
let (out, err) = CombinedOnly.drain_logs_split(&m);
assert_eq!(out, vec!["one", "two"]);
assert!(err.is_empty());
}
#[cfg(feature = "backend-oci")]
#[test]
fn create_body_empty_spec_is_bare() {
let body = Engine::create_body("scratch", &[], &[], &[]);
assert_eq!(body.image.as_deref(), Some("scratch"));
assert!(body.cmd.is_none());
assert!(body.env.is_none());
assert!(body.exposed_ports.is_none());
assert!(body.host_config.is_none());
}
}