Files
mesh-controller/internal/inventory/catalogue.go
T
jschoubben 421fe73dce The mesh builds: a machine takes the work, and the catalogue shows it
A build is work, not state. Everything else the control plane sends a
node is a declaration — this is what you should be — reconciled forever.
A build happens once and is finished. Putting it in a declaration would
mean rebuilding on every reconcile, or a declaration carrying "and I
already did this", which is state about an event rather than about a
machine.

So it travels on its own queue and the answer comes back correlated. One
queue, so several build machines share the work and each request is done
exactly once — which a per-machine routing key would not give.

mesh-builder is the program a build machine runs. Not the control plane,
which must not run commands on a machine; not the host, which would then
need a container runtime and git everywhere to do something almost no
machine will ever do. It holds its own broker credential and nothing
else.

Three properties that are decisions:

- a request is acknowledged only once the answer is away, so a builder
  that dies mid-build leaves the work for another machine rather than
  losing it with nobody ever hearing why
- one build at a time. Five at once against one runtime finishes all five
  slower than it would have finished the first, and the queue is what
  shares work between machines
- a failure is a RESULT. A build that fails silently is
  indistinguishable from a builder that is not running, and those want
  different responses

And `module list` is a catalogue: what exists, at which version, built
from which commit or handed over by hand or shipped with the control
plane, whether it is behind its source, and which machines run it. All of
that was recorded from the first build and none of it was shown, so "is
this current?" could only be answered by reading the database.

Proven against a real broker, registry and store: the mesh asked, a
builder consumed, built, published, answered; the manifest was recorded
with its commit; the source moved and the catalogue said "behind";
rebuilding caught it up with a new digest because the content changed.
2026-08-30 03:46:02 +02:00

578 lines
19 KiB
Go

package inventory
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"github.com/jackc/pgx/v5"
"github.com/novox/mesh-control/internal/catalogue"
)
// ErrNoSuchModule is what the mesh says about a module it has never been told about.
var ErrNoSuchModule = errors.New("no module of that name")
// ErrStillAssigned is why a module cannot be forgotten.
//
// Its own error because it is not a fault: it means a machine is running that module now, and
// removing the record would leave the mesh unable to describe what is on it.
var ErrStillAssigned = errors.New("that module is still assigned to nodes")
// Source is where a module comes from and what has been built from it.
type Source struct {
Repository string
Ref string
// BuiltFrom is the commit the manifest the mesh holds was read at.
BuiltFrom string
// Head is the newest commit the source is known to have.
Head string
}
// Current reports whether what the mesh holds is what the source last had.
//
// A module with no source is always current: it was handed over directly, and there is nothing
// it could be behind. Saying "out of date" about it would be inventing a comparison.
func (s Source) Current() bool {
if s.Repository == "" || s.Head == "" {
return true
}
return s.BuiltFrom == s.Head
}
// RegisterModule records a module, replacing what was there.
//
// Replacing rather than refusing, because a manifest changing is the ordinary case -- a module
// gains a requirement, a claim, a resource. What matters is that the change is visible the next
// time a node is resolved, which it is.
func (i *Inventory) RegisterModule(ctx context.Context, m catalogue.Manifest, from Source) error {
raw, err := json.Marshal(m)
if err != nil {
return err
}
// A module registered without provenance keeps whatever it had. Handing over a manifest by
// hand is a legitimate way to fix something in a hurry, and it should not silently erase the
// record of where the module normally comes from — which is the only thing that would say,
// afterwards, that the machine is running something nobody can rebuild.
_, err = i.store.Pool().Exec(ctx,
`insert into module (name, manifest, version, source, ref, built_from, source_head)
values ($1, $2, nullif($3,''), nullif($4,''), nullif($5,''), nullif($6,''), nullif($6,''))
on conflict (name) do update set
manifest = excluded.manifest,
version = excluded.version,
registered = now(),
source = coalesce(excluded.source, module.source),
ref = coalesce(excluded.ref, module.ref),
built_from = coalesce(excluded.built_from, module.built_from),
source_head = coalesce(excluded.built_from, module.source_head)`,
m.Module, raw, m.Version, from.Repository, from.Ref, from.BuiltFrom)
return err
}
// SourceMoved records that a module's source has a newer commit than the mesh has built.
//
// This is the whole of noticing. Nothing here builds anything — it writes down that the two
// halves differ, which is what makes *is this current?* answerable without building, and what
// makes a module that nobody rebuilt visible rather than silent.
func (i *Inventory) SourceMoved(ctx context.Context, module, head string) error {
tag, err := i.store.Pool().Exec(ctx,
`update module set source_head = $2, source_seen = now() where name = $1`, module, head)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("%w: %s", ErrNoSuchModule, module)
}
return nil
}
// SourceOf is where a module came from and whether the mesh is behind it.
func (i *Inventory) SourceOf(ctx context.Context, module string) (Source, error) {
var s Source
var repo, ref, built, head *string
err := i.store.Pool().QueryRow(ctx,
`select source, ref, built_from, source_head from module where name = $1`,
module).Scan(&repo, &ref, &built, &head)
if errors.Is(err, pgx.ErrNoRows) {
return Source{}, fmt.Errorf("%w: %s", ErrNoSuchModule, module)
}
if err != nil {
return Source{}, err
}
for _, pair := range []struct {
from *string
to *string
}{{repo, &s.Repository}, {ref, &s.Ref}, {built, &s.BuiltFrom}, {head, &s.Head}} {
if pair.from != nil {
*pair.to = *pair.from
}
}
return s, nil
}
// Behind is every module the mesh has not built from what its source now has, with the nodes
// running the old one.
//
// The nodes are the point. "Is this module out of date" is a fact about the catalogue; "which
// machines are running last week's version" is the question somebody actually has, and it is the
// one novox/hq ADR 0010 names as the thing that must not be lost.
func (i *Inventory) Behind(ctx context.Context) (map[string][]string, error) {
rows, err := i.store.Pool().Query(ctx,
`select m.name, coalesce(n.name, '')
from module m
left join assignment a on a.module = m.name
left join node n on n.id = a.node
where m.source is not null
and m.source_head is not null
and coalesce(m.built_from, '') is distinct from m.source_head
order by m.name, n.name`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string][]string{}
for rows.Next() {
var module, node string
if err := rows.Scan(&module, &node); err != nil {
return nil, err
}
if _, seen := out[module]; !seen {
out[module] = nil
}
if node != "" {
out[module] = append(out[module], node)
}
}
return out, rows.Err()
}
// Catalogue is every module the mesh knows about, which is what resolution needs: the question
// "how many modules provide this" cannot be asked of a subset.
func (i *Inventory) Catalogue(ctx context.Context) (map[string]catalogue.Manifest, error) {
rows, err := i.store.Pool().Query(ctx, `select manifest from module order by name`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]catalogue.Manifest{}
for rows.Next() {
var raw []byte
if err := rows.Scan(&raw); err != nil {
return nil, err
}
var m catalogue.Manifest
if err := json.Unmarshal(raw, &m); err != nil {
return nil, err
}
out[m.Module] = m
}
return out, rows.Err()
}
// Provide records a module the control plane ships with itself.
//
// The mesh's own private network is one: what computes its files is in this binary, so the
// manifest is too. Marked so it can be told apart from a module somebody wrote — not because it
// behaves differently, but because "where did this come from" must have an answer for everything
// in the catalogue, and "it came with the control plane" is that answer.
func (i *Inventory) Provide(ctx context.Context, m catalogue.Manifest) error {
raw, err := json.Marshal(m)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into module (name, manifest, version, source)
values ($1, $2, $3, 'the control plane')
on conflict (name) do update set manifest = excluded.manifest,
version = excluded.version, source = 'the control plane'`,
m.Module, raw, m.Version)
return err
}
// Provided reports whether a module came with the control plane rather than from a repository.
func (i *Inventory) Provided(ctx context.Context, name string) (bool, error) {
var source *string
err := i.store.Pool().QueryRow(ctx,
`select source from module where name = $1`, name).Scan(&source)
if err != nil {
return false, err
}
return source != nil && *source == "the control plane", nil
}
// ForgetModule removes a module, unless a machine is running it, and never one the control plane
// provides.
func (i *Inventory) ForgetModule(ctx context.Context, name string) error {
provided, err := i.Provided(ctx, name)
if err != nil {
return err
}
if provided {
// Refused rather than removed and re-created on the next migrate, which would look like
// it worked and quietly come back. Taking it off a machine is what "I do not want this"
// means, and that is what unassign is for.
return fmt.Errorf(
"%s comes with the control plane and cannot be forgotten; unassign it instead", name)
}
var on []string
rows, err := i.store.Pool().Query(ctx,
`select n.name from assignment a join node n on n.id = a.node where a.module = $1
order by n.name`, name)
if err != nil {
return err
}
for rows.Next() {
var node string
if err := rows.Scan(&node); err != nil {
rows.Close()
return err
}
on = append(on, node)
}
rows.Close()
if len(on) > 0 {
return fmt.Errorf("%w: %s. Unassign it first", ErrStillAssigned, strings.Join(on, ", "))
}
tag, err := i.store.Pool().Exec(ctx, `delete from module where name = $1`, name)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("%w: %s", ErrNoSuchModule, name)
}
return nil
}
// Assign puts a module on a node.
//
// Records the intention and checks nothing. Whether the set of assignments can actually become a
// declaration is resolution's question, asked over the whole set at once — and asking it here,
// one module at a time, would let an assignment look accepted and then refuse when a second
// arrives.
func (i *Inventory) Assign(ctx context.Context, nodeName, module string) error {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into assignment (node, module) values ($1, $2) on conflict do nothing`,
node.ID, module)
if err != nil && strings.Contains(err.Error(), "assignment_module_fkey") {
return fmt.Errorf("%w: %s", ErrNoSuchModule, module)
}
return err
}
// Unassign takes a module off a node.
func (i *Inventory) Unassign(ctx context.Context, nodeName, module string) error {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
tag, err := i.store.Pool().Exec(ctx,
`delete from assignment where node = $1 and module = $2`, node.ID, module)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("%s is not assigned to %s", module, nodeName)
}
return nil
}
// Assigned is what a person put on this node, which is not the same as what it runs: resolution
// adds whatever those modules require.
func (i *Inventory) Assigned(ctx context.Context, nodeName string) ([]string, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
`select module from assignment where node = $1 order by module`, node.ID)
if err != nil {
return nil, err
}
defer rows.Close()
var out []string
for rows.Next() {
var m string
if err := rows.Scan(&m); err != nil {
return nil, err
}
out = append(out, m)
}
return out, rows.Err()
}
// ProfileOf is what a node last said it can do, as resolution needs it: the capabilities that are
// present, and nothing else.
func (i *Inventory) ProfileOf(ctx context.Context, nodeName string) (map[string]bool, error) {
var raw []byte
err := i.store.Pool().QueryRow(ctx,
`select profile from node where name = $1`, nodeName).Scan(&raw)
if errors.Is(err, pgx.ErrNoRows) {
return nil, fmt.Errorf("%w: %s", ErrNoSuchNode, nodeName)
}
if err != nil {
return nil, err
}
out := map[string]bool{}
if len(raw) == 0 {
// A node that has never reported. Not an error, and not an empty machine either — every
// capability will read as absent, so anything requiring one is refused with "the wrong
// machine", which is wrong but visible. Better than assuming it can do everything.
return out, nil
}
var reported struct {
Capabilities []struct {
Name string `json:"name"`
Present bool `json:"present"`
} `json:"capabilities"`
}
if err := json.Unmarshal(raw, &reported); err != nil {
return nil, err
}
for _, c := range reported.Capabilities {
if c.Present {
out[c.Name] = true
}
}
return out, nil
}
// SetSettings records what somebody wants a module's configuration to say.
//
// An empty node name means the whole mesh. Replacing rather than merging what is already there:
// this is a statement of the whole layer, so removing a key is done by leaving it out, which is
// the only way removing one could work at all.
func (i *Inventory) SetSettings(ctx context.Context, nodeName, module string, values map[string]any) error {
raw, err := json.Marshal(values)
if err != nil {
return err
}
if nodeName == "" {
_, err = i.store.Pool().Exec(ctx,
`insert into settings (node, module, values) values (null, $1, $2)
on conflict (module) where node is null
do update set values = excluded.values, set_at = now()`, module, raw)
return wrapModule(err, module)
}
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into settings (node, module, values) values ($1, $2, $3)
on conflict (node, module) where node is not null
do update set values = excluded.values, set_at = now()`, node.ID, module, raw)
return wrapModule(err, module)
}
func wrapModule(err error, module string) error {
if err != nil && strings.Contains(err.Error(), "settings_module_fkey") {
return fmt.Errorf("%w: %s", ErrNoSuchModule, module)
}
return err
}
// ClearSettings removes a layer.
func (i *Inventory) ClearSettings(ctx context.Context, nodeName, module string) error {
if nodeName == "" {
_, err := i.store.Pool().Exec(ctx,
`delete from settings where module = $1 and node is null`, module)
return err
}
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`delete from settings where module = $1 and node = $2`, module, node.ID)
return err
}
// SettingsFor is the layers that apply to one module on one node, in the order they are applied.
//
// The mesh's first, then the node's, so a node that differs is expressed by differing rather
// than by restating everything the rest of the mesh already says.
func (i *Inventory) SettingsFor(ctx context.Context, nodeName, module string) ([]catalogue.Layer, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
`select node is null, values from settings
where module = $1 and (node is null or node = $2)
order by node is null desc`, module, node.ID)
if err != nil {
return nil, err
}
defer rows.Close()
var layers []catalogue.Layer
for rows.Next() {
var meshWide bool
var raw []byte
if err := rows.Scan(&meshWide, &raw); err != nil {
return nil, err
}
values := map[string]any{}
if err := json.Unmarshal(raw, &values); err != nil {
return nil, err
}
from := nodeName
if meshWide {
from = "the mesh"
}
layers = append(layers, catalogue.Layer{From: from, Values: values})
}
return layers, rows.Err()
}
// PinProvision records which node a machine gets a provision from.
//
// Only needed when more than one node could answer. Recordable before that, because a mesh with
// one database should not change where an existing machine gets its data the day a second
// arrives.
func (i *Inventory) PinProvision(ctx context.Context, nodeName, provision, provider string) error {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
from, err := i.NodeByName(ctx, provider)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into provision_pin (node, name, provider) values ($1, $2, $3)
on conflict (node, name) do update set provider = excluded.provider, pinned_at = now()`,
node.ID, provision, from.ID)
return err
}
// UnpinProvision removes a choice, putting the question back.
func (i *Inventory) UnpinProvision(ctx context.Context, nodeName, provision string) error {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
tag, err := i.store.Pool().Exec(ctx,
`delete from provision_pin where node = $1 and name = $2`, node.ID, provision)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("%s was not told where to get %q from", nodeName, provision)
}
return nil
}
// PinsFor is what a node was told about where its provisions come from.
func (i *Inventory) PinsFor(ctx context.Context, nodeName string) (map[string]string, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
`select p.name, n.name from provision_pin p join node n on n.id = p.provider
where p.node = $1`, node.ID)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]string{}
for rows.Next() {
var name, provider string
if err := rows.Scan(&name, &provider); err != nil {
return nil, err
}
out[name] = provider
}
return out, rows.Err()
}
// pinRows is how many pins a node holds, counted in the table rather than through the join that
// reads them.
//
// For a test that would otherwise pass for the wrong reason: PinsFor joins on the provider, so a
// pin left behind by a departed node is invisible through it whether it was cleaned up or not.
func (i *Inventory) pinRows(ctx context.Context, nodeName string) (int, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return 0, err
}
var n int
err = i.store.Pool().QueryRow(ctx,
`select count(*) from provision_pin where node = $1`, node.ID).Scan(&n)
return n, err
}
// Entry is one module as a catalogue shows it: what it is, where it came from, and who runs it.
type Entry struct {
Manifest catalogue.Manifest
Source Source
// On is every node this module is assigned to, sorted.
On []string
// Provided is true when the module came with the control plane rather than from a repository.
Provided bool
}
// Catalogued is every module the mesh knows about, with everything a person asks about one.
//
// **One query rather than a call per module.** A catalogue that costs a round trip per row is a
// catalogue nobody lists, and the questions here — what is this, where did it come from, who is
// running it — are asked together every time.
func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) {
rows, err := i.store.Pool().Query(ctx,
`select m.name, m.manifest,
coalesce(m.source, ''), coalesce(m.ref, ''),
coalesce(m.built_from, ''), coalesce(m.source_head, ''),
coalesce(array_agg(n.name order by n.name) filter (where n.name is not null), '{}')
from module m
left join assignment a on a.module = m.name
left join node n on n.id = a.node
group by m.name, m.manifest, m.source, m.ref, m.built_from, m.source_head
order by m.name`)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Entry
for rows.Next() {
var raw []byte
var name string
var source Source
var on []string
if err := rows.Scan(&name, &raw, &source.Repository, &source.Ref,
&source.BuiltFrom, &source.Head, &on); err != nil {
return nil, err
}
var m catalogue.Manifest
if err := json.Unmarshal(raw, &m); err != nil {
return nil, err
}
entry := Entry{Manifest: m, Source: source, On: on}
if source.Repository == providedBy {
// It came with the control plane. Not a repository, and showing it as one would have
// somebody go looking for it.
entry.Provided = true
entry.Source = Source{}
}
out = append(out, entry)
}
return out, rows.Err()
}
// providedBy is what the source column says for a module the control plane ships.
const providedBy = "the control plane"