From 8312f5bdee5cfbba2e6419bf15ff0c22f4a3de6b Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 9 Oct 2026 11:03:53 +0200 Subject: [PATCH] Compose and raise the bus of the lab's proof of the operator's answers, as this controller would (hq ADR 0259) mesh-lab's asks proof runs the router, the Telegram channel and an asker on a real bus. Its accounts, streams, workers, buckets and memberships come from this test at the controller's commit, so the lab proves the composition and not a copy of it. Skipped unless the lab asks. --- internal/inventory/asks_lab_test.go | 243 ++++++++++++++++++++++++++++ 1 file changed, 243 insertions(+) create mode 100644 internal/inventory/asks_lab_test.go diff --git a/internal/inventory/asks_lab_test.go b/internal/inventory/asks_lab_test.go new file mode 100644 index 00000000..62593d54 --- /dev/null +++ b/internal/inventory/asks_lab_test.go @@ -0,0 +1,243 @@ +package inventory + +// The bus of the lab's proof of the operator's answers (mesh-lab `asks/`, novox/hq ADR 0259). +// +// The proof runs the router, the Telegram channel and an asker against a real bus, and the bus must be the +// one this controller would compose — not a copy of its rules written again in the lab, which would prove +// the copy. So the lab asks this test, at the controller's commit, for both halves: +// +// 1. **Composed** (MESH_LAB_ASKS_OUT and MESH_LAB_ASKS_CATALOGUE set): one machine, `anchor`, running the +// router (messenger), the Telegram channel, the desk channel and the machine's runtime as the catalogue +// declares them, beside two modules of the lab's own — `lab-asker`, which uses `operator-channel`, and +// `lab-bystander`, which does not. Written to the directory: the accounts block exactly as Users and +// ComposeAccounts make it, each user's credential, and every membership as MembershipFor makes it. +// 2. **Raised** (MESH_LAB_ASKS_BUS set as well): on the lab's running bus, as the controller, what a send +// asserts — the mesh's streams and consumers, the seats' work queues and workers, the modules' buckets — +// and every membership published where the runtime reads it. +// +// Without those words it skips: the controller's own suite has nothing to raise. + +import ( + "encoding/json" + "os" + "path/filepath" + "sort" + "testing" + "time" + + "github.com/nats-io/nats.go" + "golang.org/x/crypto/bcrypt" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/catalogue" +) + +// labMachine is the one machine of the lab's bus. +const labMachine = "anchor" + +// The lab's own modules: one that asks, one that may not. +var labManifests = []string{ + `{"module": "lab-asker", "version": "1", "uses": ["operator-channel"], "state": ["acted"], + "own-secrets": {"broker": "${dir:state}/broker"}, + "resources": [{"id": "state", "type": "directory", "mode": "0700", "place": "."}]}`, + `{"module": "lab-bystander", "version": "1", "own-secrets": {"broker": "${dir:state}/broker"}, + "resources": [{"id": "state", "type": "directory", "mode": "0700", "place": "."}]}`, +} + +// labCredential is what a lab process connects as: the runtime's credential shape (mesh-tools bus.Credential). +type labCredential struct { + URL string `json:"url"` + Node string `json:"node,omitempty"` + Module string `json:"module,omitempty"` + User string `json:"user"` + Password string `json:"password"` +} + +func TestTheAsksLabBus(t *testing.T) { + out, modules := os.Getenv("MESH_LAB_ASKS_OUT"), os.Getenv("MESH_LAB_ASKS_CATALOGUE") + if out == "" || modules == "" { + t.Skip("the lab did not ask for its bus (MESH_LAB_ASKS_OUT, MESH_LAB_ASKS_CATALOGUE)") + } + var manifests []catalogue.Manifest + read := func(raw []byte, from string) { + m, err := catalogue.ParseManifest(raw) + if err != nil { + t.Fatalf("%s: %v", from, err) + } + manifests = append(manifests, m) + } + for _, name := range []string{"messenger", "telegram", "desk-channel"} { + path := filepath.Join(modules, name, "module.json") + raw, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + read(raw, path) + } + if path := os.Getenv("MESH_LAB_ASKS_RUNTIME"); path != "" { + raw, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + read(raw, path) + } + for i, raw := range labManifests { + read([]byte(raw), "the lab's module "+string(rune('1'+i))) + } + + // As BusRecords reads the store: every seat any module declares, the mesh's own beside them. + seats := map[string]catalogue.SeatDeclaration{} + declarers := map[string]string{} + for _, m := range manifests { + for _, s := range m.DefinesSeats { + seats[s.Name], declarers[s.Name] = s, m.Module + } + } + for _, own := range catalogue.SeatsWithAProtocol() { + seats[own.Name] = catalogue.SeatDeclaration{Name: own.Name, Scope: own.Scope, Accepts: own.Accepts, + Emits: own.Emits, Serves: own.Serves} + } + records := broker.Records{Nodes: []string{labMachine}, Assigned: map[string][]broker.Declared{}, + People: map[string][]string{}, Interchangeable: map[string]bool{}} + var buckets []broker.Bucket + var trafficSeats []broker.Seat + for _, m := range manifests { + records.Assigned[labMachine] = append(records.Assigned[labMachine], declaredFor(m, seats, declarers)) + buckets = append(buckets, bucketsOf(m)...) + for _, s := range m.DefinesSeats { + if seat := asSeat(s, m.Module); seat.Kinded || len(seat.ByCaller) > 0 { + trafficSeats = append(trafficSeats, seat) + } + } + } + users, err := broker.Users(records) + if err != nil { + t.Fatal(err) + } + + // Each user a password of the lab's, the hash in the composition. + passwords := map[string]string{} + if raw, err := os.ReadFile(filepath.Join(out, "passwords.json")); err == nil { + _ = json.Unmarshal(raw, &passwords) + } + for i, u := range users { + name := u.Username() + if passwords[name] == "" { + passwords[name] = "lab-" + name + "-" + time.Now().Format("150405.000000") + } + hash, err := bcrypt.GenerateFromPassword([]byte(passwords[name]), bcrypt.MinCost) + if err != nil { + t.Fatal(err) + } + users[i].PasswordHash = string(hash) + } + + bus := os.Getenv("MESH_LAB_ASKS_BUS") + if bus == "" { + accounts, err := broker.ComposeAccounts(users) + if err != nil { + t.Fatal(err) + } + creds := map[string]labCredential{} + for _, u := range users { + creds[u.Username()] = labCredential{Node: u.Node, Module: u.Module, User: u.Username(), + Password: passwords[u.Username()]} + } + where := broker.PlacementsOf(records, records.Interchangeable) + memberships := map[string]broker.Membership{} + for _, d := range records.Assigned[labMachine] { + memberships[d.Module] = broker.MembershipFor(labMachine, d, where) + } + write(t, filepath.Join(out, "accounts.conf"), []byte(accounts)) + writeJSON(t, filepath.Join(out, "passwords.json"), passwords) + writeJSON(t, filepath.Join(out, "credentials.json"), creds) + writeJSON(t, filepath.Join(out, "memberships.json"), memberships) + return + } + + // Raised on the lab's bus, as the controller, as a send asserts it (cmd/mesh-controller busobjects.go). + js, err := broker.Dial(bus, nats.UserInfo("controller", passwords["controller"]), nats.CustomInboxPrefix("_INBOX.controller")) + if err != nil { + t.Fatalf("the lab's bus, as the controller: %v", err) + } + defer js.Close() + if err := broker.Raise(js, records.Nodes); err != nil { + t.Fatal(err) + } + holders := map[string]broker.Holder{} + for _, d := range records.Assigned[labMachine] { + for _, s := range d.Holds { + if _, taken := holders[s.Name]; !taken { + holders[s.Name] = broker.Holder{Node: labMachine, Module: d.Module} + } + } + } + if err := broker.RaiseSeats(js, MeshSeats(), holders); err != nil { + t.Fatal(err) + } + streams, workers := broker.SeatTrafficObjects(users) + have := map[string]bool{} + for _, s := range streams { + have[s.Name] = true + } + for _, s := range broker.TrafficQueues(trafficSeats) { + if !have[s.Name] { + streams, have[s.Name] = append(streams, s), true + } + } + for _, s := range streams { + if err := js.EnsureStream(s); err != nil { + t.Fatalf("the work queue %s: %v", s.Name, err) + } + } + for _, c := range workers { + if err := js.EnsureConsumer(c); err != nil { + t.Fatalf("the worker %s: %v", c.Name, err) + } + } + for _, c := range broker.ConsumersOf(users) { + if err := js.EnsureConsumer(c.Consumer); err != nil { + t.Fatalf("how %s hears what it consumes: %v", c.Module, err) + } + } + if _, err := broker.RaiseBuckets(js, buckets); err != nil { + t.Fatal(err) + } + if err := js.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + where := broker.PlacementsOf(records, records.Interchangeable) + names := make([]string, 0) + for _, d := range records.Assigned[labMachine] { + body, err := json.Marshal(broker.MembershipFor(labMachine, d, where)) + if err != nil { + t.Fatal(err) + } + if _, err := js.Context().Publish(broker.MembershipSubject(labMachine, d.Module), body); err != nil { + t.Fatalf("issuing %s its membership: %v", d.Module, err) + } + names = append(names, d.Module) + } + sort.Strings(names) + t.Logf("raised on %s: %d streams of seats, %d workers, %d buckets, memberships for %v", bus, len(streams), + len(workers), len(buckets), names) +} + +func write(t *testing.T, path string, body []byte) { + t.Helper() + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(path, body, 0o600); err != nil { + t.Fatal(err) + } +} + +func writeJSON(t *testing.T, path string, v any) { + t.Helper() + body, err := json.MarshalIndent(v, "", " ") + if err != nil { + t.Fatal(err) + } + write(t, path, body) +}