diff --git a/modules/nats/Dockerfile b/modules/nats/Dockerfile index 418fa11..9f57db4 100644 --- a/modules/nats/Dockerfile +++ b/modules/nats/Dockerfile @@ -18,13 +18,30 @@ # # Unlike every other module's Dockerfile, this builds no TypeScript and uses no mesh base image: # the module's code is the server, which upstream already built. There is no BUILD_BASE here on -# purpose — nothing is compiled. The upstream image is declared in the manifest under build.on and -# arrives as NATS_BASE, like every other base the mesh copies into its own store before a build -# (novox/hq ADR 0097); the digest above is the index one for the reason given. +# purpose — the server is not compiled; only the snapshot program below is. The upstream image is +# declared in the manifest under build.on and arrives as NATS_BASE, like every other base the mesh +# copies into its own store before a build (novox/hq ADR 0097); the digest above is the index one +# for the reason given. So does GO_BASE, the toolchain the snapshot program is built with. +# +# **One thing is compiled: the bus's snapshot program** (novox/hq ADR 0235). The machine's backup +# holder runs the module's declared dump — `docker exec` into this container — and the program asks +# the server for each stream through the snapshot API, writing a tar on stdout. It is here, beside +# the server, for two reasons: the dump runs where the server is without mounting anything into it, +# and a restore needs nats-server itself — the very binary this image already carries — to fill a new +# store. Its own Go module (snapshot/go.mod), so this build fetches only public modules. ARG NATS_BASE +ARG GO_BASE +FROM ${GO_BASE} AS snapshot +WORKDIR /src +COPY snapshot/go.mod snapshot/go.sum ./ +RUN go mod download +COPY snapshot/*.go ./ +RUN CGO_ENABLED=0 go build -trimpath -ldflags '-s -w' -o /mesh-nats-snapshot . + FROM ${NATS_BASE} COPY entrypoint.sh /usr/local/bin/mesh-nats-entrypoint RUN chmod 0755 /usr/local/bin/mesh-nats-entrypoint +COPY --from=snapshot /mesh-nats-snapshot /usr/local/bin/mesh-nats-snapshot ENTRYPOINT ["/usr/local/bin/mesh-nats-entrypoint"] diff --git a/modules/nats/README.md b/modules/nats/README.md new file mode 100644 index 0000000..f07245f --- /dev/null +++ b/modules/nats/README.md @@ -0,0 +1,87 @@ +# nats — the mesh's bus + +The bus server (novox/hq design 25), the configuration the module owns (`nats.conf`: ports, TLS, +JetStream's store) and the user list the controller composes (`accounts.conf`), reloaded in place by +the entrypoint. Its tools (`cmd/nats-tools`) are served by the node's runtime and only read. + +## Its data, and how it is backed up + +Everything the bus keeps is in JetStream: every stream, and every key-value bucket (a stream named +`KV_`) — conditions and their history, calls, hand-acts, the controller's lease, assignments, +events, every module's state. + +**The store's files are never what is backed up.** They are written all the time, and a copy taken +while the server writes them may not restore. The `jetstream` data item is protected by a dump +instead (novox/hq ADR 0235, to-be 43): each night the machine's backup holder runs + + docker exec -i mesh-broker-nats mesh-nats-snapshot snapshot < > /bus.tar + +and backs up the `snapshots` directory. `mesh-nats-snapshot` (`snapshot/`, built into this image) +asks the server for each stream through JetStream's snapshot API, one at a time, acknowledging each +chunk so the server's flow control keeps moving; the server goes on taking writes and no stream is +paused or reconfigured. It writes one tar on stdout: + +- `manifest.json` — when it was taken and how long it took, the server's version, and per stream its + subjects, messages, bytes, first and last sequence, consumers, and the size and SHA-256 of its files; + a memory stream is listed as skipped, with why; +- `streams//backup.json` — the stream's configuration and state, as the server gave them; +- `streams//stream.tar.s2` — the server's archive of the stream, consumers included. + +It runs as the module's own bus account, which the controller grants the snapshot API and nothing +else — no stream may be defined, changed, purged or written by it (mesh-controller +`internal/broker`, `BusSnapshotGrants`). A failure fails that night for the module, leaves the +previous `bus.tar` in place, and reaches the operator as `backup-stale` once the last good night is +past its bound. + +**What a snapshot promises.** Every message up to the last sequence its manifest gives for a stream is +in it, exactly. A stream written to while it is read out may hold some messages past that sequence +too — the server reads its blocks a moment after it states the stream — and a restore of it ends at +or after the manifest's sequence, never before. + +## Restoring + +Nothing is restored over the live bus: a stream is only ever restored where it does not exist, and +the controller asserts every stream on start. So the bus is restored into a **new store beside the +live one**, which a person swaps in. + +1. **Take the night back, beside the live data.** `node-backup.restore` for `nats` on the bus's + machine (optionally naming a restore point) restores the snapshots directory as + `.restored-/`, holding that night's `bus.tar`. +2. **Check it.** `docker exec -i mesh-broker-nats mesh-nats-snapshot verify < /bus.tar` — + every file against the manifest's checksum, every archive read to its end. +3. **Build the new store**, with the bus's own image (the server that fills it is the same release + the bus runs): + + image=$(docker inspect -f '{{.Config.Image}}' mesh-broker-nats) + docker run --rm -v :/in:ro -v :/out \ + --entrypoint mesh-nats-snapshot "$image" restore --into /out --from /in/bus.tar + + It starts a server of its own on the container's loopback with the bus's one account (`MESH`), + restores every stream, holds each to the manifest, stops it, and leaves + `/jetstream/MESH/streams/…` — the bus's layout. +4. **Swap it in — a person's act, a planned bus step** (novox/hq to-be 45): say it as a hand act, + then in one command line, because the node-engine starts a stopped container again at its next + reconcile (within five minutes): + + docker stop mesh-broker-nats && \ + mv /jetstream /jetstream.before-restore- && \ + mv /jetstream /jetstream && \ + docker start mesh-broker-nats + + The store kept aside is removed by a person once the bus is seen whole (`nats_streams`, the + controller's `doctor`). + +A single stream can be restored into any server that does not hold it with +`mesh-nats-snapshot restore --server --credential `; the archive is also what the nats CLI +reads (`nats stream restore ` on an unpacked `streams//`). + +## Tests + +- `snapshot/`: `MESH_TEST_DOCKER=1 go test -race ./...` — a throwaway `nats:2.11` filled like the + mesh's bus (events with a gap deleted, a work queue part worked, buckets with history, deletes and a + purge, an object, durable consumers, a memory stream) and written to throughout the snapshot; the + snapshot restored into a fresh server and into a new store served by a third, each compared message + by message; a damaged archive refused; and the snapshot user refused every write. + `MESH_TEST_NATS_IMAGE=` also runs the manifest's declared dump, exactly, against the + module's image and configuration. +- `cmd/nats-tools`: `go test ./...`. diff --git a/modules/nats/module.json b/modules/nats/module.json index d59f407..7ecb8b3 100644 --- a/modules/nats/module.json +++ b/modules/nats/module.json @@ -18,6 +18,9 @@ } ], "bus-users": "/var/lib/nats-module/conf/accounts.conf", + "own-secrets": { + "broker": "${dir:mesh-state}/broker" + }, "capabilities": [ "container-runtime" ], @@ -52,7 +55,17 @@ "path": "${dir:jetstream-data}", "class": "valuable", "active": "1d", - "why": "the bus's streams and key-value buckets: the hand-act log, conditions, every module's state; written all the time" + "backup": { + "dump": "docker exec -i mesh-broker-nats mesh-nats-snapshot snapshot < ${dir:mesh-state}/broker > ${dir:snapshots}/bus.tar.partial && mv ${dir:snapshots}/bus.tar.partial ${dir:snapshots}/bus.tar || { rm -f ${dir:snapshots}/bus.tar.partial; exit 1; }", + "into": "snapshots" + }, + "why": "the bus's streams and key-value buckets: the hand-act log, conditions, every module's state; written all the time, so copied by the server's own snapshot of each stream (novox/hq ADR 0235), never as live files" + }, + { + "id": "snapshots", + "path": "${dir:snapshots}", + "class": "rebuildable", + "why": "last night's snapshot of every stream with its manifest, made again every night" } ] }, @@ -69,6 +82,17 @@ "path": "/var/lib/nats-module/conf", "mode": "0700" }, + { + "id": "snapshots", + "type": "directory", + "mode": "0700" + }, + { + "id": "mesh-state", + "type": "directory", + "mode": "0700", + "place": "mesh" + }, { "id": "server-conf", "type": "file", @@ -104,6 +128,10 @@ { "arg": "NATS_BASE", "image": "nats@sha256:e4bf19f15fd3218814a4e3c9e0064e1334bd8aa20d5984b9f1a0afd084f8cc00" + }, + { + "arg": "GO_BASE", + "image": "golang@sha256:8ac98ca534ac3f51e1f420a1dd2c15e74c75cfa0f23f3ad27eb5d7236c349a0c" } ], "artifacts": [ diff --git a/modules/nats/snapshot/archive.go b/modules/nats/snapshot/archive.go new file mode 100644 index 0000000..c894aa4 --- /dev/null +++ b/modules/nats/snapshot/archive.go @@ -0,0 +1,162 @@ +package main + +import ( + "archive/tar" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "path/filepath" +) + +// Archive is a snapshot read back: its manifest, and each file beside it unpacked into a scratch +// directory of its own, removed by Close. +type Archive struct { + Manifest Manifest + dir string + files map[string]string +} + +// ReadArchive unpacks a snapshot. Only the files the manifest names are kept, and only under names +// that are one stream's directory: an archive is data, and a path in it is never trusted to say where +// on the machine a file goes. +func ReadArchive(r io.Reader, scratch string) (*Archive, error) { + dir, err := os.MkdirTemp(scratch, "mesh-nats-restore-*") + if err != nil { + return nil, err + } + a := &Archive{dir: dir, files: map[string]string{}} + tr := tar.NewReader(r) + first := true + for { + hdr, err := tr.Next() + if err == io.EOF { + break + } + if err != nil { + a.Close() + return nil, fmt.Errorf("the archive does not read: %w", err) + } + if first { + if hdr.Name != ManifestName { + a.Close() + return nil, fmt.Errorf("the archive starts with %q, not its manifest: it is not a snapshot of the bus", hdr.Name) + } + raw, err := io.ReadAll(io.LimitReader(tr, 64<<20)) + if err != nil { + a.Close() + return nil, err + } + if err := json.Unmarshal(raw, &a.Manifest); err != nil { + a.Close() + return nil, fmt.Errorf("its manifest: %w", err) + } + if a.Manifest.Format != 1 { + a.Close() + return nil, fmt.Errorf("its manifest is format %d; this program reads format 1", a.Manifest.Format) + } + first = false + continue + } + if !a.named(hdr.Name) { + a.Close() + return nil, fmt.Errorf("the archive holds %q, which its manifest does not name", hdr.Name) + } + local := filepath.Join(dir, fmt.Sprintf("%04d", len(a.files))) + f, err := os.OpenFile(local, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0o600) + if err != nil { + a.Close() + return nil, err + } + _, err = io.Copy(f, tr) + f.Close() + if err != nil { + a.Close() + return nil, err + } + a.files[hdr.Name] = local + } + if first { + a.Close() + return nil, errors.New("the archive is empty") + } + for _, s := range a.Manifest.Streams { + if !validName(s.Name) || s.Meta != "streams/"+s.Name+"/backup.json" || s.Archive != "streams/"+s.Name+"/stream.tar.s2" { + a.Close() + return nil, fmt.Errorf("the manifest names a stream %q at paths that are not its own", s.Name) + } + } + return a, nil +} + +func (a *Archive) named(name string) bool { + for _, s := range a.Manifest.Streams { + if name == s.Meta || name == s.Archive { + return true + } + } + return false +} + +// Close removes the unpacked files. +func (a *Archive) Close() { os.RemoveAll(a.dir) } + +// Verify holds every file to the manifest's size and checksum, and reads every stream's archive to +// its end. +func (a *Archive) Verify() error { + for _, s := range a.Manifest.Streams { + for _, f := range []struct{ name, sum string }{{s.Meta, s.MetaSHA256}, {s.Archive, s.ArchiveSHA256}} { + local, ok := a.files[f.name] + if !ok { + return fmt.Errorf("%s is named in the manifest and missing from the archive", f.name) + } + got, size, err := sha256Of(local) + if err != nil { + return err + } + if got != f.sum { + return fmt.Errorf("%s does not match the checksum its manifest gives", f.name) + } + if f.name == s.Archive && size != s.ArchiveBytes { + return fmt.Errorf("%s is %d bytes; its manifest says %d", f.name, size, s.ArchiveBytes) + } + } + af, err := os.Open(a.files[s.Archive]) + if err != nil { + return err + } + err = readsThrough(af) + af.Close() + if err != nil { + return fmt.Errorf("%s does not read: %w", s.Archive, err) + } + } + return nil +} + +func sha256Of(path string) (string, int64, error) { + f, err := os.Open(path) + if err != nil { + return "", 0, err + } + defer f.Close() + h := sha256.New() + n, err := io.Copy(h, f) + if err != nil { + return "", 0, err + } + return hex.EncodeToString(h.Sum(nil)), n, nil +} + +// open is one stream's meta and archive. +func (a *Archive) open(s StreamEntry) (meta []byte, archive *os.File, err error) { + meta, err = os.ReadFile(a.files[s.Meta]) + if err != nil { + return nil, nil, err + } + archive, err = os.Open(a.files[s.Archive]) + return meta, archive, err +} diff --git a/modules/nats/snapshot/connect.go b/modules/nats/snapshot/connect.go new file mode 100644 index 0000000..97aad5b --- /dev/null +++ b/modules/nats/snapshot/connect.go @@ -0,0 +1,105 @@ +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..>` 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 + }, + } +} diff --git a/modules/nats/snapshot/go.mod b/modules/nats/snapshot/go.mod new file mode 100644 index 0000000..b089934 --- /dev/null +++ b/modules/nats/snapshot/go.mod @@ -0,0 +1,15 @@ +module nats-snapshot + +go 1.25.0 + +require ( + github.com/klauspost/compress v1.18.5 + github.com/nats-io/nats.go v1.51.0 +) + +require ( + github.com/nats-io/nkeys v0.4.15 // indirect + github.com/nats-io/nuid v1.0.1 // indirect + golang.org/x/crypto v0.49.0 // indirect + golang.org/x/sys v0.42.0 // indirect +) diff --git a/modules/nats/snapshot/go.sum b/modules/nats/snapshot/go.sum new file mode 100644 index 0000000..9050eeb --- /dev/null +++ b/modules/nats/snapshot/go.sum @@ -0,0 +1,12 @@ +github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE= +github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/nats-io/nats.go v1.51.0 h1:ByW84XTz6W03GSSsygsZcA+xgKK8vPGaa/FCAAEHnAI= +github.com/nats-io/nats.go v1.51.0/go.mod h1:26HypzazeOkyO3/mqd1zZd53STJN0EjCYF9Uy2ZOBno= +github.com/nats-io/nkeys v0.4.15 h1:JACV5jRVO9V856KOapQ7x+EY8Jo3qw1vJt/9Jpwzkk4= +github.com/nats-io/nkeys v0.4.15/go.mod h1:CpMchTXC9fxA5zrMo4KpySxNjiDVvr8ANOSZdiNfUrs= +github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= +github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= +golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= +golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo= +golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= diff --git a/modules/nats/snapshot/image_test.go b/modules/nats/snapshot/image_test.go new file mode 100644 index 0000000..8e9600a --- /dev/null +++ b/modules/nats/snapshot/image_test.go @@ -0,0 +1,152 @@ +package main + +// The declared dump, exactly as the backup holder runs it, against the bus's own image. +// +// docker build --build-arg NATS_BASE= --build-arg GO_BASE= -t mesh-nats:check .. +// MESH_TEST_NATS_IMAGE=mesh-nats:check go test -run TestTheDeclaredDump ./... + +import ( + "encoding/json" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// manifestOf is the nats module's own definition, beside this program. +func manifestOf(t *testing.T) map[string]any { + t.Helper() + raw, err := os.ReadFile("../module.json") + if err != nil { + t.Fatal(err) + } + var m map[string]any + if err := json.Unmarshal(raw, &m); err != nil { + t.Fatal(err) + } + return m +} + +func declaredDump(t *testing.T, m map[string]any) string { + t.Helper() + for _, item := range m["data"].(map[string]any)["own"].([]any) { + it := item.(map[string]any) + if it["id"] != "jetstream" { + continue + } + backup, _ := it["backup"].(map[string]any) + if backup == nil || backup["into"] != "snapshots" { + t.Fatalf("the bus's streams are not protected by a dump into its snapshots: %v", it["backup"]) + } + return backup["dump"].(string) + } + t.Fatal("the nats module declares no jetstream data item") + return "" +} + +func serverConf(t *testing.T, m map[string]any) string { + t.Helper() + for _, r := range m["resources"].([]any) { + res := r.(map[string]any) + if res["id"] == "server-conf" { + return res["content"].(string) + } + } + t.Fatal("the nats module declares no server configuration") + return "" +} + +// The dump is the one the manifest declares, run by `sh -c` as the holder runs it, with the module's +// directories filled in; the server is the module's image with the module's own configuration and a +// user list holding the bus's own module with the controller's grants for it. +func TestTheDeclaredDumpRunsInTheBussImage(t *testing.T) { + image := os.Getenv("MESH_TEST_NATS_IMAGE") + if image == "" { + t.Skip("MESH_TEST_NATS_IMAGE unset: the bus's image, built from ../Dockerfile") + } + const name = "mesh-broker-nats" // the container the dump names + if out, _ := exec.Command("docker", "ps", "-a", "-q", "--filter", "name=^"+name+"$").Output(); len(strings.TrimSpace(string(out))) > 0 { + t.Skipf("a container called %s already runs here; this test will not touch it", name) + } + m := manifestOf(t) + dump := declaredDump(t, m) + + conf, tlsDir, data, state, snapshots := t.TempDir(), t.TempDir(), t.TempDir(), t.TempDir(), t.TempDir() + pin := selfSigned(t, tlsDir) + _ = os.WriteFile(filepath.Join(conf, "nats.conf"), []byte(serverConf(t, m)), 0o644) + accounts := liveBusConf() + accounts = accounts[strings.Index(accounts, "accounts {"):] + _ = os.WriteFile(filepath.Join(conf, "accounts.conf"), []byte(accounts), 0o644) + + docker(t, "run", "-d", "--rm", "--name", name, "--user", owner(), "-p", "127.0.0.1::4222", + "-v", data+":/data", "-v", conf+":/etc/nats:ro", "-v", tlsDir+":/tls:ro", image) + t.Cleanup(func() { _ = exec.Command("docker", "rm", "-f", name).Run() }) + _, port, _ := strings.Cut(docker(t, "port", name, "4222/tcp"), ":") + + admin := dial(t, "tls://127.0.0.1:"+port, Credential{Fingerprint: pin, User: "admin", Password: "admin"}, nil) + js, _ := admin.JetStream() + kv, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "mesh-controller_hand-acts"}) + if err != nil { + t.Fatal(err) + } + for _, k := range []string{"a", "b", "c"} { + if _, err := kv.Put(k, []byte("act "+k)); err != nil { + t.Fatal(err) + } + } + + run := func(cred Credential) (string, error) { + raw, _ := json.Marshal(cred) + if err := os.WriteFile(filepath.Join(state, "broker"), raw, 0o600); err != nil { + t.Fatal(err) + } + command := strings.NewReplacer("${dir:mesh-state}", state, "${dir:snapshots}", snapshots).Replace(dump) + out, err := exec.Command("sh", "-c", command).CombinedOutput() + return string(out), err + } + + began := time.Now() + out, err := run(Credential{URL: "nats://bus.example:4222", Fingerprint: pin, User: snapshotUser, Password: "snap"}) + if err != nil { + t.Fatalf("the declared dump failed: %v\n%s", err, out) + } + t.Logf("the declared dump took %s:\n%s", time.Since(began).Round(time.Millisecond), out) + tarPath := filepath.Join(snapshots, "bus.tar") + f, err := os.Open(tarPath) + if err != nil { + t.Fatalf("the dump left no bus.tar: %v", err) + } + a, err := ReadArchive(f, "") + f.Close() + if err != nil { + t.Fatal(err) + } + defer a.Close() + if err := a.Verify(); err != nil { + t.Fatal(err) + } + if len(a.Manifest.Streams) != 1 || a.Manifest.Streams[0].Name != "KV_mesh-controller_hand-acts" || a.Manifest.Streams[0].Messages != 3 { + t.Fatalf("the snapshot holds %+v", a.Manifest.Streams) + } + before, _ := os.ReadFile(tarPath) + + // A night that cannot snapshot fails, says why, and leaves last night's snapshot as it was. + out, err = run(Credential{Fingerprint: pin, User: snapshotUser, Password: "wrong"}) + if err == nil { + t.Fatalf("a dump the bus refused succeeded:\n%s", out) + } + if !strings.Contains(out, "Authorization Violation") && !strings.Contains(strings.ToLower(out), "authorization") { + t.Errorf("the failure does not say the bus refused the credential: %s", out) + } + after, _ := os.ReadFile(tarPath) + if string(after) != string(before) { + t.Error("a failed night replaced the last good snapshot") + } + if _, err := os.Stat(tarPath + ".partial"); !os.IsNotExist(err) { + t.Error("a failed night left its partial file behind") + } +} diff --git a/modules/nats/snapshot/live_test.go b/modules/nats/snapshot/live_test.go new file mode 100644 index 0000000..520cd82 --- /dev/null +++ b/modules/nats/snapshot/live_test.go @@ -0,0 +1,587 @@ +package main + +// The snapshot and the way back, against real servers: the bus's own release in throwaway containers. +// +// MESH_TEST_DOCKER=1 go test -race ./... +// +// Every container and file it makes is removed when it ends. MESH_TEST_SNAPSHOT_MESSAGES sets how +// many events the source holds (default 6000), to time a snapshot at the live bus's size. + +import ( + "bytes" + "context" + "crypto/ecdsa" + "crypto/elliptic" + "crypto/rand" + "crypto/sha256" + "crypto/x509" + "crypto/x509/pkix" + "encoding/hex" + "encoding/pem" + "errors" + "fmt" + "math/big" + mrand "math/rand" + "os" + "os/exec" + "path/filepath" + "strconv" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// image is the bus's release (the Dockerfile's `# upstream:` line names the same one). +const image = "nats:2.11-alpine" + +// snapshotUser is the bus's own module on a machine, as the controller names it. +const snapshotUser = "anchor.nats" + +// snapshotGrants are what mesh-controller composes for the module holding mesh-broker +// (internal/broker BusSnapshotGrants), written out here because this repository does not import the +// controller. Its own test proves the same list against a real server; this one proves the program +// needs no more than it. +var snapshotGrants = struct{ publish, subscribe []string }{ + publish: []string{ + "$JS.API.STREAM.NAMES", + "$JS.API.STREAM.INFO.*", + "$JS.API.STREAM.SNAPSHOT.*", + "$JS.SNAPSHOT.ACK.>", + }, + subscribe: []string{"_INBOX." + snapshotUser + ".>"}, +} + +func needDocker(t *testing.T) { + t.Helper() + if os.Getenv("MESH_TEST_DOCKER") != "1" { + t.Skip("MESH_TEST_DOCKER is not 1: these tests start throwaway nats containers") + } +} + +func docker(t *testing.T, args ...string) string { + t.Helper() + out, err := exec.Command("docker", args...).CombinedOutput() + if err != nil { + t.Fatalf("docker %s: %v\n%s", strings.Join(args, " "), err, out) + } + return strings.TrimSpace(string(out)) +} + +func owner() string { return fmt.Sprintf("%d:%d", os.Getuid(), os.Getgid()) } + +// server is one throwaway container. +type server struct { + name string + port string +} + +// startServer runs the bus's release with a configuration, as this user so what it writes into a +// mounted directory is removable by the test. +func startServer(t *testing.T, conf string, mounts ...string) server { + t.Helper() + dir := t.TempDir() + if err := os.WriteFile(filepath.Join(dir, "nats.conf"), []byte(conf), 0o644); err != nil { + t.Fatal(err) + } + name := fmt.Sprintf("mesh-snapshot-test-%d-%d", os.Getpid(), time.Now().UnixNano()) + args := []string{"run", "-d", "--rm", "--name", name, "--user", owner(), "-p", "127.0.0.1::4222", + "-v", dir + ":/etc/nats:ro"} + data := false + for _, m := range mounts { + args = append(args, "-v", m) + data = data || strings.HasSuffix(m, ":/data") + } + if !data { + // The store on a directory of the test's, so the server writes it as this user and the + // test removes it. + args = append(args, "-v", t.TempDir()+":/data") + } + args = append(args, image, "-c", "/etc/nats/nats.conf") + docker(t, args...) + t.Cleanup(func() { _ = exec.Command("docker", "rm", "-f", name).Run() }) + mapped := docker(t, "port", name, "4222/tcp") + _, port, _ := strings.Cut(strings.Split(mapped, "\n")[0], ":") + return server{name: name, port: port} +} + +// selfSigned is a certificate with no names, the ordinary case for a mesh's bus, and its pin. +func selfSigned(t *testing.T, dir string) string { + t.Helper() + key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader) + if err != nil { + t.Fatal(err) + } + tmpl := &x509.Certificate{SerialNumber: big.NewInt(1), Subject: pkix.Name{CommonName: "bus"}, + NotBefore: time.Now().Add(-time.Hour), NotAfter: time.Now().Add(time.Hour)} + der, err := x509.CreateCertificate(rand.Reader, tmpl, tmpl, &key.PublicKey, key) + if err != nil { + t.Fatal(err) + } + keyDER, _ := x509.MarshalECPrivateKey(key) + _ = os.WriteFile(filepath.Join(dir, "tls.crt"), pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}), 0o644) + _ = os.WriteFile(filepath.Join(dir, "tls.key"), pem.EncodeToMemory(&pem.Block{Type: "EC PRIVATE KEY", Bytes: keyDER}), 0o644) + sum := sha256.Sum256(der) + return "sha256:" + hex.EncodeToString(sum[:]) +} + +func quotedList(subjects []string) string { + q := make([]string, len(subjects)) + for i, s := range subjects { + q[i] = strconv.Quote(s) + } + return strings.Join(q, ", ") +} + +// liveBusConf is a server shaped like the mesh's: TLS, the one account, an administrator to fill it, +// and the bus's own module with exactly the grants the controller composes for it. +func liveBusConf() string { + return fmt.Sprintf(`listen: "0.0.0.0:4222" +tls { cert_file: "/tls/tls.crt", key_file: "/tls/tls.key" } +jetstream { store_dir: "/data" } +accounts { MESH { jetstream: enabled, users: [ + { user: "admin", password: "admin" }, + { user: %q, password: "snap", permissions: { + publish: { allow: [%s] }, + subscribe: { allow: [%s] } } } +] } } +`, snapshotUser, quotedList(snapshotGrants.publish), quotedList(snapshotGrants.subscribe)) +} + +// plainConf is a fresh server with the one account and an administrator, nothing in it. +func plainConf() string { + return `listen: "0.0.0.0:4222" +jetstream { store_dir: "/data" } +accounts { MESH { jetstream: enabled, users: [ { user: "admin", password: "admin" } ] } } +` +} + +func dial(t *testing.T, url string, c Credential, errs chan<- error) *nats.Conn { + t.Helper() + var nc *nats.Conn + var err error + for deadline := time.Now().Add(20 * time.Second); ; { + nc, err = connect(url, c) + if err == nil { + break + } + if time.Now().After(deadline) { + t.Fatal(err) + } + time.Sleep(200 * time.Millisecond) + } + if errs != nil { + nc.SetErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) { + select { + case errs <- err: + default: + } + }) + } + t.Cleanup(nc.Close) + return nc +} + +func messageCount() int { + if n, err := strconv.Atoi(os.Getenv("MESH_TEST_SNAPSHOT_MESSAGES")); err == nil && n > 0 { + return n + } + return 6000 +} + +// fill gives the source what the mesh's bus holds: events with headers and a gap deleted from the +// middle, a work queue partly worked, key-value buckets with history, deletes and a purge, an +// object in chunks, durable consumers part way through, and a memory stream. +func fill(t *testing.T, js nats.JetStreamContext) { + t.Helper() + must := func(err error) { + t.Helper() + if err != nil { + t.Fatal(err) + } + } + _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}, + MaxMsgsPerSubject: 10000, MaxAge: 7 * 24 * time.Hour}) + must(err) + r := mrand.New(mrand.NewSource(1)) + payload := make([]byte, 2048) + total := messageCount() + var acks []nats.PubAckFuture + for i := 0; i < total; i++ { + r.Read(payload) + m := nats.NewMsg(fmt.Sprintf("mesh.mod.m%d.event.happened", i%30)) + m.Header.Set("Mesh-Seq", strconv.Itoa(i)) + m.Data = append([]byte(nil), payload[:200+r.Intn(1800)]...) + f, err := js.PublishMsgAsync(m) + must(err) + acks = append(acks, f) + if len(acks) == 2000 { + <-js.PublishAsyncComplete() + acks = acks[:0] + } + } + <-js.PublishAsyncComplete() + for seq := uint64(100); seq < 120; seq++ { + must(js.DeleteMsg("EVENTS", seq)) + } + _, err = js.AddConsumer("EVENTS", &nats.ConsumerConfig{Durable: "audit", AckPolicy: nats.AckExplicitPolicy}) + must(err) + sub, err := js.PullSubscribe("mesh.mod.*.event.>", "audit", nats.Bind("EVENTS", "audit")) + must(err) + msgs, err := sub.Fetch(500, nats.MaxWait(5*time.Second)) + must(err) + for _, m := range msgs { + must(m.Ack()) + } + + _, err = js.AddStream(&nats.StreamConfig{Name: "CONTROL", Subjects: []string{"mesh.control.*.report"}, + Retention: nats.WorkQueuePolicy}) + must(err) + for i := 0; i < 50; i++ { + _, err := js.Publish(fmt.Sprintf("mesh.control.n%d.report", i%4), []byte(fmt.Sprintf("report %d", i))) + must(err) + } + _, err = js.AddConsumer("CONTROL", &nats.ConsumerConfig{Durable: "controller", AckPolicy: nats.AckExplicitPolicy}) + must(err) + wsub, err := js.PullSubscribe("mesh.control.*.report", "controller", nats.Bind("CONTROL", "controller")) + must(err) + worked, err := wsub.Fetch(20, nats.MaxWait(5*time.Second)) + must(err) + for _, m := range worked { + must(m.AckSync()) + } + + conditions, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "mesh-controller_conditions"}) + must(err) + history, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "mesh-controller_condition-history", History: 5}) + must(err) + for i := 0; i < 200; i++ { + key := fmt.Sprintf("machine.n%d.silent", i%40) + _, err := conditions.Put(key, []byte(fmt.Sprintf(`{"raised":%d}`, i))) + must(err) + _, err = history.Put(key, []byte(fmt.Sprintf(`{"at":%d}`, i))) + must(err) + } + for i := 0; i < 5; i++ { + must(conditions.Delete(fmt.Sprintf("machine.n%d.silent", i))) + } + must(history.Purge("machine.n7.silent")) + _, err = js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "mesh-controller_lease", TTL: time.Hour}) + must(err) + + objects, err := js.CreateObjectStore(&nats.ObjectStoreConfig{Bucket: "artifacts"}) + must(err) + big := make([]byte, 3<<20) + r.Read(big) + _, err = objects.PutBytes("facts/latest", big) + must(err) + + _, err = js.AddStream(&nats.StreamConfig{Name: "SCRATCH", Subjects: []string{"scratch.>"}, Storage: nats.MemoryStorage}) + must(err) + _, err = js.Publish("scratch.x", []byte("gone at a restart")) + must(err) +} + +// sameContent holds every stream the manifest names, in the restored server, to the source: the same +// state at the snapshot, every message by sequence (subject, headers, body, and a deleted one still +// deleted), and the durable consumers where they were. +func sameContent(t *testing.T, m *Manifest, src, dst nats.JetStreamContext) { + t.Helper() + for _, s := range m.Streams { + info, err := dst.StreamInfo(s.Name) + if err != nil { + t.Errorf("%s was not restored: %v", s.Name, err) + continue + } + if err := heldTo(streamState{Messages: info.State.Msgs, LastSeq: info.State.LastSeq}, s); err != nil { + t.Errorf("%s: %v", s.Name, err) + } + // **Up to the snapshot's last sequence, exactly the source**: every message by sequence — + // subject, headers and body — and a message deleted before the snapshot still deleted. Past it, + // what was written while the blocks were read out: a message there may or may not have made it + // (the server's state says the stream reached it before the block holding it was read), and + // one that did is the source's. The source only grows after the snapshot. + for seq := info.State.FirstSeq; seq <= info.State.LastSeq && seq > 0; seq++ { + want, werr := src.GetMsg(s.Name, seq) + got, gerr := dst.GetMsg(s.Name, seq) + switch { + case errors.Is(werr, nats.ErrMsgNotFound): + if gerr == nil { + t.Errorf("%s %d was deleted and came back", s.Name, seq) + } + continue + case errors.Is(gerr, nats.ErrMsgNotFound): + if seq <= s.LastSeq { + t.Errorf("%s %d was in the stream at its snapshot and is missing from the restore", s.Name, seq) + } + continue + case werr != nil || gerr != nil: + t.Fatalf("%s %d: source %v, restored %v", s.Name, seq, werr, gerr) + } + if got.Subject != want.Subject || !bytes.Equal(got.Data, want.Data) || fmt.Sprint(got.Header) != fmt.Sprint(want.Header) { + t.Fatalf("%s %d differs: %s %q vs %s %q", s.Name, seq, got.Subject, truncate(got.Data), want.Subject, truncate(want.Data)) + } + } + if info.State.LastSeq > s.LastSeq { + t.Logf("%s: restored to sequence %d, written while the snapshot read it; the manifest said %d", s.Name, info.State.LastSeq, s.LastSeq) + } + for _, consumer := range []string{"audit", "controller"} { + wantC, err := src.ConsumerInfo(s.Name, consumer) + if err != nil { + continue + } + gotC, err := dst.ConsumerInfo(s.Name, consumer) + if err != nil { + t.Errorf("%s's consumer %s was not restored: %v", s.Name, consumer, err) + continue + } + if gotC.AckFloor.Stream != wantC.AckFloor.Stream { + t.Errorf("%s's consumer %s acknowledged up to %d, the source up to %d", s.Name, consumer, + gotC.AckFloor.Stream, wantC.AckFloor.Stream) + } + } + } + for _, s := range m.Skipped { + if _, err := dst.StreamInfo(s.Name); err == nil { + t.Errorf("%s was skipped and yet exists", s.Name) + } + } + // And the buckets work as buckets, not only as streams. + kv, err := dst.KeyValue("mesh-controller_condition-history") + if err != nil { + t.Fatalf("the condition history is not a bucket after the restore: %v", err) + } + h, err := kv.History("machine.n3.silent") + if err != nil || len(h) != 5 { + t.Errorf("the history of a key came back as %d entries (%v), not 5", len(h), err) + } + if _, err := kv.Get("machine.n7.silent"); !errors.Is(err, nats.ErrKeyNotFound) { + t.Errorf("a purged key came back: %v", err) + } + cond, _ := dst.KeyValue("mesh-controller_conditions") + if _, err := cond.Get("machine.n1.silent"); !errors.Is(err, nats.ErrKeyNotFound) { + t.Errorf("a deleted key came back: %v", err) + } + if e, err := cond.Get("machine.n10.silent"); err != nil || string(e.Value()) != `{"raised":170}` { + t.Errorf("a condition came back wrong: %v", err) + } + srcObj, _ := src.ObjectStore("artifacts") + dstObj, err := dst.ObjectStore("artifacts") + if err != nil { + t.Fatalf("the object store was not restored: %v", err) + } + want, _ := srcObj.GetBytes("facts/latest") + got, err := dstObj.GetBytes("facts/latest") + if err != nil || !bytes.Equal(got, want) { + t.Errorf("the object came back different (%d bytes, %v)", len(got), err) + } +} + +func TestASnapshotOfALiveBusRestoresToTheSameContent(t *testing.T) { + needDocker(t) + tlsDir := t.TempDir() + pin := selfSigned(t, tlsDir) + source := startServer(t, liveBusConf(), tlsDir+":/tls:ro") + url := "tls://127.0.0.1:" + source.port + + admin := dial(t, url, Credential{Fingerprint: pin, User: "admin", Password: "admin"}, nil) + srcJS, _ := admin.JetStream() + fill(t, srcJS) + + // The bus goes on taking writes while it is snapshotted, and every one is accepted. + writer := dial(t, url, Credential{Fingerprint: pin, User: "admin", Password: "admin"}, nil) + writerJS, _ := writer.JetStream() + stop := make(chan struct{}) + var wrote, failed atomic.Int64 + var slowest atomic.Int64 + var wg sync.WaitGroup + wg.Add(1) + go func() { + defer wg.Done() + for { + select { + case <-stop: + return + default: + } + began := time.Now() + if _, err := writerJS.Publish("mesh.mod.writer.event.tick", []byte("during")); err != nil { + failed.Add(1) + } else { + wrote.Add(1) + } + if d := time.Since(began).Milliseconds(); d > slowest.Load() { + slowest.Store(d) + } + time.Sleep(2 * time.Millisecond) + } + }() + + snap := dial(t, url, Credential{Fingerprint: pin, User: snapshotUser, Password: "snap"}, nil) + var archive bytes.Buffer + var said []string + var sayMu sync.Mutex + m, err := Snapshot(context.Background(), snap, &archive, SnapshotOptions{StreamTimeout: time.Minute, Pause: 50 * time.Millisecond, + Log: func(f string, a ...any) { sayMu.Lock(); said = append(said, fmt.Sprintf(f, a...)); sayMu.Unlock() }}) + time.Sleep(100 * time.Millisecond) + close(stop) + wg.Wait() + if err != nil { + t.Fatalf("the snapshot failed: %v\n%s", err, strings.Join(said, "\n")) + } + t.Logf("snapshot: %d streams, %d messages, %s of streams → %s archived in %.2fs; %d writes during it, %d refused, slowest %dms", + len(m.Streams), m.Messages, human(m.StreamBytes), human(m.ArchiveBytes), m.Seconds, wrote.Load(), failed.Load(), slowest.Load()) + for _, line := range said { + t.Log(line) + } + if failed.Load() > 0 { + t.Errorf("%d writes to the bus failed while it was snapshotted", failed.Load()) + } + if len(m.Skipped) != 1 || m.Skipped[0].Name != "SCRATCH" { + t.Errorf("the memory stream was not said as skipped: %+v", m.Skipped) + } + var names []string + for _, s := range m.Streams { + names = append(names, s.Name) + } + wantNames := "CONTROL EVENTS KV_mesh-controller_condition-history KV_mesh-controller_conditions KV_mesh-controller_lease OBJ_artifacts" + if strings.Join(names, " ") != wantNames { + t.Errorf("snapshotted %v, want %s", names, wantNames) + } + + tarPath := filepath.Join(t.TempDir(), "bus.tar") + if err := os.WriteFile(tarPath, archive.Bytes(), 0o644); err != nil { + t.Fatal(err) + } + a, err := ReadArchive(bytes.NewReader(archive.Bytes()), "") + if err != nil { + t.Fatal(err) + } + defer a.Close() + if err := a.Verify(); err != nil { + t.Fatalf("a fresh snapshot does not verify: %v", err) + } + + t.Run("nothing is restored over a live stream", func(t *testing.T) { + err := RestoreAll(context.Background(), admin, a, t.Logf) + if err == nil || !strings.Contains(err.Error(), "already holds") { + t.Fatalf("restoring over the live streams was not refused: %v", err) + } + }) + + t.Run("into a fresh server", func(t *testing.T) { + fresh := startServer(t, plainConf()) + nc := dial(t, "nats://127.0.0.1:"+fresh.port, Credential{User: "admin", Password: "admin"}, nil) + if err := RestoreAll(context.Background(), nc, a, t.Logf); err != nil { + t.Fatal(err) + } + dstJS, _ := nc.JetStream() + sameContent(t, m, srcJS, dstJS) + }) + + t.Run("into a new store, swapped in and served", func(t *testing.T) { + binary := filepath.Join(t.TempDir(), "mesh-nats-snapshot") + build := exec.Command("go", "build", "-o", binary, ".") + build.Env = append(os.Environ(), "CGO_ENABLED=0") + if out, err := build.CombinedOutput(); err != nil { + t.Fatalf("building the program: %v\n%s", err, out) + } + out := t.TempDir() + // The program in the bus's own image, as a person runs it: the archive in, a new store out. + docker(t, "run", "--rm", "--user", owner(), + "-v", binary+":/usr/local/bin/mesh-nats-snapshot:ro", + "-v", tarPath+":/in/bus.tar:ro", "-v", out+":/out", + "--entrypoint", "/usr/local/bin/mesh-nats-snapshot", image, + "restore", "--into", "/out/store", "--from", "/in/bus.tar") + if _, err := os.Stat(filepath.Join(out, "store", "jetstream", "MESH", "streams", "EVENTS")); err != nil { + t.Fatalf("the new store is not laid out as the bus's: %v", err) + } + served := startServer(t, plainConf(), filepath.Join(out, "store")+":/data") + nc := dial(t, "nats://127.0.0.1:"+served.port, Credential{User: "admin", Password: "admin"}, nil) + dstJS, _ := nc.JetStream() + sameContent(t, m, srcJS, dstJS) + }) + + t.Run("a damaged archive is refused", func(t *testing.T) { + damaged := append([]byte(nil), archive.Bytes()...) + // A byte in the last stream's archive, well past every header. + damaged[len(damaged)-2048] ^= 0xff + d, err := ReadArchive(bytes.NewReader(damaged), "") + if err == nil { + defer d.Close() + err = d.Verify() + } + if err == nil { + t.Fatal("a damaged archive verified") + } + }) +} + +// The bus's own module may snapshot and do nothing else: every API that changes a stream, every +// publish into one, and every subscription beyond its own inbox is refused by the server. +func TestTheSnapshotUserCannotChangeTheBus(t *testing.T) { + needDocker(t) + tlsDir := t.TempDir() + pin := selfSigned(t, tlsDir) + source := startServer(t, liveBusConf(), tlsDir+":/tls:ro") + url := "tls://127.0.0.1:" + source.port + admin := dial(t, url, Credential{Fingerprint: pin, User: "admin", Password: "admin"}, nil) + js, _ := admin.JetStream() + if _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}}); err != nil { + t.Fatal(err) + } + if _, err := js.Publish("mesh.mod.a.event.b", []byte("x")); err != nil { + t.Fatal(err) + } + if _, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "mesh-controller_conditions"}); err != nil { + t.Fatal(err) + } + + errs := make(chan error, 64) + snap := dial(t, url, Credential{Fingerprint: pin, User: snapshotUser, Password: "snap"}, errs) + for _, subject := range []string{ + "$JS.API.STREAM.CREATE.NEW", "$JS.API.STREAM.UPDATE.EVENTS", "$JS.API.STREAM.DELETE.EVENTS", + "$JS.API.STREAM.PURGE.EVENTS", "$JS.API.STREAM.MSG.DELETE.EVENTS", "$JS.API.STREAM.RESTORE.NEW", + "$JS.API.CONSUMER.CREATE.EVENTS", "$JS.API.STREAM.MSG.GET.EVENTS", + "mesh.mod.a.event.b", "$KV.mesh-controller_conditions.k", + } { + _, err := snap.Request(subject, []byte(`{}`), 500*time.Millisecond) + if err == nil { + t.Errorf("%s was answered for the snapshot user", subject) + } + select { + case e := <-errs: + if !strings.Contains(strings.ToLower(e.Error()), "permissions violation") { + t.Errorf("%s: %v", subject, e) + } + case <-time.After(2 * time.Second): + t.Errorf("the server did not refuse %s", subject) + } + } + for _, subject := range []string{"mesh.>", "_INBOX.other.>", "$JS.API.>"} { + if _, err := snap.SubscribeSync(subject); err != nil { + t.Fatal(err) + } + _ = snap.Flush() + select { + case e := <-errs: + if !strings.Contains(strings.ToLower(e.Error()), "permissions violation") { + t.Errorf("subscribing %s: %v", subject, e) + } + case <-time.After(2 * time.Second): + t.Errorf("the server let the snapshot user subscribe %s", subject) + } + } + info, err := js.StreamInfo("EVENTS") + if err != nil || info.State.Msgs != 1 { + t.Fatalf("the stream changed: %v %+v", err, info) + } + // And what it may do, it can. + names, err := StreamNames(context.Background(), snap) + if err != nil || len(names) != 2 { + t.Fatalf("the snapshot user cannot list the streams: %v %v", names, err) + } +} diff --git a/modules/nats/snapshot/main.go b/modules/nats/snapshot/main.go new file mode 100644 index 0000000..81c670b --- /dev/null +++ b/modules/nats/snapshot/main.go @@ -0,0 +1,209 @@ +// mesh-nats-snapshot: a consistent copy of the mesh bus's streams, and the way back from one +// (novox/hq ADR 0235, to-be 43). +// +// The bus keeps everything it holds in JetStream: every stream, and every key-value bucket, which is +// a stream too (conditions and their history, calls, hand-acts, the controller's lease, assignments, +// events, every module's state). Copying the store's files while the server writes them is not a +// copy that can be trusted to restore — a block half-written, an index from before the block it +// indexes. The server has its own way to hand a stream out whole: the snapshot API, a chunked +// transfer with flow control that the server answers from a consistent view of the stream while +// it goes on taking messages. This program asks for it, one stream at a time, and nothing else. +// +// It lives in the bus's image because that is where it runs: the machine's backup holder runs the +// nats module's declared dump — `docker exec` into the server's container — and the archive comes +// out on stdout, so nothing is mounted into the server for it and the credential arrives on stdin. +// +// mesh-nats-snapshot snapshot [--server URL] [--credential FILE|-] > bus.tar +// mesh-nats-snapshot verify [--from FILE] +// mesh-nats-snapshot restore --into DIR [--account MESH] [--from FILE] +// mesh-nats-snapshot restore --server URL --credential FILE [--from FILE] +// +// **Read-only on the live bus.** A snapshot neither pauses nor reconfigures a stream; the server +// keeps accepting writes while it streams the copy out. The credential it uses is the bus's own +// module's, granted the snapshot API and nothing that writes (mesh-controller internal/broker). +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "io" + "os" + "time" +) + +func main() { + if len(os.Args) < 2 { + usage() + os.Exit(2) + } + var err error + switch os.Args[1] { + case "snapshot": + err = runSnapshot(os.Args[2:]) + case "verify": + err = runVerify(os.Args[2:]) + case "restore": + err = runRestore(os.Args[2:]) + case "-h", "--help", "help": + usage() + return + default: + usage() + os.Exit(2) + } + if err != nil { + fmt.Fprintln(os.Stderr, "mesh-nats-snapshot:", err) + os.Exit(1) + } +} + +func usage() { + fmt.Fprint(os.Stderr, `mesh-nats-snapshot — a consistent copy of the bus's streams, and the way back + + snapshot [--server tls://127.0.0.1:4222] [--credential -] every stream, one at a time, as a tar on stdout + verify [--from bus.tar] every archive against the manifest + restore --into DIR [--account MESH] [--from bus.tar] into a new store directory, by a server of its own + restore --server URL --credential FILE [--from bus.tar] into a running server that holds none of them + +The credential is the JSON the mesh seals to the bus's own module: url, fingerprint, user, password. +`) +} + +func runSnapshot(args []string) error { + fs := flag.NewFlagSet("snapshot", flag.ContinueOnError) + server := fs.String("server", "tls://127.0.0.1:4222", "the bus, as seen from where this runs") + credential := fs.String("credential", "-", "the bus's own module's credential, a file or - for stdin") + streamTimeout := fs.Duration("stream-timeout", 10*time.Minute, "the longest one stream may take") + timeout := fs.Duration("timeout", 30*time.Minute, "the longest the whole snapshot may take") + pause := fs.Duration("pause", 250*time.Millisecond, "a rest between two streams") + tmp := fs.String("tmp", "", "where each stream's archive waits before it is written out (default: the system's)") + if err := fs.Parse(args); err != nil { + return err + } + cred, err := readCredential(*credential) + if err != nil { + return err + } + ctx, cancel := context.WithTimeout(context.Background(), *timeout) + defer cancel() + nc, err := connect(*server, cred) + if err != nil { + return err + } + defer nc.Close() + opts := SnapshotOptions{StreamTimeout: *streamTimeout, Pause: *pause, TempDir: *tmp, Log: logf} + m, err := Snapshot(ctx, nc, os.Stdout, opts) + if err != nil { + return err + } + logf("snapshot: %d stream(s), %d message(s), %s of streams in %s of archives, %.1fs", + len(m.Streams), m.Messages, human(m.StreamBytes), human(m.ArchiveBytes), m.Seconds) + return nil +} + +func runVerify(args []string) error { + fs := flag.NewFlagSet("verify", flag.ContinueOnError) + from := fs.String("from", "-", "the archive, a file or - for stdin") + if err := fs.Parse(args); err != nil { + return err + } + in, closer, err := openFrom(*from) + if err != nil { + return err + } + defer closer() + a, err := ReadArchive(in, "") + if err != nil { + return err + } + defer a.Close() + if err := a.Verify(); err != nil { + return err + } + logf("verified: %d stream(s), taken %s, every archive matches the manifest and reads to its end", + len(a.Manifest.Streams), a.Manifest.Finished.Format(time.RFC3339)) + return nil +} + +func runRestore(args []string) error { + fs := flag.NewFlagSet("restore", flag.ContinueOnError) + from := fs.String("from", "-", "the archive, a file or - for stdin") + into := fs.String("into", "", "a new store directory, filled by a server of this program's own on loopback") + account := fs.String("account", "MESH", "the account the streams belong to (the bus's one account)") + server := fs.String("server", "", "a running server to restore into instead; it must hold none of the streams") + credential := fs.String("credential", "", "the credential for --server, a file") + timeout := fs.Duration("timeout", time.Hour, "the longest the whole restore may take") + if err := fs.Parse(args); err != nil { + return err + } + if (*into == "") == (*server == "") { + return errors.New("restore needs exactly one of --into (a new store directory) or --server (a running server)") + } + in, closer, err := openFrom(*from) + if err != nil { + return err + } + defer closer() + a, err := ReadArchive(in, "") + if err != nil { + return err + } + defer a.Close() + if err := a.Verify(); err != nil { + return fmt.Errorf("the archive does not match its own manifest, so nothing is restored: %w", err) + } + ctx, cancel := context.WithTimeout(context.Background(), *timeout) + defer cancel() + if *into != "" { + where, err := RestoreOffline(ctx, a, *into, *account, logf) + if err != nil { + return err + } + logf("restored %d stream(s) into %s. With the bus stopped, put this directory in place of the "+ + "jetstream directory in the bus's store, then start it (novox/hq to-be 43).", len(a.Manifest.Streams), where) + return nil + } + if *credential == "" { + return errors.New("--server needs --credential") + } + cred, err := readCredential(*credential) + if err != nil { + return err + } + nc, err := connect(*server, cred) + if err != nil { + return err + } + defer nc.Close() + return RestoreAll(ctx, nc, a, logf) +} + +func openFrom(path string) (io.Reader, func(), error) { + if path == "-" || path == "" { + return os.Stdin, func() {}, nil + } + f, err := os.Open(path) + if err != nil { + return nil, nil, err + } + return f, func() { f.Close() }, nil +} + +func logf(format string, args ...any) { + fmt.Fprintf(os.Stderr, format+"\n", args...) +} + +func human(n int64) string { + const unit = 1024 + if n < unit { + return fmt.Sprintf("%d B", n) + } + div, exp := int64(unit), 0 + for v := n / unit; v >= unit; v /= unit { + div *= unit + exp++ + } + return fmt.Sprintf("%.1f %ciB", float64(n)/float64(div), "KMGTPE"[exp]) +} diff --git a/modules/nats/snapshot/restore.go b/modules/nats/snapshot/restore.go new file mode 100644 index 0000000..069325b --- /dev/null +++ b/modules/nats/snapshot/restore.go @@ -0,0 +1,261 @@ +package main + +import ( + "context" + "crypto/rand" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net" + "os" + "os/exec" + "path/filepath" + "strings" + "syscall" + "time" + + "github.com/nats-io/nats.go" +) + +// restoreChunk is what is sent at a time; the server answers each before the next goes. +const restoreChunk = 64 * 1024 + +// RestoreAll puts every stream of the archive into the server nc reaches, one at a time, and holds +// each restored stream's state to the manifest. **Nothing is restored over anything**: a server that +// already holds any of the streams is refused before the first is sent. +func RestoreAll(ctx context.Context, nc *nats.Conn, a *Archive, log func(string, ...any)) error { + var present []string + for _, s := range a.Manifest.Streams { + err := request(ctx, nc, apiStreamInfo+s.Name, nil, nil) + var api *apiError + switch { + case err == nil: + present = append(present, s.Name) + case errors.As(err, &api) && api.Code == 404: + default: + return err + } + } + if len(present) > 0 { + return fmt.Errorf("the server already holds %s; nothing is restored over a live stream — restore into a new "+ + "store (--into) and swap it in with the bus stopped", strings.Join(present, ", ")) + } + for _, s := range a.Manifest.Streams { + started := time.Now() + if err := restoreStream(ctx, nc, a, s); err != nil { + return fmt.Errorf("stream %s: %w", s.Name, err) + } + log("%s: restored, at least %d message(s) to sequence %d, in %.2fs", s.Name, s.Messages, s.LastSeq, time.Since(started).Seconds()) + } + return nil +} + +func restoreStream(ctx context.Context, nc *nats.Conn, a *Archive, s StreamEntry) error { + meta, archive, err := a.open(s) + if err != nil { + return err + } + defer archive.Close() + var req map[string]json.RawMessage + if err := json.Unmarshal(meta, &req); err != nil || req["config"] == nil { + return errors.New("its meta is not a stream's configuration and state") + } + var resp struct { + DeliverSubject string `json:"deliver_subject"` + } + if err := request(ctx, nc, apiStreamRestore+s.Name, req, &resp); err != nil { + return err + } + if resp.DeliverSubject == "" { + return errors.New("the server gave nowhere to send the archive") + } + buf := make([]byte, restoreChunk) + for { + n, err := archive.Read(buf) + if n > 0 { + cctx, cancel := context.WithTimeout(ctx, 30*time.Second) + reply, rerr := nc.RequestWithContext(cctx, resp.DeliverSubject, buf[:n]) + cancel() + if rerr != nil { + return fmt.Errorf("sending the archive: %w", rerr) + } + if len(reply.Data) > 0 { + return fmt.Errorf("the server refused the archive: %s", truncate(reply.Data)) + } + } + if err == io.EOF { + break + } + if err != nil { + return err + } + } + // The end: an empty message, answered once the server has rebuilt the stream from what it was + // sent — which for a large stream takes a while. + fctx, cancel := context.WithTimeout(ctx, 30*time.Minute) + defer cancel() + final, err := nc.RequestWithContext(fctx, resp.DeliverSubject, nil) + if err != nil { + return fmt.Errorf("finishing the restore: %w", err) + } + var created struct { + Error *apiError `json:"error"` + State streamState `json:"state"` + } + if err := json.Unmarshal(final.Data, &created); err != nil { + return fmt.Errorf("the server's last answer is not JSON: %q", truncate(final.Data)) + } + if created.Error != nil { + return created.Error + } + return heldTo(created.State, s) +} + +// heldTo is whether a restored stream holds at least what its snapshot said it held. +// +// **At least, not exactly — found by the restore test, not assumed.** The state the server answers a +// snapshot request with is taken when the snapshot starts; the stream goes on taking messages while +// its blocks are read out, and a block read a moment later carries some of what was appended in that +// moment. So a stream written during its snapshot restores with a later last sequence than the +// manifest says (90075 where it said 90024, with a writer publishing every 2ms through it), and of +// the messages in that tail some may be there and some not. Every message up to the manifest's last +// sequence is there, exactly; that is the snapshot's promise, and what the test holds it to message +// by message. What must never happen is a restored stream that ends before the manifest said it did. +func heldTo(restored streamState, s StreamEntry) error { + if restored.LastSeq < s.LastSeq || (s.Messages > 0 && restored.Messages == 0) { + return fmt.Errorf("restored with %d message(s) up to sequence %d; the snapshot held %d up to %d", + restored.Messages, restored.LastSeq, s.Messages, s.LastSeq) + } + return nil +} + +// RestoreOffline fills a new store directory with the archive's streams, by a server of this +// program's own: nats-server, the same binary the bus runs, on loopback, with the bus's one account +// and a user that exists only for this run. It is stopped when the streams are in. What it leaves +// is a store directory laid out exactly as the bus's — `jetstream//streams/` — which +// a person swaps in for the bus's own with the bus stopped. It answers that jetstream directory. +// +// Why a server of its own rather than the live bus: a stream is restored only where it does not +// exist, and on the live bus every stream exists — the controller asserts them on every start. So a +// restore beside the live data, swapped in by a person, is the only one that never touches it. +func RestoreOffline(ctx context.Context, a *Archive, into, account string, log func(string, ...any)) (string, error) { + if !validName(account) { + return "", fmt.Errorf("%q is not an account name", account) + } + if entries, err := os.ReadDir(into); err == nil && len(entries) > 0 { + return "", fmt.Errorf("%s is not empty; a restore goes into a new directory, never over one", into) + } + if err := os.MkdirAll(into, 0o700); err != nil { + return "", err + } + store, err := filepath.Abs(into) + if err != nil { + return "", err + } + binary, err := exec.LookPath("nats-server") + if err != nil { + return "", errors.New("no nats-server here to restore with; run this in the bus's own image") + } + work, err := os.MkdirTemp("", "mesh-nats-restore-server-*") + if err != nil { + return "", err + } + defer os.RemoveAll(work) + port, err := freePort() + if err != nil { + return "", err + } + password := randomHex(24) + conf := filepath.Join(work, "restore.conf") + config := fmt.Sprintf(`listen: "127.0.0.1:%d" +jetstream { store_dir: %q } +accounts { %s { jetstream: enabled, users: [ { user: "restore", password: %q } ] } } +`, port, store, account, password) + if err := os.WriteFile(conf, []byte(config), 0o600); err != nil { + return "", err + } + cmd := exec.Command(binary, "-c", conf) + logFile, err := os.Create(filepath.Join(work, "server.log")) + if err != nil { + return "", err + } + defer logFile.Close() + cmd.Stdout, cmd.Stderr = logFile, logFile + if err := cmd.Start(); err != nil { + return "", err + } + exited := make(chan error, 1) + go func() { exited <- cmd.Wait() }() + stop := func() error { + _ = cmd.Process.Signal(syscall.SIGTERM) + select { + case err := <-exited: + return err + case <-time.After(30 * time.Second): + _ = cmd.Process.Kill() + <-exited + return errors.New("the restore's server did not stop within 30s and was killed") + } + } + serverLog := func() string { + raw, _ := os.ReadFile(filepath.Join(work, "server.log")) + return strings.TrimSpace(string(raw)) + } + + url := fmt.Sprintf("nats://127.0.0.1:%d", port) + var nc *nats.Conn + for deadline := time.Now().Add(20 * time.Second); ; { + nc, err = connect(url, Credential{User: "restore", Password: password}) + if err == nil { + break + } + if time.Now().After(deadline) { + _ = stop() + return "", fmt.Errorf("the restore's server did not come up: %v\n%s", err, serverLog()) + } + select { + case err := <-exited: + return "", fmt.Errorf("the restore's server stopped: %v\n%s", err, serverLog()) + case <-time.After(200 * time.Millisecond): + } + } + err = RestoreAll(ctx, nc, a, log) + nc.Close() + if stopErr := stop(); err == nil && stopErr != nil && !isTerminated(stopErr) { + err = stopErr + } + if err != nil { + return "", err + } + return filepath.Join(store, "jetstream"), nil +} + +// isTerminated is a server that stopped because it was asked to. +func isTerminated(err error) bool { + var exit *exec.ExitError + if errors.As(err, &exit) { + if status, ok := exit.Sys().(syscall.WaitStatus); ok && status.Signaled() && status.Signal() == syscall.SIGTERM { + return true + } + } + return false +} + +func freePort() (int, error) { + l, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + return 0, err + } + defer l.Close() + return l.Addr().(*net.TCPAddr).Port, nil +} + +func randomHex(n int) string { + b := make([]byte, n) + if _, err := rand.Read(b); err != nil { + panic(err) + } + return hex.EncodeToString(b) +} diff --git a/modules/nats/snapshot/snapshot.go b/modules/nats/snapshot/snapshot.go new file mode 100644 index 0000000..9793a67 --- /dev/null +++ b/modules/nats/snapshot/snapshot.go @@ -0,0 +1,446 @@ +package main + +import ( + "archive/tar" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "sort" + "strings" + "sync" + "time" + + "github.com/klauspost/compress/s2" + "github.com/nats-io/nats.go" +) + +// The JetStream API's subjects this program speaks. Each is granted to the bus's own module and +// nothing else is (mesh-controller internal/broker, BusSnapshotGrants); a test holds the two lists +// equal against a real server. +const ( + apiStreamNames = "$JS.API.STREAM.NAMES" + apiStreamInfo = "$JS.API.STREAM.INFO." + apiStreamSnapshot = "$JS.API.STREAM.SNAPSHOT." + apiStreamRestore = "$JS.API.STREAM.RESTORE." +) + +// ChunkSize is what the server is asked to send at a time: its own default. The server sends up to +// an 8 MiB window ahead of the acknowledgements and waits for them past it, so one stream costs the +// server at most that much in flight, whatever its size. +const ChunkSize = 128 * 1024 + +// SnapshotOptions bounds a snapshot. +type SnapshotOptions struct { + // StreamTimeout is the longest one stream may take; past it the snapshot fails, naming it. + StreamTimeout time.Duration + // Pause is a rest between two streams, so a night's snapshot is never one long burst. + Pause time.Duration + // TempDir is where one stream's archive waits until it is written out; the system's when empty. + TempDir string + Log func(format string, args ...any) +} + +// Manifest is the first entry of an archive: what was taken, when, how long it took, and the checksum +// of every file beside it. It is what `verify` and `restore` hold the archive to. +type Manifest struct { + Format int `json:"format"` + Server string `json:"server_version"` + Started time.Time `json:"started"` + Finished time.Time `json:"finished"` + Seconds float64 `json:"seconds"` + // Messages and StreamBytes are the streams' own totals at their snapshots; ArchiveBytes what the + // archives take, compressed. + Messages uint64 `json:"messages"` + StreamBytes int64 `json:"stream_bytes"` + ArchiveBytes int64 `json:"archive_bytes"` + Streams []StreamEntry `json:"streams"` + // Skipped are streams that hold nothing across a restart (memory storage), said rather than + // silently left out. + Skipped []Skipped `json:"skipped,omitempty"` +} + +// StreamEntry is one stream as it was snapshotted. +type StreamEntry struct { + Name string `json:"name"` + Subjects []string `json:"subjects,omitempty"` + Retention string `json:"retention,omitempty"` + Messages uint64 `json:"messages"` + Bytes uint64 `json:"bytes"` + FirstSeq uint64 `json:"first_seq"` + LastSeq uint64 `json:"last_seq"` + Deleted int `json:"deleted,omitempty"` + Consumers int `json:"consumers"` + // Meta is the server's own description of the stream (its configuration and state), in the + // shape the restore API and the nats CLI take back: backup.json. + Meta string `json:"meta"` + MetaSHA256 string `json:"meta_sha256"` + // Archive is the server's snapshot of the stream, as the server sent it: stream.tar.s2. + Archive string `json:"archive"` + ArchiveBytes int64 `json:"archive_bytes"` + ArchiveSHA256 string `json:"archive_sha256"` + Seconds float64 `json:"seconds"` +} + +// Skipped is a stream not snapshotted, and why. +type Skipped struct { + Name string `json:"name"` + Why string `json:"why"` +} + +// ManifestName is the archive's first entry. +const ManifestName = "manifest.json" + +// apiError is the error every JetStream API answer may carry. +type apiError struct { + Code int `json:"code"` + ErrCode int `json:"err_code"` + Description string `json:"description"` +} + +func (e *apiError) Error() string { return fmt.Sprintf("%s (%d)", e.Description, e.Code) } + +// streamState is what this program reads of a stream's state; the whole of it is kept as the server +// gave it, in the stream's meta. +type streamState struct { + Messages uint64 `json:"messages"` + Bytes uint64 `json:"bytes"` + FirstSeq uint64 `json:"first_seq"` + LastSeq uint64 `json:"last_seq"` + NumDelete int `json:"num_deleted"` + Consumers int `json:"consumer_count"` +} + +type streamConfig struct { + Name string `json:"name"` + Subjects []string `json:"subjects"` + Retention string `json:"retention"` + Storage string `json:"storage"` +} + +// request asks the JetStream API and decodes its answer, refusing an answer that carries an error. +func request(ctx context.Context, nc *nats.Conn, subject string, body any, into any) error { + var payload []byte + if body != nil { + var err error + if payload, err = json.Marshal(body); err != nil { + return err + } + } + rctx, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + msg, err := nc.RequestWithContext(rctx, subject, payload) + if err != nil { + if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, nats.ErrTimeout) || errors.Is(err, nats.ErrNoResponders) { + return fmt.Errorf("%s was not answered (%v): the bus refuses a subject this user is not granted by not answering it", subject, err) + } + return fmt.Errorf("%s: %w", subject, err) + } + var wrapped struct { + Error *apiError `json:"error"` + } + if err := json.Unmarshal(msg.Data, &wrapped); err != nil { + return fmt.Errorf("%s answered something that is not JSON: %q", subject, truncate(msg.Data)) + } + if wrapped.Error != nil { + return fmt.Errorf("%s: %w", subject, wrapped.Error) + } + if into != nil { + return json.Unmarshal(msg.Data, into) + } + return nil +} + +func truncate(b []byte) string { + if len(b) > 200 { + return string(b[:200]) + "…" + } + return string(b) +} + +// StreamNames is every stream in the account, in name order, all pages of them. +func StreamNames(ctx context.Context, nc *nats.Conn) ([]string, error) { + var names []string + for { + var page struct { + Total int `json:"total"` + Offset int `json:"offset"` + Limit int `json:"limit"` + Streams []string `json:"streams"` + } + if err := request(ctx, nc, apiStreamNames, map[string]int{"offset": len(names)}, &page); err != nil { + return nil, err + } + names = append(names, page.Streams...) + if len(page.Streams) == 0 || len(names) >= page.Total { + break + } + } + sort.Strings(names) + return names, nil +} + +// Snapshot writes every stream, one at a time, as a tar: the manifest first, then each stream's +// meta and archive under streams//. Any stream failing fails the whole: a partial copy of the +// bus written where a whole one is expected is the silent failure this exists to prevent. +func Snapshot(ctx context.Context, nc *nats.Conn, out io.Writer, o SnapshotOptions) (*Manifest, error) { + if o.StreamTimeout <= 0 { + o.StreamTimeout = 10 * time.Minute + } + if o.Log == nil { + o.Log = func(string, ...any) {} + } + m := &Manifest{Format: 1, Server: nc.ConnectedServerVersion(), Started: time.Now().UTC()} + names, err := StreamNames(ctx, nc) + if err != nil { + return nil, err + } + if len(names) == 0 { + // The mesh's own streams always exist once the controller has run; a bus with none is not + // one to call backed up. + return nil, errors.New("the bus holds no streams; a snapshot of nothing would pass for a backup") + } + + type taken struct { + entry StreamEntry + meta []byte + archive *os.File + } + var files []taken + defer func() { + for _, t := range files { + t.archive.Close() + os.Remove(t.archive.Name()) + } + }() + + for i, name := range names { + if i > 0 && o.Pause > 0 { + select { + case <-time.After(o.Pause): + case <-ctx.Done(): + return nil, ctx.Err() + } + } + var info struct { + Config streamConfig `json:"config"` + } + if err := request(ctx, nc, apiStreamInfo+name, nil, &info); err != nil { + return nil, err + } + if info.Config.Storage == "memory" { + m.Skipped = append(m.Skipped, Skipped{Name: name, Why: "memory storage: it holds nothing across a restart, and the server cannot snapshot it"}) + o.Log("%s: skipped, memory storage", name) + continue + } + f, err := os.CreateTemp(o.TempDir, "mesh-nats-snapshot-*") + if err != nil { + return nil, err + } + files = append(files, taken{archive: f}) + sctx, cancel := context.WithTimeout(ctx, o.StreamTimeout) + entry, meta, err := snapshotStream(sctx, nc, name, f) + cancel() + if err != nil { + return nil, fmt.Errorf("stream %s: %w", name, err) + } + files[len(files)-1].entry, files[len(files)-1].meta = entry, meta + m.Messages += entry.Messages + m.StreamBytes += int64(entry.Bytes) + m.ArchiveBytes += entry.ArchiveBytes + o.Log("%s: %d message(s), %s, last sequence %d, %d consumer(s) → %s in %.2fs", + name, entry.Messages, human(int64(entry.Bytes)), entry.LastSeq, entry.Consumers, human(entry.ArchiveBytes), entry.Seconds) + } + m.Finished = time.Now().UTC() + m.Seconds = m.Finished.Sub(m.Started).Seconds() + for _, t := range files { + m.Streams = append(m.Streams, t.entry) + } + + tw := tar.NewWriter(out) + manifest, err := json.MarshalIndent(m, "", " ") + if err != nil { + return nil, err + } + if err := writeEntry(tw, ManifestName, append(manifest, '\n'), m.Finished); err != nil { + return nil, err + } + for _, t := range files { + if err := writeEntry(tw, t.entry.Meta, t.meta, m.Finished); err != nil { + return nil, err + } + if _, err := t.archive.Seek(0, io.SeekStart); err != nil { + return nil, err + } + hdr := &tar.Header{Name: t.entry.Archive, Mode: 0o600, Size: t.entry.ArchiveBytes, ModTime: m.Finished, Format: tar.FormatPAX} + if err := tw.WriteHeader(hdr); err != nil { + return nil, err + } + if _, err := io.Copy(tw, t.archive); err != nil { + return nil, err + } + } + if err := tw.Close(); err != nil { + return nil, err + } + return m, nil +} + +func writeEntry(tw *tar.Writer, name string, body []byte, at time.Time) error { + if err := tw.WriteHeader(&tar.Header{Name: name, Mode: 0o600, Size: int64(len(body)), ModTime: at, Format: tar.FormatPAX}); err != nil { + return err + } + _, err := tw.Write(body) + return err +} + +// snapshotStream takes one stream's snapshot into f: the server's chunks, each acknowledged as it +// arrives so the server's window keeps moving, until the server's last, empty message. It answers +// the stream's entry (the archive's size and checksum included) and its meta. +// +// The protocol is the server's (nats-server jetstream_api.go, streamSnapshot): the request names a +// subject to deliver to; the server answers with the stream's configuration and its state at the +// snapshot, then sends the archive in chunks whose reply subject is the acknowledgement it waits for +// past its window — two seconds without one and it gives up with "408 No Flow Response". The end is +// an empty message whose status header says 204, or the error that ended it. +func snapshotStream(ctx context.Context, nc *nats.Conn, name string, f *os.File) (StreamEntry, []byte, error) { + started := time.Now() + deliver := nc.NewRespInbox() + + hash := sha256.New() + var written int64 + var mu sync.Mutex + var failed error + done := make(chan struct{}) + var once sync.Once + finish := func(err error) { + mu.Lock() + if failed == nil { + failed = err + } + mu.Unlock() + once.Do(func() { close(done) }) + } + sub, err := nc.Subscribe(deliver, func(msg *nats.Msg) { + if len(msg.Data) == 0 { + if status := msg.Header.Get("Status"); status != "" && status != "204" { + finish(fmt.Errorf("the server ended the snapshot: %s %s", status, msg.Header.Get("Description"))) + return + } + finish(nil) + return + } + // Acknowledged before it is written: the chunk is already in memory, and a slow disk must not + // read to the server as a reader that went away. + if msg.Reply != "" { + if err := nc.Publish(msg.Reply, nil); err != nil { + finish(fmt.Errorf("acknowledging a chunk: %w", err)) + return + } + } + if _, err := f.Write(msg.Data); err != nil { + finish(fmt.Errorf("writing the archive: %w", err)) + return + } + hash.Write(msg.Data) + written += int64(len(msg.Data)) + }) + if err != nil { + return StreamEntry{}, nil, err + } + defer sub.Unsubscribe() + // Every chunk is held until the callback has it; the window bounds how many there can be. + if err := sub.SetPendingLimits(-1, -1); err != nil { + return StreamEntry{}, nil, err + } + if err := nc.FlushWithContext(ctx); err != nil { + return StreamEntry{}, nil, err + } + + var resp struct { + Config json.RawMessage `json:"config"` + State json.RawMessage `json:"state"` + } + req := map[string]any{"deliver_subject": deliver, "chunk_size": ChunkSize} + if err := request(ctx, nc, apiStreamSnapshot+name, req, &resp); err != nil { + return StreamEntry{}, nil, err + } + select { + case <-done: + case <-ctx.Done(): + return StreamEntry{}, nil, fmt.Errorf("not finished within its bound: %w", ctx.Err()) + } + // The subscription is drained before the archive is read back, so no callback is still writing. + _ = sub.Unsubscribe() + mu.Lock() + err = failed + mu.Unlock() + if err != nil { + return StreamEntry{}, nil, err + } + if written == 0 { + return StreamEntry{}, nil, errors.New("the server sent an empty archive") + } + if err := f.Sync(); err != nil { + return StreamEntry{}, nil, err + } + + var cfg streamConfig + var state streamState + if err := json.Unmarshal(resp.Config, &cfg); err != nil { + return StreamEntry{}, nil, fmt.Errorf("its configuration: %w", err) + } + if err := json.Unmarshal(resp.State, &state); err != nil { + return StreamEntry{}, nil, fmt.Errorf("its state: %w", err) + } + // The archive reads to its end — s2 frames, then a tar — before it is called a snapshot. + if _, err := f.Seek(0, io.SeekStart); err != nil { + return StreamEntry{}, nil, err + } + if err := readsThrough(f); err != nil { + return StreamEntry{}, nil, fmt.Errorf("the archive the server sent does not read: %w", err) + } + + // The meta, in the shape the restore API takes and the nats CLI writes (backup.json): the + // server's own configuration and state, kept as raw JSON so a field this program does not know + // is never dropped on the way back. + meta, err := json.MarshalIndent(map[string]json.RawMessage{"config": resp.Config, "state": resp.State}, "", " ") + if err != nil { + return StreamEntry{}, nil, err + } + metaSum := sha256.Sum256(meta) + return StreamEntry{ + Name: name, Subjects: cfg.Subjects, Retention: cfg.Retention, + Messages: state.Messages, Bytes: state.Bytes, FirstSeq: state.FirstSeq, LastSeq: state.LastSeq, + Deleted: state.NumDelete, Consumers: state.Consumers, + Meta: "streams/" + name + "/backup.json", MetaSHA256: hex.EncodeToString(metaSum[:]), + Archive: "streams/" + name + "/stream.tar.s2", ArchiveBytes: written, + ArchiveSHA256: hex.EncodeToString(hash.Sum(nil)), Seconds: time.Since(started).Seconds(), + }, meta, nil +} + +// readsThrough is whether a stream's archive decompresses and lists to its end. +func readsThrough(r io.Reader) error { + tr := tar.NewReader(s2.NewReader(r)) + for { + _, err := tr.Next() + if err == io.EOF { + return nil + } + if err != nil { + return err + } + if _, err := io.Copy(io.Discard, tr); err != nil { + return err + } + } +} + +// validName is a stream name that is one safe path component. +func validName(name string) bool { + return name != "" && name != "." && name != ".." && !strings.ContainsAny(name, "/\\\x00 *>") +}