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..>` 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 } }