diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go new file mode 100644 index 0000000..ed6b8df --- /dev/null +++ b/cmd/mesh-builder/main.go @@ -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 +} diff --git a/cmd/mesh-control/main.go b/cmd/mesh-control/main.go index 8edc035..3613f2c 100644 --- a/cmd/mesh-control/main.go +++ b/cmd/mesh-control/main.go @@ -20,7 +20,6 @@ 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" @@ -134,7 +133,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 + build [--ref R] have a build machine build it, and record what came out pin which node this one gets a provision from unpin put that question back plan [--files|--json] what that node would run, and why @@ -869,29 +868,62 @@ func moduleCommand(ctx context.Context, args []string) error { return nil 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 { return err } - if len(shelf) == 0 { + if len(entries) == 0 { fmt.Println("this mesh knows about no modules yet") return nil } - var names []string - for n := range shelf { - names = append(names, n) - } - sort.Strings(names) - for _, n := range names { - m := shelf[n] - fmt.Printf("%-20s", m.Module) - if len(m.Provides) > 0 { - fmt.Printf(" provides %s", describeOffers(m.Provides)) + var stale int + for _, e := range entries { + m := e.Manifest + fmt.Printf("%-18s %-8s", m.Module, m.Version) + + switch { + case e.Provided: + fmt.Printf(" %-22s", "with the control plane") + case e.Source.Repository == "": + // 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() + + 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 ` to catch up\n", stale) } 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", 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 case "forget": @@ -1640,33 +1672,65 @@ func pinCommand(ctx context.Context, args []string, setting bool) error { 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") + wait := set.Duration("wait", 10*time.Minute, "how long to wait for a builder to answer") 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") + return errors.New("build [--ref R] [--wait D] [--dry-run]") } - publisher := builder.Registry{Address: *registry, Run: builder.Command} - result, err := builder.Build(ctx, builder.Command, publisher, positionals[0], *ref, *workspace) + ident, err := openIdentity(ctx) if err != nil { 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) } + + // 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 { - body, err := json.MarshalIndent(result.Manifest, "", " ") + body, err := json.MarshalIndent(manifest, "", " ") if err != nil { 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 // 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, + if err := inv.RegisterModule(ctx, manifest, inventory.Source{ + Repository: result.Repository, Ref: result.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) + fmt.Printf("\n%s %s, built on %s from %s\n", + manifest.Module, manifest.Version, result.On, short(result.Commit)) + fmt.Printf(" run `assign %s` to put it somewhere\n", manifest.Module) return nil } diff --git a/internal/inventory/catalogue.go b/internal/inventory/catalogue.go index 80e8fcd..476cb01 100644 --- a/internal/inventory/catalogue.go +++ b/internal/inventory/catalogue.go @@ -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) 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" diff --git a/internal/inventory/catalogue_test.go b/internal/inventory/catalogue_test.go index 481271a..b054993 100644 --- a/internal/inventory/catalogue_test.go +++ b/internal/inventory/catalogue_test.go @@ -3,6 +3,7 @@ package inventory import ( "context" "errors" + "strings" "testing" "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) } } + +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") + } +} diff --git a/internal/link/build.go b/internal/link/build.go new file mode 100644 index 0000000..e162c00 --- /dev/null +++ b/internal/link/build.go @@ -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 + } + } +}