Files
mesh-catalog/modules/nats/snapshot/main.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

210 lines
7.2 KiB
Go

// 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])
}