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