diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 1cafa57..f499281 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -713,6 +713,12 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) if err := broker.Raise(js, names); err != nil { return err } + // The work queues of the mesh's own roles (novox/hq ADR 0121). The queue before the holder, + // deliberately: work queues until somebody arrives to do it, so assigning a build machine a week + // after something started asking for builds flushes the backlog instead of having lost it. + if err := broker.RaiseSeats(js, inventory.MeshSeats(), nil); err != nil { + return err + } fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n", address, len(names)) return nil diff --git a/internal/broker/derived.go b/internal/broker/derived.go index b5e15c0..b4b0cfa 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -89,8 +89,10 @@ type DeclaredSeat struct { Accepts []string // Emits are the verbs the seat's holder publishes under the seat's own name. An event about a // role belongs here rather than in the holder's namespace, because the name then outlives - // whoever fills it (novox/hq 04-ISSUES/127). - Emits []string + // whoever fills it (novox/hq ADR 0121, 04-ISSUES/127). + Emits []string + // Serves are the verbs the holder answers, request and reply. + Serves []string RetainSeconds int } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 4e0a7aa..45da6b0 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -76,6 +76,11 @@ type Principal struct { PasswordHash string } +// meshSeatsTheControllerUses are the roles the mesh's own flows submit work to. Named rather than +// derived from the seat set: the controller is not a module and declares no `uses`, so its side of a +// seat has to be stated, and a list is what makes "which roles does the mesh itself talk to" answerable. +var meshSeatsTheControllerUses = []string{"mesh-build-machine"} + // enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the // controller may answer an enrolment is derived from the same constant the user is named from. const enrolmentPrefix = "enrol" @@ -154,8 +159,16 @@ func PermissionsFor(p Principal) (Permissions, error) { case KindController: // The controller owns the mesh's own traffic and the streams. It is the only writer of // stream definitions (design 25 §3), so it alone reaches the JetStream API. - pub = []string{"mesh.control.>", "mesh.node.>", "mesh.build.>", "$JS.API.>"} - sub = []string{"mesh.control.>", "mesh.build.>", "$JS.API.>"} + pub = []string{"mesh.control.>", "mesh.node.>", "$JS.API.>"} + sub = []string{"mesh.control.>", "$JS.API.>"} + + // Work the mesh's own flows submit to a role, and the outcomes they wait on (ADR 0121). A + // build is the one today: the controller asks, and reads the answer from the seat's event + // like the catalogue does — which is why no holder needs to publish into anybody's inbox. + for _, seat := range meshSeatsTheControllerUses { + pub = append(pub, "mesh.seat."+seat+".accept.>") + sub = append(sub, "mesh.seat."+seat+".event.>") + } // The two events it reacts to, and its ack subject on the stream they arrive from // (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the diff --git a/internal/broker/raise_live_test.go b/internal/broker/raise_live_test.go index d58cf07..6281679 100644 --- a/internal/broker/raise_live_test.go +++ b/internal/broker/raise_live_test.go @@ -99,3 +99,47 @@ func TestTheControlConsumerDoesNotDeadLetterBeforeTheControllerGivesUp(t *testin "before the controller finished deciding about it", info.Config.MaxDeliver) } } + +// A role's work queue exists before anybody holds it, against a real server. +// +// **The queue before the holder is the point** (novox/hq ADR 0121): work queues until somebody arrives +// to do it, so assigning a build machine a week after something started asking for builds flushes the +// backlog instead of having lost it. A stream created at assignment would make "the holder is not here +// yet" mean "your requests are gone". +func TestRaisingAMeshRolesWorkQueue(t *testing.T) { + js := aLiveBus(t) + seats := []DeclaredSeat{{Name: "mesh-build-machine", Accepts: []string{"build"}, + Emits: []string{"built"}}} + t.Cleanup(func() { _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") }) + + if err := RaiseSeats(js, seats, nil); err != nil { + t.Fatalf("a real server refused a role's work queue: %v", err) + } + info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE") + if err != nil { + t.Fatalf("the role has no work queue: %v", err) + } + if info.Config.Retention != nats.WorkQueuePolicy { + t.Errorf("the queue retains as %v: work a holder took must leave it, or the next holder does "+ + "it again", info.Config.Retention) + } + if len(info.Config.Subjects) != 1 || info.Config.Subjects[0] != "mesh.seat.mesh-build-machine.accept.>" { + t.Errorf("it carries %v rather than the role's own inbound subjects", info.Config.Subjects) + } + // Nobody holds it, so there is no worker — and asserting again changes nothing, because this runs + // on every start. + if err := RaiseSeats(js, seats, nil); err != nil { + t.Fatalf("asserting a role's queue a second time failed, so a restart would: %v", err) + } + + // And once somebody holds it, the worker appears on that same queue. + if err := RaiseSeats(js, seats, map[string]Holder{ + "mesh-build-machine": {Node: "anchor", Module: "builder"}, + }); err != nil { + t.Fatal(err) + } + if _, err := js.Context().ConsumerInfo("SEAT_MESH_BUILD_MACHINE", + "SEAT_MESH_BUILD_MACHINE_worker"); err != nil { + t.Fatalf("the holder got no worker on the role's queue: %v", err) + } +} diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 7b0d11a..87748fb 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -63,8 +63,10 @@ type Stream struct { func MeshStreams() []Stream { return []Stream{ { - Name: "CONTROL", - Subjects: []string{"mesh.control.*.report", "mesh.control.enrol", "mesh.control.built"}, + Name: "CONTROL", + // A build's outcome is no longer here: it is the build-machine seat's own event, so one + // publish reaches whoever asked, the controller and the catalogue (novox/hq ADR 0121). + Subjects: []string{"mesh.control.*.report", "mesh.control.enrol"}, Retention: RetentionWorkQueue, Why: "the store-window guarantee (ADR 0083): the controller naks with a delay while its " + "store is away and the message is redelivered; nothing is dropped", @@ -76,12 +78,6 @@ 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: "BUILDS", - Subjects: []string{"mesh.build.request"}, - Retention: RetentionWorkQueue, - Why: "at least once, one builder at a time; a builder that dies mid-build has its message redelivered", - }, { Name: "EVENTS", // A seat's own events ride here too: they are 1:many like any event, and the @@ -167,15 +163,20 @@ const ControllerName = "controller" // ControllerFollows are the events the controller reacts to: the catalogue saying a module's // current version moved, and a catalogue that has just started saying it may have missed builds. // -// **These carry the local names the manifests hold today**, which still spell an event the way a -// routing key on the bus the mesh has does — `module..` rather than design 29's bare -// verb — so the derived subject names the module twice. It is consistent, and it is what the -// catalogue actually publishes, so it is what the controller must listen to. It changes when those -// names are converted, and not before: a subscription written against the name design 29 specifies -// would be a controller listening to a subject nothing publishes. +// **Derived the same way a module's subscription is**, from the emitter and the bare local event +// name, rather than written out. They were written out while the catalogue still spelled its events +// as the old bus's routing keys, and the moment those were converted (novox/hq 04-ISSUES/127) a +// hard-coded pair became a controller listening to a subject nothing publishes — the same fault, from +// the other side. Deriving them means the conversion could not leave these behind. var ControllerFollows = []string{ - "mesh.mod.mesh-catalog.event.module.mesh-catalog.upgraded", - "mesh.mod.mesh-catalog.event.module.mesh-catalog.catching-up", + moduleEventSubject("mesh-catalog", "upgraded"), + moduleEventSubject("mesh-catalog", "catching-up"), +} + +// moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so +// what the controller subscribes and what the emitter is permitted to publish cannot drift apart. +func moduleEventSubject(module, event string) string { + return "mesh.mod." + module + ".event." + event } // MeshConsumers is what the controller consumes, in the order a person reads it. diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index dc70bb8..ca450a3 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -125,7 +125,6 @@ func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) { want := map[string]Retention{ "CONTROL": RetentionWorkQueue, "NODES": RetentionLastPerSubject, - "BUILDS": RetentionWorkQueue, "EVENTS": RetentionLimits, } got := map[string]Retention{} diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 6bcd013..0a34b53 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -23,8 +23,8 @@ accounts { MESH { users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.build.>", "mesh.control.>", "mesh.node.>"] } - subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.build.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.module.mesh-catalog.catching-up", "mesh.mod.mesh-catalog.event.module.mesh-catalog.upgraded"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>"] } + subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.>"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index bfb3f7e..345a017 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -27,6 +27,17 @@ type Seat struct { // provision may only be held by a module providing it at the seat's scope, and its holder is // what a requirement for that provision resolves to when several modules provide it. Delivers string + // Accepts, Emits and Serves are the protocol of the role, as local verbs — the same three a + // module declares for a seat of its own (novox/hq ADR 0118), and empty for most of these: a seat + // is usually about who does a job and not about what may be said to them. + // + // **Named here so the mesh has no role it cannot describe** (ADR 0121). Without them a build + // machine had three audiences for one outcome and nothing derived a grant for any of them, and an + // event about a role had nowhere to live but the namespace of whichever module held that role + // today — which the bus refuses, because a namespace belongs to who it is named for. + Accepts []string + Emits []string + Serves []string // Decision is the record that made it a seat. Decision string } @@ -53,7 +64,12 @@ var seats = []Seat{ {Name: "mesh-catalog", Scope: ScopeMesh, Decision: "novox/hq ADR 0110"}, {Name: "mesh-npm-package-registry", Scope: ScopeMesh, Delivers: "npm-package-registry", Decision: "novox/hq ADR 0109"}, {Name: "mesh-git", Scope: ScopeMesh, Delivers: "git", Decision: "novox/hq ADR 0111"}, - {Name: "mesh-build-machine", Scope: ScopeNode, Decision: "novox/hq ADR 0110"}, + // A build is work submitted to this role, and its outcome is the role's own event (ADR 0121). + // One publish reaches whoever asked, the controller that records it, and the catalogue that + // places it in the module graph — which is what the old bus's shared exchange did for free, and + // what a dedicated build branch was doing a second way. + {Name: "mesh-build-machine", Scope: ScopeNode, + Accepts: []string{"build"}, Emits: []string{"built"}, Decision: "novox/hq ADR 0110"}, {Name: "mesh-dns-port", Scope: ScopeNode, Decision: "novox/hq ADR 0110"}, {Name: "mesh-intrusion-prevention", Scope: ScopeNode, Decision: "novox/hq ADR 0110"}, {Name: "mesh-packet-filter", Scope: ScopeNode, Decision: "novox/hq ADR 0110"}, @@ -197,3 +213,18 @@ var renamedSeats = map[string]string{ "the-resolver-configuration": "mesh-resolver-configuration", "the-showcase": "mesh-showcase", } + +// SeatsWithAProtocol are the mesh's own seats that say something about what may be said to them or by +// them, which is the set the bus derives streams, consumers and permissions from. +// +// Most of the set is not here, and that is the ordinary case: a seat saying only who does a job grants +// nothing on the bus and needs no queue. +func SeatsWithAProtocol() []Seat { + var out []Seat + for _, s := range seats { + if len(s.Accepts) > 0 || len(s.Emits) > 0 || len(s.Serves) > 0 { + out = append(out, s) + } + } + return out +} diff --git a/internal/inventory/busrecords.go b/internal/inventory/busrecords.go index a9a6ae1..955394f 100644 --- a/internal/inventory/busrecords.go +++ b/internal/inventory/busrecords.go @@ -38,6 +38,15 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) { seats[s.Name] = s } } + // And the mesh's own, which carry protocol too (novox/hq ADR 0121). Added after the modules' + // rather than before, because a `mesh-*` name is the mesh's and registration refuses a module + // declaring one — so this cannot be shadowed, and if it ever were, the mesh's own would win. + 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, + } + } out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{}} for _, n := range nodes { @@ -88,9 +97,9 @@ func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaratio Serves: m.Tools, } for _, c := range m.Claims { - // A seat the mesh defines for itself carries no protocol, so holding one grants nothing here: - // those seats say who does a job, not who may say what. - if s, declaredSomewhere := seats[c.Name]; declaredSomewhere { + // Every seat with a protocol, the mesh's own included. One that says only who does a job is + // not here and grants nothing, which is most of them. + if s, hasAProtocol := seats[c.Name]; hasAProtocol { d.Holds = append(d.Holds, asSeat(s)) } } @@ -106,6 +115,18 @@ func asSeat(s catalogue.SeatDeclaration) broker.Seat { return broker.Seat{Name: s.Name, Accepts: s.Accepts, Emits: s.Emits, Serves: s.Serves} } +// MeshSeats are the mesh's own seats that carry a protocol, as the bus needs them: what to make a work +// queue for, and whose holder gets a worker on it (novox/hq ADR 0121). +func MeshSeats() []broker.DeclaredSeat { + var out []broker.DeclaredSeat + for _, s := range catalogue.SeatsWithAProtocol() { + out = append(out, broker.DeclaredSeat{ + Name: s.Name, Accepts: s.Accepts, Emits: s.Emits, Serves: s.Serves, + }) + } + return out +} + // NodesWithALiveToken is every machine holding a token that could still be presented — issued, not // expired, not redeemed. //