use anyhow::{Context, Result, anyhow, bail};
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use std::io::{Read, Write};
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
const KIND_ROOTFS: &str = "heyvm.rootfs.v1";
const SCHEMA_VERSION: u32 = 1;
const ROOTFS_FILENAME: &str = "rootfs.ext4";
const ANN_PRIMITIVE: &str = "heyvm.primitive";
const ANN_IMAGE: &str = "heyvm.image";
const ANN_NOMINAL_SIZE: &str = "heyvm.nominal_size";
const PRIMITIVE_EXT4_RAW: &str = "ext4_raw";
const KIND_DOCKERFILE: &str = "heyvm.dockerfile.v1";
const DOCKERFILE_ENTRY: &str = "Dockerfile";
const CONTEXT_ENTRY: &str = "context.tar.gz";
const ANN_SIZE_MB: &str = "heyvm.size_mb";
const ANN_SOURCE: &str = "dockerfile.source";
const MAX_DOCKERFILE_BYTES: u64 = 1 << 20;
const HASH_CHUNK: usize = 1 << 20;
pub struct RegistryClient {
agent: ureq::Agent,
base: String,
api_key: Option<String>,
}
impl RegistryClient {
pub fn new(url: &str, api_key: Option<&str>, insecure: bool, timeout: Duration) -> Result<Self> {
let mut tls = ureq::native_tls::TlsConnector::builder();
if insecure {
tls.danger_accept_invalid_certs(true);
tls.danger_accept_invalid_hostnames(true);
}
let connector = tls.build().context("building the TLS connector")?;
let agent = ureq::AgentBuilder::new()
.timeout_connect(timeout)
.user_agent(concat!("heyctl/", env!("CARGO_PKG_VERSION")))
.tls_connector(Arc::new(connector))
.build();
Ok(Self {
agent,
base: normalize_url(url),
api_key: api_key.map(str::to_string),
})
}
pub fn url(&self) -> &str {
&self.base
}
pub fn has_credentials(&self) -> bool {
self.api_key.is_some()
}
fn request(&self, method: &str, path: &str) -> ureq::Request {
let req = self.agent.request(method, &format!("{}{}", self.base, path));
match &self.api_key {
Some(k) => req.set("Authorization", &format!("Bearer {k}")),
None => req,
}
}
fn raw(&self, req: ureq::Request) -> Result<(u16, String)> {
match req.call() {
Ok(resp) => {
let status = resp.status();
Ok((status, resp.into_string().unwrap_or_default()))
}
Err(ureq::Error::Status(code, resp)) => {
Ok((code, resp.into_string().unwrap_or_default()))
}
Err(ureq::Error::Transport(t)) => Err(anyhow!(
"cannot reach the artifact store at {} ({t}) — is `art serve` running, and is \
the url right?",
self.base
)),
}
}
fn send(&self, method: &str, path: &str) -> Result<(u16, String)> {
let (code, body) = self.raw(self.request(method, path))?;
if code >= 400 {
return Err(self.api_error(code, &body));
}
Ok((code, body))
}
fn json(&self, method: &str, path: &str) -> Result<Value> {
let (_, text) = self.send(method, path)?;
if text.trim().is_empty() {
return Ok(Value::Null);
}
serde_json::from_str(&text)
.with_context(|| format!("the store's answer to {method} {path} was not JSON"))
}
fn api_error(&self, code: u16, body: &str) -> anyhow::Error {
let detail = serde_json::from_str::<Value>(body)
.ok()
.and_then(|v| v.get("error")?.as_str().map(str::to_string))
.unwrap_or_else(|| body.trim().chars().take(300).collect());
match code {
401 if self.api_key.is_some() => anyhow!(
"the store rejected this API key (HTTP 401) — check it against the store's \
ART_API_KEY, or re-run `heyctl artifact login`"
),
401 => anyhow!(
"this store requires an API key (HTTP 401) — run `heyctl artifact login`, \
or pass --api-key"
),
403 => anyhow!(
"the store refused the write (HTTP 403){} — a store started with \
ART_READ_ONLY rejects every mutating route",
if detail.is_empty() { String::new() } else { format!(": {detail}") }
),
_ if detail.is_empty() => anyhow!("the store returned HTTP {code}"),
_ => anyhow!("{detail} (HTTP {code})"),
}
}
pub fn healthz(&self) -> Result<()> {
let (code, _) = self.send("GET", "/healthz")?;
if code == 200 {
Ok(())
} else {
bail!("GET /healthz answered HTTP {code}; that does not look like an artifact store")
}
}
pub fn tags(&self) -> Result<Value> {
self.json("GET", "/tags")
}
pub fn usage(&self) -> Result<Value> {
self.json("GET", "/usage")
}
pub fn manifest(&self, reference: &str) -> Result<Value> {
self.json("GET", &format!("/manifests/{}", escape(reference)))
}
pub fn blob_exists(&self, digest: &str) -> Result<Option<u64>> {
let req = self.request("HEAD", &format!("/blobs/{}", escape(digest)));
match req.call() {
Ok(resp) => Ok(Some(
resp.header("Content-Length")
.and_then(|v| v.parse().ok())
.unwrap_or(0),
)),
Err(ureq::Error::Status(404, _)) => Ok(None),
Err(ureq::Error::Status(code, resp)) => {
Err(self.api_error(code, &resp.into_string().unwrap_or_default()))
}
Err(ureq::Error::Transport(t)) => Err(anyhow!(
"cannot reach the artifact store at {} ({t})",
self.base
)),
}
}
pub fn put_blob(&self, digest: &str, path: &Path, size: u64) -> Result<bool> {
let file = std::fs::File::open(path)
.with_context(|| format!("opening {}", path.display()))?;
let req = self
.request("PUT", &format!("/blobs/{}", escape(digest)))
.set("Content-Type", "application/octet-stream")
.set("Content-Length", &size.to_string());
match req.send(file) {
Ok(resp) => Ok(resp.status() == 201),
Err(ureq::Error::Status(code, resp)) => {
Err(self.api_error(code, &resp.into_string().unwrap_or_default()))
}
Err(ureq::Error::Transport(t)) => Err(anyhow!(
"the upload to {} failed ({t}) — nothing was tagged, so the store is unchanged",
self.base
)),
}
}
pub fn put_manifest(&self, manifest: &Value) -> Result<String> {
let req = self
.request("PUT", "/manifests")
.set("Content-Type", "application/json");
let (code, body) = match req.send_json(manifest) {
Ok(resp) => (resp.status(), resp.into_string().unwrap_or_default()),
Err(ureq::Error::Status(code, resp)) => {
return Err(self.api_error(code, &resp.into_string().unwrap_or_default()));
}
Err(ureq::Error::Transport(t)) => {
return Err(anyhow!("cannot reach the artifact store at {} ({t})", self.base));
}
};
if code >= 400 {
return Err(self.api_error(code, &body));
}
let v: Value = serde_json::from_str(&body)
.context("the store's answer to PUT /manifests was not JSON")?;
v.get("digest")
.and_then(|d| d.as_str())
.map(str::to_string)
.ok_or_else(|| anyhow!("the store stored the manifest but reported no digest"))
}
pub fn put_tag(&self, name: &str, digest: &str) -> Result<()> {
let req = self
.request("PUT", &format!("/tags/{}", escape(name)))
.set("Content-Type", "text/plain");
match req.send_string(digest) {
Ok(_) => Ok(()),
Err(ureq::Error::Status(code, resp)) => {
Err(self.api_error(code, &resp.into_string().unwrap_or_default()))
}
Err(ureq::Error::Transport(t)) => {
Err(anyhow!("cannot reach the artifact store at {} ({t})", self.base))
}
}
}
pub fn delete_tag(&self, name: &str) -> Result<()> {
self.send("DELETE", &format!("/tags/{}", escape(name))).map(|_| ())
}
}
pub fn rootfs_manifest(digest: &str, size: u64, image: &str) -> Value {
json!({
"schema": SCHEMA_VERSION,
"kind": KIND_ROOTFS,
"entries": [{ "name": ROOTFS_FILENAME, "digest": digest, "size": size }],
"annotations": {
ANN_IMAGE: image,
ANN_NOMINAL_SIZE: size.to_string(),
ANN_PRIMITIVE: PRIMITIVE_EXT4_RAW,
}
})
}
pub fn dockerfile_manifest(
dockerfile: (&str, u64),
context: Option<(&str, u64)>,
image_name: Option<&str>,
size_mb: Option<u64>,
source: Option<&str>,
) -> Value {
let mut entries = vec![json!({
"name": DOCKERFILE_ENTRY,
"digest": dockerfile.0,
"size": dockerfile.1,
})];
if let Some((digest, size)) = context {
entries.push(json!({ "name": CONTEXT_ENTRY, "digest": digest, "size": size }));
}
let mut annotations = serde_json::Map::new();
if let Some(n) = image_name {
annotations.insert(ANN_IMAGE.into(), Value::String(n.to_string()));
}
if let Some(mb) = size_mb {
annotations.insert(ANN_SIZE_MB.into(), Value::String(mb.to_string()));
}
if let Some(s) = source {
annotations.insert(ANN_SOURCE.into(), Value::String(s.to_string()));
}
json!({
"schema": SCHEMA_VERSION,
"kind": KIND_DOCKERFILE,
"entries": entries,
"annotations": annotations,
})
}
pub fn pack_context(dir: &Path, dest: &Path) -> Result<u64> {
if !dir.is_dir() {
bail!("build context {} is not a directory", dir.display());
}
let out = std::fs::File::create(dest)
.with_context(|| format!("creating {}", dest.display()))?;
let gz = flate2::write::GzEncoder::new(out, flate2::Compression::default());
let mut builder = tar::Builder::new(gz);
builder.mode(tar::HeaderMode::Deterministic);
builder
.append_dir_all(".", dir)
.with_context(|| format!("packing {}", dir.display()))?;
let gz = builder.into_inner().context("finishing the context archive")?;
let mut out = gz.finish().context("compressing the context archive")?;
out.flush().context("flushing the context archive")?;
out.sync_all().context("syncing the context archive")?;
Ok(std::fs::metadata(dest)
.with_context(|| format!("reading {}", dest.display()))?
.len())
}
pub fn check_dockerfile_size(path: &Path, size: u64) -> Result<()> {
if size > MAX_DOCKERFILE_BYTES {
bail!(
"{} is {}; a Dockerfile is expected to be under {}. Did you mean --context?",
path.display(),
crate::output::bytes(size),
crate::output::bytes(MAX_DOCKERFILE_BYTES),
);
}
Ok(())
}
pub fn hash_file(path: &Path, mut progress: impl FnMut(u64, u64)) -> Result<(String, u64)> {
let mut file =
std::fs::File::open(path).with_context(|| format!("opening {}", path.display()))?;
let total = file.metadata().map(|m| m.len()).unwrap_or(0);
let mut hasher = Sha256::new();
let mut buf = vec![0u8; HASH_CHUNK];
let mut read_total: u64 = 0;
loop {
let n = file
.read(&mut buf)
.with_context(|| format!("reading {}", path.display()))?;
if n == 0 {
break;
}
hasher.update(&buf[..n]);
read_total += n as u64;
progress(read_total, total);
}
Ok((hex(&hasher.finalize()), read_total))
}
pub fn heyvm_image_path(name: &str) -> Result<std::path::PathBuf> {
let dir = match std::env::var("MVM_DATA_DIR") {
Ok(d) if !d.trim().is_empty() => std::path::PathBuf::from(d),
_ => dirs::home_dir()
.context("cannot find a home directory to look for heyvm's images in")?
.join(".heyo"),
}
.join("images")
.join("firecracker");
let stem = name.strip_suffix(".ext4").unwrap_or(name);
let path = dir.join(format!("{stem}.ext4"));
if !path.exists() {
bail!(
"no heyvm image called {stem:?} — looked for {}. `heyvm mvm images` lists what \
is built, or give a path to an .ext4 file instead",
path.display()
);
}
Ok(path)
}
pub fn is_valid_tag(name: &str) -> bool {
if name.is_empty() || name.len() > 128 {
return false;
}
let first = name.as_bytes()[0];
if first == b'-' || first == b'.' {
return false;
}
name.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'_' || b == b'-' || b == b'.')
}
pub fn default_tag_for(path: &Path) -> Option<String> {
let stem = path.file_stem()?.to_str()?;
is_valid_tag(stem).then(|| stem.to_string())
}
fn hex(bytes: &[u8]) -> String {
let mut out = String::with_capacity(bytes.len() * 2);
for b in bytes {
out.push_str(&format!("{b:02x}"));
}
out
}
fn normalize_url(s: &str) -> String {
let s = s.trim();
let with_scheme = if s.contains("://") {
s.to_string()
} else {
format!("http://{s}")
};
with_scheme.trim_end_matches('/').to_string()
}
fn escape(segment: &str) -> String {
let mut out = String::with_capacity(segment.len());
for b in segment.bytes() {
match b {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
out.push(b as char)
}
_ => out.push_str(&format!("%{b:02X}")),
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_store_url_gets_a_scheme_and_loses_its_trailing_slash() {
assert_eq!(normalize_url("localhost:8080"), "http://localhost:8080");
assert_eq!(normalize_url("http://art:8080/"), "http://art:8080");
assert_eq!(normalize_url(" https://art.example.com "), "https://art.example.com");
}
#[test]
fn a_manifest_matches_what_art_heyvm_import_writes() {
let m = rootfs_manifest("c74abee2", 4096, "debian-hermes");
assert_eq!(m["kind"], KIND_ROOTFS);
assert_eq!(m["schema"], 1);
assert_eq!(m["entries"][0]["name"], ROOTFS_FILENAME);
assert_eq!(m["entries"][0]["digest"], "c74abee2");
assert_eq!(m["entries"][0]["size"], 4096);
assert_eq!(m["annotations"][ANN_IMAGE], "debian-hermes");
assert_eq!(m["annotations"][ANN_PRIMITIVE], PRIMITIVE_EXT4_RAW);
assert_eq!(m["annotations"][ANN_NOMINAL_SIZE], "4096");
}
#[test]
fn a_dockerfile_manifest_matches_what_art_dockerfile_put_writes() {
let m = dockerfile_manifest(
("a1b2", 812),
Some(("c3d4", 40213)),
Some("web"),
Some(4096),
None,
);
assert_eq!(m["kind"], KIND_DOCKERFILE);
assert_eq!(m["schema"], 1);
assert_eq!(m["entries"][0]["name"], DOCKERFILE_ENTRY);
assert_eq!(m["entries"][0]["digest"], "a1b2");
assert_eq!(m["entries"][1]["name"], CONTEXT_ENTRY);
assert_eq!(m["entries"][1]["size"], 40213);
assert_eq!(m["annotations"][ANN_IMAGE], "web");
assert_eq!(m["annotations"][ANN_SIZE_MB], "4096");
}
#[test]
fn an_unset_annotation_is_absent_rather_than_empty() {
let m = dockerfile_manifest(("a1b2", 812), None, None, None, None);
assert_eq!(m["annotations"].as_object().unwrap().len(), 0);
assert_eq!(m["entries"].as_array().unwrap().len(), 1);
assert!(m["annotations"].get(ANN_IMAGE).is_none());
}
#[test]
fn packing_the_same_tree_twice_gives_the_same_bytes() {
let dir = std::env::temp_dir().join(format!("heyctl-pack-{}", std::process::id()));
let ctx = dir.join("ctx");
std::fs::create_dir_all(ctx.join("app")).unwrap();
std::fs::write(ctx.join("app/main.rs"), b"fn main() {}").unwrap();
std::fs::write(ctx.join("README"), b"hi").unwrap();
let first = dir.join("1.tar.gz");
let second = dir.join("2.tar.gz");
assert!(pack_context(&ctx, &first).unwrap() > 0);
pack_context(&ctx, &second).unwrap();
assert_eq!(
std::fs::read(&first).unwrap(),
std::fs::read(&second).unwrap()
);
let gz = flate2::read::GzDecoder::new(std::fs::File::open(&first).unwrap());
let names: Vec<String> = tar::Archive::new(gz)
.entries()
.unwrap()
.map(|e| e.unwrap().path().unwrap().display().to_string())
.collect();
assert!(names.iter().any(|n| n.ends_with("app/main.rs")), "{names:?}");
assert!(names.iter().any(|n| n.ends_with("README")), "{names:?}");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_context_passed_as_the_dockerfile_is_caught_before_the_upload() {
let p = Path::new("Dockerfile");
assert!(check_dockerfile_size(p, 4096).is_ok());
let e = check_dockerfile_size(p, MAX_DOCKERFILE_BYTES + 1).unwrap_err();
assert!(e.to_string().contains("--context"), "{e}");
}
#[test]
fn tags_follow_the_stores_rules() {
assert!(is_valid_tag("debian-hermes"));
assert!(is_valid_tag("ubuntu-24.04"));
assert!(is_valid_tag("web_v2"));
assert!(!is_valid_tag(""));
assert!(!is_valid_tag("-leading-dash"));
assert!(!is_valid_tag(".hidden"));
assert!(!is_valid_tag("a/b"));
assert!(!is_valid_tag("has space"));
}
#[test]
fn a_files_own_name_is_its_default_tag() {
assert_eq!(
default_tag_for(Path::new("/home/x/.heyo/images/firecracker/artifacts.ext4")),
Some("artifacts".to_string())
);
assert_eq!(default_tag_for(Path::new("/tmp/.hidden.ext4")), None);
}
#[test]
fn hashing_a_file_reports_the_digest_and_the_bytes_it_covered() {
let dir = std::env::temp_dir().join(format!("heyctl-art-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("blob");
std::fs::write(&path, b"hello").unwrap();
let mut last = (0, 0);
let (digest, size) = hash_file(&path, |done, total| last = (done, total)).unwrap();
assert_eq!(
digest,
"2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"
);
assert_eq!(size, 5);
assert_eq!(last, (5, 5));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn path_segments_are_escaped() {
assert_eq!(escape("debian-hermes"), "debian-hermes");
assert_eq!(escape("a/b"), "a%2Fb");
}
}