The bus on NATS: both transports behind seams, and the rollout switch #87
@@ -4,6 +4,7 @@ import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
@@ -122,3 +123,114 @@ func theCataloguesEvents(t *testing.T) ([]AnEmitter, []AConsumer, []DeclaredSeat
|
||||
}
|
||||
return emitters, consumers, seats
|
||||
}
|
||||
|
||||
// **Do the derived subjects meet, not just the names?**
|
||||
//
|
||||
// The check above compares what a consumer asks for against what an emitter says it emits, by name. It
|
||||
// passed while the catalogue's subscription pointed at `mesh.mod.mesh-build-machine.event.built` — a
|
||||
// module namespace for a role's event, which no emitter owns. The names agreed; the subjects did not,
|
||||
// and the graph stayed empty.
|
||||
//
|
||||
// So this compares the thing that actually has to match: the subject a consumer subscribes against the
|
||||
// subject an emitter publishes. It is the last place the two halves can be held together, because
|
||||
// after this the server is the only thing that knows and it says nothing — a subscription that matches
|
||||
// nothing is silence.
|
||||
func TestTheCataloguesDerivedSubjectsMeet(t *testing.T) {
|
||||
emitters, consumers, seats := theCataloguesEvents(t)
|
||||
|
||||
// Every subject something publishes: a module's own events, and the events of every role.
|
||||
published := map[string]bool{}
|
||||
for _, e := range emitters {
|
||||
for _, name := range e.Emits {
|
||||
published["mesh.mod."+e.Module+".event."+name] = true
|
||||
}
|
||||
}
|
||||
for _, s := range seats {
|
||||
for _, name := range s.Emits {
|
||||
published["mesh.seat."+s.Name+".event."+name] = true
|
||||
}
|
||||
}
|
||||
|
||||
byName := map[string]DeclaredSeat{}
|
||||
for _, s := range seats {
|
||||
byName[s.Name] = s
|
||||
}
|
||||
|
||||
var lonely []string
|
||||
for _, c := range consumers {
|
||||
principal := Principal{Kind: KindModule, Node: "one", Module: c.Module, PasswordHash: "x"}
|
||||
for _, want := range c.Consumes {
|
||||
emitter, event, named := strings.Cut(want, ".")
|
||||
if named {
|
||||
if s, isASeat := byName[emitter]; isASeat {
|
||||
principal.Watches = append(principal.Watches,
|
||||
Seat{Name: s.Name, Emits: []string{event}})
|
||||
continue
|
||||
}
|
||||
}
|
||||
principal.Consumes = append(principal.Consumes, want)
|
||||
}
|
||||
perms, err := PermissionsFor(principal)
|
||||
if err != nil {
|
||||
t.Fatalf("%s: %v", c.Module, err)
|
||||
}
|
||||
for _, subject := range perms.Subscribe {
|
||||
if !strings.Contains(subject, ".event.") {
|
||||
continue
|
||||
}
|
||||
if reaches(subject, published) {
|
||||
continue
|
||||
}
|
||||
// A wildcard over emitters reaches whatever arrives later, and an emitter that is not
|
||||
// installed is ordinary — both are already excused by the check above, so only a subject
|
||||
// that can never match anything gets here.
|
||||
if strings.Contains(subject, "*") || strings.Contains(subject, ">") {
|
||||
continue
|
||||
}
|
||||
lonely = append(lonely, c.Module+" subscribes "+subject+", which nothing publishes")
|
||||
}
|
||||
}
|
||||
if len(lonely) > 0 {
|
||||
sort.Strings(lonely)
|
||||
t.Fatalf("%d subscription(s) derive to a subject no emitter owns:\n %s",
|
||||
len(lonely), strings.Join(lonely, "\n "))
|
||||
}
|
||||
}
|
||||
|
||||
// And it catches the thing it exists for: a role's event read as a module's.
|
||||
func TestTheDerivedSubjectCheckCatchesARolesEventReadAsAModules(t *testing.T) {
|
||||
published := map[string]bool{"mesh.seat.mesh-build-machine.event.built": true}
|
||||
// What the derivation produced before a consumed seat name was resolved as one.
|
||||
if reaches("mesh.mod.mesh-build-machine.event.built", published) {
|
||||
t.Fatal("a module namespace was treated as reaching a role's event, which is the bug")
|
||||
}
|
||||
// And the corrected one does reach it.
|
||||
if !reaches("mesh.seat.mesh-build-machine.event.built", published) {
|
||||
t.Fatal("the role's own subject does not reach the role's event")
|
||||
}
|
||||
}
|
||||
|
||||
// reaches says whether a subscribed subject admits any published one.
|
||||
func reaches(subject string, published map[string]bool) bool {
|
||||
for p := range published {
|
||||
if admitsSubject(strings.Split(subject, "."), strings.Split(p, ".")) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func admitsSubject(pattern, subject []string) bool {
|
||||
for i, token := range pattern {
|
||||
if token == ">" {
|
||||
return i < len(subject)
|
||||
}
|
||||
if i >= len(subject) {
|
||||
return false
|
||||
}
|
||||
if token != "*" && token != subject[i] {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return len(pattern) == len(subject)
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package broker
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// Streams and consumers derived from what modules declare.
|
||||
@@ -103,16 +104,22 @@ type DeclaredSeat struct {
|
||||
// from its name: a module with a consumer per event would need an ack permission per consumer,
|
||||
// and the permission list would stop being derivable from the declaration.
|
||||
func ConsumerFor(p Principal) (Consumer, bool) {
|
||||
if p.Kind != KindModule || len(p.Consumes) == 0 {
|
||||
// A module that reacts to anything — a module's events or a role's (novox/hq ADR 0121). Watching
|
||||
// a role was missing here, so the one module that does it got no consumer at all: it started,
|
||||
// connected, and its graph stayed empty with nothing anywhere reporting why.
|
||||
if p.Kind != KindModule || (len(p.Consumes) == 0 && len(p.Watches) == 0) {
|
||||
return Consumer{}, false
|
||||
}
|
||||
perms, err := PermissionsFor(p)
|
||||
if err != nil {
|
||||
return Consumer{}, false
|
||||
}
|
||||
// Events, wherever they live: a module's own namespace, and the namespace of any role it watches
|
||||
// (novox/hq ADR 0121). Tool subjects and inboxes are subscribed directly and are not a consumer's
|
||||
// business, which is why this is a filter and not the whole list.
|
||||
var filters []string
|
||||
for _, s := range perms.Subscribe {
|
||||
if len(s) > 9 && s[:9] == "mesh.mod." {
|
||||
if strings.Contains(s, ".event.") {
|
||||
filters = append(filters, s)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -60,6 +60,15 @@ type Principal struct {
|
||||
|
||||
Holds []Seat
|
||||
Uses []Seat
|
||||
// Watches are seats whose events this principal consumes. Separate from Consumes because a
|
||||
// role's event lives under the seat's namespace and not a module's, and this package cannot tell
|
||||
// a seat's name from a module's by looking at it — whoever resolved the declaration can, and
|
||||
// does (novox/hq ADR 0121).
|
||||
//
|
||||
// **Found by a consumer reading nothing.** The catalogue consumes the build machine's outcome;
|
||||
// with that name read as a module's, its subscription pointed at `mesh.mod.mesh-build-machine.…`,
|
||||
// a namespace no such module owns. Every service started and the graph stayed empty.
|
||||
Watches []Seat
|
||||
|
||||
// Invokes are the tools a person may call, as `<module>.<tool>`; a single `*` is every tool,
|
||||
// for an administrator. Only meaningful for KindPerson.
|
||||
@@ -260,6 +269,14 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
sub = append(sub, subject)
|
||||
}
|
||||
|
||||
// 2b. Events of a role it watches, under the seat's own namespace. Subscribe only: watching a
|
||||
// role is hearing what it announced, not taking part in it.
|
||||
for _, w := range p.Watches {
|
||||
for _, e := range w.Emits {
|
||||
sub = append(sub, seatSubject(w, "event", e))
|
||||
}
|
||||
}
|
||||
|
||||
// 3. Seats it holds: full participation.
|
||||
for _, s := range p.Holds {
|
||||
for _, a := range s.Accepts {
|
||||
|
||||
@@ -143,3 +143,54 @@ func TestRaisingAMeshRolesWorkQueue(t *testing.T) {
|
||||
t.Fatalf("the holder got no worker on the role's queue: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// **A consumer created after the fact still sees what came before it**, which is why the mesh needs no
|
||||
// catch-up at all on this bus (novox/hq 04-ISSUES/050).
|
||||
//
|
||||
// On the bus the mesh runs on today a queue receives only what is published after it is bound, so
|
||||
// everything built before the catalogue existed was announced to nobody — and on a fresh mesh that is
|
||||
// always the foundation, because those are the things the catalogue needed in order to exist. A whole
|
||||
// mechanism was built for it: the catalogue asks, the controller re-publishes.
|
||||
//
|
||||
// A stream is a log and a consumer is a position in it. A consumer created later starts at the
|
||||
// beginning by default, so the builds are simply there. Asked of a real server rather than assumed,
|
||||
// because the whole decision about whether to keep that mechanism rests on it.
|
||||
func TestAConsumerCreatedAfterwardsStillSeesWhatCameBefore(t *testing.T) {
|
||||
js := aLiveBus(t)
|
||||
if err := AssertMeshStreams(js); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := js.Context().PurgeStream("EVENTS"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Genesis: things are built before anything is listening.
|
||||
built := []string{"base", "store", "mesh-catalog"}
|
||||
for _, m := range built {
|
||||
if _, err := js.Context().Publish("mesh.seat.mesh-build-machine.event.built",
|
||||
[]byte(`{"module":"`+m+`"}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
// Now the catalogue is installed and the controller creates its consumer.
|
||||
c, ok := ConsumerFor(Principal{Kind: KindModule, Node: "one", Module: "mesh-catalog",
|
||||
Watches: []Seat{{Name: "mesh-build-machine", Emits: []string{"built"}}}, PasswordHash: "x"})
|
||||
if !ok {
|
||||
t.Fatal("a module that watches a role got no consumer")
|
||||
}
|
||||
t.Cleanup(func() { _ = js.Context().DeleteConsumer(c.Stream, c.Name) })
|
||||
if err := js.EnsureConsumer(c); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
info, err := js.Context().ConsumerInfo(c.Stream, c.Name)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if info.NumPending != uint64(len(built)) {
|
||||
t.Fatalf("a consumer created after %d builds has %d waiting for it — if this is 0 the mesh "+
|
||||
"does need a catch-up after all, and the reasoning for deleting it is wrong",
|
||||
len(built), info.NumPending)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -29,6 +29,8 @@ type Declared struct {
|
||||
Holds []Seat
|
||||
// Uses are the seats this module sends to.
|
||||
Uses []Seat
|
||||
// Watches are the seats whose events it consumes.
|
||||
Watches []Seat
|
||||
}
|
||||
|
||||
// Records is what composing a user list needs to know about the mesh, and nothing more.
|
||||
@@ -59,7 +61,7 @@ func Users(r Records) ([]Principal, error) {
|
||||
out = append(out, Principal{
|
||||
Kind: KindModule, Node: node, Module: d.Module,
|
||||
Emits: d.Emits, Consumes: d.Consumes, Serves: d.Serves,
|
||||
Holds: d.Holds, Uses: d.Uses,
|
||||
Holds: d.Holds, Uses: d.Uses, Watches: d.Watches,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package inventory
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
@@ -93,10 +94,27 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
|
||||
// declaredFor is one module's manifest as the composer needs it: what it says about itself, and the
|
||||
// protocol of every seat it holds or uses.
|
||||
func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaration) broker.Declared {
|
||||
// A consumed name is a module's event unless it names a seat, and only somebody holding the seat
|
||||
// set can tell (novox/hq ADR 0121). Split here, because the composer cannot look at a name and
|
||||
// know — and a role's event read as a module's is a subscription to a namespace nobody owns.
|
||||
var fromModules []string
|
||||
var watches []broker.Seat
|
||||
for _, c := range m.Consumes {
|
||||
emitter, event, named := strings.Cut(c, ".")
|
||||
if named {
|
||||
if s, isASeat := seats[emitter]; isASeat {
|
||||
watches = append(watches, broker.Seat{Name: s.Name, Emits: []string{event}})
|
||||
continue
|
||||
}
|
||||
}
|
||||
fromModules = append(fromModules, c)
|
||||
}
|
||||
|
||||
d := broker.Declared{
|
||||
Module: m.Module,
|
||||
Emits: m.Emits,
|
||||
Consumes: m.Consumes,
|
||||
Consumes: fromModules,
|
||||
Watches: watches,
|
||||
// The tools it answers, which is `tools` and not `serves`: the manifest's `serves` is the
|
||||
// facts a consumer needs to reach a provision, a different meaning under a similar word.
|
||||
Serves: m.Tools,
|
||||
|
||||
Reference in New Issue
Block a user