The night's backup of the bus takes each stream through JetStream's snapshot API, run by the nats module under its own account. The module holding mesh-broker is composed that account: stream names and info, the snapshot request, its flow-control acks, its own inbox — no write, which the writers table checks. A bus module declaring anything else to say on the bus is refused by module check rather than silently granted nothing. The genesis user list is unchanged: the controller's grants are.
184 lines
6.2 KiB
Go
184 lines
6.2 KiB
Go
package broker
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
"golang.org/x/crypto/bcrypt"
|
|
)
|
|
|
|
// The bus's own module's composed authority, against the bus's release in a throwaway container
|
|
// (novox/hq ADR 0235): it can take a whole snapshot of a stream through the server's chunked
|
|
// protocol, and every write — to a stream's definition, its messages, a bucket — is refused.
|
|
//
|
|
// MESH_TEST_DOCKER=1 go test ./internal/broker/ -run TestTheComposedSnapshotUser
|
|
//
|
|
// The program that takes the night's snapshot lives in the nats module (mesh-catalog
|
|
// modules/nats/snapshot) and is tested there against these same grants; this proves the grants are
|
|
// what the controller composes, by composing them.
|
|
func TestTheComposedSnapshotUserCanSnapshotAndCannotWrite(t *testing.T) {
|
|
if os.Getenv("MESH_TEST_DOCKER") != "1" {
|
|
t.Skip("MESH_TEST_DOCKER is not 1: this starts a throwaway nats container")
|
|
}
|
|
hash := func(pw string) string {
|
|
h, err := bcrypt.GenerateFromPassword([]byte(pw), bcrypt.MinCost)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return string(h)
|
|
}
|
|
accounts, err := ComposeAccounts([]Principal{
|
|
{Kind: KindController, PasswordHash: hash("controller")},
|
|
{Kind: KindModule, Node: "anchor", Module: "nats", SnapshotsTheBus: true, PasswordHash: hash("snap")},
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
conf, data := t.TempDir(), t.TempDir()
|
|
_ = os.WriteFile(filepath.Join(conf, "accounts.conf"), []byte(accounts), 0o644)
|
|
_ = os.WriteFile(filepath.Join(conf, "nats.conf"), []byte("port: 4222\njetstream { store_dir: \"/data\" }\ninclude accounts.conf\n"), 0o644)
|
|
name := fmt.Sprintf("mesh-controller-snapshot-test-%d", time.Now().UnixNano())
|
|
run := exec.Command("docker", "run", "-d", "--rm", "--name", name, "--user", fmt.Sprintf("%d:%d", os.Getuid(), os.Getgid()),
|
|
"-p", "127.0.0.1::4222", "-v", conf+":/etc/nats:ro", "-v", data+":/data", "nats:2.11-alpine", "-c", "/etc/nats/nats.conf")
|
|
if out, err := run.CombinedOutput(); err != nil {
|
|
t.Fatalf("%v\n%s", err, out)
|
|
}
|
|
t.Cleanup(func() { _ = exec.Command("docker", "rm", "-f", name).Run() })
|
|
out, err := exec.Command("docker", "port", name, "4222/tcp").Output()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
_, port, _ := strings.Cut(strings.TrimSpace(strings.Split(string(out), "\n")[0]), ":")
|
|
url := "nats://127.0.0.1:" + port
|
|
|
|
dial := func(user, pw string, errs chan error) *nats.Conn {
|
|
var nc *nats.Conn
|
|
var err error
|
|
for deadline := time.Now().Add(15 * time.Second); ; time.Sleep(200 * time.Millisecond) {
|
|
nc, err = nats.Connect(url, nats.UserInfo(user, pw), nats.CustomInboxPrefix("_INBOX."+user),
|
|
nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) {
|
|
if errs != nil {
|
|
select {
|
|
case errs <- err:
|
|
default:
|
|
}
|
|
}
|
|
}))
|
|
if err == nil || time.Now().After(deadline) {
|
|
break
|
|
}
|
|
}
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(nc.Close)
|
|
return nc
|
|
}
|
|
|
|
// The controller defines the streams and fills them, as it does.
|
|
controller := dial("controller", "controller", nil)
|
|
js, _ := controller.JetStream()
|
|
if _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
kv, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: HandActsBucket})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := kv.Put("act", []byte("done by hand")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
errs := make(chan error, 32)
|
|
snap := dial("anchor.nats", "snap", errs)
|
|
|
|
// What it may do: list, read a stream's information, and take a whole snapshot.
|
|
names, err := snap.Request("$JS.API.STREAM.NAMES", nil, 5*time.Second)
|
|
if err != nil || !strings.Contains(string(names.Data), "KV_"+HandActsBucket) {
|
|
t.Fatalf("it cannot list the streams: %v", err)
|
|
}
|
|
deliver := snap.NewRespInbox()
|
|
done := make(chan string, 1)
|
|
var bytes int
|
|
sub, err := snap.Subscribe(deliver, func(m *nats.Msg) {
|
|
if len(m.Data) == 0 {
|
|
done <- m.Header.Get("Status")
|
|
return
|
|
}
|
|
bytes += len(m.Data)
|
|
if m.Reply != "" {
|
|
_ = snap.Publish(m.Reply, nil)
|
|
}
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer sub.Unsubscribe()
|
|
req, _ := json.Marshal(map[string]any{"deliver_subject": deliver})
|
|
answer, err := snap.Request("$JS.API.STREAM.SNAPSHOT.KV_"+HandActsBucket, req, 5*time.Second)
|
|
if err != nil || strings.Contains(string(answer.Data), `"error"`) {
|
|
t.Fatalf("the snapshot was not granted: %v", err)
|
|
}
|
|
select {
|
|
case status := <-done:
|
|
if status != "" && status != "204" {
|
|
t.Fatalf("the snapshot ended with %s", status)
|
|
}
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatal("the snapshot never finished")
|
|
}
|
|
if bytes == 0 {
|
|
t.Fatal("the snapshot delivered nothing")
|
|
}
|
|
select {
|
|
case e := <-errs:
|
|
t.Fatalf("taking a snapshot met a refusal: %v", e)
|
|
default:
|
|
}
|
|
|
|
// What it may not: every write, and every subscription beyond its own inbox.
|
|
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", "$KV." + HandActsBucket + ".act",
|
|
"mesh.mod.nats.event.anything", "mesh.node.anchor.declare"} {
|
|
if _, err := snap.Request(subject, []byte(`{}`), 300*time.Millisecond); err == nil {
|
|
t.Errorf("%s was answered", 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 bus did not refuse %s", subject)
|
|
}
|
|
}
|
|
for _, subject := range []string{"mesh.>", "_INBOX.controller.>", "$KV.>"} {
|
|
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 bus let it subscribe %s", subject)
|
|
}
|
|
}
|
|
if info, err := js.StreamInfo("EVENTS"); err != nil || info == nil {
|
|
t.Fatalf("the stream is gone: %v", err)
|
|
}
|
|
if e, err := kv.Get("act"); err != nil || string(e.Value()) != "done by hand" {
|
|
t.Fatalf("the bucket changed: %v", err)
|
|
}
|
|
}
|