use crate::{Boot, BootSpec, Error, ImageSource, Lifecycle, Machine, PowerState, Result};
use std::path::Path;
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();
let outcome =
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),
},
);
crate::functional_status(
"draupnir/container",
"wait_ready",
outcome.is_ok(),
&match &outcome {
Ok(()) => format!("container `{name}` reached {readiness:?}"),
Err(e) => format!("container `{name}` never became ready: {e}"),
},
);
outcome
}
#[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_with_binds_and_net(
image,
&name,
&env,
&spec.cmd,
&spec.ports,
&spec.port_maps,
&[],
spec.net.oci_value(),
spec.cpus,
spec.mem_limit_mb,
&spec.security_opt,
)?;
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,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct ExecOutcome {
pub exit_code: Option<i64>,
pub stdout: Vec<String>,
pub stderr: Vec<String>,
}
pub fn exec_argv(id: &str, argv: &[&str]) -> Vec<String> {
let mut command = Vec::with_capacity(argv.len() + 2);
command.push("exec".to_string());
command.push(id.to_string());
command.extend(argv.iter().map(|a| a.to_string()));
command
}
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);
fn list_containers(&self, prefix: &str) -> Result<Vec<(String, ContainerState)>> {
let _ = prefix;
Err(Error::Unsupported(
"this backend cannot list containers (needs the OCI engine)".into(),
))
}
fn container_env(&self, machine: &Machine) -> Result<Vec<String>> {
let _ = machine;
Err(Error::Unsupported(
"this backend cannot inspect a container's environment (needs the OCI engine)".into(),
))
}
fn exec(&self, machine: &Machine, argv: &[&str]) -> Result<ExecOutcome> {
if machine.backend != crate::Backend::Container {
return Err(Error::Spec(format!(
"exec is container-only; a {:?} machine has no container to exec into",
machine.backend
)));
}
if argv.is_empty() {
return Err(Error::Spec(
"a container exec needs a non-empty argv (nothing to run)".into(),
));
}
let command = exec_argv(&machine.id, argv);
self.exec_command(&command)
}
fn exec_command(&self, command: &[String]) -> Result<ExecOutcome> {
let _ = command;
Err(Error::Unsupported(
"this backend cannot exec into a container (needs the OCI engine)".into(),
))
}
}
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;
}
}
fn list_containers(&self, prefix: &str) -> Result<Vec<(String, ContainerState)>> {
#[cfg(feature = "backend-oci")]
{
self.engine()?.list_containers(prefix)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = prefix;
Err(Error::Unsupported(
"this build has no OCI engine, so it cannot list containers".into(),
))
}
}
fn container_env(&self, machine: &Machine) -> Result<Vec<String>> {
#[cfg(feature = "backend-oci")]
{
self.engine()?.container_env(&machine.id)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = machine;
Err(Error::Unsupported(
"this build has no OCI engine, so it cannot inspect a container".into(),
))
}
}
fn exec_command(&self, command: &[String]) -> Result<ExecOutcome> {
#[cfg(feature = "backend-oci")]
{
let id = &command[1];
let argv: Vec<String> = command[2..].to_vec();
self.engine()?.exec(id, &argv)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = command;
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
}
#[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)?;
drive_to_completion(backend, &machine, opts)
}
pub fn drive_to_completion<B>(
backend: &B,
machine: &Machine,
opts: &RunOptions,
) -> Result<RunOutcome>
where
B: ContainerControl,
{
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(),
))
}
}
pub fn boot_with_binds(&self, spec: &BootSpec, binds: &[String]) -> 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_with_binds_and_net(
image,
&name,
&env,
&spec.cmd,
&spec.ports,
&spec.port_maps,
binds,
spec.net.oci_value(),
spec.cpus,
spec.mem_limit_mb,
&spec.security_opt,
)?;
Ok(Machine::started(name, spec))
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = (image, binds);
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
pub fn run_to_completion_with_binds(
&self,
spec: &BootSpec,
opts: &RunOptions,
binds: &[String],
) -> Result<RunOutcome> {
#[cfg(feature = "backend-oci")]
{
let machine = self.boot_with_binds(spec, binds)?;
drive_to_completion(self, &machine, opts)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = (spec, opts, binds);
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
pub fn build_image(
&self,
context_dir: &Path,
containerfile: &str,
tag: &str,
) -> Result<String> {
#[cfg(feature = "backend-oci")]
{
self.engine()?.build_image(context_dir, containerfile, tag)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = (context_dir, containerfile, tag);
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
pub fn image_present(&self, image: &str) -> Result<bool> {
#[cfg(feature = "backend-oci")]
{
image_present_on(self.engine()?.as_ref(), image)
}
#[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(),
))
}
}
pub fn extract_path(&self, image: &str, container_path: &str, host_dest: &Path) -> Result<()> {
#[cfg(feature = "backend-oci")]
{
self.engine()?
.extract_path(image, container_path, host_dest)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = (image, container_path, host_dest);
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
pub fn export_image(&self, image: &str, dest: &Path) -> Result<ImageArchive> {
#[cfg(feature = "backend-oci")]
{
export_image_on(self.engine()?.as_ref(), image, dest)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = (image, dest);
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
pub fn import_image(&self, src: &Path) -> Result<ImageArchive> {
#[cfg(feature = "backend-oci")]
{
import_image_on(self.engine()?.as_ref(), src)
}
#[cfg(not(feature = "backend-oci"))]
{
let _ = src;
Err(Error::Unsupported(
"container backend needs the `backend-oci` feature (drives the podman/Docker REST API via bollard)"
.into(),
))
}
}
}
#[cfg(feature = "backend-oci")]
fn context_tar(context_dir: &Path) -> Result<Vec<u8>> {
let mut buf: Vec<u8> = Vec::new();
{
let mut builder = tar::Builder::new(&mut buf);
builder.append_dir_all(".", context_dir).map_err(|e| {
Error::Backend(format!("tar build context {}: {e}", context_dir.display()))
})?;
builder
.finish()
.map_err(|e| Error::Backend(format!("finish build-context tar: {e}")))?;
}
Ok(buf)
}
#[cfg(feature = "backend-oci")]
trait ImageDaemon {
fn build_stream(
&self,
context: Vec<u8>,
containerfile: &str,
tag: &str,
) -> Vec<Result<Option<String>>>;
fn download_path_tar(&self, image: &str, container_path: &str) -> Result<Vec<u8>>;
fn inspect_image(&self, image: &str) -> ImagePresence;
fn export_to(&self, image: &str, sink: &mut dyn std::io::Write) -> Result<u64>;
fn load_archive(&self, src: &Path) -> Result<()>;
}
#[cfg(feature = "backend-oci")]
#[derive(Debug, Clone, PartialEq, Eq)]
enum ImagePresence {
Present,
Absent,
Failed(String),
}
#[cfg(feature = "backend-oci")]
fn image_present_on<D: ImageDaemon>(daemon: &D, image: &str) -> Result<bool> {
match daemon.inspect_image(image) {
ImagePresence::Present => Ok(true),
ImagePresence::Absent => Ok(false),
ImagePresence::Failed(msg) => Err(Error::Backend(format!("inspect image {image}: {msg}"))),
}
}
#[cfg(feature = "backend-oci")]
fn build_image_on<D: ImageDaemon>(
daemon: &D,
context_dir: &Path,
containerfile: &str,
tag: &str,
) -> Result<String> {
let tar = context_tar(context_dir)?;
for step in daemon.build_stream(tar, containerfile, tag) {
if let Some(msg) = step? {
return Err(Error::Backend(format!("build image {tag}: {msg}")));
}
}
Ok(tag.to_string())
}
#[cfg(feature = "backend-oci")]
fn extract_path_on<D: ImageDaemon>(
daemon: &D,
image: &str,
container_path: &str,
host_dest: &Path,
) -> Result<()> {
let tar_bytes = daemon.download_path_tar(image, container_path)?;
std::fs::create_dir_all(host_dest)
.map_err(|e| Error::Backend(format!("create extract dest {}: {e}", host_dest.display())))?;
tar::Archive::new(std::io::Cursor::new(tar_bytes))
.unpack(host_dest)
.map_err(|e| {
Error::Backend(format!(
"unpack extracted tar into {}: {e}",
host_dest.display()
))
})?;
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ImageArchive {
pub path: std::path::PathBuf,
pub bytes: u64,
pub tags: Vec<String>,
pub config_digest: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ArchiveFault {
Unreadable(String),
Truncated {
declared: u64,
actual: u64,
},
NotAnImageArchive,
Malformed(String),
}
impl std::fmt::Display for ArchiveFault {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ArchiveFault::Unreadable(m) => write!(f, "unreadable: {m}"),
ArchiveFault::Truncated { declared, actual } => write!(
f,
"truncated: a tar entry declares bytes out to {declared} but the file is only {actual} bytes"
),
ArchiveFault::NotAnImageArchive => {
write!(f, "not an image archive: no manifest.json at the tar root")
}
ArchiveFault::Malformed(m) => write!(f, "malformed manifest.json: {m}"),
}
}
}
#[cfg(feature = "backend-oci")]
const IMPORT_CHUNK: usize = 4 * 1024 * 1024;
#[cfg(feature = "backend-oci")]
const MAX_MANIFEST_BYTES: u64 = 8 * 1024 * 1024;
#[cfg(feature = "backend-oci")]
pub fn inspect_image_archive(path: &Path) -> Result<ImageArchive> {
let file = std::fs::File::open(path)
.map_err(|e| archive_error(path, &ArchiveFault::Unreadable(e.to_string())))?;
let len = file
.metadata()
.map_err(|e| archive_error(path, &ArchiveFault::Unreadable(e.to_string())))?
.len();
let (tags, config_digest) =
scan_image_archive(file, len).map_err(|fault| archive_error(path, &fault))?;
Ok(ImageArchive {
path: path.to_path_buf(),
bytes: len,
tags,
config_digest,
})
}
#[cfg(not(feature = "backend-oci"))]
pub fn inspect_image_archive(path: &Path) -> Result<ImageArchive> {
let _ = path;
Err(Error::Unsupported(
"reading an image archive needs the `backend-oci` feature (tar + serde_json)".into(),
))
}
#[cfg(feature = "backend-oci")]
fn archive_error(path: &Path, fault: &ArchiveFault) -> Error {
Error::Backend(format!("image archive {}: {fault}", path.display()))
}
#[cfg(feature = "backend-oci")]
fn scan_image_archive<R: std::io::Read + std::io::Seek>(
reader: R,
len: u64,
) -> std::result::Result<(Vec<String>, String), ArchiveFault> {
let mut archive = tar::Archive::new(reader);
let entries = archive
.entries_with_seek()
.map_err(|e| ArchiveFault::Unreadable(e.to_string()))?;
let mut manifest: Option<Vec<u8>> = None;
let mut declared_end: u64 = 0;
for entry in entries {
let mut entry = entry.map_err(|e| {
ArchiveFault::Unreadable(format!(
"tar entry could not be read (archive is {len} bytes): {e}"
))
})?;
let size = entry.size();
let data_at = entry.raw_file_position();
let padded = size.div_ceil(512).saturating_mul(512);
declared_end = declared_end.max(data_at.saturating_add(padded));
let is_manifest = entry
.path()
.map(|p| p.as_ref() == Path::new("manifest.json"))
.unwrap_or(false);
if is_manifest {
if size > MAX_MANIFEST_BYTES {
return Err(ArchiveFault::Malformed(format!(
"manifest.json declares {size} bytes, over the {MAX_MANIFEST_BYTES}-byte \
sanity bound — this is not a real image manifest"
)));
}
let mut buf = Vec::with_capacity(size as usize);
std::io::Read::read_to_end(&mut entry, &mut buf)
.map_err(|e| ArchiveFault::Unreadable(format!("read manifest.json: {e}")))?;
manifest = Some(buf);
}
}
if declared_end > len {
return Err(ArchiveFault::Truncated {
declared: declared_end,
actual: len,
});
}
let manifest = manifest.ok_or(ArchiveFault::NotAnImageArchive)?;
parse_archive_manifest(&manifest)
}
#[cfg(feature = "backend-oci")]
fn parse_archive_manifest(
bytes: &[u8],
) -> std::result::Result<(Vec<String>, String), ArchiveFault> {
let value: serde_json::Value =
serde_json::from_slice(bytes).map_err(|e| ArchiveFault::Malformed(e.to_string()))?;
let entries = value
.as_array()
.ok_or_else(|| ArchiveFault::Malformed("top level is not a JSON array".into()))?;
if entries.is_empty() {
return Err(ArchiveFault::Malformed(
"manifest.json is an empty array — the archive names no image".into(),
));
}
let mut tags: Vec<String> = Vec::new();
for entry in entries {
if let Some(list) = entry.get("RepoTags").and_then(|t| t.as_array()) {
tags.extend(list.iter().filter_map(|t| t.as_str()).map(str::to_string));
}
}
let config = entries[0]
.get("Config")
.and_then(|c| c.as_str())
.ok_or_else(|| ArchiveFault::Malformed("first entry has no Config".into()))?;
let digest = config
.rsplit('/')
.next()
.unwrap_or(config)
.trim_end_matches(".json")
.to_string();
if digest.is_empty() {
return Err(ArchiveFault::Malformed(format!(
"Config {config:?} yields no digest"
)));
}
Ok((tags, digest))
}
#[cfg(feature = "backend-oci")]
fn export_image_on<D: ImageDaemon>(daemon: &D, image: &str, dest: &Path) -> Result<ImageArchive> {
match daemon.inspect_image(image) {
ImagePresence::Present => {}
ImagePresence::Absent => {
return Err(Error::Backend(format!(
"export image {image}: no such image locally — nothing was written to {}",
dest.display()
)));
}
ImagePresence::Failed(msg) => {
return Err(Error::Backend(format!(
"export image {image}: container engine unreachable: {msg}"
)));
}
}
if let Some(parent) = dest.parent().filter(|p| !p.as_os_str().is_empty()) {
std::fs::create_dir_all(parent).map_err(|e| {
Error::Backend(format!("create export dir {}: {e}", parent.display()))
})?;
}
{
let file = std::fs::File::create(dest).map_err(|e| {
Error::Backend(format!("create export file {}: {e}", dest.display()))
})?;
let mut sink = std::io::BufWriter::new(file);
daemon.export_to(image, &mut sink)?;
std::io::Write::flush(&mut sink).map_err(|e| {
Error::Backend(format!("flush export file {}: {e}", dest.display()))
})?;
}
let archive = inspect_image_archive(dest)?;
if !archive.tags.is_empty() && !archive.tags.iter().any(|t| t == image) {
return Err(Error::Backend(format!(
"export image {image}: the archive at {} carries {:?}, not {image} — refusing to \
record an image under a tag it does not have",
dest.display(),
archive.tags
)));
}
Ok(archive)
}
#[cfg(feature = "backend-oci")]
fn import_image_on<D: ImageDaemon>(daemon: &D, src: &Path) -> Result<ImageArchive> {
let archive = inspect_image_archive(src)?;
daemon.load_archive(src)?;
Ok(archive)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Ownership {
pub user: String,
pub group: String,
pub mode: u32,
}
impl std::fmt::Display for Ownership {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}:{} {:04o}", self.user, self.group, self.mode)
}
}
fn name_for_id(db: &str, id: u32) -> Option<String> {
let text = std::fs::read_to_string(db).ok()?;
for line in text.lines() {
let mut f = line.split(':');
let name = f.next()?;
let _passwd = f.next();
if f.next().and_then(|n| n.trim().parse::<u32>().ok()) == Some(id) {
return Some(name.to_string());
}
}
None
}
fn ownership_of(md: &std::fs::Metadata) -> Ownership {
use std::os::unix::fs::MetadataExt;
let (uid, gid) = (md.uid(), md.gid());
Ownership {
user: name_for_id("/etc/passwd", uid).unwrap_or_else(|| uid.to_string()),
group: name_for_id("/etc/group", gid).unwrap_or_else(|| gid.to_string()),
mode: md.mode() & 0o7777,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SocketProbe {
Reachable,
Absent,
Denied {
blocked_at: String,
owner: Option<Ownership>,
},
Other(String),
}
pub fn probe_socket_path(path: &Path) -> SocketProbe {
match std::fs::metadata(path) {
Ok(_) => SocketProbe::Reachable,
Err(e) => match e.kind() {
std::io::ErrorKind::NotFound => SocketProbe::Absent,
std::io::ErrorKind::PermissionDenied => {
for anc in path.ancestors().skip(1) {
if let Ok(md) = std::fs::metadata(anc) {
return SocketProbe::Denied {
blocked_at: anc.display().to_string(),
owner: Some(ownership_of(&md)),
};
}
}
SocketProbe::Denied {
blocked_at: path.display().to_string(),
owner: None,
}
}
_ => SocketProbe::Other(e.to_string()),
},
}
}
pub fn is_user_scope_socket(path: &Path) -> bool {
path.starts_with("/run/user/") || path.starts_with("/var/run/user/")
}
fn enable_hint(path: &Path) -> &'static str {
if is_user_scope_socket(path) {
"systemctl --user enable --now podman.socket"
} else {
"sudo systemctl enable --now podman.socket"
}
}
pub fn socket_unreachable_reason(path: &Path, probe: &SocketProbe) -> Option<String> {
let p = path.display();
match probe {
SocketProbe::Reachable => None,
SocketProbe::Absent => Some(format!(
"podman/Docker API socket not found at {p} — nothing exists at that path \
(enable with `{}`, or point DOCKER_HOST at a running socket)",
enable_hint(path)
)),
SocketProbe::Denied { blocked_at, owner } => {
let mut m = format!(
"podman/Docker API socket at {p} EXISTS but this process may not reach it \
(permission denied) — do NOT enable the unit, it is already there. \
Refused at {blocked_at}"
);
match owner {
Some(o) => {
m.push_str(&format!(" ({o})."));
if o.group == "root" || o.group == "0" {
m.push_str(
" Its group is root, so joining a group CANNOT fix this: \
either the host gives the socket a non-root group \
(a `SocketGroup=` drop-in on podman.socket, a HOST decision) \
or the service runs as root. Adding a service account to the \
root group is not a fix.",
);
} else {
m.push_str(&format!(
" Grant access by making the service account a member of `{}` — \
for a systemd unit that is `SupplementaryGroups={}`, \
which skidbladnir renders from the service's own declaration.",
o.group, o.group
));
}
}
None => m.push('.'),
}
Some(m)
}
SocketProbe::Other(e) => Some(format!(
"podman/Docker API socket at {p} could not be examined: {e}"
)),
}
}
pub fn socket_unreachable(path: &Path) -> Option<String> {
socket_unreachable_reason(path, &probe_socket_path(path))
}
pub fn rootless_socket_path(uid: u32) -> std::path::PathBuf {
std::path::PathBuf::from(format!("/run/user/{uid}/podman/podman.sock"))
}
pub fn current_uid() -> Option<u32> {
let status = std::fs::read_to_string("/proc/self/status").ok()?;
parse_uid_line(&status)
}
fn parse_uid_line(status: &str) -> Option<u32> {
for line in status.lines() {
if let Some(rest) = line.strip_prefix("Uid:") {
return rest.split_whitespace().next()?.parse().ok();
}
}
None
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MigrateVerdict {
NotOwed(String),
Owed(String),
}
impl MigrateVerdict {
pub fn detail(&self) -> &str {
match self {
MigrateVerdict::NotOwed(s) | MigrateVerdict::Owed(s) => s,
}
}
}
pub fn migrate_verdict(account: &str, storage_exists: bool) -> MigrateVerdict {
if !storage_exists {
MigrateVerdict::NotOwed(format!(
"podman has never run as `{account}` (no libpod storage), so there are no \
stored user-namespace mappings to refresh — the first run adopts the \
/etc/subuid grant by itself. `podman system migrate` is not needed."
))
} else {
MigrateVerdict::Owed(format!(
"podman has already run as `{account}` and holds stored user-namespace \
mappings. If the /etc/subuid grant was made AFTER that, they are stale and \
every container keeps the old single-id mapping. Only `podman system \
migrate`, run AS `{account}`, refreshes them — and draupnir will not run \
it: there is no libpod REST endpoint for it, so running it means a podman \
subprocess, which this crate does not do. Run it once, as the service \
account, or delete the account's containers and let them be re-created."
))
}
}
pub fn libpod_storage_dir(home: &str) -> std::path::PathBuf {
std::path::Path::new(home).join(".local/share/containers/storage/libpod")
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RootlessReadiness {
pub uid: u32,
pub socket: std::path::PathBuf,
pub probe: SocketProbe,
pub obstacle: Option<String>,
}
impl RootlessReadiness {
pub fn is_ready(&self) -> bool {
self.obstacle.is_none()
}
}
pub fn rootless_readiness(uid: u32) -> RootlessReadiness {
let socket = rootless_socket_path(uid);
let probe = probe_socket_path(&socket);
let runtime_dir = std::path::PathBuf::from(format!("/run/user/{uid}"));
let obstacle = match &probe {
SocketProbe::Reachable => None,
SocketProbe::Absent if !runtime_dir.exists() => Some(format!(
"no rootless engine for uid {uid}: {} does not exist, so there is no user \
manager to enable `podman.socket` in and nowhere for the socket to live. \
The account must LINGER first — `skidbladnir::rootless` performs that over \
org.freedesktop.login1 — and the runtime directory then appears by itself.",
runtime_dir.display()
)),
_ => socket_unreachable_reason(&socket, &probe),
};
RootlessReadiness {
uid,
socket,
probe,
obstacle,
}
}
#[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;
}
if let Ok(xdg) = std::env::var("XDG_RUNTIME_DIR") {
return format!("unix://{xdg}/podman/podman.sock");
}
match current_uid() {
Some(uid) => format!("unix://{}", rootless_socket_path(uid).display()),
None => "unix:///run/podman/podman.sock".to_string(),
}
}
fn connect() -> Result<Self> {
let url = Self::socket_url();
let path = url.strip_prefix("unix://").unwrap_or(&url);
if url.starts_with("unix://") {
if let Some(reason) = socket_unreachable(Path::new(path)) {
return Err(Error::Backend(reason));
}
}
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()),
})
}
#[allow(dead_code)]
pub fn create_body(
image: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
) -> bollard::models::ContainerCreateBody {
Self::create_body_with_binds(image, env, cmd, ports, &[])
}
pub fn create_body_with_binds(
image: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
binds: &[String],
) -> bollard::models::ContainerCreateBody {
Self::create_body_with_binds_and_net(image, env, cmd, ports, binds, None)
}
pub fn create_body_with_binds_and_net(
image: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
binds: &[String],
network_mode: Option<&str>,
) -> bollard::models::ContainerCreateBody {
Self::create_body_full(
image,
env,
cmd,
ports,
&[],
binds,
network_mode,
None,
None,
&[],
)
}
#[allow(clippy::too_many_arguments)]
pub fn create_body_full(
image: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
port_maps: &[crate::PortMap],
binds: &[String],
network_mode: Option<&str>,
cpus: Option<f64>,
mem_mb: Option<u32>,
security_opt: &[String],
) -> 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();
let mut publish = |host: u16, container: u16| {
let key = format!("{container}/tcp");
if !exposed.contains(&key) {
exposed.push(key.clone());
}
bindings.insert(
key,
Some(vec![PortBinding {
host_ip: Some("0.0.0.0".to_string()),
host_port: Some(host.to_string()),
}]),
);
};
for &p in ports {
publish(p, p);
}
for pm in port_maps {
publish(pm.host, pm.container);
}
let nano_cpus = cpus.filter(|c| *c > 0.0).map(|c| (c * 1e9) as i64);
let memory = mem_mb
.filter(|m| *m > 0)
.map(|m| i64::from(m) * 1024 * 1024);
let host_config = if bindings.is_empty()
&& binds.is_empty()
&& network_mode.is_none()
&& nano_cpus.is_none()
&& memory.is_none()
&& security_opt.is_empty()
{
None
} else {
Some(HostConfig {
port_bindings: if bindings.is_empty() {
None
} else {
Some(bindings)
},
binds: if binds.is_empty() {
None
} else {
Some(binds.to_vec())
},
network_mode: network_mode.map(str::to_string),
nano_cpus,
memory,
security_opt: if security_opt.is_empty() {
None
} else {
Some(security_opt.to_vec())
},
..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()
}
}
#[allow(dead_code)]
#[allow(clippy::too_many_arguments)]
pub fn create_body_with_res(
image: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
binds: &[String],
network_mode: Option<&str>,
cpus: Option<f64>,
mem_mb: Option<u32>,
) -> bollard::models::ContainerCreateBody {
Self::create_body_full(
image,
env,
cmd,
ports,
&[],
binds,
network_mode,
cpus,
mem_mb,
&[],
)
}
#[allow(dead_code)]
fn create_and_start(
&self,
image: &str,
name: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
) -> Result<()> {
self.create_and_start_with_binds(image, name, env, cmd, ports, &[])
}
#[allow(dead_code)]
fn create_and_start_with_binds(
&self,
image: &str,
name: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
binds: &[String],
) -> Result<()> {
self.create_and_start_with_binds_and_net(
image,
name,
env,
cmd,
ports,
&[],
binds,
None,
None,
None,
&[],
)
}
#[allow(clippy::too_many_arguments)]
fn create_and_start_with_binds_and_net(
&self,
image: &str,
name: &str,
env: &[String],
cmd: &[String],
ports: &[u16],
port_maps: &[crate::PortMap],
binds: &[String],
network_mode: Option<&str>,
cpus: Option<f64>,
mem_mb: Option<u32>,
security_opt: &[String],
) -> 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_full(
image,
env,
cmd,
ports,
port_maps,
binds,
network_mode,
cpus,
mem_mb,
security_opt,
);
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 build_image(&self, context_dir: &Path, containerfile: &str, tag: &str) -> Result<String> {
build_image_on(self, context_dir, containerfile, tag)
}
fn extract_path(&self, image: &str, container_path: &str, host_dest: &Path) -> Result<()> {
extract_path_on(self, image, container_path, host_dest)
}
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 list_containers(&self, prefix: &str) -> Result<Vec<(String, ContainerState)>> {
use bollard::query_parameters::ListContainersOptionsBuilder;
self.rt.block_on(async {
let summaries = self
.docker
.list_containers(Some(ListContainersOptionsBuilder::new().all(true).build()))
.await
.map_err(|e| Error::Backend(format!("listing containers: {e}")))?;
let mut out = Vec::new();
for c in summaries {
for name in c.names.iter().flatten() {
let name = name.strip_prefix('/').unwrap_or(name);
if !name.starts_with(prefix) {
continue;
}
let state = match c.state {
Some(bollard::models::ContainerSummaryStateEnum::RUNNING) => {
ContainerState::Running
}
Some(bollard::models::ContainerSummaryStateEnum::EXITED) => {
ContainerState::Exited(0)
}
_ => ContainerState::Gone,
};
out.push((name.to_owned(), state));
}
}
out.sort_by(|a, b| a.0.cmp(&b.0));
Ok(out)
})
}
fn container_env(&self, name: &str) -> Result<Vec<String>> {
use bollard::query_parameters::InspectContainerOptions;
self.rt.block_on(async {
let info = self
.docker
.inspect_container(name, None::<InspectContainerOptions>)
.await
.map_err(|e| Error::Backend(format!("inspecting container `{name}`: {e}")))?;
Ok(info
.config
.and_then(|c| c.env)
.unwrap_or_default()
.into_iter()
.collect())
})
}
fn exec(&self, name: &str, argv: &[String]) -> Result<ExecOutcome> {
use bollard::container::LogOutput;
use bollard::exec::{CreateExecOptions, StartExecOptions, StartExecResults};
use futures::StreamExt;
let docker = &self.docker;
self.rt.block_on(async {
let config = CreateExecOptions::<String> {
cmd: Some(argv.to_vec()),
attach_stdout: Some(true),
attach_stderr: Some(true),
..Default::default()
};
let created = docker
.create_exec(name, config)
.await
.map_err(|e| Error::Backend(format!("create exec in {name}: {e}")))?;
let mut stdout: Vec<String> = Vec::new();
let mut stderr: Vec<String> = Vec::new();
match docker
.start_exec(&created.id, None::<StartExecOptions>)
.await
.map_err(|e| Error::Backend(format!("start exec in {name}: {e}")))?
{
StartExecResults::Attached { mut output, .. } => {
while let Some(item) = output.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() {
if is_err {
stderr.push(line);
} else {
stdout.push(line);
}
}
}
Err(e) => {
return Err(Error::Backend(format!(
"read exec output in {name}: {e}"
)))
}
}
}
}
StartExecResults::Detached => {}
}
let inspect = docker
.inspect_exec(&created.id)
.await
.map_err(|e| Error::Backend(format!("inspect exec in {name}: {e}")))?;
Ok(ExecOutcome {
exit_code: inspect.exit_code,
stdout,
stderr,
})
})
}
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(feature = "backend-oci")]
impl ImageDaemon for Engine {
fn build_stream(
&self,
context: Vec<u8>,
containerfile: &str,
tag: &str,
) -> Vec<Result<Option<String>>> {
use bollard::body_full;
use bollard::query_parameters::BuildImageOptionsBuilder;
use futures::StreamExt;
let docker = &self.docker;
self.rt.block_on(async {
let opts = BuildImageOptionsBuilder::default()
.dockerfile(containerfile)
.t(tag)
.rm(true)
.build();
let mut stream = docker.build_image(opts, None, Some(body_full(context.into())));
let mut steps: Vec<Result<Option<String>>> = Vec::new();
while let Some(item) = stream.next().await {
steps.push(match item {
Ok(info) => Ok(info.error_detail.map(|d| d.message.unwrap_or_default())),
Err(e) => Err(Error::Backend(format!("build image {tag}: {e}"))),
});
}
steps
})
}
fn inspect_image(&self, image: &str) -> ImagePresence {
let docker = &self.docker;
self.rt.block_on(async {
match docker.inspect_image(image).await {
Ok(_) => ImagePresence::Present,
Err(bollard::errors::Error::DockerResponseServerError {
status_code: 404, ..
}) => ImagePresence::Absent,
Err(bollard::errors::Error::DockerResponseServerError { message, .. })
if message.to_ascii_lowercase().contains("no such image") =>
{
ImagePresence::Absent
}
Err(e) => ImagePresence::Failed(e.to_string()),
}
})
}
fn export_to(&self, image: &str, sink: &mut dyn std::io::Write) -> Result<u64> {
use futures::StreamExt;
let docker = &self.docker;
self.rt.block_on(async move {
let mut stream = docker.export_image(image);
let mut written: u64 = 0;
while let Some(chunk) = stream.next().await {
let chunk =
chunk.map_err(|e| Error::Backend(format!("export image {image}: {e}")))?;
sink.write_all(chunk.as_ref()).map_err(|e| {
Error::Backend(format!("write exported image {image} to sink: {e}"))
})?;
written = written.saturating_add(chunk.len() as u64);
}
Ok(written)
})
}
fn load_archive(&self, src: &Path) -> Result<()> {
use bollard::query_parameters::ImportImageOptionsBuilder;
use futures::StreamExt;
use std::io::Read;
let file = std::fs::File::open(src)
.map_err(|e| Error::Backend(format!("open image archive {}: {e}", src.display())))?;
let label = src.display().to_string();
let docker = &self.docker;
self.rt.block_on(async move {
let body = futures::stream::try_unfold(file, |mut file| async move {
let mut buf = vec![0u8; IMPORT_CHUNK];
let n = file.read(&mut buf)?;
if n == 0 {
Ok::<_, std::io::Error>(None)
} else {
buf.truncate(n);
Ok(Some((bytes::Bytes::from(buf), file)))
}
});
let mut stream =
docker.import_image_stream(ImportImageOptionsBuilder::new().build(), body, None);
while let Some(item) = stream.next().await {
item.map_err(|e| Error::Backend(format!("load image archive {label}: {e}")))?;
}
Ok(())
})
}
fn download_path_tar(&self, image: &str, container_path: &str) -> Result<Vec<u8>> {
use bollard::models::ContainerCreateBody;
use bollard::query_parameters::{
CreateContainerOptionsBuilder, DownloadFromContainerOptionsBuilder,
RemoveContainerOptionsBuilder,
};
use futures::StreamExt;
let docker = &self.docker;
let safe: String = image
.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '-' })
.collect();
let name = format!("draupnir-extract-{safe}-{}", std::process::id());
self.rt.block_on(async {
let _ = docker
.remove_container(
&name,
Some(RemoveContainerOptionsBuilder::new().force(true).build()),
)
.await;
let body = ContainerCreateBody {
image: Some(image.to_string()),
..Default::default()
};
docker
.create_container(
Some(
CreateContainerOptionsBuilder::new()
.name(name.as_str())
.build(),
),
body,
)
.await
.map_err(|e| {
Error::Backend(format!("create extract container from {image}: {e}"))
})?;
let opts = DownloadFromContainerOptionsBuilder::default()
.path(container_path)
.build();
let mut stream = docker.download_from_container(&name, Some(opts));
let mut buf: Vec<u8> = Vec::new();
let mut dl_err: Option<Error> = None;
while let Some(item) = stream.next().await {
match item {
Ok(chunk) => buf.extend_from_slice(&chunk),
Err(e) => {
dl_err = Some(Error::Backend(format!(
"download {container_path} from {image}: {e}"
)));
break;
}
}
}
let _ = docker
.remove_container(
&name,
Some(RemoveContainerOptionsBuilder::new().force(true).build()),
)
.await;
match dl_err {
Some(e) => Err(e),
None => Ok(buf),
}
})
}
}
#[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 a_permission_denied_socket_is_not_reported_as_missing() {
let path = Path::new("/run/podman/podman.sock");
let probe = SocketProbe::Denied {
blocked_at: "/run/podman".into(),
owner: Some(Ownership {
user: "root".into(),
group: "root".into(),
mode: 0o700,
}),
};
let m = socket_unreachable_reason(path, &probe).expect("denied is not reachable");
assert!(
!m.contains("not found"),
"EACCES must not be worded as ENOENT — that is the whole defect: {m}"
);
assert!(
m.contains("EXISTS") && m.contains("permission denied"),
"it must say the socket is there and the process is refused: {m}"
);
assert!(
!m.contains("enable --now"),
"never advise enabling a unit that is already running: {m}"
);
assert!(
m.contains("/run/podman") && m.contains("root:root") && m.contains("0700"),
"it must name the blocking component and its ownership: {m}"
);
assert!(
m.contains("SocketGroup=") && m.contains("HOST decision"),
"a root-group socket needs a host decision, and the message must say so: {m}"
);
assert!(
!m.contains("SupplementaryGroups=root"),
"'add the service to the root group' is advice nobody should take: {m}"
);
}
#[test]
fn an_absent_socket_still_says_it_is_absent_and_how_to_enable_it() {
let path = Path::new("/run/podman/podman.sock");
let m = socket_unreachable_reason(path, &SocketProbe::Absent).expect("absent");
assert!(m.contains("not found"), "ENOENT is genuinely 'not found': {m}");
assert!(
m.contains("enable --now podman.socket"),
"and enabling the unit is the right advice here: {m}"
);
}
#[test]
fn the_enable_hint_matches_the_sockets_scope() {
let system = Path::new("/run/podman/podman.sock");
let user = Path::new("/run/user/1000/podman/podman.sock");
assert!(!is_user_scope_socket(system), "/run/podman is system scope");
assert!(is_user_scope_socket(user), "/run/user/1000 is user scope");
let sys_msg = socket_unreachable_reason(system, &SocketProbe::Absent).unwrap();
assert!(
sys_msg.contains("sudo systemctl enable") && !sys_msg.contains("--user"),
"a SYSTEM socket must not be pointed at `systemctl --user`: {sys_msg}"
);
let usr_msg = socket_unreachable_reason(user, &SocketProbe::Absent).unwrap();
assert!(
usr_msg.contains("systemctl --user enable"),
"a rootless socket keeps the --user scope: {usr_msg}"
);
}
#[test]
fn a_non_root_socket_group_is_advised_as_a_supplementary_group() {
let m = socket_unreachable_reason(
Path::new("/run/podman/podman.sock"),
&SocketProbe::Denied {
blocked_at: "/run/podman/podman.sock".into(),
owner: Some(Ownership {
user: "root".into(),
group: "podman".into(),
mode: 0o660,
}),
},
)
.unwrap();
assert!(
m.contains("SupplementaryGroups=podman"),
"a joinable group IS the fix, and the directive must be named: {m}"
);
assert!(
!m.contains("HOST decision"),
"this one does not need a host decision — the group already exists: {m}"
);
}
#[test]
fn probe_separates_absent_from_reachable_on_the_real_filesystem() {
assert_eq!(
probe_socket_path(Path::new(
"/nonexistent-draupnir-probe/podman/podman.sock"
)),
SocketProbe::Absent
);
assert_eq!(
probe_socket_path(Path::new("/etc/passwd")),
SocketProbe::Reachable
);
assert!(
socket_unreachable(Path::new("/etc/passwd")).is_none(),
"a reachable path yields no complaint"
);
}
#[test]
fn ids_resolve_to_names_without_ffi_or_a_subprocess() {
assert_eq!(name_for_id("/etc/passwd", 0).as_deref(), Some("root"));
assert_eq!(name_for_id("/etc/group", 0).as_deref(), Some("root"));
assert_eq!(name_for_id("/etc/passwd", 4_294_967_294), None);
}
#[test]
fn the_rootless_socket_path_is_derived_from_the_uid() {
assert_eq!(
rootless_socket_path(111),
Path::new("/run/user/111/podman/podman.sock")
);
assert!(is_user_scope_socket(&rootless_socket_path(111)));
}
#[test]
fn the_current_uid_is_read_from_proc_without_shell_or_ffi() {
assert_eq!(parse_uid_line("Name:\tsh\nUid:\t111\t111\t111\t111\n"), Some(111));
assert_eq!(parse_uid_line("Uid:\t0\t0\t0\t0"), Some(0));
assert_eq!(parse_uid_line("Gid:\t111\t111\t111\t111\n"), None);
assert!(current_uid().is_some(), "Linux always publishes /proc/self/status");
}
#[test]
fn the_socket_fallback_follows_the_processs_own_uid_not_a_hardcoded_1000() {
let uid = current_uid().expect("Linux");
let derived = format!("unix://{}", rootless_socket_path(uid).display());
assert!(
derived.contains(&format!("/run/user/{uid}/")),
"the fallback must name THIS process's runtime dir: {derived}"
);
assert_eq!(
rootless_socket_path(111),
Path::new("/run/user/111/podman/podman.sock"),
"uid 111 must never be sent to /run/user/1000"
);
}
#[test]
fn system_migrate_is_owed_only_where_podman_already_ran_and_is_never_performed() {
let fresh = migrate_verdict("korp", false);
assert!(matches!(fresh, MigrateVerdict::NotOwed(_)), "{fresh:?}");
assert!(
fresh.detail().contains("not needed"),
"a fresh account must be told to skip it: {}",
fresh.detail()
);
let stale = migrate_verdict("korp", true);
assert!(matches!(stale, MigrateVerdict::Owed(_)), "{stale:?}");
let m = stale.detail();
assert!(m.contains("podman system migrate"), "{m}");
assert!(
m.contains("as `korp`"),
"it must say WHICH account has to run it — as root it does nothing: {m}"
);
assert!(
m.contains("no libpod REST endpoint") && m.contains("will not run it"),
"and it must say plainly that draupnir refuses, and why: {m}"
);
assert_eq!(
libpod_storage_dir("/home/korp"),
Path::new("/home/korp/.local/share/containers/storage/libpod")
);
}
#[test]
fn an_absent_rootless_socket_asks_about_lingering_before_the_unit() {
let r = rootless_readiness(4_294_967_200);
assert!(!r.is_ready());
let m = r.obstacle.expect("not ready ⇒ an obstacle");
assert!(m.contains("does not exist"), "{m}");
assert!(
m.contains("LINGER") && m.contains("org.freedesktop.login1"),
"the first obstacle is lingering, and it must name the mechanism: {m}"
);
assert!(
!m.contains("systemctl --user enable"),
"…and must NOT advise enabling a unit in a user manager that does not \
exist yet: {m}"
);
}
#[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"));
}
#[cfg(feature = "backend-oci")]
#[test]
fn create_body_with_binds_attaches_host_mounts_and_stays_parity_when_empty() {
let with = Engine::create_body_with_binds(
"wix:4",
&[],
&["build".to_string()],
&[],
&[
"/host/in:/work/in:ro".to_string(),
"/host/out:/work/out".to_string(),
],
);
let hc = with.host_config.expect("host config for binds");
assert_eq!(
hc.binds,
Some(vec![
"/host/in:/work/in:ro".to_string(),
"/host/out:/work/out".to_string()
])
);
assert!(hc.port_bindings.is_none(), "no ports => no port bindings");
let plain = Engine::create_body("wix:4", &[], &["build".to_string()], &[]);
let via_empty =
Engine::create_body_with_binds("wix:4", &[], &["build".to_string()], &[], &[]);
assert!(plain.host_config.is_none());
assert_eq!(
plain, via_empty,
"empty binds is byte-parity with create_body"
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn create_body_sets_network_mode_none_and_default_is_byte_identical() {
let airgap = Engine::create_body_with_binds_and_net(
"job:1",
&[],
&["run".to_string()],
&[],
&[],
Some("none"),
);
let hc = airgap
.host_config
.expect("host config minted for network mode");
assert_eq!(
hc.network_mode.as_deref(),
Some("none"),
"airgap => --network none"
);
assert!(hc.port_bindings.is_none(), "no ports => no port bindings");
assert!(hc.binds.is_none(), "no binds => no binds");
let plain = Engine::create_body_with_binds("job:1", &[], &["run".to_string()], &[], &[]);
let via_none = Engine::create_body_with_binds_and_net(
"job:1",
&[],
&["run".to_string()],
&[],
&[],
None,
);
assert!(
plain.host_config.is_none(),
"default net + no ports/binds => no host_config"
);
assert_eq!(
plain, via_none,
"None network_mode is byte-parity with the net-less builder"
);
let ported = Engine::create_body_with_binds("web:1", &[], &[], &[8080], &[]);
let ported_none =
Engine::create_body_with_binds_and_net("web:1", &[], &[], &[8080], &[], None);
assert_eq!(
ported, ported_none,
"None net leaves a ported spec byte-identical"
);
assert!(
ported.host_config.unwrap().network_mode.is_none(),
"default net sets no network_mode"
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn published_ports_bind_to_the_host_on_a_non_isolated_network() {
let body = Engine::create_body_with_binds_and_net(
"docker.io/falkordb/falkordb:v4.20.0",
&[],
&[],
&[6379],
&[],
None,
);
let hc = body.host_config.expect("host config for published port");
assert!(
hc.network_mode.is_none(),
"non-isolated: no --network none (reachable)"
);
let bindings = hc.port_bindings.expect("port bindings");
let b = bindings
.get("6379/tcp")
.and_then(|v| v.clone())
.expect("6379 published");
assert_eq!(
b[0].host_ip.as_deref(),
Some("0.0.0.0"),
"bound on the host"
);
assert_eq!(
b[0].host_port.as_deref(),
Some("6379"),
"host port == container port"
);
assert!(
body.exposed_ports.unwrap().iter().any(|s| s == "6379/tcp"),
"6379 exposed",
);
crate::functional_status(
"draupnir/container",
"published_port_non_isolated",
hc.network_mode.is_none(),
"published ports bind 0.0.0.0:host on the default (non-isolated) net → host-reachable local-zone container",
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn resource_caps_reach_the_create_body_and_none_is_byte_identical() {
let hot = Engine::create_body_with_res(
"docker.io/falkordb/falkordb:v4.20.0",
&[],
&[],
&[6379],
&[],
None,
Some(12.0),
Some(16384),
);
let hc = hot.host_config.expect("host config for resource caps");
assert_eq!(hc.nano_cpus, Some(12_000_000_000), "12 cores → nano_cpus");
assert_eq!(
hc.memory,
Some(16384_i64 * 1024 * 1024),
"16384 MiB → bytes"
);
assert!(
hc.port_bindings.unwrap().contains_key("6379/tcp"),
"port still published"
);
let plain = Engine::create_body_with_binds_and_net("redis:7", &[], &[], &[], &[], None);
let via_none =
Engine::create_body_with_res("redis:7", &[], &[], &[], &[], None, None, None);
assert!(
plain.host_config.is_none(),
"no caps + no ports/binds → no host_config"
);
assert_eq!(
plain, via_none,
"None caps is byte-parity with the resource-less builder"
);
crate::functional_status(
"draupnir/container",
"resource_caps_reach_create_body",
hc.nano_cpus == Some(12_000_000_000) && hc.memory == Some(16384_i64 * 1024 * 1024),
"cpus/mem reach HostConfig.nano_cpus/memory; None = all host cores (hot-infra never throttled)",
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn distinct_host_container_port_map_publishes_host_to_container() {
let body = Engine::create_body_full(
"docker.io/falkordb/falkordb:v4.20.0",
&[],
&[],
&[],
&[crate::PortMap::new(6380, 6379)],
&[],
None,
None,
None,
&[],
);
let hc = body.host_config.expect("host config for published map");
assert!(
hc.network_mode.is_none(),
"non-isolated: reachable (no --network none)"
);
let bindings = hc.port_bindings.expect("port bindings");
let b = bindings
.get("6379/tcp")
.and_then(|v| v.clone())
.expect("container 6379 published");
assert_eq!(
b[0].host_ip.as_deref(),
Some("0.0.0.0"),
"bound on the host"
);
assert_eq!(
b[0].host_port.as_deref(),
Some("6380"),
"host 6380 -> container 6379"
);
assert!(
!bindings.contains_key("6380/tcp"),
"the HOST port is NOT the container key"
);
assert!(
body.exposed_ports
.as_ref()
.unwrap()
.iter()
.any(|s| s == "6379/tcp"),
"the CONTAINER port is exposed, not the host port",
);
let demo = Engine::create_body_full(
"docker.io/falkordb/falkordb:v4.20.0",
&[],
&[],
&[6379],
&[],
&[],
None,
None,
None,
&[],
);
let demo_via_map = Engine::create_body_full(
"docker.io/falkordb/falkordb:v4.20.0",
&[],
&[],
&[],
&[crate::PortMap::same(6379)],
&[],
None,
None,
None,
&[],
);
let db = demo.host_config.clone().unwrap().port_bindings.unwrap();
let sb = db
.get("6379/tcp")
.and_then(|v| v.clone())
.expect("demo 6379 published");
assert_eq!(
sb[0].host_port.as_deref(),
Some("6379"),
"single-port form: host==container==6379"
);
assert_eq!(
demo, demo_via_map,
"[6379] and PortMap::same(6379) render the same wire"
);
let bare = Engine::create_body_full("redis:7", &[], &[], &[], &[], &[], None, None, None, &[]);
assert!(
bare.host_config.is_none(),
"no ports/maps/net/caps => no host_config (byte-parity)"
);
crate::functional_status(
"draupnir/container",
"distinct_host_container_port_map",
sb[0].host_port.as_deref() == Some("6379")
&& b[0].host_port.as_deref() == Some("6380")
&& bindings.contains_key("6379/tcp"),
"PortMap{host,container} publishes host:container (-p 6380:6379) so a per-zone host port reaches a fixed in-container port",
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn a_security_option_reaches_the_create_body_and_an_empty_one_changes_nothing() {
let opts = vec!["seccomp=unconfined".to_string()];
let with = Engine::create_body_full(
"redis:7",
&[],
&[],
&[],
&[],
&[],
None,
None,
None,
&opts,
);
let hc = with
.host_config
.clone()
.expect("a security option alone mints a host config");
assert_eq!(
hc.security_opt.as_deref(),
Some(opts.as_slice()),
"the security option did not reach HostConfig.security_opt, so the daemon would \
apply its default profile while the caller believes otherwise"
);
let bare =
Engine::create_body_full("redis:7", &[], &[], &[], &[], &[], None, None, None, &[]);
assert!(
bare.host_config.is_none(),
"an empty security_opt minted a host config — every existing create body would \
stop being byte-identical to what shipped"
);
let spec = crate::BootSpec::container("probe", "redis:7")
.with_security_opt("seccomp=unconfined")
.with_security_opt("no-new-privileges");
assert_eq!(
spec.security_opt,
vec![
"seccomp=unconfined".to_string(),
"no-new-privileges".to_string()
],
"with_security_opt replaced instead of appending"
);
crate::functional_status(
"draupnir/container",
"security_opt",
hc.security_opt.is_some() && bare.host_config.is_none(),
"BootSpec::with_security_opt renders --security-opt onto HostConfig.security_opt; empty stays byte-identical",
);
}
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());
}
#[derive(Default)]
struct ExecRecorder {
seen: RefCell<Vec<Vec<String>>>,
outcome: ExecOutcome,
}
impl ContainerControl for ExecRecorder {
fn container_state(&self, _m: &Machine) -> ContainerState {
ContainerState::Running
}
fn drain_logs(&self, _m: &Machine) -> Vec<String> {
Vec::new()
}
fn stop(&self, _m: &Machine) {}
fn exec_command(&self, command: &[String]) -> Result<ExecOutcome> {
self.seen.borrow_mut().push(command.to_vec());
Ok(self.outcome.clone())
}
}
#[test]
fn exec_argv_builds_the_podman_exec_command() {
assert_eq!(
exec_argv("draupnir-cache", &["redis-cli", "ping"]),
vec!["exec", "draupnir-cache", "redis-cli", "ping"]
);
assert_eq!(
exec_argv("draupnir-x", &["true"]),
vec!["exec", "draupnir-x", "true"]
);
}
#[test]
fn exec_assembles_the_command_and_returns_the_outcome() {
let recorder = ExecRecorder {
outcome: ExecOutcome {
exit_code: Some(0),
stdout: vec!["PONG".into()],
stderr: vec![],
},
..Default::default()
};
let m = Machine::started("draupnir-cache", &BootSpec::container("cache", "redis:7"));
let out = recorder.exec(&m, &["redis-cli", "ping"]).unwrap();
assert_eq!(
recorder.seen.borrow().as_slice(),
&[vec![
"exec".to_string(),
"draupnir-cache".to_string(),
"redis-cli".to_string(),
"ping".to_string(),
]]
);
assert_eq!(out.exit_code, Some(0));
assert_eq!(out.stdout, vec!["PONG"]);
}
#[test]
fn exec_is_rejected_on_a_non_container_machine() {
let recorder = ExecRecorder::default();
let kvm = Machine::started(
"vm-1",
&BootSpec::kvm_kernel_rootfs("appliance", "/bzImage", "/rootfs.cpio.gz"),
);
assert!(matches!(recorder.exec(&kvm, &["ls"]), Err(Error::Spec(_))));
assert!(
recorder.seen.borrow().is_empty(),
"guard runs before the engine"
);
}
#[test]
fn exec_rejects_an_empty_argv() {
let recorder = ExecRecorder::default();
let m = Machine::started("draupnir-cache", &BootSpec::container("cache", "redis:7"));
assert!(matches!(recorder.exec(&m, &[]), Err(Error::Spec(_))));
assert!(recorder.seen.borrow().is_empty());
}
#[test]
fn exec_without_an_engine_is_unsupported_not_faked() {
struct NoEngine;
impl ContainerControl for NoEngine {
fn container_state(&self, _m: &Machine) -> ContainerState {
ContainerState::Gone
}
fn drain_logs(&self, _m: &Machine) -> Vec<String> {
Vec::new()
}
fn stop(&self, _m: &Machine) {}
}
let m = Machine::started("draupnir-x", &BootSpec::container("x", "img:1"));
assert!(matches!(
NoEngine.exec(&m, &["true"]),
Err(Error::Unsupported(_))
));
}
#[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());
}
#[cfg(feature = "backend-oci")]
#[test]
fn context_tar_packs_the_build_context() {
let dir = std::env::temp_dir().join(format!("draupnir-ctx-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(dir.join("sub")).unwrap();
std::fs::write(dir.join("Containerfile"), b"FROM scratch\n").unwrap();
std::fs::write(dir.join("sub").join("app.txt"), b"hi").unwrap();
let bytes = context_tar(&dir).unwrap();
assert!(!bytes.is_empty(), "the tar carries the context");
let mut ar = tar::Archive::new(std::io::Cursor::new(bytes));
let names: Vec<String> = ar
.entries()
.unwrap()
.map(|e| {
e.unwrap()
.path()
.unwrap()
.to_string_lossy()
.replace('\\', "/")
})
.collect();
assert!(
names.iter().any(|n| n.ends_with("Containerfile")),
"Containerfile packed: {names:?}"
);
assert!(
names.iter().any(|n| n.ends_with("sub/app.txt")),
"nested file packed: {names:?}"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[cfg(feature = "backend-oci")]
#[derive(Default)]
struct ScriptedDaemon {
build_steps: RefCell<Vec<Result<Option<String>>>>,
download: RefCell<Option<Result<Vec<u8>>>>,
saw_build: RefCell<Vec<(usize, String)>>,
saw_download: RefCell<Vec<(String, String)>>,
presence: RefCell<Option<ImagePresence>>,
saw_inspect: RefCell<Vec<String>>,
export_bytes: RefCell<Option<Result<Vec<u8>>>>,
saw_export: RefCell<Vec<String>>,
load_result: RefCell<Option<Result<()>>>,
saw_load: RefCell<Vec<String>>,
}
#[cfg(feature = "backend-oci")]
impl ImageDaemon for ScriptedDaemon {
fn build_stream(
&self,
context: Vec<u8>,
_containerfile: &str,
tag: &str,
) -> Vec<Result<Option<String>>> {
self.saw_build
.borrow_mut()
.push((context.len(), tag.to_string()));
std::mem::take(&mut *self.build_steps.borrow_mut())
}
fn download_path_tar(&self, image: &str, container_path: &str) -> Result<Vec<u8>> {
self.saw_download
.borrow_mut()
.push((image.to_string(), container_path.to_string()));
self.download
.borrow_mut()
.take()
.unwrap_or_else(|| Ok(Vec::new()))
}
fn inspect_image(&self, image: &str) -> ImagePresence {
self.saw_inspect.borrow_mut().push(image.to_string());
self.presence
.borrow_mut()
.take()
.unwrap_or(ImagePresence::Present)
}
fn export_to(&self, image: &str, sink: &mut dyn std::io::Write) -> Result<u64> {
self.saw_export.borrow_mut().push(image.to_string());
let payload = self.export_bytes.borrow_mut().take();
match payload {
Some(Err(e)) => Err(e),
Some(Ok(bytes)) => {
sink.write_all(&bytes).unwrap();
Ok(bytes.len() as u64)
}
None => Ok(0),
}
}
fn load_archive(&self, src: &Path) -> Result<()> {
self.saw_load
.borrow_mut()
.push(src.display().to_string());
self.load_result.borrow_mut().take().unwrap_or(Ok(()))
}
}
#[cfg(feature = "backend-oci")]
#[test]
fn image_present_reads_true_when_the_daemon_resolves_it() {
let daemon = ScriptedDaemon::default();
*daemon.presence.borrow_mut() = Some(ImagePresence::Present);
assert!(image_present_on(&daemon, "localhost/nordisk-falkordb-valkey:v4.20.1").unwrap());
assert_eq!(
daemon.saw_inspect.borrow().as_slice(),
["localhost/nordisk-falkordb-valkey:v4.20.1"],
"the probe asked the daemon about the image it was handed"
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn image_present_reads_false_when_the_daemon_has_no_such_image() {
let daemon = ScriptedDaemon::default();
*daemon.presence.borrow_mut() = Some(ImagePresence::Absent);
assert!(!image_present_on(&daemon, "localhost/nope:1").unwrap());
}
#[cfg(feature = "backend-oci")]
#[test]
fn image_present_errors_when_the_daemon_cannot_be_asked() {
let daemon = ScriptedDaemon::default();
*daemon.presence.borrow_mut() = Some(ImagePresence::Failed("connection refused".into()));
let err = image_present_on(&daemon, "localhost/spark:1").unwrap_err();
assert!(
matches!(&err, Error::Backend(m) if m.contains("localhost/spark:1") && m.contains("connection refused")),
"a transport failure surfaces as a Backend error naming the image + cause, not as absence: {err:?}"
);
}
#[cfg(feature = "backend-oci")]
fn make_cp_tar(root: &str, files: &[(&str, &[u8])]) -> Vec<u8> {
let mut buf: Vec<u8> = Vec::new();
{
let mut b = tar::Builder::new(&mut buf);
for (rel, data) in files {
let mut header = tar::Header::new_gnu();
header.set_size(data.len() as u64);
header.set_mode(0o644);
header.set_cksum();
b.append_data(&mut header, format!("{root}/{rel}"), *data)
.unwrap();
}
b.finish().unwrap();
}
buf
}
#[cfg(feature = "backend-oci")]
#[test]
fn build_image_returns_the_tag_on_a_clean_stream() {
let dir = std::env::temp_dir().join(format!("draupnir-build-ok-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("Containerfile"), b"FROM scratch\n").unwrap();
let daemon = ScriptedDaemon::default();
*daemon.build_steps.borrow_mut() = vec![Ok(None), Ok(None)]; let tag = build_image_on(&daemon, &dir, "Containerfile", "app:test").unwrap();
assert_eq!(tag, "app:test", "a clean build returns its tag");
let saw = daemon.saw_build.borrow();
assert_eq!(saw.len(), 1, "the daemon was streamed exactly one build");
assert!(
saw[0].0 > 0,
"a non-empty context tar was sent: {}",
saw[0].0
);
assert_eq!(saw[0].1, "app:test", "the tag reached the daemon");
let _ = std::fs::remove_dir_all(&dir);
}
#[cfg(feature = "backend-oci")]
#[test]
fn build_image_surfaces_an_error_detail_as_err() {
let dir = std::env::temp_dir().join(format!("draupnir-build-fail-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("Containerfile"), b"FROM scratch\n").unwrap();
let daemon = ScriptedDaemon::default();
*daemon.build_steps.borrow_mut() = vec![
Ok(None), Ok(Some("no such file or directory: /nope".into())), ];
let err = build_image_on(&daemon, &dir, "Containerfile", "app:bad").unwrap_err();
match err {
Error::Backend(m) => {
assert!(m.contains("app:bad"), "names the tag: {m}");
assert!(
m.contains("no such file"),
"carries the daemon message: {m}"
);
}
other => panic!("expected Error::Backend, got {other:?}"),
}
crate::functional_status(
"draupnir/container",
"build_image_error_detail_fails",
true,
"a BuildInfo error_detail surfaces as Err(Error::Backend) naming the tag + message — never a fake image",
);
let _ = std::fs::remove_dir_all(&dir);
}
#[cfg(feature = "backend-oci")]
#[test]
fn build_image_propagates_a_transport_error() {
let dir = std::env::temp_dir().join(format!("draupnir-build-tx-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("Containerfile"), b"FROM scratch\n").unwrap();
let daemon = ScriptedDaemon::default();
*daemon.build_steps.borrow_mut() = vec![Err(Error::Backend(
"build image app:tx: connection reset".into(),
))];
let err = build_image_on(&daemon, &dir, "Containerfile", "app:tx").unwrap_err();
assert!(matches!(err, Error::Backend(m) if m.contains("connection reset")));
let _ = std::fs::remove_dir_all(&dir);
}
#[cfg(feature = "backend-oci")]
#[test]
fn extract_path_unpacks_the_downloaded_tar_with_basename_layout() {
let tar = make_cp_tar(
"out",
&[
("app.msi", b"windows-msi-bytes"),
("logs/build.log", b"built ok"),
],
);
let dest = std::env::temp_dir().join(format!("draupnir-extract-ok-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dest);
let daemon = ScriptedDaemon::default();
*daemon.download.borrow_mut() = Some(Ok(tar));
extract_path_on(&daemon, "app:test", "/opt/out", &dest).unwrap();
assert_eq!(
std::fs::read(dest.join("out").join("app.msi")).unwrap(),
b"windows-msi-bytes",
"the top-level extracted file lands at host_dest/out/app.msi"
);
assert_eq!(
std::fs::read(dest.join("out").join("logs").join("build.log")).unwrap(),
b"built ok",
"the nested file keeps its subtree under host_dest/out/logs/"
);
assert_eq!(
daemon.saw_download.borrow().as_slice(),
&[("app:test".to_string(), "/opt/out".to_string())]
);
crate::functional_status(
"draupnir/container",
"extract_path_basename_unpack",
dest.join("out").join("app.msi").exists(),
"extract_path unpacks the downloaded tar into host_dest with the podman-cp basename layout (host_dest/<basename>/…)",
);
let _ = std::fs::remove_dir_all(&dest);
}
#[cfg(feature = "backend-oci")]
#[test]
fn extract_path_surfaces_a_download_error_and_does_not_unpack() {
let dest =
std::env::temp_dir().join(format!("draupnir-extract-err-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dest);
let daemon = ScriptedDaemon::default();
*daemon.download.borrow_mut() = Some(Err(Error::Backend(
"download /out from app:x: broken pipe".into(),
)));
let err = extract_path_on(&daemon, "app:x", "/out", &dest).unwrap_err();
assert!(matches!(err, Error::Backend(m) if m.contains("broken pipe")));
assert!(
!dest.exists(),
"nothing is unpacked (dest never created) on a download error"
);
}
#[cfg(not(feature = "backend-oci"))]
#[test]
fn build_and_extract_are_unsupported_without_the_engine() {
let boot = ContainerBoot::new();
assert!(matches!(
boot.build_image(Path::new("."), "Containerfile", "x:test"),
Err(Error::Unsupported(_))
));
assert!(matches!(
boot.extract_path("x:test", "/out", Path::new("/tmp/draupnir-x")),
Err(Error::Unsupported(_))
));
}
#[cfg(feature = "backend-oci")]
fn make_image_archive(config: &str, tags: &[&str], layer_len: usize) -> Vec<u8> {
let repo_tags = tags
.iter()
.map(|t| format!("{t:?}"))
.collect::<Vec<_>>()
.join(",");
let manifest = format!(
r#"[{{"Config":"{config}","RepoTags":[{repo_tags}],"Layers":["blobs/sha256/layer0"]}}]"#
);
let mut buf: Vec<u8> = Vec::new();
{
let mut b = tar::Builder::new(&mut buf);
let layer = vec![0xABu8; layer_len];
let mut h = tar::Header::new_gnu();
h.set_size(layer.len() as u64);
h.set_mode(0o644);
h.set_cksum();
b.append_data(&mut h, "blobs/sha256/layer0", layer.as_slice())
.unwrap();
let mut h = tar::Header::new_gnu();
h.set_size(manifest.len() as u64);
h.set_mode(0o644);
h.set_cksum();
b.append_data(&mut h, "manifest.json", manifest.as_bytes())
.unwrap();
b.finish().unwrap();
}
buf
}
#[cfg(feature = "backend-oci")]
#[test]
fn image_archive_reads_its_tags_and_config_digest() {
let tar = make_image_archive(
"blobs/sha256/deadbeefcafe",
&["localhost/nordisk-falkordb-valkey:9-20260724"],
4096,
);
let len = tar.len() as u64;
let (tags, digest) = scan_image_archive(std::io::Cursor::new(tar), len).unwrap();
assert_eq!(tags, ["localhost/nordisk-falkordb-valkey:9-20260724"]);
assert_eq!(digest, "deadbeefcafe");
}
#[cfg(feature = "backend-oci")]
#[test]
fn image_archive_rejects_a_truncated_archive() {
let mut tar = make_image_archive("blobs/sha256/abc123", &["x:1"], 64 * 1024);
let full = tar.len() as u64;
tar.truncate(2048);
let len = tar.len() as u64;
let fault = scan_image_archive(std::io::Cursor::new(tar), len).unwrap_err();
match fault {
ArchiveFault::Truncated { declared, actual } => {
assert!(
declared > actual,
"the cut is reported as declared({declared}) past actual({actual})"
);
assert_eq!(actual, len);
assert!(declared <= full);
}
other => panic!("a cut archive must read as Truncated, got {other:?}"),
}
}
#[cfg(feature = "backend-oci")]
#[test]
fn image_archive_rejects_a_tar_that_is_not_an_image() {
let mut buf: Vec<u8> = Vec::new();
{
let mut b = tar::Builder::new(&mut buf);
let mut h = tar::Header::new_gnu();
h.set_size(3);
h.set_mode(0o644);
h.set_cksum();
b.append_data(&mut h, "hello.txt", &b"hi\n"[..]).unwrap();
b.finish().unwrap();
}
let len = buf.len() as u64;
assert_eq!(
scan_image_archive(std::io::Cursor::new(buf), len).unwrap_err(),
ArchiveFault::NotAnImageArchive
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn image_archive_rejects_an_empty_file() {
assert_eq!(
scan_image_archive(std::io::Cursor::new(Vec::new()), 0).unwrap_err(),
ArchiveFault::NotAnImageArchive
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn manifest_config_digest_is_the_same_for_both_spellings() {
let podman =
br#"[{"Config":"blobs/sha256/05455e6bc39e","RepoTags":["a:1"],"Layers":[]}]"#.to_vec();
let docker = br#"[{"Config":"05455e6bc39e.json","RepoTags":["a:1"],"Layers":[]}]"#.to_vec();
assert_eq!(
parse_archive_manifest(&podman).unwrap(),
parse_archive_manifest(&docker).unwrap()
);
assert_eq!(parse_archive_manifest(&podman).unwrap().1, "05455e6bc39e");
}
#[cfg(feature = "backend-oci")]
#[test]
fn manifest_that_is_not_the_archive_shape_is_malformed() {
assert!(matches!(
parse_archive_manifest(br#"{"Config":"x"}"#).unwrap_err(),
ArchiveFault::Malformed(m) if m.contains("not a JSON array")
));
assert!(matches!(
parse_archive_manifest(b"[]").unwrap_err(),
ArchiveFault::Malformed(m) if m.contains("empty array")
));
assert!(matches!(
parse_archive_manifest(br#"[{"RepoTags":["a:1"]}]"#).unwrap_err(),
ArchiveFault::Malformed(m) if m.contains("no Config")
));
}
#[cfg(feature = "backend-oci")]
#[test]
fn export_refuses_a_missing_image_and_writes_no_file() {
let dest =
std::env::temp_dir().join(format!("draupnir-export-absent-{}.tar", std::process::id()));
let _ = std::fs::remove_file(&dest);
let daemon = ScriptedDaemon::default();
*daemon.presence.borrow_mut() = Some(ImagePresence::Absent);
let err = export_image_on(&daemon, "localhost/nordisk-spark-iceberg:4.1.2", &dest)
.unwrap_err();
assert!(
matches!(&err, Error::Backend(m) if m.contains("no such image locally")),
"got {err:?}"
);
assert!(!dest.exists(), "a missing image leaves no stub archive");
assert!(
daemon.saw_export.borrow().is_empty(),
"the export stream is never opened for an image that is not there"
);
}
#[cfg(feature = "backend-oci")]
#[test]
fn export_reports_an_unreachable_engine_as_such() {
let dest = std::env::temp_dir()
.join(format!("draupnir-export-nosock-{}.tar", std::process::id()));
let _ = std::fs::remove_file(&dest);
let daemon = ScriptedDaemon::default();
*daemon.presence.borrow_mut() =
Some(ImagePresence::Failed("connection refused".into()));
let err = export_image_on(&daemon, "localhost/x:1", &dest).unwrap_err();
assert!(
matches!(&err, Error::Backend(m)
if m.contains("container engine unreachable") && m.contains("connection refused")),
"got {err:?}"
);
assert!(!dest.exists());
}
#[cfg(feature = "backend-oci")]
#[test]
fn export_rejects_an_empty_stream_instead_of_claiming_success() {
let dest =
std::env::temp_dir().join(format!("draupnir-export-empty-{}.tar", std::process::id()));
let _ = std::fs::remove_file(&dest);
let daemon = ScriptedDaemon::default();
let err = export_image_on(&daemon, "localhost/x:1", &dest).unwrap_err();
assert!(
matches!(&err, Error::Backend(m) if m.contains("not an image archive")),
"got {err:?}"
);
let _ = std::fs::remove_file(&dest);
}
#[cfg(feature = "backend-oci")]
#[test]
fn export_refuses_an_archive_tagged_as_a_different_image() {
let dest =
std::env::temp_dir().join(format!("draupnir-export-wrong-{}.tar", std::process::id()));
let _ = std::fs::remove_file(&dest);
let daemon = ScriptedDaemon::default();
*daemon.export_bytes.borrow_mut() = Some(Ok(make_image_archive(
"blobs/sha256/f00d",
&["localhost/nordisk-falkordb-redis:9-20260724"],
512,
)));
let err = export_image_on(
&daemon,
"localhost/nordisk-falkordb-valkey:9-20260724",
&dest,
)
.unwrap_err();
assert!(
matches!(&err, Error::Backend(m) if m.contains("refusing to record an image under a tag it does not have")),
"got {err:?}"
);
let _ = std::fs::remove_file(&dest);
}
#[cfg(feature = "backend-oci")]
#[test]
fn import_verifies_the_archive_before_touching_the_daemon() {
let src =
std::env::temp_dir().join(format!("draupnir-import-cut-{}.tar", std::process::id()));
let mut tar = make_image_archive("blobs/sha256/abc", &["x:1"], 64 * 1024);
tar.truncate(2048);
std::fs::write(&src, &tar).unwrap();
let daemon = ScriptedDaemon::default();
let err = import_image_on(&daemon, &src).unwrap_err();
assert!(
matches!(&err, Error::Backend(m) if m.contains("truncated")),
"got {err:?}"
);
assert!(
daemon.saw_load.borrow().is_empty(),
"the daemon is never asked to load an archive that failed verification"
);
let _ = std::fs::remove_file(&src);
}
#[cfg(feature = "backend-oci")]
#[test]
fn import_reports_a_missing_archive_as_unreadable() {
let daemon = ScriptedDaemon::default();
let err =
import_image_on(&daemon, Path::new("/nonexistent/draupnir-no-such.tar")).unwrap_err();
assert!(
matches!(&err, Error::Backend(m) if m.contains("unreadable")),
"got {err:?}"
);
assert!(daemon.saw_load.borrow().is_empty());
}
#[cfg(feature = "backend-oci")]
#[test]
fn export_then_import_round_trip_preserves_the_digest() {
let dest =
std::env::temp_dir().join(format!("draupnir-roundtrip-{}.tar", std::process::id()));
let _ = std::fs::remove_file(&dest);
let image = "localhost/nordisk-spark-iceberg:4.1.2";
let daemon = ScriptedDaemon::default();
*daemon.export_bytes.borrow_mut() = Some(Ok(make_image_archive(
"blobs/sha256/05455e6bc39e",
&[image],
32 * 1024,
)));
let exported = export_image_on(&daemon, image, &dest).unwrap();
assert_eq!(exported.config_digest, "05455e6bc39e");
assert_eq!(exported.tags, [image]);
assert_eq!(exported.path, dest);
assert_eq!(
exported.bytes,
std::fs::metadata(&dest).unwrap().len(),
"the recorded size is the size on disk"
);
let target = ScriptedDaemon::default();
let imported = import_image_on(&target, &dest).unwrap();
assert_eq!(
imported.config_digest, exported.config_digest,
"the digest survives the round trip — same image in, same image out"
);
assert_eq!(imported.tags, exported.tags);
assert_eq!(
target.saw_load.borrow().as_slice(),
[dest.display().to_string()],
"the verified archive is the exact file handed to the daemon"
);
let _ = std::fs::remove_file(&dest);
}
#[cfg(feature = "backend-oci")]
#[test]
fn import_surfaces_a_daemon_load_failure() {
let src =
std::env::temp_dir().join(format!("draupnir-import-fail-{}.tar", std::process::id()));
std::fs::write(
&src,
make_image_archive("blobs/sha256/abc", &["x:1"], 1024),
)
.unwrap();
let daemon = ScriptedDaemon::default();
*daemon.load_result.borrow_mut() =
Some(Err(Error::Backend("no space left on device".into())));
let err = import_image_on(&daemon, &src).unwrap_err();
assert!(
matches!(&err, Error::Backend(m) if m.contains("no space left")),
"got {err:?}"
);
let _ = std::fs::remove_file(&src);
}
#[cfg(not(feature = "backend-oci"))]
#[test]
fn image_archive_verbs_are_unsupported_without_the_engine() {
let boot = ContainerBoot::new();
assert!(matches!(
boot.export_image("x:test", Path::new("/tmp/draupnir-x.tar")),
Err(Error::Unsupported(_))
));
assert!(matches!(
boot.import_image(Path::new("/tmp/draupnir-x.tar")),
Err(Error::Unsupported(_))
));
assert!(matches!(
inspect_image_archive(Path::new("/tmp/draupnir-x.tar")),
Err(Error::Unsupported(_))
));
}
}