package broker import ( "errors" "os" "path/filepath" "strings" "testing" "time" "github.com/nats-io/nats-server/v2/server" "github.com/nats-io/nats.go" "golang.org/x/crypto/bcrypt" ) // The view against a real server, over WebSocket (novox/hq research 036): the composed user list is // what the server reads, the listener is the shape the bus module declares (no TLS, compression on, // reached across the overlay only), and the view with its credential binds the issue tracker's bucket, // lists it with one batch read, follows a tracker event to re-read the key it names — and is refused // every write: a put, a delete, an event, and a consumer delivering onto the tracker's event subject. // // go test ./internal/broker/ -run TestTheView func TestTheViewReadsTheIssuesOverWebSocketAndWritesNothing(t *testing.T) { hash := func(password string) string { h, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.MinCost) if err != nil { t.Fatal(err) } return string(h) } const node = "anchor" tracker := Principal{Kind: KindModule, Node: node, Module: "mesh-issues", State: []string{"issues"}, Emits: []string{"opened", "moved"}, PasswordHash: hash("tracker")} accounts, err := ComposeAccounts([]Principal{ {Kind: KindController, PasswordHash: hash("controller")}, tracker, {Kind: KindView, PasswordHash: hash("view")}, }) if err != nil { t.Fatal(err) } conf := filepath.Join(t.TempDir(), "accounts.conf") if err := os.WriteFile(conf, []byte(accounts), 0o600); err != nil { t.Fatal(err) } // The server reads the composed file as the bus does — through its own parser — and listens as the // bus module's configuration says: a WebSocket listener without TLS and with compression, beside // the client port. Ports chosen by the system, so this runs beside a live bus. opts, err := server.ProcessConfigFile(conf) if err != nil { t.Fatalf("the server refused the composed user list: %v", err) } opts.Host, opts.Port = "127.0.0.1", server.RANDOM_PORT opts.JetStream, opts.StoreDir = true, t.TempDir() opts.NoLog, opts.NoSigs = true, true opts.Websocket = server.WebsocketOpts{Host: "127.0.0.1", Port: server.RANDOM_PORT, NoTLS: true, Compression: true} s, err := server.NewServer(opts) if err != nil { t.Fatal(err) } go s.Start() if !s.ReadyForConnections(30 * time.Second) { s.Shutdown() t.Fatal("the server did not come up") } t.Cleanup(func() { s.Shutdown(); s.WaitForShutdown() }) dial := func(url, user, password string, refused chan<- string) *nats.Conn { t.Helper() nc, err := nats.Connect(url, nats.UserInfo(user, password), nats.CustomInboxPrefix("_INBOX."+user), nats.Compression(true), nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) { if refused != nil && errors.Is(err, nats.ErrPermissionViolation) { refused <- err.Error() } })) if err != nil { t.Fatalf("%s could not connect to %s: %v", user, url, err) } t.Cleanup(nc.Close) return nc } // The controller defines the bucket, as it does for every module's state; the tracker writes it. controller := dial(s.ClientURL(), "controller", "controller", nil) cjs, _ := controller.JetStream() if _, err := cjs.CreateKeyValue(&nats.KeyValueConfig{Bucket: ViewBucket, History: 8}); err != nil { t.Fatalf("the controller could not define %s: %v", ViewBucket, err) } trackerConn := dial(s.ClientURL(), tracker.Username(), "tracker", nil) tjs, _ := trackerConn.JetStream() tkv, err := tjs.KeyValue(ViewBucket) if err != nil { t.Fatal(err) } if _, err := tkv.Put("365", []byte(`{"number":365,"status":"open"}`)); err != nil { t.Fatalf("the tracker could not write its own bucket: %v", err) } // The view, over WebSocket with its credential. refused := make(chan string, 8) view := dial(s.WebsocketURL(), ViewUser, "view", refused) if !strings.HasPrefix(view.ConnectedUrl(), "ws://") { t.Fatalf("the view is connected to %s, not over WebSocket", view.ConnectedUrl()) } vjs, _ := view.JetStream(nats.MaxWait(3 * time.Second)) vkv, err := vjs.KeyValue(ViewBucket) if err != nil { t.Fatalf("the view could not bind %s: %v", ViewBucket, err) } if got, err := vkv.Get("365"); err != nil { t.Fatalf("the view could not read a key: %v", err) } else if !strings.Contains(string(got.Value()), `"number":365`) { t.Fatalf("the view read %q", got.Value()) } if _, err := tkv.Put("366", []byte(`{"number":366,"status":"open"}`)); err != nil { t.Fatal(err) } // The list: one batch read, the newest value of every key, into the view's own inbox, ended by the // server's end-of-batch status (204). listed := map[string]string{} inbox := view.NewRespInbox() batch, err := view.SubscribeSync(inbox) if err != nil { t.Fatal(err) } if err := view.PublishRequest("$JS.API.DIRECT.GET.KV_"+ViewBucket, inbox, []byte(`{"multi_last":["$KV.`+ViewBucket+`.>"]}`)); err != nil { t.Fatal(err) } for { m, err := batch.NextMsg(5 * time.Second) if err != nil { t.Fatalf("the view's batch read ended without its end-of-batch (%d keys so far): %v", len(listed), err) } if m.Header.Get("Status") == "204" { break } if status := m.Header.Get("Status"); status != "" { t.Fatalf("the batch read answered %s %s", status, m.Header.Get("Description")) } listed[m.Header.Get("Nats-Subject")] = string(m.Data) } _ = batch.Unsubscribe() for _, key := range []string{"365", "366"} { if !strings.Contains(listed["$KV."+ViewBucket+"."+key], `"number":`+key) { t.Errorf("the batch read did not list %s: %v", key, listed) } } // A change followed: the tracker moves 365 and says so; the view hears the event and reads the key // it names. moved, err := view.SubscribeSync("mesh.mod.mesh-issues.event.moved") if err != nil { t.Fatal(err) } _ = view.Flush() if _, err := tkv.Put("365", []byte(`{"number":365,"status":"located"}`)); err != nil { t.Fatal(err) } if err := trackerConn.Publish("mesh.mod.mesh-issues.event.moved", []byte(`{"number":365,"to":"located"}`)); err != nil { t.Fatal(err) } if _, err := moved.NextMsg(5 * time.Second); err != nil { t.Fatalf("the view did not hear the tracker's event: %v", err) } if got, err := vkv.Get("365"); err != nil || !strings.Contains(string(got.Value()), "located") { t.Fatalf("after the event the view read %v (%v)", got, err) } // **The hole a watch would open, shut** (ViewReads): a consumer delivering onto the tracker's event // subject would have the server republish the bucket there, as the tracker. Refused, and nothing // reaches a module listening on that subject. listener, err := trackerConn.SubscribeSync("mesh.mod.mesh-issues.event.opened") if err != nil { t.Fatal(err) } _ = trackerConn.Flush() if _, err := vjs.AddConsumer("KV_"+ViewBucket, &nats.ConsumerConfig{Name: "w", DeliverSubject: "mesh.mod.mesh-issues.event.opened", AckPolicy: nats.AckNonePolicy, FilterSubject: "$KV." + ViewBucket + ".>"}); err == nil { t.Error("the view made a consumer") } if _, err := vkv.WatchAll(); err == nil { t.Error("the view made a watch, which is a consumer") } if m, err := listener.NextMsg(2 * time.Second); err == nil { t.Errorf("a message reached the tracker's event subject from the view: %q", m.Data) } // The refusals of the consumer create land as permission violations too; drained before the writes. drain := time.After(500 * time.Millisecond) for draining := true; draining; { select { case <-refused: case <-drain: draining = false } } // And every write is refused: the server says so, and the bucket is unchanged. if _, err := vkv.Put("367", []byte(`{"number":367}`)); err == nil { t.Error("the view put a key") } if err := vkv.Delete("365"); err == nil { t.Error("the view deleted a key") } if err := view.Publish("mesh.mod.mesh-issues.event.opened", []byte(`{"number":367}`)); err != nil { t.Fatal(err) } _ = view.Flush() violations := map[string]bool{} deadline := time.After(10 * time.Second) for len(violations) < 3 { select { case v := <-refused: for _, subject := range []string{"$KV." + ViewBucket + ".367", "$KV." + ViewBucket + ".365", "mesh.mod.mesh-issues.event.opened"} { if strings.Contains(v, subject) { violations[subject] = true } } case <-deadline: t.Fatalf("the server refused %d of the view's 3 writes as permission violations", len(violations)) } } if _, err := tkv.Get("367"); !errors.Is(err, nats.ErrKeyNotFound) { t.Errorf("after the view's put, 367 is %v", err) } if _, err := tkv.Get("365"); err != nil { t.Errorf("after the view's delete, 365 is gone: %v", err) } }