Files
mesh-controller/internal/link/builds_nats.go
T
jschoubben 076e0ae259 A build is taken in where its outcome is heard, and the build tool answers at once (issue 176)
The console's `build` tool answered "no build machine answered within 0s", handed a forge path to
git as written, and a build heard afterwards was recorded and never registered: recording and
registration lived only in the waiting caller, and the tool did not wait.

Now one function takes a build's outcome in — records it, parses the manifest, refuses a definition
naming an installation, registers the module with its source as the seat and path the request
carried — and both the waiting command and the daemon that follows the role's `built` event call
it. `build --wait 0` asks and returns with the id; `builds --log <id>` follows it. The seat verb
says `--self` for a repository given without a scheme.
2026-10-01 01:27:04 +02:00

255 lines
9.2 KiB
Go

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()
}
}
// 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(BuildWork(), 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(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
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(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() }
// 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(BuildStarted(), 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(BuildLog(b.request.ID), body)
}
func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) }