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)) } // 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 &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 }, }, 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 { want.DeliverSubject = "_DELIVER." + c.Name } 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 } }