The mesh builds: a machine takes the work, and the catalogue shows it

A build is work, not state. Everything else the control plane sends a
node is a declaration — this is what you should be — reconciled forever.
A build happens once and is finished. Putting it in a declaration would
mean rebuilding on every reconcile, or a declaration carrying "and I
already did this", which is state about an event rather than about a
machine.

So it travels on its own queue and the answer comes back correlated. One
queue, so several build machines share the work and each request is done
exactly once — which a per-machine routing key would not give.

mesh-builder is the program a build machine runs. Not the control plane,
which must not run commands on a machine; not the host, which would then
need a container runtime and git everywhere to do something almost no
machine will ever do. It holds its own broker credential and nothing
else.

Three properties that are decisions:

- a request is acknowledged only once the answer is away, so a builder
  that dies mid-build leaves the work for another machine rather than
  losing it with nobody ever hearing why
- one build at a time. Five at once against one runtime finishes all five
  slower than it would have finished the first, and the queue is what
  shares work between machines
- a failure is a RESULT. A build that fails silently is
  indistinguishable from a builder that is not running, and those want
  different responses

And `module list` is a catalogue: what exists, at which version, built
from which commit or handed over by hand or shipped with the control
plane, whether it is behind its source, and which machines run it. All of
that was recorded from the first build and none of it was shown, so "is
this current?" could only be answered by reading the database.

Proven against a real broker, registry and store: the mesh asked, a
builder consumed, built, published, answered; the manifest was recorded
with its commit; the source moved and the catalogue said "behind";
rebuilding caught it up with a new digest because the content changed.
This commit is contained in:
2026-08-30 03:46:02 +02:00
parent 7d033ad9f6
commit 421fe73dce
5 changed files with 606 additions and 34 deletions
+217
View File
@@ -0,0 +1,217 @@
// mesh-builder — the thing a build machine runs.
//
// It takes work from the mesh, turns a repository into artifacts, publishes them, and says what
// came out. It is **not** the control plane and it is **not** the host:
//
// - the control plane decides and never touches a machine. Building runs commands on one, and
// what the control plane may send a machine is bounded by the declaration language
// (novox/hq ADR 0005). "Run this build" is not in it, and widening the language so it could
// be would make the control plane able to run anything anywhere.
// - the host applies declarations and holds no opinion about what they contain. A host that
// also built things would need a container runtime and git, on every machine, to do something
// almost none of them will ever do.
//
// So it is a module: a program a machine runs because the mesh told it to, holding its own broker
// credential and nothing else. Compromise of a build machine is compromise of a build machine.
package main
import (
"context"
"encoding/json"
"fmt"
"os"
"os/signal"
"strings"
"syscall"
"time"
amqp "github.com/rabbitmq/amqp091-go"
"github.com/novox/mesh-control/internal/builder"
"github.com/novox/mesh-control/internal/link"
)
// version is set at build time.
var version = "development"
func main() {
if err := run(); err != nil {
fmt.Fprintf(os.Stderr, "mesh-builder: %v\n", err)
os.Exit(1)
}
}
const usage = `mesh-builder — builds modules for the mesh
It consumes build requests and answers with what it made. Nothing is listened on and nothing
is dialled except the broker.
MESH_BROKER_AMQP where the broker is, with this builder's own credential
MESH_REGISTRY host:port to publish artifacts to
MESH_WORKSPACE where to clone and build (default: a temporary directory)
`
func run() error {
if len(os.Args) > 1 {
switch os.Args[1] {
case "version":
fmt.Println(version)
return nil
default:
fmt.Print(usage)
return nil
}
}
amqpURL := strings.TrimSpace(os.Getenv("MESH_BROKER_AMQP"))
if amqpURL == "" {
return fmt.Errorf("no MESH_BROKER_AMQP: a builder with no broker has nothing to build")
}
registry := strings.TrimSpace(os.Getenv("MESH_REGISTRY"))
if registry == "" {
return fmt.Errorf(
"no MESH_REGISTRY: a built artifact nobody can fetch is not built")
}
workspace := os.Getenv("MESH_WORKSPACE")
if workspace == "" {
workspace = os.TempDir() + "/mesh-builder"
}
on, err := os.Hostname()
if err != nil {
on = "a build machine"
}
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
conn, err := amqp.Dial(amqpURL)
if err != nil {
// Not quoted back: the URL carries this builder's broker password.
return fmt.Errorf("cannot reach the broker: %w", err)
}
defer conn.Close()
channel, err := conn.Channel()
if err != nil {
return err
}
defer channel.Close()
if _, err := channel.QueueDeclare(link.BuildQueue, true, false, false, false, nil); err != nil {
return err
}
// One at a time. A build machine that took five requests at once would run five container
// builds against one runtime and finish all of them slower than it would have finished the
// first — and the queue is what shares work between machines, so nothing is lost by it.
if err := channel.Qos(1, 0, false); err != nil {
return err
}
// Not auto-acknowledged. A request acknowledged on arrival is a build that vanishes if this
// process dies mid-way, with nobody waiting on it ever hearing why.
requests, err := channel.ConsumeWithContext(ctx, link.BuildQueue, "mesh-builder",
false, false, false, false, nil)
if err != nil {
return err
}
fmt.Printf("building for the mesh, publishing to %s\n", registry)
publisher := builder.Registry{Address: registry, Run: builder.Command}
for {
select {
case <-ctx.Done():
fmt.Println("stopping")
return nil
case delivery, ok := <-requests:
if !ok {
return fmt.Errorf("the broker closed the connection")
}
answer(ctx, channel, publisher, on, workspace, delivery)
}
}
}
// answer does one build and says what happened, whichever way it went.
func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publisher,
on, workspace string, delivery amqp.Delivery) {
var request link.BuildRequest
if err := json.Unmarshal(delivery.Body, &request); err != nil {
// Unreadable. Acknowledged and dropped rather than requeued: a message this builder
// cannot parse will not become parseable by being delivered again, and requeueing it
// would put it in front of every real request for ever.
fmt.Fprintf(os.Stderr, "a request could not be read and was dropped: %v\n", err)
_ = delivery.Ack(false)
return
}
result := link.BuildResult{
ID: request.ID, Repository: request.Repository, Ref: request.Ref, On: on,
}
fmt.Printf("building %s", request.Repository)
if request.Ref != "" {
fmt.Printf(" at %s", request.Ref)
}
fmt.Println()
built, err := builder.Build(ctx, builder.Command, publisher,
request.Repository, request.Ref, workspace)
if err != nil {
// A failure is a result. A build that fails and says nothing is indistinguishable from a
// builder that is not running, and those want completely different responses.
result.Failed = err.Error()
fmt.Fprintf(os.Stderr, " failed: %v\n", err)
} else {
manifest, marshalErr := json.Marshal(built.Manifest)
if marshalErr != nil {
result.Failed = marshalErr.Error()
} else {
result.Commit = built.Commit
result.Manifest = manifest
for _, made := range built.Built {
result.Made = append(result.Made, link.MadeArtifact{
Name: made.Name, Kind: made.Kind, Reference: made.Reference,
})
}
fmt.Printf(" built %s from %s\n", built.Manifest.Module, short(built.Commit))
}
}
body, err := json.Marshal(result)
if err != nil {
fmt.Fprintf(os.Stderr, "cannot report a build: %v\n", err)
_ = delivery.Ack(false)
return
}
replyTo := delivery.ReplyTo
if replyTo == "" {
// Nobody is waiting. Still reported, to the exchange, so a control plane that records
// builds hears about it — a build whose outcome exists nowhere is one nobody can audit.
if err := channel.PublishWithContext(ctx, link.Exchange, link.KeyBuilt, false, false,
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
fmt.Fprintf(os.Stderr, "cannot publish a build result: %v\n", err)
}
_ = delivery.Ack(false)
return
}
publishCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
if err := channel.PublishWithContext(publishCtx, "", replyTo, false, false, amqp.Publishing{
ContentType: "application/json",
CorrelationId: result.ID,
Body: body,
}); err != nil {
fmt.Fprintf(os.Stderr, "cannot answer a build request: %v\n", err)
}
// Acknowledged only once the answer is away, so a builder that dies before answering leaves
// the request for another machine rather than losing it.
_ = delivery.Ack(false)
}
func short(commit string) string {
if len(commit) > 8 {
return commit[:8]
}
return commit
}
+99 -34
View File
@@ -20,7 +20,6 @@ import (
"time" "time"
"github.com/novox/mesh-control/internal/broker" "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/catalogue"
"github.com/novox/mesh-control/internal/identity" "github.com/novox/mesh-control/internal/identity"
"github.com/novox/mesh-control/internal/inventory" "github.com/novox/mesh-control/internal/inventory"
@@ -134,7 +133,7 @@ func usage() {
settings set <module> <file> what a module's config should say, for the whole mesh 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 set <module> <file> --node <n> ...or for one machine
settings clear <module> [--node <n>] take a layer away settings clear <module> [--node <n>] take a layer away
build <repository> [--ref R] build a module from its source and record it build <repository> [--ref R] have a build machine build it, and record what came out
pin <node> <provision> <from> which node this one gets a provision from pin <node> <provision> <from> which node this one gets a provision from
unpin <node> <provision> put that question back unpin <node> <provision> put that question back
plan <node> [--files|--json] what that node would run, and why plan <node> [--files|--json] what that node would run, and why
@@ -869,29 +868,62 @@ func moduleCommand(ctx context.Context, args []string) error {
return nil return nil
case "list": case "list":
shelf, err := inv.Catalogue(ctx) // The catalogue: what exists, where it came from, whether it is current, and who runs it.
// The provenance was recorded from the first build and nothing showed it, which made
// "is this current?" a question you could only answer by reading the database.
entries, err := inv.Catalogued(ctx)
if err != nil { if err != nil {
return err return err
} }
if len(shelf) == 0 { if len(entries) == 0 {
fmt.Println("this mesh knows about no modules yet") fmt.Println("this mesh knows about no modules yet")
return nil return nil
} }
var names []string var stale int
for n := range shelf { for _, e := range entries {
names = append(names, n) m := e.Manifest
} fmt.Printf("%-18s %-8s", m.Module, m.Version)
sort.Strings(names)
for _, n := range names { switch {
m := shelf[n] case e.Provided:
fmt.Printf("%-20s", m.Module) fmt.Printf(" %-22s", "with the control plane")
if len(m.Provides) > 0 { case e.Source.Repository == "":
fmt.Printf(" provides %s", describeOffers(m.Provides)) // Handed over by hand. Legitimate — it is how a module is fixed in a hurry — and
// worth saying, because nothing can rebuild it.
fmt.Printf(" %-22s", "handed over")
case !e.Source.Current():
stale++
fmt.Printf(" %-22s", "behind "+short(e.Source.BuiltFrom)+" < "+short(e.Source.Head))
default:
fmt.Printf(" %-22s", "built "+short(e.Source.BuiltFrom))
} }
for _, c := range m.Claims {
fmt.Printf(" claims %s/%s", c.At(), c.Name) if len(e.On) > 0 {
fmt.Printf(" on %s", strings.Join(e.On, ", "))
} else {
fmt.Printf(" on nothing")
} }
fmt.Println() fmt.Println()
var says []string
if len(m.Provides) > 0 {
says = append(says, "provides "+describeOffers(m.Provides))
}
if len(m.Requires) > 0 {
says = append(says, "requires "+strings.Join(m.Requires, ", "))
}
for _, c := range m.Claims {
says = append(says, "claims "+c.At()+"/"+c.Name)
}
if len(m.Capabilities) > 0 {
says = append(says, "needs "+strings.Join(m.Capabilities, ", "))
}
if len(says) > 0 {
fmt.Printf(" %s\n", strings.Join(says, " · "))
}
}
if stale > 0 {
fmt.Printf("\n%d module(s) behind their source — `build <repository>` to catch up\n", stale)
} }
return nil return nil
@@ -912,7 +944,7 @@ func moduleCommand(ctx context.Context, args []string) error {
} }
fmt.Printf("%s is behind: the mesh holds %s and the source has %s\n", fmt.Printf("%s is behind: the mesh holds %s and the source has %s\n",
args[1], short(from.BuiltFrom), short(from.Head)) args[1], short(from.BuiltFrom), short(from.Head))
fmt.Println(" build it and `module add` the result to catch up") fmt.Printf(" run `build %s` to catch up\n", from.Repository)
return nil return nil
case "forget": case "forget":
@@ -1640,33 +1672,65 @@ func pinCommand(ctx context.Context, args []string, setting bool) error {
func buildCommand(ctx context.Context, args []string) error { func buildCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("build", flag.ContinueOnError) set := flag.NewFlagSet("build", flag.ContinueOnError)
ref := set.String("ref", "", "the branch, tag or commit to build") ref := set.String("ref", "", "the branch, tag or commit to build")
registry := set.String("registry", os.Getenv("MESH_REGISTRY"), wait := set.Duration("wait", 10*time.Minute, "how long to wait for a builder to answer")
"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") dryRun := set.Bool("dry-run", false, "build and print the manifest, recording nothing")
positionals, err := parseAround(set, args) positionals, err := parseAround(set, args)
if err != nil { if err != nil {
return err return err
} }
if len(positionals) != 1 { if len(positionals) != 1 {
return errors.New("build <repository> [--ref R] [--registry host:port]") return errors.New("build <repository> [--ref R] [--wait D] [--dry-run]")
}
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} ident, err := openIdentity(ctx)
result, err := builder.Build(ctx, builder.Command, publisher, positionals[0], *ref, *workspace)
if err != nil { if err != nil {
return err return err
} }
defer ident.Close()
for _, made := range result.Built { server, err := link.Connect(nil, nil)
if err != nil {
return err
}
defer server.Close()
// Correlated by something the control plane makes, not by the module's name: two builds of one
// module can be in flight, and the second answer is not the first one's.
request := link.BuildRequest{
ID: fmt.Sprintf("%s-%d", "build", time.Now().UnixNano()),
Repository: positionals[0],
Ref: *ref,
}
fmt.Printf("asked for %s", request.Repository)
if *ref != "" {
fmt.Printf(" at %s", *ref)
}
fmt.Println()
result, err := link.RequestBuild(ctx, server.Channel(), request, *wait)
if err != nil {
return err
}
if result.Failed != "" {
// The builder's own words. Wrapping them in something about the control plane would put
// two explanations between a person and a build log.
return fmt.Errorf("%s could not build %s:\n%s", result.On, result.Repository, result.Failed)
}
for _, made := range result.Made {
fmt.Printf(" %-12s %s %s\n", made.Name, made.Kind, made.Reference) fmt.Printf(" %-12s %s %s\n", made.Name, made.Kind, made.Reference)
} }
// Parsed with the same parser a hand-written manifest goes through. A second path would be a
// second thing to disagree about what a manifest is.
manifest, err := catalogue.ParseManifest(result.Manifest)
if err != nil {
return fmt.Errorf("%s built %s and what came back is not a manifest: %w",
result.On, result.Repository, err)
}
if *dryRun { if *dryRun {
body, err := json.MarshalIndent(result.Manifest, "", " ") body, err := json.MarshalIndent(manifest, "", " ")
if err != nil { if err != nil {
return err return err
} }
@@ -1682,13 +1746,14 @@ func buildCommand(ctx context.Context, args []string) error {
// Recorded with where it came from, so "is this current?" is answerable without building it // Recorded with where it came from, so "is this current?" is answerable without building it
// again (novox/hq ADR 0009). // again (novox/hq ADR 0009).
if err := inv.RegisterModule(ctx, result.Manifest, inventory.Source{ if err := inv.RegisterModule(ctx, manifest, inventory.Source{
Repository: positionals[0], Ref: *ref, BuiltFrom: result.Commit, Head: result.Commit, Repository: result.Repository, Ref: result.Ref,
BuiltFrom: result.Commit, Head: result.Commit,
}); err != nil { }); err != nil {
return err return err
} }
fmt.Printf("\n%s %s, built from %s\n", fmt.Printf("\n%s %s, built on %s from %s\n",
result.Manifest.Module, result.Manifest.Version, short(result.Commit)) manifest.Module, manifest.Version, result.On, short(result.Commit))
fmt.Printf(" run `assign <node> %s` to put it somewhere\n", result.Manifest.Module) fmt.Printf(" run `assign <node> %s` to put it somewhere\n", manifest.Module)
return nil return nil
} }
+60
View File
@@ -515,3 +515,63 @@ func (i *Inventory) pinRows(ctx context.Context, nodeName string) (int, error) {
`select count(*) from provision_pin where node = $1`, node.ID).Scan(&n) `select count(*) from provision_pin where node = $1`, node.ID).Scan(&n)
return n, err return n, err
} }
// Entry is one module as a catalogue shows it: what it is, where it came from, and who runs it.
type Entry struct {
Manifest catalogue.Manifest
Source Source
// On is every node this module is assigned to, sorted.
On []string
// Provided is true when the module came with the control plane rather than from a repository.
Provided bool
}
// Catalogued is every module the mesh knows about, with everything a person asks about one.
//
// **One query rather than a call per module.** A catalogue that costs a round trip per row is a
// catalogue nobody lists, and the questions here — what is this, where did it come from, who is
// running it — are asked together every time.
func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
rows, err := i.store.Pool().Query(ctx,
`select m.name, m.manifest,
coalesce(m.source, ''), coalesce(m.ref, ''),
coalesce(m.built_from, ''), coalesce(m.source_head, ''),
coalesce(array_agg(n.name order by n.name) filter (where n.name is not null), '{}')
from module m
left join assignment a on a.module = m.name
left join node n on n.id = a.node
group by m.name, m.manifest, m.source, m.ref, m.built_from, m.source_head
order by m.name`)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Entry
for rows.Next() {
var raw []byte
var name string
var source Source
var on []string
if err := rows.Scan(&name, &raw, &source.Repository, &source.Ref,
&source.BuiltFrom, &source.Head, &on); err != nil {
return nil, err
}
var m catalogue.Manifest
if err := json.Unmarshal(raw, &m); err != nil {
return nil, err
}
entry := Entry{Manifest: m, Source: source, On: on}
if source.Repository == providedBy {
// It came with the control plane. Not a repository, and showing it as one would have
// somebody go looking for it.
entry.Provided = true
entry.Source = Source{}
}
out = append(out, entry)
}
return out, rows.Err()
}
// providedBy is what the source column says for a module the control plane ships.
const providedBy = "the control plane"
+89
View File
@@ -3,6 +3,7 @@ package inventory
import ( import (
"context" "context"
"errors" "errors"
"strings"
"testing" "testing"
"github.com/novox/mesh-control/internal/catalogue" "github.com/novox/mesh-control/internal/catalogue"
@@ -428,3 +429,91 @@ func TestAPinGoesWhenTheProviderLeavesTheMesh(t *testing.T) {
t.Fatalf("a choice outlived the machine it named: %d row(s) left", rows) t.Fatalf("a choice outlived the machine it named: %d row(s) left", rows)
} }
} }
func TestTheCatalogueSaysWhereEachModuleCameFromAndWhoRunsIt(t *testing.T) {
// The provenance was recorded from the first build and nothing showed it, which made "is this
// current?" a question you could only answer by reading the database.
inv := fresh(t)
ctx := context.Background()
for _, n := range []string{"workstation", "laptop"} {
if _, err := inv.AddNode(ctx, n); err != nil {
t.Fatal(err)
}
}
if err := inv.RegisterModule(ctx, manifest("shell", []string{"login-shell"}, nil),
Source{Repository: "https://forge.invalid/shell.git", BuiltFrom: "aaa", Head: "aaa"}); err != nil {
t.Fatal(err)
}
if err := inv.RegisterModule(ctx, manifest("byhand", nil, nil), Source{}); err != nil {
t.Fatal(err)
}
if err := inv.Provide(ctx, manifest("networking", nil, []string{"login-shell"})); err != nil {
t.Fatal(err)
}
for _, n := range []string{"workstation", "laptop"} {
if err := inv.Assign(ctx, n, "shell"); err != nil {
t.Fatal(err)
}
}
entries, err := inv.Catalogued(ctx)
if err != nil {
t.Fatal(err)
}
by := map[string]Entry{}
for _, e := range entries {
by[e.Manifest.Module] = e
}
if len(by) != 3 {
t.Fatalf("the catalogue has %d modules", len(by))
}
// Sorted, and both nodes, so a person reading it twice sees the same thing.
if got := strings.Join(by["shell"].On, ","); got != "laptop,workstation" {
t.Fatalf("shell runs on %q", got)
}
if by["shell"].Source.Repository != "https://forge.invalid/shell.git" {
t.Fatalf("shell came from %q", by["shell"].Source.Repository)
}
if by["shell"].Provided {
t.Fatal("a module built from a repository was reported as shipped with the control plane")
}
// A module nobody runs is in the catalogue: the catalogue is what EXISTS, and what runs is a
// different question the same row answers.
if len(by["byhand"].On) != 0 {
t.Fatalf("byhand runs on %v", by["byhand"].On)
}
// Handed over by hand is its own state. Nothing can rebuild it, and showing it as a
// repository would send somebody looking for one.
if by["byhand"].Source.Repository != "" || by["byhand"].Provided {
t.Fatalf("byhand: %+v", by["byhand"])
}
if !by["networking"].Provided {
t.Fatal("a module the control plane ships was not marked as such")
}
if by["networking"].Source.Repository != "" {
// It is not a repository, and showing it as one would have somebody go looking for it.
t.Fatalf("networking claims to come from %q", by["networking"].Source.Repository)
}
}
func TestACatalogueEntryKnowsWhetherItIsBehind(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
if err := inv.RegisterModule(ctx, manifest("shell", nil, nil),
Source{Repository: "https://forge.invalid/shell.git", BuiltFrom: "aaa", Head: "aaa"}); err != nil {
t.Fatal(err)
}
if err := inv.SourceMoved(ctx, "shell", "bbb"); err != nil {
t.Fatal(err)
}
entries, err := inv.Catalogued(ctx)
if err != nil {
t.Fatal(err)
}
if entries[0].Source.Current() {
t.Fatal("a module whose source moved reported itself current")
}
}
+141
View File
@@ -0,0 +1,141 @@
package link
import (
"context"
"encoding/json"
"fmt"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
// Asking a machine to build a module, and hearing what came out.
//
// **A build is work, not state.** Everything else the control plane sends a node is a declaration
// — *this is what you should be* — and the node reconciles toward it forever. A build happens
// once, produces something, and is finished. Putting it in a declaration would mean a machine
// rebuilding on every reconcile, or the declaration carrying "and I already did this", which is
// state about an event rather than about the machine.
//
// So it travels on its own queue, and the reply comes back correlated. That is also why the
// builder is a **separate consumer** rather than the host: the host applies declarations and
// holds no opinion about what they contain, and a host that also built things would be a host
// with a container runtime requirement and a git dependency (novox/hq ADR 0005).
// BuildQueue is where build requests wait. One queue, so several build machines can share the
// work and each request is done exactly once — which is what a queue is for and what a
// per-machine routing key would not give.
const BuildQueue = "builds"
// KeyBuilt is what a builder publishes when it has finished, successfully or not.
const KeyBuilt = "built"
// BuildRequest is one module to build.
type BuildRequest struct {
// ID correlates the answer with the asking. Not the module name: two builds of one module can
// be in flight, and the second answer is not the first one's.
ID string `json:"id"`
// Repository is where the source is, as git would clone it.
Repository string `json:"repository"`
// Ref is the branch, tag or commit. Empty means whatever the repository's default is, which
// is the only case where the mesh does not know what it built until it has built it.
Ref string `json:"ref,omitempty"`
}
// BuildResult is what a builder says back.
//
// **Failure is a result, not an absence.** A build that fails and says nothing is
// indistinguishable from a builder that is not running, and those want completely different
// responses — the same rule the host follows about a service that does not exist.
type BuildResult struct {
ID string `json:"id"`
Repository string `json:"repository"`
Ref string `json:"ref,omitempty"`
// On is the machine that did it, so a failure that is about one machine can be told from one
// about the source.
On string `json:"on"`
// Commit is what was actually built. The mesh records it, which is what makes "is this
// current?" answerable without building again.
Commit string `json:"commit,omitempty"`
// Manifest is the module as the mesh should hold it, artifacts resolved to digests. Raw,
// because the control plane parses it with the same parser it uses for one handed over by
// hand — a second path would be a second thing to disagree.
Manifest json.RawMessage `json:"manifest,omitempty"`
// Made is each artifact, for reporting.
Made []MadeArtifact `json:"made,omitempty"`
// Failed is why, when it did.
Failed string `json:"failed,omitempty"`
}
// MadeArtifact is one thing a build produced, as a person would want it reported.
type MadeArtifact struct {
Name string `json:"name"`
Kind string `json:"kind"`
Reference string `json:"reference"`
}
// RequestBuild asks for a module to be built and waits for the answer.
//
// Waiting rather than returning immediately, because the thing a person wants after asking for a
// build is to know whether it worked. A build that is dispatched and forgotten needs somewhere to
// look afterwards, and there is nowhere yet.
func RequestBuild(ctx context.Context, channel *amqp.Channel, request BuildRequest,
timeout time.Duration) (BuildResult, error) {
// Its own queue for the answer, declared before the ask. Consuming from the shared exchange
// would mean competing with the control plane's own consumer for a message meant for this
// caller — which is the fault this package's own doc comment records having had.
replies, err := channel.QueueDeclare("", false, true, true, false, nil)
if err != nil {
return BuildResult{}, err
}
answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil)
if err != nil {
return BuildResult{}, err
}
body, err := json.Marshal(request)
if err != nil {
return BuildResult{}, err
}
if err := channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{
ContentType: "application/json",
DeliveryMode: amqp.Persistent,
CorrelationId: request.ID,
ReplyTo: replies.Name,
Body: body,
}); err != nil {
return BuildResult{}, err
}
waiting, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
for {
select {
case <-waiting.Done():
return BuildResult{}, fmt.Errorf(
"no builder answered within %s. Something must be consuming %q, and nothing is "+
"— or it is building something that takes longer than this",
timeout, BuildQueue)
case delivery, ok := <-answers:
if !ok {
return BuildResult{}, fmt.Errorf("the connection closed while waiting for a build")
}
var result BuildResult
if err := json.Unmarshal(delivery.Body, &result); err != nil {
return BuildResult{}, fmt.Errorf("a builder answered with something unreadable: %w", err)
}
if result.ID != request.ID {
// Somebody else's answer on this queue. Ignored rather than returned, because
// returning it would attribute one build's outcome to another's.
continue
}
return result, nil
}
}
}