The pipeline was observable from a merge to an artifact and went dark where it touched a machine: a node's report is control traffic only the control plane reads, so nothing said which version a machine runs, or that it refused to (novox/hq ADR 0134). The control plane now states both under the seat it holds — a role's events belong to the role and keep their address when the holder is replaced — and only when the report is news, because a machine reconciles every minute and a fact per report would be a fact per minute per machine. Whether a report is news is the store's answer: it holds the previous one, so the listener returns it and the server states the fact. That also gives the catch-up replay a subject the controller may publish: it was published as a module's event from a module called "control-plane", which does not exist, so the controller's own account refused it and every catalogue that asked what it missed was answered with nothing.
182 lines
7.8 KiB
Go
182 lines
7.8 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"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"
|
|
)
|
|
|
|
// 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" }
|
|
|
|
// 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
|
|
}
|
|
|
|
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
|
|
}
|