Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
525f10b858 | ||
|
|
0547316cf2 | ||
|
|
9fe9b5349c | ||
|
|
d1e488efaf | ||
|
|
c5dc7e732a | ||
|
|
7efcccd013 | ||
|
|
964285f08c | ||
|
|
5698dda11f | ||
|
|
6005a8471f | ||
|
|
e2ee0dfe98 | ||
|
|
70341cfbc7 | ||
|
|
77643aa3f4 | ||
|
|
1fd6194ff8 | ||
|
|
c37018fdd2 | ||
|
|
3907ea0db0 | ||
|
|
40f5e9a41c | ||
|
|
2c2eb51878 | ||
|
|
aa2d0b51ea | ||
|
|
64d154d9d7 | ||
|
|
ffa390f916 | ||
|
|
4d62e6caf1 | ||
|
|
386ae676ca | ||
|
|
f8a9c3d6bc | ||
|
|
83671fae5f |
@@ -135,6 +135,16 @@ func run() error {
|
||||
// machine told about both would take work from one and answer on the other, and every log line would
|
||||
// say it was fine.
|
||||
func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
|
||||
// **The credential decides, before any variable does.** A machine moved to the new bus was
|
||||
// handed a credential for it and nothing else changed in its environment; that credential
|
||||
// names the bus by scheme, so it is enough to know which bus to take work from.
|
||||
if credential.onTheNewBus() {
|
||||
js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.MachineOverNATS(js, on), nil
|
||||
}
|
||||
address, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -460,10 +470,26 @@ func brokerFrom() (Credential, error) {
|
||||
// **The same shape a node gets, for the same reason** (novox/hq ADR 0004): the fingerprint travels
|
||||
// out of band — here, sealed with the credential — and the endpoint is verified once at connect.
|
||||
type Credential struct {
|
||||
URL string `json:"url"`
|
||||
// Fingerprint is SHA-256 over the broker certificate's DER bytes, or empty to verify the
|
||||
// ordinary way.
|
||||
URL string `json:"url"`
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
// User and Password ride beside the address on the bus being built (design 25): a credential
|
||||
// embedded in a URL leaks into every log line that prints a connection, so the mesh seals them
|
||||
// as two fields and this machine joins them once, here, to dial.
|
||||
User string `json:"user,omitempty"`
|
||||
Password string `json:"password,omitempty"`
|
||||
}
|
||||
|
||||
// onTheNewBus is whether a credential is for the bus being built: its address says so, and the
|
||||
// mesh only ever seals such a credential with the user and password beside it.
|
||||
func (c Credential) onTheNewBus() bool { return strings.HasPrefix(strings.TrimSpace(c.URL), "nats://") }
|
||||
|
||||
// natsURL is the address with this machine's credential in it, for the one dial that needs it.
|
||||
func (c Credential) natsURL() string {
|
||||
rest := strings.TrimPrefix(strings.TrimSpace(c.URL), "nats://")
|
||||
if c.User == "" {
|
||||
return "nats://" + rest
|
||||
}
|
||||
return "nats://" + c.User + ":" + c.Password + "@" + rest
|
||||
}
|
||||
|
||||
// dial opens the connection, pinning the broker's certificate when there is one to pin.
|
||||
|
||||
@@ -14,8 +14,8 @@ func TestBuilderDiagnosticsStayOffStdout(t *testing.T) {
|
||||
allowed := map[string]bool{
|
||||
"string(body)": true, // once.go: the result JSON, which IS stdout
|
||||
"version)": true, // --version
|
||||
`"stopping")`: true, // the loop.s shutdown line
|
||||
"usage)": true, // --help text, for a human
|
||||
`"stopping")`: true, // the loop.s shutdown line
|
||||
"usage)": true, // --help text, for a human
|
||||
}
|
||||
for _, file := range []string{"once.go", "main.go"} {
|
||||
src, err := os.ReadFile(file)
|
||||
|
||||
@@ -37,7 +37,7 @@ func askCommand(ctx context.Context, args []string) error {
|
||||
arguments = json.RawMessage(positionals[2])
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -31,6 +31,53 @@ import (
|
||||
// control plane may send a machine is bounded by the declaration language. This is the shape the
|
||||
// builder module will take when it is given work over the broker; today a person runs it, and the
|
||||
// mesh records the result the same way either way.
|
||||
// buildOn rebuilds every module the mesh holds that stands on the named module's artifacts — the
|
||||
// rebuild a changed base needs, which nothing else asks for: their sources did not move, and
|
||||
// "behind" does not see a base that did (novox/hq 04-ISSUES/131). Bases first among them too.
|
||||
func buildOn(ctx context.Context, base string, wait time.Duration) error {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
held, err := open.inventory.Catalogued(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var on []inventory.Entry
|
||||
for _, e := range held {
|
||||
if e.Manifest.Build == nil {
|
||||
continue
|
||||
}
|
||||
for _, b := range e.Manifest.Build.On {
|
||||
if b.Module == base {
|
||||
on = append(on, e)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(on) == 0 {
|
||||
fmt.Printf("nothing the mesh holds stands on %s\n", base)
|
||||
return nil
|
||||
}
|
||||
on = orderByBases(on)
|
||||
fmt.Printf("%d module(s) stand on %s:\n", len(on), base)
|
||||
var failed []string
|
||||
for _, e := range on {
|
||||
fmt.Printf("--- %s\n", e.Manifest.Module)
|
||||
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
|
||||
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, wait); err != nil {
|
||||
fmt.Printf(" %v\n", err)
|
||||
failed = append(failed, e.Manifest.Module)
|
||||
}
|
||||
}
|
||||
if len(failed) > 0 {
|
||||
return fmt.Errorf("%d of %d could not be built: %s", len(failed), len(on), strings.Join(failed, ", "))
|
||||
}
|
||||
fmt.Printf("\n%d module(s) rebuilt on %s. `push --behind` sends them on\n", len(on), base)
|
||||
return nil
|
||||
}
|
||||
|
||||
func buildCommand(ctx context.Context, args []string) error {
|
||||
set := flag.NewFlagSet("build", flag.ContinueOnError)
|
||||
ref := set.String("ref", "", "the branch, tag or commit to build")
|
||||
@@ -46,6 +93,7 @@ func buildCommand(ctx context.Context, args []string) error {
|
||||
// retype each repository is asking them to be the loop. Naming a repository and asking which
|
||||
// ones need building are different requests, so they are not combined.
|
||||
behind := set.Bool("behind", false, "every module the mesh holds older than its source has")
|
||||
on := set.String("on", "", "rebuild every module that stands on this module's artifacts — the rebuild a changed base needs")
|
||||
// A repository on the mesh's own forge, named by its path there (novox/hq ADR 0111). Without it
|
||||
// the repository is external, cloned exactly as given — see source.go.
|
||||
self := set.Bool("self", false, "the repository is a path on the forge holding the git seat")
|
||||
@@ -53,6 +101,13 @@ func buildCommand(ctx context.Context, args []string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if *on != "" {
|
||||
if len(positionals) != 0 || *behind || *self {
|
||||
return errors.New("build --on <module> names a base and nothing else")
|
||||
}
|
||||
return buildOn(ctx, *on, *wait)
|
||||
}
|
||||
|
||||
if *behind {
|
||||
if len(positionals) != 0 || *self {
|
||||
return errors.New("build <repository> or build --behind, not both: one names a " +
|
||||
@@ -61,7 +116,7 @@ func buildCommand(ctx context.Context, args []string) error {
|
||||
return buildBehind(ctx, *wait)
|
||||
}
|
||||
if len(positionals) != 1 {
|
||||
return errors.New("build <repository> [--self] [--path P] [--ref R] [--wait D] [--dry-run]")
|
||||
return errors.New("build <repository> [--self] [--path P] [--ref R] [--wait D] [--dry-run] | build --behind | build --on <module>")
|
||||
}
|
||||
source := buildSource{Repository: positionals[0]}
|
||||
if *self {
|
||||
@@ -329,6 +384,10 @@ func buildBehind(ctx context.Context, wait time.Duration) error {
|
||||
}
|
||||
fmt.Println()
|
||||
|
||||
// Bases first: a module built before the module it stands on is built against the old one
|
||||
// and reports success (novox/hq 04-ISSUES/131).
|
||||
stale = orderByBases(stale)
|
||||
|
||||
var failed []string
|
||||
for _, e := range stale {
|
||||
fmt.Printf("--- %s\n", e.Manifest.Module)
|
||||
@@ -368,7 +427,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
|
||||
}
|
||||
defer ident.Close()
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -472,7 +531,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
|
||||
return err
|
||||
}
|
||||
defer ident.Close()
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -184,6 +184,7 @@ func usage() {
|
||||
operator key show the operator key, and what it can recover
|
||||
build <repository> [--ref R] have a build machine build it, and record what came out
|
||||
build --behind build every module the mesh holds older than its source
|
||||
build --on <module> rebuild every module that stands on this module's artifacts, bases first
|
||||
builds [<module>] what has been built lately, and what came of it
|
||||
builder issue <name> a broker account for a build machine, scoped to build work,
|
||||
delivered as the builder module's broker secret (module add it first)
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
func entry(module string, on ...string) inventory.Entry {
|
||||
b := &catalogue.Build{}
|
||||
for _, o := range on {
|
||||
b.On = append(b.On, catalogue.BuildsOn{Arg: "X", Module: o, Artifact: "runtime"})
|
||||
}
|
||||
return inventory.Entry{Manifest: catalogue.Manifest{Module: module, Build: b}}
|
||||
}
|
||||
|
||||
// A module built before the module it stands on is built against the old one and reports success
|
||||
// (novox/hq 04-ISSUES/131). So bases come first, however the set arrived.
|
||||
func TestBasesAreBuiltBeforeWhatStandsOnThem(t *testing.T) {
|
||||
in := []inventory.Entry{entry("app", "runtime"), entry("runtime", "base"), entry("other"), entry("base")}
|
||||
got := orderByBases(in)
|
||||
pos := map[string]int{}
|
||||
for i, e := range got {
|
||||
pos[e.Manifest.Module] = i
|
||||
}
|
||||
if !(pos["base"] < pos["runtime"] && pos["runtime"] < pos["app"]) {
|
||||
t.Fatalf("bases not first: %v", pos)
|
||||
}
|
||||
if len(got) != 4 {
|
||||
t.Fatalf("an entry was lost or doubled: %d", len(got))
|
||||
}
|
||||
// A base outside the set is not waited for: it is not being rebuilt.
|
||||
got = orderByBases([]inventory.Entry{entry("app", "elsewhere")})
|
||||
if len(got) != 1 {
|
||||
t.Fatalf("a dependency outside the set changed the set: %v", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A merge names a repository the way the forge does; a source is recorded the way a build was
|
||||
// asked for. The two meet on owner/repo and branch, whichever form the record took.
|
||||
func TestAMergeMatchesTheSourcesBuiltFromIt(t *testing.T) {
|
||||
m := link.SourceMoved{Owner: "novox", Repo: "mesh-controller", Base: "main",
|
||||
CloneURL: "http://forge.internal:20000/novox/mesh-controller.git"}
|
||||
for _, s := range []inventory.Source{
|
||||
{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "main"},
|
||||
{Repository: "novox/mesh-controller", Seat: "git", Ref: ""},
|
||||
{Repository: "https://elsewhere.example/novox/mesh-controller", Ref: "main"},
|
||||
} {
|
||||
if !sourceIs(s, m) {
|
||||
t.Errorf("%+v was not matched by the merge", s)
|
||||
}
|
||||
}
|
||||
for _, s := range []inventory.Source{
|
||||
{Repository: "novox/mesh-host", Seat: "git"},
|
||||
{Repository: "http://forge.internal:20000/novox/mesh-controller.git", Ref: "release"},
|
||||
} {
|
||||
if sourceIs(s, m) {
|
||||
t.Errorf("%+v was matched by a merge that is not its", s)
|
||||
}
|
||||
}
|
||||
}
|
||||
+80
-16
@@ -38,6 +38,33 @@ func reportUnhostable(node string, plan catalogue.Resolution) {
|
||||
// nothing in it was wrong, and no one edit was the one that should have been a new file.
|
||||
|
||||
// serve is the control plane running: one connection to the broker, one queue, one consumer.
|
||||
// connectLink opens the controller's link over whichever bus this process is on (design 25: one
|
||||
// variable moves it). The streams and this controller's consumers are raised first on the new bus,
|
||||
// so nothing served here finds them missing.
|
||||
func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) {
|
||||
busAddress, onNATS, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !onNATS {
|
||||
return link.Connect(enroller, listener)
|
||||
}
|
||||
if inv != nil {
|
||||
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
js, err := broker.Dial(busAddress)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
||||
broker.BareAddress(busAddress), err)
|
||||
}
|
||||
return link.ConnectNats(js, enroller, listener), nil
|
||||
}
|
||||
|
||||
func serve(ctx context.Context) error {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
@@ -90,7 +117,7 @@ func serve(ctx context.Context) error {
|
||||
|
||||
work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known,
|
||||
OnNATS: onNATS}
|
||||
server, err := link.Connect(work, work)
|
||||
server, err := connectLink(ctx, inv, work, work)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -100,11 +127,6 @@ func serve(ctx context.Context) error {
|
||||
// somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was
|
||||
// replaced all have records and no objects — and a node whose consumer is missing hears nothing
|
||||
// while everything else about it looks correct.
|
||||
if onNATS {
|
||||
if err := raiseTheBus(ctx, inv, busAddress); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
// And build results nobody was waiting for. A build triggered any other way than `build`
|
||||
// would otherwise be reported into the void, which is the same as not reporting it.
|
||||
server.Records(builds{inv})
|
||||
@@ -158,13 +180,13 @@ func declare(ctx context.Context, args []string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer server.Close()
|
||||
|
||||
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil {
|
||||
if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
|
||||
@@ -271,7 +293,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
return err
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -332,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
||||
if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
|
||||
return err
|
||||
}
|
||||
// After it is away, not before. A digest recorded for something that failed to send would
|
||||
@@ -415,7 +437,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
||||
},
|
||||
func(s readyNode, body []byte) error {
|
||||
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body,
|
||||
if err := link.Declare(ctx, server.Bus(), ident, s.node, body,
|
||||
15*time.Second); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -628,7 +650,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
||||
len(refusals), strings.Join(refusals, "\n\n"))
|
||||
}
|
||||
|
||||
server, err := link.Connect(nil, nil)
|
||||
server, err := connectLink(ctx, nil, nil, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -639,7 +661,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil {
|
||||
if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
|
||||
return err
|
||||
}
|
||||
record, err := inv.NodeByName(ctx, s.node)
|
||||
@@ -704,7 +726,7 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
||||
js, err := broker.Dial(address)
|
||||
if err != nil {
|
||||
return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w",
|
||||
address, err)
|
||||
broker.BareAddress(address), err)
|
||||
}
|
||||
defer js.Close()
|
||||
|
||||
@@ -739,10 +761,52 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string)
|
||||
// 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 {
|
||||
// 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 err
|
||||
}
|
||||
if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); 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))
|
||||
broker.BareAddress(address), len(names))
|
||||
return nil
|
||||
}
|
||||
|
||||
// seatHolders is who holds each of the mesh's seats, by seat name: the record where a handover
|
||||
// wrote one, and the assigned module claiming the seat otherwise — the same derivation the
|
||||
// resolver makes, read from the catalogue rather than re-resolved.
|
||||
func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]broker.Holder, error) {
|
||||
out := map[string]broker.Holder{}
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, e := range entries {
|
||||
if len(e.On) == 0 {
|
||||
continue
|
||||
}
|
||||
for _, c := range e.Manifest.Claims {
|
||||
seat, known := catalogue.SeatNamed(c.Name)
|
||||
if !known {
|
||||
continue
|
||||
}
|
||||
if _, taken := out[seat.Name]; !taken {
|
||||
out[seat.Name] = broker.Holder{Node: e.On[0], Module: e.Manifest.Module}
|
||||
}
|
||||
}
|
||||
}
|
||||
recorded, err := inv.Holdings(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, h := range recorded {
|
||||
if seat, known := catalogue.SeatNamed(h.Claim); known {
|
||||
out[seat.Name] = broker.Holder{Node: h.Node, Module: h.Module}
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -33,7 +34,7 @@ import (
|
||||
// ability to change things, not the services its modules are serving — measured on 2026-09-27, when
|
||||
// a seat emptied mid-change and the control plane looped for two hours while every service stayed up.
|
||||
|
||||
const rolloutUsage = "rollout check | rollout mint [--again] | rollout --confirm"
|
||||
const rolloutUsage = "rollout check | rollout mint [--again] | rollout hand <node> | rollout --confirm"
|
||||
|
||||
func rolloutCommand(ctx context.Context, args []string) error {
|
||||
switch {
|
||||
@@ -41,6 +42,8 @@ func rolloutCommand(ctx context.Context, args []string) error {
|
||||
return rolloutCheck(ctx)
|
||||
case len(args) == 1 && args[0] == "mint":
|
||||
return rolloutMint(ctx, false)
|
||||
case len(args) == 2 && args[0] == "hand":
|
||||
return rolloutHand(ctx, args[1])
|
||||
case len(args) == 2 && args[0] == "mint" && args[1] == "--again":
|
||||
// Every credential minted afresh, whether or not one exists — for a mint that was wrong
|
||||
// before anything was pushed. Afterwards nothing that received the old one still works,
|
||||
@@ -340,6 +343,14 @@ func rolloutMint(ctx context.Context, again bool) error {
|
||||
machines++
|
||||
|
||||
case broker.KindModule:
|
||||
if p.Module == "mesh-controller" {
|
||||
// The control plane is a module too, and its `broker` secret is the old bus's
|
||||
// credential it is still using while this runs. Writing the new bus's blob there
|
||||
// cut the mesh off from its own old bus mid-move (2026-09-28). Its new-bus credential
|
||||
// is the controller principal's `bus` secret above; nothing else is needed here.
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
m, inShelf := shelf[p.Module]
|
||||
if !inShelf {
|
||||
skipped++
|
||||
@@ -379,3 +390,62 @@ func providesBus(m catalogue.Manifest) bool {
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// rolloutHand mints a machine its credential for the new bus afresh and prints its membership
|
||||
// once, for an operator to carry by hand — the rescue for a machine that cannot be reached over
|
||||
// any bus: rotated while it still held the old password, or reachable only by ssh. The plaintext
|
||||
// exists on this terminal and then only where it is written; the store keeps the hash, and the
|
||||
// sealed copy in the machine's declaration is replaced too, so the next push says the same.
|
||||
func rolloutHand(ctx context.Context, node string) error {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
inv := open.inventory
|
||||
|
||||
known, err := broker.FromEnvironment()
|
||||
if err != nil {
|
||||
return fmt.Errorf("the bus's certificate is not known to this process: %w", err)
|
||||
}
|
||||
busAddress, _, err := broker.OnNATS()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if busAddress == "" {
|
||||
return errors.New("this control plane is not on the new bus, so there is no membership to hand out")
|
||||
}
|
||||
_, _, bare := broker.CredentialIn(busAddress)
|
||||
if _, after, has := strings.Cut(bare, "://"); has {
|
||||
bare = after
|
||||
}
|
||||
if _, err := inv.NodeByName(ctx, node); err != nil {
|
||||
return err
|
||||
}
|
||||
p := broker.Principal{Kind: broker.KindNode, Node: node}
|
||||
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusNode, Node: node})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
membership, _ := json.Marshal(map[string]string{
|
||||
"broker": bare, "fingerprint": known.Fingerprint, "password": password, "transport": "nats",
|
||||
})
|
||||
key, err := inv.SealingKeyOf(ctx, node)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
sealed, err := secrets.Seal(key, membership)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := inv.PutBusMembership(ctx, node, sealed); err != nil {
|
||||
return err
|
||||
}
|
||||
// The one line of output is the membership itself, so it can be piped to the machine without
|
||||
// being read on the way. Everything else goes to stderr.
|
||||
fmt.Fprintf(os.Stderr, "%s's credential is minted afresh. Write this to %s on it and restart its host; "+
|
||||
"then push the machine running the bus so the user list carries the new hash.\n",
|
||||
node, catalogue.BusMembershipPath)
|
||||
fmt.Println(string(membership))
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
@@ -218,3 +219,127 @@ func notNow(err error) error {
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// SourceMoved is the forge announcing a merge: every module recorded as built from that
|
||||
// repository and branch is marked as moved to the merge commit, and built — bases first, so a
|
||||
// module that stands on another's artifact is built after it and not against the old one
|
||||
// (novox/hq 04-ISSUES/131). Nothing is pushed here: what a finished build does to the machines
|
||||
// running the module is the upgrade's decision, taken when the catalogue announces it.
|
||||
func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
|
||||
inv := f.open.inventory
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
if err != nil {
|
||||
return notNow(err)
|
||||
}
|
||||
var moved []inventory.Entry
|
||||
for _, e := range entries {
|
||||
if !sourceIs(e.Source, m) {
|
||||
continue
|
||||
}
|
||||
if e.Source.BuiltFrom == m.Commit {
|
||||
continue
|
||||
}
|
||||
if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil {
|
||||
return notNow(err)
|
||||
}
|
||||
moved = append(moved, e)
|
||||
}
|
||||
if len(moved) == 0 {
|
||||
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds is built from it\n",
|
||||
m.Owner, m.Repo, m.Base, m.Commit)
|
||||
return nil
|
||||
}
|
||||
ordered := orderByBases(moved)
|
||||
names := make([]string, 0, len(ordered))
|
||||
for _, e := range ordered {
|
||||
names = append(names, e.Manifest.Module)
|
||||
}
|
||||
fmt.Printf("%s/%s merged into %s (%.8s); building %s\n",
|
||||
m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", "))
|
||||
var failed []string
|
||||
for _, e := range ordered {
|
||||
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
|
||||
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 20*time.Minute); err != nil {
|
||||
fmt.Printf(" %s: %v\n", e.Manifest.Module, err)
|
||||
failed = append(failed, e.Manifest.Module)
|
||||
// A base that failed is a reason to stop: what stands on it would be built against
|
||||
// the old one, and report success (novox/hq 04-ISSUES/131).
|
||||
if standsOn(ordered, e.Manifest.Module) {
|
||||
fmt.Printf(" stopping: %s is a base of what was still to build\n", e.Manifest.Module)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(failed) > 0 {
|
||||
fmt.Printf("%d of %d not built: %s\n", len(failed), len(ordered), strings.Join(failed, ", "))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// sourceIs is whether a recorded source is the repository and branch a merge announced. A source on
|
||||
// the git seat is recorded as its path on the forge; one elsewhere as the URL it was cloned from.
|
||||
// An empty recorded ref is the repository's default branch, which is what a merge into the base
|
||||
// branch of the forge's default means.
|
||||
func sourceIs(s inventory.Source, m link.SourceMoved) bool {
|
||||
want := strings.ToLower(m.Owner + "/" + m.Repo)
|
||||
repo := strings.ToLower(strings.TrimSuffix(s.Repository, ".git"))
|
||||
matches := repo == want || strings.HasSuffix(repo, "/"+want) ||
|
||||
(m.CloneURL != "" && strings.EqualFold(strings.TrimSuffix(s.Repository, ".git"), strings.TrimSuffix(m.CloneURL, ".git")))
|
||||
if !matches {
|
||||
return false
|
||||
}
|
||||
return s.Ref == "" || s.Ref == m.Base
|
||||
}
|
||||
|
||||
// orderByBases is the entries with every base before what stands on it: a module whose build names
|
||||
// another's artifact under build.on comes after that module. Entries outside the set are not
|
||||
// waited for — they are not being rebuilt. Stable for what has no order between it.
|
||||
func orderByBases(entries []inventory.Entry) []inventory.Entry {
|
||||
inSet := map[string]bool{}
|
||||
for _, e := range entries {
|
||||
inSet[e.Manifest.Module] = true
|
||||
}
|
||||
var out []inventory.Entry
|
||||
placed := map[string]bool{}
|
||||
var place func(e inventory.Entry, seen map[string]bool)
|
||||
place = func(e inventory.Entry, seen map[string]bool) {
|
||||
name := e.Manifest.Module
|
||||
if placed[name] || seen[name] {
|
||||
return
|
||||
}
|
||||
seen[name] = true
|
||||
if e.Manifest.Build != nil {
|
||||
for _, on := range e.Manifest.Build.On {
|
||||
if on.Module == "" || on.Module == name || !inSet[on.Module] {
|
||||
continue
|
||||
}
|
||||
for _, base := range entries {
|
||||
if base.Manifest.Module == on.Module {
|
||||
place(base, seen)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
placed[name] = true
|
||||
out = append(out, e)
|
||||
}
|
||||
for _, e := range entries {
|
||||
place(e, map[string]bool{})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// standsOn is whether anything in the set is built on the named module's artifacts.
|
||||
func standsOn(entries []inventory.Entry, module string) bool {
|
||||
for _, e := range entries {
|
||||
if e.Manifest.Build == nil {
|
||||
continue
|
||||
}
|
||||
for _, on := range e.Manifest.Build.On {
|
||||
if on.Module == module {
|
||||
return true
|
||||
}
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -1,8 +1,14 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
@@ -24,21 +30,78 @@ type JetStream struct {
|
||||
|
||||
// Dial connects and returns the controller's JetStream handle.
|
||||
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
|
||||
// A name, because a connection nobody can identify in the server's own monitoring is one
|
||||
// nobody can attribute a problem to.
|
||||
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
|
||||
// **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a
|
||||
// public verifier would accept (design 25 §4: a host pins the server's exact certificate and
|
||||
// checks nothing else, and so does this). Without this, the first connection failed with
|
||||
// "certificate is not valid for any names" against a bus that was answering (2026-09-28).
|
||||
if path := strings.TrimSpace(os.Getenv(CertificateVar)); path != "" {
|
||||
pinned, err := pinnedTo(path)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
opts = append(opts, nats.Secure(pinned))
|
||||
}
|
||||
// **Its own inbox, and nothing wider.** Every principal is granted `_INBOX.<its user>.>` and
|
||||
// no other inbox; the client's default prefix is random, and the server refused the first
|
||||
// subscription to it (2026-09-28). The user is in the URL, so the prefix follows from it.
|
||||
if user, _, _ := CredentialIn(url); user != "" {
|
||||
opts = append(opts, nats.CustomInboxPrefix("_INBOX."+user))
|
||||
}
|
||||
// The address in an error is the address alone. The URL carries this controller's password,
|
||||
// and an error here is written on the assumption it will be logged.
|
||||
where := BareAddress(url)
|
||||
conn, err := nats.Connect(url, opts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("connecting to the bus at %s: %w", url, err)
|
||||
return nil, fmt.Errorf("connecting to the bus at %s: %w", where, err)
|
||||
}
|
||||
js, err := conn.JetStream()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", url, err)
|
||||
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", where, err)
|
||||
}
|
||||
return &JetStream{conn: conn, js: js}, nil
|
||||
}
|
||||
|
||||
// pinnedTo is a TLS configuration that accepts exactly the certificate in the file and no other:
|
||||
// the leaf's SHA-256, compared on every handshake, with the name and the chain deliberately not
|
||||
// consulted — a self-signed certificate with no names is the ordinary case for a mesh's bus.
|
||||
func pinnedTo(path string) (*tls.Config, error) {
|
||||
want, err := FingerprintOf(path)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return PinnedToFingerprint(want), nil
|
||||
}
|
||||
|
||||
// DialPinned is Dial with the server's certificate pinned by a fingerprint the caller already holds
|
||||
// — a module or a build machine that was handed one beside its credential, and has no file.
|
||||
func DialPinned(url, fingerprint string, opts ...nats.Option) (*JetStream, error) {
|
||||
if strings.TrimSpace(fingerprint) != "" {
|
||||
opts = append(opts, nats.Secure(PinnedToFingerprint(fingerprint)))
|
||||
}
|
||||
return Dial(url, opts...)
|
||||
}
|
||||
|
||||
// PinnedToFingerprint accepts exactly the certificate with this SHA-256 and no other.
|
||||
func PinnedToFingerprint(want string) *tls.Config {
|
||||
return &tls.Config{
|
||||
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin below, which is stricter
|
||||
MinVersion: tls.VersionTLS12,
|
||||
VerifyPeerCertificate: func(rawCerts [][]byte, _ [][]*x509.Certificate) error {
|
||||
if len(rawCerts) == 0 {
|
||||
return errors.New("the bus presented no certificate")
|
||||
}
|
||||
sum := sha256.Sum256(rawCerts[0])
|
||||
got := "sha256:" + hex.EncodeToString(sum[:])
|
||||
if got != want {
|
||||
return fmt.Errorf("the bus presented a certificate this mesh does not know (%s…), expected %s…", got[:23], want[:23])
|
||||
}
|
||||
return nil
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
|
||||
// a tool call — where a lost message is answered by the next one or by a timeout the caller
|
||||
// already handles (design 25 §3).
|
||||
|
||||
+29
-4
@@ -169,7 +169,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
// 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.>", "$JS.API.>"}
|
||||
sub = []string{"mesh.control.>", "$JS.API.>"}
|
||||
// **And where its consumers deliver.** A push consumer delivers on `_DELIVER.<its name>`,
|
||||
// and a client bound to it subscribes exactly that; the server refused it for every
|
||||
// principal the first time one bound a consumer (2026-09-28). Each kind below is granted
|
||||
// its own consumers' delivery subjects and no other's.
|
||||
sub = []string{"mesh.control.>", "$JS.API.>", "_DELIVER." + ControllerName, "_DELIVER." + ControllerName + ".>"}
|
||||
|
||||
// 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
|
||||
@@ -244,8 +248,15 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
case KindNode:
|
||||
// A host publishes its own node's control traffic and subscribes its own declaration —
|
||||
// and nothing of any other node's.
|
||||
pub = []string{"mesh.control." + p.Node + ".>"}
|
||||
sub = []string{"mesh.node." + p.Node + ".declare"}
|
||||
// And binding to its consumer, which asks the server about it (CONSUMER.INFO) — the one
|
||||
// thing the host does that nothing granted. Found the first time a machine dialled a
|
||||
// permissioned server: "this node cannot read its declarations" (2026-09-28). The ack and
|
||||
// the inbox are granted below with every principal's.
|
||||
pub = []string{
|
||||
"mesh.control." + p.Node + ".>",
|
||||
"$JS.API.CONSUMER.INFO.NODES." + p.Node,
|
||||
}
|
||||
sub = []string{"mesh.node." + p.Node + ".declare", "_DELIVER." + p.Node}
|
||||
|
||||
case KindModule:
|
||||
// 1. Its own namespace: it publishes its events there and serves its tools there. Nothing
|
||||
@@ -278,7 +289,18 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
}
|
||||
|
||||
// 3. Seats it holds: full participation.
|
||||
// Its consumer's name, not ConsumerFor: that asks for these permissions to build the
|
||||
// consumer, and would ask forever. A subject for a consumer that turns out not to exist
|
||||
// grants nothing anybody can use.
|
||||
sub = append(sub, "_DELIVER."+consumerDurable(p))
|
||||
for _, s := range p.Holds {
|
||||
// Taking work from the role's queue: the worker consumer it binds (asked about,
|
||||
// delivered on, acknowledged), each on the seat's own stream. The first machine to
|
||||
// take work over the new bus was refused the asking (2026-09-28).
|
||||
worker := "SEAT_" + upperSnake(s.Name) + "_worker"
|
||||
stream := seatStreamName(s.Name)
|
||||
sub = append(sub, "_DELIVER."+worker)
|
||||
pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>")
|
||||
for _, a := range s.Accepts {
|
||||
sub = append(sub, seatSubject(s, "accept", a))
|
||||
}
|
||||
@@ -514,7 +536,10 @@ func ComposeAccounts(principals []Principal) (string, error) {
|
||||
// One account for the mesh: accounts in NATS isolate subject spaces entirely, and the mesh is
|
||||
// one space (design 25 §4). The cost of that — that permissions are the only isolation — is
|
||||
// paid in the scoping of every inbox and every ack subject.
|
||||
b.WriteString("accounts {\n MESH {\n users = [\n")
|
||||
// JetStream is enabled per account once accounts exist at all: with only the global block set,
|
||||
// a user in MESH is told "JetStream not enabled for account" the first time it binds a
|
||||
// consumer, which is the first thing every host does (2026-09-28).
|
||||
b.WriteString("accounts {\n MESH {\n jetstream: enabled\n users = [\n")
|
||||
for _, p := range sorted {
|
||||
perms, err := PermissionsFor(p)
|
||||
if err != nil {
|
||||
|
||||
+8
-7
@@ -21,10 +21,11 @@ jetstream {
|
||||
|
||||
accounts {
|
||||
MESH {
|
||||
jetstream: enabled
|
||||
users = [
|
||||
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
|
||||
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.built"] }
|
||||
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
||||
@@ -32,21 +33,21 @@ accounts {
|
||||
subscribe: { allow: ["_INBOX.enrol.one.>"] }
|
||||
} }
|
||||
{ user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: {
|
||||
publish: { allow: ["$JS.ACK.NODES.one.>", "mesh.control.one.>"] }
|
||||
subscribe: { allow: ["_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||
publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] }
|
||||
subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||
} }
|
||||
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
||||
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
|
||||
subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "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.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
|
||||
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.one_telegram", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.status", "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.>"] }
|
||||
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] }
|
||||
subscribe: { allow: ["_DELIVER.two_audit", "_INBOX.two.audit.>", "mesh.mod.shop.event.order.placed"] }
|
||||
} }
|
||||
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
|
||||
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
|
||||
subscribe: { allow: ["_INBOX.two.shop.>"] }
|
||||
subscribe: { allow: ["_DELIVER.two_shop", "_INBOX.two.shop.>"] }
|
||||
} }
|
||||
]
|
||||
}
|
||||
|
||||
@@ -186,7 +186,13 @@ func TestWhatTheMeshWritesIsUsersAndNothingAboutTheServer(t *testing.T) {
|
||||
}
|
||||
// None of the server's own settings. Each of these in the mesh's file is a value the controller
|
||||
// would then own, and the module could no longer change its own image without the mesh agreeing.
|
||||
for _, absent := range []string{"port:", "http:", "jetstream", "tls {", "store_dir", "cert_file"} {
|
||||
// `jetstream {` is the server's block (its store, its limits); `jetstream: enabled` inside the
|
||||
// account is the account's, and the mesh owns the account — a user in it is told "JetStream
|
||||
// not enabled for account" without it (2026-09-28).
|
||||
if !strings.Contains(got, "jetstream: enabled") {
|
||||
t.Errorf("the account does not enable JetStream, so no user in it can bind a consumer")
|
||||
}
|
||||
for _, absent := range []string{"port:", "http:", "jetstream {", "tls {", "store_dir", "cert_file"} {
|
||||
if strings.Contains(got, absent) {
|
||||
t.Errorf("the accounts file contains %q, which belongs to the module that raises the "+
|
||||
"server, not to the mesh", absent)
|
||||
|
||||
@@ -143,3 +143,18 @@ func TestAMembershipForTheNewBusIsComposedAsASealedFile(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The store's seat rows have no protocol columns yet; loading them must not drop the protocol the
|
||||
// bus is derived from, or no role's work queue is ever raised (found live, 2026-09-28).
|
||||
func TestAStoreRowWithoutAProtocolKeepsTheCompiledOne(t *testing.T) {
|
||||
was := Seats()
|
||||
t.Cleanup(func() { UseSeats(was) })
|
||||
UseSeats([]Seat{{Name: "mesh-build-machine", Scope: ScopeMesh, Decision: "row"}})
|
||||
got, ok := SeatNamed("mesh-build-machine")
|
||||
if !ok || len(got.Accepts) == 0 {
|
||||
t.Fatalf("the build machine's seat lost what it accepts when loaded from the store: %+v", got)
|
||||
}
|
||||
if got.Decision != "row" {
|
||||
t.Fatalf("the store's own columns were not kept: %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
|
||||
"type": "container", "id": "server", "name": "mesh-controller",
|
||||
"env": map[string]any{
|
||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
|
||||
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
||||
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||
},
|
||||
@@ -27,7 +26,6 @@ func TestASeatPlaceholderAnswersWhereThisMachinePutTheHolder(t *testing.T) {
|
||||
env := control["env"].(map[string]any)
|
||||
for key, want := range map[string]string{
|
||||
"MESH_STORE_INVENTORY_PORT": "6852",
|
||||
"MESH_BROKER_AMQP_PORT": "5679",
|
||||
"MESH_BROKER_ADDRESS_PORT": "5671",
|
||||
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
|
||||
} {
|
||||
@@ -146,7 +144,6 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
|
||||
"MESH_STORE_INVENTORY_PORT": "6852",
|
||||
"MESH_STORE_IDENTITY_PORT": "6852",
|
||||
"MESH_STORE_LICENCES_PORT": "6852",
|
||||
"MESH_BROKER_AMQP_PORT": "5679",
|
||||
"MESH_BROKER_MANAGEMENT_PORT": "15673",
|
||||
"MESH_BROKER_ADDRESS_PORT": "5671",
|
||||
} {
|
||||
@@ -164,7 +161,7 @@ func TestTheControlPlanesOwnAddressesFollowTheNodesPorts(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
env, _ = fileNamed(out, "mesh-controller.server")["env"].(map[string]any)
|
||||
if env["MESH_STORE_INVENTORY_PORT"] != "" || env["MESH_BROKER_AMQP_PORT"] != "" {
|
||||
if env["MESH_STORE_INVENTORY_PORT"] != "" {
|
||||
t.Errorf("with no settings, the control plane is told %v", env)
|
||||
}
|
||||
}
|
||||
@@ -175,7 +172,6 @@ var SeatPorts = map[string]string{
|
||||
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
|
||||
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
|
||||
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}",
|
||||
"MESH_BROKER_AMQP_PORT": "${seat:mesh-broker:5672}",
|
||||
"MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}",
|
||||
"MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}",
|
||||
}
|
||||
|
||||
@@ -109,9 +109,30 @@ func DefaultSeats() []Seat { return append([]Seat(nil), defaultSeats...) }
|
||||
// than running on the set the binary shipped with. So the store can only ever *replace* the set with
|
||||
// a non-empty one, never erase it.
|
||||
func UseSeats(s []Seat) {
|
||||
if len(s) > 0 {
|
||||
seats = s
|
||||
if len(s) == 0 {
|
||||
return
|
||||
}
|
||||
// **The store's rows carry no protocol yet, and the protocol is what the bus is derived
|
||||
// from.** ADR 0129 gives a seat what it accepts, emits and serves; ADR 0122 moved the set into
|
||||
// a table that has name, scope, delivers and decision and nothing else, and the columns for
|
||||
// the rest are not there yet. So a row replacing a compiled entry would silently drop the
|
||||
// protocol, and the roles' work queues would never be raised — found live as "no response
|
||||
// from stream" the first time a build was submitted over the new bus (2026-09-28). Until the
|
||||
// table gains the columns, a row without a protocol keeps the compiled one of the same name.
|
||||
byName := map[string]Seat{}
|
||||
for _, d := range defaultSeats {
|
||||
byName[d.Name] = d
|
||||
}
|
||||
merged := make([]Seat, 0, len(s))
|
||||
for _, row := range s {
|
||||
if len(row.Accepts)+len(row.Emits)+len(row.Serves) == 0 {
|
||||
if d, known := byName[row.Name]; known {
|
||||
row.Accepts, row.Emits, row.Serves = d.Accepts, d.Emits, d.Serves
|
||||
}
|
||||
}
|
||||
merged = append(merged, row)
|
||||
}
|
||||
seats = merged
|
||||
}
|
||||
|
||||
// aliases maps a seat's former names to its current canonical name (novox/hq ADR 0122). Loaded from
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// OverNats is the controller's outbound on the bus being built: the same three acts the other
|
||||
// transport has, on the subjects the permissions were derived for (design 25). A declaration is a
|
||||
// JetStream publish into NODES, where the node's own consumer waits for it; an event is announced on
|
||||
// the subject its name derives to; a tool is asked by request and reply on the module's tool subject.
|
||||
type OverNats struct{ JS *broker.JetStream }
|
||||
|
||||
// declareSubject is where one node's declaration lands — the NODES stream's subject for it, and the
|
||||
// only subject that node's consumer delivers. The host subscribes exactly this.
|
||||
func declareSubject(node string) string { return "mesh.node." + node + ".declare" }
|
||||
|
||||
func (b OverNats) PublishDeclaration(ctx context.Context, node string, body []byte) error {
|
||||
publish, cancel := context.WithTimeout(ctx, 15*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.JS.Context().Publish(declareSubject(node), body, nats.Context(publish)); err != nil {
|
||||
return fmt.Errorf("declaring to %s: %w", node, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNats) PublishEvent(ctx context.Context, key, source, node string, body []byte) error {
|
||||
// The key is the subject: the controller's own events are named in full, and what a module
|
||||
// emits is derived before it reaches here. Headers carry the envelope the other transport put
|
||||
// in message properties (ADR 0042), so a consumer reads who and when without the payload.
|
||||
msg := nats.NewMsg(key)
|
||||
msg.Data = body
|
||||
msg.Header.Set("x-source", source)
|
||||
msg.Header.Set("x-node", node)
|
||||
msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano))
|
||||
if err := b.JS.Conn().PublishMsg(msg); err != nil {
|
||||
return fmt.Errorf("announcing %s: %w", key, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNats) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) {
|
||||
ask, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
reply, err := b.JS.Conn().RequestWithContext(ask, "mesh.mod."+module+".tool."+tool, args)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err)
|
||||
}
|
||||
return reply.Data, nil
|
||||
}
|
||||
|
||||
// ConnectNats is Connect for the bus being built: the controller's inbound and outbound over one
|
||||
// JetStream connection the caller has already raised the streams on. Nothing is declared here —
|
||||
// the streams and the controller's consumers are asserted by Raise, before anything is served.
|
||||
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
|
||||
return &Server{
|
||||
inbound: Nats(js),
|
||||
bus: OverNats{JS: js},
|
||||
js: js,
|
||||
enroller: enroller,
|
||||
listener: listener,
|
||||
log: newLog(),
|
||||
}
|
||||
}
|
||||
|
||||
// Bus is the controller's outbound, whichever transport it connected over. Callers that send a
|
||||
// declaration or ask a tool use this rather than the channel, which one transport does not have.
|
||||
func (s *Server) Bus() Bus { return s.bus }
|
||||
@@ -117,6 +117,19 @@ type Announcement struct {
|
||||
}
|
||||
|
||||
// Upgraded is what the catalogue says when a module's current version moves.
|
||||
// SourceMoved is what the forge announces when a pull request is merged: which repository, into
|
||||
// which branch, producing which commit. The mesh matches it against every module's recorded
|
||||
// source and builds what moved, bases first.
|
||||
type SourceMoved struct {
|
||||
Owner string `json:"owner"`
|
||||
Repo string `json:"repo"`
|
||||
Base string `json:"base"`
|
||||
Head string `json:"head"`
|
||||
Commit string `json:"merge_commit_sha"`
|
||||
CloneURL string `json:"clone_url"`
|
||||
HTMLURL string `json:"html_url"`
|
||||
}
|
||||
|
||||
type Upgraded struct {
|
||||
Module string `json:"module"`
|
||||
Commit string `json:"commit"`
|
||||
|
||||
@@ -30,6 +30,9 @@ const (
|
||||
KindHeartbeat = "heartbeat"
|
||||
KindBuilt = "built"
|
||||
KindModuleMoved = "module-moved"
|
||||
// KindSourceMoved is the forge announcing a merge: a source moved, and what it produces is
|
||||
// built without anybody telling the mesh (novox/hq 04-ISSUES/131).
|
||||
KindSourceMoved = "source-moved"
|
||||
KindCatchUp = "catch-up"
|
||||
)
|
||||
|
||||
|
||||
@@ -67,6 +67,10 @@ func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound {
|
||||
// are events get their own queue each, and only when something is listening.
|
||||
func (c *currentInbound) Also(kind string) error {
|
||||
switch kind {
|
||||
case KindSourceMoved:
|
||||
// Not followed on the bus the mesh is leaving: the forge's merges are announced on the
|
||||
// new one, and this transport goes with the move (design 28, task 5.5).
|
||||
return nil
|
||||
case KindModuleMoved:
|
||||
if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil {
|
||||
return err
|
||||
|
||||
@@ -48,7 +48,7 @@ func Nats(js *broker.JetStream) Inbound {
|
||||
// whatever was asked for — and not at all when nothing was.
|
||||
func (n *natsInbound) Also(kind string) error {
|
||||
switch kind {
|
||||
case KindModuleMoved, KindCatchUp:
|
||||
case KindModuleMoved, KindCatchUp, KindSourceMoved:
|
||||
n.follows[kind] = true
|
||||
return nil
|
||||
default:
|
||||
@@ -180,6 +180,8 @@ func kindOfSubject(subject string) (string, bool) {
|
||||
return KindModuleMoved, true
|
||||
case broker.ControllerFollows[1]:
|
||||
return KindCatchUp, true
|
||||
case broker.ControllerFollows[3]:
|
||||
return KindSourceMoved, true
|
||||
case BuildOutcome():
|
||||
// A build's outcome is the role's event now, so it arrives on the events stream rather than
|
||||
// the control branch — and is acted on by the same handler, because what the controller does
|
||||
|
||||
@@ -374,3 +374,6 @@ func TestNatsAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
|
||||
t.Fatalf("a report was recorded by a store that never came back")
|
||||
}
|
||||
}
|
||||
|
||||
// A merge announcement is not what these tests are about; taken and forgotten.
|
||||
func (t *toldAbout) SourceMoved(context.Context, SourceMoved) error { return nil }
|
||||
|
||||
+36
-1
@@ -7,6 +7,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"log"
|
||||
"os"
|
||||
"time"
|
||||
@@ -57,6 +58,8 @@ type Upgrader interface {
|
||||
// stop every upgrade behind it — except the store unreachable for the moment, which is asked
|
||||
// again for a bounded time (novox/hq issue 083).
|
||||
Upgraded(ctx context.Context, u Upgraded) error
|
||||
// SourceMoved is a merge on the forge: build what that source produces, bases first.
|
||||
SourceMoved(ctx context.Context, m SourceMoved) error
|
||||
}
|
||||
|
||||
// Server acts on what nodes and modules say.
|
||||
@@ -70,6 +73,7 @@ type Server struct {
|
||||
bus Bus
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
js *broker.JetStream
|
||||
|
||||
enroller Enroller
|
||||
listener Listener
|
||||
@@ -99,6 +103,9 @@ func (s *Server) Records(r Recorder) { s.recorder = r }
|
||||
|
||||
// Follows says what to do about upgrades, and asks for them to be delivered.
|
||||
func (s *Server) Follows(u Upgrader) error {
|
||||
if err := s.inbound.Also(KindSourceMoved); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.inbound.Also(KindModuleMoved); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -180,10 +187,12 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
||||
channel: channel,
|
||||
enroller: enroller,
|
||||
listener: listener,
|
||||
log: log.New(os.Stdout, "", log.LstdFlags),
|
||||
log: newLog(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) }
|
||||
|
||||
// Channel is the controller's channel, for the command line's own publishing.
|
||||
func (s *Server) Channel() *amqp.Channel { return s.channel }
|
||||
|
||||
@@ -197,6 +206,9 @@ func (s *Server) Close() {
|
||||
if s.conn != nil {
|
||||
_ = s.conn.Close()
|
||||
}
|
||||
if s.js != nil {
|
||||
s.js.Close()
|
||||
}
|
||||
}
|
||||
|
||||
// Serve acts on what arrives until the context ends.
|
||||
@@ -225,6 +237,8 @@ func (s *Server) act(ctx context.Context, m Control) {
|
||||
s.wasBuilt(ctx, m)
|
||||
case KindModuleMoved:
|
||||
s.moved(ctx, m)
|
||||
case KindSourceMoved:
|
||||
s.sourceMoved(ctx, m)
|
||||
case KindCatchUp:
|
||||
s.catchingUp(ctx, m)
|
||||
default:
|
||||
@@ -596,3 +610,24 @@ func short(commit string) string {
|
||||
}
|
||||
return commit
|
||||
}
|
||||
|
||||
// sourceMoved acts on the forge's announcement of a merge. Taken whatever happens: a build that
|
||||
// fails is reported by the build itself, and re-delivering the merge would only re-fail it.
|
||||
func (s *Server) sourceMoved(ctx context.Context, m Control) {
|
||||
var moved SourceMoved
|
||||
if err := json.Unmarshal(m.Body(), &moved); err != nil {
|
||||
s.log.Printf("a merge announcement could not be read: %v", err)
|
||||
_ = m.Took()
|
||||
return
|
||||
}
|
||||
m.About("merge " + moved.Owner + "/" + moved.Repo + " into " + moved.Base)
|
||||
if moved.Commit == "" || moved.Repo == "" {
|
||||
s.log.Printf("a merge announcement named no repository or no commit; ignored")
|
||||
_ = m.Took()
|
||||
return
|
||||
}
|
||||
if err := s.upgrader.SourceMoved(ctx, moved); err != nil {
|
||||
s.log.Printf("%s/%s moved to %.8s and the mesh could not act on it: %v", moved.Owner, moved.Repo, moved.Commit, err)
|
||||
}
|
||||
_ = m.Took()
|
||||
}
|
||||
|
||||
@@ -152,3 +152,5 @@ func TestAnUpgradeHandledDuringShutdownIsLeftForTheBus(t *testing.T) {
|
||||
t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to)
|
||||
}
|
||||
}
|
||||
|
||||
func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return nil }
|
||||
|
||||
@@ -62,6 +62,7 @@
|
||||
"/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro",
|
||||
"/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro",
|
||||
"/var/lib/mesh/mesh-controller/broker:/run/secrets/broker:ro",
|
||||
"/var/lib/mesh/mesh-controller/bus:/run/secrets/bus:ro",
|
||||
"/var/lib/mesh/mesh-controller/broker-management:/run/secrets/broker-management:ro",
|
||||
"/var/lib/mesh/mesh-controller/broker-address:/run/secrets/broker-address:ro"
|
||||
],
|
||||
|
||||
Reference in New Issue
Block a user