Files
mesh-controller/internal/inventory/busrecords.go
T
jschoubben cb77f35a27 A module may watch a role's events, and the catch-up turns out to be unnecessary
Moving the build outcome onto its role broke the one module that consumes it, and my own
agreement check passed anyway. The catalogue's subscription derived
`mesh.mod.mesh-build-machine.event.built` — a module namespace for a role's event, which no
such module owns — so it started, connected, and its graph stayed empty. The check compared
names, and the names agreed: the build machine does emit `built`. Only the subjects
disagreed, and a subscription that matches nothing is silence.

A consumed name is a module's event unless it names a role, and this package cannot tell by
looking — so whoever resolved the declaration says which, the way it already does for a seat
held or used. A module that watches a role gets the role's event subject and a consumer
filtered on it; watching grants subscribe and nothing else, because hearing what a role
announced is not taking part in it.

The check now compares the two halves that actually have to match — the subject a consumer
subscribes against the subject an emitter publishes — with a case pinning that it catches
this exact confusion. Comparing names was checking the easy half.

**And that answered the open question about catch-up: there is nothing to build.** The
mechanism exists because a queue on the old bus receives only what is published after it is
bound, so everything built before the catalogue existed was announced to nobody. A stream is
a log and a consumer is a position in it: a consumer created afterwards starts at the
beginning, so the builds are simply there. Asked of a real server, since the whole decision
rested on it — three builds published with nothing listening, then a consumer created, and
all three waiting for it.
2026-09-27 17:22:30 +02:00

179 lines
6.6 KiB
Go

package inventory
import (
"context"
"fmt"
"strings"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
)
// What the bus's user list is derived from, read out of the mesh's records.
//
// The deriving itself is pure and lives in the broker package; this is the reading, and it is kept
// apart for the reason that package keeps its own types: a permission must be a function of what a
// module declared, and a query that decided anything would be a second place authority came from.
// BusRecords is every fact the composer needs about who may reach the bus.
//
// **A module's authority comes from the manifest, not from the assignment.** The assignment says
// *where* it runs; what it may say is in what it declared, so the two are read together and the
// manifest is the one that decides.
func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
nodes, err := i.Nodes(ctx)
if err != nil {
return broker.Records{}, fmt.Errorf("cannot read the mesh's machines: %w", err)
}
declared, err := i.Catalogue(ctx)
if err != nil {
return broker.Records{}, fmt.Errorf("cannot read the catalogue: %w", err)
}
// Every seat any module declares, by name, so a module's claim can be resolved to the protocol
// that seat promises. **Across the whole catalogue, not one manifest**: a seat is declared by
// one module and held by another, which is the whole reason a seat exists (ADR 0118).
seats := map[string]catalogue.SeatDeclaration{}
for _, m := range declared {
for _, s := range m.Seats {
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 {
out.Nodes = append(out.Nodes, n.Name)
modules, err := i.Assigned(ctx, n.Name)
if err != nil {
return broker.Records{}, fmt.Errorf("cannot read what %s runs: %w", n.Name, err)
}
for _, module := range modules {
m, known := declared[module]
if !known {
// Assigned and not in the catalogue. Said rather than composed with no authority:
// a user with an empty permission list is a module that starts, connects, and is
// refused by the server on its first publish — an authorisation error that says
// nothing about a missing manifest.
//
// **The catalogue refuses to forget an assigned module, so this is the second line
// and not the first.** It earns its place there anyway: relying on another
// package's invariant is how a rule ends up enforced by nothing.
return broker.Records{}, fmt.Errorf(
"%s is assigned to %s and is not in the catalogue, so what it may say cannot "+
"be derived", module, n.Name)
}
out.Assigned[n.Name] = append(out.Assigned[n.Name], declaredFor(m, seats))
}
}
enrolling, err := i.NodesWithALiveToken(ctx)
if err != nil {
return broker.Records{}, err
}
out.Enrolling = enrolling
people, err := i.People(ctx)
if err != nil {
return broker.Records{}, err
}
for _, p := range people {
out.People[p.Name] = p.Invokes
}
return out, nil
}
// 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: 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,
}
for _, c := range m.Claims {
// 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))
}
}
for _, name := range m.Uses {
if s, declaredSomewhere := seats[name]; declaredSomewhere {
d.Uses = append(d.Uses, asSeat(s))
}
}
return d
}
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.
//
// **One enrolment user per such token** (design 25 §6): the inbox an answer goes to is scoped to the
// token, because an answer carries that machine's credentials sealed to it and a shared inbox is one
// machine able to read another's.
func (i *Inventory) NodesWithALiveToken(ctx context.Context) ([]string, error) {
rows, err := i.store.Pool().Query(ctx,
`select distinct n.name
from enrolment_token t join node n on n.id = t.node
where t.redeemed is null and t.expires > now()
order by n.name`)
if err != nil {
return nil, fmt.Errorf("cannot read which machines hold a live token: %w", err)
}
defer rows.Close()
var out []string
for rows.Next() {
var name string
if err := rows.Scan(&name); err != nil {
return nil, err
}
out = append(out, name)
}
return out, rows.Err()
}