Publish to the registry, and a command that builds a repository

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/<name>/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 <repository>` 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.
This commit is contained in:
2026-08-30 03:36:04 +02:00
parent 604b04886b
commit 7d033ad9f6
4 changed files with 397 additions and 0 deletions
+67
View File
@@ -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 <module> <file> what a module's config should say, for the whole mesh
settings set <module> <file> --node <n> ...or for one machine
settings clear <module> [--node <n>] take a layer away
build <repository> [--ref R] build a module from its source and record it
pin <node> <provision> <from> which node this one gets a provision from
unpin <node> <provision> put that question back
plan <node> [--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 <repository> [--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 <node> %s` to put it somewhere\n", result.Manifest.Module)
return nil
}
+5
View File
@@ -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
+151
View File
@@ -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/<name>/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 "?"
}
+174
View File
@@ -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)
}
}