novox/hq 04-ISSUES/156. Issue 146 put the stream into a push consumer's delivery subject. The server will not move that subject while a subscriber is bound, 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 every node consumer in a running mesh became one the assertion could not bring to match, and the control plane crash-looped on the assertion it makes before it serves. A fresh mesh showed nothing, because 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. On the live mesh the nodes are granted `_DELIVER.<node>` and not `_DELIVER.<node>.>`, so re-making would have silenced every machine — worse than the collision it fixes, and harder to undo. Kept rather than fatal, which is what 146's change intended and did not do. The bare subject still delivers, and 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. Also: an existing consumer's deliver policy is carried across rather than reasserted, because the server refuses to change it and where a consumer starts is its history. Two tests against a real server: a consumer with a subscriber bound keeps its subject, is reported, and still delivers; one with nothing bound moves, so 146's fix still applies where it matters.
272 lines
11 KiB
Go
272 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,
|
|
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
|
|
}
|
|
}
|