Skip to main content

isb_core/registry/
oci.rs

1//! OCI images on the wire: an OCI image layout in a tar (what BuildKit's
2//! `type=oci` exporter writes) and the distribution API to push it, resolve
3//! tags and delete manifests.
4//!
5//! The layout comes out of a build sandbox, so it is untrusted: names are
6//! looked up only at the paths a digest dictates, sizes are bounded by the
7//! file, manifests by [`MAX_MANIFEST`], and every manifest's digest is
8//! checked here while the registry checks every blob's on upload.
9
10use std::collections::HashMap;
11use std::fs::File;
12use std::io::{Read, Seek, SeekFrom};
13use std::path::{Path, PathBuf};
14use std::time::Duration;
15
16use serde::Deserialize;
17use serde_json::Value;
18
19use crate::error::{Error, Result};
20
21pub const MT_OCI_INDEX: &str = "application/vnd.oci.image.index.v1+json";
22pub const MT_OCI_MANIFEST: &str = "application/vnd.oci.image.manifest.v1+json";
23pub const MT_DOCKER_LIST: &str = "application/vnd.docker.distribution.manifest.list.v2+json";
24pub const MT_DOCKER_MANIFEST: &str = "application/vnd.docker.distribution.manifest.v2+json";
25
26/// What a manifest request accepts: every kind isb pushes.
27const ACCEPT: &str = "application/vnd.oci.image.index.v1+json, application/vnd.oci.image.manifest.v1+json, application/vnd.docker.distribution.manifest.list.v2+json, application/vnd.docker.distribution.manifest.v2+json";
28
29/// The largest manifest or index read into memory.
30pub const MAX_MANIFEST: u64 = 4 << 20;
31/// How deep an index may nest.
32const MAX_DEPTH: usize = 3;
33
34/// `sha256:` and 64 lowercase hex digits: the only digest isb handles, and
35/// safe to put in a URL or a path.
36pub fn valid_digest(d: &str) -> bool {
37    d.strip_prefix("sha256:").is_some_and(|h| {
38        h.len() == 64
39            && h.bytes()
40                .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
41    })
42}
43
44/// The `sha256:` digest of some bytes.
45pub fn digest_of(b: &[u8]) -> String {
46    let d = ring::digest::digest(&ring::digest::SHA256, b);
47    let mut s = String::from("sha256:");
48    for x in d.as_ref() {
49        s.push_str(&format!("{x:02x}"));
50    }
51    s
52}
53
54/// A content descriptor.
55#[derive(Debug, Clone, Deserialize)]
56pub struct Descriptor {
57    #[serde(rename = "mediaType", default)]
58    pub media_type: String,
59    pub digest: String,
60    pub size: u64,
61}
62
63#[derive(Deserialize)]
64struct Index {
65    #[serde(default)]
66    manifests: Vec<Descriptor>,
67}
68
69#[derive(Deserialize)]
70struct Manifest {
71    config: Descriptor,
72    #[serde(default)]
73    layers: Vec<Descriptor>,
74}
75
76fn is_index(mt: &str) -> bool {
77    mt == MT_OCI_INDEX || mt == MT_DOCKER_LIST
78}
79
80fn is_manifest(mt: &str) -> bool {
81    mt == MT_OCI_MANIFEST || mt == MT_DOCKER_MANIFEST
82}
83
84// ---------------------------------------------------------------------------
85// The layout tar
86// ---------------------------------------------------------------------------
87
88/// An OCI image layout inside a tar file, indexed without unpacking it.
89pub struct Layout {
90    path: PathBuf,
91    /// Regular files: normalized path -> (data offset, size).
92    entries: HashMap<String, (u64, u64)>,
93}
94
95fn octal(field: &[u8]) -> Option<u64> {
96    // GNU base-256 for sizes over 8 GiB.
97    if field.first().is_some_and(|b| b & 0x80 != 0) {
98        let mut v: u64 = (field[0] & 0x7f) as u64;
99        for b in &field[1..] {
100            v = v.checked_mul(256)?.checked_add(*b as u64)?;
101        }
102        return Some(v);
103    }
104    let s: String = field
105        .iter()
106        .take_while(|b| **b != 0)
107        .map(|b| *b as char)
108        .collect();
109    let s = s.trim();
110    if s.is_empty() {
111        return Some(0);
112    }
113    u64::from_str_radix(s, 8).ok()
114}
115
116fn cstr(field: &[u8]) -> String {
117    let end = field.iter().position(|b| *b == 0).unwrap_or(field.len());
118    String::from_utf8_lossy(&field[..end]).into_owned()
119}
120
121fn normalize(p: &str) -> String {
122    p.trim_start_matches("./")
123        .trim_start_matches('/')
124        .to_string()
125}
126
127/// `path=` from a pax extended header.
128fn pax_path(data: &[u8]) -> Option<String> {
129    let text = String::from_utf8_lossy(data);
130    let mut rest: &str = &text;
131    let mut found = None;
132    while !rest.is_empty() {
133        let (len, _) = rest.split_once(' ')?;
134        let n: usize = len.parse().ok()?;
135        if n == 0 || n > rest.len() {
136            return found;
137        }
138        let rec = &rest[..n];
139        if let Some((_, kv)) = rec.split_once(' ') {
140            if let Some(v) = kv.strip_prefix("path=") {
141                found = Some(v.trim_end_matches('\n').to_string());
142            }
143        }
144        rest = &rest[n..];
145    }
146    found
147}
148
149impl Layout {
150    /// Index the tar at `path`.
151    pub fn open(path: &Path) -> Result<Layout> {
152        let bad = |why: String| Error::invalid(format!("image archive {}: {why}", path.display()));
153        let mut f = File::open(path)?;
154        let len = f.metadata()?.len();
155        let mut entries = HashMap::new();
156        let mut pos: u64 = 0;
157        let mut long_name: Option<String> = None;
158        let mut hdr = [0u8; 512];
159        loop {
160            if pos + 512 > len {
161                break;
162            }
163            f.seek(SeekFrom::Start(pos))?;
164            f.read_exact(&mut hdr)?;
165            if hdr.iter().all(|b| *b == 0) {
166                break;
167            }
168            let size = octal(&hdr[124..136]).ok_or_else(|| bad("bad size field".into()))?;
169            let data = pos + 512;
170            if data.checked_add(size).is_none_or(|end| end > len) {
171                return Err(bad("an entry runs past the end".into()));
172            }
173            let kind = hdr[156];
174            let mut name = cstr(&hdr[0..100]);
175            if &hdr[257..262] == b"ustar" {
176                let prefix = cstr(&hdr[345..500]);
177                if !prefix.is_empty() {
178                    name = format!("{prefix}/{name}");
179                }
180            }
181            match kind {
182                b'x' | b'L' => {
183                    if size > 64 << 10 {
184                        return Err(bad("an oversized extended header".into()));
185                    }
186                    let mut buf = vec![0u8; size as usize];
187                    f.seek(SeekFrom::Start(data))?;
188                    f.read_exact(&mut buf)?;
189                    long_name = if kind == b'L' {
190                        Some(cstr(&buf))
191                    } else {
192                        pax_path(&buf)
193                    };
194                }
195                b'g' => {}
196                b'0' | 0 => {
197                    let n = long_name.take().unwrap_or(name);
198                    entries.insert(normalize(&n), (data, size));
199                }
200                _ => {
201                    long_name = None;
202                }
203            }
204            pos = data + size.div_ceil(512) * 512;
205        }
206        Ok(Layout {
207            path: path.to_path_buf(),
208            entries,
209        })
210    }
211
212    fn entry(&self, name: &str) -> Result<(u64, u64)> {
213        self.entries.get(name).copied().ok_or_else(|| {
214            Error::invalid(format!("image archive {}: no {name}", self.path.display()))
215        })
216    }
217
218    fn blob_name(digest: &str) -> Result<String> {
219        if !valid_digest(digest) {
220            return Err(Error::invalid(format!(
221                "image archive: unsupported digest {digest:?}"
222            )));
223        }
224        Ok(format!("blobs/sha256/{}", &digest[7..]))
225    }
226
227    fn read_at(&self, (off, size): (u64, u64), limit: u64) -> Result<Vec<u8>> {
228        if size > limit {
229            return Err(Error::invalid(format!(
230                "image archive {}: a {size}-byte manifest is over the {limit}-byte limit",
231                self.path.display()
232            )));
233        }
234        let mut f = File::open(&self.path)?;
235        f.seek(SeekFrom::Start(off))?;
236        let mut buf = vec![0u8; size as usize];
237        f.read_exact(&mut buf)?;
238        Ok(buf)
239    }
240
241    /// A blob read into memory (manifests and configs), checked against its
242    /// digest.
243    pub fn read_blob(&self, d: &Descriptor) -> Result<Vec<u8>> {
244        let e = self.entry(&Self::blob_name(&d.digest)?)?;
245        let b = self.read_at(e, MAX_MANIFEST)?;
246        if digest_of(&b) != d.digest {
247            return Err(Error::invalid(format!(
248                "image archive: blob {} does not match its digest",
249                d.digest
250            )));
251        }
252        Ok(b)
253    }
254
255    /// A blob's bytes as a reader, and its size.
256    pub fn blob_reader(&self, digest: &str) -> Result<(std::io::Take<File>, u64)> {
257        let (off, size) = self.entry(&Self::blob_name(digest)?)?;
258        let mut f = File::open(&self.path)?;
259        f.seek(SeekFrom::Start(off))?;
260        Ok((f.take(size), size))
261    }
262
263    /// The image the layout holds: `index.json`'s one manifest.
264    pub fn top(&self) -> Result<Descriptor> {
265        let raw = self.read_at(self.entry("index.json")?, MAX_MANIFEST)?;
266        let idx: Index = serde_json::from_slice(&raw)
267            .map_err(|e| Error::invalid(format!("image archive: index.json: {e}")))?;
268        match idx.manifests.as_slice() {
269            [one] => Ok(one.clone()),
270            [] => Err(Error::invalid("image archive: index.json lists no image")),
271            more => Err(Error::invalid(format!(
272                "image archive: index.json lists {} images; a build exports one",
273                more.len()
274            ))),
275        }
276    }
277}
278
279// ---------------------------------------------------------------------------
280// The registry API
281// ---------------------------------------------------------------------------
282
283/// A registry, reached at `base` (`https://127.0.0.1:5480`).
284#[derive(Clone)]
285pub struct Remote {
286    agent: ureq::Agent,
287    base: String,
288}
289
290/// Repositories are `<org>/<app>`, both validated by their owners; check
291/// once more before a name goes into a URL.
292fn check_repo(repo: &str) -> Result<()> {
293    let ok = !repo.is_empty()
294        && repo.len() <= 255
295        && repo.split('/').all(|p| {
296            !p.is_empty()
297                && p.starts_with(|c: char| c.is_ascii_lowercase() || c.is_ascii_digit())
298                && p.chars()
299                    .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || "._-".contains(c))
300        });
301    if ok {
302        Ok(())
303    } else {
304        Err(Error::invalid(format!("invalid repository name {repo:?}")))
305    }
306}
307
308/// A tag or a digest.
309fn check_reference(r: &str) -> Result<()> {
310    if valid_digest(r) || super::valid_tag(r) {
311        Ok(())
312    } else {
313        Err(Error::invalid(format!("invalid tag or digest {r:?}")))
314    }
315}
316
317impl Remote {
318    /// `ca_pem`: the one CA to trust (none: the platform's roots, or plain
319    /// http in tests).
320    pub fn new(base: &str, ca_pem: Option<&str>, timeout: Duration) -> Result<Remote> {
321        let mut cfg = ureq::Agent::config_builder()
322            .timeout_global(Some(timeout))
323            .http_status_as_error(false)
324            .max_redirects(0)
325            .user_agent(concat!("isb/", env!("CARGO_PKG_VERSION")));
326        if let Some(pem) = ca_pem {
327            let cert = ureq::tls::Certificate::from_pem(pem.as_bytes())
328                .map_err(|e| Error::invalid(format!("registry CA: {e}")))?;
329            cfg = cfg.tls_config(
330                ureq::tls::TlsConfig::builder()
331                    .root_certs(ureq::tls::RootCerts::new_with_certs(&[cert]))
332                    .build(),
333            );
334        }
335        Ok(Remote {
336            agent: cfg.build().into(),
337            base: base.trim_end_matches('/').to_string(),
338        })
339    }
340
341    pub fn base(&self) -> &str {
342        &self.base
343    }
344
345    fn err(&self, step: &str, e: impl std::fmt::Display) -> Error {
346        Error::invalid(format!("registry {}: {step}: {e}", self.base))
347    }
348
349    fn status_err(&self, step: &str, mut r: ureq::http::Response<ureq::Body>) -> Error {
350        let status = r.status().as_u16();
351        let body = r
352            .body_mut()
353            .with_config()
354            .limit(64 << 10)
355            .read_to_string()
356            .unwrap_or_default();
357        let detail = serde_json::from_str::<Value>(&body)
358            .ok()
359            .and_then(|v| {
360                v["errors"].as_array().map(|a| {
361                    a.iter()
362                        .map(|e| {
363                            format!(
364                                "{} {}",
365                                e["code"].as_str().unwrap_or(""),
366                                e["message"].as_str().unwrap_or("")
367                            )
368                        })
369                        .collect::<Vec<_>>()
370                        .join("; ")
371                })
372            })
373            .unwrap_or(body);
374        self.err(step, format!("HTTP {status} {}", detail.trim()))
375    }
376
377    /// Whether `repo` has the blob.
378    fn blob_exists(&self, repo: &str, digest: &str) -> Result<bool> {
379        let step = format!("check blob {digest}");
380        let r = self
381            .agent
382            .head(format!("{}/v2/{repo}/blobs/{digest}", self.base))
383            .call()
384            .map_err(|e| self.err(&step, e))?;
385        match r.status().as_u16() {
386            200 => Ok(true),
387            404 => Ok(false),
388            _ => Err(self.status_err(&step, r)),
389        }
390    }
391
392    /// Where an upload continues: the registry's `Location`, which must stay
393    /// on this registry.
394    fn location(&self, r: &ureq::http::Response<ureq::Body>, step: &str) -> Result<String> {
395        let loc = r
396            .headers()
397            .get("location")
398            .and_then(|v| v.to_str().ok())
399            .ok_or_else(|| self.err(step, "no Location"))?;
400        if loc.starts_with('/') {
401            Ok(format!("{}{loc}", self.base))
402        } else if loc.starts_with(&format!("{}/", self.base)) {
403            Ok(loc.to_string())
404        } else {
405            Err(self.err(step, format!("refusing to upload to {loc}")))
406        }
407    }
408
409    /// Upload one blob, monolithically, unless the repository has it.
410    /// Returns whether it was uploaded.
411    pub fn put_blob(
412        &self,
413        repo: &str,
414        digest: &str,
415        size: u64,
416        body: &mut dyn Read,
417    ) -> Result<bool> {
418        check_repo(repo)?;
419        if !valid_digest(digest) {
420            return Err(self.err("upload", format!("unsupported digest {digest:?}")));
421        }
422        if self.blob_exists(repo, digest)? {
423            return Ok(false);
424        }
425        let step = format!("start upload of {digest}");
426        let r = self
427            .agent
428            .post(format!("{}/v2/{repo}/blobs/uploads/", self.base))
429            .header("Content-Length", "0")
430            .send_empty()
431            .map_err(|e| self.err(&step, e))?;
432        if r.status().as_u16() != 202 {
433            return Err(self.status_err(&step, r));
434        }
435        let loc = self.location(&r, &step)?;
436        let sep = if loc.contains('?') { '&' } else { '?' };
437        let step = format!("upload {digest} ({size} bytes)");
438        let r = self
439            .agent
440            .put(format!("{loc}{sep}digest={digest}"))
441            .header("Content-Type", "application/octet-stream")
442            .header("Content-Length", size.to_string())
443            .send(ureq::SendBody::from_reader(body))
444            .map_err(|e| self.err(&step, e))?;
445        if r.status().as_u16() != 201 {
446            return Err(self.status_err(&step, r));
447        }
448        Ok(true)
449    }
450
451    /// Put a manifest under `reference` (a tag or its digest). Returns the
452    /// digest the registry gives it.
453    pub fn put_manifest(
454        &self,
455        repo: &str,
456        reference: &str,
457        media_type: &str,
458        bytes: &[u8],
459    ) -> Result<String> {
460        check_repo(repo)?;
461        check_reference(reference)?;
462        let step = format!("put manifest {repo}:{reference}");
463        let r = self
464            .agent
465            .put(format!("{}/v2/{repo}/manifests/{reference}", self.base))
466            .header("Content-Type", media_type)
467            .send(bytes)
468            .map_err(|e| self.err(&step, e))?;
469        if r.status().as_u16() != 201 {
470            return Err(self.status_err(&step, r));
471        }
472        let ours = digest_of(bytes);
473        match r
474            .headers()
475            .get("docker-content-digest")
476            .and_then(|v| v.to_str().ok())
477        {
478            Some(d) if d != ours => Err(self.err(
479                &step,
480                format!("the registry stored {d}, but the manifest is {ours}"),
481            )),
482            _ => Ok(ours),
483        }
484    }
485
486    /// Push the layout's image as `repo:tag`. Returns its manifest digest.
487    pub fn push(
488        &self,
489        layout: &Layout,
490        repo: &str,
491        tag: &str,
492        log: &mut dyn FnMut(&str),
493    ) -> Result<String> {
494        check_repo(repo)?;
495        check_reference(tag)?;
496        let top = layout.top()?;
497        let bytes = self.push_children(layout, repo, &top, 0, log)?;
498        let digest = self.put_manifest(repo, tag, &top.media_type, &bytes)?;
499        if digest != top.digest {
500            return Err(self.err("push", "the top manifest does not match index.json"));
501        }
502        log(&format!("pushed {repo}:{tag} ({digest})"));
503        Ok(digest)
504    }
505
506    /// Push everything `d` refers to; returns `d`'s own bytes.
507    fn push_children(
508        &self,
509        layout: &Layout,
510        repo: &str,
511        d: &Descriptor,
512        depth: usize,
513        log: &mut dyn FnMut(&str),
514    ) -> Result<Vec<u8>> {
515        let bytes = layout.read_blob(d)?;
516        // Some exporters leave the media type to the document itself.
517        let mt = if d.media_type.is_empty() {
518            serde_json::from_slice::<Value>(&bytes)
519                .ok()
520                .and_then(|v| v["mediaType"].as_str().map(String::from))
521                .unwrap_or_default()
522        } else {
523            d.media_type.clone()
524        };
525        if is_index(&mt) {
526            if depth >= MAX_DEPTH {
527                return Err(self.err("push", "image indexes nest too deeply"));
528            }
529            let idx: Index = serde_json::from_slice(&bytes)
530                .map_err(|e| self.err("push", format!("index {}: {e}", d.digest)))?;
531            for m in &idx.manifests {
532                let child = self.push_children(layout, repo, m, depth + 1, log)?;
533                let cm = if m.media_type.is_empty() {
534                    MT_OCI_MANIFEST
535                } else {
536                    &m.media_type
537                };
538                self.put_manifest(repo, &m.digest, cm, &child)?;
539            }
540        } else if is_manifest(&mt) {
541            let m: Manifest = serde_json::from_slice(&bytes)
542                .map_err(|e| self.err("push", format!("manifest {}: {e}", d.digest)))?;
543            for b in std::iter::once(&m.config).chain(m.layers.iter()) {
544                let (mut r, size) = layout.blob_reader(&b.digest)?;
545                if size != b.size {
546                    return Err(self.err(
547                        "push",
548                        format!("blob {} is {size} bytes, not {}", b.digest, b.size),
549                    ));
550                }
551                if self.put_blob(repo, &b.digest, size, &mut r)? {
552                    log(&format!("pushed blob {} ({size} bytes)", short(&b.digest)));
553                }
554            }
555        } else {
556            return Err(self.err("push", format!("unsupported media type {mt:?}")));
557        }
558        Ok(bytes)
559    }
560
561    /// The digest `reference` (a tag or a digest) names, or `None`.
562    pub fn resolve(&self, repo: &str, reference: &str) -> Result<Option<String>> {
563        check_repo(repo)?;
564        check_reference(reference)?;
565        let step = format!("resolve {repo}:{reference}");
566        let r = self
567            .agent
568            .head(format!("{}/v2/{repo}/manifests/{reference}", self.base))
569            .header("Accept", ACCEPT)
570            .call()
571            .map_err(|e| self.err(&step, e))?;
572        match r.status().as_u16() {
573            200 => r
574                .headers()
575                .get("docker-content-digest")
576                .and_then(|v| v.to_str().ok())
577                .filter(|d| valid_digest(d))
578                .map(|d| Some(d.to_string()))
579                .ok_or_else(|| self.err(&step, "no Docker-Content-Digest")),
580            404 => Ok(None),
581            _ => Err(self.status_err(&step, r)),
582        }
583    }
584
585    /// A manifest's media type and bytes, checked against its digest.
586    pub fn get_manifest(&self, repo: &str, digest: &str) -> Result<Option<(String, Vec<u8>)>> {
587        check_repo(repo)?;
588        if !valid_digest(digest) {
589            return Err(self.err("get manifest", format!("not a digest: {digest:?}")));
590        }
591        let step = format!("get {repo}@{digest}");
592        let mut r = self
593            .agent
594            .get(format!("{}/v2/{repo}/manifests/{digest}", self.base))
595            .header("Accept", ACCEPT)
596            .call()
597            .map_err(|e| self.err(&step, e))?;
598        match r.status().as_u16() {
599            200 => {
600                let mt = r
601                    .headers()
602                    .get("content-type")
603                    .and_then(|v| v.to_str().ok())
604                    .unwrap_or(MT_OCI_MANIFEST)
605                    .to_string();
606                let b = r
607                    .body_mut()
608                    .with_config()
609                    .limit(MAX_MANIFEST)
610                    .read_to_vec()
611                    .map_err(|e| self.err(&step, e))?;
612                if digest_of(&b) != digest {
613                    return Err(self.err(&step, "the manifest does not match its digest"));
614                }
615                Ok(Some((mt, b)))
616            }
617            404 => Ok(None),
618            _ => Err(self.status_err(&step, r)),
619        }
620    }
621
622    fn get_json(&self, url: &str, step: &str) -> Result<Option<Value>> {
623        let mut r = self.agent.get(url).call().map_err(|e| self.err(step, e))?;
624        match r.status().as_u16() {
625            200 => {
626                let b = r
627                    .body_mut()
628                    .with_config()
629                    .limit(16 << 20)
630                    .read_to_vec()
631                    .map_err(|e| self.err(step, e))?;
632                Ok(Some(
633                    serde_json::from_slice(&b).map_err(|e| self.err(step, e))?,
634                ))
635            }
636            404 => Ok(None),
637            _ => Err(self.status_err(step, r)),
638        }
639    }
640
641    /// A repository's tags (none when it does not exist).
642    pub fn tags(&self, repo: &str) -> Result<Vec<String>> {
643        check_repo(repo)?;
644        let mut t = self.paged(
645            &format!("/v2/{repo}/tags/list"),
646            "tags",
647            &format!("list tags of {repo}"),
648        )?;
649        t.sort();
650        Ok(t)
651    }
652
653    /// Every repository.
654    pub fn catalog(&self) -> Result<Vec<String>> {
655        self.paged("/v2/_catalog", "repositories", "list repositories")
656    }
657
658    /// A paginated list (`n` and `last`, as the distribution API pages).
659    fn paged(&self, path: &str, key: &str, step: &str) -> Result<Vec<String>> {
660        const PAGE: usize = 1000;
661        let mut out: Vec<String> = Vec::new();
662        loop {
663            let mut url = format!("{}{path}?n={PAGE}", self.base);
664            if let Some(l) = out.last() {
665                url.push_str(&format!("&last={}", crate::client::encode_query(l)));
666            }
667            let page: Vec<String> = self
668                .get_json(&url, step)?
669                .as_ref()
670                .and_then(|v| v[key].as_array())
671                .map(|a| {
672                    a.iter()
673                        .filter_map(|x| x.as_str().map(String::from))
674                        .collect()
675                })
676                .unwrap_or_default();
677            let n = page.len();
678            out.extend(page);
679            if n < PAGE || out.len() > 1_000_000 {
680                return Ok(out);
681            }
682        }
683    }
684
685    /// Delete a manifest (and with it every tag pointing at it).
686    pub fn delete_manifest(&self, repo: &str, digest: &str) -> Result<()> {
687        check_repo(repo)?;
688        if !valid_digest(digest) {
689            return Err(self.err("delete", format!("not a digest: {digest:?}")));
690        }
691        let step = format!("delete {repo}@{digest}");
692        let r = self
693            .agent
694            .delete(format!("{}/v2/{repo}/manifests/{digest}", self.base))
695            .call()
696            .map_err(|e| self.err(&step, e))?;
697        match r.status().as_u16() {
698            202 | 200 | 404 => Ok(()),
699            _ => Err(self.status_err(&step, r)),
700        }
701    }
702}
703
704/// `sha256:0123456789ab` for logs.
705pub fn short(d: &str) -> &str {
706    &d[..d.len().min(19)]
707}
708
709#[cfg(test)]
710pub(crate) mod tests {
711    use super::*;
712    use std::collections::BTreeMap;
713    use std::io::{BufRead, BufReader, Write};
714    use std::net::{TcpListener, TcpStream};
715    use std::sync::{Arc, Mutex};
716
717    /// A tar file with these entries (ustar, short names).
718    pub fn tar(entries: &[(&str, &[u8])]) -> Vec<u8> {
719        let mut out = Vec::new();
720        for (name, data) in entries {
721            let mut h = [0u8; 512];
722            h[..name.len()].copy_from_slice(name.as_bytes());
723            h[100..108].copy_from_slice(b"0000644\0");
724            h[124..136].copy_from_slice(format!("{:011o}\0", data.len()).as_bytes());
725            h[136..148].copy_from_slice(b"00000000000\0");
726            h[156] = b'0';
727            h[257..263].copy_from_slice(b"ustar\0");
728            h[263..265].copy_from_slice(b"00");
729            h[148..156].copy_from_slice(b"        ");
730            let sum: u32 = h.iter().map(|b| *b as u32).sum();
731            h[148..156].copy_from_slice(format!("{sum:06o}\0 ").as_bytes());
732            out.extend_from_slice(&h);
733            out.extend_from_slice(data);
734            out.resize(out.len().div_ceil(512) * 512, 0);
735        }
736        out.extend_from_slice(&[0u8; 1024]);
737        out
738    }
739
740    /// An OCI layout tar with one image: (tar bytes, manifest digest).
741    pub fn image_tar(layer: &[u8]) -> (Vec<u8>, String) {
742        let config = br#"{"architecture":"amd64","os":"linux","config":{}}"#.to_vec();
743        let (cd, ld) = (digest_of(&config), digest_of(layer));
744        let manifest = format!(
745            r#"{{"schemaVersion":2,"mediaType":"{MT_OCI_MANIFEST}","config":{{"mediaType":"application/vnd.oci.image.config.v1+json","digest":"{cd}","size":{}}},"layers":[{{"mediaType":"application/vnd.oci.image.layer.v1.tar","digest":"{ld}","size":{}}}]}}"#,
746            config.len(),
747            layer.len()
748        );
749        let md = digest_of(manifest.as_bytes());
750        let index = format!(
751            r#"{{"schemaVersion":2,"manifests":[{{"mediaType":"{MT_OCI_MANIFEST}","digest":"{md}","size":{}}}]}}"#,
752            manifest.len()
753        );
754        let p = |d: &str| format!("blobs/sha256/{}", &d[7..]);
755        let (pc, pl, pm) = (p(&cd), p(&ld), p(&md));
756        let t = tar(&[
757            ("oci-layout", br#"{"imageLayoutVersion":"1.0.0"}"#),
758            ("index.json", index.as_bytes()),
759            (&pc, &config),
760            (&pl, layer),
761            (&pm, manifest.as_bytes()),
762        ]);
763        (t, md)
764    }
765
766    /// A registry that keeps blobs and manifests in memory and records
767    /// every request, speaking plain HTTP on loopback.
768    #[derive(Default)]
769    pub struct Fake {
770        pub blobs: BTreeMap<String, Vec<u8>>,
771        /// (repo, reference) -> (media type, bytes)
772        pub manifests: BTreeMap<(String, String), (String, Vec<u8>)>,
773        pub log: Vec<String>,
774    }
775
776    pub fn fake() -> (String, Arc<Mutex<Fake>>) {
777        let l = TcpListener::bind("127.0.0.1:0").unwrap();
778        let addr = l.local_addr().unwrap();
779        let state = Arc::new(Mutex::new(Fake::default()));
780        let s = state.clone();
781        std::thread::spawn(move || {
782            for c in l.incoming() {
783                let Ok(c) = c else { continue };
784                let s = s.clone();
785                std::thread::spawn(move || serve(c, &s));
786            }
787        });
788        (format!("http://{addr}"), state)
789    }
790
791    fn serve(c: TcpStream, s: &Mutex<Fake>) {
792        let mut r = BufReader::new(c.try_clone().unwrap());
793        let mut w = c;
794        loop {
795            let mut line = String::new();
796            if r.read_line(&mut line).unwrap_or(0) == 0 {
797                return;
798            }
799            let mut parts = line.split_whitespace();
800            let (method, target) = (
801                parts.next().unwrap_or("").to_string(),
802                parts.next().unwrap_or("").to_string(),
803            );
804            let mut headers = BTreeMap::new();
805            loop {
806                let mut h = String::new();
807                r.read_line(&mut h).unwrap();
808                let h = h.trim_end();
809                if h.is_empty() {
810                    break;
811                }
812                let (k, v) = h.split_once(':').unwrap();
813                headers.insert(k.to_ascii_lowercase(), v.trim().to_string());
814            }
815            assert!(
816                !headers.contains_key("transfer-encoding"),
817                "uploads are length-delimited"
818            );
819            let len: usize = headers
820                .get("content-length")
821                .map(|v| v.parse().unwrap())
822                .unwrap_or(0);
823            let mut body = vec![0u8; len];
824            r.read_exact(&mut body).unwrap();
825            let (path, query) = target.split_once('?').unwrap_or((&target, ""));
826            let mut st = s.lock().unwrap();
827            st.log.push(format!("{method} {path}"));
828            let (status, extra, out): (u16, Vec<(String, String)>, Vec<u8>) =
829                route(&mut st, &method, path, query, &headers, body);
830            drop(st);
831            let mut resp = format!("HTTP/1.1 {status} X\r\nContent-Length: {}\r\n", out.len());
832            for (k, v) in extra {
833                resp.push_str(&format!("{k}: {v}\r\n"));
834            }
835            resp.push_str("\r\n");
836            // A HEAD response has the headers of a GET and no body.
837            let out = if method == "HEAD" { Vec::new() } else { out };
838            if w.write_all(resp.as_bytes()).is_err() || w.write_all(&out).is_err() {
839                return;
840            }
841        }
842    }
843
844    type Reply = (u16, Vec<(String, String)>, Vec<u8>);
845
846    #[expect(
847        clippy::too_many_lines,
848        reason = "predates the lint ratchet; split it when next changed"
849    )]
850    fn route(
851        st: &mut Fake,
852        method: &str,
853        path: &str,
854        query: &str,
855        headers: &BTreeMap<String, String>,
856        body: Vec<u8>,
857    ) -> Reply {
858        let p = path.trim_start_matches("/v2/");
859        let none = || {
860            (
861                404u16,
862                vec![],
863                b"{\"errors\":[{\"code\":\"NOT_FOUND\"}]}".to_vec(),
864            )
865        };
866        if let Some((repo, rest)) = p.split_once("/blobs/uploads/") {
867            return match method {
868                "POST" => (
869                    202,
870                    vec![(
871                        "Location".into(),
872                        format!("/v2/{repo}/blobs/uploads/u1?_state=x"),
873                    )],
874                    vec![],
875                ),
876                "PUT" => {
877                    assert_eq!(rest, "u1");
878                    let d = query
879                        .split('&')
880                        .find_map(|kv| kv.strip_prefix("digest="))
881                        .unwrap()
882                        .replace("%3A", ":");
883                    if digest_of(&body) != d {
884                        return (
885                            400,
886                            vec![],
887                            b"{\"errors\":[{\"code\":\"DIGEST_INVALID\"}]}".to_vec(),
888                        );
889                    }
890                    st.blobs.insert(format!("{repo}@{d}"), body);
891                    (201, vec![], vec![])
892                }
893                _ => none(),
894            };
895        }
896        if let Some((repo, d)) = p.split_once("/blobs/") {
897            return if st.blobs.contains_key(&format!("{repo}@{d}")) {
898                (200, vec![], vec![])
899            } else {
900                none()
901            };
902        }
903        if let Some((repo, reference)) = p.split_once("/manifests/") {
904            let key = (repo.to_string(), reference.to_string());
905            return match method {
906                "PUT" => {
907                    let d = digest_of(&body);
908                    let mt = headers.get("content-type").cloned().unwrap_or_default();
909                    st.manifests.insert(key, (mt.clone(), body.clone()));
910                    st.manifests
911                        .insert((repo.to_string(), d.clone()), (mt, body));
912                    (201, vec![("Docker-Content-Digest".into(), d)], vec![])
913                }
914                "HEAD" | "GET" => match st.manifests.get(&key) {
915                    Some((mt, b)) => (
916                        200,
917                        vec![
918                            ("Docker-Content-Digest".into(), digest_of(b)),
919                            ("Content-Type".into(), mt.clone()),
920                        ],
921                        b.clone(),
922                    ),
923                    None => none(),
924                },
925                "DELETE" => {
926                    let before = st.manifests.len();
927                    st.manifests
928                        .retain(|(r, _), (_, b)| !(r == repo && digest_of(b) == reference));
929                    if st.manifests.len() < before {
930                        (202, vec![], vec![])
931                    } else {
932                        none()
933                    }
934                }
935                _ => none(),
936            };
937        }
938        if let Some(repo) = p.strip_suffix("/tags/list") {
939            let mut tags: Vec<String> = st
940                .manifests
941                .keys()
942                .filter(|(r, t)| r == repo && !t.starts_with("sha256:"))
943                .map(|(_, t)| t.clone())
944                .collect();
945            tags.dedup();
946            return (
947                200,
948                vec![],
949                serde_json::to_vec(&serde_json::json!({"name": repo, "tags": tags})).unwrap(),
950            );
951        }
952        if p == "_catalog" {
953            let mut repos: Vec<String> = st.manifests.keys().map(|(r, _)| r.clone()).collect();
954            repos.dedup();
955            return (
956                200,
957                vec![],
958                serde_json::to_vec(&serde_json::json!({"repositories": repos})).unwrap(),
959            );
960        }
961        none()
962    }
963
964    fn layout_of(bytes: &[u8]) -> (tempfile::TempDir, Layout) {
965        let d = tempfile::tempdir().unwrap();
966        let p = d.path().join("image.tar");
967        std::fs::write(&p, bytes).unwrap();
968        let l = Layout::open(&p).unwrap();
969        (d, l)
970    }
971
972    #[test]
973    fn digests() {
974        assert!(valid_digest(&digest_of(b"x")));
975        assert!(!valid_digest("sha256:../../etc"));
976        assert!(!valid_digest("sha512:00"));
977        assert!(!valid_digest(&format!("sha256:{}", "A".repeat(64))));
978    }
979
980    #[test]
981    fn layout_is_read_without_unpacking() {
982        let (t, md) = image_tar(b"layer bytes");
983        let (_d, l) = layout_of(&t);
984        let top = l.top().unwrap();
985        assert_eq!(top.digest, md);
986        assert!(
987            l.read_blob(&top)
988                .unwrap()
989                .starts_with(b"{\"schemaVersion\"")
990        );
991        let (mut r, n) = l.blob_reader(&digest_of(b"layer bytes")).unwrap();
992        let mut s = String::new();
993        r.read_to_string(&mut s).unwrap();
994        assert_eq!((s.as_str(), n), ("layer bytes", 11));
995        assert!(l.blob_reader(&digest_of(b"other")).is_err());
996    }
997
998    #[test]
999    fn layout_rejects_lies() {
1000        // A blob whose content does not match the digest it is filed under.
1001        let fake = digest_of(b"claimed");
1002        let name = format!("blobs/sha256/{}", &fake[7..]);
1003        let t = tar(&[(&name, b"actual")]);
1004        let (_d, l) = layout_of(&t);
1005        let d = Descriptor {
1006            media_type: MT_OCI_MANIFEST.into(),
1007            digest: fake,
1008            size: 6,
1009        };
1010        assert!(l.read_blob(&d).is_err());
1011        // A size field past the end of the file.
1012        let mut t = tar(&[("index.json", b"{}")]);
1013        t[124..136].copy_from_slice(b"77777777777\0");
1014        let d = tempfile::tempdir().unwrap();
1015        std::fs::write(d.path().join("x.tar"), &t).unwrap();
1016        assert!(Layout::open(&d.path().join("x.tar")).is_err());
1017        // Two images in one export.
1018        let t = tar(&[(
1019            "index.json",
1020            br#"{"manifests":[{"digest":"sha256:00","size":1},{"digest":"sha256:01","size":1}]}"#,
1021        )]);
1022        let (_d, l) = layout_of(&t);
1023        assert!(l.top().is_err());
1024    }
1025
1026    #[test]
1027    fn push_uploads_missing_blobs_then_the_manifest() {
1028        let (base, st) = fake();
1029        let r = Remote::new(&base, None, Duration::from_secs(10)).unwrap();
1030        let (t, md) = image_tar(b"some layer");
1031        let (_d, l) = layout_of(&t);
1032        let mut lines = Vec::new();
1033        let d = r
1034            .push(&l, "acme/web", "v1", &mut |s| lines.push(s.to_string()))
1035            .unwrap();
1036        assert_eq!(d, md);
1037        assert_eq!(
1038            r.resolve("acme/web", "v1").unwrap().as_deref(),
1039            Some(md.as_str())
1040        );
1041        assert_eq!(r.resolve("acme/web", "v2").unwrap(), None);
1042        assert_eq!(r.tags("acme/web").unwrap(), vec!["v1".to_string()]);
1043        let log = st.lock().unwrap().log.clone();
1044        let uploads = log
1045            .iter()
1046            .filter(|l| l.starts_with("PUT /v2/acme/web/blobs/uploads"))
1047            .count();
1048        assert_eq!(uploads, 2, "config and layer: {log:?}");
1049        let put = log
1050            .iter()
1051            .position(|l| l == "PUT /v2/acme/web/manifests/v1")
1052            .unwrap();
1053        let last_upload = log
1054            .iter()
1055            .rposition(|l| l.contains("/blobs/uploads/"))
1056            .unwrap();
1057        assert!(put > last_upload, "the manifest goes last: {log:?}");
1058
1059        // A second push of the same image uploads nothing new.
1060        st.lock().unwrap().log.clear();
1061        r.push(&l, "acme/web", "v2", &mut |_| {}).unwrap();
1062        let log = st.lock().unwrap().log.clone();
1063        assert!(!log.iter().any(|l| l.contains("uploads")), "{log:?}");
1064        // Blobs are per repository: another one uploads its own.
1065        r.push(&l, "other/web", "v1", &mut |_| {}).unwrap();
1066        assert!(
1067            st.lock()
1068                .unwrap()
1069                .blobs
1070                .contains_key(&format!("other/web@{}", digest_of(b"some layer")))
1071        );
1072
1073        let (mt, b) = r.get_manifest("acme/web", &md).unwrap().unwrap();
1074        assert_eq!((mt.as_str(), digest_of(&b)), (MT_OCI_MANIFEST, md.clone()));
1075        r.delete_manifest("acme/web", &md).unwrap();
1076        assert!(r.get_manifest("acme/web", &md).unwrap().is_none());
1077        assert_eq!(r.resolve("acme/web", "v1").unwrap(), None);
1078        assert!(r.resolve("acme/../x", "v1").is_err());
1079        assert!(r.resolve("acme/web", "bad tag").is_err());
1080    }
1081}