diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index f321436..e27d4b3 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -29,10 +29,10 @@ import ( "os/signal" "strings" "syscall" - "time" amqp "github.com/rabbitmq/amqp091-go" + "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/builder" "github.com/novox/mesh-controller/internal/link" ) @@ -115,72 +115,63 @@ func run() error { ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() - conn, err := dial(credential) - 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) + machine, err := takeWorkFrom(credential, on) if err != nil { return err } + defer machine.Close() fmt.Fprintf(os.Stderr, "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) - } + return machine.Take(ctx, func(ctx context.Context, work link.Build) { + answer(ctx, publisher, on, workspace, work) + }) +} + +// takeWorkFrom opens this machine's link to whichever bus the mesh is on. +// +// **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5): a build +// machine told about both would take work from one and answer on the other, and every log line would +// say it was fine. +func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) { + address, onNATS, err := broker.OnNATS() + if err != nil { + return nil, err } + if err := broker.MustBeOneBus(credential.URL, address); err != nil { + return nil, err + } + if onNATS { + js, err := broker.Dial(address) + if err != nil { + return nil, fmt.Errorf("cannot reach the bus at %s: %w", address, err) + } + return link.MachineOverNATS(js, on), nil + } + + conn, err := dial(credential) + if err != nil { + // Not quoted back: the URL carries this builder's broker password. + return nil, fmt.Errorf("cannot reach the broker: %w", err) + } + channel, err := conn.Channel() + if err != nil { + conn.Close() + return nil, err + } + return link.MachineOverCurrent(conn, channel, on), nil } // 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) { +func answer(ctx context.Context, publisher builder.Publisher, on, workspace string, work link.Build) { + request := work.Request() - // **First thing, and to stdout.** A build request that arrives and produces no visible line - // until it either finishes or fails is indistinguishable from one that never arrived — which - // cost a long diagnosis against a running mesh, chasing "the handler never fired" when the - // truth was only that the handler said nothing until the end. - fmt.Fprintf(os.Stderr, "a build request arrived (%d bytes)\n", len(delivery.Body)) - - 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 - } + // **First thing, and to stdout.** A build request that arrives and produces no visible line until + // it either finishes or fails is indistinguishable from one that never arrived — which cost a long + // diagnosis against a running mesh, chasing "the handler never fired" when the truth was only that + // the handler said nothing until the end. + fmt.Fprintf(os.Stderr, "a build request arrived for %s\n", request.Repository) result := link.BuildResult{ ID: request.ID, Repository: request.Repository, Path: request.Path, @@ -199,8 +190,8 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis var built builder.Result if err == nil { // The package-registry credential is a build input, so it is resolved before the clone: a - // build that could not have resolved its dependencies is refused in front of the reason, - // not after a clone that then fails at npm ci. + // build that could not have resolved its dependencies is refused in front of the reason, not + // after a clone that then fails at npm ci. built, err = builder.Build(ctx, builder.Command, publisher, request.Repository, request.Path, request.Ref, workspace, request.Held, npmrc, forgeFrom(), @@ -230,67 +221,19 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis } } - body, err := json.Marshal(result) - if err != nil { - fmt.Fprintf(os.Stderr, "cannot report a build: %v\n", err) - _ = delivery.Ack(false) + if err := work.Announce(ctx, result); err != nil { + // Said, not fatal: the build happened. A build reported as failed because announcing it + // failed is a lie about work that was done — and the request stays unsettled below only if + // nothing was said at all, so another machine can try. + fmt.Fprintf(os.Stderr, "cannot say what came of a build: %v\n", err) return } - // Always through the exchange, whether or not somebody is waiting. - // - // **Never the default exchange.** Permission there is granted per exchange rather than per - // queue, so a builder allowed to use it could publish into any node's queue — the privilege a - // build machine most obviously should not have. An asker binds its own reply queue to this - // key and filters by correlation; a control plane that records builds is bound to it too, so - // a result nobody asked for is still kept rather than reported into the void. - publishCtx, cancel := context.WithTimeout(ctx, 30*time.Second) - defer cancel() - if err := channel.PublishWithContext(publishCtx, link.Exchange, link.KeyBuilt, 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) + // Settled only once the outcome is away, so a machine that dies before answering leaves the work + // for another rather than losing it. + if err := work.Done(); err != nil { + fmt.Fprintf(os.Stderr, "the outcome is away and the request could not be settled: %v\n", err) } - // **And announced, which is a different act from answering.** The reply goes to whoever asked - // and is correlated to their request; this says to the whole mesh that a module now exists at - // a commit, and the catalogue places it in the module graph (novox/hq ADR 0072). A build - // nobody asked for still has to be announced, or the graph knows less than the registry does. - // - // Only on success: a failed build produced no module-version, and announcing one would put - // something in the graph that was never made. - if result.Failed == "" && result.Commit != "" { - announced := map[string]any{ - "module": moduleOf(result.Manifest), "commit": result.Commit, - "repository": result.Repository, "path": result.Path, "ref": result.Ref, - "manifest": json.RawMessage(result.Manifest), "against": result.Against, - "made": result.Made, - } - if err := link.EmitEvent(publishCtx, link.OverCurrent{Channel: channel}, link.KeyModuleBuilt, "builder", on, announced); err != nil { - // Said, not fatal: the build happened and was answered. A module the catalogue has not - // heard of is a gap somebody can close; a build reported as failed because announcing - // it failed is a lie about work that was done. - fmt.Fprintf(os.Stderr, " built, but could not announce it: %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) -} - -// moduleOf reads the module's name out of the manifest it just built, which is the only place it is -// authoritative — the request named a repository and a path, not a module. -func moduleOf(manifest json.RawMessage) string { - var named struct { - Module string `json:"module"` - } - if err := json.Unmarshal(manifest, &named); err != nil { - return "" - } - return named.Module } // packagesFrom is where a build resolves the mesh's own published packages — the SDK above all diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 582b04c..600d36a 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -395,7 +395,13 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti } fmt.Println() - result, err := link.RequestBuild(ctx, server.Channel(), request, wait) + ask, err := askOver(server) + if err != nil { + return err + } + defer ask.Close() + + result, err := ask.Submit(ctx, request, wait) if err != nil { return err } @@ -472,7 +478,13 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai } defer server.Close() - result, err := link.RequestBuild(ctx, server.Channel(), link.BuildRequest{ + ask, err := askOver(server) + if err != nil { + return err + } + defer ask.Close() + + result, err := ask.Submit(ctx, link.BuildRequest{ ID: fmt.Sprintf("%s-%d", "build", time.Now().UnixNano()), Repository: repository, Path: path, Ref: ref, Held: heldBy(ctx), @@ -557,3 +569,22 @@ func heldBy(ctx context.Context) map[string]string { } return routed } + +// askOver opens the way a build is asked for, on whichever bus the mesh is on. +// +// **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5). On the bus +// the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus +// being built it dials, because a build request is a one-shot and holds nothing else. +func askOver(server *link.Server) (link.Builders, error) { + address, onNATS, err := broker.OnNATS() + if err != nil { + return nil, err + } + if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), address); err != nil { + return nil, err + } + if onNATS { + return link.BuildsOverNATS(address) + } + return link.BuildsOverCurrent(server.Channel()), nil +} diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 87748fb..9b31266 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -171,6 +171,10 @@ const ControllerName = "controller" var ControllerFollows = []string{ moduleEventSubject("mesh-catalog", "upgraded"), moduleEventSubject("mesh-catalog", "catching-up"), + // A build's outcome, which is the build-machine role's own event now (ADR 0121) rather than a + // message on the control branch. Same three audiences, one publish: whoever asked, this, and the + // catalogue. + seatEventSubject("mesh-build-machine", "built"), } // moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so @@ -179,6 +183,11 @@ func moduleEventSubject(module, event string) string { return "mesh.mod." + module + ".event." + event } +// seatEventSubject is where a role's own event lands, derived the same way a holder's permission is. +func seatEventSubject(seat, verb string) string { + return "mesh.seat." + seat + ".event." + verb +} + // MeshConsumers is what the controller consumes, in the order a person reads it. // // **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's, diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 0a34b53..8e44a6e 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] } - subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.>"] } + subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.>", "mesh.seat.mesh-build-machine.event.built"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/link/build.go b/internal/link/build.go index 2ebe602..f097233 100644 --- a/internal/link/build.go +++ b/internal/link/build.go @@ -82,6 +82,16 @@ type BuildResult struct { // about the source. On string `json:"on"` + // Module is what was built, read out of the manifest — the only place it is authoritative, since + // a request names a repository and a path. + // + // **Here because one message now reaches three audiences** (novox/hq ADR 0121). On the bus the + // mesh runs on today the answer and the announcement were two publishes to two topologies, so a + // result needed no module name and the announcement carried one. On the bus being built the + // outcome is the role's own event, and the catalogue reading it needs to know what was built. + // Empty on a failed build, which produced no module version. + Module string `json:"module,omitempty"` + // 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"` diff --git a/internal/link/builds.go b/internal/link/builds.go new file mode 100644 index 0000000..0542f61 --- /dev/null +++ b/internal/link/builds.go @@ -0,0 +1,100 @@ +package link + +import ( + "context" + "encoding/json" + "fmt" + "time" +) + +// Asking a role to build something, and being told what came of it. +// +// A build is work submitted to a role, not a message to a machine (novox/hq ADR 0121). The +// build-machine seat accepts a build and emits an outcome, so the same publish that answers whoever +// asked also reaches the controller that records it and the catalogue that places it in the module +// graph — and no build machine needs permission to publish into anybody's inbox. +// +// **This is the one flow whose shape differs from every other**, which is why it has its own seam +// rather than living in `Bus`. Everything else the controller sends is either an event nobody must +// act on or a declaration a node reconciles toward; a build is a request that takes minutes and has +// exactly one answer. Too long for request/reply, too particular to be an event. + +// TheBuildMachine is the role a build is submitted to. +const TheBuildMachine = "mesh-build-machine" + +// BuildWork is where a build request lands, and BuildOutcome is where its result does. Derived from +// the seat, so both sides name the role and neither names the other. +func BuildWork() string { return "mesh.seat." + TheBuildMachine + ".accept.build" } +func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.built" } + +// Builders is how work reaches a build machine and how the outcome comes back. +type Builders interface { + // Submit asks for one build and waits for its outcome. + // + // The wait is long by nature. A build clones, pulls a base image and runs a container build, so + // a timeout here says "nothing is doing builds" rather than "this build is slow" — and the two + // need different remedies, which is why the message distinguishes them. + Submit(ctx context.Context, request BuildRequest, wait time.Duration) (BuildResult, error) + + // Close lets go of whatever was dialled. + Close() +} + +// BuildMachine is a machine taking work from the role it holds. +type BuildMachine interface { + // Take hands each request to do until the context ends, and says why it stopped. + Take(ctx context.Context, do func(context.Context, Build)) error + Close() +} + +// Build is one request a machine has been handed. +type Build interface { + // Request is what to build. + Request() BuildRequest + + // Announce publishes the outcome as the role's own event. + // + // One publish, three audiences: whoever asked matches it by the id their request carried, the + // controller records it, and the catalogue places it. On the bus the mesh runs on today that + // fan-out came from a shared exchange; here the mesh derived the subject. + Announce(ctx context.Context, result BuildResult) error + + // Done settles the request. Called only after the outcome is away, so a machine that dies + // before announcing leaves the work for another rather than losing it. + Done() error + + // Hold hands the work back for another attempt after the delay. + Hold(after time.Duration) error +} + +// waitingFor is the message a caller gets when nothing answered. Its own function because both +// transports say it, and saying it differently in two places is how one of them ends up vague. +func waitingFor(wait time.Duration) error { + return fmt.Errorf( + "no build machine answered within %s. Either nothing holds %s — in which case the work is "+ + "queued and will be done when something does — or a build is taking longer than this", + wait, TheBuildMachine) +} + +// theOutcomeOf reads a result and says whether it is the answer to this request. +func theOutcomeOf(body []byte, id string) (BuildResult, bool, error) { + var result BuildResult + if err := json.Unmarshal(body, &result); err != nil { + return BuildResult{}, false, fmt.Errorf("a build machine answered with something unreadable: %w", err) + } + // Somebody else's build. Skipped rather than returned, because returning it would attribute one + // build's outcome to another's. + return result, result.ID == id, nil +} + +// ModuleOf reads the module's name out of a manifest a build produced, which is the only place it is +// authoritative — a request named a repository and a path, not a module. +func ModuleOf(manifest json.RawMessage) string { + var named struct { + Module string `json:"module"` + } + if err := json.Unmarshal(manifest, &named); err != nil { + return "" + } + return named.Module +} diff --git a/internal/link/builds_current.go b/internal/link/builds_current.go new file mode 100644 index 0000000..1f60d41 --- /dev/null +++ b/internal/link/builds_current.go @@ -0,0 +1,184 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// The build flow on the bus the mesh runs on today. +// +// Moved behind the seam rather than changed. The queue, the reply binding and the correlation are +// what they were, because the mesh is running on this. + +// currentBuilds asks for builds over a channel. +type currentBuilds struct{ channel *amqp.Channel } + +// BuildsOverCurrent is the asking side on the bus the mesh has. +func BuildsOverCurrent(channel *amqp.Channel) Builders { return currentBuilds{channel: channel} } + +func (b currentBuilds) Close() {} + +func (b currentBuilds) Submit(ctx context.Context, request BuildRequest, + wait time.Duration) (BuildResult, error) { + + // Its own queue for the answer, declared before the ask. Consuming from the shared exchange + // would mean competing with the controller's own consumer for a message meant for this caller. + replies, err := b.channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil) + if err != nil { + return BuildResult{}, err + } + // **A builder never publishes to the default exchange**, because permission there is per + // exchange and not per queue — a builder allowed to use it could publish into any node's queue, + // which is the privilege a build machine most obviously should not have. The cost is that every + // asker sees every result, which is why the correlation is checked below rather than assumed. + if err := b.channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil { + return BuildResult{}, err + } + answers, err := b.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 := b.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, wait) + defer cancel() + for { + select { + case <-waiting.Done(): + return BuildResult{}, waitingFor(wait) + case delivery, ok := <-answers: + if !ok { + return BuildResult{}, errors.New("the connection closed while waiting for a build") + } + result, mine, err := theOutcomeOf(delivery.Body, request.ID) + if err != nil { + return BuildResult{}, err + } + if mine { + return result, nil + } + } + } +} + +// --- the machine's side --------------------------------------------------------------------- + +type currentMachine struct { + conn *amqp.Connection + channel *amqp.Channel + on string +} + +// MachineOverCurrent takes build work over a channel. +func MachineOverCurrent(conn *amqp.Connection, channel *amqp.Channel, on string) BuildMachine { + return ¤tMachine{conn: conn, channel: channel, on: on} +} + +func (m *currentMachine) Close() {} + +func (m *currentMachine) Take(ctx context.Context, do func(context.Context, Build)) error { + if _, err := m.channel.QueueDeclare(BuildQueue, true, false, false, false, nil); err != nil { + return err + } + // One at a time. A 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 := m.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 := m.channel.ConsumeWithContext(ctx, BuildQueue, "mesh-builder", + false, false, false, false, nil) + if err != nil { + return err + } + for { + select { + case <-ctx.Done(): + return nil + case delivery, ok := <-requests: + if !ok { + return errors.New("the broker closed the connection") + } + var request BuildRequest + if err := json.Unmarshal(delivery.Body, &request); err != nil { + // Unreadable: rejected rather than retried, because the next attempt reads the same + // bytes. Nobody waiting hears an answer, which is correct — there was no request. + _ = delivery.Reject(false) + continue + } + do(ctx, ¤tBuild{request: request, delivery: delivery, on: m.on, channel: m.channel}) + } + } +} + +type currentBuild struct { + request BuildRequest + delivery amqp.Delivery + on string + channel *amqp.Channel +} + +func (b *currentBuild) Request() BuildRequest { return b.request } + +// Announce answers and announces, which on this bus are two publishes to two exchanges. +// +// The reply goes to whoever asked, correlated to their request; the announcement says to the whole +// mesh that a module now exists at a commit (novox/hq ADR 0072). Only a successful build is +// announced: a failed one produced no module version, and announcing one would put something in the +// graph that was never made. +func (b *currentBuild) Announce(ctx context.Context, result BuildResult) error { + body, err := json.Marshal(result) + if err != nil { + return err + } + if err := b.channel.PublishWithContext(ctx, Exchange, KeyBuilt, false, false, amqp.Publishing{ + ContentType: "application/json", + CorrelationId: result.ID, + Body: body, + }); err != nil { + return fmt.Errorf("cannot answer a build request: %w", err) + } + if result.Failed != "" || result.Commit == "" { + return nil + } + return EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyModuleBuilt, "builder", b.on, + announcementOf(result)) +} + +func (b *currentBuild) Done() error { return b.delivery.Ack(false) } + +func (b *currentBuild) Hold(time.Duration) error { + // No delayed redelivery on this bus: handed back at once, which is what it has always done. + return b.delivery.Nack(false, true) +} + +// announcementOf is what the mesh is told about a finished build. One function, so the two +// transports cannot describe the same build differently. +func announcementOf(result BuildResult) map[string]any { + return map[string]any{ + "module": ModuleOf(result.Manifest), "commit": result.Commit, + "repository": result.Repository, "path": result.Path, "ref": result.Ref, + "manifest": json.RawMessage(result.Manifest), "against": result.Against, + "made": result.Made, + } +} diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go new file mode 100644 index 0000000..039d340 --- /dev/null +++ b/internal/link/builds_nats.go @@ -0,0 +1,206 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The build flow on the bus being built. +// +// **One publish where the old bus needed two** (novox/hq ADR 0121). There, the answer went to a +// reply queue and the announcement to an events exchange, because the two audiences were reached by +// two topologies. Here the outcome is the role's own event: whoever asked matches it by the id their +// request carried, the controller records it, the catalogue places it in the graph. So a build +// machine publishes once and needs permission for nothing but its own role's subjects — no reply +// queue to declare, and no grant over anybody's inbox. + +// natsBuilds asks for builds over a connection. +type natsBuilds struct { + js *broker.JetStream + owned bool +} + +// BuildsOverNATS is the asking side on the bus being built. It dials, because the command that asks +// for a build is a one-shot and holds nothing else. +func BuildsOverNATS(address string) (Builders, error) { + js, err := broker.Dial(address) + if err != nil { + return nil, fmt.Errorf("cannot reach the bus at %s to ask for a build: %w", address, err) + } + return &natsBuilds{js: js, owned: true}, nil +} + +func (b *natsBuilds) Close() { + if b.owned && b.js != nil { + b.js.Close() + } +} + +func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest, + wait time.Duration) (BuildResult, error) { + + // Subscribed before the ask, so an outcome cannot arrive before there is anywhere for it to + // land. Core, not the stream: the asker is waiting now, and the durable copy of this outcome is + // the same event on EVENTS, which the controller records. + outcomes, err := b.js.Conn().SubscribeSync(BuildOutcome()) + if err != nil { + return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err) + } + defer func() { _ = outcomes.Unsubscribe() }() + if err := b.js.Conn().Flush(); err != nil { + return BuildResult{}, err + } + + body, err := json.Marshal(request) + if err != nil { + return BuildResult{}, err + } + // Into the role's work queue and awaited: work the bus never accepted must fail here rather than + // be assumed, because nothing else will ever say so. + publish, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil { + return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err) + } + + waiting, cancelWait := context.WithTimeout(ctx, wait) + defer cancelWait() + for { + msg, err := outcomes.NextMsgWithContext(waiting) + switch { + case errors.Is(err, context.DeadlineExceeded): + return BuildResult{}, waitingFor(wait) + case errors.Is(err, context.Canceled): + return BuildResult{}, ctx.Err() + case err != nil: + return BuildResult{}, fmt.Errorf("waiting for a build's outcome: %w", err) + } + result, mine, err := theOutcomeOf(msg.Data, request.ID) + if err != nil { + return BuildResult{}, err + } + if mine { + return result, nil + } + } +} + +// --- the machine's side --------------------------------------------------------------------- + +type natsMachine struct { + js *broker.JetStream + on string + seat string + sub *nats.Subscription +} + +// MachineOverNATS takes build work from the role this machine holds. +func MachineOverNATS(js *broker.JetStream, on string) BuildMachine { + return &natsMachine{js: js, on: on, seat: TheBuildMachine} +} + +func (m *natsMachine) Close() { + if m.sub != nil { + _ = m.sub.Unsubscribe() + } +} + +// Take binds to the role's worker and hands each request over, one at a time. +// +// **Bound, never created.** The work queue and the worker on it are the controller's to define +// (design 25 §3), and a build machine reaches no part of the JetStream API — so a missing one is said +// as the mesh's to answer rather than quietly created with whatever this client defaults to. +func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) error { + worker, found := broker.HolderConsumerFor(m.on, "builder", + broker.DeclaredSeat{Name: m.seat, Accepts: []string{"build"}}) + if !found { + return fmt.Errorf("%s accepts no work, so there is nothing for this machine to take", m.seat) + } + + // One at a time, which the consumer's own ack-pending limit enforces rather than a prefetch + // setting: a machine that took five requests at once would run five container builds against one + // runtime and finish all of them slower than the first. + work := make(chan *nats.Msg, 1) + // **The consumer's own filter, not the one subject this machine cares about.** The client checks + // what is asked for against the consumer's filter and refuses anything that is not the same — + // "subject does not match consumer" — so subscribing `…accept.build` against a consumer filtered + // on `…accept.>` is rejected even though it is narrower. Learned twice now, on two different + // consumers, which is why it is written down here. + filter := worker.Filters[0] + sub, err := m.js.Context().ChanQueueSubscribe(filter, worker.Queue, work, + nats.Bind(worker.Stream, worker.Name), nats.ManualAck()) + if err != nil { + return fmt.Errorf( + "this machine cannot take work from %s: %w. The mesh creates that queue and this "+ + "machine's worker on it, and a build machine may not create one itself — so this is "+ + "the mesh's to answer, not this machine's", m.seat, err) + } + m.sub = sub + + for { + select { + case <-ctx.Done(): + return nil + case msg, ok := <-work: + if !ok { + return errors.New("the bus stopped delivering build work") + } + var request BuildRequest + if err := json.Unmarshal(msg.Data, &request); err != nil { + // Unreadable: terminated rather than retried, because the next attempt reads the same + // bytes. Nobody waiting hears an answer, which is right — there was no request. + _ = msg.Term() + continue + } + do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js}) + } + } +} + +type natsBuild struct { + request BuildRequest + msg *nats.Msg + on string + js *broker.JetStream +} + +func (b *natsBuild) Request() BuildRequest { return b.request } + +// Announce publishes the outcome as the role's own event, once, for all three audiences. +// +// Into the stream, so a controller that was restarting still records it and a catalogue that was +// down still catches up. The asker is listening on core for the same subject and gets it either way: +// a stream delivers to its durable consumers and the plain subscribers both. +func (b *natsBuild) Announce(ctx context.Context, result BuildResult) error { + // One body, three readers. Whoever asked matches the id; the controller records it; the catalogue + // needs to know what was built, which only the manifest says — so the result carries it rather + // than a second message carrying a second shape. + // + // **A failed build names no module**, because it produced no module version and the catalogue + // would otherwise put something in the graph that was never made. The asker still gets its + // answer: a failure is the answer. + if result.Failed == "" && result.Commit != "" { + result.Module = ModuleOf(result.Manifest) + } + body, err := json.Marshal(result) + if err != nil { + return err + } + publish, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + if _, err := b.js.Context().Publish(BuildOutcome(), body, nats.Context(publish)); err != nil { + return fmt.Errorf("cannot announce a build's outcome: %w", err) + } + return nil +} + +func (b *natsBuild) Done() error { return b.msg.Ack() } + +func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) } diff --git a/internal/link/builds_nats_test.go b/internal/link/builds_nats_test.go new file mode 100644 index 0000000..2672258 --- /dev/null +++ b/internal/link/builds_nats_test.go @@ -0,0 +1,215 @@ +package link + +import ( + "context" + "encoding/json" + "io" + "log" + "os" + "sync" + "testing" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// A build, end to end, against a real server. +// +// The claim worth checking is the one ADR 0121 rests on: **one publish reaches three audiences**. +// Whoever asked matches the outcome by the id their request carried; the controller records it; the +// catalogue places it. On the old bus that fan-out came from a shared exchange, and it would be easy +// to write a version where only the asker hears it and nobody notices for weeks. +// +// docker run -d --rm --name t -p 14230:4222 nats:2.10-alpine -js +// MESH_TEST_NATS=nats://127.0.0.1:14230 go test ./internal/link/ -run TestNatsABuild + +func aBusWithTheBuildRole(t *testing.T) *broker.JetStream { + t.Helper() + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + js, err := broker.Dial(url) + if err != nil { + t.Fatal(err) + } + t.Cleanup(js.Close) + + seats := []broker.DeclaredSeat{{Name: TheBuildMachine, Accepts: []string{"build"}, + Emits: []string{"built"}}} + if err := broker.AssertMeshStreams(js); err != nil { + t.Fatal(err) + } + if err := broker.RaiseSeats(js, seats, map[string]broker.Holder{ + TheBuildMachine: {Node: "anchor", Module: "builder"}, + }); err != nil { + t.Fatal(err) + } + clean := func() { + _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") + for _, s := range broker.MeshStreams() { + _ = js.Context().PurgeStream(s.Name) + } + } + t.Cleanup(clean) + return js +} + +// The whole round trip: asked, taken, built, and the outcome heard by the asker and by a consumer of +// the role's event who never asked for anything. +func TestNatsABuildIsTakenAndItsOutcomeReachesEverybody(t *testing.T) { + js := aBusWithTheBuildRole(t) + ctx, stop := context.WithCancel(context.Background()) + defer stop() + + // A third party on the role's event — what the catalogue is. Subscribed first, so nothing is + // missed. + heard := make(chan BuildResult, 4) + watching, err := js.Conn().Subscribe(BuildOutcome(), func(msg *nats.Msg) { + var r BuildResult + if json.Unmarshal(msg.Data, &r) == nil { + heard <- r + } + }) + if err != nil { + t.Fatal(err) + } + defer func() { _ = watching.Unsubscribe() }() + _ = js.Conn().Flush() + + // A build machine holding the role. + machine := MachineOverNATS(js, "anchor") + defer machine.Close() + failed := make(chan error, 1) + go func() { + failed <- machine.Take(ctx, func(ctx context.Context, work Build) { + r := work.Request() + _ = work.Announce(ctx, BuildResult{ + ID: r.ID, Repository: r.Repository, On: "anchor", Commit: "abc1234", + Manifest: json.RawMessage(`{"module":"shop"}`), + }) + _ = work.Done() + }) + }() + + ask := &natsBuilds{js: js} + result, err := ask.Submit(ctx, BuildRequest{ID: "b-1", Repository: "/r"}, 15*time.Second) + if err != nil { + select { + case why := <-failed: + t.Fatalf("the machine could not take work: %v", why) + default: + } + t.Fatalf("the asker never got an outcome: %v", err) + } + if result.ID != "b-1" || result.Commit != "abc1234" { + t.Fatalf("the asker got %+v", result) + } + // Named in the outcome, because only the manifest says what was built and the catalogue reading + // this event needs to know. + if result.Module != "shop" { + t.Errorf("the outcome names module %q, so a catalogue reading it cannot place the build", + result.Module) + } + + select { + case also := <-heard: + if also.ID != "b-1" { + t.Fatalf("a third party heard %+v", also) + } + case <-time.After(5 * time.Second): + t.Fatal("nobody but the asker heard the outcome, so the catalogue would never place the build") + } + + // And the work left the queue: a request a machine took and settled must not be given to another. + deadline := time.Now().Add(5 * time.Second) + for time.Now().Before(deadline) { + info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE") + if err == nil && info.State.Msgs == 0 { + return + } + time.Sleep(20 * time.Millisecond) + } + t.Fatal("the request is still queued after being settled, so another machine would build it again") +} + +// Work submitted with no machine holding the role waits rather than failing, and is done when one +// arrives. **That is what a queue is for**, and the alternative — refusing because nobody is there +// yet — would make installing a build machine an ordering problem. +func TestNatsABuildWaitsForAMachineRatherThanFailing(t *testing.T) { + js := aBusWithTheBuildRole(t) + + body, _ := json.Marshal(BuildRequest{ID: "b-2", Repository: "/r"}) + if _, err := js.Context().Publish(BuildWork(), body); err != nil { + t.Fatal(err) + } + info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE") + if err != nil || info.State.Msgs != 1 { + t.Fatalf("the work did not queue: %+v %v", info, err) + } + + // Now a machine arrives and finds it waiting. + ctx, stop := context.WithCancel(context.Background()) + defer stop() + took := make(chan string, 1) + machine := MachineOverNATS(js, "anchor") + defer machine.Close() + go func() { + _ = machine.Take(ctx, func(ctx context.Context, work Build) { + took <- work.Request().ID + _ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: "anchor", Failed: "no"}) + _ = work.Done() + }) + }() + select { + case id := <-took: + if id != "b-2" { + t.Fatalf("the machine took %q", id) + } + case <-time.After(10 * time.Second): + t.Fatal("a machine that arrived after the work did never got it, so the backlog was lost") + } +} + +// A machine that dies before saying anything leaves the work for another. +func TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue(t *testing.T) { + js := aBusWithTheBuildRole(t) + ctx, stop := context.WithCancel(context.Background()) + defer stop() + + body, _ := json.Marshal(BuildRequest{ID: "b-3", Repository: "/r"}) + if _, err := js.Context().Publish(BuildWork(), body); err != nil { + t.Fatal(err) + } + + var once sync.Once + handed := make(chan struct{}, 2) + machine := MachineOverNATS(js, "anchor") + defer machine.Close() + go func() { + _ = machine.Take(ctx, func(ctx context.Context, work Build) { + handed <- struct{}{} + // The first time, hand it straight back — a machine that stopped mid-build. + var settled bool + once.Do(func() { _ = work.Hold(200 * time.Millisecond); settled = true }) + if !settled { + _ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: "anchor"}) + _ = work.Done() + } + }) + }() + + for i := 0; i < 2; i++ { + select { + case <-handed: + case <-time.After(10 * time.Second): + t.Fatalf("the work was handed over %d time(s); unanswered work must come back", i) + } + } +} + +func quietLog() *log.Logger { return log.New(io.Discard, "", 0) } + +var _ = quietLog diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index 3acac49..24ab322 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -180,6 +180,11 @@ func kindOfSubject(subject string) (string, bool) { return KindModuleMoved, true case broker.ControllerFollows[1]: return KindCatchUp, true + case BuildOutcome(): + // A build's outcome is the role's event now, so it arrives on the events stream rather than + // the control branch — and is acted on by the same handler, because what the controller does + // with it did not change (novox/hq ADR 0121). + return KindBuilt, true } return "", false }