From 7d033ad9f6998772cfa5fae56d0d6269ed0685c7 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 30 Aug 2026 03:36:04 +0200 Subject: [PATCH] Publish to the registry, and a command that builds a repository MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One store, and it is the registry the bootstrap already pulls from. An OCI registry is a content-addressed blob store that also understands images: PUT a blob and it is retrievable at /v2//blobs/sha256:… for ever, by digest. An archive is a content-addressed blob. A second store beside it was considered and is the right answer for objects that are mutable, need per-reader access, or are not build output — somebody's uploads, a backup, a thing with a lifecycle. None of that describes a digest-pinned archive, and running a second service to hold one kind of immutable blob is two things to run, two to back up, and two ways for an artifact to be missing. Overturnable by reading: the manifest carries a URL and a digest, and neither says what served it. `build ` clones, reads module.json, builds what it declares, publishes, and records the manifest with the commit it came from. It is a command rather than something the control plane does on its own, because building runs things on a machine and what the control plane may send a machine is bounded by the declaration language. This is the shape the builder module takes when it is given work over the broker. Proven end to end on a real repository and a real registry: a shell module with a package, a user and a dotfile archive built, published, fetched back at the digest it declared, rebuilt to the same digest, and its manifest accepted by the host's own parser — including `user` and `archive`, which did not exist this morning. A tag is never accepted as a pin, and a blob already stored is not sent again — it is named by its content, so re-uploading asks the registry to store what it already has under the name it already has. --- cmd/mesh-control/main.go | 67 ++++++++++++ internal/builder/builder.go | 5 + internal/builder/registry.go | 151 ++++++++++++++++++++++++++ internal/builder/registry_test.go | 174 ++++++++++++++++++++++++++++++ 4 files changed, 397 insertions(+) create mode 100644 internal/builder/registry.go create mode 100644 internal/builder/registry_test.go diff --git a/cmd/mesh-control/main.go b/cmd/mesh-control/main.go index 977dc7d..8edc035 100644 --- a/cmd/mesh-control/main.go +++ b/cmd/mesh-control/main.go @@ -20,6 +20,7 @@ import ( "time" "github.com/novox/mesh-control/internal/broker" + "github.com/novox/mesh-control/internal/builder" "github.com/novox/mesh-control/internal/catalogue" "github.com/novox/mesh-control/internal/identity" "github.com/novox/mesh-control/internal/inventory" @@ -63,6 +64,8 @@ func run() error { defer stop() switch args[0] { + case "build": + return buildCommand(ctx, args[1:]) case "pin": return pinCommand(ctx, args[1:], true) case "unpin": @@ -131,6 +134,7 @@ func usage() { settings set what a module's config should say, for the whole mesh settings set --node ...or for one machine settings clear [--node ] take a layer away + build [--ref R] build a module from its source and record it pin which node this one gets a provision from unpin put that question back plan [--files|--json] what that node would run, and why @@ -1625,3 +1629,66 @@ func pinCommand(ctx context.Context, args []string, setting bool) error { fmt.Printf(" run `push %s` to send it\n", args[0]) return nil } + +// buildCommand builds a module from its source and records what came out. +// +// **Run where there is a container runtime**, which is why it is a command rather than something +// the control plane does on its own: building needs to run things on a machine, and what the +// control plane may send a machine is bounded by the declaration language. This is the shape the +// builder module will take when it is given work over the broker; today a person runs it, and the +// mesh records the result the same way either way. +func buildCommand(ctx context.Context, args []string) error { + set := flag.NewFlagSet("build", flag.ContinueOnError) + ref := set.String("ref", "", "the branch, tag or commit to build") + registry := set.String("registry", os.Getenv("MESH_REGISTRY"), + "host:port of the registry to publish to") + workspace := set.String("workspace", os.TempDir(), "where to clone and build") + dryRun := set.Bool("dry-run", false, "build and print the manifest, recording nothing") + positionals, err := parseAround(set, args) + if err != nil { + return err + } + if len(positionals) != 1 { + return errors.New("build [--ref R] [--registry host:port]") + } + if strings.TrimSpace(*registry) == "" { + return errors.New( + "no --registry and no MESH_REGISTRY: a built artifact nobody can fetch is not built") + } + + publisher := builder.Registry{Address: *registry, Run: builder.Command} + result, err := builder.Build(ctx, builder.Command, publisher, positionals[0], *ref, *workspace) + if err != nil { + return err + } + + for _, made := range result.Built { + fmt.Printf(" %-12s %s %s\n", made.Name, made.Kind, made.Reference) + } + if *dryRun { + body, err := json.MarshalIndent(result.Manifest, "", " ") + if err != nil { + return err + } + fmt.Println(string(body)) + return nil + } + + inv, err := openInventory(ctx) + if err != nil { + return err + } + defer inv.Close() + + // Recorded with where it came from, so "is this current?" is answerable without building it + // again (novox/hq ADR 0009). + if err := inv.RegisterModule(ctx, result.Manifest, inventory.Source{ + Repository: positionals[0], Ref: *ref, BuiltFrom: result.Commit, Head: result.Commit, + }); err != nil { + return err + } + fmt.Printf("\n%s %s, built from %s\n", + result.Manifest.Module, result.Manifest.Version, short(result.Commit)) + fmt.Printf(" run `assign %s` to put it somewhere\n", result.Manifest.Module) + return nil +} diff --git a/internal/builder/builder.go b/internal/builder/builder.go index 3ca6495..a558711 100644 --- a/internal/builder/builder.go +++ b/internal/builder/builder.go @@ -59,6 +59,11 @@ type Result struct { func Build(ctx context.Context, run Runner, publish Publisher, repository, ref, workspace string) (Result, error) { + // Made rather than required. A builder that fails because the directory it was told to work + // in does not exist is a builder that needs a setup step nobody documented. + if err := os.MkdirAll(workspace, 0o755); err != nil { + return Result{}, err + } tree := filepath.Join(workspace, "source") if err := os.RemoveAll(tree); err != nil { return Result{}, err diff --git a/internal/builder/registry.go b/internal/builder/registry.go new file mode 100644 index 0000000..481ce53 --- /dev/null +++ b/internal/builder/registry.go @@ -0,0 +1,151 @@ +package builder + +import ( + "bytes" + "context" + "fmt" + "io" + "net/http" + "strings" +) + +// Where built artifacts go. +// +// **One store, and it is the registry the bootstrap already pulls from.** An OCI registry is a +// content-addressed blob store that happens to also understand images: `PUT` a blob and it is +// retrievable at `/v2//blobs/sha256:…` for ever, by digest, over plain HTTP. An archive is +// a content-addressed blob. So it goes there. +// +// The alternative considered was a second store beside it — S3-shaped, buckets, signed URLs. It +// is the right answer for objects that are *mutable*, or need per-reader access, or are not build +// output: somebody's uploads, a backup, a thing with a lifecycle. None of that describes a +// digest-pinned archive, and standing up a second service to hold one kind of immutable blob +// means two things to run, two things to back up and two ways for an artifact to be missing. +// +// **This is a decision that can be overturned by reading**: if something needs an object store +// for reasons other than build artifacts, it should have one, and archives can move to it without +// anything else changing — the manifest carries a URL and a digest, and neither says what served +// it. + +// Registry publishes to an OCI registry over plain HTTP. +// +// Plain HTTP because the registry the mesh runs is reached over the private network, which is +// already the encrypted, authenticated thing. A second layer inside it would be certificates to +// issue and rotate for no property the first does not have. +type Registry struct { + // Address is host:port, as a machine will fetch from. + Address string + // Run is how docker is invoked, so a test does not need one. + Run Runner + // HTTP is the client used for blobs. + HTTP *http.Client +} + +// PublishImage pushes a locally built image and returns a reference pinned by digest. +// +// The digest is read back from the registry's own answer rather than computed here. What matters +// is what the registry will serve for that reference, and only it can say. +func (r Registry) PublishImage(ctx context.Context, localTag, repository string) (string, error) { + remote := r.Address + "/" + repository + if _, err := r.Run(ctx, "", "docker", "tag", localTag, remote); err != nil { + return "", err + } + if _, err := r.Run(ctx, "", "docker", "push", remote); err != nil { + return "", err + } + out, err := r.Run(ctx, "", "docker", "inspect", "--format", "{{index .RepoDigests 0}}", remote) + if err != nil { + return "", fmt.Errorf("pushed %s and cannot read back what it was pinned as: %w", remote, err) + } + pinned := strings.TrimSpace(out) + if !strings.Contains(pinned, "@sha256:") { + // A tag is not a pin. It can be made to point at something else, and this file is applied + // on machines with no mesh to ask about anything (novox/hq ADR 0006). + return "", fmt.Errorf("%s came back as %q, which is not pinned by digest", remote, pinned) + } + return pinned, nil +} + +// PublishArchive stores bytes as a blob and returns where to fetch them from. +// +// Two steps, which is the registry's own protocol: ask for somewhere to put it, then put it there +// naming the digest. The registry verifies the digest itself, so a blob that arrived corrupted is +// refused by the thing storing it rather than by the machine unpacking it a week later. +func (r Registry) PublishArchive(ctx context.Context, repository string, body []byte, digest string) (string, error) { + base := "http://" + r.Address + "/v2/" + repository + final := base + "/blobs/" + digest + + // Already there. Blobs are immutable and named by their content, so this is not an + // optimisation — re-uploading would be asking the registry to store what it already has under + // the name it already has. + if there, err := r.has(ctx, final); err != nil { + return "", err + } else if there { + return final, nil + } + + start, err := http.NewRequestWithContext(ctx, http.MethodPost, base+"/blobs/uploads/", nil) + if err != nil { + return "", err + } + begun, err := r.client().Do(start) + if err != nil { + return "", fmt.Errorf("cannot start an upload to %s: %w", base, err) + } + defer begun.Body.Close() + if begun.StatusCode != http.StatusAccepted { + return "", fmt.Errorf("%s answered %s when asked where to put a blob", base, begun.Status) + } + where := begun.Header.Get("Location") + if where == "" { + return "", fmt.Errorf("%s accepted an upload and said nowhere to put it", base) + } + if strings.HasPrefix(where, "/") { + where = "http://" + r.Address + where + } + + put, err := http.NewRequestWithContext(ctx, http.MethodPut, + where+separator(where)+"digest="+digest, bytes.NewReader(body)) + if err != nil { + return "", err + } + put.Header.Set("Content-Type", "application/octet-stream") + done, err := r.client().Do(put) + if err != nil { + return "", err + } + defer done.Body.Close() + if done.StatusCode != http.StatusCreated { + said, _ := io.ReadAll(io.LimitReader(done.Body, 4096)) + return "", fmt.Errorf("%s refused the blob: %s %s", base, done.Status, strings.TrimSpace(string(said))) + } + return final, nil +} + +func (r Registry) has(ctx context.Context, url string) (bool, error) { + request, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil) + if err != nil { + return false, err + } + response, err := r.client().Do(request) + if err != nil { + return false, fmt.Errorf("cannot reach the registry at %s: %w", r.Address, err) + } + defer response.Body.Close() + return response.StatusCode == http.StatusOK, nil +} + +func (r Registry) client() *http.Client { + if r.HTTP != nil { + return r.HTTP + } + return http.DefaultClient +} + +// separator is whether the upload location already carries a query. +func separator(where string) string { + if strings.Contains(where, "?") { + return "&" + } + return "?" +} diff --git a/internal/builder/registry_test.go b/internal/builder/registry_test.go new file mode 100644 index 0000000..20d4f00 --- /dev/null +++ b/internal/builder/registry_test.go @@ -0,0 +1,174 @@ +package builder + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" +) + +// An OCI registry as a content-addressed blob store, which is what it is. +// +// Against a stand-in rather than a real registry because what is under test is this side of the +// protocol — that the two steps happen in order, that the digest travels, that a blob already +// there is not sent again, and that a tag is never accepted as a pin. + +type fakeRegistry struct { + blobs map[string][]byte + uploads int + location string +} + +func (f *fakeRegistry) serve(t *testing.T) *httptest.Server { + t.Helper() + if f.blobs == nil { + f.blobs = map[string][]byte{} + } + mux := http.NewServeMux() + server := httptest.NewServer(mux) + mux.HandleFunc("/v2/", func(w http.ResponseWriter, r *http.Request) { + switch { + case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/sha256:"): + digest := r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:] + if _, ok := f.blobs[digest]; ok { + w.WriteHeader(http.StatusOK) + return + } + w.WriteHeader(http.StatusNotFound) + case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/blobs/uploads/"): + f.uploads++ + where := f.location + if where == "" { + where = "/v2/upload/" + hex.EncodeToString([]byte("session")) + } + w.Header().Set("Location", where) + w.WriteHeader(http.StatusAccepted) + case r.Method == http.MethodPut: + body, _ := io.ReadAll(r.Body) + digest := r.URL.Query().Get("digest") + sum := sha256.Sum256(body) + if digest != "sha256:"+hex.EncodeToString(sum[:]) { + // A registry verifies what it is given, which is why a corrupt blob is refused + // here rather than by a machine unpacking it a week later. + w.WriteHeader(http.StatusBadRequest) + return + } + f.blobs[digest] = body + w.WriteHeader(http.StatusCreated) + default: + w.WriteHeader(http.StatusNotFound) + } + }) + t.Cleanup(server.Close) + return server +} + +func registryFor(t *testing.T, f *fakeRegistry) Registry { + t.Helper() + server := f.serve(t) + return Registry{Address: strings.TrimPrefix(server.URL, "http://")} +} + +func TestAnArchiveIsStoredAndFetchableByItsDigest(t *testing.T) { + f := &fakeRegistry{} + r := registryFor(t, f) + body := []byte("a theme") + sum := sha256.Sum256(body) + digest := "sha256:" + hex.EncodeToString(sum[:]) + + where, err := r.PublishArchive(context.Background(), "shell/config", body, digest) + if err != nil { + t.Fatal(err) + } + if !strings.HasSuffix(where, "/blobs/"+digest) { + t.Fatalf("what a machine is told to fetch is %q, which does not name the digest", where) + } + if string(f.blobs[digest]) != "a theme" { + t.Fatalf("the registry holds %q", f.blobs[digest]) + } +} + +func TestABlobAlreadyThereIsNotSentAgain(t *testing.T) { + // Not an optimisation: blobs are named by their content, so re-uploading is asking the + // registry to store what it already has under the name it already has. + f := &fakeRegistry{} + r := registryFor(t, f) + body := []byte("a theme") + sum := sha256.Sum256(body) + digest := "sha256:" + hex.EncodeToString(sum[:]) + + if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil { + t.Fatal(err) + } + if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil { + t.Fatal(err) + } + if f.uploads != 1 { + t.Fatalf("the blob was uploaded %d times", f.uploads) + } +} + +func TestAnAbsoluteUploadLocationIsFollowed(t *testing.T) { + // A registry may answer with a full URL or with a path. Both happen in the wild, and a client + // that handles one produces a confusing failure against the other. + f := &fakeRegistry{} + server := f.serve(t) + f.location = server.URL + "/v2/upload/session?state=abc" + r := Registry{Address: strings.TrimPrefix(server.URL, "http://")} + + body := []byte("x") + sum := sha256.Sum256(body) + digest := "sha256:" + hex.EncodeToString(sum[:]) + if _, err := r.PublishArchive(context.Background(), "a/b", body, digest); err != nil { + t.Fatal(err) + } + if _, ok := f.blobs[digest]; !ok { + t.Fatal("the blob did not arrive when the location carried a query") + } +} + +func TestATagIsNeverAcceptedAsAPin(t *testing.T) { + // A tag can be made to point at something else, and a bundle is applied on machines with no + // mesh to ask about anything. + r := Registry{Address: "registry.invalid", Run: func( + _ context.Context, _, name string, args ...string) (string, error) { + if name == "docker" && len(args) > 0 && args[0] == "inspect" { + return "registry.invalid/shell/server:latest\n", nil + } + return "", nil + }} + _, err := r.PublishImage(context.Background(), "local", "shell/server") + if err == nil { + t.Fatal("an image referred to by tag was accepted") + } + if !strings.Contains(err.Error(), "pinned by digest") { + t.Fatalf("refused for the wrong reason: %v", err) + } +} + +func TestAPushedImageComesBackPinned(t *testing.T) { + pinned := "registry.invalid/shell/server@sha256:" + strings.Repeat("a", 64) + var ran []string + r := Registry{Address: "registry.invalid", Run: func( + _ context.Context, _, name string, args ...string) (string, error) { + ran = append(ran, name+" "+strings.Join(args, " ")) + if name == "docker" && len(args) > 0 && args[0] == "inspect" { + return pinned + "\n", nil + } + return "", nil + }} + got, err := r.PublishImage(context.Background(), "local", "shell/server") + if err != nil { + t.Fatal(err) + } + if got != pinned { + t.Fatalf("got %q", got) + } + if len(ran) < 2 || !strings.HasPrefix(ran[0], "docker tag") || !strings.HasPrefix(ran[1], "docker push") { + t.Fatalf("it did not tag and then push: %v", ran) + } +}