diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index d61e53d..095086f 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -4,6 +4,7 @@ import ( "context" "crypto/sha256" "encoding/hex" + "encoding/json" "errors" "flag" "fmt" @@ -690,6 +691,39 @@ func sendTo(ctx context.Context, open *stores, names []string) error { } fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) } + // And every assignment on those machines its membership (novox/hq ADR 0160): composed from the + // same records the bus's accounts are, so what a runtime serves and what its account may are one + // composition. Issued after the declaration, because the runtime it is for arrives with it. + return issueMemberships(ctx, open, server, names) +} + +// issueMemberships publishes the membership of every module on the named machines. +func issueMemberships(ctx context.Context, open *stores, server *link.Server, names []string) error { + records, err := open.inventory.BusRecords(ctx) + if err != nil { + return err + } + where := broker.PlacementsOf(records, records.Interchangeable) + bus, ok := server.Bus().(link.OverNATS) + if !ok { + return nil + } + issued := 0 + for _, node := range names { + for _, d := range records.Assigned[node] { + body, err := json.Marshal(broker.MembershipFor(node, d, where)) + if err != nil { + return err + } + if err := bus.PublishMembership(ctx, node, d.Module, body); err != nil { + return err + } + issued++ + } + } + if issued > 0 { + fmt.Printf(" issued %d membership(s)\n", issued) + } return nil } diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 4fb76b6..8b4a26e 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -144,6 +144,7 @@ func (j *JetStream) EnsureStream(s Stream) error { MaxMsgsPerSubject: int64(s.MaxMsgsPerSubject), Description: s.Why, } + want.AllowDirect = s.Direct if s.Retention == RetentionLastPerSubject { // Last-per-subject is a limits stream with one message kept per subject, not a // retention policy of its own — the state shape, spelled the way the server spells it. diff --git a/internal/broker/membership.go b/internal/broker/membership.go new file mode 100644 index 0000000..dbd99c5 --- /dev/null +++ b/internal/broker/membership.go @@ -0,0 +1,131 @@ +package broker + +import ( + "sort" + "strings" +) + +// What the mesh issues an assignment to serve and to reach (novox/hq ADR 0160). +// +// A module's code names its tools and its events; **where they land is the mesh's to decide**, and +// it decided it twice — once in the runtime, once here, by one rule compiled into both. Now the +// controller composes a membership for every module on every machine and publishes it to a subject +// only that assignment reads; the runtime serves exactly what the membership says, and the account's +// grant is the same composition read the other way. The shape issued today is the shape the mesh +// already had, so nothing moves when a membership first arrives; only who decides it moves. + +// Membership is one assignment's subjects: what this instance of a module on this machine serves, +// and what it may reach. +type Membership struct { + Node string `json:"node"` + Module string `json:"module"` + // Serves is every address a tool of this instance answers on. `{tool}` stands for the tool's + // own name, which the module knows and the mesh does not need to: the mesh issues the address, + // the runtime fills the name. An address with a queue is shared with the module's other + // instances, and the bus hands each call to one of them; an address without is this instance's. + Serves []Served `json:"serves"` + // Seats is every verb of a seat this instance holds, at the subject the seat's callers use. + Seats []SeatServed `json:"seats,omitempty"` + // Emits is where an event of this module lands; `{event}` stands for the event's name. + Emits string `json:"emits"` + // Reaches is each tool this module may call, `.`, to the subjects that reach it: + // the first is whichever instance answers, when the mesh issued one; the rest name a machine. + Reaches map[string][]string `json:"reaches,omitempty"` + // Tools is where this instance answers what it serves — the runtime's one verb of its own. + Tools string `json:"tools"` +} + +// Served is one address a tool is answered on. +type Served struct { + Subject string `json:"subject"` + Queue string `json:"queue,omitempty"` +} + +// SeatServed is one verb of a held seat, where its callers ask. +type SeatServed struct { + Seat string `json:"seat"` + Verb string `json:"verb"` + Subject string `json:"subject"` +} + +// MembershipSubject is the one address a runtime derives for itself: where its own membership is +// published, from the two names its credential carries. Everything else is in the membership. +func MembershipSubject(node, module string) string { + return "mesh.assignment." + node + "." + module +} + +// Placements is where every module runs, for deciding which instance answers for the module. +type Placements struct { + // Nodes is each module's machines. + Nodes map[string][]string + // Interchangeable is each module whose definition says its instances are the same anywhere, + // so the module's plain subject is issued to all of them in one queue. + Interchangeable map[string]bool +} + +// AnswersForTheModule says whether an instance of a module on one machine is issued the module's +// plain subject: when it is the only instance, or when the definition says instances are +// interchangeable. A stateful module on two machines gets only its machines' subjects, so a call +// that names none reaches nothing rather than the wrong store. +func (p Placements) AnswersForTheModule(module string) bool { + return len(p.Nodes[module]) <= 1 || p.Interchangeable[module] +} + +// MembershipFor composes one assignment's membership from what it declared and where everything +// runs. The subjects are the ones PermissionsFor grants, derived here once more only until the +// grant itself is read from the membership — which is the next step, not this one. +func MembershipFor(node string, d Declared, where Placements) Membership { + own := "mesh.mod." + d.Module + m := Membership{ + Node: node, Module: d.Module, + Emits: own + ".event.{event}", + Tools: own + ".tool.tools", + } + // This machine's address always; the module's when this instance answers for the module. + m.Serves = append(m.Serves, Served{Subject: own + ".tool.{tool}." + node}) + if where.AnswersForTheModule(d.Module) { + m.Serves = append(m.Serves, Served{Subject: own + ".tool.{tool}", Queue: "serve." + d.Module}) + } + for _, s := range d.Holds { + for _, verb := range s.Serves { + m.Seats = append(m.Seats, SeatServed{Seat: s.Name, Verb: verb, Subject: seatToolSubject(s, verb, node)}) + } + } + if len(d.Invokes) > 0 { + m.Reaches = map[string][]string{} + for _, t := range d.Invokes { + if t == "*" || strings.HasPrefix(t, "seat:") { + continue // every tool, or a role's: addressed by name, not resolved per instance + } + module, tool, ok := strings.Cut(t, ".") + if !ok { + continue + } + var reach []string + if where.AnswersForTheModule(module) { + reach = append(reach, "mesh.mod."+module+".tool."+tool) + } + nodes := append([]string{}, where.Nodes[module]...) + sort.Strings(nodes) + for _, n := range nodes { + reach = append(reach, "mesh.mod."+module+".tool."+tool+"."+n) + } + m.Reaches[t] = reach + } + } + return m +} + +// PlacementsOf reads where everything runs from the records the bus's accounts are composed from. +func PlacementsOf(r Records, interchangeable map[string]bool) Placements { + p := Placements{Nodes: map[string][]string{}, Interchangeable: interchangeable} + for node, declared := range r.Assigned { + for _, d := range declared { + p.Nodes[d.Module] = append(p.Nodes[d.Module], node) + } + } + for _, nodes := range p.Nodes { + sort.Strings(nodes) + } + return p +} diff --git a/internal/broker/membership_test.go b/internal/broker/membership_test.go new file mode 100644 index 0000000..055438c --- /dev/null +++ b/internal/broker/membership_test.go @@ -0,0 +1,75 @@ +package broker + +import ( + "reflect" + "testing" +) + +// The mesh issues an assignment's subjects (novox/hq ADR 0160): a module alone on one machine +// answers for the module and for its machine; a stateful module on two machines answers only for +// each machine; one that says its instances are interchangeable answers for the module everywhere; +// a holder serves its seat's verbs; and what a module may reach is resolved the same way. +func TestAMembershipIsIssuedFromWhereEverythingRuns(t *testing.T) { + records := Records{Assigned: map[string][]Declared{ + "anchor": { + {Module: "postgres", Serves: []string{"query"}, Holds: []Seat{{Name: "mesh-store", Scope: "mesh", Serves: []string{"databases", "query"}}}}, + {Module: "catalog", Invokes: []string{"postgres.query", "search.find"}}, + }, + "home-server": { + {Module: "postgres"}, + {Module: "search"}, + {Module: "dashboard", Invokes: []string{"postgres.query"}}, + }, + "laptop": {{Module: "search"}}, + }, Interchangeable: map[string]bool{"search": true}} + where := PlacementsOf(records, records.Interchangeable) + + pg := MembershipFor("anchor", records.Assigned["anchor"][0], where) + if !reflect.DeepEqual(pg.Serves, []Served{{Subject: "mesh.mod.postgres.tool.{tool}.anchor"}}) { + t.Fatalf("a stateful module on two machines answers only for its machine: %+v", pg.Serves) + } + if len(pg.Seats) != 2 || pg.Seats[0].Subject != "mesh.seat.mesh-store.tool.databases" { + t.Fatalf("the holder serves the seat's verbs at the seat's subjects: %+v", pg.Seats) + } + if pg.Emits != "mesh.mod.postgres.event.{event}" || pg.Tools != "mesh.mod.postgres.tool.tools" { + t.Fatalf("events and the tools verb: %+v", pg) + } + + search := MembershipFor("laptop", records.Assigned["laptop"][0], where) + if !reflect.DeepEqual(search.Serves, []Served{ + {Subject: "mesh.mod.search.tool.{tool}.laptop"}, + {Subject: "mesh.mod.search.tool.{tool}", Queue: "serve.search"}, + }) { + t.Fatalf("an interchangeable module answers for the module in the queue too: %+v", search.Serves) + } + + dashboard := MembershipFor("home-server", records.Assigned["home-server"][2], where) + if !reflect.DeepEqual(dashboard.Serves, []Served{ + {Subject: "mesh.mod.dashboard.tool.{tool}.home-server"}, + {Subject: "mesh.mod.dashboard.tool.{tool}", Queue: "serve.dashboard"}, + }) { + t.Fatalf("a module alone on one machine answers for the module: %+v", dashboard.Serves) + } + if !reflect.DeepEqual(dashboard.Reaches["postgres.query"], + []string{"mesh.mod.postgres.tool.query.anchor", "mesh.mod.postgres.tool.query.home-server"}) { + t.Fatalf("reaching a stateful module names each machine and no plain subject: %v", dashboard.Reaches) + } + catalog := MembershipFor("anchor", records.Assigned["anchor"][1], where) + if !reflect.DeepEqual(catalog.Reaches["search.find"], + []string{"mesh.mod.search.tool.find", "mesh.mod.search.tool.find.home-server", "mesh.mod.search.tool.find.laptop"}) { + t.Fatalf("reaching an interchangeable module offers the plain subject first: %v", catalog.Reaches) + } + if MembershipSubject("anchor", "postgres") != "mesh.assignment.anchor.postgres" { + t.Fatal("the one subject a runtime derives for itself") + } +} + +func TestAnAccountMayReadItsOwnMembershipAndNoOthers(t *testing.T) { + perms, err := PermissionsFor(Principal{Kind: KindModule, Node: "anchor", Module: "postgres", PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + has(t, perms.Subscribe, "mesh.assignment.anchor.postgres") + has(t, perms.Publish, "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.anchor.postgres") + hasNot(t, perms.Subscribe, "mesh.assignment.>") +} diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 9032888..17bca40 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -298,6 +298,10 @@ func PermissionsFor(p Principal) (Permissions, error) { // away — no other principal may subscribe this namespace, and a caller's authority is // still granted per tool, by name, on the publish side. sub = append(sub, own+".tool.>") + // Its own membership (ADR 0160): the one subject a runtime derives for itself, read + // directly from the stream and followed live. Nothing else's. + sub = append(sub, MembershipSubject(p.Node, p.Module)) + pub = append(pub, "$JS.API.DIRECT.GET."+AssignmentsStream+"."+MembershipSubject(p.Node, p.Module)) // 1b. The tools it calls, if its manifest says it calls any (novox/hq ADR 0152). The same // grant a person gets and derived the same way, so "what may this module ask" is diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 8b588d0..fd9f7f9 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -47,8 +47,14 @@ type Stream struct { // Why is carried into the assertion so an operator reading the server's own state finds the // reason there, rather than only in a repository they may not have. Why string + // Direct lets a client read a subject's last message without a consumer, which is how a + // runtime reads its own membership with no JetStream API beyond one request (ADR 0160). + Direct bool } +// AssignmentsStream holds every assignment's membership, the newest per subject. +const AssignmentsStream = "ASSIGNMENTS" + // MeshStreams is the foundation set, in the order a person reads it. // // **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live @@ -78,6 +84,14 @@ func MeshStreams() []Stream { Why: "one declaration per node, always the newest; a node that sees sequence n refuses " + "n-1 by construction (issue 107)", }, + { + Name: AssignmentsStream, + Subjects: []string{"mesh.assignment.*.*"}, + Retention: RetentionLastPerSubject, + Direct: true, + Why: "one membership per assignment, always the newest: what the mesh issued this module " + + "on this machine to serve and to reach (ADR 0160); read directly by the runtime it is for", + }, { Name: EventsStream, // A seat's own events ride here too: they are 1:many like any event, and the diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index ae94143..9a2592e 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -123,9 +123,10 @@ func subjectMatches(filter, subject string) bool { // Each relationship's retention is the thing that makes it what it is (design 29 §4). func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) { want := map[string]Retention{ - "CONTROL": RetentionWorkQueue, - "NODES": RetentionLastPerSubject, - "EVENTS": RetentionLimits, + "CONTROL": RetentionWorkQueue, + "NODES": RetentionLastPerSubject, + "EVENTS": RetentionLimits, + "ASSIGNMENTS": RetentionLastPerSubject, } got := map[string]Retention{} for _, s := range MeshStreams() { diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index e55b370..10f86b9 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -37,18 +37,18 @@ accounts { subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] } } } { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { - publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } - subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } + publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } + subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { - publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit"] } - subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] } + publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.audit"] } + subscribe: { allow: ["_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: { - publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] } - subscribe: { allow: ["_INBOX.two.shop.>", "mesh.mod.shop.tool.>"] } + publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] } + subscribe: { allow: ["_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] } allow_responses: { max: 1, ttl: "1m" } } } ] diff --git a/internal/broker/users.go b/internal/broker/users.go index de181e4..8c82dc7 100644 --- a/internal/broker/users.go +++ b/internal/broker/users.go @@ -47,6 +47,9 @@ type Records struct { Enrolling []string // People is each person's name against the tools they may invoke, `*` for an administrator. People map[string][]string + // Interchangeable is each module whose definition says its instances are the same anywhere + // (ADR 0160), which decides whether the module's plain subject is issued to every instance. + Interchangeable map[string]bool } // Users is every user the composed file should contain, in the order it will be written. diff --git a/internal/catalogue/manifest.go b/internal/catalogue/manifest.go index 60e926f..f75fed6 100644 --- a/internal/catalogue/manifest.go +++ b/internal/catalogue/manifest.go @@ -293,6 +293,12 @@ type Manifest struct { // module claiming a seat answers what that seat's protocol promises (novox/hq ADR 0118). Tools []string `json:"tools,omitempty"` + // Instances says whether this module's instances are the same anywhere — `interchangeable` — + // so a call that names no machine may be answered by any of them (novox/hq ADR 0160). A fact + // about the software, not about the bus: a stateless web tool says it; a database does not, + // and its instances are then each addressed by machine, never confused for one another. + Instances string `json:"instances,omitempty"` + // Invokes are the tools this module calls, each `.` or a role's `seat:.`, // or the single entry `*` for every tool on the mesh (novox/hq ADR 0152, ADR 0154). // @@ -1173,6 +1179,11 @@ func ParseManifest(raw []byte) (Manifest, error) { m.Module, r)) } } + if m.Instances != "" && m.Instances != InstancesInterchangeable { + problems = append(problems, fmt.Sprintf( + "%s says its instances are %q; the one word is %q, for a module that is the same on every machine", + m.Module, m.Instances, InstancesInterchangeable)) + } for _, offer := range m.Provides { p := offer.Name if !name.MatchString(p) { @@ -1960,3 +1971,7 @@ func (o OwnSecrets) Paths() map[string]string { } return out } + +// InstancesInterchangeable is the one value of a definition's `instances`: the module is the same +// on every machine, so any instance may answer for the module. +const InstancesInterchangeable = "interchangeable" diff --git a/internal/inventory/busrecords.go b/internal/inventory/busrecords.go index 3efdc65..3fa241b 100644 --- a/internal/inventory/busrecords.go +++ b/internal/inventory/busrecords.go @@ -49,7 +49,8 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) { } } - out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{}} + out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{}, + Interchangeable: map[string]bool{}} for _, n := range nodes { out.Nodes = append(out.Nodes, n.Name) modules, err := i.Assigned(ctx, n.Name) @@ -72,6 +73,9 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) { "be derived", module, n.Name) } out.Assigned[n.Name] = append(out.Assigned[n.Name], declaredFor(m, seats)) + if m.Instances == catalogue.InstancesInterchangeable { + out.Interchangeable[m.Module] = true + } } } diff --git a/internal/inventory/people_test.go b/internal/inventory/people_test.go index 1085d57..a3a4ac6 100644 --- a/internal/inventory/people_test.go +++ b/internal/inventory/people_test.go @@ -42,8 +42,11 @@ func TestAPersonMayCallToolsAndNothingElse(t *testing.T) { if err != nil { t.Fatal(err) } - if len(perms.Publish) != 1 || perms.Publish[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" { - t.Errorf("ada may publish %v, which should be the one tool and nothing else", perms.Publish) + // The one tool, both ways it is addressed (novox/hq ADR 0159): to whichever instance + // answers, and to the instance on one machine. Nothing else. + if len(perms.Publish) != 2 || perms.Publish[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" || + perms.Publish[1] != "mesh.mod.mesh-catalog.tool.catalog_tools.*" { + t.Errorf("ada may publish %v, which should be the one tool, both ways addressed, and nothing else", perms.Publish) } for _, s := range perms.Publish { if strings.HasPrefix(s, "mesh.control") || strings.HasPrefix(s, "mesh.node") || diff --git a/internal/link/bus.go b/internal/link/bus.go index 5640c79..3e83ce7 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "github.com/novox/mesh-controller/internal/broker" "time" "github.com/nats-io/nats.go" @@ -145,6 +146,16 @@ func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body return nil } +// PublishMembership issues one assignment what it serves and reaches (novox/hq ADR 0160), last per +// subject, so the runtime that connects later reads the current one and one that is running follows. +func (b OverNATS) PublishMembership(ctx context.Context, node, module string, body []byte) error { + _, err := b.JS.Publish(broker.MembershipSubject(node, module), body, nats.Context(ctx)) + if err != nil { + return fmt.Errorf("issuing %s on %s its membership: %w", module, node, err) + } + return nil +} + func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error { _, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx)) if err != nil {