novox/hq 04-ISSUES/146. A push consumer delivers onto an ordinary subject and everything subscribed to it 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. Measured: one enrolment published, one message in the stream, one delivery, no redelivery, and the controller enrolled the machine twice — the second minting a credential that replaced the one the machine had just been handed, which is why it then reconnected for ever as a user whose password the mesh had rotated. Every report and every followed event doubled the same way, silently. The stream goes in the subject because the pair is what identifies a consumer. A subscriber's permission gains the same shape, keeping the bare name so an existing consumer keeps working until the next assertion moves it.
227 lines
8.7 KiB
Go
227 lines
8.7 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
|
|
}
|
|
|
|
// 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 _, err := j.js.ConsumerInfo(c.Stream, c.Name); {
|
|
case err == nil:
|
|
if _, err := j.js.UpdateConsumer(c.Stream, want); err != 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
|
|
}
|
|
}
|