A controller that asked node-build-agent from its first run would queue every build where nothing pulls, and the build that registers build-agent — the first holder — would be among them. So the role is chosen at ask time from the catalogue: the current role when any assigned module claims it, the retired one while only the builder does, the current one when neither. Outcomes are followed on both seats, the controller may publish to both, and a build's log is read under whichever role did it; a machine on the retired role is proven on the bus to take that role's asks. The switch order is written where the role is named, and the retired half is marked for removal with the seat row.
304 lines
12 KiB
Go
304 lines
12 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
|
|
// 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 {
|
|
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 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}
|
|
}
|
|
|
|
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
|
|
}
|
|
// 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.
|
|
fetched, err := sub.Fetch(1, nats.Context(ctx))
|
|
switch {
|
|
case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded):
|
|
return nil
|
|
case errors.Is(err, nats.ErrTimeout):
|
|
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 {
|
|
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
|
|
}
|
|
// 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, &natsBuild{request: request, msg: msg, on: m.on, js: m.js, seat: m.seat})
|
|
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) }
|