Commit

g1t push: images with layers of any size, uploaded in chunks under the edge's limit

Cloudflare refuses a request body over 100 MB before g1t sees it, and docker, crane and oras send each layer in one request. g1t push reads an image with docker save (or --archive, either save layout), skips what the registry has, and uploads each layer in chunks of at most 90 MiB with retries and resume, then the manifest. Tokens from G1T_TOKEN, --token-stdin or docker's own stored login. Checked on g1t.sh: a 143 MiB layer docker push could not send (413) pushed in two chunks, pulled and ran identically. The containers guide says how to push large layers with it.

syntaqxcommitted Parent93eb962Browse files
9 files+2465−170/9 viewed
+66−13
11 ---
22 title: Container images
3−description: Push and pull container images on g1t.sh with docker, in workflows with G1T_TOKEN, and what to do about large layers.
3+description: Push and pull container images on g1t.sh with docker, in workflows with G1T_TOKEN, and push layers over 100 MB with g1t push.
44 ---
55
66 g1t.sh is a container registry. Images are named after their workspace,
7979 ## The 100 MB limit
8080
8181 A single request to g1t.sh may carry at most 100 MB. `docker push` sends
82−each layer whole, in one request, and so do the other common clients
83−(`crane push` and `oras push` included), so a layer over 100 MB, as
84−compressed for the push, is refused.
82+each layer whole, in one request, so a layer over 100 MB, as compressed for
83+the push, is refused. [`g1t push`](#push-large-layers-with-g1t-push) sends
84+it in chunks instead, and has no limit.
8585
8686 What you see depends on where it is refused. Usually it is before the
8787 request reaches g1t, and `docker push` stops with a bare status:
9595 the request itself, the error is `SIZE_INVALID`, with a message naming the
9696 limit and this page.
9797
98−A client that uploads a layer in chunks, each its own request under
99−100 MB, is not limited: the layer can then be any size.
98+An installation you [run yourself](/guides/self-hosting/) has no such
99+limit.
100+
101+### Push large layers with g1t push
102+
103+`g1t push` takes an image from your local docker and pushes it to g1t.sh,
104+each layer in chunks of 90 MiB, each chunk its own request. A layer can be
105+any size.
100106
101−To stay under the limit, keep each layer under 100 MB:
107+```sh
108+docker build -t g1t.sh/acme/model:1 .
109+g1t push g1t.sh/acme/model:1
110+```
102111
103−- build in stages, and copy only what the image needs into the last one;
104−- split a large `RUN` or `COPY` into several, so each makes its own layer;
105−- leave caches, build tools and test data out of the image (`.dockerignore`).
112+An image with a local name is pushed to the address after `--as`:
106113
107−An installation you [run yourself](/guides/self-hosting/) has no such
108−limit.
114+```sh
115+g1t push model:dev --as g1t.sh/acme/model:1
116+```
109117
118+```text
119+Reading model:dev from docker
120+Pushing to g1t.sh/acme/model:1
121+ config sha256:040e744c070b 851 B already on the registry, skipped
122+ layer sha256:25f1d6b1951a 3.5 MiB already on the registry, skipped
123+ chunk 1/2 90.0 MiB 63%
124+ chunk 2/2 53.1 MiB 100%
125+ layer sha256:c74595c4a2cd 143.1 MiB uploaded in 2 chunks
126+Pushed g1t.sh/acme/model:1
127+digest: sha256:39b972d91774d58b5fbd27cfdce73fe58034d9e5772a261d183c11bcf78c1dba
128+```
129+
130+It pushes the image docker has: the same config, layers and manifest, so
131+`docker pull` gets back exactly what you built. When docker keeps a layer
132+uncompressed, `g1t push` gzips it for the push, as `docker push` does.
133+Layers already on g1t are not sent again. A request refused with `429` or
134+an error on g1t's side is tried again after a wait, and an interrupted
135+layer goes on from where it stopped.
136+
137+It signs in with the first of:
138+
139+1. a token on stdin, with `--token-stdin`
140+ (`echo "$G1T_TOKEN" | g1t push … --token-stdin`);
141+2. the `G1T_TOKEN` environment variable, as in [workflows](#in-workflows);
142+3. what `docker login g1t.sh` stored, in `~/.docker/config.json` or the
143+ credential store it names.
144+
145+| Option | |
146+| --- | --- |
147+| `--as <address>` | Where to push, `g1t.sh/<workspace>/<name>:<tag>`. Without it, the image's own name must be such an address. The tag defaults to `latest`. |
148+| `--chunk-size <size>` | How much each request carries: `5MB` to `95MB`, `90MB` unless set. `MB` and `MiB` both mean 1,048,576 bytes. |
149+| `--token-stdin` | Read the token from stdin. |
150+| `--archive <file>` | Push a tarball written by `docker save`, instead of asking docker. Needs `--as`. |
151+
152+`g1t --version` prints its version, and `g1t help` its options.
153+
154+Release binaries of the g1t command line are coming. Until then, build it
155+from the g1t source with [Rust](https://www.rust-lang.org/tools/install):
156+
157+```sh
158+git clone https://g1t.sh/flagon-io/g1t
159+cd g1t
160+cargo install --path crates/g1t
161+```
162+
110163 ## Storage and pull limits
111164
112165 Without the [g1t plan](/guides/usage-and-billing/#the-g1t-plan), a
146199 | `NAME_UNKNOWN` | No such image, or one you cannot see. |
147200 | `MANIFEST_UNKNOWN`, `BLOB_UNKNOWN` | No such tag, digest or layer in that image. |
148201 | `NAME_INVALID` | Names are lowercase letters and digits, separated by `.`, `_`, `__`, `-` or `/`, and start with a workspace. |
149−| `SIZE_INVALID`, or a bare `413` | A request over [the 100 MB limit](#the-100-mb-limit). |
202+| `SIZE_INVALID`, or a bare `413` | A request over [the 100 MB limit](#the-100-mb-limit). Push the image with [`g1t push`](#push-large-layers-with-g1t-push). |
150203 | `TOOMANYREQUESTS` | Too many requests in a minute; see [the limits](#storage-and-pull-limits). |
151204 | `DIGEST_INVALID` | What was uploaded does not have the digest the client said. Push again. |
+8−0
44 edition.workspace = true
55 license.workspace = true
66 repository.workspace = true
7+description = "The g1t command line: g1t push sends a local image to g1t.sh in chunks, so a layer can be any size."
78
89 [dependencies]
10+base64 = "0.22"
11+flate2 = "1"
12+hex = "0.4"
13+serde_json = { workspace = true }
14+sha2 = "0.10"
15+tar = { version = "0.4", default-features = false }
16+ureq = { version = "2", default-features = false, features = ["tls"] }
+187−0
1+//! Sizes, chunk plans and the range headers of a chunked upload.
2+
3+pub const MIB: u64 = 1024 * 1024;
4+/// What a chunk is unless `--chunk-size` says otherwise.
5+pub const DEFAULT_CHUNK: u64 = 90 * MIB;
6+/// The largest chunk: well under the 100 MB a request to g1t.sh may carry.
7+pub const MAX_CHUNK: u64 = 95 * MIB;
8+/// The smallest: every chunk but the last becomes a part in storage.
9+pub const MIN_CHUNK: u64 = 5 * MIB;
10+
11+/// Reads a size like `90MB`, `90MiB`, `90M`, `512K` or `94371840`. K, M and
12+/// G count in 1,024s, with or without a trailing `B` or `iB`.
13+pub fn parse_size(text: &str) -> Result<u64, String> {
14+ let text = text.trim();
15+ let split = text
16+ .find(|c: char| !c.is_ascii_digit() && c != '.')
17+ .unwrap_or(text.len());
18+ let (number, unit) = text.split_at(split);
19+ let number: f64 = number
20+ .parse()
21+ .map_err(|_| format!("`{text}` is not a size, like 90MB"))?;
22+ let unit = unit.trim().to_ascii_lowercase();
23+ let unit = unit
24+ .strip_suffix("ib")
25+ .or_else(|| unit.strip_suffix('b'))
26+ .unwrap_or(&unit);
27+ let scale = match unit {
28+ "" => 1,
29+ "k" => 1024,
30+ "m" => MIB,
31+ "g" => 1024 * MIB,
32+ _ => return Err(format!("`{text}` is not a size, like 90MB")),
33+ };
34+ Ok((number * scale as f64) as u64)
35+}
36+
37+/// Checks a chunk size against the limits.
38+pub fn check_chunk(size: u64) -> Result<u64, String> {
39+ if size > MAX_CHUNK {
40+ return Err(format!(
41+ "A chunk may be at most 95MiB ({MAX_CHUNK} bytes), to stay under the 100 MB a request may carry."
42+ ));
43+ }
44+ if size < MIN_CHUNK {
45+ return Err(format!(
46+ "A chunk must be at least 5MiB ({MIN_CHUNK} bytes)."
47+ ));
48+ }
49+ Ok(size)
50+}
51+
52+/// The chunks that send bytes `from..total`: `(start, length)` each.
53+pub fn plan(total: u64, chunk: u64, from: u64) -> Vec<(u64, u64)> {
54+ let mut chunks = Vec::new();
55+ let mut start = from;
56+ while start < total {
57+ let length = chunk.min(total - start);
58+ chunks.push((start, length));
59+ start += length;
60+ }
61+ chunks
62+}
63+
64+/// How many chunks a blob of `total` bytes takes.
65+pub fn count(total: u64, chunk: u64) -> u64 {
66+ total.div_ceil(chunk).max(1)
67+}
68+
69+/// A PATCH's `Content-Range`: `<first byte>-<last byte>`.
70+pub fn content_range(start: u64, length: u64) -> String {
71+ format!("{start}-{}", start + length.max(1) - 1)
72+}
73+
74+/// The offset an upload's `Range` header (`0-<last byte>`) says comes next.
75+/// `0-0` reads as nothing received yet: registries say it before the first
76+/// byte, and a one-byte upload that was in fact received is caught by the
77+/// registry refusing the resent byte.
78+pub fn next_offset(range: &str) -> Option<u64> {
79+ let text = range.trim();
80+ let text = text
81+ .strip_prefix("bytes=")
82+ .or_else(|| text.strip_prefix("bytes "))
83+ .unwrap_or(text);
84+ let (start, end) = text.split_once('-')?;
85+ let (start, end): (u64, u64) = (start.trim().parse().ok()?, end.trim().parse().ok()?);
86+ if start != 0 {
87+ return None;
88+ }
89+ Some(if end == 0 { 0 } else { end + 1 })
90+}
91+
92+/// A size for people: `612 B`, `3.5 MiB`.
93+pub fn human(bytes: u64) -> String {
94+ const UNITS: [&str; 4] = ["KiB", "MiB", "GiB", "TiB"];
95+ if bytes < 1024 {
96+ return format!("{bytes} B");
97+ }
98+ let mut value = bytes as f64 / 1024.0;
99+ let mut unit = 0;
100+ while value >= 1024.0 && unit < UNITS.len() - 1 {
101+ value /= 1024.0;
102+ unit += 1;
103+ }
104+ format!("{value:.1} {}", UNITS[unit])
105+}
106+
107+#[cfg(test)]
108+mod tests {
109+ use super::*;
110+
111+ #[test]
112+ fn sizes_read_in_1024s_with_any_suffix() {
113+ assert_eq!(parse_size("90MB"), Ok(90 * MIB));
114+ assert_eq!(parse_size("90MiB"), Ok(90 * MIB));
115+ assert_eq!(parse_size("90m"), Ok(90 * MIB));
116+ assert_eq!(parse_size("512K"), Ok(512 * 1024));
117+ assert_eq!(parse_size("94371840"), Ok(94_371_840));
118+ assert_eq!(parse_size("1.5G"), Ok(1536 * MIB));
119+ assert!(parse_size("ninety").is_err());
120+ assert!(parse_size("90XB").is_err());
121+ }
122+
123+ #[test]
124+ fn chunks_stay_between_5_and_95_mib() {
125+ assert_eq!(check_chunk(95 * MIB), Ok(95 * MIB));
126+ assert!(check_chunk(95 * MIB + 1).is_err());
127+ assert!(check_chunk(100_000_000).is_err());
128+ assert!(check_chunk(MIN_CHUNK - 1).is_err());
129+ }
130+
131+ #[test]
132+ fn a_plan_covers_every_byte_once() {
133+ let total = 150_000_000 + 4_096;
134+ let chunks = plan(total, DEFAULT_CHUNK, 0);
135+ assert_eq!(
136+ chunks,
137+ vec![(0, DEFAULT_CHUNK), (DEFAULT_CHUNK, total - DEFAULT_CHUNK)]
138+ );
139+ assert_eq!(chunks.iter().map(|c| c.1).sum::<u64>(), total);
140+ assert_eq!(count(total, DEFAULT_CHUNK), 2);
141+ }
142+
143+ #[test]
144+ fn a_plan_on_an_exact_multiple_has_no_empty_tail() {
145+ let chunks = plan(3 * MIB * 10, 10 * MIB, 0);
146+ assert_eq!(chunks.len(), 3);
147+ assert!(chunks.iter().all(|c| c.1 == 10 * MIB));
148+ }
149+
150+ #[test]
151+ fn a_plan_resumes_from_an_offset() {
152+ let chunks = plan(25 * MIB, 10 * MIB, 10 * MIB);
153+ assert_eq!(chunks, vec![(10 * MIB, 10 * MIB), (20 * MIB, 5 * MIB)]);
154+ assert!(plan(25 * MIB, 10 * MIB, 25 * MIB).is_empty());
155+ }
156+
157+ #[test]
158+ fn small_and_empty_blobs_are_one_chunk_or_none() {
159+ assert_eq!(plan(612, DEFAULT_CHUNK, 0), vec![(0, 612)]);
160+ assert!(plan(0, DEFAULT_CHUNK, 0).is_empty());
161+ assert_eq!(count(0, DEFAULT_CHUNK), 1);
162+ }
163+
164+ #[test]
165+ fn content_ranges_name_the_first_and_last_byte() {
166+ assert_eq!(content_range(0, 1024), "0-1023");
167+ assert_eq!(
168+ content_range(DEFAULT_CHUNK, 10),
169+ format!("{}-{}", DEFAULT_CHUNK, DEFAULT_CHUNK + 9)
170+ );
171+ }
172+
173+ #[test]
174+ fn upload_ranges_give_the_next_offset() {
175+ assert_eq!(next_offset("0-0"), Some(0));
176+ assert_eq!(next_offset("0-94371839"), Some(94_371_840));
177+ assert_eq!(next_offset("bytes=0-1023"), Some(1024));
178+ assert_eq!(next_offset("5-10"), None);
179+ assert_eq!(next_offset("junk"), None);
180+ }
181+
182+ #[test]
183+ fn sizes_print_for_people() {
184+ assert_eq!(human(612), "612 B");
185+ assert_eq!(human(150_000_000), "143.1 MiB");
186+ }
187+}
+228−0
1+//! Where the push's g1t token comes from: `G1T_TOKEN`, `--token-stdin`, or
2+//! what `docker login` stored.
3+
4+use base64::Engine as _;
5+use serde_json::Value;
6+use std::io::{Read, Write};
7+use std::path::PathBuf;
8+use std::process::{Command, Stdio};
9+
10+/// A username and a g1t token, for the registry's token endpoint.
11+#[derive(Clone, Debug, PartialEq, Eq)]
12+pub struct Credentials {
13+ pub username: String,
14+ pub secret: String,
15+}
16+
17+/// Any username will do with a token; this is the one g1t sends.
18+const USERNAME: &str = "g1t";
19+
20+pub fn find(host: &str, token_stdin: bool) -> Result<Credentials, String> {
21+ if token_stdin {
22+ let mut token = String::new();
23+ std::io::stdin()
24+ .read_to_string(&mut token)
25+ .map_err(|e| format!("Could not read the token from stdin: {e}"))?;
26+ return token_credentials(token.trim())
27+ .ok_or_else(|| "--token-stdin read an empty token.".to_owned());
28+ }
29+ if let Some(credentials) = std::env::var("G1T_TOKEN")
30+ .ok()
31+ .and_then(|t| token_credentials(t.trim()))
32+ {
33+ return Ok(credentials);
34+ }
35+ let config = docker_config_path().and_then(|path| std::fs::read_to_string(path).ok());
36+ if let Some(found) = config.and_then(|config| from_docker_config(&config, host, &run_helper)) {
37+ return Ok(found);
38+ }
39+ Err(format!(
40+ "No token for {host}. Set G1T_TOKEN, pass one with --token-stdin, or sign in with `docker login {host}`."
41+ ))
42+}
43+
44+fn token_credentials(token: &str) -> Option<Credentials> {
45+ (!token.is_empty()).then(|| Credentials {
46+ username: USERNAME.to_owned(),
47+ secret: token.to_owned(),
48+ })
49+}
50+
51+fn docker_config_path() -> Option<PathBuf> {
52+ if let Some(dir) = std::env::var_os("DOCKER_CONFIG") {
53+ return Some(PathBuf::from(dir).join("config.json"));
54+ }
55+ let home = std::env::var_os("HOME").or_else(|| std::env::var_os("USERPROFILE"))?;
56+ Some(PathBuf::from(home).join(".docker").join("config.json"))
57+}
58+
59+/// Looks `host` up in a docker `config.json` the way docker does: its own
60+/// credential helper first, then the credential store, then a plain `auth`.
61+/// `helper` runs `docker-credential-<name> get`.
62+pub fn from_docker_config(
63+ config: &str,
64+ host: &str,
65+ helper: &dyn Fn(&str, &str) -> Option<Credentials>,
66+) -> Option<Credentials> {
67+ let config: Value = serde_json::from_str(config).ok()?;
68+ let servers = [host.to_owned(), format!("https://{host}")];
69+ if let Some(name) = config
70+ .get("credHelpers")
71+ .and_then(|h| h.get(host))
72+ .and_then(Value::as_str)
73+ {
74+ return servers.iter().find_map(|server| helper(name, server));
75+ }
76+ if let Some(name) = config
77+ .get("credsStore")
78+ .and_then(Value::as_str)
79+ .filter(|n| !n.is_empty())
80+ && let Some(found) = servers.iter().find_map(|server| helper(name, server))
81+ {
82+ return Some(found);
83+ }
84+ let auths = config.get("auths")?.as_object()?;
85+ auths
86+ .iter()
87+ .find(|(key, _)| server_host(key) == host)
88+ .and_then(|(_, entry)| {
89+ let auth = entry.get("auth")?.as_str()?;
90+ let decoded = base64::engine::general_purpose::STANDARD
91+ .decode(auth.trim())
92+ .ok()?;
93+ let decoded = String::from_utf8(decoded).ok()?;
94+ let (username, secret) = decoded.split_once(':')?;
95+ (!secret.is_empty()).then(|| Credentials {
96+ username: username.to_owned(),
97+ secret: secret.to_owned(),
98+ })
99+ })
100+}
101+
102+/// `https://g1t.sh/v2/` and `g1t.sh` are the same server.
103+fn server_host(key: &str) -> &str {
104+ let key = key
105+ .strip_prefix("https://")
106+ .or_else(|| key.strip_prefix("http://"))
107+ .unwrap_or(key);
108+ key.split('/').next().unwrap_or(key)
109+}
110+
111+/// `docker-credential-<name> get`, with the server on stdin.
112+fn run_helper(name: &str, server: &str) -> Option<Credentials> {
113+ let mut child = Command::new(format!("docker-credential-{name}"))
114+ .arg("get")
115+ .stdin(Stdio::piped())
116+ .stdout(Stdio::piped())
117+ .stderr(Stdio::null())
118+ .spawn()
119+ .ok()?;
120+ child.stdin.take()?.write_all(server.as_bytes()).ok()?;
121+ let output = child.wait_with_output().ok()?;
122+ if !output.status.success() {
123+ return None;
124+ }
125+ parse_helper_output(&output.stdout)
126+}
127+
128+pub fn parse_helper_output(stdout: &[u8]) -> Option<Credentials> {
129+ let found: Value = serde_json::from_slice(stdout).ok()?;
130+ let secret = found.get("Secret")?.as_str()?.to_owned();
131+ let username = found
132+ .get("Username")
133+ .and_then(Value::as_str)
134+ .unwrap_or(USERNAME)
135+ .to_owned();
136+ (!secret.is_empty()).then_some(Credentials { username, secret })
137+}
138+
139+#[cfg(test)]
140+mod tests {
141+ use super::*;
142+ use std::cell::RefCell;
143+
144+ fn no_helper(_: &str, _: &str) -> Option<Credentials> {
145+ None
146+ }
147+
148+ fn creds(user: &str, secret: &str) -> Credentials {
149+ Credentials {
150+ username: user.into(),
151+ secret: secret.into(),
152+ }
153+ }
154+
155+ #[test]
156+ fn a_plain_auth_entry_is_decoded() {
157+ // "syntaqx:g1t_secret"
158+ let config = r#"{"auths":{"g1t.sh":{"auth":"c3ludGFxeDpnMXRfc2VjcmV0"}}}"#;
159+ assert_eq!(
160+ from_docker_config(config, "g1t.sh", &no_helper),
161+ Some(creds("syntaqx", "g1t_secret"))
162+ );
163+ }
164+
165+ #[test]
166+ fn auth_keys_may_be_urls() {
167+ let config = r#"{"auths":{"https://g1t.sh/v2/":{"auth":"c3ludGFxeDpnMXRfc2VjcmV0"}}}"#;
168+ assert_eq!(
169+ from_docker_config(config, "g1t.sh", &no_helper),
170+ Some(creds("syntaqx", "g1t_secret"))
171+ );
172+ assert_eq!(
173+ from_docker_config(config, "other.example", &no_helper),
174+ None
175+ );
176+ }
177+
178+ #[test]
179+ fn a_credential_store_is_asked_for_the_host() {
180+ let config = r#"{"auths":{"g1t.sh":{}},"credsStore":"desktop"}"#;
181+ let asked = RefCell::new(Vec::new());
182+ let helper = |name: &str, server: &str| {
183+ asked.borrow_mut().push(format!("{name} {server}"));
184+ (server == "g1t.sh").then(|| creds("syntaqx", "from-store"))
185+ };
186+ assert_eq!(
187+ from_docker_config(config, "g1t.sh", &helper),
188+ Some(creds("syntaqx", "from-store"))
189+ );
190+ assert_eq!(asked.borrow().as_slice(), ["desktop g1t.sh"]);
191+ }
192+
193+ #[test]
194+ fn a_credential_store_without_the_host_falls_back_to_auths() {
195+ let config =
196+ r#"{"auths":{"g1t.sh":{"auth":"c3ludGFxeDpnMXRfc2VjcmV0"}},"credsStore":"desktop"}"#;
197+ assert_eq!(
198+ from_docker_config(config, "g1t.sh", &no_helper),
199+ Some(creds("syntaqx", "g1t_secret"))
200+ );
201+ }
202+
203+ #[test]
204+ fn a_host_helper_comes_before_the_store() {
205+ let config = r#"{"credsStore":"desktop","credHelpers":{"g1t.sh":"g1t"}}"#;
206+ let helper = |name: &str, _: &str| (name == "g1t").then(|| creds("u", "from-helper"));
207+ assert_eq!(
208+ from_docker_config(config, "g1t.sh", &helper),
209+ Some(creds("u", "from-helper"))
210+ );
211+ }
212+
213+ #[test]
214+ fn nothing_stored_is_none() {
215+ assert_eq!(
216+ from_docker_config(r#"{"auths":{}}"#, "g1t.sh", &no_helper),
217+ None
218+ );
219+ assert_eq!(from_docker_config("not json", "g1t.sh", &no_helper), None);
220+ }
221+
222+ #[test]
223+ fn helper_output_is_read() {
224+ let out = br#"{"ServerURL":"g1t.sh","Username":"syntaqx","Secret":"s"}"#;
225+ assert_eq!(parse_helper_output(out), Some(creds("syntaqx", "s")));
226+ assert_eq!(parse_helper_output(b"credentials not found"), None);
227+ }
228+}
+943−0
1+//! A local image, read from what `docker save` writes: the OCI layout
2+//! (`index.json`, `oci-layout`, `blobs/sha256/*`) or the classic layout
3+//! (`manifest.json` and a tar per layer), turned into the blobs and the
4+//! manifest a registry takes.
5+//!
6+//! Layers already compressed are pushed as they are, and so is the image's
7+//! own manifest when nothing had to change. An uncompressed layer is gzipped,
8+//! and the manifest names the gzipped blob; the config, and so its
9+//! `diff_ids` (the digests of the uncompressed layers), stays as it was.
10+
11+use flate2::Compression;
12+use flate2::write::GzEncoder;
13+use serde_json::{Value, json};
14+use sha2::{Digest as _, Sha256};
15+use std::collections::HashMap;
16+use std::fs::File;
17+use std::io::{self, BufWriter, Read, Seek, SeekFrom, Write};
18+use std::path::{Path, PathBuf};
19+
20+pub const OCI_INDEX: &str = "application/vnd.oci.image.index.v1+json";
21+pub const OCI_MANIFEST: &str = "application/vnd.oci.image.manifest.v1+json";
22+pub const DOCKER_LIST: &str = "application/vnd.docker.distribution.manifest.list.v2+json";
23+pub const DOCKER_MANIFEST: &str = "application/vnd.docker.distribution.manifest.v2+json";
24+pub const DOCKER_CONFIG: &str = "application/vnd.docker.container.image.v1+json";
25+pub const DOCKER_LAYER_GZIP: &str = "application/vnd.docker.image.rootfs.diff.tar.gzip";
26+pub const OCI_LAYER_GZIP: &str = "application/vnd.oci.image.layer.v1.tar+gzip";
27+
28+/// Where a blob's bytes are.
29+#[derive(Clone, Debug, PartialEq, Eq)]
30+pub enum Source {
31+ /// `size` bytes at `offset` in the saved tar.
32+ Tar { offset: u64 },
33+ /// A file this push wrote (a gzipped layer).
34+ File(PathBuf),
35+}
36+
37+#[derive(Clone, Debug, PartialEq, Eq)]
38+pub struct Blob {
39+ pub digest: String,
40+ pub size: u64,
41+ pub source: Source,
42+}
43+
44+#[derive(Debug)]
45+pub struct Image {
46+ pub config: Blob,
47+ pub layers: Vec<Blob>,
48+ pub manifest: Vec<u8>,
49+ pub media_type: String,
50+ /// Whether any layer was compressed for the push (so the manifest is new).
51+ pub converted: bool,
52+}
53+
54+impl Image {
55+ pub fn manifest_digest(&self) -> String {
56+ sha256_of(&self.manifest)
57+ }
58+}
59+
60+pub fn sha256_of(bytes: &[u8]) -> String {
61+ format!("sha256:{}", hex::encode(Sha256::digest(bytes)))
62+}
63+
64+type Result<T> = std::result::Result<T, String>;
65+
66+/// The files in a tar: where each one's bytes start, and how many.
67+struct TarIndex {
68+ files: HashMap<String, (u64, u64)>,
69+ links: HashMap<String, String>,
70+}
71+
72+/// `./a/../b//c` → `b/c`.
73+fn normalize(path: &str) -> String {
74+ let mut parts: Vec<&str> = Vec::new();
75+ for part in path.split(['/', '\\']) {
76+ match part {
77+ "" | "." => {}
78+ ".." => {
79+ parts.pop();
80+ }
81+ part => parts.push(part),
82+ }
83+ }
84+ parts.join("/")
85+}
86+
87+impl TarIndex {
88+ fn read(file: &mut File) -> Result<TarIndex> {
89+ file.seek(SeekFrom::Start(0)).map_err(|e| e.to_string())?;
90+ let mut archive = tar::Archive::new(&mut *file);
91+ let mut files = HashMap::new();
92+ let mut links = HashMap::new();
93+ let entries = archive
94+ .entries_with_seek()
95+ .map_err(|e| format!("The saved image is not a tar: {e}"))?;
96+ for entry in entries {
97+ let entry = entry.map_err(|e| format!("The saved image is not a readable tar: {e}"))?;
98+ let path = normalize(&entry.path().map_err(|e| e.to_string())?.to_string_lossy());
99+ let kind = entry.header().entry_type();
100+ if kind.is_file() {
101+ files.insert(path, (entry.raw_file_position(), entry.size()));
102+ } else if (kind.is_symlink() || kind.is_hard_link())
103+ && let Ok(Some(target)) = entry.link_name()
104+ {
105+ let target = target.to_string_lossy().into_owned();
106+ // A symlink is relative to its own directory; a hard link
107+ // to the top of the archive.
108+ let resolved = if kind.is_symlink() && !target.starts_with('/') {
109+ let dir = path.rsplit_once('/').map_or("", |(dir, _)| dir);
110+ normalize(&format!("{dir}/{target}"))
111+ } else {
112+ normalize(&target)
113+ };
114+ links.insert(path, resolved);
115+ }
116+ }
117+ Ok(TarIndex { files, links })
118+ }
119+
120+ fn find(&self, path: &str) -> Option<(u64, u64)> {
121+ let mut path = normalize(path);
122+ for _ in 0..16 {
123+ if let Some(found) = self.files.get(&path) {
124+ return Some(*found);
125+ }
126+ path = self.links.get(&path)?.clone();
127+ }
128+ None
129+ }
130+}
131+
132+/// Reads the saved image in `tar`, writing any layer it has to compress
133+/// into `work`.
134+pub fn load(tar: &Path, work: &Path) -> Result<Image> {
135+ let mut file = File::open(tar).map_err(|e| format!("Could not open {}: {e}", tar.display()))?;
136+ let index = TarIndex::read(&mut file)?;
137+ let mut saved = Saved {
138+ file,
139+ index,
140+ work: work.to_owned(),
141+ };
142+ let has_oci = saved.index.find("index.json").is_some();
143+ let has_classic = saved.index.find("manifest.json").is_some();
144+ match (has_oci, has_classic) {
145+ (true, false) => saved.oci(),
146+ (true, true) => saved.oci().or_else(|oci| saved.classic().map_err(|_| oci)),
147+ (false, true) => saved.classic(),
148+ (false, false) => {
149+ Err("The saved image has neither index.json nor manifest.json.".to_owned())
150+ }
151+ }
152+}
153+
154+struct Saved {
155+ file: File,
156+ index: TarIndex,
157+ work: PathBuf,
158+}
159+
160+/// What a layer's first bytes say it is.
161+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
162+enum Compressed {
163+ Gzip,
164+ Zstd,
165+ No,
166+}
167+
168+impl Saved {
169+ fn read_path(&mut self, path: &str) -> Result<Vec<u8>> {
170+ let (offset, size) = self
171+ .index
172+ .find(path)
173+ .ok_or_else(|| format!("The saved image has no {path}."))?;
174+ if size > 64 * 1024 * 1024 {
175+ return Err(format!("{path} is too large to be a manifest or config."));
176+ }
177+ let mut bytes = vec![0; size as usize];
178+ self.file
179+ .seek(SeekFrom::Start(offset))
180+ .map_err(|e| e.to_string())?;
181+ self.file
182+ .read_exact(&mut bytes)
183+ .map_err(|e| format!("Could not read {path}: {e}"))?;
184+ Ok(bytes)
185+ }
186+
187+ fn read_json(&mut self, path: &str) -> Result<(Vec<u8>, Value)> {
188+ let bytes = self.read_path(path)?;
189+ let value =
190+ serde_json::from_slice(&bytes).map_err(|e| format!("{path} is not JSON: {e}"))?;
191+ Ok((bytes, value))
192+ }
193+
194+ fn blob_path(digest: &str) -> Result<String> {
195+ let (algorithm, hex) = digest
196+ .split_once(':')
197+ .ok_or_else(|| format!("`{digest}` is not a digest."))?;
198+ Ok(format!("blobs/{algorithm}/{hex}"))
199+ }
200+
201+ fn has_blob(&self, digest: &str) -> bool {
202+ Self::blob_path(digest).is_ok_and(|path| self.index.find(&path).is_some())
203+ }
204+
205+ fn sniff(&mut self, offset: u64, size: u64) -> Result<Compressed> {
206+ let mut magic = [0u8; 4];
207+ let n = size.min(4) as usize;
208+ self.file
209+ .seek(SeekFrom::Start(offset))
210+ .map_err(|e| e.to_string())?;
211+ self.file
212+ .read_exact(&mut magic[..n])
213+ .map_err(|e| e.to_string())?;
214+ Ok(match magic {
215+ [0x1f, 0x8b, ..] => Compressed::Gzip,
216+ [0x28, 0xb5, 0x2f, 0xfd] => Compressed::Zstd,
217+ _ => Compressed::No,
218+ })
219+ }
220+
221+ /// Hashes `size` bytes at `offset`.
222+ fn hash_range(&mut self, offset: u64, size: u64) -> Result<String> {
223+ self.file
224+ .seek(SeekFrom::Start(offset))
225+ .map_err(|e| e.to_string())?;
226+ let mut hasher = Sha256::new();
227+ io::copy(
228+ &mut (&mut self.file).take(size),
229+ &mut HashWriter(&mut hasher),
230+ )
231+ .map_err(|e| e.to_string())?;
232+ Ok(format!("sha256:{}", hex::encode(hasher.finalize())))
233+ }
234+
235+ /// Gzips `size` bytes at `offset` into the work directory. Returns the
236+ /// gzipped blob and the digest of the bytes before compression.
237+ fn gzip(&mut self, offset: u64, size: u64, n: usize) -> Result<(Blob, String)> {
238+ let path = self.work.join(format!("layer-{n}.tar.gz"));
239+ let out =
240+ File::create(&path).map_err(|e| format!("Could not write {}: {e}", path.display()))?;
241+ let mut compressed_hash = Sha256::new();
242+ let mut counted = Counted {
243+ inner: BufWriter::with_capacity(1 << 20, out),
244+ hasher: &mut compressed_hash,
245+ bytes: 0,
246+ };
247+ let mut plain_hash = Sha256::new();
248+ {
249+ let mut encoder = GzEncoder::new(&mut counted, Compression::default());
250+ self.file
251+ .seek(SeekFrom::Start(offset))
252+ .map_err(|e| e.to_string())?;
253+ let mut reader = (&mut self.file).take(size);
254+ let mut buffer = vec![0u8; 1 << 20];
255+ loop {
256+ let read = reader.read(&mut buffer).map_err(|e| e.to_string())?;
257+ if read == 0 {
258+ break;
259+ }
260+ plain_hash.update(&buffer[..read]);
261+ encoder
262+ .write_all(&buffer[..read])
263+ .map_err(|e| format!("Could not compress a layer: {e}"))?;
264+ }
265+ encoder
266+ .finish()
267+ .map_err(|e| format!("Could not compress a layer: {e}"))?;
268+ }
269+ counted.inner.flush().map_err(|e| e.to_string())?;
270+ let written = counted.bytes;
271+ drop(counted);
272+ let blob = Blob {
273+ digest: format!("sha256:{}", hex::encode(compressed_hash.finalize())),
274+ size: written,
275+ source: Source::File(path),
276+ };
277+ Ok((
278+ blob,
279+ format!("sha256:{}", hex::encode(plain_hash.finalize())),
280+ ))
281+ }
282+
283+ /// A layer as the registry takes it: as it is when compressed, gzipped
284+ /// when not. `diff_id` is checked against what gets compressed.
285+ fn layer(&mut self, path: &str, n: usize, diff_id: Option<&str>) -> Result<(Blob, Compressed)> {
286+ let (offset, size) = self
287+ .index
288+ .find(path)
289+ .ok_or_else(|| format!("The saved image has no layer {path}."))?;
290+ match self.sniff(offset, size)? {
291+ Compressed::No => {
292+ let (blob, plain) = self.gzip(offset, size, n)?;
293+ if let Some(diff_id) = diff_id
294+ && diff_id != plain
295+ {
296+ return Err(format!(
297+ "Layer {n} is {plain}, but the image's config says {diff_id}."
298+ ));
299+ }
300+ Ok((blob, Compressed::No))
301+ }
302+ compressed => {
303+ let digest = match path.strip_prefix("blobs/sha256/") {
304+ Some(hex) if hex.len() == 64 => format!("sha256:{hex}"),
305+ _ => self.hash_range(offset, size)?,
306+ };
307+ Ok((
308+ Blob {
309+ digest,
310+ size,
311+ source: Source::Tar { offset },
312+ },
313+ compressed,
314+ ))
315+ }
316+ }
317+ }
318+
319+ fn config_blob(&mut self, path: &str) -> Result<(Blob, Value)> {
320+ let (offset, size) = self
321+ .index
322+ .find(path)
323+ .ok_or_else(|| format!("The saved image has no config {path}."))?;
324+ let (bytes, config) = self.read_json(path)?;
325+ Ok((
326+ Blob {
327+ digest: sha256_of(&bytes),
328+ size,
329+ source: Source::Tar { offset },
330+ },
331+ config,
332+ ))
333+ }
334+
335+ /// The OCI layout: index.json, down through any image index, to one
336+ /// image manifest whose blobs were saved.
337+ fn oci(&mut self) -> Result<Image> {
338+ let (_, index) = self.read_json("index.json")?;
339+ let mut descriptor = first_manifest(&index).ok_or("index.json lists no image.")?;
340+ let (bytes, manifest, media_type) = loop {
341+ let digest = descriptor
342+ .get("digest")
343+ .and_then(Value::as_str)
344+ .ok_or("A descriptor has no digest.")?
345+ .to_owned();
346+ let (bytes, value) = self.read_json(&Self::blob_path(&digest)?)?;
347+ let media_type = value
348+ .get("mediaType")
349+ .and_then(Value::as_str)
350+ .or_else(|| descriptor.get("mediaType").and_then(Value::as_str))
351+ .unwrap_or(OCI_MANIFEST)
352+ .to_owned();
353+ if media_type == OCI_INDEX
354+ || media_type == DOCKER_LIST
355+ || value.get("manifests").is_some()
356+ {
357+ descriptor = self
358+ .pick_platform(&value)
359+ .ok_or("No image in the saved index has all its layers saved.")?;
360+ continue;
361+ }
362+ break (bytes, value, media_type);
363+ };
364+ let config_digest = manifest
365+ .pointer("/config/digest")
366+ .and_then(Value::as_str)
367+ .ok_or("The manifest has no config.")?;
368+ let (config, config_json) = self.config_blob(&Self::blob_path(config_digest)?)?;
369+ if config.digest != config_digest {
370+ return Err(format!(
371+ "The saved config is {}, not {config_digest}.",
372+ config.digest
373+ ));
374+ }
375+ let diff_ids = diff_ids(&config_json);
376+ let descriptors = manifest
377+ .get("layers")
378+ .and_then(Value::as_array)
379+ .cloned()
380+ .unwrap_or_default();
381+ let mut new_manifest = manifest.clone();
382+ let mut layers = Vec::new();
383+ let mut converted = false;
384+ for (n, descriptor) in descriptors.iter().enumerate() {
385+ let digest = descriptor
386+ .get("digest")
387+ .and_then(Value::as_str)
388+ .ok_or("A layer has no digest.")?;
389+ let declared = descriptor
390+ .get("mediaType")
391+ .and_then(Value::as_str)
392+ .unwrap_or("");
393+ let (blob, compressed) = self.layer(
394+ &Self::blob_path(digest)?,
395+ n,
396+ diff_ids.get(n).map(String::as_str),
397+ )?;
398+ let media_type = layer_media_type(declared, compressed, &media_type);
399+ if blob.digest != digest || media_type != declared {
400+ converted = true;
401+ let entry = &mut new_manifest["layers"][n];
402+ entry["digest"] = json!(blob.digest);
403+ entry["size"] = json!(blob.size);
404+ entry["mediaType"] = json!(media_type);
405+ }
406+ layers.push(blob);
407+ }
408+ let manifest = if converted {
409+ serde_json::to_vec(&new_manifest).map_err(|e| e.to_string())?
410+ } else {
411+ bytes
412+ };
413+ Ok(Image {
414+ config,
415+ layers,
416+ manifest,
417+ media_type,
418+ converted,
419+ })
420+ }
421+
422+ /// Of an index's images, the one for this machine whose blobs were all
423+ /// saved; attestations are not images.
424+ fn pick_platform(&self, index: &Value) -> Option<Value> {
425+ let arch = match std::env::consts::ARCH {
426+ "x86_64" => "amd64",
427+ "aarch64" => "arm64",
428+ other => other,
429+ };
430+ let candidates: Vec<&Value> = index
431+ .get("manifests")?
432+ .as_array()?
433+ .iter()
434+ .filter(|d| d.pointer("/platform/os").and_then(Value::as_str) != Some("unknown"))
435+ .filter(|d| {
436+ d.pointer("/annotations/vnd.docker.reference.type")
437+ .is_none()
438+ })
439+ .filter(|d| {
440+ d.get("digest")
441+ .and_then(Value::as_str)
442+ .is_some_and(|digest| self.complete(digest))
443+ })
444+ .collect();
445+ let native = candidates.iter().find(|d| {
446+ d.pointer("/platform/os").and_then(Value::as_str) == Some("linux")
447+ && d.pointer("/platform/architecture").and_then(Value::as_str) == Some(arch)
448+ });
449+ native.or(candidates.first()).map(|d| (*d).clone())
450+ }
451+
452+ /// Whether a manifest (or index) and everything it names were saved.
453+ fn complete(&self, digest: &str) -> bool {
454+ let Ok(path) = Self::blob_path(digest) else {
455+ return false;
456+ };
457+ let Some((offset, size)) = self.index.find(&path) else {
458+ return false;
459+ };
460+ let Ok(mut file) = self.file.try_clone() else {
461+ return false;
462+ };
463+ let mut bytes = vec![0; size.min(64 * 1024 * 1024) as usize];
464+ if file.seek(SeekFrom::Start(offset)).is_err() || file.read_exact(&mut bytes).is_err() {
465+ return false;
466+ }
467+ let Ok(value) = serde_json::from_slice::<Value>(&bytes) else {
468+ return false;
469+ };
470+ if let Some(children) = value.get("manifests").and_then(Value::as_array) {
471+ return children.iter().any(|c| {
472+ c.get("digest")
473+ .and_then(Value::as_str)
474+ .is_some_and(|d| self.complete(d))
475+ });
476+ }
477+ let config = value.pointer("/config/digest").and_then(Value::as_str);
478+ let layers = value.get("layers").and_then(Value::as_array);
479+ config.is_some_and(|c| self.has_blob(c))
480+ && layers.is_some_and(|l| {
481+ l.iter().all(|d| {
482+ d.get("digest")
483+ .and_then(Value::as_str)
484+ .is_some_and(|d| self.has_blob(d))
485+ })
486+ })
487+ }
488+
489+ /// The classic layout: manifest.json names the config and the layer
490+ /// tars, and the push gets a new Docker manifest.
491+ fn classic(&mut self) -> Result<Image> {
492+ let (_, saved) = self.read_json("manifest.json")?;
493+ let entry = saved
494+ .as_array()
495+ .and_then(|images| images.first())
496+ .ok_or("manifest.json lists no image.")?
497+ .clone();
498+ let config_path = entry
499+ .get("Config")
500+ .and_then(Value::as_str)
501+ .ok_or("manifest.json names no config.")?;
502+ let (config, config_json) = self.config_blob(config_path)?;
503+ let diff_ids = diff_ids(&config_json);
504+ let paths: Vec<String> = entry
505+ .get("Layers")
506+ .and_then(Value::as_array)
507+ .map(|l| {
508+ l.iter()
509+ .filter_map(|p| p.as_str().map(str::to_owned))
510+ .collect()
511+ })
512+ .unwrap_or_default();
513+ if !diff_ids.is_empty() && diff_ids.len() != paths.len() {
514+ return Err(format!(
515+ "The image has {} layers, but its config lists {}.",
516+ paths.len(),
517+ diff_ids.len()
518+ ));
519+ }
520+ let mut layers = Vec::new();
521+ let mut descriptors = Vec::new();
522+ for (n, path) in paths.iter().enumerate() {
523+ let (blob, compressed) = self.layer(path, n, diff_ids.get(n).map(String::as_str))?;
524+ if compressed == Compressed::Zstd {
525+ return Err(format!(
526+ "Layer {n} is zstd, which a Docker manifest cannot name."
527+ ));
528+ }
529+ descriptors.push(
530+ json!({ "mediaType": DOCKER_LAYER_GZIP, "size": blob.size, "digest": blob.digest }),
531+ );
532+ layers.push(blob);
533+ }
534+ let manifest = json!({
535+ "schemaVersion": 2,
536+ "mediaType": DOCKER_MANIFEST,
537+ "config": { "mediaType": DOCKER_CONFIG, "size": config.size, "digest": config.digest },
538+ "layers": descriptors,
539+ });
540+ let manifest = serde_json::to_vec(&manifest).map_err(|e| e.to_string())?;
541+ Ok(Image {
542+ config,
543+ layers,
544+ manifest,
545+ media_type: DOCKER_MANIFEST.to_owned(),
546+ converted: true,
547+ })
548+ }
549+}
550+
551+fn first_manifest(index: &Value) -> Option<Value> {
552+ index.get("manifests")?.as_array()?.first().cloned()
553+}
554+
555+fn diff_ids(config: &Value) -> Vec<String> {
556+ config
557+ .pointer("/rootfs/diff_ids")
558+ .and_then(Value::as_array)
559+ .map(|ids| {
560+ ids.iter()
561+ .filter_map(|d| d.as_str().map(str::to_owned))
562+ .collect()
563+ })
564+ .unwrap_or_default()
565+}
566+
567+/// The media type a layer is pushed with: what it said, when that matches
568+/// its bytes; gzip or zstd in the manifest's family when it does not.
569+fn layer_media_type(declared: &str, compressed: Compressed, manifest_type: &str) -> String {
570+ let docker = manifest_type == DOCKER_MANIFEST;
571+ let says_gzip = declared.ends_with("+gzip") || declared.ends_with(".gzip");
572+ let says_zstd = declared.ends_with("+zstd") || declared.ends_with(".zstd");
573+ match compressed {
574+ Compressed::Gzip if says_gzip => declared.to_owned(),
575+ Compressed::Zstd if says_zstd => declared.to_owned(),
576+ Compressed::Zstd => "application/vnd.oci.image.layer.v1.tar+zstd".to_owned(),
577+ // Gzipped by the push, or gzip that said otherwise.
578+ _ if docker || declared.starts_with("application/vnd.docker.") => {
579+ DOCKER_LAYER_GZIP.to_owned()
580+ }
581+ _ => OCI_LAYER_GZIP.to_owned(),
582+ }
583+}
584+
585+struct HashWriter<'a>(&'a mut Sha256);
586+
587+impl Write for HashWriter<'_> {
588+ fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
589+ self.0.update(buf);
590+ Ok(buf.len())
591+ }
592+ fn flush(&mut self) -> io::Result<()> {
593+ Ok(())
594+ }
595+}
596+
597+/// Writes through, hashing and counting what passes.
598+struct Counted<'a, W: Write> {
599+ inner: W,
600+ hasher: &'a mut Sha256,
601+ bytes: u64,
602+}
603+
604+impl<W: Write> Write for Counted<'_, W> {
605+ fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
606+ let n = self.inner.write(buf)?;
607+ self.hasher.update(&buf[..n]);
608+ self.bytes += n as u64;
609+ Ok(n)
610+ }
611+ fn flush(&mut self) -> io::Result<()> {
612+ self.inner.flush()
613+ }
614+}
615+
616+#[cfg(test)]
617+mod tests {
618+ use super::*;
619+ use flate2::read::GzDecoder;
620+
621+ /// A scratch directory, gone when dropped.
622+ struct Scratch(PathBuf);
623+
624+ impl Scratch {
625+ fn new(name: &str) -> Scratch {
626+ let nanos = std::time::SystemTime::now()
627+ .duration_since(std::time::UNIX_EPOCH)
628+ .unwrap()
629+ .as_nanos();
630+ let dir = std::env::temp_dir()
631+ .join(format!("g1t-test-{name}-{}-{nanos}", std::process::id()));
632+ std::fs::create_dir_all(&dir).unwrap();
633+ Scratch(dir)
634+ }
635+ }
636+
637+ impl Drop for Scratch {
638+ fn drop(&mut self) {
639+ let _ = std::fs::remove_dir_all(&self.0);
640+ }
641+ }
642+
643+ /// A layer: a tar holding one file.
644+ fn layer_tar(name: &str, contents: &[u8]) -> Vec<u8> {
645+ let mut builder = tar::Builder::new(Vec::new());
646+ let mut header = tar::Header::new_gnu();
647+ header.set_size(contents.len() as u64);
648+ header.set_mode(0o644);
649+ header.set_cksum();
650+ builder.append_data(&mut header, name, contents).unwrap();
651+ builder.into_inner().unwrap()
652+ }
653+
654+ fn gzip(bytes: &[u8]) -> Vec<u8> {
655+ let mut encoder = GzEncoder::new(Vec::new(), Compression::fast());
656+ encoder.write_all(bytes).unwrap();
657+ encoder.finish().unwrap()
658+ }
659+
660+ fn hex_of(bytes: &[u8]) -> String {
661+ hex::encode(Sha256::digest(bytes))
662+ }
663+
664+ enum Entry<'a> {
665+ File(&'a str, Vec<u8>),
666+ Symlink(&'a str, &'a str),
667+ }
668+
669+ fn write_tar(path: &Path, entries: Vec<Entry>) {
670+ let mut builder = tar::Builder::new(File::create(path).unwrap());
671+ for entry in entries {
672+ let mut header = tar::Header::new_gnu();
673+ header.set_mode(0o644);
674+ match entry {
675+ Entry::File(name, bytes) => {
676+ header.set_size(bytes.len() as u64);
677+ header.set_cksum();
678+ builder
679+ .append_data(&mut header, name, bytes.as_slice())
680+ .unwrap();
681+ }
682+ Entry::Symlink(name, target) => {
683+ header.set_entry_type(tar::EntryType::Symlink);
684+ header.set_size(0);
685+ builder.append_link(&mut header, name, target).unwrap();
686+ }
687+ }
688+ }
689+ builder.finish().unwrap();
690+ }
691+
692+ fn config_for(layers: &[&[u8]]) -> Vec<u8> {
693+ let ids: Vec<String> = layers
694+ .iter()
695+ .map(|l| format!("sha256:{}", hex_of(l)))
696+ .collect();
697+ serde_json::to_vec(&json!({
698+ "architecture": "amd64", "os": "linux",
699+ "config": { "Cmd": ["sh"] },
700+ "rootfs": { "type": "layers", "diff_ids": ids },
701+ }))
702+ .unwrap()
703+ }
704+
705+ fn read_blob(blob: &Blob, tar: &Path) -> Vec<u8> {
706+ match &blob.source {
707+ Source::File(path) => std::fs::read(path).unwrap(),
708+ Source::Tar { offset } => {
709+ let mut file = File::open(tar).unwrap();
710+ file.seek(SeekFrom::Start(*offset)).unwrap();
711+ let mut bytes = vec![0; blob.size as usize];
712+ file.read_exact(&mut bytes).unwrap();
713+ bytes
714+ }
715+ }
716+ }
717+
718+ fn gunzip(bytes: &[u8]) -> Vec<u8> {
719+ let mut out = Vec::new();
720+ GzDecoder::new(bytes).read_to_end(&mut out).unwrap();
721+ out
722+ }
723+
724+ /// Every blob's digest and size match its bytes, and the manifest names
725+ /// exactly the blobs that will be pushed.
726+ fn assert_consistent(image: &Image, tar: &Path) {
727+ let manifest: Value = serde_json::from_slice(&image.manifest).unwrap();
728+ assert_eq!(manifest["mediaType"], image.media_type);
729+ let config_bytes = read_blob(&image.config, tar);
730+ assert_eq!(sha256_of(&config_bytes), image.config.digest);
731+ assert_eq!(manifest["config"]["digest"], image.config.digest);
732+ assert_eq!(manifest["config"]["size"], image.config.size);
733+ let config: Value = serde_json::from_slice(&config_bytes).unwrap();
734+ let listed = manifest["layers"].as_array().unwrap();
735+ assert_eq!(listed.len(), image.layers.len());
736+ for (n, (blob, descriptor)) in image.layers.iter().zip(listed).enumerate() {
737+ let bytes = read_blob(blob, tar);
738+ assert_eq!(bytes.len() as u64, blob.size);
739+ assert_eq!(sha256_of(&bytes), blob.digest);
740+ assert_eq!(descriptor["digest"], blob.digest);
741+ assert_eq!(descriptor["size"], blob.size);
742+ assert!(descriptor["mediaType"].as_str().unwrap().ends_with("gzip"));
743+ // The uncompressed layer is the one the config's diff_ids name.
744+ assert_eq!(config["rootfs"]["diff_ids"][n], sha256_of(&gunzip(&bytes)));
745+ }
746+ }
747+
748+ #[test]
749+ fn the_classic_layout_gets_gzipped_layers_and_a_docker_manifest() {
750+ let scratch = Scratch::new("classic");
751+ let one = layer_tar("hello", b"hello world");
752+ let two = layer_tar("blob", &vec![7u8; 300_000]);
753+ let config = config_for(&[&one, &two]);
754+ let config_name = format!("{}.json", hex_of(&config));
755+ let saved = json!([{ "Config": config_name, "RepoTags": ["web:1"], "Layers": ["aaa/layer.tar", "bbb/layer.tar"] }]);
756+ let tar = scratch.0.join("save.tar");
757+ write_tar(
758+ &tar,
759+ vec![
760+ Entry::File(&config_name, config.clone()),
761+ Entry::File("aaa/layer.tar", one.clone()),
762+ Entry::File("bbb/layer.tar", two.clone()),
763+ Entry::File("manifest.json", serde_json::to_vec(&saved).unwrap()),
764+ ],
765+ );
766+ let image = load(&tar, &scratch.0).unwrap();
767+ assert!(image.converted);
768+ assert_eq!(image.media_type, DOCKER_MANIFEST);
769+ assert_eq!(image.config.digest, format!("sha256:{}", hex_of(&config)));
770+ assert_consistent(&image, &tar);
771+ let manifest: Value = serde_json::from_slice(&image.manifest).unwrap();
772+ assert_eq!(manifest["config"]["mediaType"], DOCKER_CONFIG);
773+ assert_eq!(manifest["layers"][0]["mediaType"], DOCKER_LAYER_GZIP);
774+ assert_eq!(image.manifest_digest(), sha256_of(&image.manifest));
775+ }
776+
777+ #[test]
778+ fn the_classic_layout_through_docker_25_symlinks() {
779+ let scratch = Scratch::new("symlinks");
780+ let one = layer_tar("hello", b"hi");
781+ let config = config_for(&[&one]);
782+ let (config_hex, layer_hex) = (hex_of(&config), hex_of(&one));
783+ let saved = json!([{ "Config": format!("{config_hex}.json"), "Layers": [format!("{layer_hex}/layer.tar")] }]);
784+ let tar = scratch.0.join("save.tar");
785+ write_tar(
786+ &tar,
787+ vec![
788+ Entry::File(&format!("blobs/sha256/{config_hex}"), config.clone()),
789+ Entry::File(&format!("blobs/sha256/{layer_hex}"), one.clone()),
790+ Entry::Symlink(
791+ &format!("{config_hex}.json"),
792+ &format!("blobs/sha256/{config_hex}"),
793+ ),
794+ Entry::Symlink(
795+ &format!("{layer_hex}/layer.tar"),
796+ &format!("../blobs/sha256/{layer_hex}"),
797+ ),
798+ Entry::File("manifest.json", serde_json::to_vec(&saved).unwrap()),
799+ ],
800+ );
801+ let image = load(&tar, &scratch.0).unwrap();
802+ assert_consistent(&image, &tar);
803+ }
804+
805+ #[test]
806+ fn a_layer_that_is_not_its_diff_id_is_refused() {
807+ let scratch = Scratch::new("mismatch");
808+ let one = layer_tar("hello", b"hello");
809+ let config = config_for(&[&layer_tar("other", b"other")]);
810+ let saved = json!([{ "Config": "c.json", "Layers": ["l/layer.tar"] }]);
811+ let tar = scratch.0.join("save.tar");
812+ write_tar(
813+ &tar,
814+ vec![
815+ Entry::File("c.json", config),
816+ Entry::File("l/layer.tar", one),
817+ Entry::File("manifest.json", serde_json::to_vec(&saved).unwrap()),
818+ ],
819+ );
820+ let error = load(&tar, &scratch.0).unwrap_err();
821+ assert!(error.contains("config says"), "{error}");
822+ }
823+
824+ /// The OCI layout Docker 25+ writes from its own image store:
825+ /// uncompressed layers, which the push gzips.
826+ #[test]
827+ fn the_oci_layout_with_uncompressed_layers_is_gzipped() {
828+ let scratch = Scratch::new("oci-plain");
829+ let one = layer_tar("hello", b"hello world");
830+ let config = config_for(&[&one]);
831+ let manifest = serde_json::to_vec(&json!({
832+ "schemaVersion": 2, "mediaType": OCI_MANIFEST,
833+ "config": { "mediaType": "application/vnd.oci.image.config.v1+json", "digest": format!("sha256:{}", hex_of(&config)), "size": config.len() },
834+ "layers": [{ "mediaType": "application/vnd.oci.image.layer.v1.tar", "digest": format!("sha256:{}", hex_of(&one)), "size": one.len() }],
835+ "annotations": { "org.opencontainers.image.created": "2026-10-06T00:00:00Z" },
836+ }))
837+ .unwrap();
838+ let index = json!({ "schemaVersion": 2, "manifests": [{ "mediaType": OCI_MANIFEST, "digest": format!("sha256:{}", hex_of(&manifest)), "size": manifest.len() }] });
839+ let tar = scratch.0.join("save.tar");
840+ write_tar(
841+ &tar,
842+ vec![
843+ Entry::File(&format!("blobs/sha256/{}", hex_of(&config)), config.clone()),
844+ Entry::File(&format!("blobs/sha256/{}", hex_of(&one)), one.clone()),
845+ Entry::File(
846+ &format!("blobs/sha256/{}", hex_of(&manifest)),
847+ manifest.clone(),
848+ ),
849+ Entry::File("index.json", serde_json::to_vec(&index).unwrap()),
850+ Entry::File("oci-layout", br#"{"imageLayoutVersion":"1.0.0"}"#.to_vec()),
851+ Entry::File("manifest.json", b"[]".to_vec()),
852+ ],
853+ );
854+ let image = load(&tar, &scratch.0).unwrap();
855+ assert!(image.converted);
856+ assert_eq!(image.media_type, OCI_MANIFEST);
857+ assert_consistent(&image, &tar);
858+ let pushed: Value = serde_json::from_slice(&image.manifest).unwrap();
859+ assert_eq!(pushed["layers"][0]["mediaType"], OCI_LAYER_GZIP);
860+ // Everything else in the manifest stays.
861+ assert_eq!(
862+ pushed["annotations"]["org.opencontainers.image.created"],
863+ "2026-10-06T00:00:00Z"
864+ );
865+ }
866+
867+ /// The containerd image store: an index of indexes, gzipped layers, and
868+ /// only this machine's platform saved. The manifest goes up unchanged.
869+ #[test]
870+ fn the_oci_layout_with_compressed_layers_keeps_its_manifest() {
871+ let scratch = Scratch::new("oci-gzip");
872+ let plain = layer_tar("hello", b"hello world");
873+ let one = gzip(&plain);
874+ let config = config_for(&[&plain]);
875+ let manifest = serde_json::to_vec(&json!({
876+ "schemaVersion": 2, "mediaType": OCI_MANIFEST,
877+ "config": { "mediaType": "application/vnd.oci.image.config.v1+json", "digest": format!("sha256:{}", hex_of(&config)), "size": config.len() },
878+ "layers": [{ "mediaType": OCI_LAYER_GZIP, "digest": format!("sha256:{}", hex_of(&one)), "size": one.len() }],
879+ }))
880+ .unwrap();
881+ let arch = match std::env::consts::ARCH {
882+ "x86_64" => "amd64",
883+ "aarch64" => "arm64",
884+ other => other,
885+ };
886+ let platforms = serde_json::to_vec(&json!({
887+ "schemaVersion": 2, "mediaType": OCI_INDEX,
888+ "manifests": [
889+ { "mediaType": OCI_MANIFEST, "digest": format!("sha256:{}", "1".repeat(64)), "size": 1, "platform": { "os": "linux", "architecture": "s390x" } },
890+ { "mediaType": OCI_MANIFEST, "digest": format!("sha256:{}", hex_of(&manifest)), "size": manifest.len(), "platform": { "os": "linux", "architecture": arch } },
891+ { "mediaType": OCI_MANIFEST, "digest": format!("sha256:{}", "2".repeat(64)), "size": 1, "platform": { "os": "unknown", "architecture": "unknown" } },
892+ ],
893+ }))
894+ .unwrap();
895+ let index = json!({ "schemaVersion": 2, "mediaType": OCI_INDEX, "manifests": [{ "mediaType": OCI_INDEX, "digest": format!("sha256:{}", hex_of(&platforms)), "size": platforms.len() }] });
896+ let tar = scratch.0.join("save.tar");
897+ write_tar(
898+ &tar,
899+ vec![
900+ Entry::File(&format!("blobs/sha256/{}", hex_of(&config)), config.clone()),
901+ Entry::File(&format!("blobs/sha256/{}", hex_of(&one)), one.clone()),
902+ Entry::File(
903+ &format!("blobs/sha256/{}", hex_of(&manifest)),
904+ manifest.clone(),
905+ ),
906+ Entry::File(
907+ &format!("blobs/sha256/{}", hex_of(&platforms)),
908+ platforms.clone(),
909+ ),
910+ Entry::File("index.json", serde_json::to_vec(&index).unwrap()),
911+ ],
912+ );
913+ let image = load(&tar, &scratch.0).unwrap();
914+ assert!(!image.converted);
915+ assert_eq!(image.manifest, manifest);
916+ assert_eq!(
917+ image.manifest_digest(),
918+ format!("sha256:{}", hex_of(&manifest))
919+ );
920+ assert_eq!(
921+ image.layers[0].source,
922+ Source::Tar {
923+ offset: image.layers[0].source_offset()
924+ }
925+ );
926+ assert_consistent(&image, &tar);
927+ }
928+
929+ impl Blob {
930+ fn source_offset(&self) -> u64 {
931+ match self.source {
932+ Source::Tar { offset } => offset,
933+ Source::File(_) => panic!("expected a blob in the tar"),
934+ }
935+ }
936+ }
937+
938+ #[test]
939+ fn paths_normalize() {
940+ assert_eq!(normalize("./a/../b//c"), "b/c");
941+ assert_eq!(normalize("blobs\\sha256\\x"), "blobs/sha256/x");
942+ }
943+}
+299−2
1−fn main() {
2− println!("Hello, world!");
1+//! The g1t command line.
2+
3+mod chunk;
4+mod credentials;
5+mod image;
6+mod reference;
7+mod registry;
8+
9+use reference::Reference;
10+use std::path::{Path, PathBuf};
11+use std::process::{Command, ExitCode, Stdio};
12+
13+const HELP: &str = "\
14+g1t, the command line for g1t.sh
15+
16+Usage:
17+ g1t push <image>[:<tag>] [--as g1t.sh/<workspace>/<name>:<tag>]
18+ [--chunk-size 90MB] [--token-stdin]
19+ g1t push --archive <file.tar> --as g1t.sh/<workspace>/<name>:<tag>
20+ g1t help
21+ g1t --version
22+
23+Commands:
24+ push Send a local docker image to g1t.sh. Each layer goes up in
25+ chunks, so a layer can be any size.
26+
27+Push options:
28+ --as <address> Where to push, when the image's own name is not an
29+ address on g1t.sh. The tag defaults to latest.
30+ --chunk-size <n> How much of a layer each request carries: 5MB to
31+ 95MB, 90MB unless set. MB and MiB both mean 1,048,576
32+ bytes.
33+ --token-stdin Read the g1t token from stdin.
34+ --archive <file> Push a tarball `docker save` wrote, instead of asking
35+ docker for the image.
36+
37+The token is the first of: --token-stdin, the G1T_TOKEN environment
38+variable, or what `docker login g1t.sh` stored.
39+";
40+
41+fn main() -> ExitCode {
42+ let args: Vec<String> = std::env::args().skip(1).collect();
43+ let result = match args.first().map(String::as_str) {
44+ None | Some("help" | "-h" | "--help") => {
45+ print!("{HELP}");
46+ Ok(())
47+ }
48+ Some("-V" | "--version" | "version") => {
49+ println!("g1t {}", env!("CARGO_PKG_VERSION"));
50+ Ok(())
51+ }
52+ Some("push") if args.iter().any(|a| a == "-h" || a == "--help") => {
53+ print!("{HELP}");
54+ Ok(())
55+ }
56+ Some("push") => PushArgs::parse(&args[1..]).and_then(|push| push.run()),
57+ Some(other) => Err(format!("Unknown command `{other}`. See `g1t help`.")),
58+ };
59+ match result {
60+ Ok(()) => ExitCode::SUCCESS,
61+ Err(error) => {
62+ eprintln!("g1t: {error}");
63+ ExitCode::FAILURE
64+ }
65+ }
66+}
67+
68+#[derive(Debug, PartialEq, Eq)]
69+struct PushArgs {
70+ /// The local image, or the tarball it was saved to.
71+ image: Option<String>,
72+ archive: Option<PathBuf>,
73+ target: Reference,
74+ chunk_size: u64,
75+ token_stdin: bool,
76+}
77+
78+impl PushArgs {
79+ fn parse(args: &[String]) -> Result<PushArgs, String> {
80+ let mut image = None;
81+ let mut target = None;
82+ let mut chunk_size = chunk::DEFAULT_CHUNK;
83+ let mut token_stdin = false;
84+ let mut archive = None;
85+ let mut args = args.iter();
86+ while let Some(arg) = args.next() {
87+ let (flag, inline) = match arg.split_once('=') {
88+ Some((flag, value)) if flag.starts_with("--") => (flag, Some(value.to_owned())),
89+ _ => (arg.as_str(), None),
90+ };
91+ let mut value = |name: &str| {
92+ inline
93+ .clone()
94+ .or_else(|| args.next().cloned())
95+ .ok_or_else(|| format!("{name} needs a value."))
96+ };
97+ match flag {
98+ "--as" => target = Some(value("--as")?),
99+ "--chunk-size" => {
100+ chunk_size = chunk::check_chunk(chunk::parse_size(&value("--chunk-size")?)?)?
101+ }
102+ "--token-stdin" => token_stdin = true,
103+ "--archive" => archive = Some(PathBuf::from(value("--archive")?)),
104+ flag if flag.starts_with('-') => {
105+ return Err(format!("Unknown option `{flag}`. See `g1t help`."));
106+ }
107+ _ if image.is_none() => image = Some(arg.clone()),
108+ _ => return Err(format!("`{arg}`: push takes one image.")),
109+ }
110+ }
111+ if image.is_some() && archive.is_some() {
112+ return Err("Push an image or an --archive, not both.".to_owned());
113+ }
114+ let target = match (target, &image) {
115+ (Some(address), _) => Reference::parse(&address)?
116+ .ok_or_else(|| format!("`{address}` does not name a registry, like g1t.sh/<workspace>/<name>:<tag>."))?,
117+ (None, Some(image)) => Reference::parse(image)?.ok_or_else(|| {
118+ format!("`{image}` is not an address on g1t.sh. Say where to push it: --as g1t.sh/<workspace>/<name>:<tag>.")
119+ })?,
120+ (None, None) if archive.is_some() => return Err("Say where to push the archive: --as g1t.sh/<workspace>/<name>:<tag>.".to_owned()),
121+ (None, None) => return Err("Name the image to push: g1t push <image>[:<tag>].".to_owned()),
122+ };
123+ Ok(PushArgs {
124+ image,
125+ archive,
126+ target,
127+ chunk_size,
128+ token_stdin,
129+ })
130+ }
131+
132+ fn run(self) -> Result<(), String> {
133+ let credentials = credentials::find(&self.target.host, self.token_stdin)?;
134+ let work = WorkDir::new()?;
135+ let tar = match (&self.archive, &self.image) {
136+ (Some(archive), _) => {
137+ println!("Reading {}", archive.display());
138+ archive.clone()
139+ }
140+ (None, image) => {
141+ let image = image.as_deref().unwrap_or_default();
142+ println!("Reading {image} from docker");
143+ let tar = work.0.join("image.tar");
144+ docker_save(image, &tar)?;
145+ tar
146+ }
147+ };
148+ let image = image::load(&tar, &work.0)?;
149+ let compressed = image
150+ .layers
151+ .iter()
152+ .filter(|l| matches!(l.source, image::Source::File(_)))
153+ .count();
154+ if image.converted {
155+ println!(
156+ "Compressed {compressed} layer(s) saved uncompressed, and wrote the image a manifest naming them"
157+ );
158+ }
159+ let mut registry = registry::Registry::new(
160+ self.target.base_url(),
161+ self.target.host.clone(),
162+ self.target.name.clone(),
163+ credentials,
164+ );
165+ registry.sign_in()?;
166+ println!("Pushing to {}", self.target);
167+ let blobs = std::iter::once(("config", &image.config))
168+ .chain(image.layers.iter().map(|l| ("layer", l)));
169+ let mut seen = std::collections::HashSet::new();
170+ for (kind, blob) in blobs {
171+ if !seen.insert(blob.digest.clone()) {
172+ continue;
173+ }
174+ let short = &blob.digest[..blob.digest.len().min(19)];
175+ let size = chunk::human(blob.size);
176+ let outcome = registry.push_blob(blob, &tar, self.chunk_size)?;
177+ let said = match outcome {
178+ registry::Outcome::Exists => "already on the registry, skipped".to_owned(),
179+ registry::Outcome::Uploaded { chunks: 1 } => "uploaded".to_owned(),
180+ registry::Outcome::Uploaded { chunks } => format!("uploaded in {chunks} chunks"),
181+ };
182+ println!(" {kind:<6} {short} {size:>10} {said}");
183+ }
184+ let ours = image.manifest_digest();
185+ let theirs =
186+ registry.push_manifest(&self.target.tag, &image.media_type, &image.manifest)?;
187+ if !theirs.is_empty() && theirs != ours {
188+ return Err(format!(
189+ "The registry kept the manifest as {theirs}, not {ours}."
190+ ));
191+ }
192+ println!("Pushed {}", self.target);
193+ println!("digest: {ours}");
194+ Ok(())
195+ }
196+}
197+
198+/// `docker save <image>`, streamed into `out`.
199+fn docker_save(image: &str, out: &Path) -> Result<(), String> {
200+ let mut file = std::fs::File::create(out)
201+ .map_err(|e| format!("Could not write {}: {e}", out.display()))?;
202+ let mut child = Command::new("docker")
203+ .args(["save", image])
204+ .stdout(Stdio::piped())
205+ .stderr(Stdio::piped())
206+ .spawn()
207+ .map_err(|e| format!("Could not run docker: {e}"))?;
208+ let mut stdout = child.stdout.take().ok_or("docker gave no output.")?;
209+ let copied = std::io::copy(&mut stdout, &mut file);
210+ let output = child
211+ .wait_with_output()
212+ .map_err(|e| format!("docker save failed: {e}"))?;
213+ if !output.status.success() {
214+ let stderr = String::from_utf8_lossy(&output.stderr);
215+ return Err(format!("docker save {image} failed: {}", stderr.trim()));
216+ }
217+ copied.map_err(|e| format!("Could not save the image: {e}"))?;
218+ Ok(())
219+}
220+
221+/// A directory for the saved image and compressed layers, removed after.
222+struct WorkDir(PathBuf);
223+
224+impl WorkDir {
225+ fn new() -> Result<WorkDir, String> {
226+ let nanos = std::time::SystemTime::now()
227+ .duration_since(std::time::UNIX_EPOCH)
228+ .map(|d| d.as_nanos())
229+ .unwrap_or_default();
230+ let dir = std::env::temp_dir().join(format!("g1t-push-{}-{nanos}", std::process::id()));
231+ std::fs::create_dir_all(&dir)
232+ .map_err(|e| format!("Could not make {}: {e}", dir.display()))?;
233+ Ok(WorkDir(dir))
234+ }
235+}
236+
237+impl Drop for WorkDir {
238+ fn drop(&mut self) {
239+ let _ = std::fs::remove_dir_all(&self.0);
240+ }
241+}
242+
243+#[cfg(test)]
244+mod tests {
245+ use super::*;
246+
247+ fn parse(args: &[&str]) -> Result<PushArgs, String> {
248+ PushArgs::parse(&args.iter().map(|s| s.to_string()).collect::<Vec<_>>())
249+ }
250+
251+ #[test]
252+ fn an_image_named_for_g1t_pushes_to_itself() {
253+ let push = parse(&["g1t.sh/acme/web:1"]).unwrap();
254+ assert_eq!(push.target.to_string(), "g1t.sh/acme/web:1");
255+ assert_eq!(push.chunk_size, chunk::DEFAULT_CHUNK);
256+ assert!(!push.token_stdin);
257+ }
258+
259+ #[test]
260+ fn as_names_the_destination() {
261+ let push = parse(&[
262+ "web:dev",
263+ "--as",
264+ "g1t.sh/acme/web:2",
265+ "--chunk-size=50MB",
266+ "--token-stdin",
267+ ])
268+ .unwrap();
269+ assert_eq!(push.image.as_deref(), Some("web:dev"));
270+ assert_eq!(push.target.to_string(), "g1t.sh/acme/web:2");
271+ assert_eq!(push.chunk_size, 50 * chunk::MIB);
272+ assert!(push.token_stdin);
273+ }
274+
275+ #[test]
276+ fn bad_pushes_are_explained() {
277+ assert!(parse(&["web:dev"]).unwrap_err().contains("--as"));
278+ assert!(parse(&[]).is_err());
279+ assert!(
280+ parse(&["g1t.sh/acme/web", "--chunk-size", "200MB"])
281+ .unwrap_err()
282+ .contains("95MiB")
283+ );
284+ assert!(parse(&["g1t.sh/acme/web", "--bogus"]).is_err());
285+ assert!(parse(&["a", "b", "--as", "g1t.sh/acme/web"]).is_err());
286+ assert!(
287+ parse(&["--archive", "web.tar"])
288+ .unwrap_err()
289+ .contains("--as")
290+ );
291+ assert!(parse(&["web:dev", "--archive", "web.tar", "--as", "g1t.sh/acme/web"]).is_err());
292+ }
293+
294+ #[test]
295+ fn an_archive_pushes_without_docker() {
296+ let push = parse(&["--archive", "web.tar", "--as", "g1t.sh/acme/web:3"]).unwrap();
297+ assert_eq!(push.archive, Some(PathBuf::from("web.tar")));
298+ assert_eq!(push.image, None);
299+ }
3300 }
+120−0
1+//! Image addresses: `g1t.sh/<workspace>/<name>[:<tag>]`.
2+
3+#[derive(Clone, Debug, PartialEq, Eq)]
4+pub struct Reference {
5+ /// The registry, `g1t.sh` or a host of your own (`localhost:8790`).
6+ pub host: String,
7+ /// The image's name in the registry, `<workspace>/<name>`.
8+ pub name: String,
9+ pub tag: String,
10+}
11+
12+impl Reference {
13+ /// Reads an address that names its registry. `None` when the first part
14+ /// is not a host (a local name like `web:1.0`).
15+ pub fn parse(text: &str) -> Result<Option<Reference>, String> {
16+ let text = text.trim();
17+ if text.contains('@') {
18+ return Err(format!("`{text}`: push to a tag, not a digest."));
19+ }
20+ let Some((host, rest)) = text.split_once('/') else {
21+ return Ok(None);
22+ };
23+ if !(host.contains('.') || host.contains(':') || host == "localhost") {
24+ return Ok(None);
25+ }
26+ // A colon after the last slash starts the tag.
27+ let (name, tag) = match rest.rsplit_once(':') {
28+ Some((name, tag)) if !tag.contains('/') => (name, tag),
29+ _ => (rest, "latest"),
30+ };
31+ if name.is_empty() || tag.is_empty() {
32+ return Err(format!(
33+ "`{text}` is not an image address, like g1t.sh/<workspace>/<name>:<tag>."
34+ ));
35+ }
36+ if !name.contains('/') && host == "g1t.sh" {
37+ return Err(format!(
38+ "`{text}`: images on g1t.sh start with a workspace, like g1t.sh/<workspace>/<name>:<tag>."
39+ ));
40+ }
41+ let valid = |c: char| c.is_ascii_lowercase() || c.is_ascii_digit() || "._-/".contains(c);
42+ if !name.chars().all(valid) {
43+ return Err(format!(
44+ "`{name}`: image names are lowercase letters and digits, separated by `.`, `_`, `-` or `/`."
45+ ));
46+ }
47+ let valid_tag = |c: char| c.is_ascii_alphanumeric() || "._-".contains(c);
48+ if tag.len() > 128 || !tag.chars().all(valid_tag) {
49+ return Err(format!(
50+ "`{tag}` is not a tag: letters, digits, `.`, `_` and `-`, at most 128."
51+ ));
52+ }
53+ Ok(Some(Reference {
54+ host: host.to_owned(),
55+ name: name.to_owned(),
56+ tag: tag.to_owned(),
57+ }))
58+ }
59+
60+ /// Where the registry answers: plain HTTP on this machine, HTTPS elsewhere.
61+ pub fn base_url(&self) -> String {
62+ let local = self.host.starts_with("localhost") || self.host.starts_with("127.0.0.1");
63+ format!("{}://{}", if local { "http" } else { "https" }, self.host)
64+ }
65+}
66+
67+impl std::fmt::Display for Reference {
68+ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
69+ write!(f, "{}/{}:{}", self.host, self.name, self.tag)
70+ }
71+}
72+
73+#[cfg(test)]
74+mod tests {
75+ use super::*;
76+
77+ fn parse(text: &str) -> Reference {
78+ Reference::parse(text).unwrap().unwrap()
79+ }
80+
81+ #[test]
82+ fn addresses_split_into_host_name_and_tag() {
83+ let r = parse("g1t.sh/acme/web:1.4.0");
84+ assert_eq!(
85+ (r.host.as_str(), r.name.as_str(), r.tag.as_str()),
86+ ("g1t.sh", "acme/web", "1.4.0")
87+ );
88+ assert_eq!(r.base_url(), "https://g1t.sh");
89+ assert_eq!(r.to_string(), "g1t.sh/acme/web:1.4.0");
90+ }
91+
92+ #[test]
93+ fn the_tag_defaults_to_latest() {
94+ assert_eq!(parse("g1t.sh/acme/tools/web").tag, "latest");
95+ }
96+
97+ #[test]
98+ fn a_host_with_a_port_is_not_a_tag() {
99+ let r = parse("localhost:8790/acme/web");
100+ assert_eq!(
101+ (r.host.as_str(), r.name.as_str(), r.tag.as_str()),
102+ ("localhost:8790", "acme/web", "latest")
103+ );
104+ assert_eq!(r.base_url(), "http://localhost:8790");
105+ }
106+
107+ #[test]
108+ fn local_names_have_no_registry() {
109+ assert_eq!(Reference::parse("web:1.0"), Ok(None));
110+ assert_eq!(Reference::parse("library/alpine:3.20"), Ok(None));
111+ }
112+
113+ #[test]
114+ fn bad_addresses_are_refused() {
115+ assert!(Reference::parse("g1t.sh/web:1").is_err());
116+ assert!(Reference::parse("g1t.sh/Acme/Web:1").is_err());
117+ assert!(Reference::parse("g1t.sh/acme/web@sha256:00").is_err());
118+ assert!(Reference::parse("g1t.sh/acme/web:bad tag").is_err());
119+ }
120+}
+610−0
1+//! The registry side of a push: the token, blobs in chunks, the manifest.
2+//!
3+//! Every blob is checked with `HEAD` first and skipped when the registry
4+//! has it. Otherwise it is uploaded in chunks, each its own `PATCH` under
5+//! the 100 MB a request may carry, then closed with `PUT ?digest=`. A
6+//! `429` or `5xx` is retried after a wait (the `Retry-After` when there is
7+//! one), and an interrupted upload continues from the offset the registry
8+//! says it holds; one the registry let go is started again.
9+
10+use crate::chunk;
11+use crate::credentials::Credentials;
12+use crate::image::{Blob, Source};
13+use base64::Engine as _;
14+use serde_json::Value;
15+use std::fs::File;
16+use std::io::{Read, Seek, SeekFrom};
17+use std::path::Path;
18+use std::time::Duration;
19+
20+/// How many times a request is tried before the push gives up.
21+const ATTEMPTS: u32 = 6;
22+/// How many times one blob's upload may start over.
23+const RESTARTS: u32 = 3;
24+
25+pub struct Registry {
26+ agent: ureq::Agent,
27+ base: String,
28+ host: String,
29+ name: String,
30+ credentials: Credentials,
31+ token: Option<String>,
32+}
33+
34+/// A registry's answer, whatever its status.
35+pub struct Reply {
36+ pub status: u16,
37+ location: Option<String>,
38+ range: Option<String>,
39+ digest: Option<String>,
40+ retry_after: Option<u64>,
41+ challenge: Option<String>,
42+ body: String,
43+}
44+
45+pub enum Outcome {
46+ Exists,
47+ Uploaded { chunks: u64 },
48+}
49+
50+enum Failed {
51+ /// The registry no longer has the upload: start the blob again.
52+ Restart(String),
53+ Fatal(String),
54+}
55+
56+type Result<T> = std::result::Result<T, String>;
57+
58+impl Registry {
59+ pub fn new(base: String, host: String, name: String, credentials: Credentials) -> Registry {
60+ let agent = ureq::AgentBuilder::new()
61+ .timeout_connect(Duration::from_secs(30))
62+ .timeout_read(Duration::from_secs(300))
63+ .user_agent(concat!("g1t/", env!("CARGO_PKG_VERSION")))
64+ .try_proxy_from_env(true)
65+ .build();
66+ Registry {
67+ agent,
68+ base,
69+ host,
70+ name,
71+ credentials,
72+ token: None,
73+ }
74+ }
75+
76+ /// Trades the g1t token for a registry token that may pull and push
77+ /// this image, at the token endpoint the registry names.
78+ pub fn sign_in(&mut self) -> Result<()> {
79+ let ping = self.raw("GET", &format!("{}/v2/", self.base), &[], None, false)?;
80+ let (realm, service) = ping
81+ .challenge
82+ .as_deref()
83+ .and_then(parse_challenge)
84+ .unwrap_or_else(|| (format!("{}/v2/token", self.base), self.host.clone()));
85+ let url = format!(
86+ "{realm}?service={service}&scope=repository:{}:pull,push",
87+ self.name
88+ );
89+ let basic = base64::engine::general_purpose::STANDARD.encode(format!(
90+ "{}:{}",
91+ self.credentials.username, self.credentials.secret
92+ ));
93+ let authorization = format!("Basic {basic}");
94+ let mut reply = None;
95+ for attempt in 0..ATTEMPTS {
96+ match self.raw(
97+ "GET",
98+ &url,
99+ &[("authorization", &authorization)],
100+ None,
101+ false,
102+ ) {
103+ Ok(r) if retry_wait(&r, attempt).is_none() => {
104+ reply = Some(r);
105+ break;
106+ }
107+ Ok(r) => wait(retry_wait(&r, attempt)),
108+ Err(error) if attempt + 1 == ATTEMPTS => return Err(error),
109+ Err(_) => wait(Some(backoff(attempt))),
110+ }
111+ }
112+ let reply = reply.ok_or_else(|| format!("{} did not answer the sign-in.", self.host))?;
113+ if reply.status != 200 {
114+ return Err(format!(
115+ "{} refused the token: {}",
116+ self.host,
117+ explain(&reply)
118+ ));
119+ }
120+ let body: Value = serde_json::from_str(&reply.body).map_err(|_| {
121+ format!(
122+ "{} answered the sign-in with something that is not JSON.",
123+ self.host
124+ )
125+ })?;
126+ let token = body
127+ .get("token")
128+ .or_else(|| body.get("access_token"))
129+ .and_then(Value::as_str);
130+ self.token = Some(
131+ token
132+ .ok_or_else(|| format!("{} answered the sign-in without a token.", self.host))?
133+ .to_owned(),
134+ );
135+ Ok(())
136+ }
137+
138+ /// One request, as it went. Transport failures are errors.
139+ fn raw(
140+ &self,
141+ method: &str,
142+ url: &str,
143+ headers: &[(&str, &str)],
144+ body: Option<&[u8]>,
145+ bearer: bool,
146+ ) -> Result<Reply> {
147+ let mut request = self.agent.request(method, url);
148+ for (name, value) in headers {
149+ request = request.set(name, value);
150+ }
151+ if bearer && let Some(token) = &self.token {
152+ request = request.set("authorization", &format!("Bearer {token}"));
153+ }
154+ let result = match body {
155+ Some(bytes) => request.send_bytes(bytes),
156+ None => request.call(),
157+ };
158+ let response = match result {
159+ Ok(response) | Err(ureq::Error::Status(_, response)) => response,
160+ Err(ureq::Error::Transport(error)) => return Err(format!("{method} {url}: {error}")),
161+ };
162+ let header = |name: &str| response.header(name).map(str::to_owned);
163+ let mut reply = Reply {
164+ status: response.status(),
165+ location: header("location"),
166+ range: header("range"),
167+ digest: header("docker-content-digest"),
168+ retry_after: header("retry-after").and_then(|s| s.trim().parse().ok()),
169+ challenge: header("www-authenticate"),
170+ body: String::new(),
171+ };
172+ if method != "HEAD" {
173+ let mut body = String::new();
174+ let _ = response
175+ .into_reader()
176+ .take(1 << 20)
177+ .read_to_string(&mut body);
178+ reply.body = body;
179+ }
180+ Ok(reply)
181+ }
182+
183+ /// One request with the registry token, signing in again once if the
184+ /// token has run out.
185+ fn send(
186+ &mut self,
187+ method: &str,
188+ url: &str,
189+ headers: &[(&str, &str)],
190+ body: Option<&[u8]>,
191+ ) -> Result<Reply> {
192+ let reply = self.raw(method, url, headers, body, true)?;
193+ if reply.status == 401 {
194+ self.sign_in()?;
195+ return self.raw(method, url, headers, body, true);
196+ }
197+ Ok(reply)
198+ }
199+
200+ /// A request that may be sent again as it is, tried until it gets an
201+ /// answer that is not a `429` or `5xx`.
202+ fn call(
203+ &mut self,
204+ method: &str,
205+ url: &str,
206+ headers: &[(&str, &str)],
207+ body: Option<&[u8]>,
208+ ) -> Result<Reply> {
209+ let mut last = String::new();
210+ for attempt in 0..ATTEMPTS {
211+ match self.send(method, url, headers, body) {
212+ Ok(reply) => match retry_wait(&reply, attempt) {
213+ None => return Ok(reply),
214+ Some(delay) => {
215+ last = explain(&reply);
216+ wait(Some(delay));
217+ }
218+ },
219+ Err(error) => {
220+ last = error;
221+ wait(Some(backoff(attempt)));
222+ }
223+ }
224+ }
225+ Err(format!(
226+ "{method} {url} failed {ATTEMPTS} times; the last: {last}"
227+ ))
228+ }
229+
230+ pub fn blob_exists(&mut self, digest: &str) -> Result<bool> {
231+ let url = format!("{}/v2/{}/blobs/{digest}", self.base, self.name);
232+ let reply = self.call("HEAD", &url, &[], None)?;
233+ match reply.status {
234+ 200 => Ok(true),
235+ 404 => Ok(false),
236+ _ => Err(format!("Could not check for {digest}: {}", explain(&reply))),
237+ }
238+ }
239+
240+ /// Puts a blob in the registry unless it is there already.
241+ pub fn push_blob(&mut self, blob: &Blob, tar: &Path, chunk_size: u64) -> Result<Outcome> {
242+ if self.blob_exists(&blob.digest)? {
243+ return Ok(Outcome::Exists);
244+ }
245+ let mut restarts = 0;
246+ loop {
247+ match self.upload(blob, tar, chunk_size) {
248+ Ok(chunks) => return Ok(Outcome::Uploaded { chunks }),
249+ Err(Failed::Restart(reason)) if restarts < RESTARTS => {
250+ restarts += 1;
251+ eprintln!(" {reason}; starting this blob again");
252+ }
253+ Err(Failed::Restart(reason) | Failed::Fatal(reason)) => return Err(reason),
254+ }
255+ }
256+ }
257+
258+ fn resolve(&self, location: Option<&str>, fallback: &str) -> String {
259+ match location {
260+ Some(l) if l.starts_with("http://") || l.starts_with("https://") => l.to_owned(),
261+ Some(l) if l.starts_with('/') => format!("{}{l}", self.base),
262+ Some(l) if !l.is_empty() => format!("{}/{l}", self.base),
263+ _ => fallback.to_owned(),
264+ }
265+ }
266+
267+ /// Where an interrupted upload stands: the offset to continue from,
268+ /// or a restart when the registry let it go.
269+ fn upload_offset(&mut self, location: &mut String) -> std::result::Result<u64, Failed> {
270+ let reply = self
271+ .call("GET", location, &[], None)
272+ .map_err(Failed::Fatal)?;
273+ match reply.status {
274+ 200 | 204 => {
275+ *location = self.resolve(reply.location.as_deref(), location);
276+ reply
277+ .range
278+ .as_deref()
279+ .and_then(chunk::next_offset)
280+ .ok_or_else(|| Failed::Restart("the upload's progress is unknown".to_owned()))
281+ }
282+ 404 => Err(Failed::Restart("the registry let the upload go".to_owned())),
283+ _ => Err(Failed::Fatal(format!(
284+ "Could not check the upload: {}",
285+ explain(&reply)
286+ ))),
287+ }
288+ }
289+
290+ fn upload(
291+ &mut self,
292+ blob: &Blob,
293+ tar: &Path,
294+ chunk_size: u64,
295+ ) -> std::result::Result<u64, Failed> {
296+ let start = format!("{}/v2/{}/blobs/uploads/", self.base, self.name);
297+ let reply = self
298+ .call("POST", &start, &[], Some(&[]))
299+ .map_err(Failed::Fatal)?;
300+ if reply.status != 202 {
301+ return Err(Failed::Fatal(format!(
302+ "Could not start an upload: {}",
303+ explain(&reply)
304+ )));
305+ }
306+ let mut location = self.resolve(reply.location.as_deref(), "");
307+ if location.is_empty() {
308+ return Err(Failed::Fatal(
309+ "The registry started an upload without saying where.".to_owned(),
310+ ));
311+ }
312+ let (mut source, base_offset) = open(blob, tar).map_err(Failed::Fatal)?;
313+ let total_chunks = chunk::count(blob.size, chunk_size);
314+ let mut offset = 0u64;
315+ let mut chunks = 0u64;
316+ let mut failures = 0u32;
317+ let mut buffer = Vec::new();
318+ while let Some(&(_, length)) = chunk::plan(blob.size, chunk_size, offset).first() {
319+ read_at(&mut source, base_offset + offset, length, &mut buffer)
320+ .map_err(Failed::Fatal)?;
321+ let range = chunk::content_range(offset, length);
322+ let headers = [
323+ ("content-type", "application/octet-stream"),
324+ ("content-range", range.as_str()),
325+ ];
326+ let sent = self.send("PATCH", &location.clone(), &headers, Some(&buffer));
327+ let retry = match &sent {
328+ Ok(reply) if reply.status == 202 => None,
329+ Ok(reply) => retry_wait(reply, failures),
330+ Err(_) => Some(backoff(failures)),
331+ };
332+ match sent {
333+ Ok(reply) if reply.status == 202 => {
334+ location = self.resolve(reply.location.as_deref(), &location);
335+ let next = reply
336+ .range
337+ .as_deref()
338+ .and_then(chunk::next_offset)
339+ .unwrap_or(offset + length);
340+ if next > blob.size {
341+ return Err(Failed::Fatal(format!(
342+ "The registry says it holds {next} bytes of a {}-byte blob.",
343+ blob.size
344+ )));
345+ }
346+ offset = if next == 0 { offset + length } else { next };
347+ chunks += 1;
348+ failures = 0;
349+ if total_chunks > 1 {
350+ eprintln!(
351+ " chunk {}/{total_chunks} {} {}",
352+ chunks.min(total_chunks),
353+ chunk::human(length),
354+ percent(offset, blob.size)
355+ );
356+ }
357+ }
358+ Ok(reply) if reply.status == 413 => {
359+ return Err(Failed::Fatal(format!(
360+ "A {} chunk was refused as too large ({}). Try a smaller --chunk-size.",
361+ chunk::human(length),
362+ explain(&reply)
363+ )));
364+ }
365+ Ok(reply) if reply.status == 404 || reply.status == 416 => {
366+ offset = self.upload_offset(&mut location)?;
367+ }
368+ Ok(reply) if retry.is_none() => {
369+ return Err(Failed::Fatal(format!(
370+ "A chunk was refused: {}",
371+ explain(&reply)
372+ )));
373+ }
374+ sent => {
375+ failures += 1;
376+ let why = match sent {
377+ Ok(reply) => explain(&reply),
378+ Err(error) => error,
379+ };
380+ if failures >= ATTEMPTS {
381+ return Err(Failed::Fatal(format!(
382+ "A chunk failed {ATTEMPTS} times; the last: {why}"
383+ )));
384+ }
385+ eprintln!(
386+ " chunk at {} failed ({why}); trying again",
387+ chunk::human(offset)
388+ );
389+ wait(retry);
390+ offset = self.upload_offset(&mut location)?;
391+ }
392+ }
393+ }
394+ // Every byte is in: close the upload with its digest.
395+ let separator = if location.contains('?') { '&' } else { '?' };
396+ let finish = format!(
397+ "{location}{separator}digest={}",
398+ blob.digest.replace(':', "%3A")
399+ );
400+ let mut failures = 0u32;
401+ loop {
402+ let sent = self.send(
403+ "PUT",
404+ &finish,
405+ &[("content-type", "application/octet-stream")],
406+ Some(&[]),
407+ );
408+ let why = match sent {
409+ Ok(reply) if reply.status == 201 => return Ok(chunks.max(1)),
410+ Ok(reply) if reply.status != 404 && retry_wait(&reply, failures).is_none() => {
411+ return Err(Failed::Fatal(format!(
412+ "The registry did not take {}: {}",
413+ blob.digest,
414+ explain(&reply)
415+ )));
416+ }
417+ Ok(reply) => explain(&reply),
418+ Err(error) => error,
419+ };
420+ failures += 1;
421+ if failures >= ATTEMPTS {
422+ return Err(Failed::Fatal(format!(
423+ "Closing the upload failed {ATTEMPTS} times; the last: {why}"
424+ )));
425+ }
426+ wait(Some(backoff(failures)));
427+ // The close may have landed before the answer was lost.
428+ if self.blob_exists(&blob.digest).map_err(Failed::Fatal)? {
429+ return Ok(chunks.max(1));
430+ }
431+ if self.upload_offset(&mut location.clone())? != blob.size {
432+ return Err(Failed::Restart("the upload lost bytes".to_owned()));
433+ }
434+ }
435+ }
436+
437+ /// Puts the manifest under the tag. Returns the registry's digest.
438+ pub fn push_manifest(
439+ &mut self,
440+ tag: &str,
441+ media_type: &str,
442+ manifest: &[u8],
443+ ) -> Result<String> {
444+ let url = format!("{}/v2/{}/manifests/{tag}", self.base, self.name);
445+ let reply = self.call("PUT", &url, &[("content-type", media_type)], Some(manifest))?;
446+ if reply.status != 201 && reply.status != 200 {
447+ return Err(format!(
448+ "The registry did not take the manifest: {}",
449+ explain(&reply)
450+ ));
451+ }
452+ Ok(reply.digest.unwrap_or_default())
453+ }
454+}
455+
456+fn open(blob: &Blob, tar: &Path) -> Result<(File, u64)> {
457+ let (path, offset) = match &blob.source {
458+ Source::Tar { offset } => (tar, *offset),
459+ Source::File(path) => (path.as_path(), 0),
460+ };
461+ let file = File::open(path).map_err(|e| format!("Could not open {}: {e}", path.display()))?;
462+ Ok((file, offset))
463+}
464+
465+fn read_at(file: &mut File, offset: u64, length: u64, buffer: &mut Vec<u8>) -> Result<()> {
466+ buffer.resize(length as usize, 0);
467+ file.seek(SeekFrom::Start(offset))
468+ .map_err(|e| e.to_string())?;
469+ file.read_exact(buffer)
470+ .map_err(|e| format!("Could not read the image: {e}"))
471+}
472+
473+fn percent(done: u64, total: u64) -> String {
474+ format!(
475+ "{:.0}%",
476+ if total == 0 {
477+ 100.0
478+ } else {
479+ done as f64 * 100.0 / total as f64
480+ }
481+ )
482+}
483+
484+/// `Bearer realm="https://g1t.sh/v2/token",service="g1t.sh",scope="…"`.
485+fn parse_challenge(header: &str) -> Option<(String, String)> {
486+ let rest = header
487+ .trim()
488+ .strip_prefix("Bearer ")
489+ .or_else(|| header.trim().strip_prefix("bearer "))?;
490+ let mut realm = None;
491+ let mut service = None;
492+ for part in rest.split(',') {
493+ let (key, value) = part.trim().split_once('=')?;
494+ let value = value.trim().trim_matches('"').to_owned();
495+ match key.trim() {
496+ "realm" => realm = Some(value),
497+ "service" => service = Some(value),
498+ _ => {}
499+ }
500+ }
501+ Some((realm?, service.unwrap_or_default()))
502+}
503+
504+/// How long to wait before trying again, when the answer is worth
505+/// another try: `429` and `5xx`.
506+fn retry_wait(reply: &Reply, attempt: u32) -> Option<Duration> {
507+ if reply.status != 429 && reply.status < 500 {
508+ return None;
509+ }
510+ Some(match reply.retry_after {
511+ Some(seconds) => Duration::from_secs(seconds.min(300)),
512+ None => backoff(attempt),
513+ })
514+}
515+
516+/// 1s, 2s, 4s … at most 30s.
517+pub fn backoff(attempt: u32) -> Duration {
518+ Duration::from_secs((1u64 << attempt.min(5)).min(30))
519+}
520+
521+fn wait(delay: Option<Duration>) {
522+ if let Some(delay) = delay {
523+ std::thread::sleep(delay);
524+ }
525+}
526+
527+/// An error reply for people: the registry's own code and message when it
528+/// sent them, the status otherwise.
529+fn explain(reply: &Reply) -> String {
530+ if let Ok(body) = serde_json::from_str::<Value>(&reply.body)
531+ && let Some(error) = body
532+ .get("errors")
533+ .and_then(Value::as_array)
534+ .and_then(|e| e.first())
535+ {
536+ let code = error.get("code").and_then(Value::as_str).unwrap_or("ERROR");
537+ let message = error.get("message").and_then(Value::as_str).unwrap_or("");
538+ return format!("{code}: {message} ({})", reply.status);
539+ }
540+ let text = reply.body.trim();
541+ if text.is_empty() || text.starts_with('<') {
542+ format!("HTTP {}", reply.status)
543+ } else {
544+ format!(
545+ "HTTP {}: {}",
546+ reply.status,
547+ text.chars().take(200).collect::<String>()
548+ )
549+ }
550+}
551+
552+#[cfg(test)]
553+mod tests {
554+ use super::*;
555+
556+ fn reply(status: u16, retry_after: Option<u64>, body: &str) -> Reply {
557+ Reply {
558+ status,
559+ location: None,
560+ range: None,
561+ digest: None,
562+ retry_after,
563+ challenge: None,
564+ body: body.to_owned(),
565+ }
566+ }
567+
568+ #[test]
569+ fn challenges_name_the_token_endpoint() {
570+ let header = r#"Bearer realm="https://g1t.sh/v2/token",service="g1t.sh",scope="repository:acme/web:pull""#;
571+ assert_eq!(
572+ parse_challenge(header),
573+ Some(("https://g1t.sh/v2/token".into(), "g1t.sh".into()))
574+ );
575+ assert_eq!(parse_challenge("Basic realm=\"x\""), None);
576+ }
577+
578+ #[test]
579+ fn only_429_and_5xx_are_tried_again() {
580+ assert_eq!(
581+ retry_wait(&reply(429, Some(7), ""), 0),
582+ Some(Duration::from_secs(7))
583+ );
584+ assert_eq!(
585+ retry_wait(&reply(503, None, ""), 2),
586+ Some(Duration::from_secs(4))
587+ );
588+ assert_eq!(retry_wait(&reply(400, None, ""), 0), None);
589+ assert_eq!(retry_wait(&reply(413, None, ""), 0), None);
590+ }
591+
592+ #[test]
593+ fn backoff_doubles_up_to_30s() {
594+ let waits: Vec<u64> = (0..8).map(|n| backoff(n).as_secs()).collect();
595+ assert_eq!(waits, [1, 2, 4, 8, 16, 30, 30, 30]);
596+ }
597+
598+ #[test]
599+ fn errors_read_as_the_registry_wrote_them() {
600+ let body = r#"{"errors":[{"code":"DENIED","message":"Pushing needs Write."}]}"#;
601+ assert_eq!(
602+ explain(&reply(403, None, body)),
603+ "DENIED: Pushing needs Write. (403)"
604+ );
605+ assert_eq!(
606+ explain(&reply(413, None, "<html>Too large</html>")),
607+ "HTTP 413"
608+ );
609+ }
610+}
+4−2
7777 Uploads are written to R2 as multipart parts, so a layer may be any size if it arrives in
7878 chunks under the limit. `docker push` sends a layer in one request, so on Cloudflare a layer
7979 over the limit is refused with a message naming the limit and the way around it: `g1t push`
80− (the CLI), which asks for signed part URLs and uploads straight to R2 at any size, and is what
81− g1t Actions and the runner use. Self-hosted, there is no such limit.
80+ (the CLI, `crates/g1t`), which reads the image with `docker save` (OCI layout or classic,
81+ gzipping uncompressed layers and writing a matching manifest) and sends every blob as OCI
82+ chunks of 90 MiB (at most 95 MiB), resuming from the upload's `Range` after a `429`/`5xx`.
83+ Self-hosted, there is no such limit.
8284
8385 ### npm
8486