For every module on every machine the controller composes what that instance serves — its machine's address always, the module's plain address in a queue when it is alone or its definition says its instances are interchangeable — the verbs of the seats it holds at the seats' subjects, where its events land, and what it may reach, resolved the same way for the modules it invokes. Published beside the node's declaration on `mesh.assignment.<node>.<module>`, last per subject in a stream that allows direct reads, and the account may read exactly its own. Composed from the same records the bus's accounts are, so what a runtime serves and what its account may are one composition. `instances: interchangeable` is the one fact a definition states for it. The shape issued is the shape the mesh already had, so nothing moves when the membership arrives; the runtime that reads it instead of deriving it is the next piece.
274 lines
11 KiB
Go
274 lines
11 KiB
Go
package broker
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"crypto/tls"
|
|
"crypto/x509"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// The JetStream side of the controller: the one place the mesh's streams and consumers are
|
|
// actually created.
|
|
//
|
|
// Everything that decides *what* they are is pure and lives beside this (streams.go, derived.go).
|
|
// This is only the part that talks to a server, kept small on purpose: a bug in a subject filter
|
|
// should be findable in a unit test, and only a bug in "did the server accept it" should need one
|
|
// running.
|
|
|
|
// A JetStream is a connection to the bus, as the controller uses it.
|
|
type JetStream struct {
|
|
conn *nats.Conn
|
|
js nats.JetStreamContext
|
|
// Note is how this says something it decided not to fail over. Nil is silent, which is only
|
|
// right for a caller that has no way to report; the controller sets it.
|
|
Note func(string, ...any)
|
|
}
|
|
|
|
// note reports without requiring a caller to have set one.
|
|
func (j *JetStream) note(format string, args ...any) {
|
|
if j.Note != nil {
|
|
j.Note(format, args...)
|
|
}
|
|
}
|
|
|
|
// Dial connects and returns the controller's JetStream handle.
|
|
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
|
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
|
|
// **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a
|
|
// public verifier would accept (design 25 §4: a host pins the server's exact certificate and
|
|
// checks nothing else, and so does this). Without this, the first connection failed with
|
|
// "certificate is not valid for any names" against a bus that was answering (2026-09-28).
|
|
if path := strings.TrimSpace(os.Getenv(CertificateVar)); path != "" {
|
|
pinned, err := pinnedTo(path)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
opts = append(opts, nats.Secure(pinned))
|
|
}
|
|
// **Its own inbox, and nothing wider.** Every principal is granted `_INBOX.<its user>.>` and
|
|
// no other inbox; the client's default prefix is random, and the server refused the first
|
|
// subscription to it (2026-09-28). The user is in the URL, so the prefix follows from it.
|
|
if user, _, _ := CredentialIn(url); user != "" {
|
|
opts = append(opts, nats.CustomInboxPrefix("_INBOX."+user))
|
|
}
|
|
// The address in an error is the address alone. The URL carries this controller's password,
|
|
// and an error here is written on the assumption it will be logged.
|
|
where := BareAddress(url)
|
|
conn, err := nats.Connect(url, opts...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("connecting to the bus at %s: %w", where, err)
|
|
}
|
|
js, err := conn.JetStream()
|
|
if err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", where, err)
|
|
}
|
|
return &JetStream{conn: conn, js: js}, nil
|
|
}
|
|
|
|
// pinnedTo is a TLS configuration that accepts exactly the certificate in the file and no other:
|
|
// the leaf's SHA-256, compared on every handshake, with the name and the chain deliberately not
|
|
// consulted — a self-signed certificate with no names is the ordinary case for a mesh's bus.
|
|
func pinnedTo(path string) (*tls.Config, error) {
|
|
want, err := FingerprintOf(path)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return PinnedToFingerprint(want), nil
|
|
}
|
|
|
|
// DialPinned is Dial with the server's certificate pinned by a fingerprint the caller already holds
|
|
// — a module or a build machine that was handed one beside its credential, and has no file.
|
|
func DialPinned(url, fingerprint string, opts ...nats.Option) (*JetStream, error) {
|
|
if strings.TrimSpace(fingerprint) != "" {
|
|
opts = append(opts, nats.Secure(PinnedToFingerprint(fingerprint)))
|
|
}
|
|
return Dial(url, opts...)
|
|
}
|
|
|
|
// PinnedToFingerprint accepts exactly the certificate with this SHA-256 and no other.
|
|
func PinnedToFingerprint(want string) *tls.Config {
|
|
return &tls.Config{
|
|
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter
|
|
MinVersion: tls.VersionTLS12,
|
|
VerifyPeerCertificate: func(rawCerts [][]byte, _ [][]*x509.Certificate) error {
|
|
if len(rawCerts) == 0 {
|
|
return errors.New("the bus presented no certificate")
|
|
}
|
|
sum := sha256.Sum256(rawCerts[0])
|
|
got := "sha256:" + hex.EncodeToString(sum[:])
|
|
if got != want {
|
|
return fmt.Errorf("the bus presented a certificate this mesh does not know (%s…), expected %s…", got[:23], want[:23])
|
|
}
|
|
return nil
|
|
},
|
|
}
|
|
}
|
|
|
|
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
|
|
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
|
// already handles (design 25 §3).
|
|
func (j *JetStream) Conn() *nats.Conn { return j.conn }
|
|
|
|
// Context is the JetStream handle, for subscribing to what the consumers above define.
|
|
func (j *JetStream) Context() nats.JetStreamContext { return j.js }
|
|
|
|
func (j *JetStream) Close() {
|
|
if j.conn != nil {
|
|
j.conn.Close()
|
|
}
|
|
}
|
|
|
|
// EnsureStream creates the stream if it is absent and brings it to match if it is present.
|
|
//
|
|
// **Idempotent, because the controller asserts on every start** rather than creating once at
|
|
// genesis: a stream somebody deleted, or a mesh raised from a restored backup, has to converge
|
|
// rather than run without the guarantee its messages assume.
|
|
//
|
|
// An update, not a delete and recreate. Recreating would discard every message the stream holds
|
|
// and every consumer's position in it — which for CONTROL means the pushes being held through a
|
|
// store restart, exactly the guarantee the stream exists for.
|
|
func (j *JetStream) EnsureStream(s Stream) error {
|
|
want := &nats.StreamConfig{
|
|
Name: s.Name,
|
|
Subjects: s.Subjects,
|
|
Retention: retentionOf(s.Retention),
|
|
MaxAge: time.Duration(s.MaxAge) * time.Second,
|
|
MaxMsgsPerSubject: int64(s.MaxMsgsPerSubject),
|
|
Description: s.Why,
|
|
}
|
|
want.AllowDirect = s.Direct
|
|
if s.Retention == RetentionLastPerSubject {
|
|
// Last-per-subject is a limits stream with one message kept per subject, not a
|
|
// retention policy of its own — the state shape, spelled the way the server spells it.
|
|
want.Retention = nats.LimitsPolicy
|
|
want.MaxMsgsPerSubject = 1
|
|
want.MaxAge = 0
|
|
}
|
|
|
|
switch _, err := j.js.StreamInfo(s.Name); {
|
|
case err == nil:
|
|
if _, err := j.js.UpdateStream(want); err != nil {
|
|
return fmt.Errorf("bringing stream %s to match: %w", s.Name, err)
|
|
}
|
|
return nil
|
|
case errors.Is(err, nats.ErrStreamNotFound):
|
|
if _, err := j.js.AddStream(want); err != nil {
|
|
return fmt.Errorf("creating stream %s: %w", s.Name, err)
|
|
}
|
|
return nil
|
|
default:
|
|
return fmt.Errorf("asking about stream %s: %w", s.Name, err)
|
|
}
|
|
}
|
|
|
|
// EnsureConsumer creates or updates one durable consumer.
|
|
//
|
|
// Explicit acknowledgement throughout: a consumer that acknowledges on delivery cannot redeliver
|
|
// work its holder died in the middle of, which is the whole difference between a queue and a
|
|
// firehose.
|
|
func (j *JetStream) EnsureConsumer(c Consumer) error {
|
|
want := &nats.ConsumerConfig{
|
|
Durable: c.Name,
|
|
AckPolicy: nats.AckExplicitPolicy,
|
|
AckWait: time.Duration(c.AckWaitSeconds) * time.Second,
|
|
MaxDeliver: c.MaxDeliver,
|
|
MaxAckPending: c.MaxAckPending,
|
|
DeliverGroup: c.Queue,
|
|
DeliverSubject: "",
|
|
Description: c.Why,
|
|
}
|
|
switch len(c.Filters) {
|
|
case 0:
|
|
case 1:
|
|
want.FilterSubject = c.Filters[0]
|
|
default:
|
|
want.FilterSubjects = c.Filters
|
|
}
|
|
// A queue group needs a delivery subject: a pull consumer has no group, and declaring one
|
|
// without the other is refused by the server with a message that does not say which half is
|
|
// missing.
|
|
if c.Queue != "" || c.Push {
|
|
// **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146).
|
|
// A push consumer delivers onto an ordinary subject, and everything subscribed to that
|
|
// subject gets a copy. The controller holds a consumer called `controller` on CONTROL and
|
|
// another called `controller` on EVENTS, and both were given `_DELIVER.controller` — so the
|
|
// one process, holding both subscriptions, acted on every message twice. It enrolled a
|
|
// joining machine twice from one request, minting a second credential that replaced the one
|
|
// the machine had just been given; the same doubling applied to every report and every
|
|
// event the controller follows.
|
|
//
|
|
// The stream is in the name because the pair is what identifies a consumer — the server
|
|
// scopes a durable's name to its stream, and this subject is the only place that scoping
|
|
// was dropped. Already within what the controller may subscribe (`_DELIVER.controller.>`),
|
|
// so no permission moves.
|
|
want.DeliverSubject = DeliverSubjectFor(c)
|
|
}
|
|
|
|
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
|
case err == nil:
|
|
// Where an existing consumer starts is its history, not something an assertion may move:
|
|
// the server refuses a changed deliver policy outright. Carried across, so asserting twice
|
|
// is the no-op a restart depends on.
|
|
want.DeliverPolicy = have.Config.DeliverPolicy
|
|
want.OptStartSeq = have.Config.OptStartSeq
|
|
want.OptStartTime = have.Config.OptStartTime
|
|
|
|
if _, err := j.js.UpdateConsumer(c.Stream, want); err != nil {
|
|
// **A consumer that works is not replaced to make its name tidier**
|
|
// (novox/hq 04-ISSUES/156).
|
|
//
|
|
// The server will not move a push consumer's delivery subject while a subscriber is
|
|
// bound to it, and answers `consumer name already in use` — a message about the name,
|
|
// for a conflict about the subject. A node is bound to its declaration consumer the
|
|
// whole time it is up; that IS a node listening. So when 04-ISSUES/146 put the stream
|
|
// into the subject, every node consumer in a running mesh became one this could not
|
|
// bring to match, and the control plane crash-looped on the assertion it makes before
|
|
// it serves. A fresh mesh showed nothing: nothing was bound.
|
|
//
|
|
// Kept rather than deleted and re-made. Re-making moves the subject, and a holder may
|
|
// not be allowed to subscribe to the new one yet — the wider grant travels in the bus's
|
|
// user list, which this same control plane composes and a machine applies minutes
|
|
// later. Re-making here would have silenced every machine in the mesh, which is worse
|
|
// than the collision it was fixing and harder to undo.
|
|
//
|
|
// Kept rather than fatal, which is what 146's change intended and did not do: the bare
|
|
// subject it replaces still delivers, and it collides only where one holder has two
|
|
// consumers of one name. That is the controller's own pair, and the controller is not
|
|
// bound to them while it asserts, so those do move. A node has one consumer and nothing
|
|
// to collide with.
|
|
if have.Config.DeliverSubject != want.DeliverSubject {
|
|
j.note("consumer %s on %s still delivers to %q and not %q: %v. It keeps working; "+
|
|
"the subject moves on an assertion made while nothing is bound to it",
|
|
c.Name, c.Stream, have.Config.DeliverSubject, want.DeliverSubject, err)
|
|
return nil
|
|
}
|
|
return fmt.Errorf("bringing consumer %s on %s to match: %w", c.Name, c.Stream, err)
|
|
}
|
|
return nil
|
|
case errors.Is(err, nats.ErrConsumerNotFound):
|
|
if _, err := j.js.AddConsumer(c.Stream, want); err != nil {
|
|
return fmt.Errorf("creating consumer %s on %s: %w", c.Name, c.Stream, err)
|
|
}
|
|
return nil
|
|
default:
|
|
return fmt.Errorf("asking about consumer %s on %s: %w", c.Name, c.Stream, err)
|
|
}
|
|
}
|
|
|
|
func retentionOf(r Retention) nats.RetentionPolicy {
|
|
switch r {
|
|
case RetentionWorkQueue:
|
|
return nats.WorkQueuePolicy
|
|
default:
|
|
return nats.LimitsPolicy
|
|
}
|
|
}
|