Every one of the 48 core failures of research 031 was found by a person looking; the mesh's answers carried the fact for whoever asked and told nobody. - The condition store (to-be 45 §2): mesh-controller_conditions, one key per open condition, written by compare-and-set so a person's silence and the watchdogs never lose each other's word; every transition kept ninety days in mesh-controller_condition-history and said as the seat's events condition-raised / condition-changed / condition-cleared (the condition at the top level, with event, at, change, why, show), offered again while the bus is away. Raised and cleared by observation only; a clearing reopened within ten minutes is the same condition with its count up, its silence kept. Verbs: conditions, conditions show, conditions silence (a hand act, at most a week), conditions history. - ADR 0224's provider standing is the first kind, provider-failing, held by the provider's events; the provider_standing table is no longer read or written (left in place: dropping it is the operator's word). - status leads with the open conditions, urgent first, and says all well only with none open; conditions it cannot read are said and not well. - The signals table compiled in, one watchdog loop over it every 30s: S1 heartbeat (3 intervals, asleep machines excepted, control node urgent after 30 min), S2 report after a send, S3 plan tier, S4 event loop deaf, S5 merge not acted, S6 ask lost, S7 call hung, S8 provider silent, S9 advisories, S10 self-check silent, S11 node tools silent, S13 stale refusals; S12, S14, S15 deferred with their reasons. A row that cannot see raises probe-failed and clears nothing. A test generated from the table suppresses each signal inside and past its bound. - The bus's advisories (maximum deliveries, a mesh consumer deleted) and the controller's own slow consumer and refused subjects, said in the mesh's words. - doctor: the probe registry D1-D10 (D5 deferred) and DW, every five minutes, each in thirty seconds; a probe that cannot run is never a pass. D1 validates with mesh-host's own validator. Every run ends with the doctor-heartbeat event mesh-watcher listens for. - The controller is granted its new buckets, events, the two advisories and $SRV.INFO; the node tools their tools-alive heartbeat. The streams and consumers the controller asserts and the ones D6/D7 expect are one derivation.
415 lines
16 KiB
Go
415 lines
16 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
)
|
|
|
|
// The consume side on the bus being built.
|
|
//
|
|
// The shape is the AMQP one's, because the seam made them comparable: one loop, one message at a
|
|
// time, and the same window deciding. What differs is where a held message lives — and that is the
|
|
// whole point of the move. On the bus the mesh has, holding one means keeping an unacknowledged
|
|
// delivery in this process, bounded by the prefetch and lost if the controller stops. Here it is a
|
|
// `nak` with a delay: the message stays the server's, the controller keeps nothing but the moment
|
|
// it first could not take it, and a controller that restarts mid-window has nothing to lose.
|
|
|
|
// natsInbound consumes what nodes and modules say over NATS.
|
|
// Prefetch is how many controls the controller holds unacknowledged at once: the store window
|
|
// (ADR 0083) is the server's, and this is what it may hand this process ahead of its acting.
|
|
const Prefetch = 64
|
|
|
|
type natsInbound struct {
|
|
js *broker.JetStream
|
|
// follows is the kinds asked for beyond what nodes say (Also). The events those are are the
|
|
// only ones the controller subscribes, and only when something is listening.
|
|
follows map[string]bool
|
|
// since is when the controller first could not take a message, by that message's place in its
|
|
// stream.
|
|
//
|
|
// **A timestamp, not a message.** This is the whole difference the move buys: the AMQP side
|
|
// keeps the delivery, and this keeps eight bytes saying when the window opened. A controller
|
|
// that restarts loses these and starts the window again, which is correct — it is holding
|
|
// nothing, and the messages are all still on the server.
|
|
since map[uint64]time.Time
|
|
}
|
|
|
|
// Nats is the consume side of the bus being built.
|
|
func Nats(js *broker.JetStream) Inbound {
|
|
return &natsInbound{js: js, follows: map[string]bool{}, since: map[uint64]time.Time{}}
|
|
}
|
|
|
|
// Also records one more kind to subscribe. Nothing is subscribed here: the controller's consumer on
|
|
// the events stream carries both of these as filters, so it is created once, in Receive, with
|
|
// whatever was asked for — and not at all when nothing was.
|
|
func (n *natsInbound) Also(kind string) error {
|
|
switch kind {
|
|
case KindModuleMoved, KindCatchUp, KindSourceMoved, KindProvisioner:
|
|
n.follows[kind] = true
|
|
return nil
|
|
default:
|
|
return fmt.Errorf("nothing subscribes %s separately on this bus", kind)
|
|
}
|
|
}
|
|
|
|
func (n *natsInbound) Close() {}
|
|
|
|
// Receive consumes until the context ends.
|
|
//
|
|
// Three subscriptions, and each is a channel the one loop selects on. **Channels rather than
|
|
// callbacks**: the library would run a handler on its own goroutine, and the window's bookkeeping —
|
|
// which message is held, and since when — is read and written without a lock because the AMQP loop
|
|
// never had two. A second goroutine would make that wrong in a way no test would catch.
|
|
func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Control)) error {
|
|
if err := broker.AssertMeshConsumers(n.js); err != nil {
|
|
return err
|
|
}
|
|
js, conn := n.js.Context(), n.js.Conn()
|
|
|
|
// What nodes say, off the CONTROL stream. Bound to the durable the controller asserted rather
|
|
// than creating one here: the consumer is an object with a configuration — ack policy, ack
|
|
// wait, redelivery — and a client that creates its own would be a second opinion about it.
|
|
control := make(chan *nats.Msg, Prefetch)
|
|
said, err := standingBy(ctx, log.Default(), "CONTROL", func() (*nats.Subscription, error) {
|
|
return js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName))
|
|
})
|
|
if err != nil {
|
|
return fmt.Errorf("subscribing to what nodes say: %w", err)
|
|
}
|
|
if said == nil {
|
|
return nil // stopped while standing by
|
|
}
|
|
defer func() { _ = said.Unsubscribe() }()
|
|
// This controller is the one acting now: its watchdogs and self-check may say what they see
|
|
// (novox/hq to-be 45 §3). One standing by hears nothing, and would call every machine silent.
|
|
holding.Store(true)
|
|
defer holding.Store(false)
|
|
|
|
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
|
|
// they are their own guarantee: a lost one is the next one.
|
|
beats := make(chan *nats.Msg, Prefetch)
|
|
alive, err := conn.ChanSubscribe(AliveSubjects, beats)
|
|
if err != nil {
|
|
return fmt.Errorf("subscribing to heartbeats: %w", err)
|
|
}
|
|
defer func() { _ = alive.Unsubscribe() }()
|
|
// And each machine's node tools, on the same channel: as cheap, and lost the same way.
|
|
toolsAlive, err := conn.ChanSubscribe(ToolsAliveSubjects, beats)
|
|
if err != nil {
|
|
return fmt.Errorf("subscribing to the node tools' heartbeats: %w", err)
|
|
}
|
|
defer func() { _ = toolsAlive.Unsubscribe() }()
|
|
|
|
// The events the controller follows, when something is listening for them.
|
|
var events chan *nats.Msg
|
|
if len(n.follows) > 0 {
|
|
events = make(chan *nats.Msg, Prefetch)
|
|
followed, err := standingBy(ctx, log.Default(), "EVENTS", func() (*nats.Subscription, error) {
|
|
return js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName))
|
|
})
|
|
if err != nil {
|
|
return fmt.Errorf("subscribing to what the catalogue says: %w", err)
|
|
}
|
|
if followed == nil {
|
|
return nil
|
|
}
|
|
defer func() { _ = followed.Unsubscribe() }()
|
|
}
|
|
|
|
// A connection that dropped is said, not discovered. A controller whose bus connection is gone
|
|
// is a mesh where nothing can be told anything.
|
|
gone := make(chan error, 1)
|
|
conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) {
|
|
select {
|
|
case gone <- err:
|
|
default:
|
|
}
|
|
})
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case err := <-gone:
|
|
return fmt.Errorf("the bus connection dropped: %w", err)
|
|
case msg := <-beats:
|
|
n.deliver(ctx, act, msg, false)
|
|
case msg := <-events:
|
|
n.deliver(ctx, act, msg, true)
|
|
case msg, ok := <-control:
|
|
if !ok {
|
|
return errors.New("the bus stopped delivering")
|
|
}
|
|
n.deliver(ctx, act, msg, true)
|
|
}
|
|
}
|
|
}
|
|
|
|
// deliver names one message and hands it to the loop, or drops it where the mesh has no name for
|
|
// its subject — which cannot happen through a filter the controller wrote, and is said rather than
|
|
// ignored for exactly that reason.
|
|
func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control),
|
|
msg *nats.Msg, streamed bool) {
|
|
|
|
kind, known := kindOfSubject(msg.Subject)
|
|
if !known {
|
|
if streamed {
|
|
_ = msg.Term()
|
|
}
|
|
return
|
|
}
|
|
m := &natsControl{kind: kind, msg: msg, on: n}
|
|
if streamed {
|
|
// The event loop took one: what S4 watches (novox/hq to-be 45 §3).
|
|
Loop.Took(time.Now())
|
|
// A message with no metadata is not from a stream, whatever it was delivered on, and the
|
|
// window has nothing to hold it by. Said by leaving the sequence at zero.
|
|
if meta, err := msg.Metadata(); err == nil {
|
|
m.seq = meta.Sequence.Stream
|
|
m.delivered = meta.NumDelivered
|
|
}
|
|
// **Work that outlives the acknowledgement window says so while it runs.**
|
|
//
|
|
// The bus waits a fixed time to be told a message was taken, and then hands it to whoever
|
|
// consumes next — which is right for a consumer that died and wrong for one that is busy.
|
|
// Acting on a merge builds every module the merge changed: minutes of work against a
|
|
// thirty-second window. So the same merge was handed over again while the first build was
|
|
// still running, and again after that — on 2026-09-28 one merge ran the mesh's whole
|
|
// catalogue five times over and exhausted a public registry's pull limit.
|
|
//
|
|
// Here rather than in each handler, because the window belongs to the transport and every
|
|
// handler would otherwise have to remember it. It changes nothing about a handler that
|
|
// dies: a message is kept alive only while this goroutine is, so a controller that stops
|
|
// stops saying so, and the bus redelivers exactly as it should.
|
|
working := make(chan struct{})
|
|
defer close(working)
|
|
go stillWorking(msg, working)
|
|
}
|
|
act(ctx, m)
|
|
}
|
|
|
|
// heartbeatWhileWorking is how often a handler still running tells the bus so — comfortably inside
|
|
// the shortest acknowledgement window the mesh gives any of its consumers.
|
|
const heartbeatWhileWorking = 10 * time.Second
|
|
|
|
// stillWorking keeps one message alive until the work on it returns.
|
|
//
|
|
// An error is not worth reporting: what the bus does when it is not told is redeliver, which is
|
|
// exactly what happens if this fails, and the handler's own outcome is the thing worth logging.
|
|
func stillWorking(msg *nats.Msg, done <-chan struct{}) {
|
|
tick := time.NewTicker(heartbeatWhileWorking)
|
|
defer tick.Stop()
|
|
for {
|
|
select {
|
|
case <-done:
|
|
return
|
|
case <-tick.C:
|
|
_ = msg.InProgress()
|
|
}
|
|
}
|
|
}
|
|
|
|
// kindOfSubject is how this transport's addressing becomes what the mesh calls a message.
|
|
//
|
|
// By subject, which is the only thing the server enforces: a body claiming to be a report does not
|
|
// make it one, and on this bus the subject an account may publish *is* its authority (design 29
|
|
// §2). The mirror of the routing-key table on the bus the mesh has.
|
|
func kindOfSubject(subject string) (string, bool) {
|
|
switch subject {
|
|
case EnrolSubject:
|
|
return KindEnrolment, true
|
|
case BuiltSubject:
|
|
return KindBuilt, true
|
|
}
|
|
if node, rest, ok := strings.Cut(strings.TrimPrefix(subject, "mesh.control."), "."); ok &&
|
|
node != "" && !strings.Contains(node, ".") {
|
|
switch rest {
|
|
case "report":
|
|
return KindReport, true
|
|
case "alive":
|
|
return KindHeartbeat, true
|
|
case "tools-alive":
|
|
return KindToolsHeartbeat, true
|
|
}
|
|
}
|
|
switch subject {
|
|
case broker.ControllerFollows[0]:
|
|
return KindModuleMoved, true
|
|
case broker.ControllerFollows[1]:
|
|
return KindCatchUp, true
|
|
case broker.ControllerFollows[3]:
|
|
return KindSourceMoved, true
|
|
case BuildOutcome(), BuildOutcomeOf(TheBuildMachineBefore):
|
|
// A build's outcome is the role's event now, so it arrives on the events stream rather than
|
|
// the control branch — and is acted on by the same handler, because what the controller does
|
|
// with it did not change (novox/hq ADR 0121). From either build role while the handover
|
|
// runs (ADR 0190): the old builder still answers on the retired seat until it is unassigned.
|
|
return KindBuilt, true
|
|
}
|
|
if _, ok := ProvisionerEmitter(subject); ok {
|
|
return KindProvisioner, true
|
|
}
|
|
return "", false
|
|
}
|
|
|
|
// ProvisionerEmitter is the module a provider's standing event came from, read from its subject
|
|
// (`mesh.mod.<module>.event.provisioner.<failing|recovered>`); false for any other subject. The
|
|
// controller's own follow pattern, with `*` for the module, decodes too.
|
|
func ProvisionerEmitter(subject string) (string, bool) {
|
|
rest, ok := strings.CutPrefix(subject, "mesh.mod.")
|
|
if !ok {
|
|
return "", false
|
|
}
|
|
module, event, ok := strings.Cut(rest, ".event.")
|
|
if !ok || module == "" || strings.Contains(module, ".") {
|
|
return "", false
|
|
}
|
|
if event != broker.ProvisionerFailing && event != broker.ProvisionerRecovered {
|
|
return "", false
|
|
}
|
|
return module, true
|
|
}
|
|
|
|
// natsControl is one message from the bus being built, as the controller reads it.
|
|
type natsControl struct {
|
|
kind string
|
|
msg *nats.Msg
|
|
on *natsInbound
|
|
// seq is this message's place in its stream; zero for a core message, which has none and
|
|
// cannot be held.
|
|
seq uint64
|
|
// delivered is how many times the server has handed this message over, this time included.
|
|
delivered uint64
|
|
}
|
|
|
|
func (m *natsControl) Kind() string { return m.kind }
|
|
func (m *natsControl) Body() []byte { return m.msg.Data }
|
|
func (m *natsControl) Subject() string { return m.msg.Subject }
|
|
|
|
// Redelivered is what the server counted, not what the controller remembers. Which is the answer to
|
|
// a question the AMQP side could only guess at across a restart: an enrolment redelivered because
|
|
// the controller stopped mid-answer reads as redelivered to the controller that comes back.
|
|
func (m *natsControl) Redelivered() bool { return m.delivered > 1 }
|
|
|
|
func (m *natsControl) HeldFor() time.Duration {
|
|
if m.seq == 0 {
|
|
return 0
|
|
}
|
|
first, held := m.on.since[m.seq]
|
|
if !held {
|
|
return 0
|
|
}
|
|
return time.Since(first)
|
|
}
|
|
|
|
// Answer publishes to the reply subject the request carries **in its payload**.
|
|
//
|
|
// Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had
|
|
// that field claimed for the consumer's own ack address, so answering it would send the reply to
|
|
// `$JS.ACK.CONTROL.controller.…` and the enrolling node would wait out its timeout. Verified
|
|
// against a running server (design 25 §2), which is why it is a field of the request and this reads
|
|
// it from there.
|
|
func (m *natsControl) Answer(ctx context.Context, body []byte) error {
|
|
var addressed replyAddressed
|
|
if err := json.Unmarshal(m.msg.Data, &addressed); err != nil {
|
|
return fmt.Errorf("that request cannot be read, so its reply address cannot be: %w", err)
|
|
}
|
|
if addressed.ReplyTo == "" {
|
|
return errors.New("that request named no reply subject in its payload, so nothing can be " +
|
|
"told the answer")
|
|
}
|
|
return m.on.js.Conn().PublishMsg(&nats.Msg{Subject: addressed.ReplyTo, Data: body})
|
|
}
|
|
|
|
func (m *natsControl) Took() error {
|
|
m.forget()
|
|
if m.seq == 0 {
|
|
// Core NATS: nothing is keeping it, so there is nothing to settle.
|
|
return nil
|
|
}
|
|
return m.msg.Ack(nats.Context(context.Background()))
|
|
}
|
|
|
|
// Drop terminates the delivery: understood, and the server is told not to send it again. Different
|
|
// from an ack only in the server's own accounting, which is where somebody asking "what happened to
|
|
// that message" will look.
|
|
func (m *natsControl) Drop() error {
|
|
m.forget()
|
|
if m.seq == 0 {
|
|
return nil
|
|
}
|
|
return m.msg.Term()
|
|
}
|
|
|
|
// Hold hands the message back with a delay, and remembers when the window opened.
|
|
func (m *natsControl) Hold(after time.Duration) error {
|
|
if m.seq == 0 {
|
|
return errors.New("a message that is not in a stream cannot be held: nothing is keeping it")
|
|
}
|
|
if _, already := m.on.since[m.seq]; !already {
|
|
m.on.since[m.seq] = time.Now()
|
|
}
|
|
return m.msg.NakWithDelay(after)
|
|
}
|
|
|
|
func (m *natsControl) forget() {
|
|
if m.on != nil && m.seq != 0 {
|
|
delete(m.on.since, m.seq)
|
|
}
|
|
}
|
|
|
|
// replyAddressed is the one field every message that expects an answer carries.
|
|
type replyAddressed struct {
|
|
ReplyTo string `json:"reply_to,omitempty"`
|
|
}
|
|
|
|
// StandbyPoll is how often a controller standing by looks again for its consumers. A variable so a
|
|
// test need not wait.
|
|
var StandbyPoll = 2 * time.Second
|
|
|
|
// standingBy binds one of the controller's consumers, waiting while another controller holds it.
|
|
//
|
|
// **Two controllers, one consumer** (novox/hq issue 213). The controller's consumers are push
|
|
// consumers with no delivery group, so the server lets one subscription bind each — on purpose:
|
|
// two would each act on every message (issue 146). When a machine hands its controller over from
|
|
// the container to the process, the host starts the process first and removes the container only
|
|
// once the process is up; the process then finds the consumers bound. Exiting on that would never
|
|
// be up, so the container would never go. It stands by instead — the seat's verbs are already
|
|
// served from a queue group, and the plans wait on their lock — and binds as soon as the other lets
|
|
// go. Nil and no error is ctx ending while it waited.
|
|
func standingBy(ctx context.Context, logger interface{ Printf(string, ...any) }, stream string,
|
|
bind func() (*nats.Subscription, error)) (*nats.Subscription, error) {
|
|
said := false
|
|
for {
|
|
sub, err := bind()
|
|
if err == nil {
|
|
if said {
|
|
logger.Printf("took the controller's consumer on %s: the controller that held it let go", stream)
|
|
}
|
|
return sub, nil
|
|
}
|
|
if !strings.Contains(err.Error(), "already bound") {
|
|
return nil, err
|
|
}
|
|
if !said {
|
|
logger.Printf("another controller holds the controller's consumer on %s; standing by "+
|
|
"until it lets go", stream)
|
|
said = true
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, nil
|
|
case <-time.After(StandbyPoll):
|
|
}
|
|
}
|
|
}
|