Merge pull request 'Back up the bus by the server's own snapshot of each stream, not its live files (hq ADR 0235)' (#93) from feat/bus-snapshot into main

This commit was merged in pull request #93.
This commit is contained in:
2026-10-06 17:33:13 +00:00
12 changed files with 2085 additions and 4 deletions
+20 -3
View File
@@ -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"]
+87
View File
@@ -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_<bucket>`) — 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 < <the module's broker credential> > <snapshots>/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/<name>/backup.json` — the stream's configuration and state, as the server gave them;
- `streams/<name>/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
`<snapshots>.restored-<stamp>/`, holding that night's `bus.tar`.
2. **Check it.** `docker exec -i mesh-broker-nats mesh-nats-snapshot verify < <restored>/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 <restored>:/in:ro -v <new-store>:/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
`<new-store>/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-data>/jetstream <jetstream-data>/jetstream.before-restore-<stamp> && \
mv <new-store>/jetstream <jetstream-data>/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 <url> --credential <file>`; the archive is also what the nats CLI
reads (`nats stream restore <dir>` on an unpacked `streams/<name>/`).
## 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=<this image>` also runs the manifest's declared dump, exactly, against the
module's image and configuration.
- `cmd/nats-tools`: `go test ./...`.
+29 -1
View File
@@ -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": [
+162
View File
@@ -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
}
+105
View File
@@ -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.<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
},
}
}
+15
View File
@@ -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
)
+12
View File
@@ -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=
+152
View File
@@ -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=<digest> --build-arg GO_BASE=<digest> -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")
}
}
+587
View File
@@ -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)
}
}
+209
View File
@@ -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])
}
+261
View File
@@ -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/<account>/streams/<name>` — 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)
}
+446
View File
@@ -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/<name>/. 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 *>")
}