The gate judged a module by what the mesh saw from outside, so a container that crash-looped after it applied passed it. Each machine's node-engine now states the health of every long-running resource it runs; the controller keeps the newest statement per machine, raises module.<module>.<machine>.unhealthy on the second statement in a row, clears it on the first that does not say it, and the gate passes a module only when every long-running resource of it is stated healthy since the send. An engine that states nothing is judged as before.
216 lines
9.6 KiB
Go
216 lines
9.6 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// Bus is what the controller needs of the mesh's bus, **in the mesh's own words rather than a
|
|
// transport's** (novox/hq ADR 0116 step 3).
|
|
//
|
|
// Until now every one of these functions took the transport's own channel type, so the transport
|
|
// reached every caller and changing it meant touching all of them. The seam is small — the
|
|
// controller sends exactly two kinds of message that expect no answer, and asks two kinds of
|
|
// question — which is why the bus can be replaced at all.
|
|
//
|
|
// Two implementations live below, and both ship until the rollout (ADR 0116: nothing moves a
|
|
// node's bus before step 5). Both shipping is what makes them comparable — the same caller, the
|
|
// same arguments, and one conformance fixture holding them to one envelope.
|
|
type Bus interface {
|
|
// PublishEvent announces something that happened, under the emitter's own name. 1:many, and
|
|
// nobody is obliged to act (ADR 0041).
|
|
PublishEvent(ctx context.Context, key, source, node string, body []byte) error
|
|
|
|
// PublishDeclaration delivers one node what it should be. Addressed to that node alone: a
|
|
// declaration is not an event, and replaying yesterday's is actively harmful
|
|
// (design 29 §4, the *state* shape).
|
|
PublishDeclaration(ctx context.Context, node string, body []byte) error
|
|
|
|
// PublishSeatEvent states a fact under a role's own name, for the holder of that role. A
|
|
// module's event is addressed to the module; a role's is addressed to the role, so it keeps
|
|
// meaning when the holder changes (novox/hq ADR 0121, ADR 0129).
|
|
PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error
|
|
|
|
// AskTool sends one question to a module's tool and awaits one answer. A tool nobody serves
|
|
// must say so **at once** rather than after the whole wait: the difference between "that
|
|
// module is down" and "that tool is slow" is the first thing a person asking wants.
|
|
AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error)
|
|
}
|
|
|
|
// --- The bus the mesh runs on today -----------------------------------------------------
|
|
|
|
// --- NATS, the bus being built ----------------------------------------------------------------
|
|
|
|
// OverNATS is the bus as a JetStream context.
|
|
type OverNATS struct {
|
|
Conn *nats.Conn
|
|
JS nats.JetStreamContext
|
|
}
|
|
|
|
// The subjects a node publishes on, and the controller listens to.
|
|
//
|
|
// One tree, and each name says who it is about: `mesh.control.<node>.…` is a node's own, which is
|
|
// what lets a node's account be granted exactly its own prefix and nothing of any other node's
|
|
// (design 25 §2, §4). The two that belong to no node — an enrolment, because a machine enrolling
|
|
// has no name the mesh has agreed to yet, and a build's outcome, because a builder is not
|
|
// reporting about itself — are named directly.
|
|
const (
|
|
// EnrolSubject is where a joining machine asks. Its enrolment user may publish here and
|
|
// nowhere else, so a leaked token buys nothing but the chance to enrol.
|
|
EnrolSubject = "mesh.control.enrol"
|
|
|
|
// BuiltSubject is where a build's outcome lands, for results nobody was waiting for.
|
|
BuiltSubject = "mesh.control.built"
|
|
|
|
// AliveSubjects is every node's heartbeat. Core NATS, never a stream: a lost heartbeat is the
|
|
// next heartbeat, and a stream of them is the mesh's least valuable message competing for
|
|
// retention with its most valuable (design 25 §3).
|
|
AliveSubjects = "mesh.control.*.alive"
|
|
|
|
// ToolsAliveSubjects is every machine's node tools saying they are there (novox/hq to-be 45 §3,
|
|
// S11): core NATS like the host's, for the same reason.
|
|
ToolsAliveSubjects = "mesh.control.*.tools-alive"
|
|
|
|
// HealthSubjects is every machine's health statement between its reports (novox/hq ADR 0240): core
|
|
// NATS like the heartbeat, because a statement lost is said again within a minute while anything is
|
|
// not healthy, and the next report carries it whatever happens.
|
|
HealthSubjects = "mesh.control.*.health"
|
|
)
|
|
|
|
// HealthSubject is one machine's health statement.
|
|
func HealthSubject(node string) string { return "mesh.control." + node + ".health" }
|
|
|
|
// ReportSubject is where one node says what it did. On the CONTROL stream, because it is the
|
|
// message the store-window guarantee is about (ADR 0083).
|
|
func ReportSubject(node string) string { return "mesh.control." + node + ".report" }
|
|
|
|
// AliveSubject is one node's heartbeat.
|
|
func AliveSubject(node string) string { return "mesh.control." + node + ".alive" }
|
|
|
|
// ToolsAliveSubject is one machine's node tools' heartbeat.
|
|
func ToolsAliveSubject(node string) string { return "mesh.control." + node + ".tools-alive" }
|
|
|
|
// EventSubject is where a module's event lands. Derived from the emitter, never taken from the
|
|
// caller: a source that could differ from the subject is an envelope that can lie about its
|
|
// origin, and on NATS the account's permissions make the subject the authority (design 29 §2).
|
|
func EventSubject(source, key string) string {
|
|
return "mesh.mod." + source + ".event." + key
|
|
}
|
|
|
|
// SeatEventSubject is where a role's own event lands. Derived from the role, never from its holder:
|
|
// a fact about the build machine or about the control plane keeps its address when the module holding
|
|
// that role is replaced (novox/hq ADR 0121, ADR 0129).
|
|
func SeatEventSubject(seat, event string) string {
|
|
return "mesh.seat." + seat + ".event." + event
|
|
}
|
|
|
|
// DeclareSubject is where one node's declaration lands. Last-per-subject on the NODES stream, so
|
|
// a node that was away gets exactly the current one and a replayed older one is refused by
|
|
// sequence — the wire-level answer to novox/hq issue 107.
|
|
func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" }
|
|
|
|
func (b OverNATS) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
|
|
id, err := eventID()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
h := nats.Header{}
|
|
h.Set("x-event-id", id)
|
|
h.Set("x-source", source)
|
|
h.Set("x-node", node)
|
|
h.Set("x-time", time.Now().UTC().Format(time.RFC3339))
|
|
h.Set("content-type", "application/json")
|
|
|
|
// The id is also the publish's message id, so the server refuses a duplicate inside its
|
|
// window. That narrows the window a consumer must deduplicate in; it does not remove the
|
|
// requirement, because the window is finite (design 19, delivery).
|
|
_, err = b.JS.PublishMsg(&nats.Msg{
|
|
Subject: EventSubject(source, key),
|
|
Header: h,
|
|
Data: body,
|
|
}, nats.MsgId(id), nats.Context(ctx))
|
|
if err != nil {
|
|
return fmt.Errorf("emitting %s: %w", key, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// PublishSeatEvent states a role's own fact. Same envelope as a module's event and a different
|
|
// address: the source header is the role, because that is what the fact is about.
|
|
func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error {
|
|
id, err := eventID()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
h := nats.Header{}
|
|
h.Set("x-event-id", id)
|
|
h.Set("x-source", seat)
|
|
_, err = b.JS.PublishMsg(&nats.Msg{
|
|
Subject: SeatEventSubject(seat, event),
|
|
Header: h,
|
|
Data: body,
|
|
}, nats.MsgId(id), nats.Context(ctx))
|
|
if err != nil {
|
|
return fmt.Errorf("stating %s of the %s seat: %w", event, seat, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// PublishMembership issues one assignment what it serves and reaches (novox/hq ADR 0160), last per
|
|
// subject, so the runtime that connects later reads the current one and one that is running follows.
|
|
// MembershipWait bounds how long issuing one membership may take. A publish the server refuses is
|
|
// never acknowledged, and a stream publish waits for its acknowledgement for as long as its
|
|
// context lives: on 2026-10-01 the daemon's own context was that long, and one refused membership
|
|
// held the controller's receive loop for good (novox/hq issue 185).
|
|
const MembershipWait = 10 * time.Second
|
|
|
|
func (b OverNATS) PublishMembership(ctx context.Context, node, module string, body []byte) error {
|
|
ctx, cancel := context.WithTimeout(ctx, MembershipWait)
|
|
defer cancel()
|
|
_, err := b.JS.Publish(broker.MembershipSubject(node, module), body, nats.Context(ctx))
|
|
if err != nil {
|
|
return fmt.Errorf("issuing %s on %s its membership: %w", module, node, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
|
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
|
|
if err != nil {
|
|
return fmt.Errorf("declaring to %s: %w", node, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ToolSubject is where a module answers. Derived from the module and the tool, so a caller names
|
|
// what it wants rather than where it lives.
|
|
func ToolSubject(module, tool string) string { return "mesh.mod." + module + ".tool." + tool }
|
|
|
|
func (b OverNATS) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
|
|
if len(args) == 0 {
|
|
args = []byte(`{}`)
|
|
}
|
|
ask, cancel := context.WithTimeout(ctx, timeout)
|
|
defer cancel()
|
|
|
|
// **No reply queue, and no correlation to check.** The caller's inbox is its own — each
|
|
// account is granted one prefix and no other (design 25 §4) — so an answer cannot reach the
|
|
// wrong asker and there is nothing to correlate against. That also settles a cost recorded
|
|
// in build.go: on a shared reply exchange every asker saw every result.
|
|
msg, err := b.Conn.RequestWithContext(ask, ToolSubject(module, tool), args)
|
|
if err != nil {
|
|
if errors.Is(err, nats.ErrNoResponders) {
|
|
// Said at once rather than after the whole wait: nothing is subscribed to that
|
|
// subject, which is a different fact from a tool being slow.
|
|
return nil, fmt.Errorf("nothing serves %s.%s", module, tool)
|
|
}
|
|
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
|
|
}
|
|
return msg.Data, nil
|
|
}
|