A merge's handler builds for minutes and keeps its own delivery alive; the announcements handed over behind it timed out on the client and came back, and a merge that came back rebuilt what it had just built, five times over (novox/hq issue 175). MaxAckPending 1 on the events consumer: the server holds the rest.
273 lines
11 KiB
Go
273 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,
|
|
}
|
|
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
|
|
}
|
|
}
|