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.
255 lines
9.2 KiB
Go
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) }
|