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

153 lines
5.3 KiB
Go

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")
}
}