diff --git a/modules/distribution/README.md b/modules/distribution/README.md new file mode 100644 index 00000000..70d49c73 --- /dev/null +++ b/modules/distribution/README.md @@ -0,0 +1,76 @@ +# distribution + +The mesh's artifact store: an OCI registry that holds every image, mirrored upstream image, bundle and +archive the mesh delivers, by digest (novox/hq ADR 0156). It claims the mesh seat `mesh-artifact-store` +and provides `artifact-store`. + +## What it declares + +| resource | what | +|---|---| +| `state`, `registry-data` | the module's state directory and the store's files | +| `store` | the registry, with deletion enabled on its one door (ADR 0189 §1) | +| `collect` | the registry's own collector, nightly at 03:30, with `store` held still while it runs (ADR 0189 §4) | +| `tools-go` | the Go bundle `store-tools`, below | + +The mesh decides what the store may let go of, from its build records, and lets go of it after each +build it records (ADR 0189). The store's collector reclaims the bytes each night. + +## Tools (novox/hq ADR 0251, to-be 51) + +| tool | | what | +|---|---|---| +| `store_repositories` | r | every repository: tags and what each names, how many manifests (tagged or not), size | +| `store_usage` | r | the store's size, the largest repositories, bytes shared between repositories, bytes no manifest marks (what the nightly collector frees next) | +| `store_references` | r | what the controller's records say of each manifest the store holds, counted and sized by state and repository; each manifest listed when one repository is asked; what the records keep that the store does not hold | +| `store_collect` | a | a dry run unless `dry_run` is false: what the controller would let go of, and the bytes the nightly collector would free then and now. A real run needs `why` and asks the controller's `collect` | + +### How the store is read + +The store's door lists repositories and tags, but not a manifest no tag names, and the mesh pins every +machine by digest, so that is most of them. So the bundle lists the store's own files through its own +container (`docker exec mesh-registry`, busybox `find` and `stat`), read-only: every blob with its size, +every manifest each repository holds, what each tag names. It reads each manifest's content through +the door, eight at a time and within a budget per call, and keeps what it read (a manifest never +changes: it is named by its content). A manifest that could not be read is counted and said; what it +marks beyond itself is then unknown, so sizes may read low and freed bytes high, and the answer says so. + +The docker command runs as the tool runner's account; a socket that refuses it is asked again through +`sudo -n`, never with a prompt, as the container runtime's own tools do. + +A manifest *marks* its own content, its configuration and its layers, and an index marks the manifests +it lists. The store's collector removes every blob no manifest marks, so these numbers are its own. + +### The states of a manifest + +| state | what | removable | +|---|---|---| +| `kept` | a definition names it, or one of the five most recent builds of a module the mesh holds; `why` says which | no | +| `holder-of-kept-archive` | the manifest that keeps a kept archive's blob (hq issue 253) | no | +| `eligible` | the mesh made it and keeps it for no reason | through the controller's `collect` | +| `holder-of-eligible-archive` | the manifest that keeps an eligible archive's blob | with its archive | +| `let-go-yet-present` | the controller recorded letting go of it, and the store still holds it | no — a finding | +| `named-document` | a one-layer manifest the controller keeps under a tag (hq to-be 45 §9) | no | +| `unrecorded` | no record names it | **never**, by any tool (ADR 0189 §3) | + +The records come from the controller's `artifacts` verb, asked with the references already let go of; +when that answer is too large to carry, it is asked without them and the answer says so. + +### Collecting + +`store_collect` never deletes through the store's door. The controller decides and records what it +lets go of (ADR 0189 §2), so a real run asks its `collect` with `why` and `confirm`; the controller holds +every kept archive first and records the run as a hand-act. Bytes come back at the nightly collection. + +## Tests + +``` +go test ./... +``` + +Against a fake listing and a fake door: the listing parsed (a repository name with slashes, a cut +listing refused), the mark (an index, a shared blob, a blob nothing marks), every state, what the +records keep that the store lacks, bytes freed counting a blob shared with a kept manifest as kept, a +dry run never confirming, a real run without `why` refused before anything is asked, the fall-back when +the references let go of cannot be carried, escalation through `sudo -n`, and that the tools served are +exactly the manifest's `tools` and its `invokes` exactly the two controller verbs. diff --git a/modules/distribution/client.ts b/modules/distribution/client.ts deleted file mode 100644 index 5ff3457e..00000000 --- a/modules/distribution/client.ts +++ /dev/null @@ -1,81 +0,0 @@ -// The Docker Registry v2 client — registry's own code, living in the module (novox/hq ADR 0039). -// Ported from the shared hal sdk, where a change to the registry API rebuilt everything; here it -// rebuilds only registry. Both this module's tools and its events entrypoint import it. - -export interface RegistryImage { - repo: string; - tag: string; - digest: string; -} - -export class RegistryClient { - readonly baseUrl: string; - - // Auth is optional: a mesh-internal registry often runs open on the node, so a Basic header is - // sent only when credentials were configured — an empty one would look like a failed login. - constructor( - url: string, - private readonly authHeader?: string, - ) { - this.baseUrl = url.replace(/\/+$/, ""); - } - - /** - * Build from the module's resolved environment. The URL is MESH_REGISTRY_URL (or the local - * registry port), and credentials — if the registry requires them — are MESH_REGISTRY_USER and - * MESH_REGISTRY_PASSWORD. Throws when no URL is configured, so a misconfigured module exposes - * nothing rather than talking to the wrong place. - */ - static fromEnv(env: NodeJS.ProcessEnv = process.env): RegistryClient { - const url = env.MESH_REGISTRY_URL ?? `http://127.0.0.1:${env.REGISTRY_PORT ?? "5000"}`; - if (!url) throw new Error("no registry URL — set MESH_REGISTRY_URL"); - const user = env.MESH_REGISTRY_USER; - const password = env.MESH_REGISTRY_PASSWORD; - const authHeader = - user && password ? `Basic ${Buffer.from(`${user}:${password}`).toString("base64")}` : undefined; - return new RegistryClient(url, authHeader); - } - - private headers(extra: Record = {}): Record { - return { ...(this.authHeader ? { Authorization: this.authHeader } : {}), ...extra }; - } - - private async getJson(path: string): Promise { - const res = await fetch(`${this.baseUrl}${path}`, { headers: this.headers() }); - if (!res.ok) throw new Error(`Registry ${path}: ${res.status} ${await res.text()}`); - return res.json() as Promise; - } - - /** The catalog — every repository the registry holds. */ - async listRepositories(): Promise { - const data = await this.getJson<{ repositories: string[] | null }>("/v2/_catalog"); - return data.repositories ?? []; - } - - /** The tags of one repository. */ - async listTags(repo: string): Promise { - const data = await this.getJson<{ tags: string[] | null }>(`/v2/${repo}/tags/list`); - return data.tags ?? []; - } - - /** The content digest of a repo:tag — the stable identity a tag currently points at. */ - async getManifestDigest(repo: string, tag: string): Promise { - const res = await fetch(`${this.baseUrl}/v2/${repo}/manifests/${tag}`, { - method: "HEAD", - headers: this.headers({ Accept: "application/vnd.docker.distribution.manifest.v2+json" }), - }); - if (!res.ok) throw new Error(`Registry manifest ${repo}:${tag}: ${res.status} ${await res.text()}`); - const digest = res.headers.get("docker-content-digest"); - if (!digest) throw new Error(`no Docker-Content-Digest for ${repo}:${tag}`); - return digest; - } - - /** Delete a manifest by digest. Garbage collection reclaims the storage later. */ - async deleteManifest(repo: string, digest: string): Promise { - const res = await fetch(`${this.baseUrl}/v2/${repo}/manifests/${digest}`, { - method: "DELETE", - headers: this.headers({ Accept: "application/vnd.docker.distribution.manifest.v2+json" }), - }); - if (!res.ok) throw new Error(`Registry delete ${repo}@${digest}: ${res.status} ${await res.text()}`); - } -} diff --git a/modules/distribution/cmd/store-tools/classify.go b/modules/distribution/cmd/store-tools/classify.go new file mode 100644 index 00000000..a01f3f79 --- /dev/null +++ b/modules/distribution/cmd/store-tools/classify.go @@ -0,0 +1,203 @@ +package main + +import ( + "sort" +) + +// The states a manifest the store holds can be in, as the records see it (novox/hq to-be 51). +const ( + StateKept = "kept" + StateHolderKept = "holder-of-kept-archive" + StateEligible = "eligible" + StateHolderEligible = "holder-of-eligible-archive" + StateCollectedPresent = "let-go-yet-present" + StateNamedDocument = "named-document" + StateUnrecorded = "unrecorded" + collectionRemovesThese = "eligible and holder-of-eligible-archive, through the controller's collect; nothing else" +) + +// Entry is one manifest the store holds, and what the records say of it. +type Entry struct { + Repository string `json:"repository"` + Digest string `json:"digest"` + State string `json:"state"` + Why []string `json:"why,omitempty"` + // Archive is the archive a holder keeps. + Archive string `json:"archive,omitempty"` + Tags []string `json:"tags,omitempty"` + // Bytes are the blobs it marks, its own content included. + Bytes int64 `json:"bytes"` + // Unread says its content could not be read, so what it marks is not known beyond itself. + Unread string `json:"unread,omitempty"` +} + +// recordIndex finds a record by repository and digest. +type recordIndex map[string]*Record + +func indexRecords(r *Records) recordIndex { + idx := recordIndex{} + for i := range r.References { + rec := &r.References[i] + if rec.Repository == "" || rec.Digest == "" { + continue + } + idx[rec.Kind+" "+Key(rec.Repository, rec.Digest)] = rec + } + return idx +} + +// Classify gives every manifest the store holds one state. +func Classify(v *View, recs *Records) []Entry { + idx := indexRecords(recs) + var out []Entry + for _, name := range v.L.RepoNames() { + r := v.L.Repos[name] + tagsAt := map[string][]string{} + for t, d := range r.Tags { + tagsAt[d] = append(tagsAt[d], t) + } + digests := make([]string, 0, len(r.Revisions)) + for d := range r.Revisions { + digests = append(digests, d) + } + sort.Strings(digests) + for _, d := range digests { + e := Entry{Repository: name, Digest: d, State: StateUnrecorded, Tags: tagsAt[d]} + sort.Strings(e.Tags) + blobs := map[string]bool{} + for _, b := range v.Marks(name, d) { + blobs[b] = true + } + e.Bytes = v.Bytes(blobs) + e.Unread = v.Unread[Key(name, d)] + m := v.M[Key(name, d)] + if rec := idx["image "+Key(name, d)]; rec != nil { + e.State, e.Why = stateOf(rec, false), rec.Why + } else if m.Holder() { + archive := m.Layers[0] + if rec := idx["archive "+Key(name, archive)]; rec != nil { + e.State, e.Why, e.Archive = stateOf(rec, true), rec.Why, archive + } else if len(e.Tags) > 0 { + e.State = StateNamedDocument + } + } + out = append(out, e) + } + } + return out +} + +func stateOf(rec *Record, holder bool) string { + switch rec.State { + case "kept": + if holder { + return StateHolderKept + } + return StateKept + case "eligible": + if holder { + return StateHolderEligible + } + return StateEligible + case "collected": + return StateCollectedPresent + } + // A state this bundle does not know is kept: never offered as removable. + return StateKept +} + +// Missing are the references the records keep that the store does not hold: an image whose manifest +// is not in its repository, an archive whose blob is not in the store. +func Missing(v *View, recs *Records) []string { + var out []string + for _, rec := range recs.References { + if rec.State != "kept" || rec.Repository == "" || rec.Digest == "" { + continue + } + switch rec.Kind { + case "image": + if r := v.L.Repos[rec.Repository]; r == nil || !r.Revisions[rec.Digest] { + out = append(out, rec.Reference) + } + case "archive": + if _, ok := v.L.Blobs[rec.Digest]; !ok { + out = append(out, rec.Reference) + } + } + } + sort.Strings(out) + return out +} + +// StateSum is the manifests of one state: how many, the bytes they mark together, and the bytes only +// they mark — what the store would give back if they alone went. +type StateSum struct { + Manifests int `json:"manifests"` + Bytes int64 `json:"bytes"` + OnlyTheirBytes int64 `json:"only_their_bytes"` +} + +// Summarise sums entries by state, over the whole store or one repository's entries, the "only theirs" +// always judged against every manifest in the store. +func Summarise(v *View, entries []Entry, all []Entry) map[string]*StateSum { + statesOf := map[string]map[string]bool{} + for _, e := range all { + for _, b := range v.Marks(e.Repository, e.Digest) { + if statesOf[b] == nil { + statesOf[b] = map[string]bool{} + } + statesOf[b][e.State] = true + } + } + sums := map[string]*StateSum{} + blobs := map[string]map[string]bool{} + for _, e := range entries { + s := sums[e.State] + if s == nil { + s = &StateSum{} + sums[e.State] = s + blobs[e.State] = map[string]bool{} + } + s.Manifests++ + for _, b := range v.Marks(e.Repository, e.Digest) { + blobs[e.State][b] = true + } + } + for state, set := range blobs { + sums[state].Bytes = v.Bytes(set) + only := map[string]bool{} + for b := range set { + if len(statesOf[b]) == 1 { + only[b] = true + } + } + sums[state].OnlyTheirBytes = v.Bytes(only) + } + return sums +} + +// LetGoKeys are the manifests the store would no longer hold once the controller lets go of these +// references: an image's manifest, and an archive's holder. +func LetGoKeys(v *View, refs []string) map[string]bool { + out := map[string]bool{} + for _, ref := range refs { + repo, digest, kind := splitReference(ref) + r := v.L.Repos[repo] + if r == nil { + continue + } + switch kind { + case "image": + if r.Revisions[digest] { + out[Key(repo, digest)] = true + } + case "archive": + for d := range r.Revisions { + if m := v.M[Key(repo, d)]; m.Holder() && m.Layers[0] == digest { + out[Key(repo, d)] = true + } + } + } + } + return out +} diff --git a/modules/distribution/cmd/store-tools/main.go b/modules/distribution/cmd/store-tools/main.go new file mode 100644 index 00000000..df45280f --- /dev/null +++ b/modules/distribution/cmd/store-tools/main.go @@ -0,0 +1,31 @@ +// The distribution module's Go bundle (novox/hq ADR 0251 §1–3, to-be 51): a process the node's runtime +// launches and speaks MCP over stdio to. It says what the artifact store holds — its repositories, its +// size, and what the controller's records say of each manifest — read from the store's own files +// through its own container and from its door, and writing to neither. A collection is asked of the +// controller, which decides and records what it lets go of (ADR 0189). stdout is the protocol; what +// this bundle says, it says on stderr. +package main + +import ( + "fmt" + "os" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +func main() { + s := &Store{ + URL: os.Getenv("MESH_STORE_URL"), + Container: os.Getenv("MESH_STORE_CONTAINER"), + Run: ExecRunner, + UID: os.Getuid(), + } + if s.URL == "" { + fmt.Fprintln(os.Stderr, "[store-tools] MESH_STORE_URL is not set: every manifest will be said as unread") + } + // An empty name serves as the module the runtime names (MESH_SERVED_MODULE): distribution. + if err := stdio.Serve("", Tools(s, Controller{Ask: stdio.Ask})); err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } +} diff --git a/modules/distribution/cmd/store-tools/records.go b/modules/distribution/cmd/store-tools/records.go new file mode 100644 index 00000000..1f0e2b93 --- /dev/null +++ b/modules/distribution/cmd/store-tools/records.go @@ -0,0 +1,209 @@ +package main + +import ( + "encoding/json" + "errors" + "fmt" + "strconv" + "strings" +) + +// What the records say, asked of the controller and never copied (novox/hq ADR 0251 §2). The +// controller decides what the mesh keeps; this bundle only sets that beside what the store holds. + +// Asking is how a tool on the mesh is asked: stdio.Ask, or a test's. +type Asking func(key string, body any) (json.RawMessage, error) + +// ControllerSeat is the seat whose verbs the controller serves. +const ControllerSeat = "mesh-controller" + +// Record is one reference the mesh recorded making, as the controller's `artifacts` answers it. +type Record struct { + Reference string `json:"reference"` + Repository string `json:"repository,omitempty"` + Digest string `json:"digest,omitempty"` + Kind string `json:"kind,omitempty"` + State string `json:"state"` + Why []string `json:"why,omitempty"` +} + +// Records is the controller's `artifacts` answer. +type Records struct { + Store string `json:"store"` + KeptBuilds int `json:"kept_builds"` + Counts map[string]int `json:"counts"` + References []Record `json:"references"` + // CollectedAsked is whether the references already let go of were asked for too. + CollectedAsked bool `json:"-"` + // Note says what could not be asked, when something could not. + Note string `json:"-"` +} + +// CollectAnswer is the controller's `collect` answer. +type CollectAnswer struct { + DryRun bool `json:"dry_run"` + Store string `json:"store"` + KeptArchives int `json:"kept_archives"` + Held int `json:"held"` + HoldersWritten int `json:"holders_written"` + Missing int `json:"missing"` + Eligible int `json:"eligible"` + WouldLetGo []string `json:"would_let_go"` + LetGo []string `json:"let_go"` + Skipped int `json:"skipped"` + Left int `json:"left"` + Stopped string `json:"stopped"` + // Rest is whatever else the controller said, passed on as it said it. + Rest map[string]json.RawMessage `json:"-"` +} + +// Controller asks the controller's seat through the runtime. +type Controller struct{ Ask Asking } + +// answerOf reads a verb's answer whatever wraps it: the controller's {output, ok, answer}, a text the +// runtime handed over, or the protocol's content list. +func answerOf(raw json.RawMessage) (json.RawMessage, string, bool) { + var s string + if json.Unmarshal(raw, &s) == nil { + return answerOf(json.RawMessage(s)) + } + var m map[string]json.RawMessage + if json.Unmarshal(raw, &m) != nil { + return raw, string(raw), true + } + if content, has := m["content"]; has { + var items []struct { + Text string `json:"text"` + } + var e struct { + IsError bool `json:"isError"` + } + _ = json.Unmarshal(raw, &e) + if json.Unmarshal(content, &items) == nil && len(items) > 0 { + inner, out, ok := answerOf(json.RawMessage(items[0].Text)) + return inner, out, ok && !e.IsError + } + } + var v struct { + Output string `json:"output"` + OK *bool `json:"ok"` + Answer json.RawMessage `json:"answer"` + } + if json.Unmarshal(raw, &v) == nil && v.OK != nil { + return v.Answer, v.Output, *v.OK + } + return raw, string(raw), true +} + +func (c Controller) verb(name string, args map[string]any) (json.RawMessage, error) { + if c.Ask == nil { + return nil, errors.New("this bundle has no way to ask the controller") + } + raw, err := c.Ask("seat:"+ControllerSeat+"."+name, args) + if err != nil { + return nil, fmt.Errorf("the controller's %s could not be asked: %w", name, err) + } + answer, output, ok := answerOf(raw) + if !ok { + return nil, fmt.Errorf("the controller refused %s: %s", name, lastLine(output)) + } + if len(answer) == 0 || string(answer) == "null" { + return nil, fmt.Errorf("the controller's %s answered no data: %s", name, lastLine(output)) + } + return answer, nil +} + +func lastLine(s string) string { + lines := strings.Split(strings.TrimSpace(s), "\n") + return strings.TrimSpace(lines[len(lines)-1]) +} + +// Artifacts asks what the records say, with the references already let go of when the answer can carry +// them; when it cannot, without, and says so. +func (c Controller) Artifacts(repository string) (*Records, error) { + args := map[string]any{"collected": "true"} + if repository != "" { + args["repository"] = repository + } + raw, err := c.verb("artifacts", args) + collected := true + note := "" + if err != nil { + delete(args, "collected") + var again error + raw, again = c.verb("artifacts", args) + if again != nil { + return nil, again + } + collected = false + note = "the references the mesh already let go of could not be asked for (" + err.Error() + + "), so a manifest the mesh let go of that the store still holds is counted as unrecorded" + } + var r Records + if err := json.Unmarshal(raw, &r); err != nil { + return nil, fmt.Errorf("the controller's artifacts answer is not readable: %v", err) + } + for i := range r.References { + fill(&r.References[i]) + } + r.CollectedAsked, r.Note = collected, note + return &r, nil +} + +// Collect asks the controller's sweep: a dry run unless confirm, which needs why. +func (c Controller) Collect(why string, confirm bool, most int) (*CollectAnswer, error) { + args := map[string]any{} + if confirm { + if strings.TrimSpace(why) == "" { + return nil, errors.New("a real collection needs why: nothing was asked") + } + args["why"] = why + args["confirm"] = "true" + } + if most > 0 { + args["most"] = strconv.Itoa(most) + } + raw, err := c.verb("collect", args) + if err != nil { + return nil, err + } + var a CollectAnswer + if err := json.Unmarshal(raw, &a); err != nil { + return nil, fmt.Errorf("the controller's collect answer is not readable: %v", err) + } + _ = json.Unmarshal(raw, &a.Rest) + if !confirm && !a.DryRun && len(a.LetGo) > 0 { + return nil, fmt.Errorf("the controller let go of %d artifacts when only a dry run was asked", len(a.LetGo)) + } + return &a, nil +} + +// fill reads the repository, digest and kind out of a reference when the answer left them out. +func fill(r *Record) { + repo, digest, kind := splitReference(r.Reference) + if r.Repository == "" { + r.Repository = repo + } + if r.Digest == "" { + r.Digest = digest + } + if r.Kind == "" { + r.Kind = kind + } +} + +// splitReference reads `artifact-store://@sha256:…` (an image) and +// `artifact-store:///blobs/sha256:…` (an archive). +func splitReference(ref string) (repo, digest, kind string) { + path := ref + if _, after, ok := strings.Cut(ref, "://"); ok { + path = after + } + if before, after, ok := strings.Cut(path, "/blobs/sha256:"); ok { + return before, "sha256:" + after, "archive" + } + if before, after, ok := strings.Cut(path, "@sha256:"); ok { + return before, "sha256:" + after, "image" + } + return "", "", "" +} diff --git a/modules/distribution/cmd/store-tools/runner.go b/modules/distribution/cmd/store-tools/runner.go new file mode 100644 index 00000000..726309b6 --- /dev/null +++ b/modules/distribution/cmd/store-tools/runner.go @@ -0,0 +1,96 @@ +package main + +import ( + "bytes" + "context" + "errors" + "fmt" + "os/exec" + "regexp" + "strings" + "time" +) + +// Ran is what a command did: its output, its exit status, and why it never ran to an answer. +type Ran struct { + Stdout string + Stderr string + Status int + // Err is "ENOENT" when the program is not installed, or says it was ended for taking too long. + Err string +} + +// Runner runs one command, so every tool can be tested without a daemon. +type Runner func(ctx context.Context, name string, args ...string) Ran + +// ListTimeout is how long listing the store's files may take: below the runtime's thirty-second call +// limit, so a store that hangs is answered as such rather than as a call the runtime gave up on. +const ListTimeout = 15 * time.Second + +// ExecRunner runs a command on this machine, bounded by ListTimeout. +func ExecRunner(ctx context.Context, name string, args ...string) Ran { + ctx, cancel := context.WithTimeout(ctx, ListTimeout) + defer cancel() + cmd := exec.CommandContext(ctx, name, args...) + var out, errb bytes.Buffer + cmd.Stdout, cmd.Stderr = &out, &errb + err := cmd.Run() + r := Ran{Stdout: out.String(), Stderr: errb.String()} + var exitErr *exec.ExitError + switch { + case errors.Is(ctx.Err(), context.DeadlineExceeded): + r.Status, r.Err = 124, fmt.Sprintf("no answer within %d s", int(ListTimeout/time.Second)) + case errors.Is(err, exec.ErrNotFound): + r.Status, r.Err = 127, "ENOENT" + case errors.As(err, &exitErr): + r.Status = exitErr.ExitCode() + case err != nil: + r.Status, r.Err = 1, err.Error() + } + return r +} + +var socketRefused = regexp.MustCompile(`(?i)permission denied.*docker.*sock|docker\.sock.*permission denied`) + +// docker runs one docker command as this account and answers its stdout. A socket that refuses the +// account is asked again through `sudo -n`, as the container runtime's own tools do; never as root +// otherwise, and never with a prompt. +func docker(ctx context.Context, run Runner, uid int, args ...string) (string, error) { + r := run(ctx, "docker", args...) + program := "docker" + if r.Status != 0 && r.Err == "" && uid != 0 && socketRefused.MatchString(r.Stderr) { + program = "sudo" + r = run(ctx, "sudo", append([]string{"-n", "docker"}, args...)...) + } + if r.Status == 0 && r.Err == "" { + return r.Stdout, nil + } + said := strings.TrimSpace(r.Stderr + "\n" + r.Stdout) + switch { + case r.Err == "ENOENT" && program == "sudo": + return "", errors.New("the runtime's socket refused this account, and sudo is not installed here to escalate with") + case r.Err == "ENOENT": + return "", errors.New("docker is not installed on this machine, so the store's files cannot be listed") + case r.Err != "": + return "", fmt.Errorf("listing the store's files did not answer: %s", r.Err) + case program == "sudo" && strings.HasPrefix(said, "sudo:"): + return "", fmt.Errorf("the runtime's socket refused this account and it may not escalate without a prompt: %s", firstLine(said)) + case strings.Contains(said, "No such container"): + return "", fmt.Errorf("the store's container is not on this machine: %s", firstLine(said)) + case strings.Contains(said, "is not running"): + return "", fmt.Errorf("the store's container is not running: %s", firstLine(said)) + } + if l := firstLine(said); l != "" { + return "", fmt.Errorf("listing the store's files failed (%d): %s", r.Status, l) + } + return "", fmt.Errorf("listing the store's files failed with status %d", r.Status) +} + +func firstLine(s string) string { + for _, l := range strings.Split(s, "\n") { + if l = strings.TrimSpace(l); l != "" { + return l + } + } + return "" +} diff --git a/modules/distribution/cmd/store-tools/store.go b/modules/distribution/cmd/store-tools/store.go new file mode 100644 index 00000000..62086bac --- /dev/null +++ b/modules/distribution/cmd/store-tools/store.go @@ -0,0 +1,313 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "sort" + "strconv" + "strings" + "sync" + "time" +) + +// What the store holds, read from its own files and its own door, and nothing written to either. +// +// **Why its files.** The store's door lists repositories and the tags in each, but not a manifest no tag +// names — and since the mesh pins every machine by digest, that is almost every manifest the store +// holds. The files list them all, with every blob's size. They are read through the store's own +// container, which is the only process that has them, and only read (novox/hq ADR 0251 §1). + +// storageRoot is where the registry keeps its files, inside its container. +const storageRoot = "/var/lib/registry/docker/registry/v2" + +// listScript lists, in three sections, every blob with its size, every manifest each repository holds, +// and what each tag names now. Busybox's find and stat (the store's image is Alpine): no -printf. An +// upload in progress is not listed, and nor are a tag's former targets. +const listScript = `set -e +cd ` + storageRoot + ` +echo '#blobs' +if [ -d blobs ]; then find blobs -type f -name data | xargs -r stat -c '%s %n'; fi +echo '#revisions' +if [ -d repositories ]; then find repositories -type f -path '*/_manifests/revisions/sha256/*/link'; fi +echo '#tags' +if [ -d repositories ]; then find repositories -type f -path '*/_manifests/tags/*/current/link' -exec grep -H . {} + || true; fi +echo '#end'` + +// Repo is one repository as the store's files list it. +type Repo struct { + Name string + // Revisions are the digests of every manifest it holds, tagged or not. + Revisions map[string]bool + // Tags are what each tag names now. + Tags map[string]string +} + +// Listing is the store's files: every blob with its size, and every repository. +type Listing struct { + Blobs map[string]int64 + Repos map[string]*Repo +} + +func (l *Listing) repo(name string) *Repo { + r := l.Repos[name] + if r == nil { + r = &Repo{Name: name, Revisions: map[string]bool{}, Tags: map[string]string{}} + l.Repos[name] = r + } + return r +} + +// RepoNames are the repositories in order. +func (l *Listing) RepoNames() []string { + out := make([]string, 0, len(l.Repos)) + for n := range l.Repos { + out = append(out, n) + } + sort.Strings(out) + return out +} + +// ParseListing reads listScript's output. A section that never ended means the listing was cut short, +// which is an error: a partial listing would make held things look absent. +func ParseListing(out string) (*Listing, error) { + l := &Listing{Blobs: map[string]int64{}, Repos: map[string]*Repo{}} + section := "" + ended := false + for _, line := range strings.Split(out, "\n") { + line = strings.TrimSpace(line) + if line == "" { + continue + } + if strings.HasPrefix(line, "#") { + section = line + ended = line == "#end" + continue + } + switch section { + case "#blobs": + size, path, ok := strings.Cut(line, " ") + n, err := strconv.ParseInt(size, 10, 64) + if !ok || err != nil { + return nil, fmt.Errorf("the store's listing has a blob line it cannot read: %q", line) + } + // blobs/sha256/ab//data + parts := strings.Split(strings.TrimPrefix(path, "./"), "/") + if len(parts) != 5 || parts[0] != "blobs" || parts[4] != "data" { + return nil, fmt.Errorf("the store's listing has a blob at an unexpected place: %q", path) + } + l.Blobs[parts[1]+":"+parts[3]] = n + case "#revisions": + // repositories//_manifests/revisions/sha256//link + name, rest, ok := strings.Cut(strings.TrimPrefix(strings.TrimPrefix(line, "./"), "repositories/"), "/_manifests/revisions/") + parts := strings.Split(rest, "/") + if !ok || len(parts) != 3 || parts[2] != "link" { + return nil, fmt.Errorf("the store's listing has a manifest at an unexpected place: %q", line) + } + l.repo(name).Revisions[parts[0]+":"+parts[1]] = true + case "#tags": + // repositories//_manifests/tags//current/link:sha256: + path, digest, ok := strings.Cut(line, ":sha256:") + name, rest, ok2 := strings.Cut(strings.TrimPrefix(strings.TrimPrefix(path, "./"), "repositories/"), "/_manifests/tags/") + tag, _, ok3 := strings.Cut(rest, "/current/link") + if !ok || !ok2 || !ok3 || tag == "" { + return nil, fmt.Errorf("the store's listing has a tag line it cannot read: %q", line) + } + l.repo(name).Tags[tag] = "sha256:" + digest + default: + return nil, fmt.Errorf("the store's listing has a line outside any section: %q", line) + } + } + if !ended { + return nil, fmt.Errorf("the store's listing was cut short (it never reached its end)") + } + return l, nil +} + +// Manifest is what one manifest names. +type Manifest struct { + MediaType string + ConfigMedia string + // Refs are the blobs it names: its configuration and its layers. + Refs []string + // Layers are its layers alone, in order. + Layers []string + // Children are the manifests an index names. + Children []string +} + +// Holder is the shape of a manifest that keeps one blob: an empty configuration and one layer. The +// controller holds every archive the mesh keeps this way (novox/hq issue 253), and keeps a named +// document the same way, under a tag. +func (m *Manifest) Holder() bool { + return m != nil && m.ConfigMedia == "application/vnd.oci.empty.v1+json" && len(m.Layers) == 1 +} + +// ParseManifest reads a manifest's content, of any of the shapes a registry holds. +func ParseManifest(raw []byte) (*Manifest, error) { + type desc struct { + MediaType string `json:"mediaType"` + Digest string `json:"digest"` + } + var v struct { + MediaType string `json:"mediaType"` + Config *desc `json:"config"` + Layers []desc `json:"layers"` + Manifests []desc `json:"manifests"` + FSLayers []struct { + BlobSum string `json:"blobSum"` + } `json:"fsLayers"` + } + if err := json.Unmarshal(raw, &v); err != nil { + return nil, fmt.Errorf("not a manifest: %v", err) + } + m := &Manifest{MediaType: v.MediaType} + if v.Config != nil && v.Config.Digest != "" { + m.ConfigMedia = v.Config.MediaType + m.Refs = append(m.Refs, v.Config.Digest) + } + for _, d := range v.Layers { + m.Refs = append(m.Refs, d.Digest) + m.Layers = append(m.Layers, d.Digest) + } + for _, d := range v.FSLayers { + m.Refs = append(m.Refs, d.BlobSum) + m.Layers = append(m.Layers, d.BlobSum) + } + for _, d := range v.Manifests { + m.Children = append(m.Children, d.Digest) + } + return m, nil +} + +var manifestAccept = []string{ + "application/vnd.oci.image.manifest.v1+json", + "application/vnd.oci.image.index.v1+json", + "application/vnd.docker.distribution.manifest.v2+json", + "application/vnd.docker.distribution.manifest.list.v2+json", + "application/vnd.docker.distribution.manifest.v1+prettyjws", +} + +// Store reaches the store: its files through its container, its manifests through its door. +type Store struct { + URL string + Container string + Run Runner + UID int + HTTP *http.Client + // ReadBudget bounds reading manifests in one call; what is not read in it is said. + ReadBudget time.Duration + + mu sync.Mutex + // cache holds every manifest read: a manifest is named by its content's digest, so it never changes. + cache map[string]*Manifest +} + +// List reads the store's files. +func (s *Store) List(ctx context.Context) (*Listing, error) { + if s.Container == "" { + return nil, fmt.Errorf("this bundle was not told the store's container (MESH_STORE_CONTAINER)") + } + out, err := docker(ctx, s.Run, s.UID, "exec", s.Container, "sh", "-c", listScript) + if err != nil { + return nil, err + } + return ParseListing(out) +} + +// Key names one manifest in one repository. +func Key(repo, digest string) string { return repo + "@" + digest } + +// Manifests reads every manifest the listing names, from the cache or the door, eight at a time and +// within the read budget. Answers each one read, and each one not read with why. +func (s *Store) Manifests(ctx context.Context, l *Listing) (map[string]*Manifest, map[string]string) { + type job struct{ repo, digest string } + var jobs []job + got := map[string]*Manifest{} + unread := map[string]string{} + s.mu.Lock() + if s.cache == nil { + s.cache = map[string]*Manifest{} + } + for _, name := range l.RepoNames() { + for d := range l.Repos[name].Revisions { + if m, ok := s.cache[d]; ok { + got[Key(name, d)] = m + } else { + jobs = append(jobs, job{name, d}) + } + } + } + s.mu.Unlock() + + budget := s.ReadBudget + if budget == 0 { + budget = 12 * time.Second + } + ctx, cancel := context.WithTimeout(ctx, budget) + defer cancel() + var mu sync.Mutex + work := make(chan job) + var wg sync.WaitGroup + for w := 0; w < 8; w++ { + wg.Add(1) + go func() { + defer wg.Done() + for j := range work { + m, err := s.read(ctx, j.repo, j.digest) + mu.Lock() + if err != nil { + unread[Key(j.repo, j.digest)] = err.Error() + } else { + got[Key(j.repo, j.digest)] = m + } + mu.Unlock() + } + }() + } + for _, j := range jobs { + work <- j + } + close(work) + wg.Wait() + return got, unread +} + +func (s *Store) read(ctx context.Context, repo, digest string) (*Manifest, error) { + if err := ctx.Err(); err != nil { + return nil, fmt.Errorf("not read within this call's budget") + } + req, err := http.NewRequestWithContext(ctx, http.MethodGet, strings.TrimRight(s.URL, "/")+"/v2/"+repo+"/manifests/"+digest, nil) + if err != nil { + return nil, err + } + for _, a := range manifestAccept { + req.Header.Add("Accept", a) + } + client := s.HTTP + if client == nil { + client = &http.Client{Timeout: 10 * time.Second} + } + resp, err := client.Do(req) + if err != nil { + return nil, fmt.Errorf("the store's door did not answer: %v", err) + } + defer resp.Body.Close() + body, err := io.ReadAll(io.LimitReader(resp.Body, 4<<20)) + if err != nil { + return nil, err + } + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("the store's door answered %s", resp.Status) + } + m, err := ParseManifest(body) + if err != nil { + return nil, err + } + s.mu.Lock() + s.cache[digest] = m + s.mu.Unlock() + return m, nil +} diff --git a/modules/distribution/cmd/store-tools/store_test.go b/modules/distribution/cmd/store-tools/store_test.go new file mode 100644 index 00000000..b9e51314 --- /dev/null +++ b/modules/distribution/cmd/store-tools/store_test.go @@ -0,0 +1,421 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "os" + "reflect" + "sort" + "strings" + "sync" + "testing" +) + +// d is a digest made from a short name, so a test reads as names. +func d(name string) string { + return fmt.Sprintf("sha256:%064x", []byte(name))[:71] +} + +func hexOf(digest string) string { return strings.TrimPrefix(digest, "sha256:") } + +// fakeStore is a store's files and its manifests: what listScript would print, and a door. +type fakeStore struct { + blobs map[string]int64 + revisions map[string][]string // repo -> digests + tags map[string]map[string]string + manifests map[string]string // digest -> content +} + +func (f *fakeStore) listing() string { + var b strings.Builder + b.WriteString("#blobs\n") + for dg, n := range f.blobs { + fmt.Fprintf(&b, "%d blobs/sha256/%s/%s/data\n", n, hexOf(dg)[:2], hexOf(dg)) + } + b.WriteString("#revisions\n") + for repo, ds := range f.revisions { + for _, dg := range ds { + fmt.Fprintf(&b, "repositories/%s/_manifests/revisions/sha256/%s/link\n", repo, hexOf(dg)) + } + } + b.WriteString("#tags\n") + for repo, ts := range f.tags { + for t, dg := range ts { + fmt.Fprintf(&b, "repositories/%s/_manifests/tags/%s/current/link:%s\n", repo, t, dg) + } + } + b.WriteString("#end\n") + return b.String() +} + +func (f *fakeStore) door() *httptest.Server { + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, digest, ok := strings.Cut(r.URL.Path, "/manifests/") + if !ok || r.Method != http.MethodGet { + http.Error(w, "no", http.StatusMethodNotAllowed) + return + } + body, ok := f.manifests[digest] + if !ok { + http.NotFound(w, r) + return + } + w.Write([]byte(body)) + })) +} + +func image(config string, layers ...string) string { + var ls []string + for _, l := range layers { + ls = append(ls, fmt.Sprintf(`{"mediaType":"application/vnd.oci.image.layer.v1.tar+gzip","digest":%q,"size":1}`, l)) + } + return fmt.Sprintf(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","config":{"mediaType":"application/vnd.oci.image.config.v1+json","digest":%q,"size":1},"layers":[%s]}`, + config, strings.Join(ls, ",")) +} + +func holder(blob string) string { + return fmt.Sprintf(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","config":{"mediaType":"application/vnd.oci.empty.v1+json","digest":%q,"size":2},"layers":[{"mediaType":"application/vnd.oci.image.layer.v1.tar+gzip","digest":%q,"size":1}]}`, + d("empty"), blob) +} + +func index(children ...string) string { + var cs []string + for _, c := range children { + cs = append(cs, fmt.Sprintf(`{"mediaType":"application/vnd.oci.image.manifest.v1+json","digest":%q,"size":1}`, c)) + } + return fmt.Sprintf(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.index.v1+json","manifests":[%s]}`, strings.Join(cs, ",")) +} + +// world is one store with every state in it: +// +// app/server: m-kept (definition), m-old (eligible), m-stranger (unrecorded, tagged), an index naming m-kept +// app/tools: h-kept holds a-kept (kept archive), h-old holds a-old (eligible archive) +// mesh/facts: doc (a named document under the tag latest) +// base layer "shared" is marked by m-kept and m-old; "old-only" only by m-old. +func world() *fakeStore { + f := &fakeStore{ + blobs: map[string]int64{}, + revisions: map[string][]string{}, + tags: map[string]map[string]string{}, + manifests: map[string]string{}, + } + put := func(repo, name, content string) { + f.manifests[d(name)] = content + f.blobs[d(name)] = int64(len(content)) + f.revisions[repo] = append(f.revisions[repo], d(name)) + } + for name, size := range map[string]int64{"cfg": 10, "shared": 1000, "old-only": 500, "stranger-layer": 300, "a-kept": 2000, "a-old": 4000, "empty": 2, "doc-body": 50, "orphan": 7000} { + f.blobs[d(name)] = size + } + put("app/server", "m-kept", image(d("cfg"), d("shared"))) + put("app/server", "m-old", image(d("cfg"), d("shared"), d("old-only"))) + put("app/server", "m-stranger", image(d("cfg"), d("stranger-layer"))) + put("app/server", "idx", index(d("m-kept"))) + put("app/tools", "h-kept", holder(d("a-kept"))) + put("app/tools", "h-old", holder(d("a-old"))) + put("mesh/facts", "doc", holder(d("doc-body"))) + f.tags["app/server"] = map[string]string{"latest": d("m-stranger")} + f.tags["mesh/facts"] = map[string]string{"latest": d("doc")} + return f +} + +func records() *Records { + return &Records{Store: "store:5000", KeptBuilds: 5, Counts: map[string]int{"kept": 3, "eligible": 2}, + References: []Record{ + {Reference: "artifact-store://app/server@" + d("m-kept"), State: "kept", Why: []string{"definition"}}, + {Reference: "artifact-store://app/server@" + d("m-old"), State: "eligible"}, + {Reference: "artifact-store://app/tools/blobs/" + d("a-kept"), State: "kept", Why: []string{"recent-build"}}, + {Reference: "artifact-store://app/tools/blobs/" + d("a-old"), State: "eligible"}, + {Reference: "artifact-store://app/server@" + d("m-gone"), State: "kept", Why: []string{"definition"}}, + }} +} + +func view(t *testing.T, f *fakeStore) *View { + t.Helper() + door := f.door() + t.Cleanup(door.Close) + s := &Store{URL: door.URL, Container: "mesh-registry", UID: 1000, Run: func(_ context.Context, name string, args ...string) Ran { + if name != "docker" || args[0] != "exec" || args[1] != "mesh-registry" { + t.Fatalf("ran %s %v", name, args) + } + return Ran{Stdout: f.listing()} + }} + v, err := s.View(context.Background()) + if err != nil { + t.Fatal(err) + } + return v +} + +func withFill(r *Records) *Records { + for i := range r.References { + fill(&r.References[i]) + } + return r +} + +func TestTheListingIsRead(t *testing.T) { + l, err := ParseListing("#blobs\n12 blobs/sha256/ab/abcd/data\n#revisions\nrepositories/a/b/c/_manifests/revisions/sha256/ff/link\n#tags\nrepositories/a/b/c/_manifests/tags/v1/current/link:sha256:ff\n#end\n") + if err != nil { + t.Fatal(err) + } + if l.Blobs["sha256:abcd"] != 12 { + t.Errorf("blobs %v", l.Blobs) + } + r := l.Repos["a/b/c"] + if r == nil || !r.Revisions["sha256:ff"] || r.Tags["v1"] != "sha256:ff" { + t.Errorf("repository with slashes misread: %+v", r) + } +} + +func TestACutListingIsRefused(t *testing.T) { + for _, out := range []string{"#blobs\n12 blobs/sha256/ab/abcd/data\n", "#blobs\nnot a line\n#end\n", "stray\n#end\n"} { + if _, err := ParseListing(out); err == nil { + t.Errorf("accepted %q", out) + } + } +} + +func TestTheMarkIsTheCollectorsIncludingAnIndexAndSharedBlobs(t *testing.T) { + v := view(t, world()) + marked := v.Marked(nil) + for _, name := range []string{"cfg", "shared", "old-only", "stranger-layer", "a-kept", "a-old", "doc-body", "empty", "m-kept", "idx"} { + if !marked[d(name)] { + t.Errorf("%s not marked", name) + } + } + if marked[d("orphan")] { + t.Error("a blob no manifest names is marked") + } + freed, n := v.Unmarked(nil) + if freed != 7000 || n != 1 { + t.Errorf("the collector frees %d in %d, want 7000 in 1", freed, n) + } + // The index alone gone frees its own content and nothing m-kept still marks. + _, n = v.Unmarked(map[string]bool{Key("app/server", d("idx")): true}) + if n != 2 { + t.Errorf("without the index %d blobs unmarked, want 2 (orphan and the index itself)", n) + } + if sh := v.Shared(); !sh[d("empty")] || sh[d("cfg")] { + t.Error("shared misjudged") + } +} + +func TestEveryStateIsGiven(t *testing.T) { + v := view(t, world()) + got := map[string]string{} + for _, e := range Classify(v, withFill(records())) { + got[e.Repository+" "+e.Digest] = e.State + } + want := map[string]string{ + "app/server " + d("m-kept"): StateKept, + "app/server " + d("m-old"): StateEligible, + "app/server " + d("m-stranger"): StateUnrecorded, + "app/server " + d("idx"): StateUnrecorded, + "app/tools " + d("h-kept"): StateHolderKept, + "app/tools " + d("h-old"): StateHolderEligible, + "mesh/facts " + d("doc"): StateNamedDocument, + } + if !reflect.DeepEqual(got, want) { + t.Errorf("got %v\nwant %v", got, want) + } + if m := Missing(v, withFill(records())); len(m) != 1 || !strings.Contains(m[0], hexOf(d("m-gone"))) { + t.Errorf("missing %v", m) + } +} + +func TestALetGoReferenceStillPresentIsSaidSo(t *testing.T) { + v := view(t, world()) + r := records() + r.References = append(r.References, Record{Reference: "artifact-store://app/server@" + d("m-stranger"), State: "collected"}) + for _, e := range Classify(v, withFill(r)) { + if e.Digest == d("m-stranger") && e.State != StateCollectedPresent { + t.Errorf("state %s", e.State) + } + } +} + +func TestAStateThisBundleDoesNotKnowIsKept(t *testing.T) { + if stateOf(&Record{State: "something-new"}, false) != StateKept { + t.Fatal("an unknown state was not kept") + } +} + +func TestBytesFreedCountASharedBlobAsKept(t *testing.T) { + v := view(t, world()) + a := &CollectAnswer{DryRun: true, Eligible: 2, WouldLetGo: []string{ + "artifact-store://app/server@" + d("m-old"), "artifact-store://app/tools/blobs/" + d("a-old")}} + out := Collected(v, a, true) + now, after := out["collector_frees_now_bytes"].(int64), out["collector_frees_after_bytes"].(int64) + // m-old's own content, old-only and a-old and its holder go; shared and cfg stay with m-kept. + want := int64(7000) + 500 + 4000 + v.L.Blobs[d("m-old")] + v.L.Blobs[d("h-old")] + if now != 7000 || after != want { + t.Errorf("now %d after %d, want 7000 and %d", now, after, want) + } + if out["manifests_gone"].(int) != 2 { + t.Errorf("gone %v", out["manifests_gone"]) + } +} + +func TestOnlyTheirBytesLeaveOutWhatAnotherStateMarks(t *testing.T) { + v := view(t, world()) + entries := Classify(v, withFill(records())) + sums := Summarise(v, entries, entries) + el := sums[StateEligible] + if el.OnlyTheirBytes != 500+v.L.Blobs[d("m-old")] { + t.Errorf("eligible only-theirs %d", el.OnlyTheirBytes) + } +} + +// asked records every ask, and answers as the controller would. +type asked struct { + mu sync.Mutex + calls []string + args []map[string]any + fail map[string]bool +} + +func (a *asked) ask(key string, body any) (json.RawMessage, error) { + a.mu.Lock() + defer a.mu.Unlock() + args, _ := body.(map[string]any) + a.calls = append(a.calls, key) + a.args = append(a.args, args) + if args["collected"] == "true" && a.fail["collected"] { + return nil, fmt.Errorf("maximum payload exceeded") + } + var answer any + switch key { + case "seat:mesh-controller.artifacts": + answer = records() + case "seat:mesh-controller.collect": + if args["confirm"] == "true" { + answer = CollectAnswer{LetGo: []string{"artifact-store://app/server@" + d("m-old")}} + } else { + answer = CollectAnswer{DryRun: true, Eligible: 1, WouldLetGo: []string{"artifact-store://app/server@" + d("m-old")}} + } + default: + return nil, fmt.Errorf("no such verb %s", key) + } + raw, _ := json.Marshal(answer) + return json.Marshal(map[string]any{"ok": true, "output": "", "answer": json.RawMessage(raw)}) +} + +func tool(t *testing.T, a *asked, name string) func(map[string]any) (any, error) { + f := world() + door := f.door() + t.Cleanup(door.Close) + s := &Store{URL: door.URL, Container: "mesh-registry", Run: func(context.Context, string, ...string) Ran { return Ran{Stdout: f.listing()} }} + for _, tl := range Tools(s, Controller{Ask: a.ask}) { + if tl.Name == name { + return tl.Run + } + } + t.Fatalf("no tool %s", name) + return nil +} + +func TestADryRunNeverConfirms(t *testing.T) { + a := &asked{} + for _, args := range []map[string]any{{}, {"dry_run": true}, {"dry_run": true, "why": "tidy"}} { + if _, err := tool(t, a, "store_collect")(args); err != nil { + t.Fatal(err) + } + } + for _, args := range a.args { + if args["confirm"] != nil || args["why"] != nil { + t.Errorf("a dry run asked %v", args) + } + } +} + +func TestARealRunWithoutWhyIsRefusedAndAsksNothing(t *testing.T) { + a := &asked{} + if _, err := tool(t, a, "store_collect")(map[string]any{"dry_run": false}); err == nil { + t.Fatal("a real run without why was accepted") + } + if len(a.calls) != 0 { + t.Errorf("asked %v", a.calls) + } + out, err := tool(t, a, "store_collect")(map[string]any{"dry_run": false, "why": "the store is full", "most": float64(10)}) + if err != nil { + t.Fatal(err) + } + if a.args[0]["confirm"] != "true" || a.args[0]["why"] != "the store is full" || a.args[0]["most"] != "10" { + t.Errorf("asked %v", a.args[0]) + } + if out.(map[string]any)["dry_run"] != false { + t.Error("a real run answered as a dry run") + } +} + +func TestReferencesFallBackWhenCollectedCannotBeCarried(t *testing.T) { + a := &asked{fail: map[string]bool{"collected": true}} + out, err := tool(t, a, "store_references")(map[string]any{"repository": "app/server"}) + if err != nil { + t.Fatal(err) + } + m := out.(map[string]any) + if m["note"] == nil || m["manifests"] == nil { + t.Errorf("answer %v", m) + } + if _, err := tool(t, a, "store_references")(map[string]any{"repository": "nope"}); err == nil { + t.Error("an unknown repository was answered") + } +} + +func TestTheToolsServedAreTheToolsTheManifestNames(t *testing.T) { + raw, err := os.ReadFile("../../module.json") + if err != nil { + t.Fatal(err) + } + var m struct { + Tools []string `json:"tools"` + Invokes []string `json:"invokes"` + } + if err := json.Unmarshal(raw, &m); err != nil { + t.Fatal(err) + } + var served []string + for _, tl := range Tools(&Store{}, Controller{}) { + if !strings.HasPrefix(tl.Name, "store_") || tl.Description == "" || tl.Run == nil { + t.Errorf("tool %q", tl.Name) + } + served = append(served, tl.Name) + } + sort.Strings(served) + listed := append([]string{}, m.Tools...) + sort.Strings(listed) + if !reflect.DeepEqual(served, listed) { + t.Fatalf("served %v, manifest %v", served, listed) + } + want := []string{"seat:mesh-controller.artifacts", "seat:mesh-controller.collect"} + if !reflect.DeepEqual(m.Invokes, want) { + t.Errorf("invokes %v", m.Invokes) + } +} + +func TestASocketRefusalIsAskedAgainThroughSudo(t *testing.T) { + var ran []string + run := func(_ context.Context, name string, args ...string) Ran { + ran = append(ran, name) + if name == "docker" { + return Ran{Status: 1, Stderr: "permission denied while trying to connect to the Docker daemon socket at unix:///var/run/docker.sock"} + } + return Ran{Stdout: "#end\n"} + } + if _, err := docker(context.Background(), run, 1000, "exec", "x"); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(ran, []string{"docker", "sudo"}) { + t.Errorf("ran %v", ran) + } + if _, err := docker(context.Background(), func(context.Context, string, ...string) Ran { + return Ran{Status: 1, Stderr: "Error: No such container: mesh-registry"} + }, 1000, "exec"); err == nil || !strings.Contains(err.Error(), "not on this machine") { + t.Errorf("err %v", err) + } +} diff --git a/modules/distribution/cmd/store-tools/tools.go b/modules/distribution/cmd/store-tools/tools.go new file mode 100644 index 00000000..936ad784 --- /dev/null +++ b/modules/distribution/cmd/store-tools/tools.go @@ -0,0 +1,351 @@ +package main + +import ( + "context" + "fmt" + "math" + "sort" + "strconv" + "strings" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +// Tools are the store module's tools (novox/hq ADR 0251 §1). The first three only read; store_collect +// only reads too, unless it is a real run, and then it asks the controller, which decides and deletes. +func Tools(s *Store, c Controller) []stdio.Tool { + ctx := context.Background() + repoArg := map[string]any{"type": "string", "description": "one repository, as the store names it (/)"} + return []stdio.Tool{ + { + Name: "store_repositories", + Description: "Every repository the artifact store holds: its tags and what each names, how many manifests it holds " + + "(tagged or not — the mesh pins by digest, so most are untagged) and its size, the blobs its manifests mark. " + + "A blob two repositories share is counted in each, and shared_bytes says how much of a size that is. Read from " + + "the store's own files and its door; changes nothing. Replaces curl /v2/_catalog and /v2//tags/list. (r)", + Input: map[string]any{"repository": repoArg}, + Run: func(args map[string]any) (any, error) { + v, err := s.View(ctx) + if err != nil { + return nil, err + } + return Repositories(v, text(args, "repository")) + }, + }, + { + Name: "store_usage", + Description: "How large the artifact store is: every blob's bytes together, the largest repositories, the bytes " + + "more than one repository's manifests mark, and the bytes no manifest marks — what the store's nightly collector " + + "frees next. Read from the store's own files; changes nothing. Replaces du on the store's directory. (r)", + Input: map[string]any{"top": map[string]any{"type": "integer", "description": "how many of the largest repositories to name (default 10, at most 200)"}}, + Run: func(args map[string]any) (any, error) { + top, err := count(args, "top", 10, 1, 200) + if err != nil { + return nil, err + } + v, err := s.View(ctx) + if err != nil { + return nil, err + } + return Usage(v, top), nil + }, + }, + { + Name: "store_references", + Description: "What the controller's records say of each manifest the artifact store holds: kept (a definition " + + "names it, or one of the five most recent builds of a module the mesh holds), the holder of a kept archive, " + + "eligible (the mesh made it and keeps it for no reason), the holder of an eligible archive, let go yet present, " + + "a named document the controller keeps under a tag, or unrecorded — no record names it. Counted and sized by " + + "state and by repository; each manifest listed when one repository is asked. Then what the records keep that " + + "the store does not hold. Unrecorded manifests are never removed by any tool (novox/hq ADR 0189 §3). (r)", + Input: map[string]any{"repository": repoArg}, + Run: func(args map[string]any) (any, error) { + repo := text(args, "repository") + recs, err := c.Artifacts("") + if err != nil { + return nil, err + } + v, err := s.View(ctx) + if err != nil { + return nil, err + } + return References(v, recs, repo) + }, + }, + { + Name: "store_collect", + Description: "Collect what the mesh made and keeps for no reason (novox/hq ADR 0189, ADR 0251 §3). A dry run unless " + + "dry_run is false: it asks the controller's collect what it would let go of, and works out how many bytes the " + + "store's nightly collector would free once those are gone, beside what it frees tonight anyway. A real run needs " + + "why: it asks the controller's collect to hold every kept archive and let go of the eligible ones, recorded as a " + + "hand-act with the why; the bytes come back at the nightly collection. Never removes an unrecorded manifest, and " + + "never deletes through the store's door itself. (a)", + Input: map[string]any{ + "dry_run": map[string]any{"type": "boolean", "description": "true (the default): say what would go, change nothing"}, + "why": map[string]any{"type": "string", "description": "why — required for a real run, and recorded as the hand-act's reason"}, + "most": map[string]any{"type": "integer", "description": "the most artifacts one real run lets go of (the controller's default when absent)"}, + }, + Run: func(args map[string]any) (any, error) { + dry := true + if b, given := args["dry_run"].(bool); given { + dry = b + } else if t := text(args, "dry_run"); t == "false" { + dry = false + } + why := text(args, "why") + if !dry && why == "" { + return nil, fmt.Errorf("a real collection needs why: nothing was asked of the controller") + } + most, err := count(args, "most", 0, 1, 5000) + if err != nil { + return nil, err + } + answer, err := c.Collect(why, !dry, most) + if err != nil { + return nil, err + } + v, err := s.View(ctx) + if err != nil { + return map[string]any{"controller": answer, "bytes_not_worked_out": err.Error()}, nil + } + return Collected(v, answer, dry), nil + }, + }, + } +} + +// View reads the store now. +func (s *Store) View(ctx context.Context) (*View, error) { + l, err := s.List(ctx) + if err != nil { + return nil, err + } + m, unread := s.Manifests(ctx, l) + return &View{L: l, M: m, Unread: unread}, nil +} + +// unreadNote says what manifests that could not be read do to the numbers, when there are any. +func unreadNote(v *View) map[string]any { + if len(v.Unread) == 0 { + return nil + } + keys := make([]string, 0, len(v.Unread)) + for k := range v.Unread { + keys = append(keys, k) + } + sort.Strings(keys) + first := keys[0] + ": " + v.Unread[keys[0]] + return map[string]any{ + "manifests": len(keys), + "first": first, + "means": "what these mark beyond themselves is not known, so sizes may be low and the bytes said to be " + + "freed may be high; ask again — what was read is kept, so the next call reads only the rest", + } +} + +// Repositories answers store_repositories. +func Repositories(v *View, only string) (map[string]any, error) { + shared := v.Shared() + var repos []map[string]any + for _, name := range v.L.RepoNames() { + if only != "" && name != only { + continue + } + r := v.L.Repos[name] + blobs := v.RepoBlobs(name) + sharedHere := map[string]bool{} + for b := range blobs { + if shared[b] { + sharedHere[b] = true + } + } + repos = append(repos, map[string]any{ + "repository": name, "tags": tagsOf(r), "manifests": len(r.Revisions), + "bytes": v.Bytes(blobs), "size": human(v.Bytes(blobs)), "shared_bytes": v.Bytes(sharedHere), + }) + } + if only != "" && len(repos) == 0 { + return nil, fmt.Errorf("the store holds no repository %q", only) + } + out := map[string]any{"count": len(repos), "repositories": repos, + "note": "a repository's size is the blobs its manifests mark; a blob shared with another repository is counted in each (shared_bytes)"} + if u := unreadNote(v); u != nil { + out["unread"] = u + } + return out, nil +} + +// Usage answers store_usage. +func Usage(v *View, top int) map[string]any { + type sized struct { + name string + bytes int64 + n int + } + var all []sized + manifests := 0 + for _, name := range v.L.RepoNames() { + n := len(v.L.Repos[name].Revisions) + manifests += n + all = append(all, sized{name, v.Bytes(v.RepoBlobs(name)), n}) + } + sort.SliceStable(all, func(i, j int) bool { return all[i].bytes > all[j].bytes }) + var largest []map[string]any + for i, s := range all { + if i >= top { + break + } + largest = append(largest, map[string]any{"repository": s.name, "bytes": s.bytes, "size": human(s.bytes), "manifests": s.n}) + } + total := v.Total() + sharedBytes := v.Bytes(v.Shared()) + freed, freedBlobs := v.Unmarked(nil) + out := map[string]any{ + "total_bytes": total, "total": human(total), "blobs": len(v.L.Blobs), "manifests": manifests, + "repositories": len(v.L.Repos), "largest": largest, + "shared_bytes": sharedBytes, "collector_frees_next_bytes": freed, "collector_frees_next_blobs": freedBlobs, + "said": []string{ + fmt.Sprintf("the store holds %s in %d blobs, %d manifests in %d repositories", human(total), len(v.L.Blobs), manifests, len(v.L.Repos)), + fmt.Sprintf("%s is marked by more than one repository's manifests", human(sharedBytes)), + fmt.Sprintf("the nightly collector frees %s in %d blobs no manifest marks", human(freed), freedBlobs), + }, + } + if u := unreadNote(v); u != nil { + out["unread"] = u + } + return out +} + +// References answers store_references. +func References(v *View, recs *Records, only string) (map[string]any, error) { + if only != "" && v.L.Repos[only] == nil { + return nil, fmt.Errorf("the store holds no repository %q", only) + } + entries := Classify(v, recs) + byRepo := map[string][]Entry{} + for _, e := range entries { + byRepo[e.Repository] = append(byRepo[e.Repository], e) + } + var repos []map[string]any + for _, name := range v.L.RepoNames() { + if only != "" && name != only { + continue + } + repos = append(repos, map[string]any{"repository": name, "states": Summarise(v, byRepo[name], entries)}) + } + whole := Summarise(v, entries, entries) + said := []string{} + for _, state := range []string{StateKept, StateHolderKept, StateEligible, StateHolderEligible, StateCollectedPresent, StateNamedDocument, StateUnrecorded} { + if s := whole[state]; s != nil { + said = append(said, fmt.Sprintf("%s: %d manifests marking %s, %s of it marked by nothing else", state, s.Manifests, human(s.Bytes), human(s.OnlyTheirBytes))) + } + } + missing := Missing(v, recs) + if len(missing) > 0 { + said = append(said, fmt.Sprintf("%d references the records keep are not in the store", len(missing))) + } + out := map[string]any{ + "store": recs.Store, "kept_builds": recs.KeptBuilds, "records": recs.Counts, + "states": whole, "repositories": repos, "missing": nonNil(missing), "said": said, + "removable": collectionRemovesThese, + } + if only != "" { + out["manifests"] = byRepo[only] + } + if recs.Note != "" { + out["note"] = recs.Note + } + if u := unreadNote(v); u != nil { + out["unread"] = u + } + return out, nil +} + +// Collected answers store_collect: the controller's answer, and what the store's collector frees. +func Collected(v *View, a *CollectAnswer, dry bool) map[string]any { + refs := a.LetGo + if dry { + refs = a.WouldLetGo + } + gone := LetGoKeys(v, refs) + now, nowBlobs := v.Unmarked(nil) + after, afterBlobs := v.Unmarked(gone) + said := []string{} + if dry { + said = append(said, fmt.Sprintf("dry run: the controller would let go of %d of %d eligible artifacts; nothing was changed", len(a.WouldLetGo), a.Eligible)) + } else { + said = append(said, fmt.Sprintf("the controller let go of %d artifacts (%d left, %d skipped)", len(a.LetGo), a.Left, a.Skipped)) + if a.Stopped != "" { + said = append(said, "it stopped: "+a.Stopped) + } + } + said = append(said, + fmt.Sprintf("the nightly collector frees %s tonight as the store is now", human(now)), + fmt.Sprintf("and %s once those %d manifests are gone (%s more)", human(after), len(gone), human(after-now))) + if dry && a.Eligible > len(a.WouldLetGo) { + said = append(said, fmt.Sprintf("the controller listed %d of %d eligible: the bytes cover only those listed", len(a.WouldLetGo), a.Eligible)) + } + out := map[string]any{ + "dry_run": dry, "controller": a, + "collector_frees_now_bytes": now, "collector_frees_now_blobs": nowBlobs, + "collector_frees_after_bytes": after, "collector_frees_after_blobs": afterBlobs, + "manifests_gone": len(gone), "said": said, + "note": "a manifest let go of frees no bytes until the store's nightly collector runs with the store held still (ADR 0189 §4)", + } + if u := unreadNote(v); u != nil { + out["unread"] = u + } + return out +} + +func nonNil(s []string) []string { + if s == nil { + return []string{} + } + return s +} + +func human(n int64) string { + switch { + case n >= 1<<30: + return fmt.Sprintf("%.1f GiB", float64(n)/(1<<30)) + case n >= 1<<20: + return fmt.Sprintf("%.1f MiB", float64(n)/(1<<20)) + case n >= 1<<10: + return fmt.Sprintf("%.1f KiB", float64(n)/(1<<10)) + } + return fmt.Sprintf("%d B", n) +} + +func text(args map[string]any, key string) string { + s, _ := args[key].(string) + return strings.TrimSpace(s) +} + +// count is a whole-number argument, defaulted, refused below least, held at most. +func count(args map[string]any, key string, fallback, least, most int) (int, error) { + v, given := args[key] + if !given || v == nil { + return fallback, nil + } + var n int + switch x := v.(type) { + case float64: + if x != math.Trunc(x) { + return 0, fmt.Errorf("%s must be a whole number, not %v", key, x) + } + n = int(x) + case string: + i, err := strconv.Atoi(strings.TrimSpace(x)) + if err != nil { + return 0, fmt.Errorf("%s must be a whole number, not %q", key, x) + } + n = i + default: + return 0, fmt.Errorf("%s must be a whole number", key) + } + if n < least { + return 0, fmt.Errorf("%s must be at least %d", key, least) + } + return min(n, most), nil +} diff --git a/modules/distribution/cmd/store-tools/view.go b/modules/distribution/cmd/store-tools/view.go new file mode 100644 index 00000000..9d853a0d --- /dev/null +++ b/modules/distribution/cmd/store-tools/view.go @@ -0,0 +1,129 @@ +package main + +import ( + "sort" +) + +// View is the store at one moment: its files, and every manifest's content that could be read. +type View struct { + L *Listing + // M is each manifest read, by Key. + M map[string]*Manifest + // Unread is each manifest that could not be read, by Key, with why. + Unread map[string]string +} + +// Marks are the blobs one manifest keeps from the store's collector: its own content, its +// configuration and its layers, and for an index the manifests it lists. A manifest that could not be +// read marks only itself here, and the answers that rest on it say so. +func (v *View) Marks(repo, digest string) []string { + out := []string{digest} + if m := v.M[Key(repo, digest)]; m != nil { + out = append(out, m.Refs...) + out = append(out, m.Children...) + } + return out +} + +// Marked is every blob some manifest marks, leaving out the manifests named in without — the +// collector's own mark, and the mark it would make once those are gone. +func (v *View) Marked(without map[string]bool) map[string]bool { + marked := map[string]bool{} + for name, r := range v.L.Repos { + for d := range r.Revisions { + if without[Key(name, d)] { + continue + } + for _, b := range v.Marks(name, d) { + marked[b] = true + } + } + } + return marked +} + +// Bytes is the size of these blobs, as the store holds them; a blob the store does not hold counts +// nothing. +func (v *View) Bytes(blobs map[string]bool) int64 { + var n int64 + for b := range blobs { + n += v.L.Blobs[b] + } + return n +} + +// Total is the size of every blob the store holds. +func (v *View) Total() int64 { + var n int64 + for _, s := range v.L.Blobs { + n += s + } + return n +} + +// Unmarked is the bytes and count of the blobs no manifest marks, once the manifests in without are +// gone: what the nightly collector frees. +func (v *View) Unmarked(without map[string]bool) (int64, int) { + marked := v.Marked(without) + var n int64 + count := 0 + for b, s := range v.L.Blobs { + if !marked[b] { + n += s + count++ + } + } + return n, count +} + +// RepoBlobs are the blobs one repository's manifests mark. +func (v *View) RepoBlobs(name string) map[string]bool { + out := map[string]bool{} + if r := v.L.Repos[name]; r != nil { + for d := range r.Revisions { + for _, b := range v.Marks(name, d) { + out[b] = true + } + } + } + return out +} + +// Shared are the blobs that the manifests of more than one repository mark. +func (v *View) Shared() map[string]bool { + seen := map[string]string{} + shared := map[string]bool{} + for _, name := range v.L.RepoNames() { + for b := range v.RepoBlobs(name) { + if first, ok := seen[b]; ok && first != name { + shared[b] = true + } else { + seen[b] = name + } + } + } + return shared +} + +// tagsOf is a repository's tags, in order, with what each names. +func tagsOf(r *Repo) []map[string]string { + names := make([]string, 0, len(r.Tags)) + for t := range r.Tags { + names = append(names, t) + } + sort.Strings(names) + out := make([]map[string]string, 0, len(names)) + for _, t := range names { + out = append(out, map[string]string{"tag": t, "digest": r.Tags[t]}) + } + return out +} + +// taggedDigests are the manifests some tag names, in one repository. +func taggedDigests(r *Repo) map[string]bool { + out := map[string]bool{} + for _, d := range r.Tags { + out[d] = true + } + return out +} diff --git a/modules/distribution/go.mod b/modules/distribution/go.mod new file mode 100644 index 00000000..dc21c025 --- /dev/null +++ b/modules/distribution/go.mod @@ -0,0 +1,5 @@ +module distribution + +go 1.22 + +require git.novox.be/novox/mesh-sdk/go v0.1.7 diff --git a/modules/distribution/go.sum b/modules/distribution/go.sum new file mode 100644 index 00000000..b474419e --- /dev/null +++ b/modules/distribution/go.sum @@ -0,0 +1,2 @@ +git.novox.be/novox/mesh-sdk/go v0.1.7 h1:C0sTQmtTiyYH7bnqZb7PusXnqA37gKuT7Nqjn9gG47w= +git.novox.be/novox/mesh-sdk/go v0.1.7/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= diff --git a/modules/distribution/index.ts b/modules/distribution/index.ts deleted file mode 100644 index 949f1f98..00000000 --- a/modules/distribution/index.ts +++ /dev/null @@ -1,51 +0,0 @@ -// registry's events. The tool runtime imports this once the broker is bound. -// -// Emits (novox/hq ADR 0041/0042): -// module.registry.image.pushed — a new image (repo:tag) was published to the registry -// -// This is a genuinely useful signal: a build finished and its image is now pullable, so anything -// on the mesh that redeploys, mirrors or announces releases can react without polling the registry -// itself. It is discovered by diffing the catalog and each repo's tags — the registry has no push -// webhook of its own, so the module watches for it. -// -// The polling is deliberately unhurried: a new image a minute late is still the event, whereas -// hammering the registry's catalog for immediacy nobody asked for is not. - -import { emit } from "@novox/mesh-sdk/events"; -import { RegistryClient } from "./client.js"; - -const registry = RegistryClient.fromEnv(); - -// Every repo:tag we have already accounted for. Primed silently on the first look so a registry -// that was already full when this started does not announce its whole history as freshly pushed. -const seen = new Set(); -let primed = false; - -async function pollCatalog(): Promise { - const repos = await registry.listRepositories(); - for (const repo of repos) { - let tags: string[]; - try { - tags = await registry.listTags(repo); - } catch { - continue; // a repo can vanish between catalog and tag read — skip it, catch it next tick - } - for (const tag of tags) { - const id = `${repo}:${tag}`; - if (!seen.has(id)) { - if (primed) await emit("image.pushed", { repo, tag }); - seen.add(id); - } - } - } - primed = true; -} - -const tick = (fn: () => Promise, everyMs: number): void => { - const run = (): void => void fn().catch((err) => console.error(`[registry] ${err}`)); - setInterval(run, everyMs); - run(); -}; -tick(pollCatalog, 60_000); - -console.log("[registry] watching the catalog for newly pushed images"); diff --git a/modules/distribution/module.json b/modules/distribution/module.json index e8eb32d3..f97e9a98 100644 --- a/modules/distribution/module.json +++ b/modules/distribution/module.json @@ -16,8 +16,15 @@ "capabilities": [ "container-runtime" ], - "emits": [ - "image.pushed" + "invokes": [ + "seat:mesh-controller.artifacts", + "seat:mesh-controller.collect" + ], + "tools": [ + "store_repositories", + "store_usage", + "store_references", + "store_collect" ], "own-secrets": { "broker": "/var/lib/mesh/registry/broker" @@ -93,5 +100,24 @@ "store" ] } - ] + ], + "build": { + "artifacts": [ + { + "name": "tools-go", + "kind": "bundle", + "language": "go", + "system": "arch", + "from": "cmd/store-tools", + "binary": "store-tools", + "loads": [ + "store-tools" + ], + "env": { + "MESH_STORE_URL": "http://127.0.0.1:${port:5000}", + "MESH_STORE_CONTAINER": "mesh-registry" + } + } + ] + } } diff --git a/modules/distribution/package.json b/modules/distribution/package.json deleted file mode 100644 index 8914148f..00000000 --- a/modules/distribution/package.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "name": "@novox/module-registry", - "version": "0.1.0", - "description": "registry — private Docker image registry. Its API client, tools and events live here (novox/hq ADR 0039).", - "type": "module", - "private": true, - "dependencies": { - "@novox/mesh-sdk": "^0.1.0" - }, - "devDependencies": { - "@types/node": "^22.0.0", - "typescript": "^5.6.0" - } -} diff --git a/modules/distribution/tools/index.ts b/modules/distribution/tools/index.ts deleted file mode 100644 index 62af1159..00000000 --- a/modules/distribution/tools/index.ts +++ /dev/null @@ -1,59 +0,0 @@ -// registry's tools — moved here from the shared sdk (novox/hq ADR 0039), importing registry's own -// client. They return structured data; the mesh serves them through the sdk's tool harness. - -import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; -import { RegistryClient } from "../client.js"; - -export function getRegistryTools(registry: RegistryClient): ToolDefinition[] { - return [ - { - name: "registry_list", - description: "List every repository in the Docker registry (the catalog).", - input: {}, - run: async () => { - const repositories = await registry.listRepositories(); - return { count: repositories.length, repositories }; - }, - }, - { - name: "registry_tags", - description: "List the tags of one repository in the Docker registry.", - input: { repo: { type: "string", description: "the repository name, e.g. 'novox/mesh'" } }, - run: async (args) => { - const repo = String(args.repo); - const tags = await registry.listTags(repo); - return { repo, count: tags.length, tags }; - }, - }, - { - name: "registry_delete_image", - description: - "Delete an image tag from the registry (DESTRUCTIVE). Removes the manifest; storage is reclaimed by garbage collection later. Requires confirm: true.", - input: { - repo: { type: "string", description: "the repository name, e.g. 'novox/mesh'" }, - tag: { type: "string", description: "the tag to delete, e.g. 'latest'" }, - confirm: { type: "boolean", description: "must be true to actually delete" }, - }, - run: async (args) => { - const repo = String(args.repo); - const tag = String(args.tag); - if (args.confirm !== true) { - return { deleted: false, reason: "confirm must be true to delete an image" }; - } - const digest = await registry.getManifestDigest(repo, tag); - await registry.deleteManifest(repo, digest); - return { deleted: true, repo, tag, digest, note: "run registry garbage collection to reclaim storage" }; - }, - }, - ]; -} - -// The tools exist only when a registry URL is configured; otherwise registry contributes none -// rather than failing the whole runtime. -registerModuleTools("registry", (env) => { - try { - return getRegistryTools(RegistryClient.fromEnv(env)); - } catch { - return []; - } -}); diff --git a/modules/distribution/tsconfig.json b/modules/distribution/tsconfig.json deleted file mode 100644 index 36778592..00000000 --- a/modules/distribution/tsconfig.json +++ /dev/null @@ -1,12 +0,0 @@ -{ - "compilerOptions": { - "target": "ES2022", - "module": "NodeNext", - "moduleResolution": "NodeNext", - "strict": true, - "esModuleInterop": true, - "skipLibCheck": true, - "noEmit": true - }, - "include": ["client.ts", "index.ts", "tools/index.ts"] -}