Skip to main content

isb_core/registry/
mod.rs

1//! The local OCI registry: where builds push and incus pulls.
2//!
3//! One per host, run by isb as an OCI container (`registry:2`, pinned by
4//! digest) in the system project [`PROJECT`], never in an org:
5//!
6//! - It has no network interface. A proxy device listening on the host's
7//!   `127.0.0.1:<port>` is its only way in, so incusd (and the daemon) reach
8//!   it and no org network can, by construction.
9//! - It speaks TLS with a certificate from an isb CA kept in the daemon's
10//!   state directory ([`tls`]); `isb host setup` installs that CA where the
11//!   skopeo inside incusd looks for it, since incus pulls OCI images only
12//!   over https.
13//! - Only the daemon pushes. A build exports an OCI layout inside its
14//!   sandbox; the daemon copies it out and pushes it ([`oci`]). Org
15//!   sandboxes never get credentials for, or a route to, the registry.
16//! - Repositories are `<org>/<app>`. A compose file names an image as
17//!   `registry:<app>:<tag>` (or `@sha256:...`), always resolved in the org
18//!   the stack or sandbox lives in, so one org cannot name another's images.
19
20pub mod oci;
21pub mod tls;
22
23use std::collections::{BTreeMap, BTreeSet};
24use std::path::{Path, PathBuf};
25use std::sync::{Arc, Mutex, OnceLock};
26use std::time::{Duration, Instant};
27
28use serde::{Deserialize, Serialize};
29use serde_json::{Value, json};
30
31use crate::client::{Client, encode_segment};
32use crate::error::{Error, Result};
33use crate::org::OrgId;
34
35/// The incus project for isb's own services. Not an org: `system` is a
36/// reserved org name.
37pub const PROJECT: &str = "isb-system";
38pub const INSTANCE: &str = "registry";
39pub const VOLUME: &str = "registry-data";
40/// `registry:2.8.3`, by digest, so what runs is what was reviewed.
41pub const IMAGE: &str =
42    "docker:registry@sha256:a3d8aaa63ed8681a604f1dea0aa03f100d5895b6a58ace528858a7b332415373";
43pub const DEFAULT_PORT: u16 = 5480;
44/// Tags kept per repository by [`Registry::gc`], besides deployed ones.
45pub const DEFAULT_KEEP: usize = 10;
46
47/// Where the registry listens inside its container: a unix socket, since a
48/// container without a NIC has no loopback up either.
49const INNER_SOCKET: &str = "/tmp/registry.sock";
50const KEY_ADDR: &str = "user.isb.registry.addr";
51const KEY_CA: &str = "user.isb.registry.ca";
52const CERT_DIR: &str = "/certs";
53
54/// What a compose file's `registry:` image names: an app's repository in
55/// the org, and a tag and/or a digest.
56#[derive(Debug, Clone, PartialEq, Eq)]
57pub struct ImageRef {
58    pub app: String,
59    pub tag: Option<String>,
60    pub digest: Option<String>,
61}
62
63/// An app name as a repository component: `[a-z0-9][a-z0-9._-]*`, at most
64/// 128 characters, no `/` (the org is the only namespace).
65pub fn valid_app(a: &str) -> bool {
66    !a.is_empty()
67        && a.len() <= 128
68        && a.starts_with(|c: char| c.is_ascii_lowercase() || c.is_ascii_digit())
69        && a.chars()
70            .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || "._-".contains(c))
71}
72
73/// An OCI tag: `[A-Za-z0-9_][A-Za-z0-9_.-]{0,127}`.
74pub fn valid_tag(t: &str) -> bool {
75    !t.is_empty()
76        && t.len() <= 128
77        && !t.starts_with(['.', '-'])
78        && t.chars()
79            .all(|c| c.is_ascii_alphanumeric() || "_.-".contains(c))
80}
81
82impl ImageRef {
83    /// `APP[:TAG][@sha256:...]`, what follows `registry:`.
84    pub fn parse(s: &str) -> Result<ImageRef> {
85        let bad = |why: &str| {
86            Error::invalid(format!(
87                "registry image {s:?}: {why} (want registry:APP[:TAG][@sha256:DIGEST], the app's image in this org)"
88            ))
89        };
90        let (rest, digest) = match s.split_once('@') {
91            Some((r, d)) => {
92                if !oci::valid_digest(d) {
93                    return Err(bad("the digest must be sha256: and 64 hex digits"));
94                }
95                (r, Some(d.to_string()))
96            }
97            None => (s, None),
98        };
99        let (app, tag) = match rest.split_once(':') {
100            Some((a, t)) => (a, Some(t.to_string())),
101            None => (rest, None),
102        };
103        if app.contains('/') {
104            return Err(bad(
105                "images are named by app alone; the org is implied and other orgs' images cannot be named",
106            ));
107        }
108        if !valid_app(app) {
109            return Err(bad("the app is [a-z0-9][a-z0-9._-]*"));
110        }
111        if let Some(t) = &tag {
112            if !valid_tag(t) {
113                return Err(bad("bad tag"));
114            }
115        }
116        Ok(ImageRef {
117            app: app.to_string(),
118            tag,
119            digest,
120        })
121    }
122
123    /// The tag, `latest` by default.
124    pub fn tag_or_latest(&self) -> &str {
125        self.tag.as_deref().unwrap_or("latest")
126    }
127
128    /// `APP:TAG@DIGEST`, as written after `registry:`.
129    pub fn render(&self) -> String {
130        let mut s = self.app.clone();
131        if let Some(t) = &self.tag {
132            s.push(':');
133            s.push_str(t);
134        }
135        if let Some(d) = &self.digest {
136            s.push('@');
137            s.push_str(d);
138        }
139        s
140    }
141
142    /// This reference pinned to `digest` (the tag is kept for people).
143    pub fn pinned(&self, digest: &str) -> ImageRef {
144        ImageRef {
145            digest: Some(digest.to_string()),
146            ..self.clone()
147        }
148    }
149
150    /// The incus `alias` to pull: `<org>/<app>@digest`, or `:tag`. Skopeo
151    /// refuses a tag and a digest together, and the digest is what counts.
152    pub fn pull_alias(&self, org: &OrgId) -> String {
153        match &self.digest {
154            Some(d) => format!("{}@{d}", repo(org, &self.app)),
155            None => format!("{}:{}", repo(org, &self.app), self.tag_or_latest()),
156        }
157    }
158}
159
160/// An org's repository for an app.
161pub fn repo(org: &OrgId, app: &str) -> String {
162    format!("{org}/{app}")
163}
164
165/// Where the registry is, as recorded on the system project (so any isb
166/// process with incus access finds it, the daemon's state dir or not).
167#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
168pub struct Info {
169    /// `127.0.0.1:5480`
170    pub addr: String,
171    pub ca_pem: String,
172}
173
174impl Info {
175    pub fn url(&self) -> String {
176        format!("https://{}", self.addr)
177    }
178}
179
180fn host(base: &Client) -> Client {
181    base.clone().project("default")
182}
183
184fn sys(base: &Client) -> Client {
185    base.clone().project(PROJECT)
186}
187
188/// The registry's address and CA, if it has been set up.
189pub fn info(base: &Client) -> Result<Option<Info>> {
190    let Some(p) = host(base).get_opt(&format!("/1.0/projects/{PROJECT}"))? else {
191        return Ok(None);
192    };
193    let addr = p["config"][KEY_ADDR].as_str().unwrap_or_default();
194    let ca = p["config"][KEY_CA].as_str().unwrap_or_default();
195    if addr.is_empty() || ca.is_empty() {
196        return Ok(None);
197    }
198    Ok(Some(Info {
199        addr: addr.to_string(),
200        ca_pem: ca.to_string(),
201    }))
202}
203
204/// Where skopeo (inside incusd) looks for the CA of `addr`.
205pub fn host_ca_path(addr: &str) -> PathBuf {
206    Path::new("/etc/containers/certs.d")
207        .join(addr)
208        .join("ca.crt")
209}
210
211/// Create the registry, or bring it in line: the system project, its data
212/// volume, the TLS material (in `<state>/registry`), the container and its
213/// loopback proxy. Safe to run again; `renew` reissues the certificate.
214#[expect(
215    clippy::too_many_lines,
216    reason = "predates the lint ratchet; split it when next changed"
217)]
218pub fn setup(
219    base: &Client,
220    state: &Path,
221    port: u16,
222    renew: bool,
223    report: &mut dyn FnMut(&str),
224) -> Result<Info> {
225    let ip: std::net::IpAddr = "127.0.0.1".parse().expect("constant");
226    let addr = format!("{ip}:{port}");
227    let mat = tls::ensure(&state.join("registry"), ip, renew)?;
228    let h = host(base);
229    let other = h.get_timeouts().other;
230
231    let proj_path = format!("/1.0/projects/{PROJECT}");
232    let config = json!({
233        "features.images": "false",
234        "features.profiles": "true",
235        "features.storage.volumes": "true",
236        "features.networks": "false",
237        KEY_ADDR: addr,
238        KEY_CA: mat.ca_cert,
239    });
240    match h.get_opt(&proj_path)? {
241        None => {
242            report(&format!("creating project {PROJECT}"));
243            h.mutate(
244                "POST",
245                "/1.0/projects",
246                Some(&json!({"name": PROJECT, "description": "isb system services (not an org)", "config": config})),
247                &format!("create project {PROJECT}"),
248                other,
249            )?;
250        }
251        Some(p) => {
252            let mut merged = p["config"].clone();
253            for (k, v) in config.as_object().expect("object") {
254                merged[k] = v.clone();
255            }
256            h.mutate(
257                "PUT",
258                &proj_path,
259                Some(&json!({"description": p["description"], "config": merged})),
260                &format!("update project {PROJECT}"),
261                other,
262            )?;
263        }
264    }
265    let s = sys(base);
266    let facts = crate::sandbox::host_facts(&h)?;
267    let pool = facts.pick_pool(None)?;
268    // A root disk and nothing else: no NIC.
269    s.mutate(
270        "PUT",
271        "/1.0/profiles/default",
272        Some(&json!({
273            "description": "isb system services: no network",
274            "config": {},
275            "devices": {"root": {"type": "disk", "path": "/", "pool": pool}},
276        })),
277        "set the system project's default profile",
278        other,
279    )?;
280    let vol_path = format!(
281        "/1.0/storage-pools/{}/volumes/custom/{VOLUME}",
282        encode_segment(&pool)
283    );
284    if s.get_opt(&vol_path)?.is_none() {
285        report(&format!("creating volume {VOLUME} on {pool}"));
286        s.mutate(
287            "POST",
288            &format!("/1.0/storage-pools/{}/volumes/custom", encode_segment(&pool)),
289            Some(&json!({"name": VOLUME, "type": "custom", "content_type": "filesystem", "config": {}})),
290            &format!("create volume {VOLUME}"),
291            other,
292        )?;
293    }
294
295    let env = [
296        ("REGISTRY_HTTP_NET", "unix".into()),
297        ("REGISTRY_HTTP_ADDR", INNER_SOCKET.to_string()),
298        (
299            "REGISTRY_HTTP_TLS_CERTIFICATE",
300            format!("{CERT_DIR}/tls.crt"),
301        ),
302        ("REGISTRY_HTTP_TLS_KEY", format!("{CERT_DIR}/tls.key")),
303        ("REGISTRY_STORAGE_DELETE_ENABLED", "true".into()),
304        (
305            "REGISTRY_STORAGE_FILESYSTEM_ROOTDIRECTORY",
306            "/var/lib/registry".into(),
307        ),
308        ("REGISTRY_LOG_LEVEL", "warn".into()),
309    ];
310    let mut cfg = json!({
311        "boot.autostart": "true",
312        "limits.cpu": "2",
313        "limits.memory": "1GiB",
314        "user.isb.role": "registry",
315    });
316    for (k, v) in &env {
317        cfg[format!("environment.{k}")] = json!(v);
318    }
319    let devices = json!({
320        "data": {"type": "disk", "pool": pool, "source": VOLUME, "path": "/var/lib/registry"},
321        "https": {
322            "type": "proxy",
323            "listen": format!("tcp:{addr}"),
324            "connect": format!("unix:{INNER_SOCKET}"),
325            "bind": "host",
326        },
327    });
328    let inst_path = format!("/1.0/instances/{INSTANCE}");
329    let mut restart = false;
330    match s.get_opt(&inst_path)? {
331        None => {
332            report(&format!("creating {INSTANCE} from {IMAGE}"));
333            let src = crate::plan::ImageSource::parse(IMAGE)?;
334            s.mutate(
335                "POST",
336                "/1.0/instances",
337                Some(&json!({
338                    "name": INSTANCE,
339                    "type": "container",
340                    "source": src.to_api(None),
341                    "config": cfg,
342                    "devices": devices,
343                    "profiles": ["default"],
344                })),
345                &format!("create {INSTANCE}"),
346                s.get_timeouts().create,
347            )?;
348        }
349        Some(i) => {
350            let mut c = i["config"].clone();
351            let mut d = i["devices"].clone();
352            let mut changed = false;
353            for (k, v) in cfg.as_object().expect("object") {
354                if c[k] != *v {
355                    c[k] = v.clone();
356                    changed = true;
357                }
358            }
359            for (k, v) in devices.as_object().expect("object") {
360                if d[k] != *v {
361                    d[k] = v.clone();
362                    changed = true;
363                }
364            }
365            if changed {
366                report(&format!("updating {INSTANCE}"));
367                s.mutate(
368                    "PATCH",
369                    &inst_path,
370                    Some(&json!({"config": c, "devices": d})),
371                    &format!("update {INSTANCE}"),
372                    other,
373                )?;
374                restart = true;
375            }
376        }
377    }
378    s.make_dir(INSTANCE, CERT_DIR, 0, 0, 0o700)?;
379    for (name, text, mode) in [("tls.crt", &mat.cert, 0o644), ("tls.key", &mat.key, 0o600)] {
380        let path = format!("{CERT_DIR}/{name}");
381        if s.read_file(INSTANCE, &path)?.as_deref() != Some(text.as_bytes()) {
382            s.push_file(INSTANCE, &path, text.as_bytes(), 0, 0, mode)?;
383            restart = true;
384        }
385    }
386    let state_now = s.get(&format!("{inst_path}/state"))?;
387    let running = state_now["status"].as_str() == Some("Running");
388    let action = match (running, restart) {
389        (false, _) => Some("start"),
390        (true, true) => Some("restart"),
391        (true, false) => None,
392    };
393    if let Some(a) = action {
394        report(&format!("{a}ing {INSTANCE}"));
395        s.mutate(
396            "PUT",
397            &format!("{inst_path}/state"),
398            Some(&json!({"action": a, "timeout": 30, "force": true})),
399            &format!("{a} {INSTANCE}"),
400            other,
401        )?;
402    }
403    let info = Info {
404        addr,
405        ca_pem: mat.ca_cert,
406    };
407    wait_up(&info, Duration::from_secs(60))?;
408    Ok(info)
409}
410
411/// Wait until the registry answers `/v2/`.
412fn wait_up(info: &Info, deadline: Duration) -> Result<()> {
413    let agent: ureq::Agent = ureq::Agent::config_builder()
414        .timeout_global(Some(Duration::from_secs(5)))
415        .http_status_as_error(false)
416        .tls_config(
417            ureq::tls::TlsConfig::builder()
418                .root_certs(ureq::tls::RootCerts::new_with_certs(&[
419                    ureq::tls::Certificate::from_pem(info.ca_pem.as_bytes())
420                        .map_err(|e| Error::invalid(format!("registry CA: {e}")))?,
421                ]))
422                .build(),
423        )
424        .build()
425        .into();
426    let started = Instant::now();
427    let mut last = String::new();
428    while started.elapsed() < deadline {
429        match agent.get(format!("{}/v2/", info.url())).call() {
430            Ok(r) if r.status().as_u16() == 200 => return Ok(()),
431            Ok(r) => last = format!("HTTP {}", r.status()),
432            Err(e) => last = e.to_string(),
433        }
434        std::thread::sleep(Duration::from_millis(500));
435    }
436    Err(Error::invalid(format!(
437        "registry {} not answering after {deadline:?}: {last}",
438        info.url()
439    )))
440}
441
442/// When an app's tag was pushed, as the daemon recorded it: retention keeps
443/// the newest.
444#[derive(Debug, Clone, Default, Serialize, Deserialize)]
445struct PushIndex {
446    /// repo -> tag -> (digest, unix seconds)
447    #[serde(default)]
448    repos: BTreeMap<String, BTreeMap<String, (String, u64)>>,
449}
450
451/// One tag of a repository.
452#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
453pub struct TagInfo {
454    pub tag: String,
455    pub digest: String,
456    /// Unix seconds; 0 when the push was not recorded.
457    pub pushed_at: u64,
458}
459
460#[derive(Debug, Clone, Serialize)]
461pub struct RepoInfo {
462    pub repo: String,
463    pub org: String,
464    pub app: String,
465    pub tags: Vec<TagInfo>,
466}
467
468/// What a garbage collection did (or, dry, would do).
469#[derive(Debug, Clone, Default, Serialize)]
470pub struct GcReport {
471    /// `repo@digest` with the tags that pointed at it.
472    pub deleted: Vec<String>,
473    pub kept: usize,
474    pub dry_run: bool,
475    /// The registry's own garbage-collect output, last lines.
476    pub collect: String,
477}
478
479/// Which manifests retention deletes from one repository: everything but
480/// the newest `keep` tags and the `protected` digests (deployed ones). A
481/// digest is deleted only when no kept tag points at it.
482pub fn select_deletions(
483    tags: &[TagInfo],
484    keep: usize,
485    protected: &BTreeSet<String>,
486) -> Vec<String> {
487    // Retention's own tags ([`KEEP_TAG`]) and previews' tags
488    // ([`is_preview_tag`]) never count as one of the newest: a preview's
489    // image is kept while it is deployed (protected), and no longer.
490    let aside = |t: &TagInfo| t.tag.starts_with(KEEP_TAG) || is_preview_tag(&t.tag);
491    let mut sorted: Vec<&TagInfo> = tags.iter().collect();
492    sorted.sort_by(|a, b| {
493        aside(a)
494            .cmp(&aside(b))
495            .then_with(|| b.pushed_at.cmp(&a.pushed_at))
496            .then_with(|| b.tag.cmp(&a.tag))
497    });
498    let newest = sorted.iter().filter(|t| !aside(t)).count().min(keep);
499    let mut kept: BTreeSet<&str> = protected.iter().map(String::as_str).collect();
500    for t in sorted.iter().take(newest) {
501        kept.insert(&t.digest);
502    }
503    let mut seen = BTreeSet::new();
504    sorted
505        .iter()
506        .skip(newest)
507        .filter(|t| !kept.contains(t.digest.as_str()))
508        .map(|t| t.digest.clone())
509        .filter(|d| seen.insert(d.clone()))
510        .collect()
511}
512
513/// Tags retention puts on deployed digests (`isb-keep-<12 hex>`), so the
514/// registry's untagged-manifest collection spares them.
515pub const KEEP_TAG: &str = "isb-keep-";
516
517/// A preview's image tag: `pr-<number>-<sha>`.
518pub fn is_preview_tag(t: &str) -> bool {
519    preview_tag_number(t).is_some()
520}
521
522/// The pull request number of a preview tag (`pr-12-abc` -> 12).
523pub fn preview_tag_number(t: &str) -> Option<u64> {
524    let rest = t.strip_prefix("pr-")?;
525    let (n, sha) = rest.split_once('-')?;
526    if sha.is_empty() || n.is_empty() || !n.bytes().all(|b| b.is_ascii_digit()) {
527        return None;
528    }
529    n.parse().ok()
530}
531
532/// The daemon's handle on the registry.
533pub struct Registry {
534    base: Client,
535    info: Info,
536    remote: oci::Remote,
537    /// `<state>/registry`, for the push index; `None` read-only.
538    dir: Option<PathBuf>,
539    /// Pushes and the index against garbage collection.
540    lock: Mutex<()>,
541}
542
543static SHARED: OnceLock<Arc<Registry>> = OnceLock::new();
544
545/// Make `r` the registry [`Registry::shared`] returns (the daemon does this
546/// at start, with its own state directory).
547pub fn install(r: Arc<Registry>) {
548    let _ = SHARED.set(r);
549}
550
551impl Registry {
552    /// The registry as set up on this host, or `None` when it is not.
553    /// `state` is the daemon's state directory, for the push index.
554    pub fn open(base: &Client, state: Option<&Path>) -> Result<Option<Registry>> {
555        let Some(info) = info(base)? else {
556            return Ok(None);
557        };
558        let remote = oci::Remote::new(&info.url(), Some(&info.ca_pem), Duration::from_secs(600))?;
559        Ok(Some(Registry {
560            base: base.clone().project("default"),
561            remote,
562            info,
563            dir: state.map(|s| s.join("registry")),
564            lock: Mutex::new(()),
565        }))
566    }
567
568    /// The installed registry, else the host's with the default state dir.
569    pub fn shared(base: &Client) -> Result<Arc<Registry>> {
570        if let Some(r) = SHARED.get() {
571            return Ok(r.clone());
572        }
573        let state = crate::stack::Store::default_dir();
574        let state = std::env::var_os("ISB_SERVE_STATE_DIR")
575            .map(PathBuf::from)
576            .unwrap_or(state);
577        Registry::open(base, Some(&state))?
578            .map(Arc::new)
579            .ok_or_else(not_set_up)
580    }
581
582    pub fn info(&self) -> &Info {
583        &self.info
584    }
585
586    pub fn remote(&self) -> &oci::Remote {
587        &self.remote
588    }
589
590    fn index_path(&self) -> Option<PathBuf> {
591        self.dir.as_ref().map(|d| d.join("pushes.json"))
592    }
593
594    fn load_index(&self) -> PushIndex {
595        self.index_path()
596            .and_then(|p| std::fs::read(p).ok())
597            .and_then(|b| serde_json::from_slice(&b).ok())
598            .unwrap_or_default()
599    }
600
601    fn save_index(&self, idx: &PushIndex) -> Result<()> {
602        let Some(p) = self.index_path() else {
603            return Ok(());
604        };
605        if let Some(d) = p.parent() {
606            std::fs::create_dir_all(d)?;
607        }
608        let tmp = p.with_extension("tmp");
609        std::fs::write(&tmp, serde_json::to_vec_pretty(idx)?)?;
610        std::fs::rename(tmp, p)?;
611        Ok(())
612    }
613
614    /// Push an OCI layout tar as `<org>/<app>:<tag>`. Returns the digest.
615    pub fn push(
616        &self,
617        org: &OrgId,
618        app: &str,
619        tag: &str,
620        tar: &Path,
621        log: &mut dyn FnMut(&str),
622    ) -> Result<String> {
623        if !valid_app(app) {
624            return Err(Error::invalid(format!(
625                "app name {app:?}: [a-z0-9][a-z0-9._-]*"
626            )));
627        }
628        if !valid_tag(tag) {
629            return Err(Error::invalid(format!(
630                "tag {tag:?}: [A-Za-z0-9_][A-Za-z0-9_.-]*"
631            )));
632        }
633        let layout = oci::Layout::open(tar)?;
634        let _g = self.lock.lock().unwrap();
635        let r = repo(org, app);
636        let digest = self.remote.push(&layout, &r, tag, log)?;
637        let mut idx = self.load_index();
638        idx.repos
639            .entry(r)
640            .or_default()
641            .insert(tag.to_string(), (digest.clone(), crate::stack::now_secs()));
642        self.save_index(&idx)?;
643        Ok(digest)
644    }
645
646    /// The digest an org's `registry:` image names now.
647    pub fn resolve(&self, org: &OrgId, r: &ImageRef) -> Result<String> {
648        let name = repo(org, &r.app);
649        let (reference, sep) = match &r.digest {
650            Some(d) => (d.as_str(), '@'),
651            None => (r.tag_or_latest(), ':'),
652        };
653        self.remote.resolve(&name, reference)?.ok_or_else(|| {
654            Error::NotFound(format!(
655                "image registry:{} in org {org} (no {name}{sep}{reference} in the local registry)",
656                r.render()
657            ))
658        })
659    }
660
661    /// Repositories and their tags, of one org or all.
662    pub fn list(&self, org: Option<&OrgId>) -> Result<Vec<RepoInfo>> {
663        let idx = self.load_index();
664        let mut out = Vec::new();
665        for name in self.remote.catalog()? {
666            let Some((o, app)) = name.split_once('/') else {
667                continue;
668            };
669            if org.is_some_and(|x| x.as_str() != o) {
670                continue;
671            }
672            let mut tags = Vec::new();
673            for t in self.remote.tags(&name)? {
674                let Some(d) = self.remote.resolve(&name, &t)? else {
675                    continue;
676                };
677                let pushed_at = idx
678                    .repos
679                    .get(&name)
680                    .and_then(|m| m.get(&t))
681                    .filter(|(pd, _)| *pd == d)
682                    .map(|(_, at)| *at)
683                    .unwrap_or(0);
684                tags.push(TagInfo {
685                    tag: t,
686                    digest: d,
687                    pushed_at,
688                });
689            }
690            if tags.is_empty() {
691                continue;
692            }
693            tags.sort_by(|a, b| {
694                b.pushed_at
695                    .cmp(&a.pushed_at)
696                    .then_with(|| a.tag.cmp(&b.tag))
697            });
698            out.push(RepoInfo {
699                org: o.to_string(),
700                app: app.to_string(),
701                repo: name.clone(),
702                tags,
703            });
704        }
705        Ok(out)
706    }
707
708    /// Delete the manifests of `<org>/<app>` whose every tag `doomed`
709    /// picks, except `spare` digests (what is deployed). A digest that
710    /// another tag still names stays, and so does that tag: the registry
711    /// deletes manifests, not tags. Returns `digest (tags)` of each one
712    /// deleted. Blobs go at the next `registry gc`.
713    pub fn delete_tags(
714        &self,
715        org: &OrgId,
716        app: &str,
717        doomed: &dyn Fn(&str) -> bool,
718        spare: &BTreeSet<String>,
719    ) -> Result<Vec<String>> {
720        let _g = self.lock.lock().unwrap();
721        let r = repo(org, app);
722        let mut by_digest: BTreeMap<String, Vec<String>> = BTreeMap::new();
723        for t in self.remote.tags(&r)? {
724            if let Some(d) = self.remote.resolve(&r, &t)? {
725                by_digest.entry(d).or_default().push(t);
726            }
727        }
728        let mut out = Vec::new();
729        for (d, tags) in &by_digest {
730            if spare.contains(d) || !tags.iter().all(|t| doomed(t)) {
731                continue;
732            }
733            self.remote.delete_manifest(&r, d)?;
734            out.push(format!("{d} ({})", tags.join(", ")));
735        }
736        if !out.is_empty() {
737            let mut idx = self.load_index();
738            if let Some(m) = idx.repos.get_mut(&r) {
739                m.retain(|_, (d, _)| !out.iter().any(|x| x.starts_with(d.as_str())));
740            }
741            self.save_index(&idx)?;
742        }
743        Ok(out)
744    }
745
746    /// Retention: per repository keep the newest `keep` tags and every
747    /// `protected` `repo@digest` (what deployed stacks, current and
748    /// previous, run), delete the other manifests, then have the registry
749    /// collect unreferenced blobs. Pushes wait meanwhile.
750    pub fn gc(
751        &self,
752        keep: usize,
753        protected: &BTreeSet<String>,
754        dry_run: bool,
755        log: &mut dyn FnMut(&str),
756    ) -> Result<GcReport> {
757        let _g = self.lock.lock().unwrap();
758        let mut report = GcReport {
759            dry_run,
760            ..Default::default()
761        };
762        // A deployed digest may have lost its tag (the tag moved on): give
763        // it one of retention's own, so the untagged-manifest collection
764        // below spares it.
765        if !dry_run {
766            for p in protected {
767                let Some((repo, digest)) = p.split_once('@') else {
768                    continue;
769                };
770                let tag = format!("{KEEP_TAG}{}", &digest[7..19.min(digest.len())]);
771                if self.remote.resolve(repo, &tag)?.as_deref() == Some(digest) {
772                    continue;
773                }
774                if let Some((mt, bytes)) = self.remote.get_manifest(repo, digest)? {
775                    self.remote.put_manifest(repo, &tag, &mt, &bytes)?;
776                }
777            }
778        }
779        for r in self.list(None)? {
780            let prot: BTreeSet<String> = protected
781                .iter()
782                .filter_map(|p| p.strip_prefix(&format!("{}@", r.repo)).map(String::from))
783                .collect();
784            let del = select_deletions(&r.tags, keep, &prot);
785            report.kept += r.tags.iter().filter(|t| !del.contains(&t.digest)).count();
786            for d in del {
787                let tags: Vec<&str> = r
788                    .tags
789                    .iter()
790                    .filter(|t| t.digest == d)
791                    .map(|t| t.tag.as_str())
792                    .collect();
793                let what = format!("{}@{d} ({})", r.repo, tags.join(", "));
794                log(&format!(
795                    "{}{what}",
796                    if dry_run {
797                        "would delete "
798                    } else {
799                        "deleting "
800                    }
801                ));
802                if !dry_run {
803                    self.remote.delete_manifest(&r.repo, &d)?;
804                }
805                report.deleted.push(what);
806            }
807        }
808        if dry_run {
809            return Ok(report);
810        }
811        let mut idx = self.load_index();
812        for (repo, tags) in idx.repos.iter_mut() {
813            tags.retain(|_, (d, _)| {
814                !report
815                    .deleted
816                    .iter()
817                    .any(|x| x.starts_with(&format!("{repo}@{d}")))
818            });
819        }
820        self.save_index(&idx)?;
821        // Untagged manifests (left by moved tags) and the blobs nothing
822        // references any more. The registry is idle (pushes hold the lock),
823        // which garbage-collect requires.
824        log("collecting untagged manifests and unreferenced blobs");
825        let s = sys(&self.base);
826        let sb = crate::sandbox::Sandbox::get(&s, INSTANCE)?;
827        let out = sb
828            .exec_stream(
829                [
830                    "/bin/registry",
831                    "garbage-collect",
832                    "--delete-untagged",
833                    "/etc/docker/registry/config.yml",
834                ],
835                crate::exec::ExecOptions::default()
836                    .env(
837                        "REGISTRY_STORAGE_FILESYSTEM_ROOTDIRECTORY",
838                        "/var/lib/registry",
839                    )
840                    .env("REGISTRY_STORAGE_DELETE_ENABLED", "true")
841                    .timeout(Duration::from_secs(1800)),
842            )?
843            .collect_output()?;
844        let text = format!("{}{}", out.stdout_text(), out.stderr_text());
845        let tail: Vec<&str> = text.lines().rev().take(20).collect();
846        report.collect = tail.into_iter().rev().collect::<Vec<_>>().join("\n");
847        if !out.success() {
848            return Err(Error::invalid(format!(
849                "registry garbage-collect failed ({}): {}",
850                out.exit_code, report.collect
851            )));
852        }
853        // The registry caches blob descriptors in memory: restart it so a
854        // collected blob is not reported as present to the next push.
855        s.mutate(
856            "PUT",
857            &format!("/1.0/instances/{INSTANCE}/state"),
858            Some(&json!({"action": "restart", "timeout": 30, "force": true})),
859            &format!("restart {INSTANCE}"),
860            s.get_timeouts().other,
861        )?;
862        wait_up(&self.info, Duration::from_secs(60))?;
863        Ok(report)
864    }
865}
866
867fn not_set_up() -> Error {
868    Error::invalid(
869        "no local registry on this host: run `isb registry setup`, then `sudo isb host setup`",
870    )
871}
872
873/// The registry `base` can reach, for resolving `registry:` images outside
874/// the daemon (CLI deploys, plans).
875pub fn require_info(base: &Client) -> Result<Info> {
876    info(base)?.ok_or_else(not_set_up)
877}
878
879/// What deployed stacks reference, as `repo@digest`: retention never takes
880/// these.
881pub fn protected_by(defs: &[Arc<crate::stack::StackDef>]) -> BTreeSet<String> {
882    let mut out = BTreeSet::new();
883    for d in defs {
884        let mut cur: Option<&crate::stack::StackDef> = Some(d);
885        while let Some(def) = cur {
886            for (svc, spec) in &def.file.services {
887                let Some(r) = spec.image.strip_prefix("registry:") else {
888                    continue;
889                };
890                let Ok(r) = ImageRef::parse(r) else { continue };
891                let digest = r.digest.clone().or_else(|| def.images.get(svc).cloned());
892                if let Some(dg) = digest {
893                    out.insert(format!("{}@{dg}", repo(&def.org, &r.app)));
894                }
895            }
896            cur = def.previous.as_deref();
897        }
898    }
899    out
900}
901
902/// The registry's status for `isb registry ls` and tools.
903pub fn status_json(r: &Registry) -> Value {
904    json!({"addr": r.info.addr, "url": r.info.url()})
905}
906
907#[cfg(test)]
908mod tests;