package link import ( "context" "encoding/json" "errors" "fmt" "os" "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 // seat is the build role asked: the one that has a holder (ADR 0190 handover), chosen by the // controller from what is assigned, so an ask lands where a machine is pulling. seat string } // BuildsOverNATS is the asking side on the bus being built, asking the current build role. It dials, // because the command that asks for a build is a one-shot and holds nothing else. func BuildsOverNATS(address string) (Builders, error) { return BuildsOverNATSOn(address, TheBuildMachine) } // BuildsOverNATSOn is the asking side for one named build role — during the handover from the one // build machine to build agents, the role that has a holder (ADR 0190). func BuildsOverNATSOn(address, seat 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, seat: seat}, nil } // role is the seat asked: what the asker was made for, or the current build role for one made // without saying (a test building the struct by hand). func (b *natsBuilds) role() string { if b.seat == "" { return TheBuildMachine } return b.seat } func (b *natsBuilds) Close() { if b.owned && b.js != nil { b.js.Close() } } // Ask publishes the work and returns; see Builders. func (b *natsBuilds) Ask(ctx context.Context, request BuildRequest) error { body, err := json.Marshal(request) if err != nil { return err } publish, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil { return fmt.Errorf("cannot submit a build: %w", err) } return nil } 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(BuildOutcomeOf(b.role())) 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(BuildWorkOf(b.role()), 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 { // **A wait that hears a cancel** (novox/hq ADR 0219). A person cancelling an ask records its // failure without publishing the role's outcome — that subject is the holders', and the // controller may not speak for them — so the waiter also looks, every while, at the cancelled // set the cancel wrote first. A kill needs no look: the holder announces it as any outcome. slice, endSlice := context.WithTimeout(waiting, cancelLook) msg, err := outcomes.NextMsgWithContext(slice) endSlice() if errors.Is(err, context.DeadlineExceeded) && waiting.Err() == nil { if cancelled, _ := IsCancelled(b.js.Conn(), b.role(), request.ID); cancelled { return BuildResult{ID: request.ID, Repository: request.Repository, Path: request.Path, Ref: request.Ref, Source: request.Source, DryRun: request.DryRun, Failed: CancelledByHand}, nil } continue } 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 } } } // cancelLook is how often a waiting asker looks whether its ask was cancelled. var cancelLook = 2 * time.Second // --- the machine's side --------------------------------------------------------------------- type natsMachine struct { js *broker.JetStream on string seat string sub *nats.Subscription opts MachineOptions } // MachineOptions is what a holder tells the taking loop about itself (novox/hq ADR 0219). type MachineOptions struct { // Paused is asked before every fetch: while it says so, nothing new is taken, and a build // already running finishes. Nil is never paused. Paused func() bool } // MachineOverNATS takes build work from the current build role. func MachineOverNATS(js *broker.JetStream, on string) BuildMachine { return MachineOverNATSOn(js, on, TheBuildMachine) } // MachineOverNATSOn takes build work from the role named — the one this machine's credential claims // (ADR 0190 handover): its asks come from that seat's worker, and what it says about a build goes // out as that seat's events, so an outcome is heard where the asker listens. func MachineOverNATSOn(js *broker.JetStream, on, seat string) BuildMachine { return &natsMachine{js: js, on: on, seat: seat} } // MachineOverNATSWith is MachineOverNATSOn for a holder that can be paused (novox/hq ADR 0219). func MachineOverNATSWith(js *broker.JetStream, on, seat string, opts MachineOptions) BuildMachine { return &natsMachine{js: js, on: on, seat: seat, opts: opts} } // pausedPoll is how often a paused holder looks again whether it was resumed, and pausedFetch how // long one pull of a holder that can be paused waits. const ( pausedPoll = 2 * time.Second pausedFetch = 5 * time.Second ) func (m *natsMachine) Close() { if m.sub != nil { _ = m.sub.Unsubscribe() } } // Take binds to the role's worker and pulls one request at a time, handing each over. // // **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. // // **Pulled, one at a time, by whichever holder is free** (novox/hq ADR 0190). Every machine holding // the role binds this same worker; a machine asks for the next request only when it has finished // the last, so a slow machine never holds an ask an idle one could take, and a machine that took // five at once would run five container builds against one runtime and finish all of them slower // than the first. func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) error { worker, found := broker.HolderConsumerFor(m.on, "build-agent", 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) } // **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().PullSubscribe(filter, worker.Name, 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 { if ctx.Err() != nil { return nil } // **Paused takes nothing new** (novox/hq ADR 0219). Asked before the fetch, never during a // build: what this machine already took it finishes, and what it has not taken stays in the // queue for another holder — or for this one, resumed. if m.opts.Paused != nil && m.opts.Paused() { select { case <-ctx.Done(): return nil case <-time.After(pausedPoll): } continue } // One, and wait a while for it; an empty queue is a timeout, which is the normal state of a // machine with nothing to build, and is asked again. // Asked for a few seconds at a time when the holder can be paused, so a pause reaches a pull // already waiting within that, rather than when the client's own wait runs out. asking, endAsking := ctx, func() {} if m.opts.Paused != nil { asking, endAsking = context.WithTimeout(ctx, pausedFetch) } fetched, err := sub.Fetch(1, nats.Context(asking)) endAsking() switch { case ctx.Err() != nil: // Ours ended: the machine is being stopped. return nil case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout): // **An empty queue, not the end.** A fetch on a context without a deadline waits the // client's own while and then says the deadline passed — the client's, not ours. Read // as "stop", every idle build machine exited clean every half minute and was started // again by its supervisor, which looked like a crash loop with nothing in the log to // say why (2026-10-03, the first build agents). Asked again. continue case err != nil: if sub.IsValid() { // A transient fault in asking — a reconnect, a slow server — is asked past rather // than ending the machine; one that outlasts the ack wait redelivers nothing lost. time.Sleep(time.Second) continue } return fmt.Errorf("the bus stopped delivering build work: %w", err) } for _, msg := range fetched { // **Paused while the pull was answered** (ADR 0219): "takes no new build" holds even for // the one that arrived in that moment. Handed back at once for another holder — counted as // a delivery, which a pause landing exactly then costs and nothing else does. if m.opts.Paused != nil && m.opts.Paused() { _ = msg.Nak() continue } 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 } // **Cancelled in the moment it was fetched** (novox/hq ADR 0219): the controller wrote the // id into the seat's cancelled set before deleting the ask, so one taken in between is // ended here, terminated rather than redelivered, and its outcome said as failed — never // built. A set that cannot be read is said and the ask built: a cancel is a person's // exception, and a holder that refused every build while its grant was missing would // stop the mesh building for it. build := &natsBuild{request: request, msg: msg, on: m.on, js: m.js, seat: m.seat} if cancelled, err := IsCancelled(m.js.Conn(), m.seat, request.ID); err != nil { fmt.Fprintf(os.Stderr, "could not read whether %s was cancelled, so it is built: %v\n", request.ID, err) } else if cancelled { endCancelled(ctx, build) continue } // A build outlives the acknowledgement window many times over; said while it runs, // as the controller says it for its own long handlers, so the server neither hands // the ask to a second machine nor counts the wait against its deliveries. working := make(chan struct{}) go stillWorking(msg, working) do(ctx, build) close(working) } } } type natsBuild struct { request BuildRequest msg *nats.Msg on string js *broker.JetStream // seat is the role this build was taken from; what the machine says about it is that role's. seat string seq int } 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(BuildOutcomeOf(b.seat), 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() } // Began publishes that this machine has taken the build, into the stream like the outcome, so a // reader that asks afterwards sees when it started as well as how it ended. func (b *natsBuild) Began(ctx context.Context) error { body, err := json.Marshal(BuildStart{ ID: b.request.ID, Repository: b.request.Repository, Path: b.request.Path, Ref: b.request.Ref, On: b.on, At: time.Now().UTC().Format(time.RFC3339Nano), }) if err != nil { return err } publish, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() if _, err := b.js.Context().Publish(BuildStartedOf(b.seat), body, nats.Context(publish)); err != nil { return fmt.Errorf("cannot say a build started: %w", err) } return nil } // Say publishes one line under the build's id. Core publish, unawaited: the stream that holds the // role's events captures it on its way through, and a build must not slow to the pace of an ack // per line. A line the bus did not take is counted anyway, so the gap is visible to a reader. func (b *natsBuild) Say(step, message string) { b.seq++ body, err := json.Marshal(BuildLine{ ID: b.request.ID, Seq: b.seq, At: time.Now().UTC().Format(time.RFC3339Nano), Step: step, Message: message, }) if err != nil { return } _ = b.js.Conn().Publish(BuildLogOf(b.seat, b.request.ID), body) } func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) } // endCancelled settles an ask that was cancelled as it was taken: its outcome said as failed, as the // controller already recorded it, so anybody still waiting on it hears the answer — then // terminated, so the queue neither keeps nor redelivers it. func endCancelled(ctx context.Context, b *natsBuild) { r := b.request fmt.Fprintf(os.Stderr, "%s was cancelled by hand as this machine took it; not built\n", r.ID) result := BuildResult{ID: r.ID, Repository: r.Repository, Path: r.Path, Ref: r.Ref, On: b.on, Source: r.Source, DryRun: r.DryRun, Failed: CancelledByHand} if err := b.Announce(ctx, result); err != nil { fmt.Fprintf(os.Stderr, "cannot say %s was cancelled: %v\n", r.ID, err) } _ = b.msg.Term() }