diff --git a/cmd/mesh-controller/bus_view.go b/cmd/mesh-controller/bus_view.go index ad44797b..33dd7038 100644 --- a/cmd/mesh-controller/bus_view.go +++ b/cmd/mesh-controller/bus_view.go @@ -68,6 +68,7 @@ func busViewCredential(ctx context.Context, args []string) error { InboxPrefix string `json:"inbox_prefix"` Hears []string `json:"hears"` Reads string `json:"reads"` + HowToRead string `json:"how_to_read"` }{ WebSocket: websocket, URL: "nats://" + where.Address, Fingerprint: where.Fingerprint, User: broker.ViewUser, Password: password, @@ -75,6 +76,10 @@ func busViewCredential(ctx context.Context, args []string) error { // `_INBOX.view.>` and no wider (design 25 §4), and a client's default inbox is not under it. InboxPrefix: "_INBOX." + broker.ViewUser, Hears: broker.ViewHears, Reads: broker.ViewBucket, + HowToRead: "direct reads only, no watch (a consumer is refused): list with a request to $JS.API.DIRECT.GET.KV_" + + broker.ViewBucket + ` carrying {"multi_last":["$KV.` + broker.ViewBucket + `.>"]}, answered until a 204 status; ` + + "read one key with $JS.API.DIRECT.GET.KV_" + broker.ViewBucket + ".$KV." + broker.ViewBucket + ". " + + "(nats.js: kvm.open(bucket, {allow_direct: true}), never create); re-read the key an event's number names", }) if err != nil { return err @@ -83,8 +88,9 @@ func busViewCredential(ctx context.Context, args []string) error { fmt.Printf("issued the view, which hears %s and reads the bucket %s, and nothing else\n", strings.Join(broker.ViewHears, ", "), broker.ViewBucket) fmt.Println(" this is the only time the credential is printed; the mesh keeps a hash") - fmt.Println(" it works once the bus has been told, which is the next push to the machine holding mesh-broker;") - fmt.Println(" the WebSocket listener it connects through is the bus module's, live there after the bus step a person starts (`bus upgrade`)") + fmt.Println(" it works once the bus has been told, which is the next push to the machine holding mesh-broker —") + fmt.Println(" and while a new build of the bus module waits for that machine, the next `bus upgrade` a person starts,") + fmt.Println(" which is also what brings the WebSocket listener it connects through") fmt.Println() fmt.Println(string(held)) return nil diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 97556da7..89d17835 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -45,8 +45,8 @@ const ( // KindView is the one read-only principal a view onto the bus connects as (novox/hq research 036, // gap G1): a page in a browser, over the bus module's WebSocket listener, watching the issue tracker. // Fixed, and derived from no declaration: what it hears is ViewHears, what it reads is ViewBucket, - // and it publishes nothing but the JetStream API requests a read-only watcher of that one bucket - // makes (ViewReads), each answered in its own inbox. Composed like every other user, into the same + // and it publishes nothing but direct reads of that one bucket (ViewReads), each answered in its own + // inbox. Composed like every other user, into the same // file, once its credential is minted (`bus view-credential`); forgotten like every other user // (`bus view-revoke`), at the next composition. KindView Kind = "view" @@ -76,22 +76,27 @@ var ViewHears = []string{ moduleEventSubject("mesh-delivery", "group"), } -// ViewReads are the JetStream API requests a read-only watcher of ViewBucket makes, on that bucket's -// stream and no other: the same requests a module reading another's state is granted (stateGrants, -// measured against the server) — binding (STREAM.INFO), a key read directly (DIRECT.GET), an ordered -// consumer for a watch or a listing (CONSUMER.CREATE), deleting it, and answering its flow control -// ($JS.FC) — and CONSUMER.INFO, which the browser client asks of the consumer it just made before it -// hands a watch over (nats.js kv: `oc.info(true)`). Nothing here is a write: no `$KV..>`, which -// is what a put or a delete publishes to, and no STREAM.* that defines, purges or deletes. +// ViewReads are the JetStream API requests the view makes, on ViewBucket's stream and no other: binding +// (STREAM.INFO), and direct reads — one key by its subject (`DIRECT.GET..$KV..`), and +// the batch form on the bare subject, which answers the newest value of every key (`multi_last`) into the +// asker's inbox. The page lists the bucket with the batch, and re-reads one key when the tracker's event +// names it (every event carries the issue's `number`). +// +// **No consumer, deliberately, and so no watch.** A KV watch is a push consumer, and a push consumer's +// deliver subject is the creator's choice, delivered by the server's own client — which the server does +// not hold to the creator's permissions. Measured on 2.11.17 (2026-10-10): the view, granted +// CONSUMER.CREATE on this stream, made a consumer delivering to `mesh.mod.mesh-issues.event.opened`, and +// a module subscribed there received the bucket's entry as the tracker's event. A grant of CONSUMER.CREATE +// is a publish to any subject in the account; the view publishes nothing, so it has none (nor +// CONSUMER.DELETE, which would let it delete a module's consumer). Nothing here is a write either: no +// `$KV..>`, which is what a put or a delete publishes to, and no STREAM.* that defines, purges or +// deletes. func ViewReads() []string { stream := "KV_" + ViewBucket return []string{ "$JS.API.STREAM.INFO." + stream, + "$JS.API.DIRECT.GET." + stream, "$JS.API.DIRECT.GET." + stream + ".>", - "$JS.API.CONSUMER.CREATE." + stream + ".>", - "$JS.API.CONSUMER.INFO." + stream + ".>", - "$JS.API.CONSUMER.DELETE." + stream + ".>", - "$JS.FC." + stream + ".>", } } diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index ac99af7e..aa4620e8 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -57,7 +57,7 @@ accounts { allow_responses: { max: 1, ttl: "1m" } } } { user: "view", password: "$2a$11$vvvvvvvvvvvvvvvvvvvvvv", permissions: { - publish: { allow: ["$JS.API.CONSUMER.CREATE.KV_mesh-issues_issues.>", "$JS.API.CONSUMER.DELETE.KV_mesh-issues_issues.>", "$JS.API.CONSUMER.INFO.KV_mesh-issues_issues.>", "$JS.API.DIRECT.GET.KV_mesh-issues_issues.>", "$JS.API.STREAM.INFO.KV_mesh-issues_issues", "$JS.FC.KV_mesh-issues_issues.>"] } + publish: { allow: ["$JS.API.DIRECT.GET.KV_mesh-issues_issues", "$JS.API.DIRECT.GET.KV_mesh-issues_issues.>", "$JS.API.STREAM.INFO.KV_mesh-issues_issues"] } subscribe: { allow: ["_INBOX.view.>", "mesh.mod.mesh-delivery.event.group", "mesh.mod.mesh-delivery.event.transition", "mesh.mod.mesh-issues.event.linked", "mesh.mod.mesh-issues.event.moved", "mesh.mod.mesh-issues.event.noted", "mesh.mod.mesh-issues.event.opened", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.plan-moved"] } } } ] diff --git a/internal/broker/view_live_test.go b/internal/broker/view_live_test.go index c120edce..020d3b98 100644 --- a/internal/broker/view_live_test.go +++ b/internal/broker/view_live_test.go @@ -16,10 +16,11 @@ import ( // 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, -// reads a key, watches a change land — and is refused every write: a put, a delete, an event. +// 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 TestTheViewWatchesTheIssuesOverWebSocketAndWritesNothing(t *testing.T) { +func TestTheViewReadsTheIssuesOverWebSocketAndWritesNothing(t *testing.T) { hash := func(password string) string { h, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.MinCost) if err != nil { @@ -29,7 +30,7 @@ func TestTheViewWatchesTheIssuesOverWebSocketAndWritesNothing(t *testing.T) { } const node = "anchor" tracker := Principal{Kind: KindModule, Node: node, Module: "mesh-issues", State: []string{"issues"}, - Emits: []string{"opened"}, PasswordHash: hash("tracker")} + Emits: []string{"opened", "moved"}, PasswordHash: hash("tracker")} accounts, err := ComposeAccounts([]Principal{ {Kind: KindController, PasswordHash: hash("controller")}, tracker, @@ -111,30 +112,89 @@ func TestTheViewWatchesTheIssuesOverWebSocketAndWritesNothing(t *testing.T) { } else if !strings.Contains(string(got.Value()), `"number":365`) { t.Fatalf("the view read %q", got.Value()) } - watch, err := vkv.WatchAll() - if err != nil { - t.Fatalf("the view could not watch the bucket: %v", err) - } - t.Cleanup(func() { _ = watch.Stop() }) - seen := func(key string) { - t.Helper() - deadline := time.After(10 * time.Second) - for { - select { - case e := <-watch.Updates(): - if e != nil && e.Key() == key { - return - } - case <-deadline: - t.Fatalf("the view's watch never saw %s", key) - } - } - } - seen("365") if _, err := tkv.Put("366", []byte(`{"number":366,"status":"open"}`)); err != nil { t.Fatal(err) } - seen("366") + + // 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 { diff --git a/internal/broker/view_test.go b/internal/broker/view_test.go index b58e04ed..e3e8e4b0 100644 --- a/internal/broker/view_test.go +++ b/internal/broker/view_test.go @@ -3,7 +3,6 @@ package broker import ( "reflect" "sort" - "strings" "testing" ) @@ -34,35 +33,38 @@ func TestTheViewHearsAndReadsAndCanPublishNothingElse(t *testing.T) { t.Error("the view may answer, and nothing is ever asked of it") } - // Every publish grant is a JetStream API request about the one bucket's stream, or its flow control. - // **The mutation this holds against**: a write grant of any shape. + // Every publish grant is binding the one bucket's stream or reading it directly. **The mutation this + // holds against**: a write grant of any shape, and a consumer of any shape — a consumer's deliver + // subject is the creator's choice, so creating one is publishing anywhere (ViewReads). stream := "KV_" + ViewBucket for _, p := range perms.Publish { - readOnly := strings.HasPrefix(p, "$JS.API.STREAM.INFO."+stream) || - strings.HasPrefix(p, "$JS.API.DIRECT.GET."+stream+".") || - strings.HasPrefix(p, "$JS.API.CONSUMER.CREATE."+stream+".") || - strings.HasPrefix(p, "$JS.API.CONSUMER.INFO."+stream+".") || - strings.HasPrefix(p, "$JS.API.CONSUMER.DELETE."+stream+".") || - strings.HasPrefix(p, "$JS.FC."+stream+".") + readOnly := p == "$JS.API.STREAM.INFO."+stream || + p == "$JS.API.DIRECT.GET."+stream || + p == "$JS.API.DIRECT.GET."+stream+".>" if !readOnly { t.Errorf("the view is granted a publish on %q, which is not a read of %s", p, ViewBucket) } } for _, refused := range []string{ - "$KV." + ViewBucket + ".365", // a put or a delete - "$KV.mesh-controller_conditions.x", // another bucket - "$JS.API.STREAM.CREATE." + stream, // defining the stream - "$JS.API.STREAM.PURGE." + stream, // emptying it - "$JS.API.STREAM.DELETE." + stream, // deleting it - "$JS.API.STREAM.MSG.DELETE." + stream, // deleting a message - "$JS.API.CONSUMER.CREATE.KV_mesh-controller_conditions.x", // reading another bucket - "$JS.API.STREAM.INFO.EVENTS", // the events stream - "$JS.API.INFO", // the account - "mesh.mod.mesh-issues.event.opened", // claiming the tracker said something - "mesh.mod.mesh-issues.tool.open", // opening an issue - "mesh.seat.issue-tracker.tool.open", // through the seat - "mesh.seat.issue-tracker.tool.open.novox", // on one machine - "mesh.seat.mesh-controller.tool.status", // the controller's verbs + "$KV." + ViewBucket + ".365", // a put or a delete + "$KV.mesh-controller_conditions.x", // another bucket + "$JS.API.STREAM.CREATE." + stream, // defining the stream + "$JS.API.STREAM.PURGE." + stream, // emptying it + "$JS.API.STREAM.DELETE." + stream, // deleting it + "$JS.API.STREAM.MSG.DELETE." + stream, // deleting a message + "$JS.API.CONSUMER.CREATE.KV_mesh-controller_conditions.x", // reading another bucket + "$JS.API.CONSUMER.CREATE." + stream + ".w.$KV." + ViewBucket + ".>", // a watch: delivers anywhere + "$JS.API.CONSUMER.CREATE." + stream, // an unnamed consumer + "$JS.API.CONSUMER.DELETE." + stream + ".a_mesh-issues", // a module's consumer + "$JS.FC." + stream + ".x", + "$JS.API.DIRECT.GET.KV_mesh-controller_conditions", // another bucket, directly + "$JS.API.STREAM.INFO.EVENTS", // the events stream + "$JS.API.INFO", // the account + "mesh.mod.mesh-issues.event.opened", // claiming the tracker said something + "mesh.mod.mesh-issues.tool.open", // opening an issue + "mesh.seat.issue-tracker.tool.open", // through the seat + "mesh.seat.issue-tracker.tool.open.novox", // on one machine + "mesh.seat.mesh-controller.tool.status", // the controller's verbs "mesh.seat.mesh-controller.event.plan-moved", "$SRV.PING", "_INBOX.controller.x",