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() } } 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 } 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() } func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) }