|
|
@@ -1,8 +1,14 @@
|
|
|
|
package broker
|
|
|
|
package broker
|
|
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
import (
|
|
|
|
|
|
|
|
"crypto/sha256"
|
|
|
|
|
|
|
|
"crypto/tls"
|
|
|
|
|
|
|
|
"crypto/x509"
|
|
|
|
|
|
|
|
"encoding/hex"
|
|
|
|
"errors"
|
|
|
|
"errors"
|
|
|
|
"fmt"
|
|
|
|
"fmt"
|
|
|
|
|
|
|
|
"os"
|
|
|
|
|
|
|
|
"strings"
|
|
|
|
"time"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
@@ -24,21 +30,58 @@ type JetStream struct {
|
|
|
|
|
|
|
|
|
|
|
|
// Dial connects and returns the controller's JetStream handle.
|
|
|
|
// Dial connects and returns the controller's JetStream handle.
|
|
|
|
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
|
|
|
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
|
|
|
// A name, because a connection nobody can identify in the server's own monitoring is one
|
|
|
|
|
|
|
|
// nobody can attribute a problem to.
|
|
|
|
|
|
|
|
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
|
|
|
|
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...)
|
|
|
|
conn, err := nats.Connect(url, opts...)
|
|
|
|
if err != nil {
|
|
|
|
if err != nil {
|
|
|
|
return nil, fmt.Errorf("connecting to the bus at %s: %w", url, err)
|
|
|
|
return nil, fmt.Errorf("connecting to the bus at %s: %w", where, err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
js, err := conn.JetStream()
|
|
|
|
js, err := conn.JetStream()
|
|
|
|
if err != nil {
|
|
|
|
if err != nil {
|
|
|
|
conn.Close()
|
|
|
|
conn.Close()
|
|
|
|
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", url, err)
|
|
|
|
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", where, err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return &JetStream{conn: conn, js: js}, nil
|
|
|
|
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,
|
|
|
|
// 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
|
|
|
|
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
|
|
|
// already handles (design 25 §3).
|
|
|
|
// already handles (design 25 §3).
|
|
|
|