distribution: say what the store holds and what the records name, so collection is no longer blind (hq ADR 0251)
mesh/merge-gate fail: builds distribution → novox; no bus step; a manifest the change touches fails the module check: modules/distribution/module.json: thi…
mesh/repo-check fail: its merge-check.sh failed: long-running resources without health: 69
mesh/delivery superseded: a newer head of the same pull request
mesh/merge-gate fail: builds distribution → novox; no bus step; a manifest the change touches fails the module check: modules/distribution/module.json: thi…
mesh/repo-check fail: its merge-check.sh failed: long-running resources without health: 69
mesh/delivery superseded: a newer head of the same pull request
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.
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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://<repo>@sha256:…` (an image) and
|
||||
// `artifact-store://<repo>/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 "", "", ""
|
||||
}
|
||||
@@ -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 ""
|
||||
}
|
||||
@@ -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/<hex>/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/<repo>/_manifests/revisions/sha256/<hex>/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/<repo>/_manifests/tags/<tag>/current/link:sha256:<hex>
|
||||
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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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 (<module>/<artifact>)"}
|
||||
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/<name>/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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user