use anyhow::Context;
use zygo_core::image::{Platform, PullProgress, Reference, RegistryClient, Store};
use crate::cli::{Cli, ImageCommand};
use crate::output::{self, Style};
pub fn pull(cli: &Cli, image: &str, platform: Option<&str>) -> anyhow::Result<u8> {
let reference: Reference = image
.parse()
.with_context(|| format!("cannot pull `{image}`"))?;
let paths = super::paths(cli);
paths.ensure()?;
let store = Store::new(paths);
let mut client = RegistryClient::new(store)?;
if let Some(p) = platform {
client = client.with_platform(parse_platform(p)?);
}
let style = Style::stdout();
let quiet = cli.json;
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
let entry = runtime.block_on(client.pull(&reference, |event| {
if quiet {
return;
}
match event {
PullProgress::Resolving { reference } => {
println!("resolving {reference}");
}
PullProgress::LayerStart { digest, size } => {
println!(
" {} {} {}",
style.dim("pull"),
short(&digest),
output::human_bytes(size)
);
}
PullProgress::LayerCached { digest } => {
println!(" {} {}", style.dim("cached"), short(&digest));
}
PullProgress::Unpacking { digest } => {
println!(" {} {}", style.dim("unpack"), short(&digest));
}
PullProgress::LayerDone { .. } => {}
PullProgress::Done { layers, bytes } => {
println!(
"{} {layers} layers, {}",
style.green("done"),
output::human_bytes(bytes)
);
}
}
}))?;
#[cfg(target_os = "linux")]
{
let store = Store::new(super::paths(cli));
let compiled = zygo_core::bytecode::ensure_now(&store, &entry)?;
if compiled.built && !quiet {
println!(
"{} Python bytecode, as one more layer",
style.dim("compiled")
);
}
}
if cli.json {
output::json(&entry)?;
}
Ok(0)
}
pub fn list(cli: &Cli) -> anyhow::Result<u8> {
let store = Store::new(super::paths(cli));
let images = store.list();
if cli.json {
output::json(&images)?;
return Ok(0);
}
if images.is_empty() {
println!("no images pulled yet — try `zygo pull python:3.12-slim`");
return Ok(0);
}
let rows: Vec<Vec<String>> = images
.iter()
.map(|i| {
vec![
i.reference.clone(),
short(&i.manifest),
i.layers.len().to_string(),
output::human_bytes(i.size),
output::human_age(i.pulled_at),
]
})
.collect();
print!(
"{}",
output::table(&["reference", "digest", "layers", "size", "pulled"], &rows)
);
Ok(0)
}
pub fn maintain(cli: &Cli, command: &ImageCommand) -> anyhow::Result<u8> {
match command {
ImageCommand::Prune {
dry_run,
unused_for,
blobs,
} => prune(
cli,
&PruneOptions {
dry_run: *dry_run,
unused_for: unused_for.map(|d| d.get()),
blobs: *blobs,
},
),
ImageCommand::Rm { references } => rm(cli, references),
}
}
fn rm(cli: &Cli, references: &[String]) -> anyhow::Result<u8> {
use zygo_core::supervisor::client::Client;
use zygo_core::supervisor::protocol::{Request, Response};
let paths = super::paths(cli);
let store = Store::new(paths.clone());
let style = Style::stdout();
let wanted: Vec<Reference> = references
.iter()
.map(|r| {
r.parse::<Reference>()
.with_context(|| format!("cannot remove `{r}`"))
})
.collect::<anyhow::Result<_>>()?;
let warm = match Client::connect(&paths) {
Ok(mut client) => match client.send(&Request::List)? {
Response::Functions { functions } => functions,
_ => Vec::new(),
},
Err(_) => Vec::new(),
};
for reference in &wanted {
let key = reference.to_string();
let on_it: Vec<&str> = warm
.iter()
.filter(|f| {
f.image
.parse::<Reference>()
.is_ok_and(|r| r.to_string() == key)
})
.map(|f| f.name.as_str())
.collect();
anyhow::ensure!(
on_it.is_empty(),
"`{key}` is in use: {} running on it\n → stop it first: zygo stop {}",
on_it.join(", "),
on_it.join(" ")
);
}
let mut removed = Vec::new();
for reference in &wanted {
let gone = store.remove(reference)?;
anyhow::ensure!(
!gone.is_empty(),
"no image `{reference}` in the store\n → `zygo images` lists what is here"
);
removed.extend(gone);
}
let opts = PruneOptions {
dry_run: false,
unused_for: None,
blobs: false,
};
if cli.json {
let mut pruned = serde_json::Value::Null;
let code = prune_with(cli, &opts, Some(&mut pruned))?;
output::json(&serde_json::json!({
"removed": removed.iter().map(|e| e.reference.clone()).collect::<Vec<_>>(),
"pruned": pruned,
}))?;
return Ok(code);
}
for entry in &removed {
println!("{} {}", style.dim("removed"), entry.reference);
}
prune(cli, &opts)
}
struct PruneOptions {
dry_run: bool,
unused_for: Option<std::time::Duration>,
blobs: bool,
}
struct Category {
name: &'static str,
items: Vec<Reclaimable>,
}
struct Reclaimable {
label: String,
dir: std::path::PathBuf,
blob: Option<std::path::PathBuf>,
bytes: u64,
}
fn prune(cli: &Cli, opts: &PruneOptions) -> anyhow::Result<u8> {
prune_with(cli, opts, None)
}
fn prune_with(
cli: &Cli,
opts: &PruneOptions,
into: Option<&mut serde_json::Value>,
) -> anyhow::Result<u8> {
let dry_run = opts.dry_run;
let paths = super::paths(cli);
let store = Store::new(paths.clone());
let layers = Category {
name: "layers",
items: store
.unreferenced_layers()?
.into_iter()
.map(|digest| {
let dir = paths.layers().join(digest.trim_start_matches("sha256:"));
let blob = store.blob_path(&digest).ok();
let bytes = dir_size(&dir)
+ blob
.as_ref()
.and_then(|b| b.metadata().ok())
.map_or(0, |m| m.len());
Reclaimable {
label: short(&digest),
dir,
blob,
bytes,
}
})
.collect(),
};
let cache = |name, dirs: Vec<std::path::PathBuf>| Category {
name,
items: dirs
.into_iter()
.map(|dir| Reclaimable {
label: key_of(&dir),
bytes: dir_size(&dir),
blob: None,
dir,
})
.collect(),
};
let abandoned = abandoned_roots(&paths.tmp());
let mut categories = vec![
layers,
cache("flattened rootfs caches", store.unreferenced_flat()?),
cache("venvs", zygo_core::venv::unreferenced(&store)?),
cache(
"system layer records",
zygo_core::derive::unreferenced(&store)?,
),
cache("abandoned sandbox roots", abandoned),
];
if let Some(window) = opts.unused_for {
let cutoff = std::time::SystemTime::now()
.checked_sub(window)
.unwrap_or(std::time::UNIX_EPOCH);
categories.push(cache(
"unused venvs",
zygo_core::venv::unused_since(&store, cutoff)?,
));
categories.push(cache(
"unused flattened rootfs caches",
store.flat_unused_since(cutoff)?,
));
}
if opts.blobs {
categories.push(Category {
name: "compressed layer blobs",
items: store
.droppable_blobs()?
.into_iter()
.map(|(digest, blob)| Reclaimable {
label: short(&digest),
bytes: blob.metadata().map_or(0, |m| m.len()),
dir: blob,
blob: None,
})
.collect(),
});
}
let total: u64 = categories
.iter()
.flat_map(|c| &c.items)
.map(|i| i.bytes)
.sum();
let count: usize = categories.iter().map(|c| c.items.len()).sum();
if cli.json {
let report = serde_json::json!({
"dry_run": dry_run,
"bytes": total,
"items": count,
"categories": categories.iter().map(|c| serde_json::json!({
"name": c.name,
"count": c.items.len(),
"bytes": c.items.iter().map(|i| i.bytes).sum::<u64>(),
"items": c.items.iter().map(|i| i.label.clone()).collect::<Vec<_>>(),
})).collect::<Vec<_>>(),
});
match into {
Some(slot) => *slot = report,
None => output::json(&report)?,
}
if !dry_run {
remove(&categories);
}
return Ok(0);
}
let style = Style::stdout();
if count == 0 {
println!("nothing to prune");
print_kept(&store, opts, &style);
return Ok(0);
}
for category in &categories {
for item in &category.items {
println!(
"{} {} {}",
style.dim(if dry_run { "would remove" } else { "removing" }),
item.label,
style.dim(&output::human_bytes(item.bytes)),
);
}
}
if !dry_run {
remove(&categories);
}
for category in &categories {
if category.items.is_empty() {
continue;
}
let bytes: u64 = category.items.iter().map(|i| i.bytes).sum();
println!(
" {} {} across {} {}",
if dry_run { "would free" } else { "freed" },
output::human_bytes(bytes),
category.items.len(),
category.name,
);
}
println!(
"{} {} in total",
if dry_run { "would free" } else { "freed" },
output::human_bytes(total)
);
print_kept(&store, opts, &style);
Ok(0)
}
fn print_kept(store: &Store, opts: &PruneOptions, style: &Style) {
let now = std::time::SystemTime::now();
if opts.unused_for.is_none() {
let venvs = zygo_core::venv::unused_since(store, now).unwrap_or_default();
let flat = store.flat_unused_since(now).unwrap_or_default();
for (name, dirs) in [("venvs", venvs), ("flattened rootfs caches", flat)] {
if dirs.is_empty() {
continue;
}
let bytes: u64 = dirs.iter().map(|d| dir_size(d)).sum();
let age = dirs
.iter()
.filter_map(|d| marker_age(d, now))
.max()
.map_or_else(|| "an unknown time".to_string(), human_duration);
println!(
"{}",
style.dim(&format!(
" kept {} across {} {name}, the oldest unused for {age} \
(`--unused-for <duration>` collects those)",
output::human_bytes(bytes),
dirs.len(),
))
);
}
}
if !opts.blobs {
let blobs = store.droppable_blobs().unwrap_or_default();
if !blobs.is_empty() {
let bytes: u64 = blobs
.iter()
.map(|(_, b)| b.metadata().map_or(0, |m| m.len()))
.sum();
println!(
"{}",
style.dim(&format!(
" kept {} of compressed blobs beside {} unpacked layers \
(`--blobs` drops them)",
output::human_bytes(bytes),
blobs.len(),
))
);
}
}
}
fn abandoned_roots(tmp: &std::path::Path) -> Vec<std::path::PathBuf> {
let Ok(entries) = std::fs::read_dir(tmp) else {
return Vec::new();
};
let mut out: Vec<std::path::PathBuf> = entries
.filter_map(|e| e.ok())
.filter_map(|entry| {
let name = entry.file_name();
let pid: i32 = name
.to_str()?
.strip_prefix("root-")?
.split('-')
.next()?
.parse()
.ok()?;
let path = entry.path();
let settled = entry
.metadata()
.and_then(|m| m.modified())
.is_ok_and(|t| t.elapsed().is_ok_and(|age| age.as_secs() > 60));
(path.is_dir() && settled && !is_running(pid)).then_some(path)
})
.collect();
out.sort();
out
}
#[cfg(unix)]
fn is_running(pid: i32) -> bool {
unsafe {
libc::kill(pid, 0) == 0
|| std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
}
}
#[cfg(not(unix))]
fn is_running(_pid: i32) -> bool {
true
}
fn marker_age(dir: &std::path::Path, now: std::time::SystemTime) -> Option<std::time::Duration> {
let used = std::fs::read_dir(dir)
.ok()?
.filter_map(|e| e.ok())
.filter(|e| e.file_name().to_string_lossy().starts_with(".zygo-"))
.filter_map(|e| e.metadata().ok()?.modified().ok())
.max()?;
now.duration_since(used).ok()
}
fn human_duration(d: std::time::Duration) -> String {
let secs = d.as_secs();
let (n, unit) = match secs {
0..=119 => return "under 2 minutes".to_string(),
120..=7199 => (secs / 60, "minutes"),
7200..=172_799 => (secs / 3600, "hours"),
_ => (secs / 86_400, "days"),
};
format!("{n} {unit}")
}
fn remove(categories: &[Category]) {
for item in categories.iter().flat_map(|c| &c.items) {
std::fs::remove_dir_all(&item.dir).ok();
if let Some(blob) = &item.blob {
std::fs::remove_file(blob).ok();
}
}
}
fn key_of(dir: &std::path::Path) -> String {
dir.file_name()
.map(|n| n.to_string_lossy().chars().take(12).collect())
.unwrap_or_else(|| dir.display().to_string())
}
fn dir_size(dir: &std::path::Path) -> u64 {
let Ok(entries) = std::fs::read_dir(dir) else {
return 0;
};
entries
.filter_map(|e| e.ok())
.map(|e| match e.metadata() {
Ok(m) if m.is_dir() => dir_size(&e.path()),
Ok(m) => m.len(),
Err(_) => 0,
})
.sum()
}
fn short(digest: &str) -> String {
digest
.trim_start_matches("sha256:")
.chars()
.take(12)
.collect()
}
fn parse_platform(s: &str) -> anyhow::Result<Platform> {
let mut parts = s.split('/');
let os = parts
.next()
.filter(|p| !p.is_empty())
.context("--platform expects os/arch, e.g. linux/amd64")?;
let architecture = parts
.next()
.filter(|p| !p.is_empty())
.context("--platform expects os/arch, e.g. linux/amd64")?;
Ok(Platform {
os: os.to_string(),
architecture: architecture.to_string(),
variant: parts.next().map(str::to_string),
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_json_listing_is_the_stable_interface() {
let entry = zygo_core::image::store::ImageEntry {
reference: "node:22-alpine".into(),
manifest: format!("sha256:{}", "b".repeat(64)),
config: format!("sha256:{}", "c".repeat(64)),
layers: vec![format!("sha256:{}", "d".repeat(64))],
size: 51_200,
pulled_at: 1_700_000_000,
index: Some(format!("sha256:{}", "e".repeat(64))),
platform: Some("linux/arm64".into()),
};
let v = serde_json::to_value(vec![entry]).unwrap();
let first = &v[0];
for key in [
"reference",
"manifest",
"config",
"layers",
"size",
"pulled_at",
] {
assert!(first.get(key).is_some(), "`{key}` is part of the contract");
}
assert_eq!(
first["reference"], "node:22-alpine",
"verbatim, so `==` works and `contains` is not needed"
);
assert_eq!(first["layers"].as_array().unwrap().len(), 1);
assert_eq!(first["size"], 51_200);
assert_eq!(first["index"].as_str().unwrap().len(), 71);
let without = zygo_core::image::store::ImageEntry {
index: None,
platform: None,
..serde_json::from_value(first.clone()).unwrap()
};
let v = serde_json::to_value(without).unwrap();
assert!(v.get("index").is_none());
}
#[test]
fn digests_are_shortened_the_way_registries_print_them() {
assert_eq!(short(&format!("sha256:{}", "a".repeat(64))), "aaaaaaaaaaaa");
assert_eq!(short("abc"), "abc");
}
#[test]
fn platform_strings_parse() {
let p = parse_platform("linux/amd64").unwrap();
assert_eq!(p.os, "linux");
assert_eq!(p.architecture, "amd64");
assert_eq!(p.variant, None);
let p = parse_platform("linux/arm/v7").unwrap();
assert_eq!(p.variant.as_deref(), Some("v7"));
assert!(parse_platform("linux").is_err());
assert!(parse_platform("linux/").is_err());
}
#[test]
fn directory_sizes_are_summed_recursively() {
let tmp = tempfile::tempdir().unwrap();
std::fs::write(tmp.path().join("a"), vec![0u8; 100]).unwrap();
std::fs::create_dir(tmp.path().join("sub")).unwrap();
std::fs::write(tmp.path().join("sub/b"), vec![0u8; 50]).unwrap();
assert_eq!(dir_size(tmp.path()), 150);
assert_eq!(dir_size(&tmp.path().join("missing")), 0);
}
}