Files
mesh-catalog/modules/nats/snapshot/connect.go
T
jochen 419e82cded Back up the bus by the server's own snapshot of each stream, not its live files (hq ADR 0235)
The restic holder copied JetStream's store while the server wrote it; such a
copy may not restore. The nats image now carries mesh-nats-snapshot, run by
the declared dump under the module's own bus account (snapshot API only):
every stream one at a time, flow-controlled, into one tar with a manifest of
counts, sequences and checksums. Restore builds a new store beside the live
one with the bus's own server; a person swaps it in. Proven against
throwaway nats 2.11 servers being written to during the snapshot.
2026-10-06 18:20:51 +02:00

106 lines
3.8 KiB
Go

package main
import (
"crypto/sha256"
"crypto/tls"
"crypto/x509"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"strings"
"time"
"github.com/nats-io/nats.go"
)
// Credential is what the mesh seals to a module so it can reach the bus (mesh-controller
// cmd/mesh-controller issueWith): where, the server's certificate by fingerprint, who, and the
// password. The address is not used here — this runs beside the server, and the server is reached on
// loopback — but the fingerprint is: loopback or not, the certificate is the mesh's own.
type Credential struct {
URL string `json:"url"`
Fingerprint string `json:"fingerprint,omitempty"`
User string `json:"user"`
Password string `json:"password"`
}
func readCredential(path string) (Credential, error) {
var raw []byte
var err error
if path == "-" {
raw, err = io.ReadAll(io.LimitReader(os.Stdin, 1<<20))
} else {
raw, err = os.ReadFile(path)
}
if err != nil {
return Credential{}, fmt.Errorf("cannot read the credential: %w", err)
}
var c Credential
if err := json.Unmarshal(raw, &c); err != nil {
return Credential{}, errors.New("the credential is not the JSON the mesh seals to a module (url, fingerprint, user, password)")
}
if c.User == "" || c.Password == "" {
return Credential{}, errors.New("the credential names no user or password, so the bus would refuse it")
}
return c, nil
}
// connect dials the bus as the credential's user. A `tls://` server is held to the credential's
// fingerprint and nothing else — the mesh's bus presents a certificate of its own, in no trust store,
// so a name check could only fail or be skipped. A plain `nats://` server is accepted only on
// loopback, and is what a restore's own private server is.
func connect(server string, c Credential) (*nats.Conn, error) {
opts := []nats.Option{
nats.UserInfo(c.User, c.Password),
nats.Name("mesh-nats-snapshot"),
// Answers come to the user's own inbox and no other: the bus grants every user its own
// `_INBOX.<user>.>` and nothing wider (novox/hq design 25 §4).
nats.CustomInboxPrefix("_INBOX." + c.User),
nats.Timeout(10 * time.Second),
// One attempt: a snapshot that cannot reach the bus fails the night, loudly, rather than
// waiting for it.
nats.NoReconnect(),
}
switch {
case strings.HasPrefix(server, "tls://"):
if c.Fingerprint == "" {
return nil, errors.New("the credential carries no fingerprint, so the bus's certificate cannot be checked")
}
opts = append(opts, nats.Secure(pinnedTo(c.Fingerprint)))
case strings.HasPrefix(server, "nats://127.0.0.1:"), strings.HasPrefix(server, "nats://localhost:"):
if c.Fingerprint != "" {
// The mesh's bus requires TLS; a plain connection to it would be refused anyway.
opts = append(opts, nats.Secure(pinnedTo(c.Fingerprint)))
}
default:
return nil, fmt.Errorf("%s: the bus is reached over tls:// (or plain nats:// on loopback, for a restore's own server)", server)
}
nc, err := nats.Connect(server, opts...)
if err != nil {
return nil, fmt.Errorf("cannot reach the bus at %s as %s: %w", server, c.User, err)
}
return nc, nil
}
// pinnedTo accepts exactly the certificate with this SHA-256 — the same pin every machine holds
// (mesh-controller internal/broker PinnedToFingerprint).
func pinnedTo(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])
if got := "sha256:" + hex.EncodeToString(sum[:]); got != want {
return fmt.Errorf("the bus presented a certificate this mesh does not know (%s)", got)
}
return nil
},
}
}