211 lines
8.8 KiB
Go
211 lines
8.8 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/conditions"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
)
|
|
|
|
// The streams and durable consumers the mesh's own traffic needs, derived once (novox/hq to-be 45 §4,
|
|
// D6 and D7).
|
|
//
|
|
// **What the controller asserts at its start and what the self-check expects to find are one
|
|
// derivation**, run against the bus to make them and against a recorder to list them. Two lists would
|
|
// drift, and a self-check comparing the bus with a second opinion of what should be there would find
|
|
// the drift rather than the fault.
|
|
|
|
// assertBusObjects brings every stream and consumer into being on r, and answers the machines that
|
|
// can now hear a declaration.
|
|
func assertBusObjects(ctx context.Context, inv *inventory.Inventory, r broker.Raiser) ([]string, error) {
|
|
nodes, err := inv.Nodes(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
names := make([]string, 0, len(nodes))
|
|
for _, n := range nodes {
|
|
names = append(names, n.Name)
|
|
}
|
|
if err := broker.Raise(r, names); err != nil {
|
|
return nil, 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.
|
|
// With the seats' holders, so each role's work queue gets the consumer its holder takes
|
|
// work from. Passed as nil until the first live raise, which left the build machine bound to a
|
|
// consumer nothing had created (2026-09-28).
|
|
holders, err := seatHolders(ctx, inv)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := broker.RaiseSeats(r, inventory.MeshSeats(), holders); err != nil {
|
|
return nil, err
|
|
}
|
|
// And how every module hears what it consumes. Derived from the same records the user list is
|
|
// composed from, so a module the mesh grants a consumer's subjects has that consumer waiting.
|
|
// Done on every raise, not only when a credential is issued: every module moved onto this bus
|
|
// by the rollout was issued on the old one, and came up with nothing to bind to (2026-09-28).
|
|
consumers, err := moduleConsumers(ctx, inv)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
// And the work queues of seats that name their caller or their kind, with each holder's worker
|
|
// (novox/hq ADR 0259 §3): an ask queues until the router takes it, a channel's work until that kind
|
|
// takes it.
|
|
trafficStreams, trafficWorkers, err := seatTrafficObjects(ctx, inv)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
// Every one tried, and every failure named: one module's consumer the bus refuses is no reason
|
|
// the modules after it in the list hear nothing (novox/hq issue 208, where this runs on each send).
|
|
var failed []error
|
|
for _, s := range trafficStreams {
|
|
if err := r.EnsureStream(s); err != nil {
|
|
failed = append(failed, fmt.Errorf("the work queue %s: %w", s.Name, err))
|
|
}
|
|
}
|
|
for _, c := range trafficWorkers {
|
|
if err := r.EnsureConsumer(c); err != nil {
|
|
failed = append(failed, fmt.Errorf("the worker %s on %s: %w", c.Name, c.Stream, err))
|
|
}
|
|
}
|
|
for _, c := range consumers {
|
|
if err := r.EnsureConsumer(c.Consumer); err != nil {
|
|
failed = append(failed, fmt.Errorf("how %s on %s hears what it consumes: %w", c.Module, c.Node, err))
|
|
}
|
|
}
|
|
if len(failed) > 0 {
|
|
return nil, errors.Join(failed...)
|
|
}
|
|
return names, nil
|
|
}
|
|
|
|
// What raises the condition a send says when the objects it implies could not be asserted, and its kind.
|
|
const (
|
|
sourceBusObjects = "bus-objects"
|
|
kindBusObjectsUnasserted = "bus-objects-unasserted"
|
|
)
|
|
|
|
// assertOnSend asserts, on the bus a send is about to use, every object assertBusObjects derives —
|
|
// **whenever a declaration is sent, not only when the controller starts** (novox/hq issue 208).
|
|
//
|
|
// A module assigned after the controller started was sent its declaration and found no consumer to
|
|
// bind (`consumer not found`, messenger on 2026-10-06), and a seat holder assigned after it found no
|
|
// worker: the objects a declaration implies were asserted at start and nowhere else, so they existed
|
|
// only for what was assigned before the last restart. The same derivation, not a second list of what
|
|
// a send needs: what start asserts, the self-check expects and a send asserts are one answer. Every
|
|
// part is idempotent, so asserting the whole of it again is the no-op a restart already relies on.
|
|
//
|
|
// **A failure is said and raised, and the send goes on.** The objects are the mesh's, not the
|
|
// machines' being sent: holding every machine back for one consumer that none of them may use would
|
|
// turn one fault into all of them, and the declarations are not what is wrong. It is never silent —
|
|
// said in the send's own output and raised as a condition, which the next send that asserts them
|
|
// clears — and the start-time raise still refuses to serve without them.
|
|
func assertOnSend(ctx context.Context, inv *inventory.Inventory, r broker.Raiser, indent string) error {
|
|
_, err := assertBusObjects(ctx, inv, r)
|
|
var observed []conditions.Observation
|
|
if err != nil {
|
|
fmt.Printf("%sTHE BUS DOES NOT HOLD WHAT THIS SEND IMPLIES: %v\n", indent, err)
|
|
fmt.Printf("%s a module may find no consumer to bind, or a holder no worker; sent anyway, raised as "+
|
|
"condition %s, and asserted again by the next send\n", indent, unassertedObservation(err).Key())
|
|
observed = append(observed, unassertedObservation(err))
|
|
}
|
|
// Observed when it failed, cleared when it did not: a send that asserted everything is the
|
|
// observation that the bus holds what it should.
|
|
if kerr := withKeeper(ctx, func(k *conditions.Keeper) error {
|
|
return k.Reconcile(ctx, sourceBusObjects, observed)
|
|
}); kerr != nil {
|
|
fmt.Printf("%sand whether the bus holds what this send implies could not be kept as a condition: %v\n",
|
|
indent, kerr)
|
|
}
|
|
return err
|
|
}
|
|
|
|
// unassertedObservation is a send's failure to assert the bus's objects, as a condition.
|
|
func unassertedObservation(err error) conditions.Observation {
|
|
return conditions.Observation{Scope: conditions.ScopeBus, ID: "objects", Token: "unasserted",
|
|
Kind: kindBusObjectsUnasserted, Severity: conditions.Warning, Source: sourceBusObjects,
|
|
Summary: "the bus's streams and consumers could not be asserted when a declaration was sent: " +
|
|
"a module may find no consumer to bind, or a seat's holder no worker",
|
|
Said: err.Error()}
|
|
}
|
|
|
|
// moduleConsumers is every module's durable consumer, from the records the user list is composed from.
|
|
func moduleConsumers(ctx context.Context, inv *inventory.Inventory) ([]broker.ModuleConsumer, error) {
|
|
records, err := inv.BusRecords(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
users, err := broker.Users(records)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return broker.ConsumersOf(users), nil
|
|
}
|
|
|
|
// seatTrafficObjects is the work queues and workers of seats that name their caller or their kind, from
|
|
// the records the user list is composed from.
|
|
func seatTrafficObjects(ctx context.Context, inv *inventory.Inventory) ([]broker.Stream, []broker.Consumer, error) {
|
|
records, err := inv.BusRecords(ctx)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
users, err := broker.Users(records)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
streams, workers := broker.SeatTrafficObjects(users)
|
|
// And the queue of every such seat the catalogue declares, held or not: work queues from registration,
|
|
// so what is submitted before a holder is assigned waits for it (the correctness review of 2026-10-08).
|
|
declared, err := inv.DeclaredTrafficSeats(ctx)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
have := map[string]bool{}
|
|
for _, s := range streams {
|
|
have[s.Name] = true
|
|
}
|
|
for _, s := range broker.TrafficQueues(declared) {
|
|
if !have[s.Name] {
|
|
streams = append(streams, s)
|
|
have[s.Name] = true
|
|
}
|
|
}
|
|
return streams, workers, nil
|
|
}
|
|
|
|
// moduleConsumerCount is how many modules hear what they consume, for the raise's one line.
|
|
func moduleConsumerCount(ctx context.Context, inv *inventory.Inventory) (int, error) {
|
|
consumers, err := moduleConsumers(ctx, inv)
|
|
return len(consumers), err
|
|
}
|
|
|
|
// expectedBusObjects is every stream and consumer assertBusObjects would make, made nowhere.
|
|
func expectedBusObjects(ctx context.Context, inv *inventory.Inventory) ([]broker.Stream, []broker.Consumer, error) {
|
|
var rec recordingRaiser
|
|
if _, err := assertBusObjects(ctx, inv, &rec); err != nil {
|
|
return nil, nil, fmt.Errorf("what the bus should hold cannot be worked out: %w", err)
|
|
}
|
|
return rec.streams, rec.consumers, nil
|
|
}
|
|
|
|
// recordingRaiser keeps what it was asked to assert and asserts nothing.
|
|
type recordingRaiser struct {
|
|
streams []broker.Stream
|
|
consumers []broker.Consumer
|
|
}
|
|
|
|
func (r *recordingRaiser) EnsureStream(s broker.Stream) error {
|
|
r.streams = append(r.streams, s)
|
|
return nil
|
|
}
|
|
|
|
func (r *recordingRaiser) EnsureConsumer(c broker.Consumer) error {
|
|
r.consumers = append(r.consumers, c)
|
|
return nil
|
|
}
|