With the worker consumer's default of many deliveries in flight, every ask behind the one being built was delivered at once, left unacknowledged for the length of the build, redelivered after the ack wait and dropped after the fifth time: on 2026-10-01 twenty-six of forty-three builds asked in two minutes were never built and the queue read as empty (hq issue 186). The holder's worker now has one in flight, and a running build tells the bus it is still working, as the controller's long handlers do, so a build longer than the ack wait is neither redelivered nor counted out.
261 lines
9.5 KiB
Go
261 lines
9.5 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
|
|
}
|
|
// 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})
|
|
close(working)
|
|
}
|
|
}
|
|
}
|
|
|
|
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) }
|