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.
244 lines
8.4 KiB
Go
244 lines
8.4 KiB
Go
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)
|
|
}
|