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.
153 lines
5.3 KiB
Go
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")
|
|
}
|
|
}
|