1use 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
26const 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
29pub const MAX_MANIFEST: u64 = 4 << 20;
31const MAX_DEPTH: usize = 3;
33
34pub 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
44pub 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#[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
84pub struct Layout {
90 path: PathBuf,
91 entries: HashMap<String, (u64, u64)>,
93}
94
95fn octal(field: &[u8]) -> Option<u64> {
96 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
127fn 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 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 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 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 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#[derive(Clone)]
285pub struct Remote {
286 agent: ureq::Agent,
287 base: String,
288}
289
290fn 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
308fn 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 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 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 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 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 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 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 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 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 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 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 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 pub fn catalog(&self) -> Result<Vec<String>> {
655 self.paged("/v2/_catalog", "repositories", "list repositories")
656 }
657
658 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 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
704pub 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 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 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 #[derive(Default)]
769 pub struct Fake {
770 pub blobs: BTreeMap<String, Vec<u8>>,
771 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 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 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 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 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 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 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}