From 5f22ebbe9724ef8f9d2e021d30dcab2b8d8ff28d Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 11:59:39 +0200 Subject: [PATCH] distribution: say what the store holds and what the records name, so collection is no longer blind (hq ADR 0251) The store had no working tools: its TypeScript client was never built, and nothing could count what the store holds that no record names. A Go bundle lists the store's files through its own container, reads each manifest through its door, and sets that beside the controller's records. Collection is asked of the controller, which decides and records; the store's tools never delete. The image.pushed event was declared and never emitted, and nothing consumes it, so it goes. --- modules/distribution/README.md | 76 ++++ modules/distribution/client.ts | 81 ---- .../distribution/cmd/store-tools/classify.go | 203 +++++++++ modules/distribution/cmd/store-tools/main.go | 31 ++ .../distribution/cmd/store-tools/records.go | 209 +++++++++ .../distribution/cmd/store-tools/runner.go | 96 ++++ modules/distribution/cmd/store-tools/store.go | 313 +++++++++++++ .../cmd/store-tools/store_test.go | 421 ++++++++++++++++++ modules/distribution/cmd/store-tools/tools.go | 351 +++++++++++++++ modules/distribution/cmd/store-tools/view.go | 129 ++++++ modules/distribution/go.mod | 5 + modules/distribution/go.sum | 2 + modules/distribution/index.ts | 51 --- modules/distribution/module.json | 32 +- modules/distribution/package.json | 14 - modules/distribution/tools/index.ts | 59 --- modules/distribution/tsconfig.json | 12 - 17 files changed, 1865 insertions(+), 220 deletions(-) create mode 100644 modules/distribution/README.md delete mode 100644 modules/distribution/client.ts create mode 100644 modules/distribution/cmd/store-tools/classify.go create mode 100644 modules/distribution/cmd/store-tools/main.go create mode 100644 modules/distribution/cmd/store-tools/records.go create mode 100644 modules/distribution/cmd/store-tools/runner.go create mode 100644 modules/distribution/cmd/store-tools/store.go create mode 100644 modules/distribution/cmd/store-tools/store_test.go create mode 100644 modules/distribution/cmd/store-tools/tools.go create mode 100644 modules/distribution/cmd/store-tools/view.go create mode 100644 modules/distribution/go.mod create mode 100644 modules/distribution/go.sum delete mode 100644 modules/distribution/index.ts delete mode 100644 modules/distribution/package.json delete mode 100644 modules/distribution/tools/index.ts delete mode 100644 modules/distribution/tsconfig.json 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"] -}