Compare commits

..
Author SHA1 Message Date
jschoubben e09a14ba4c A served value may name the consumer it is served to (hq ADR 0188)
${consumer:as} and ${consumer:as:dns} in a serves block are filled per
consumer at resolution, and the one filled value reaches both ends: the
consumer's binding and its ${bound:...} substitutions, and the provider's
contributions entry as `derived`. A fact or alphabet the mesh does not have
is refused at parse; a consumer whose own file already holds the derived
value is refused at resolution, naming the placeholder to write instead.
2026-10-02 21:24:43 +02:00
56 changed files with 834 additions and 2377 deletions
+1 -21
View File
@@ -137,12 +137,7 @@ func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
if err != nil {
return nil, err
}
// **The seat this machine serves is the one its credential claims** (novox/hq ADR 0190, the
// handover): the mesh issues a build machine's credential naming the seat its module claims,
// and one binary serves the old role as `builder` and the new as `build-agent` from that alone.
seat := link.BuildSeatClaimed(credential.seatsClaimed())
fmt.Fprintf(os.Stderr, "taking build work as a holder of %s\n", seat)
return link.MachineOverNATSOn(js, on, seat), nil
return link.MachineOverNATS(js, on), nil
}
// answer does one build and says what happened, whichever way it went.
@@ -458,21 +453,6 @@ type Credential struct {
// as two fields and this machine joins them once, here, to dial.
User string `json:"user,omitempty"`
Password string `json:"password,omitempty"`
// Claims are the seats the module this credential was issued for claims, as the mesh writes
// them beside the credential (novox/hq ADR 0159). The first is the build role this machine
// serves; a credential naming none is from before claims travelled in it.
Claims []struct {
Seat string `json:"seat"`
} `json:"claims,omitempty"`
}
// seatsClaimed is the seats the credential names, in order.
func (c Credential) seatsClaimed() []string {
out := make([]string, 0, len(c.Claims))
for _, claim := range c.Claims {
out = append(out, claim.Seat)
}
return out
}
// onTheNewBus is whether a credential is for the bus being built: its address says so, and the
-27
View File
@@ -1,27 +0,0 @@
package main
import (
"encoding/json"
"testing"
"github.com/novox/mesh-controller/internal/link"
)
// The seat a build machine serves comes from its credential (novox/hq ADR 0190 handover).
func TestTheCredentialSaysWhichBuildRoleThisMachineServes(t *testing.T) {
var held Credential
if err := json.Unmarshal([]byte(`{"url":"nats://bus:4222","user":"anchor.builder","password":"x",
"claims":[{"seat":"mesh-build-machine","scope":"mesh","serves":[]}]}`), &held); err != nil {
t.Fatal(err)
}
if got := link.BuildSeatClaimed(held.seatsClaimed()); got != "mesh-build-machine" {
t.Errorf("the old builder's credential serves %q", got)
}
var bare Credential
if err := json.Unmarshal([]byte(`{"url":"nats://bus:4222","user":"anchor.build-agent","password":"x"}`), &bare); err != nil {
t.Fatal(err)
}
if got := link.BuildSeatClaimed(bare.seatsClaimed()); got != link.TheBuildMachine {
t.Errorf("a credential without claims serves %q, want %s", got, link.TheBuildMachine)
}
}
-38
View File
@@ -3,7 +3,6 @@ package main
import (
"context"
"fmt"
"github.com/novox/mesh-controller/internal/broker"
"sort"
"strings"
@@ -68,13 +67,6 @@ func assign(ctx context.Context, open *stores, node, module string) (string, err
for _, line := range settled {
said += "\n " + line
}
// Its bus credential, in the same act (novox/hq issue 203): an assignment pushed before its
// credential exists delivers a process that cannot authenticate and crash-loops until somebody
// runs a second verb and a second push. Issued here when the module speaks on the bus and has
// no credential yet; kept when it has one, so re-assigning rotates nothing.
if line := issueOnAssign(ctx, open, node, module); line != "" {
said += "\n " + line
}
plan, _, err := planFor(ctx, open, node)
if err != nil {
// Kept, and still refused. Both halves are the answer, and the rest of the mesh is still
@@ -160,33 +152,3 @@ func blockedElsewhere(ctx context.Context, open *stores, except string) string {
out.WriteString("\nThis may or may not be what just changed — it is what is true now.")
return out.String()
}
// issueOnAssign gives a newly assigned module its bus credential, the way `module issue` does, and
// says what it did in one line. Nothing for a module that declares no broker secret; nothing for one
// whose user is already minted (a credential is rotated on purpose, never by re-assigning); and when
// the bus cannot be reached from here, the line names the verb and the push that would refuse the
// module until it is run — never a silent placeholder (novox/hq issue 203).
func issueOnAssign(ctx context.Context, open *stores, node, module string) string {
inv := open.inventory
shelf, err := inv.Catalogue(ctx)
if err != nil {
return ""
}
m, known := shelf[module]
if !known || mayIssue(m) != nil {
return ""
}
user := broker.Principal{Kind: broker.KindModule, Node: node, Module: module}.Username()
if _, minted, err := inv.BusUserHash(ctx, user); err != nil || minted {
return ""
}
busAddress, err := broker.BusAddress()
if err == nil {
err = issueOnTheNewBus(ctx, inv, m, node, busAddress)
}
if err != nil {
return fmt.Sprintf("its bus credential is not issued (%v): `module issue %s --node %s` first — "+
"`push %s` refuses to send %s until it is", err, module, node, node, module)
}
return fmt.Sprintf("its bus credential is issued and sealed to %s, and arrives with the push", node)
}
-6
View File
@@ -175,12 +175,6 @@ func TestTheControlPlaneIsToldWhereTheNodePutTheStoreAndTheBroker(t *testing.T)
Guards: []int{15672},
Resources: []map[string]any{{"id": "server", "type": "container", "name": "mesh-broker",
"ports": []any{"5671:5671", "5672:5672", "127.0.0.1:15672:15672"}, "image": "mq@" + aDigest}}})
// The control plane's own bus user is the installer's, seeded at genesis before the controller
// runs (SeedBusUser); without it a push now refuses the credential nobody issued (issue 203).
if err := open.inventory.SeedBusUser(ctx, inventory.BusUser{Username: "anchor.mesh-controller",
Kind: inventory.BusController, Node: "anchor", Module: "mesh-controller"}, "bootstrap"); err != nil {
t.Fatal(err)
}
if _, err := assign(ctx, open, "anchor", "mesh-controller"); err != nil {
t.Fatal(err)
}
+6 -57
View File
@@ -430,13 +430,11 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
}
fmt.Println()
seat := buildSeatHeld(ctx)
ask, err := askOverOn(seat)
ask, err := askOver(server)
if err != nil {
return err
}
defer ask.Close()
fmt.Printf(" of %s\n", seat)
if wait == 0 {
// Asked and not waited for (novox/hq issue 176): the outcome is the role's event, and the
@@ -524,8 +522,6 @@ func takeIn(ctx context.Context, inv *inventory.Inventory, result link.BuildResu
recorded := inventory.Source{
Repository: result.Repository, Path: result.Path, Ref: result.Ref,
BuiltFrom: result.Commit, Head: result.Commit,
// What it stood on, so registration can judge a built manifest's base (to-be 38 WP2.4).
Against: kept.Against,
}
if result.Source != nil && result.Source.Seat != "" {
recorded.Repository, recorded.Seat = result.Source.Repository, result.Source.Seat
@@ -557,7 +553,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
}
defer server.Close()
ask, err := askOverOn(buildSeatHeld(ctx))
ask, err := askOver(server)
if err != nil {
return err
}
@@ -670,56 +666,12 @@ func heldBy(ctx context.Context) map[string]string {
// **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5). On the bus
// the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus
// being built it dials, because a build request is a one-shot and holds nothing else.
func askOverOn(seat string) (link.Builders, error) {
func askOver(_ *link.Server) (link.Builders, error) {
address, err := broker.BusAddress()
if err != nil {
return nil, err
}
return link.BuildsOverNATSOn(address, seat)
}
// buildSeatHeld is the build role to ask: the one some assigned module claims (novox/hq ADR 0190,
// the handover). Read from the catalogue at ask time, because the answer changes exactly once, the
// moment the first build-agent is assigned — and a controller that asked the new role before then
// would queue work nothing takes, while the outcome that registers build-agent itself has to come
// from the old builder. When the catalogue cannot be read the current role is asked, said aloud.
func buildSeatHeld(ctx context.Context) string {
open, err := openStores(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "could not read what is assigned, so the build is asked of %s: %v\n",
link.TheBuildMachine, err)
return link.TheBuildMachine
}
defer open.Close()
entries, err := open.inventory.Catalogued(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "could not read the catalogue, so the build is asked of %s: %v\n",
link.TheBuildMachine, err)
return link.TheBuildMachine
}
return buildSeatAmong(entries)
}
// buildSeatAmong is the rule, over what the catalogue holds: the current build role when any
// assigned module claims it; else the retired role while an assigned module still claims that; else
// the current role, which is where every ask goes once the handover is done.
func buildSeatAmong(entries []inventory.Entry) string {
heldBefore := false
for _, e := range entries {
if len(e.On) == 0 {
continue
}
if e.Manifest.ClaimsSeat(link.TheBuildMachine) {
return link.TheBuildMachine
}
if e.Manifest.ClaimsSeat(link.TheBuildMachineBefore) {
heldBefore = true
}
}
if heldBefore {
return link.TheBuildMachineBefore
}
return link.TheBuildMachine
return link.BuildsOverNATS(address)
}
// buildLog prints everything a build machine said about one build, read back from the bus.
@@ -739,13 +691,10 @@ func buildLog(ctx context.Context, id string) error {
}
defer js.Close()
// Under whichever build role did it: a build asked of the retired role during the handover
// (ADR 0190) said its lines as that role's events, and a reader should not have to know which.
lines := link.BuildLogOf("*", id)
sub, err := js.Context().PullSubscribe(lines, "",
sub, err := js.Context().PullSubscribe(link.BuildLog(id), "",
nats.BindStream(broker.EventsStream), nats.DeliverAll(), nats.AckNone())
if err != nil {
return fmt.Errorf("cannot read %s from the bus: %w", lines, err)
return fmt.Errorf("cannot read %s from the bus: %w", link.BuildLog(id), err)
}
defer func() { _ = sub.Unsubscribe() }()
-43
View File
@@ -1,43 +0,0 @@
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 claiming(module, seat string, on ...string) inventory.Entry {
return inventory.Entry{
Manifest: catalogue.Manifest{Module: module, Claims: []catalogue.Claim{{Name: seat}}},
On: on,
}
}
// The controller asks the build role that has a holder (novox/hq ADR 0190 handover): the retired
// one while only the builder is assigned, the current one from the first build-agent on, and the
// current one when nothing holds either — where every ask goes once the handover is done.
func TestTheControllerAsksTheBuildRoleThatHasAHolder(t *testing.T) {
onlyTheBuilder := []inventory.Entry{
claiming("builder", link.TheBuildMachineBefore, "anchor"),
claiming("build-agent", link.TheBuildMachine), // registered, assigned nowhere yet
}
if got := buildSeatAmong(onlyTheBuilder); got != link.TheBuildMachineBefore {
t.Errorf("with only the builder assigned, asked %q", got)
}
bothHeld := []inventory.Entry{
claiming("builder", link.TheBuildMachineBefore, "anchor"),
claiming("build-agent", link.TheBuildMachine, "home-server"),
}
if got := buildSeatAmong(bothHeld); got != link.TheBuildMachine {
t.Errorf("with a build-agent assigned anywhere, asked %q", got)
}
neither := []inventory.Entry{claiming("builder", link.TheBuildMachineBefore)}
if got := buildSeatAmong(neither); got != link.TheBuildMachine {
t.Errorf("with no holder of either, asked %q, want the current role", got)
}
if got := buildSeatAmong(nil); got != link.TheBuildMachine {
t.Errorf("an empty catalogue asks %q", got)
}
}
@@ -1,90 +0,0 @@
package main
import (
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// A fresh assignment is pushed before its credential exists (novox/hq issue 203): `assign` recorded
// the module, `push` sealed a random own secret where the bus credential belongs, and the process
// crash-looped until a person ran `module issue` and pushed again. Now assigning a module that speaks
// on the bus issues its credential in the same act — or, when the bus cannot be reached from here,
// says which verb to run — and a push never seals a placeholder in a credential's place.
func aTalker() catalogue.Manifest {
return catalogue.Manifest{Module: "talker", Version: "1",
OwnSecrets: catalogue.OwnSecrets{"broker": {Path: "/var/lib/mesh/talker/broker"}},
Resources: []map[string]any{
{"id": "state", "type": "directory", "path": "/var/lib/mesh/talker", "mode": "0700"},
}}
}
func TestAssigningAModuleThatSpeaksOnTheBusNamesItsCredential(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
register(t, open, aTalker())
// No bus is known to this process, so the credential cannot be issued here: the assignment
// stands and says exactly what must happen before a push — never silently.
said, err := assign(ctx, open, "laptop", "talker")
if err != nil {
t.Fatal(err)
}
if !strings.Contains(said, "module issue talker --node laptop") {
t.Fatalf("an assignment whose credential could not be issued does not name the verb:\n%s", said)
}
// And the push refuses to send it, naming the same verb, rather than sealing a placeholder.
plan, settings, err := planFor(ctx, open, "laptop")
if err != nil {
t.Fatal(err)
}
_, err = declarationFor(ctx, open, "laptop", plan, settings)
if err == nil {
t.Fatal("a push sealed a placeholder where talker's bus credential belongs")
}
if !strings.Contains(err.Error(), "module issue talker --node laptop") || !strings.Contains(err.Error(), "issue 203") {
t.Fatalf("the refusal does not say what to run: %v", err)
}
// Once the user is minted, the push goes on to the credential the mesh sealed, and re-assigning
// does not mint again: a credential rotates on purpose, never by habit.
if _, err := open.inventory.MintBusPassword(ctx, inventory.BusUser{
Username: "laptop.talker", Kind: inventory.BusModule, Node: "laptop", Module: "talker"}); err != nil {
t.Fatal(err)
}
hash, _, err := open.inventory.BusUserHash(ctx, "laptop.talker")
if err != nil {
t.Fatal(err)
}
said, err = assign(ctx, open, "laptop", "talker")
if err != nil {
t.Fatal(err)
}
if strings.Contains(said, "module issue") {
t.Fatalf("a module with a minted credential was told to issue one:\n%s", said)
}
again, _, err := open.inventory.BusUserHash(ctx, "laptop.talker")
if err != nil {
t.Fatal(err)
}
if again != hash {
t.Fatal("re-assigning rotated the credential")
}
}
// A module that declares no broker secret is left alone: nothing to issue, nothing said.
func TestAssigningAModuleThatDoesNotSpeakSaysNothingOfCredentials(t *testing.T) {
open := aMesh(t)
register(t, open, helloWeb())
said, err := assign(t.Context(), open, "laptop", "hello-web")
if err != nil {
t.Fatal(err)
}
if strings.Contains(said, "credential") {
t.Fatalf("a module without a broker secret was told about credentials:\n%s", said)
}
}
@@ -0,0 +1,33 @@
package main
// The broker opening belongs only on the node that listens on it (novox/hq: it leaked onto
// every enrolled node's declaration, opening a from-anywhere hole for a port nothing there
// serves). foundationPortsFor is the scope.
import (
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
func TestTheBrokerHostGetsTheFoundationOpening(t *testing.T) {
broker := catalogue.Manifest{Module: "lavinmq", Listens: []catalogue.Listening{
{Port: 5671, Protocol: "tcp", From: "mesh"},
{Port: 5672, Protocol: "tcp", From: "mesh"},
}}
got := foundationPortsFor(5671, []catalogue.Manifest{broker})
if len(got) != 1 || got[0] != 5671 {
t.Fatalf("the node that listens on the broker port keeps it; got %v", got)
}
}
func TestANodeThatOnlyDialsTheBrokerGetsNoOpening(t *testing.T) {
// ace's set: things that reach the broker as a client, none listening on 5671.
ace := []catalogue.Manifest{
{Module: "plex", Listens: []catalogue.Listening{{Port: 32400, Protocol: "tcp", From: "anywhere"}}},
{Module: "postgres", Listens: []catalogue.Listening{{Port: 5432, Protocol: "tcp", From: "mesh"}}},
}
if got := foundationPortsFor(5671, ace); got != nil {
t.Fatalf("a node that only dials out opens nothing for the broker; got %v", got)
}
}
+1 -11
View File
@@ -557,7 +557,7 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
user := broker.Principal{Kind: broker.KindModule, Node: node, Module: m.Module}.Username()
password, err := inv.MintBusPassword(ctx, inventory.BusUser{
Username: user, Kind: busKindOf(m.Module), Node: node, Module: m.Module,
Username: user, Kind: inventory.BusModule, Node: node, Module: m.Module,
})
if err != nil {
return err
@@ -577,16 +577,6 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
return issueWith(ctx, inv, m, node, busAddress, known, reachable, user, password)
}
// busKindOf is what a module's bus user is recorded as: the node's tool runtime where the module is
// the runtime (novox/hq ADR 0175), a module otherwise. The username is the same either way — the
// runtime is issued through this same path — and the kind is what a reader of the records sees.
func busKindOf(module string) string {
if module == catalogue.RuntimeModule {
return inventory.BusNodeTools
}
return inventory.BusModule
}
// issueWith is the delivery half: the minted password sealed to the machine as the module's broker
// secret, and the module's consumer created where the bus can be reached. Split from the minting
// so the move can issue every module against a bus whose address it worked out itself
+1 -2
View File
@@ -1,7 +1,6 @@
package main
import (
"reflect"
"strings"
"testing"
"time"
@@ -192,7 +191,7 @@ func TestWhatAHandedOverModuleRecordsAboutItsSource(t *testing.T) {
t.Fatalf("the source records as %+v", from)
}
// A manifest with no provenance at all is legitimate: fixing something in a hurry.
if from, err := whereItComesFrom("", "", "", "", false); err != nil || !reflect.DeepEqual(from, inventory.Source{}) {
if from, err := whereItComesFrom("", "", "", "", false); err != nil || from != (inventory.Source{}) {
t.Fatalf("a manifest handed over with no provenance was refused: %+v, %v", from, err)
}
for _, c := range []struct {
+37 -22
View File
@@ -16,6 +16,8 @@ import (
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/licences"
"github.com/novox/mesh-controller/internal/overlay"
"net"
"strconv"
)
// working out what one machine should be.
@@ -533,23 +535,6 @@ func renderingFor(ctx context.Context, open *stores, node string,
var sealed string
var err error
if choosing == Allocating {
// **The broker credential is never invented here** (novox/hq issue 203). Every other
// own secret is the mesh's to make — a password nobody else knows — but this one
// is an account on the bus, minted by `module issue` and sealed by it; a push that
// made a random one would deliver a file the process cannot read and report the
// machine applied. Refused by name, with the verb.
if name == "broker" {
user := broker.Principal{Kind: broker.KindModule, Node: node, Module: m.Module}.Username()
if _, minted, err := inv.BusUserHash(ctx, user); err != nil {
return catalogue.Rendering{}, inventory.Node{}, err
} else if !minted {
return catalogue.Rendering{}, inventory.Node{}, fmt.Errorf(
"%s on %s has no bus credential: nothing was issued for %s, and a push "+
"would seal a placeholder its process cannot read (novox/hq issue 203). "+
"`module issue %s --node %s`, then push again",
m.Module, node, user, m.Module, node)
}
}
sealed, err = inv.SecretForModule(ctx, node, m.Module, name)
} else {
var held bool
@@ -663,12 +648,25 @@ func renderingFor(ctx context.Context, open *stores, node string,
names[name] = at
}
// **The bus is never public** (novox/hq ADR 0169). It was a foundation port — widened from the
// broker's own `from: mesh` to from-anywhere on the broker's host, so a machine could enrol
// before it had an address on the private network. A machine joins through the tunnel now, and
// every link to the bus crosses it, so its reach is what the `nats` module declares: the mesh.
// Nothing the mesh itself needs is opened beyond what a module declares.
// The ports the mesh itself needs open, which no module declares. Read from the broker this
// control plane was told about rather than written down twice: the address a node is handed in
// its token and the port its machine must accept on are the same fact.
//
// **Only on the node that listens on it** (novox/hq issue: the broker opening leaked onto
// every node). The opening exists to WIDEN the broker's port to from-anywhere — a machine
// enrolling is not on the mesh yet, so the broker's own `from: mesh` listen would refuse its
// first dial. That widening belongs on the broker's host and nowhere else: a node that only
// dials out needs no incoming rule, and an opening for a port nothing here listens on is a
// from-anywhere hole for a dead port. So the foundation port is kept only when a module
// resolved onto THIS node actually listens on it.
var foundation []int
if b, err := broker.FromEnvironment(); err == nil {
if _, port, err := net.SplitHostPort(b.Address); err == nil {
if n, err := strconv.Atoi(port); err == nil {
foundation = foundationPortsFor(n, plan.Modules)
}
}
}
// And, for a module that keeps them, every operator-sealed secret in the mesh — the vault's
// copy, outside the store (novox/hq ADR 0085, amended). Read only; nothing here mints. The
@@ -1382,6 +1380,23 @@ func composeBusUsers(ctx context.Context, inv *inventory.Inventory,
return broker.ComposeAccounts(filled)
}
// foundationPortsFor is the broker port, kept only when a module resolved onto this node listens
// on it (novox/hq issue: the broker opening leaked onto every node). The foundation opening
// exists to WIDEN the broker's `from: mesh` port to from-anywhere, because a machine enrolling is
// not on the mesh yet and its first dial would be refused. That widening belongs on the broker's
// host alone: a node that only dials out needs no incoming rule, and an opening for a port
// nothing here listens on is a from-anywhere hole for a dead port.
func foundationPortsFor(brokerPort int, modules []catalogue.Manifest) []int {
for _, m := range modules {
for _, l := range m.Listens {
if l.Port == brokerPort {
return []int{brokerPort}
}
}
}
return nil
}
// providerModuleOf is which module answers a need on the providing node: the one in this node's
// own set when the provider is here, else the one the catalogue says offers it.
func providerModuleOf(resolved catalogue.Resolution, open *stores, ctx context.Context, n catalogue.Needed) string {
+5 -20
View File
@@ -36,29 +36,17 @@ func tiersOf(set []string, edges []inventory.Edge) [][]string {
for _, m := range set {
deps[m] = map[string]bool{}
}
// The build seat's holders follow the controller that defines their worker (EdgeWorkerOf,
// novox/hq issue 206), so the built-by edge from that controller to such a holder yields: the
// controller is built by whichever build machine is running, as the runtime image always was.
worker := map[string]map[string]bool{}
for _, e := range edges {
if e.Kind == inventory.EdgeWorkerOf && in[e.From] && in[e.To] {
if worker[e.To] == nil {
worker[e.To] = map[string]bool{}
}
worker[e.To][e.From] = true
}
}
for _, e := range edges {
// A code dependency — B packages A's source — rebuilds B with A, in the same tier: B's
// build needs nothing of A's first. The other kinds order: stands-on and declared after
// the base is built, built-by after the build machine is built and running — except for
// what the build machine itself stands on, and for the controller whose worker the build
// machine binds. The runtime image is built by the builder and the builder is built on the
// runtime image; the image comes first, built by the builder that is running.
// what the build machine itself stands on. The runtime image is built by the builder and
// the builder is built on the runtime image; the image comes first, built by the builder
// that is running, which is the only one there could be.
if !in[e.From] || !in[e.To] || e.From == e.To || e.Kind == inventory.EdgePackages {
continue
}
if e.Kind == inventory.EdgeBuiltBy && (isBaseOf(e.From, e.To, edges, in) || worker[e.From][e.To]) {
if e.Kind == inventory.EdgeBuiltBy && isBaseOf(e.From, e.To, edges, in) {
continue
}
deps[e.From][e.To] = true
@@ -134,10 +122,7 @@ func reachableFrom(moved []string, edges []inventory.Edge) []string {
for grew := true; grew; {
grew = false
for _, e := range edges {
// Built-by and worker-of order a plan; neither widens it. A new build machine changes
// nothing it builds, and a new controller changes nothing about the holder it orders —
// what packages the controller's source is already a code edge.
if e.Kind == inventory.EdgeBuiltBy || e.Kind == inventory.EdgeWorkerOf {
if e.Kind == inventory.EdgeBuiltBy {
continue
}
if in[e.To] && !in[e.From] {
+2 -4
View File
@@ -346,9 +346,7 @@ func rolloutMint(ctx context.Context, again bool) error {
}
machines++
case broker.KindModule, broker.KindNodeTools:
// The runtime is minted and delivered exactly as a module is (novox/hq ADR 0175): it is
// issued as the module it stands for, to that module's `broker` secret.
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
@@ -367,7 +365,7 @@ func rolloutMint(ctx context.Context, again bool) error {
skipped++
continue
}
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: busKindOf(p.Module), Node: p.Node, Module: p.Module})
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusModule, Node: p.Node, Module: p.Module})
if err != nil {
return err
}
-45
View File
@@ -1,45 +0,0 @@
package main
import (
"testing"
"github.com/novox/mesh-controller/internal/inventory"
)
// The holder of the build seat follows the controller that defines its worker (novox/hq issue 206).
// On 2026-10-03 a plan put the build machine in tier 0 and the controller in tier 1; the new build
// machine could not bind the worker the old controller had defined, and nothing could build the
// controller that would have redefined it. The built-by edge from the controller to its build
// machine yields to that order: the controller is built by whichever build machine is running.
func TestTheBuildSeatsHolderFollowsTheControllerThatDefinesItsWorker(t *testing.T) {
edges := []inventory.Edge{
{From: "build-agent", To: "mesh-controller", Kind: inventory.EdgePackages},
{From: "build-agent", To: "mesh-controller", Kind: inventory.EdgeWorkerOf},
{From: "mesh-controller", To: "build-agent", Kind: inventory.EdgeBuiltBy},
{From: "route-proxy", To: "mesh-controller", Kind: inventory.EdgePackages},
{From: "route-proxy", To: "build-agent", Kind: inventory.EdgeBuiltBy},
}
set := reachableFrom([]string{"mesh-controller"}, edges)
if len(set) != 3 {
t.Fatalf("the controller, what packages it, and nothing more: %v", set)
}
tiers := tiersOf(set, edges)
pos := map[string]int{}
for i, tier := range tiers {
for _, m := range tier {
pos[m] = i
}
}
if pos["mesh-controller"] != 0 {
t.Fatalf("the controller first, built by the build machine that is running: %v", tiers)
}
if pos["build-agent"] <= pos["mesh-controller"] {
t.Fatalf("the build machine after the controller that defines its worker: %v", tiers)
}
if pos["route-proxy"] <= pos["build-agent"] {
t.Fatalf("what the build machine builds comes after it: %v", tiers)
}
if hasCycle(tiers, edges) {
t.Fatalf("no cycle here: %v", tiers)
}
}
@@ -234,13 +234,3 @@ func admitsSubject(pattern, subject []string) bool {
}
return len(pattern) == len(subject)
}
// The two packages name the runtime module separately — the broker's types stay free of the
// catalogue's on purpose — so this is what holds them to one string. A rename that reached only one
// side would compose a runtime principal for a module nobody assigns, silently, and leave the one
// that is assigned with a module's own grants.
func TestTheBrokerAndTheCatalogueAgreeOnTheRuntimeModule(t *testing.T) {
if RuntimeModule != catalogue.RuntimeModule {
t.Fatalf("the broker calls the runtime %q and the catalogue %q", RuntimeModule, catalogue.RuntimeModule)
}
}
+15 -16
View File
@@ -144,18 +144,12 @@ func ConsumerFor(p Principal) (Consumer, bool) {
}, true
}
// HolderConsumerFor is the worker a seat's holders share on that seat's work queue.
// HolderConsumerFor is the worker a seat's holder gets on that seat's work queue.
//
// **One worker for every holder, and each holder pulls one ask when it is idle** (novox/hq ADR
// 0190). The seat is *authority* — who may be the telegram sender — and the worker is *delivery*,
// kept separate so that relaxing one changes nothing about the other: a node-scoped seat has a
// holder per machine, and all of them take from this one consumer, so the work is shared without
// any holder knowing about the others. Pulled rather than pushed because a push consumer hands the
// next ask to whichever subscriber the server picks, busy or not, and a pulled one is asked for by
// a holder that has just become free. Which is also what ends the race issue 186 describes — asks
// delivered behind the one being worked, expiring unacknowledged and dropped after the fifth
// redelivery: nothing is delivered that nobody asked for. A long build keeps its own ask alive
// (stillWorking); the ack wait is for a holder that died.
// **A queue group even though the seat guarantees one holder.** The seat is *authority* — who may
// be the telegram sender — and the queue group is *delivery*. Tie delivery to the seat and the
// day somebody allows two holders for throughput, every message is processed twice with nothing
// reporting it. Kept separate, relaxing one changes nothing about the other.
func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) {
if len(seat.Accepts) == 0 {
return Consumer{}, false
@@ -164,13 +158,18 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
Name: "SEAT_" + upperSnake(seat.Name) + "_worker",
Stream: seatStreamName(seat.Name),
Filters: []string{"mesh.seat." + seat.Name + ".accept.>"},
Queue: "holders",
AckWaitSeconds: 60,
MaxDeliver: 5,
// As many in flight as there are holders working, which pulling bounds by itself: a holder
// fetches one and fetches again only after it acknowledged. The server's default stands.
Why: fmt.Sprintf("%s on %s holds %s; every holder pulls one ask at a time from this worker "+
"and acknowledges after the work is done, so a crash mid-work redelivers rather than "+
"loses and an idle holder is the one that takes the next ask", module, node, seat.Name),
// **One in flight.** A holder works one ask at a time, so the server hands it one at a
// time: with the default of many, every ask behind the one being worked was delivered,
// left unacknowledged for the length of the work, redelivered after the ack wait, and
// after the fifth time dropped — on 2026-10-01 twenty-six of forty-three builds asked in
// two minutes were never built, and the queue read as empty (novox/hq issue 186).
MaxAckPending: 1,
Why: fmt.Sprintf("%s on %s holds %s; it acknowledges after the work is done, so a "+
"crash mid-work redelivers rather than loses; one in flight, so a queue of asks is a "+
"queue and not a race against the ack wait", module, node, seat.Name),
}, true
}
+12 -25
View File
@@ -88,20 +88,15 @@ func TestAModuleThatConsumesNothingGetsNoConsumer(t *testing.T) {
}
}
// The seat is authority and the worker is delivery (novox/hq ADR 0190): one worker per seat, shared
// by every holder and pulled from, so a second holder takes the next ask rather than a copy of the
// same one — which is what a queue group used to guard, and what pulling one durable gives outright.
func TestAHoldersWorkerIsOneSharedByItsHolders(t *testing.T) {
// The seat is authority and the queue group is delivery. Tie them together and the day somebody
// allows two holders, every message is processed twice with nothing reporting it.
func TestAHoldersWorkerUsesAQueueGroupAnyway(t *testing.T) {
c, ok := HolderConsumerFor("one", "telegram", telegramSeat())
if !ok {
t.Fatal("the holder of a seat with inbound work got no worker")
}
two, _ := HolderConsumerFor("two", "telegram", telegramSeat())
if c.Name != two.Name || c.Stream != two.Stream {
t.Fatal("two holders got two workers, so each would process every ask")
}
if c.Push || c.Queue != "" {
t.Fatal("the worker is pushed, so the server would hand an ask to a busy holder")
if c.Queue == "" {
t.Fatal("the worker is not in a queue group, so a second holder would double-process")
}
if c.Stream != "SEAT_TELEGRAM_SENDER" {
t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream)
@@ -159,23 +154,15 @@ func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) {
has(t, perms.Subscribe, c.Filters[0])
}
// Every holder of a seat shares one worker and pulls from it (novox/hq ADR 0190): no queue group
// and no delivery subject, because a push consumer hands the next ask to whichever subscriber the
// server picks, busy or not; and no cap of one in flight, because pulling bounds the asks in flight
// by the holders that are free — which is what ended the race of issue 186, where asks delivered
// behind the one being worked expired and were dropped.
func TestAHoldersWorkerIsPulledByEveryHolder(t *testing.T) {
c, found := HolderConsumerFor("anchor", "build-agent", DeclaredSeat{Name: "node-build-agent", Accepts: []string{"build"}})
// A holder works one ask at a time, so the server hands it one at a time (novox/hq issue 186):
// asks queued behind the one being worked wait in the stream rather than being delivered,
// left to expire and dropped after the fifth redelivery.
func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) {
c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}})
if !found {
t.Fatal("a seat that accepts work has no worker")
}
if c.Queue != "" || c.Push {
t.Fatalf("the worker is pushed (queue %q, push %v); a holder pulls when it is free", c.Queue, c.Push)
}
if c.MaxAckPending != 0 {
t.Fatalf("the worker caps asks in flight at %d; pulling bounds them by the holders working", c.MaxAckPending)
}
if c.Name != "SEAT_NODE_BUILD_AGENT_worker" || c.Stream != "SEAT_NODE_BUILD_AGENT" {
t.Fatalf("the worker is %s on %s; one per seat, shared by its holders", c.Name, c.Stream)
if c.MaxAckPending != 1 {
t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending)
}
}
-37
View File
@@ -214,43 +214,6 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); {
case err == nil:
// **The controller owns the worker's shape, type included** (novox/hq issue 206). A holder
// built for a pull worker cannot bind a push one — `cannot pull subscribe to push based
// consumer` — and on 2026-10-03 the build machine rolled before the controller that would
// have redefined its worker, restarted on that for an hour, and nothing could build the
// controller that would have ended it. The server cannot change a consumer's type in place,
// so one of the wrong type is re-made: on a work queue nothing is lost, because what was
// acknowledged is gone from the stream and what was not is delivered again from the start.
// On any other stream a re-made consumer would replay what this one acknowledged (issue
// 156), so there it is said and left, and the person re-makes it knowing the cost.
if havePush, wantPush := have.Config.DeliverSubject != "", want.DeliverSubject != ""; havePush != wantPush {
shape := func(push bool) string {
if push {
return "push"
}
return "pull"
}
info, err := j.js.StreamInfo(c.Stream)
if err != nil {
return fmt.Errorf("asking about stream %s to re-make consumer %s: %w", c.Stream, c.Name, err)
}
if info.Config.Retention != nats.WorkQueuePolicy {
j.note("consumer %s on %s is %s and should be %s; not re-made, because %s keeps its history "+
"and a re-made consumer replays what this one acknowledged (novox/hq issue 156). Re-make it by hand",
c.Name, c.Stream, shape(havePush), shape(wantPush), c.Stream)
return nil
}
j.note("consumer %s on %s changes from %s to %s delivery: re-made where it left off, nothing "+
"acknowledged comes back and nothing pending is lost (novox/hq issue 206); a holder bound to "+
"the old shape binds again", c.Name, c.Stream, shape(havePush), shape(wantPush))
if err := j.js.DeleteConsumer(c.Stream, c.Name); err != nil {
return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err)
}
if _, err := j.js.AddConsumer(c.Stream, want); err != nil {
return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err)
}
return nil
}
// Where an existing consumer starts is its history, not something an assertion may move:
// the server refuses a changed deliver policy outright. Carried across, so asserting twice
// is the no-op a restart depends on.
-34
View File
@@ -73,37 +73,3 @@ func TestAnAccountMayReadItsOwnMembershipAndNoOthers(t *testing.T) {
has(t, perms.Publish, "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.anchor.postgres")
hasNot(t, perms.Subscribe, "mesh.assignment.>")
}
// The runtime arriving on a machine changes nothing about what each module is issued (to-be 38 WP2):
// the memberships are composed as before and the runtime reads several of them. What the machine's
// user list gains is one runtime principal, and loses nothing but the runtime module's own.
func TestTheRuntimeArrivingLeavesEveryMembershipAsItWas(t *testing.T) {
filter := Seat{Name: "node-packet-filter", Scope: "node", Serves: []string{"rules", "reload"}}
three := []Declared{
{Module: "nftables", Holds: []Seat{filter}, Serves: []string{"firewall_rules"}},
{Module: "zsh", Serves: []string{"execute"}},
{Module: "systemd", Serves: []string{"units"}},
}
before := Records{Nodes: []string{"anchor"}, Assigned: map[string][]Declared{"anchor": three}}
after := Records{Nodes: []string{"anchor"}, Assigned: map[string][]Declared{
"anchor": append(append([]Declared{}, three...), Declared{Module: RuntimeModule}),
}}
for _, d := range three {
was := MembershipFor("anchor", d, PlacementsOf(before, nil))
is := MembershipFor("anchor", d, PlacementsOf(after, nil))
if !reflect.DeepEqual(was, is) {
t.Errorf("%s's membership changed when the runtime arrived:\n%+v\n%+v", d.Module, was, is)
}
}
users, err := Users(after)
if err != nil {
t.Fatal(err)
}
kinds := map[Kind]int{}
for _, p := range users {
kinds[p.Kind]++
}
if kinds[KindNodeTools] != 1 || kinds[KindModule] != 3 || kinds[KindNode] != 1 || kinds[KindController] != 1 {
t.Errorf("the machine's users are %v; one runtime, the three modules, the host and the controller", kinds)
}
}
+10 -102
View File
@@ -34,20 +34,8 @@ const (
// authority is a list of tools and nothing else — not control, not declarations, not builds,
// and no ability to answer anything, because a person asks.
KindPerson Kind = "person"
// KindNodeTools is a machine's tool runtime (novox/hq ADR 0175, to-be 38): one process per
// node, on the host side, serving every assigned module's tools and every held seat's verbs.
// Its authority is the union of what the modules it carries would each have had for their
// tools — and nothing of what they consume, because tools are what it runs, not reactions.
KindNodeTools Kind = "node-tools"
)
// RuntimeModule is the module that IS the node's tool runtime (novox/hq ADR 0175). Where it is
// assigned, the mesh composes one runtime principal for the machine in place of that module's own,
// and the per-module containers that served tools until then stop being the way tools reach a node.
// Mirrored in the catalogue package, which the agreement test holds to the same string; one
// constant, so a rename is one edit and the two packages cannot drift.
const RuntimeModule = "node-tools"
// Seat is a role on the bus as a principal relates to it: the subjects it accepts, and those it
// emits (novox/hq ADR 0118, design 29 §5).
type Seat struct {
@@ -86,13 +74,6 @@ type Principal struct {
// a namespace no such module owns. Every service started and the graph stayed empty.
Watches []Seat
// Carries are the modules whose tools this principal serves, for a KindNodeTools principal
// (novox/hq ADR 0175): every module assigned to its node, as each declares itself. Its
// serving authority is the union of theirs — each module's own tool namespace and each held
// seat's verbs on this node — derived from the same declarations the modules' own principals
// are, so the runtime can serve nothing a module could not have served for itself.
Carries []Declared
// Invokes are the tools this principal may call, as `<module>.<tool>`; a single `*` is every
// tool. A person's whole authority (design 25 §7), and a module's only if its manifest says so
// (novox/hq ADR 0152) — the console's does, and nothing else's.
@@ -110,13 +91,10 @@ type Principal struct {
PasswordHash string
}
// seatsTheControllerAsks are the roles the mesh's own flows submit work to. Named rather than
// meshSeatsTheControllerUses are the roles the mesh's own flows submit work to. Named rather than
// derived from the seat set: the controller is not a module and declares no `uses`, so its side of a
// seat has to be stated, and a list is what makes "which roles does the mesh itself talk to" answerable.
// Both build roles while the handover runs (novox/hq ADR 0190): the controller asks whichever has a
// holder, and the retired one has one until build-agent replaces the builder. The second entry
// goes with the retired seat row.
var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"}
var meshSeatsTheControllerUses = []string{"mesh-build-machine"}
// enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the
// controller may answer an enrolment is derived from the same constant the user is named from.
@@ -134,10 +112,7 @@ func (p Principal) Username() string {
switch p.Kind {
case KindPerson:
return "person." + p.Module
case KindModule, KindNodeTools:
// The runtime is named exactly as the module it stands for would have been: the mesh
// issues its credential through the same path a module's takes (`module issue`), and
// that path knows the node and the module, not the kind.
case KindModule:
return p.Node + "." + p.Module
case KindNode:
return "node." + p.Node
@@ -211,9 +186,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
// Work the mesh's own flows submit to a role, and the outcomes they wait on (ADR 0121). A
// build is the one today: the controller asks, and reads the answer from the seat's event
// like the catalogue does — which is why no holder needs to publish into anybody's inbox.
// A node-scoped seat's work subject carries no node (novox/hq ADR 0190): the ask goes to
// the role, and whichever machine holding it is idle takes it.
for _, seat := range seatsTheControllerAsks {
for _, seat := range meshSeatsTheControllerUses {
pub = append(pub, "mesh.seat."+seat+".accept.>")
}
// **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own
@@ -373,17 +346,13 @@ func PermissionsFor(p Principal) (Permissions, error) {
// 3. Seats it holds: full participation.
for _, s := range p.Holds {
// Taking work from the role's queue: the worker consumer every holder shares (asked
// about, pulled from, acknowledged), on the seat's own stream (novox/hq ADR 0190). A
// holder pulls — asks the consumer for its next message, answered on its own inbox —
// so what it needs is MSG.NEXT on that worker and nothing delivered to it. The first
// machine to take work over the new bus was refused the asking (2026-09-28).
// 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)
pub = append(pub,
"$JS.API.CONSUMER.INFO."+stream+"."+worker,
"$JS.API.CONSUMER.MSG.NEXT."+stream+"."+worker,
"$JS.ACK."+stream+"."+worker+".>")
sub = append(sub, "_DELIVER."+worker, "_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))
}
@@ -406,48 +375,6 @@ func PermissionsFor(p Principal) (Permissions, error) {
pub = append(pub, seatToolSubject(s, t, "*"))
}
}
case KindNodeTools:
// **One process serves what every module on the machine would have served for itself**
// (novox/hq ADR 0175). Each carried module's whole tool namespace — the same grant that
// module's own principal has, for the same reason: the tools a module serves are what its
// code answers, and a list here would be a second copy of it. Each held seat's verbs on
// this node, as the holder's own principal would be granted them.
for _, d := range p.Carries {
if !safeSubject.MatchString(d.Module) {
return Permissions{}, fmt.Errorf(
"%q cannot be part of a subject: a permission is a subject pattern, and this would widen it", d.Module)
}
own := "mesh.mod." + d.Module
sub = append(sub, own+".tool.>")
// A tool that emits an event is the module's code and emits under the module's name
// (ADR 0042); the runtime carrying that code may publish what the module declared it
// emits, and nothing it did not.
for _, e := range d.Emits {
pub = append(pub, own+".event."+e)
}
for _, s := range d.Holds {
for _, t := range s.Serves {
sub = append(sub, seatToolSubject(s, t, p.Node))
}
}
}
// Every assigned module's membership on this node (ADR 0160): one per module, read
// directly from the stream and followed live. This node's and no other's — the one token
// that varies is the module, so the pattern is the machine's own assignments.
sub = append(sub, "mesh.assignment."+p.Node+".*")
pub = append(pub, "$JS.API.DIRECT.GET."+AssignmentsStream+".mesh.assignment."+p.Node+".*")
// And every tool on the mesh (ADR 0175, decision 5): any node may call any tool on any
// node, as the console already could — the runtime is the console's serving mode.
invoked, err := invokedSubjects([]string{"*"})
if err != nil {
return Permissions{}, err
}
pub = append(pub, invoked...)
// Nothing about consumers: it consumes nothing. A module's reactions to events are its
// own long-lived process, which ADR 0175 leaves where it is; what moves here is tools.
sub = unique(sub)
pub = unique(pub)
}
if p.Kind == KindPerson {
@@ -455,11 +382,6 @@ func PermissionsFor(p Principal) (Permissions, error) {
// consumer, because nothing is delivered to a person — they ask and are answered.
sub = append(sub, p.inbox())
}
if p.Kind == KindNodeTools {
// Its reply space, so the answers to what its tools call come back to it. No ack subject
// for the same reason a person has none: nothing is delivered to it.
sub = append(sub, p.inbox())
}
if p.Kind == KindModule || p.Kind == KindNode || p.Kind == KindController {
// Its own reply space, and nothing wider.
@@ -481,7 +403,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
// A module answers what it was asked — a tool call reaches it on its own namespace, so the
// authority is bounded by having been asked — and so does the controller. A node and a
// person are never asked anything, and are granted nothing here.
AllowResponses: p.Kind == KindModule || p.Kind == KindController || p.Kind == KindNodeTools,
AllowResponses: p.Kind == KindModule || p.Kind == KindController,
}, nil
}
@@ -703,20 +625,6 @@ func ComposeAccounts(principals []Principal) (string, error) {
return b.String(), nil
}
// unique is a sorted list with each subject once. Two carried modules holding seats with the same
// verb, or the runtime module itself carried beside the others, would otherwise write a grant twice
// — harmless to the server, and noise in a file that is read as the mesh's authority model.
func unique(values []string) []string {
sort.Strings(values)
out := values[:0]
for i, v := range values {
if i == 0 || v != values[i-1] {
out = append(out, v)
}
}
return out
}
func quoted(values []string) string {
if len(values) == 0 {
return ""
-91
View File
@@ -371,94 +371,3 @@ func TestAModulePullsItsOwnConsumerAndNoOthers(t *testing.T) {
}
}
}
// The runtime's authority is the union of what the modules it carries would have been granted for
// their tools (novox/hq ADR 0175): every carried module's tool namespace, every held seat's verbs
// on this node, every module's membership on this node, and a call to anything. Nothing it
// consumes, because it reacts to nothing.
func TestTheRuntimeServesTheUnionAndConsumesNothing(t *testing.T) {
filter := Seat{Name: "node-packet-filter", Scope: "node", Serves: []string{"rules", "reload"}}
p := Principal{Kind: KindNodeTools, Node: "anchor", Module: RuntimeModule, Carries: []Declared{
{Module: "nftables", Holds: []Seat{filter}, Serves: []string{"firewall_rules"}},
{Module: "zsh", Emits: []string{"shell.opened"}, Consumes: []string{"shop.order.placed"}},
{Module: RuntimeModule},
}}
perms, err := PermissionsFor(p)
if err != nil {
t.Fatal(err)
}
for _, want := range []string{
"mesh.mod.nftables.tool.>", "mesh.mod.zsh.tool.>", "mesh.mod." + RuntimeModule + ".tool.>",
"mesh.seat.node-packet-filter.tool.rules.anchor", "mesh.seat.node-packet-filter.tool.reload.anchor",
"mesh.assignment.anchor.*",
"_INBOX.anchor." + RuntimeModule + ".>",
} {
if !contains(perms.Subscribe, want) {
t.Errorf("the runtime may not subscribe %s: %v", want, perms.Subscribe)
}
}
for _, want := range []string{
"mesh.mod.*.tool.>", "mesh.seat.*.tool.>",
"$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.anchor.*",
"mesh.mod.zsh.event.shell.opened",
} {
if !contains(perms.Publish, want) {
t.Errorf("the runtime may not publish %s: %v", want, perms.Publish)
}
}
// Nothing of what a carried module consumes, and no consumer of its own to ack.
for _, s := range perms.Subscribe {
if strings.Contains(s, ".event.") || strings.HasPrefix(s, "_DELIVER.") {
t.Errorf("the runtime was granted a delivery it has no consumer for: %s", s)
}
}
for _, s := range perms.Publish {
if strings.HasPrefix(s, "$JS.ACK.") || strings.Contains(s, "CONSUMER") {
t.Errorf("the runtime was granted a consumer's subject and has no consumer: %s", s)
}
}
if !perms.AllowResponses {
t.Error("the runtime answers what it is asked, and may not reply")
}
if _, needed := ConsumerFor(p); needed {
t.Error("a consumer would be made for the runtime, which consumes nothing")
}
// Each subject once: the file is read as the mesh's authority model.
seen := map[string]bool{}
for _, s := range append(append([]string{}, perms.Subscribe...), perms.Publish...) {
if seen[s] {
t.Errorf("%s is granted twice", s)
}
seen[s] = true
}
}
func contains(list []string, want string) bool {
for _, s := range list {
if s == want {
return true
}
}
return false
}
// A node-scoped seat's work is shared (novox/hq ADR 0190): its holder on any machine subscribes the
// seat's one work subject, with no node in it, so holders on several machines read one queue. The
// node token belongs to a seat's tools, which are asked of one machine (design 33 §4), not to its work.
func TestANodeSeatsWorkSubjectCarriesNoNode(t *testing.T) {
seat := Seat{Name: "node-build-agent", Scope: "node", Accepts: []string{"build"}, Serves: []string{"status"}}
perms, err := PermissionsFor(Principal{Kind: KindModule, Node: "anchor", Module: "build-agent", Holds: []Seat{seat}})
if err != nil {
t.Fatal(err)
}
has(t, perms.Subscribe, "mesh.seat.node-build-agent.accept.build")
hasNot(t, perms.Subscribe, "mesh.seat.node-build-agent.accept.build.anchor")
// And its tools still carry the machine.
has(t, perms.Subscribe, "mesh.seat.node-build-agent.tool.status.anchor")
// The controller asks the role, not a machine.
controller, err := PermissionsFor(Principal{Kind: KindController})
if err != nil {
t.Fatal(err)
}
has(t, controller.Publish, "mesh.seat.node-build-agent.accept.>")
}
+1 -6
View File
@@ -217,15 +217,10 @@ var ControllerFollows = []string{
// A build's outcome, which is the build-machine role's own event now (ADR 0121) rather than a
// message on the control branch. Same three audiences, one publish: whoever asked, this, and the
// catalogue.
seatEventSubject("node-build-agent", "built"),
seatEventSubject("mesh-build-machine", "built"),
// The forge's merges: what moved a source, so the mesh builds what that source produces
// without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name.
moduleEventSubject("gitea", "pull.merged"),
// The retired build role's outcome too, while the handover runs (novox/hq ADR 0190): the one
// build machine keeps answering on its seat until build-agent replaces it, and the outcome that
// registers build-agent itself comes from there. Appended, for the same reason as above; goes
// with the retired seat row.
seatEventSubject("mesh-build-machine", "built"),
}
// moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so
+4 -4
View File
@@ -24,8 +24,8 @@ accounts {
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.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] }
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] }
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused"] }
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>"] }
allow_responses: { max: 1, ttl: "1m" }
} }
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
@@ -37,8 +37,8 @@ accounts {
subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] }
} }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "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.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
allow_responses: { max: 1, ttl: "1m" }
} }
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
-20
View File
@@ -62,33 +62,13 @@ func Users(r Records) ([]Principal, error) {
for _, node := range sortedCopy(r.Nodes) {
out = append(out, Principal{Kind: KindNode, Node: node})
// **Where the runtime is assigned, the machine gets one runtime principal in place of the
// runtime module's own** (novox/hq ADR 0175, to-be 38). It carries every module on the
// node: its serving grants are the union of theirs. Every other module keeps its own
// principal — a module still serving tools from its own container holds its own
// credential until it moves, and the two serve side by side in the meantime.
runtimeHere := false
for _, d := range r.Assigned[node] {
if d.Module == RuntimeModule {
runtimeHere = true
}
}
for _, d := range r.Assigned[node] {
if runtimeHere && d.Module == RuntimeModule {
continue
}
out = append(out, Principal{
Kind: KindModule, Node: node, Module: d.Module,
Emits: d.Emits, Consumes: d.Consumes, Serves: d.Serves,
Holds: d.Holds, Uses: d.Uses, Watches: d.Watches, Invokes: d.Invokes,
})
}
if runtimeHere {
out = append(out, Principal{
Kind: KindNodeTools, Node: node, Module: RuntimeModule,
Carries: append([]Declared(nil), r.Assigned[node]...),
})
}
}
for _, node := range sortedCopy(r.Enrolling) {
out = append(out, Principal{Kind: KindEnrolment, Node: node})
-51
View File
@@ -245,54 +245,3 @@ func TestAUserListIsComposedBeforeAnythingMovesOntoTheBus(t *testing.T) {
t.Errorf("the composed list does not contain the machine running the bus")
}
}
// Where the runtime module is assigned, the machine gets one runtime principal in place of the
// runtime module's own (novox/hq ADR 0175, to-be 38). Every other module keeps its own: a module
// still serving tools from its own container holds its own credential until it moves.
func TestTheRuntimeModuleBecomesTheMachinesRuntimePrincipal(t *testing.T) {
r := someRecords()
r.Assigned["one"] = append(r.Assigned["one"], Declared{Module: RuntimeModule})
users, err := Users(r)
if err != nil {
t.Fatal(err)
}
var runtime *Principal
for i := range users {
p := &users[i]
if p.Node == "one" && p.Module == RuntimeModule {
if p.Kind == KindModule {
t.Fatalf("%s on one was composed as an ordinary module beside the runtime", RuntimeModule)
}
runtime = p
}
}
if runtime == nil || runtime.Kind != KindNodeTools {
t.Fatalf("one runs %s and got no runtime principal: %v", RuntimeModule, namesOf(t, r))
}
if runtime.Username() != "one."+RuntimeModule {
t.Errorf("the runtime is named %q; `module issue` names it as the module it stands for", runtime.Username())
}
carried := map[string]bool{}
for _, d := range runtime.Carries {
carried[d.Module] = true
}
if !carried["telegram"] || !carried[RuntimeModule] {
t.Errorf("the runtime carries %v; it carries every module on its node", carried)
}
// And the other node, where the runtime is not assigned, is exactly as before.
for _, p := range users {
if p.Node == "two" && p.Kind == KindNodeTools {
t.Fatal("two runs no runtime and was given a runtime principal")
}
}
// A module serving its own tools beside the runtime keeps its own principal.
found := false
for _, p := range users {
if p.Kind == KindModule && p.Node == "one" && p.Module == "telegram" {
found = true
}
}
if !found {
t.Error("telegram lost its own principal when the runtime arrived on its node")
}
}
-145
View File
@@ -1,145 +0,0 @@
package broker
import (
"os"
"testing"
"time"
"github.com/nats-io/nats.go"
)
// A seat's worker that changed from push to pull delivery strands a holder built for the new shape
// (novox/hq issue 206): the server refuses a pull subscription on a push consumer, and the controller
// that would redefine it was the build that nobody could take. The controller owns the worker's
// shape, type included: on a work queue it re-makes one of the wrong type, losing nothing, and a
// pull subscription then binds and takes what was pending.
//
// docker run -d --rm --name t -p 14231:4222 nats:2.10-alpine -js
// MESH_TEST_NATS=nats://127.0.0.1:14231 go test ./internal/broker/ -run TestAWorker
func TestAWorkerOfTheWrongTypeIsRemadeOnAWorkQueueAndAPullThenBinds(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
js, err := Dial(url)
if err != nil {
t.Fatal(err)
}
defer js.Close()
const stream, worker, filter = "SEAT_T_SHELF", "SEAT_T_SHELF_worker", "mesh.seat.t-shelf.accept.>"
_ = js.js.DeleteStream(stream)
if _, err := js.js.AddStream(&nats.StreamConfig{
Name: stream, Subjects: []string{filter}, Retention: nats.WorkQueuePolicy, Storage: nats.MemoryStorage,
}); err != nil {
t.Fatal(err)
}
defer func() { _ = js.js.DeleteStream(stream) }()
// The worker as the previous controller defined it: push, in a queue group.
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second, MaxDeliver: 5,
FilterSubject: filter, DeliverSubject: "_DELIVER." + worker, DeliverGroup: "holders",
}); err != nil {
t.Fatal(err)
}
for _, body := range []string{"one", "two", "three"} {
if _, err := js.js.Publish("mesh.seat.t-shelf.accept.build", []byte(body)); err != nil {
t.Fatal(err)
}
}
// The old holder took and acknowledged the first ask, then went away.
old, err := js.js.QueueSubscribeSync(filter, "holders", nats.Bind(stream, worker))
if err != nil {
t.Fatal(err)
}
m, err := old.NextMsg(twoSeconds)
if err != nil {
t.Fatal(err)
}
if string(m.Data) != "one" {
t.Fatalf("the first ask is %q", m.Data)
}
if err := m.AckSync(); err != nil {
t.Fatal(err)
}
if err := old.Unsubscribe(); err != nil {
t.Fatal(err)
}
// The new controller asserts the worker as the mesh derives it now: pull.
if err := js.EnsureConsumer(Consumer{
Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5,
Why: "the test's worker",
}); err != nil {
t.Fatal(err)
}
have, err := js.js.ConsumerInfo(stream, worker)
if err != nil {
t.Fatal(err)
}
if have.Config.DeliverSubject != "" || have.Config.DeliverGroup != "" {
t.Fatalf("the worker is still push: %+v", have.Config)
}
// A holder built for the new shape binds, and takes exactly what the old one left.
sub, err := js.js.PullSubscribe(filter, worker, nats.Bind(stream, worker), nats.ManualAck())
if err != nil {
t.Fatalf("a pull subscription does not bind the re-made worker: %v", err)
}
got, err := sub.Fetch(3, nats.MaxWait(twoSeconds))
if err != nil && len(got) == 0 {
t.Fatalf("nothing pending was delivered: %v", err)
}
var bodies []string
for _, g := range got {
bodies = append(bodies, string(g.Data))
_ = g.Ack()
}
if len(bodies) != 2 || bodies[0] != "two" || bodies[1] != "three" {
t.Fatalf("the pending asks after the acknowledged one, in order: %v", bodies)
}
// Asserted again, the pull worker is the no-op a restart depends on.
if err := js.EnsureConsumer(Consumer{
Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5,
}); err != nil {
t.Fatal(err)
}
}
// On a stream that keeps its history, a worker of the wrong type is said and left: re-making it would
// replay what it acknowledged (novox/hq issue 156), and that is a person's call.
func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
js, err := Dial(url)
if err != nil {
t.Fatal(err)
}
defer js.Close()
const stream, worker, filter = "EVENTS_T", "EVENTS_T_reader", "mesh.t.event.>"
_ = js.js.DeleteStream(stream)
if _, err := js.js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{filter}, Storage: nats.MemoryStorage}); err != nil {
t.Fatal(err)
}
defer func() { _ = js.js.DeleteStream(stream) }()
if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{
Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second,
FilterSubject: filter, DeliverSubject: "_DELIVER." + worker,
}); err != nil {
t.Fatal(err)
}
if err := js.EnsureConsumer(Consumer{Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60}); err != nil {
t.Fatal(err)
}
have, err := js.js.ConsumerInfo(stream, worker)
if err != nil {
t.Fatal(err)
}
if have.Config.DeliverSubject == "" {
t.Fatal("a history stream's consumer was re-made, which replays what it acknowledged")
}
}
-20
View File
@@ -968,26 +968,6 @@ func compile(ctx context.Context, run Runner, tree string, chain Toolchain,
if _, err := run(ctx, tree, "docker", invocation...); err != nil {
return "", err
}
if chain.Dependencies != "" {
// **What the bundle runs with, from the image it was compiled in** (Toolchain.Dependencies).
// A second run in the same image rather than a shell wrapped around the compiler: the
// compile line stays a plain command a reader can run by hand, and the copy is one more
// plain command beside it. Refused by name when the image carries no such directory — an
// older toolchain image — because a bundle packed without its dependencies starts nowhere
// and says so three layers away from here.
copying := []string{
"run", "--rm",
"--volume", tree + ":" + within,
"--workdir", within,
base,
"sh", "-c",
`test -d "$1" || { echo "the toolchain image carries no $1: it predates the mesh shipping a bundle's dependencies, rebuild $2 first" >&2; exit 1; }; cp -a "$1/." "$3/"`,
"dependencies", chain.Dependencies, chain.Base, out,
}
if _, err := run(ctx, tree, "docker", copying...); err != nil {
return "", fmt.Errorf("copying the %s dependencies a bundle runs with: %w", chain.Language, err)
}
}
return filepath.Join(tree, out), nil
}
+1 -40
View File
@@ -82,41 +82,6 @@ func TestABundleIsCompiledAndPackedWithNoDockerfile(t *testing.T) {
if !strings.HasPrefix(digest, "sha256:") {
t.Fatalf("the bundle was not pinned: %v", got.Manifest.Resources[0])
}
// **And what it runs with, from the image it was compiled in** (novox/hq to-be 38 WP3). A
// second run in the same toolchain image copies the toolchain's runtime directory — the
// `"type": "module"` package.json and the pruned node_modules — into the output's root, and
// refuses by name when the image carries none rather than packing a bundle that starts nowhere.
var copied string
for _, line := range r.ran {
if strings.HasPrefix(line, "docker run") && strings.Contains(line, "/app/runtime") {
copied = line
}
}
if copied == "" {
t.Fatalf("the bundle's dependencies were not copied in after the compile:\n%s", strings.Join(r.ran, "\n"))
}
if !strings.Contains(copied, "mesh-tools/build@sha256:") || !strings.Contains(copied, "predates") ||
!strings.Contains(copied, Out("code")) {
t.Fatalf("the copy does not run in the same toolchain, refuse an older image by name, or land in the artifact's output: %s", copied)
}
if strings.Index(strings.Join(r.ran, "\n"), "--outDir") > strings.Index(strings.Join(r.ran, "\n"), "/app/runtime") {
t.Fatal("the dependencies were copied before the compile wrote its output")
}
}
// A language whose bundle carries its own dependencies copies nothing in: a Go binary is static.
func TestOnlyALanguageWithARuntimeDirectoryCopiesDependenciesIn(t *testing.T) {
ts, _ := ToolchainFor("typescript")
if ts.Dependencies != "/app/runtime" {
t.Fatalf("typescript bundles run with %q", ts.Dependencies)
}
for _, language := range []string{"go", "python"} {
chain, _ := ToolchainFor(language)
if chain.Dependencies != "" {
t.Fatalf("%s copies %q into every bundle, and its bundles carry their own", language, chain.Dependencies)
}
}
}
// **Refused before anything is built, naming what to build first.** A base the mesh has not built
@@ -180,13 +145,9 @@ func TestTwoBundlesInOneModuleArePackedSeparately(t *testing.T) {
t.Fatalf("a module with two bundles did not build: %v", err)
}
// Compiled into two different places. Only the compile lines: the copy of each bundle's
// dependencies names the same directory again, deliberately.
// Compiled into two different places.
var outputs []string
for _, line := range r.ran {
if !strings.Contains(line, "--outDir") {
continue
}
for _, part := range strings.Fields(line) {
if strings.HasPrefix(part, ".mesh-build/") {
outputs = append(outputs, part)
+1 -25
View File
@@ -57,23 +57,6 @@ type Toolchain struct {
// carrying its debug info. The mistake was believing a comment rather than reading the file it
// produced (novox/hq 04-ISSUES/161).
LinkerFlags []string
// Dependencies is a directory inside the toolchain image whose contents a bundle in this
// language runs with, copied whole into the compiled output's root after the compile.
//
// **A bundle that compiles is not yet a bundle that runs.** The compiler resolves `import
// "nats"` from the toolchain image's own node_modules and the pack takes only what the compiler
// wrote, so what a machine unpacked could not find a single dependency — and no TypeScript bundle
// had ever run live to show it (novox/hq to-be 38 WP3). For TypeScript the directory holds a
// `package.json` saying `"type": "module"` — Node reads a bare `.js` as CommonJS otherwise, so a
// bundle with its dependencies and without that line still fails to start — and the pruned,
// production-only node_modules the runtime itself ships with: the SDK's and the runtime's
// dependencies, and nothing module-specific yet (novox/hq ADR 0188 §5: a skeleton; a module's
// own npm dependencies are a later step). Empty for a language whose bundle carries its own —
// a Go binary is static, a Python bundle is installed with its dependencies.
//
// A toolchain image without the directory fails the build by name rather than packing a bundle
// that starts nowhere: the image predates this and must be rebuilt first.
Dependencies string
// SystemStamp is the variable this language's linker fills with the artifact's declared system,
// for a language whose binaries are pinned to one at link time (novox/hq ADR 0005).
//
@@ -124,21 +107,14 @@ var toolchains = []Toolchain{
// symlinks to a launcher that requires its library relatively — and the base image's own
// assembly resolves them away, leaving a launcher whose relative require points nowhere.
// Every module's hand-written Dockerfile had to know this. Now none of them does.
// **Rooted at the module, so an entrypoint lands where it is named.** Without a root the
// compiler takes the common directory of the files it is given: a module compiling only
// `tools/index.ts` had its output at `index.js`, and the entrypoint it declared —
// `tools/index.js`, "named as it will be found" — named a file the bundle did not
// contain. The runtime that loads bundles by their declared entrypoints (novox/hq ADR
// 0175) is what made this visible.
Compile: []string{
"node", "/app/node_modules/typescript/bin/tsc",
"--module", "NodeNext", "--moduleResolution", "NodeNext",
"--target", "ES2022", "--rootDir", ".",
"--target", "ES2022",
},
OutputFlag: "--outDir",
Unit: UnitSources,
SourceExt: ".ts",
Dependencies: "/app/runtime",
},
{
Language: "go",
+12 -4
View File
@@ -46,7 +46,7 @@ func boundUsed(content string) [][2]string {
// Three facts the mesh states about any provision, plus whatever the provider said it serves. A
// module may not reach a binding it does not have — the same boundary as a secret, for the same
// reason.
func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]string {
func knownFor(m Manifest, needs []Needed, node string) (map[string]map[string]string, error) {
out := map[string]map[string]string{}
for _, want := range m.Wants() {
for i := range needs {
@@ -54,12 +54,20 @@ func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]str
if n.Name != want || n.For != m.Module {
continue
}
as := ConsumerIdentity(node, IdentitySource(m.Slug, m.Module))
values := map[string]string{
"at": n.At,
"from": n.From,
"as": ConsumerIdentity(node, IdentitySource(m.Slug, m.Module)),
"as": as,
}
for key, value := range n.Serves {
// What the provider derives for this consumer rather than for all of them
// (novox/hq ADR 0188). Filled here, the one place a provision and the module
// requiring it are both in hand.
served, err := ServedTo(n.Serves, as)
if err != nil {
return nil, fmt.Errorf("%s requires %s: %w", m.Module, want, err)
}
for key, value := range served {
// The provider's own vocabulary. Rendered plainly: a port is 5432, not 5432.000000,
// which is what a float would write and what a connection string would refuse.
values[key] = plainly(value)
@@ -67,7 +75,7 @@ func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]str
out[want] = values
}
}
return out
return out, nil
}
// withOwnNames adds a module's own composed names to what it may name from one binding:
+1 -46
View File
@@ -67,36 +67,6 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
out := m
out.Build = nil
out.Resources = nil
// What the build compiled, kept on the resolved manifest (novox/hq ADR 0175): a tools bundle is
// named by no resource of the module's own — the node's runtime loads it — so this is the only
// place the mesh would otherwise not have it. In artifact order, so two resolutions of one
// build compare equal.
out.Bundles = nil
if m.Build != nil {
for _, a := range m.Build.Artifacts {
if a.Kind != ArtifactBundle {
continue
}
made := by[a.Name]
// What the runtime loads: what the artifact said, else every entrypoint of a module
// that declares tools, else nothing (the field's own rule; see Artifact.Loads).
loads := append([]string(nil), a.Loads...)
if a.Loads == nil && len(m.Tools) > 0 {
loads = append([]string(nil), a.Entrypoints...)
}
// **Kept, never routed** (ADR 0155): the builder publishes to the store at the address
// it reached it by, and a manifest carrying that address names an installation —
// registration refused node-tools for exactly this on 2026-10-02. The build record
// already keeps the store-relative form; the resolved manifest keeps the same, and
// composition routes it through the store a machine reaches (Routed).
out.Bundles = append(out.Bundles, Bundle{
Name: a.Name, Source: Recorded(made.Reference), Digest: made.Digest,
Language: a.Language, Entrypoints: append([]string(nil), a.Entrypoints...),
Loads: loads,
})
}
sort.Slice(out.Bundles, func(i, j int) bool { return out.Bundles[i].Name < out.Bundles[j].Name })
}
for _, r := range m.Resources {
named, _ := r["artifact"].(string)
if named == "" {
@@ -140,8 +110,7 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
// The same on the wire: both are bytes fetched by digest and unpacked. They differ in
// how they were made — one packed as it stood, the other compiled first — and a
// machine has no reason to care which.
// Kept, not routed, for the reason the bundles above are (ADR 0155).
filled["source"] = Recorded(artifact.Reference)
filled["source"] = artifact.Reference
filled["digest"] = artifact.Digest
// **And `${version}`, so a resource can name a place that is this build's alone**
// (novox/hq ADR 0141, 04-ISSUES/142). A component is unpacked into a directory named
@@ -222,20 +191,6 @@ func (b *Build) problems(module string) []string {
"%s: %q is a bundle and says no language, so nothing can choose a compiler "+
"for it", module, a.Name))
}
// What the runtime loads is among what was compiled (ADR 0175): a name here that is
// not an entrypoint is a file the bundle does not contain, and the runtime would
// fail to import it on every machine rather than here.
for _, load := range a.Loads {
found := false
for _, e := range a.Entrypoints {
found = found || e == load
}
if !found {
problems = append(problems, fmt.Sprintf(
"%s: %q says the runtime loads %q, which is not among its entrypoints — "+
"what is loaded is compiled, so it is named there too", module, a.Name, load))
}
}
// **A system, for a language that compiles to a binary** (novox/hq ADR 0142). A binary
// is pinned to one operating system at link time so a host refuses to touch a machine
// it was not built for (novox/hq ADR 0005); an artifact that says nothing would be
+290
View File
@@ -0,0 +1,290 @@
package catalogue
import (
"fmt"
"regexp"
"sort"
"strings"
)
// What a provider derives for one consumer, said once in the provider's definition and delivered
// to both ends (novox/hq ADR 0188, issue 124).
//
// A `serves` block is otherwise literal: the same values for every consumer. Where the provider
// *names the resource* — a bucket, a database, a vhost — the name is derived from who is asking,
// and before this the mesh had no channel for it. The provider recomputed it in its own code and
// every consumer transcribed it into its own definition by hand, which is a copy of somebody
// else's rule kept in agreement by nobody. One of three transcriptions was wrong for months.
//
// **The mesh learns no protocol here; it spells its own name in an alphabet it already knows.**
// The only fact a served value may name is the identity the mesh itself minted for the consumer,
// in one of two alphabets: as it was minted, and as a DNS label. Everything a provider wants
// around it — a prefix, a suffix, a separator — it writes around the placeholder, because a
// served value is a string.
// consumerFact is `${consumer:<fact>}` or `${consumer:<fact>:<alphabet>}`.
var consumerFact = regexp.MustCompile(`\$\{consumer:([a-z][a-z0-9-]*)(?::([a-z][a-z0-9-]*))?\}`)
// consumerFacts are what a served value may name about the consumer it is being derived for.
// One entry, deliberately: the identity is the one thing about a consumer the mesh itself chose,
// so it is the one thing the mesh can hand to a provider without either end guessing.
var consumerFacts = []string{"as"}
// consumerAlphabets are the ways the mesh will write that identity. `dns` is the mesh's own
// identifier with its separator written `-` instead of `_` — the whole of the difference between
// the alphabet the mesh mints in and the one buckets, vhosts and hostnames accept.
var consumerAlphabets = []string{"dns"}
// ServedTo fills a provider's served values for one consumer.
//
// `as` is the identity the mesh minted for that consumer — the same string it is told to present
// as a login. Values with no placeholder are returned exactly as they were, and a block with no
// placeholder at all is returned unchanged, so this costs nothing for the providers that derive
// nothing.
//
// Only strings carry placeholders. A number, a boolean or a nested object is a value the provider
// stated outright, and is left alone.
func ServedTo(serves map[string]any, as string) (map[string]any, error) {
if len(serves) == 0 {
return serves, nil
}
var out map[string]any
for _, key := range sortedAnyKeys(serves) {
text, ok := serves[key].(string)
if !ok || !strings.Contains(text, "${consumer:") {
continue
}
filled, err := consumerInto(text, as)
if err != nil {
return nil, fmt.Errorf("the value served as %q: %w", key, err)
}
if out == nil {
// Copied only once something actually changes: the caller's map is the manifest's,
// and a provider that derives nothing must not have it rewritten underneath it.
out = make(map[string]any, len(serves))
for k, v := range serves {
out[k] = v
}
}
out[key] = filled
}
if out == nil {
return serves, nil
}
return out, nil
}
// consumerInto replaces every `${consumer:…}` in one value.
//
// **A fact or an alphabet the mesh does not have is refused, not left standing.** Written through,
// the literal `${consumer:as}` would reach a configuration file and be read as a bucket name,
// failing somewhere that names neither the module nor the mesh — the same reasoning `${bound:…}`
// is refused by (boundInto).
func consumerInto(value, as string) (string, error) {
var failed error
out := consumerFact.ReplaceAllStringFunc(value, func(match string) string {
parts := consumerFact.FindStringSubmatch(match)
fact, alphabet := parts[1], parts[2]
if fact != "as" {
if failed == nil {
failed = fmt.Errorf(
"says %s, and the mesh states %s about a consumer", match, orNothing(consumerFacts))
}
return match
}
switch alphabet {
case "":
return as
case "dns":
return asDNSLabel(as)
default:
if failed == nil {
failed = fmt.Errorf(
"says %s, and the mesh writes an identity as %s", match, orNothing(consumerAlphabets))
}
return match
}
})
if failed != nil {
return "", failed
}
return out, nil
}
// asDNSLabel writes a minted identity as a DNS label.
//
// The mesh's identities are already lower-case letters, digits and `_` (ConsumerIdentity), and
// already short enough for the tightest backend they reach (CheckIdentity, twenty characters). So
// this is the separator and nothing else — no lower-casing of what is already lower case, no
// truncation to a limit the identity is already inside, no padding of a name that is already long
// enough. Each of those would be the mesh guessing at a rule it has not been given.
func asDNSLabel(as string) string {
return strings.ReplaceAll(as, "_", "-")
}
// CheckServes refuses a `serves` block that names a consumer fact or an alphabet the mesh does not
// have, when the definition is parsed rather than when a consumer is resolved.
//
// A provision nobody consumes yet still has its rule read: a definition that would be refused the
// first time somebody required it is a definition that is wrong now.
func CheckServes(m Manifest) []string {
var problems []string
for _, provision := range sortedServes(m.Serves) {
for _, key := range sortedAnyKeys(m.Serves[provision]) {
text, ok := m.Serves[provision][key].(string)
if !ok {
continue
}
// A probe identity, because what is checked is the shape of the statement and not
// what any consumer is called.
if _, err := consumerInto(text, "mesh_node_module"); err != nil {
problems = append(problems, fmt.Sprintf(
"%s serves %s, and the value it serves as %q %s", m.Module, provision, key, err))
}
}
}
return problems
}
func sortedServes(serves map[string]map[string]any) []string {
out := make([]string, 0, len(serves))
for k := range serves {
out = append(out, k)
}
sort.Strings(out)
return out
}
func sortedAnyKeys(values map[string]any) []string {
out := make([]string, 0, len(values))
for k := range values {
out = append(out, k)
}
sort.Strings(out)
return out
}
// derivedFor is what the provider on this machine derives for one consumer of one provision
// (novox/hq ADR 0188).
//
// Settled first, then derived: an operator may set a prefix on what the provider serves and the
// mesh still fills the consumer's half of it ([ADR 0174]). Only the keys that actually name the
// consumer are returned — the rest of a `serves` block is the same for every consumer and is
// already in the provider's own definition, so repeating it here would be a second copy to go
// stale.
//
// The first module in the resolved order that says it serves the provision answers, which is the
// choice servedOnThisMachine makes for the consumer's half. Nothing serving it on this machine is
// not an error: a contribution can reach a machine whose provider is a record or an adapter, and
// then there is nothing derived to tell.
func (r Resolution) derivedFor(provision, as string, settings SettingsBy) (map[string]any, error) {
for _, m := range r.Modules {
serves, said := m.Serves[provision]
if !said {
continue
}
var names map[string]any
for key, value := range serves {
if text, ok := value.(string); ok && strings.Contains(text, "${consumer:") {
if names == nil {
names = map[string]any{}
}
names[key] = value
}
}
if names == nil {
return nil, nil
}
settled, err := Settle(names, settings[m.Module])
if err != nil {
return nil, fmt.Errorf("%s serving %s: %w", m.Module, provision, err)
}
derived, err := ServedTo(settled, as)
if err != nil {
return nil, fmt.Errorf("%s serving %s to %s: %w", m.Module, provision, as, err)
}
return derived, nil
}
return nil, nil
}
// notTranscribed refuses a consumer's file that writes out the value its provider derives for it,
// instead of asking for it (novox/hq ADR 0188, issue 124).
//
// **What would have caught the one wrong instance.** The object store's three consumers each wrote
// their bucket into their own configuration by hand. One of them named a predecessor's bucket, and
// nothing compared it to what the provider would actually create: the module would have
// authenticated successfully and been refused on every object, which reads like a credential fault
// and is not one. It looked authoritative for months.
//
// The test is exact and costs one string search: a definition whose file already contains the
// value the mesh is about to derive for it has written down somebody else's rule. It cannot be a
// coincidence — a derived value carries the identity the mesh minted for this very consumer on
// this very machine, which nothing else would spell out — and it cannot be checked afterwards,
// because after substitution every consumer's file contains it legitimately.
//
// Only values that actually name the consumer are judged. A provider that serves a constant under
// the same key serves the same constant to everyone, and a consumer repeating it is redundant
// rather than wrong.
func notTranscribed(resource map[string]any, known map[string]map[string]string, module string) error {
if fmt.Sprint(resource["type"]) != "file" {
return nil
}
content, ok := resource["content"].(string)
if !ok || content == "" {
return nil
}
for _, provision := range sortedKnown(known) {
values := known[provision]
identity := values["as"]
if identity == "" {
continue
}
for _, key := range sortedStringKeys(values) {
if key == "as" {
// The login is not derived from itself, and a consumer that must present it in a
// connection string legitimately has it from `${bound:…}` — which is what it will
// be after substitution, so this would judge the substitution, not the module.
continue
}
value := values[key]
if value == "" || !namesTheConsumer(value, identity) {
continue
}
if !strings.Contains(content, value) {
continue
}
return fmt.Errorf(
"%s writes %q into %v, and that is exactly what %s derives for it — a definition "+
"keeping its own copy of somebody else's naming rule is one that can disagree "+
"with it, silently. Say ${bound:%s:%s} and be told",
module, value, resource["id"], provision, provision, key)
}
}
return nil
}
// namesTheConsumer is whether a derived value was built from this consumer's identity — in the
// alphabet it was minted in, or as a DNS label. A value that does not contain it was not derived
// from it, whatever else it may be.
func namesTheConsumer(value, identity string) bool {
return strings.Contains(value, identity) || strings.Contains(value, asDNSLabel(identity))
}
func sortedKnown(known map[string]map[string]string) []string {
out := make([]string, 0, len(known))
for k := range known {
out = append(out, k)
}
sort.Strings(out)
return out
}
func sortedStringKeys(values map[string]string) []string {
out := make([]string, 0, len(values))
for k := range values {
out = append(out, k)
}
sort.Strings(out)
return out
}
+51 -127
View File
@@ -508,27 +508,10 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
return nil, fmt.Errorf(
"%s needs a secret called %q and none was made for it", m.Module, name)
}
// The runtime's credential belongs to the account the runtime runs as (novox/hq ADR 0175,
// to-be 38 WP3): its process is composed `user: <account>` where the node has one, and a
// root-owned 0600 file is one that process cannot read. Composed here rather than said in
// the manifest, because a manifest cannot say ${machine:account} safely — a node with no
// account has nothing to resolve it to, and then the runtime runs as root and the file
// stays root's.
owner := m.SecretsOwner
if m.Module == RuntimeModule && r.Account != "" {
owner = r.Account
}
first = append(first, ownedBy(owner, map[string]any{
first = append(first, ownedBy(m.SecretsOwner, map[string]any{
"id": NeedID(name), "type": "file", "path": m.OwnSecrets[name].Path, "sealed": sealed,
}))
}
// This module's tools bundles, where the machine runs the node's tool runtime (novox/hq
// ADR 0175, to-be 38 WP2). Mesh-computed like everything above it, and before the module's
// own resources for the same reason: the runtime's process names the files inside these
// and is restarted when one changes, so they are on the machine before it is.
if r.runtimeHere() {
first = append(first, bundleArchives(m)...)
}
// Operator-owned paths this module is granted use of (novox/hq ADR 0051). Written before
// the module's own resources, and so before the container that mounts them: the host must
// find each present — refusing clearly if the operator has not provided it — before it
@@ -645,7 +628,16 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
if err != nil {
return nil, err
}
file, err := boundFile(*found, m.Binds[to], ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module)), own)
as := ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module))
// What the provider derives for THIS consumer, filled here where the consumer is
// known (novox/hq ADR 0188). The same fill knownFor does below, so the binding file
// and the module's `${bound:…}` substitutions cannot say different things.
told := *found
told.Serves, err = ServedTo(told.Serves, as)
if err != nil {
return nil, fmt.Errorf("%s is told about %s: %w", m.Module, to, err)
}
file, err := boundFile(told, m.Binds[to], as, own)
if err != nil {
return nil, err
}
@@ -711,7 +703,10 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
return nil, err
}
// And what its bindings say, for the half of a connection that is not secret.
known := knownFor(m, r.Needs, r.Node)
known, err := knownFor(m, r.Needs, r.Node)
if err != nil {
return nil, err
}
// A requirement answered on this same machine is not in r.Needs — its binding file is
// written from `here` (above) — and so `${bound:…}` could not name it, though the file
// beside it said the same facts. Filled from the same answer, so the two cannot disagree.
@@ -728,7 +723,11 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
}
local := *answered
local.For = m.Module
for provision, values := range knownFor(m, []Needed{local}, r.Node) {
here, err := knownFor(m, []Needed{local}, r.Node)
if err != nil {
return nil, err
}
for provision, values := range here {
known[provision] = values
}
}
@@ -753,6 +752,17 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
// And the machine underneath, which no binding of its own can tell it.
thisMachine := machineFacts(r, with.Names, with.MeshRange)
// **A definition that already holds the answer transcribed it** (novox/hq ADR 0188).
// Judged over what the module itself declares, and before anything is substituted: the
// mesh's own generated files — the binding, the contributions — legitimately carry the
// derived value, and after substitution so does every consumer's file, so this is the one
// moment the two can be told apart.
for _, own := range m.Resources {
if err := notTranscribed(own, known, m.Module); err != nil {
return nil, err
}
}
// Which of this module's files carry a secret, for the rule that a container may not read
// one of them as its environment without saying so (ADR 0086, issue 041).
secretFiles := secretFilesOf(resources)
@@ -872,15 +882,6 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
copied["reload-on"] = renamed
}
// **What reads one of this module's own secrets is restarted when it changes** (novox/hq
// issue 203, issue 206). A credential is re-issued by the mesh, and a container that
// mounted the old file keeps the old one open: the build machine ran for an hour on a
// credential the mesh had replaced, because its manifest restarted it on its
// environment file and nobody had thought to name the credential too. Composed here so
// no manifest has to say it, for a container or a daemon that names the secret's path.
if reads := secretsReadBy(copied, m); len(reads) > 0 {
copied["restart-on"] = withRestartOn(copied["restart-on"], reads)
}
// **A version prepares its state before it runs** (novox/hq ADR 0135). Derived from the
// module's own resource rather than declared beside it: what prepares the state is the
// module's own code, so what it is given has to be what that code is given — and a
@@ -911,17 +912,6 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
out = append(out, fact)
}
}
// The node's tool runtime, last (novox/hq ADR 0175, to-be 38 WP2.3): one process loading every
// bundle delivered above and holding the credential sealed above, so both exist before it starts
// — the order written here is the order the machine applies.
if r.runtimeHere() {
process, err := r.runtimeProcess(with)
if err != nil {
return nil, err
}
owner[fmt.Sprint(process["id"])] = RuntimeModule
out = append(out, process)
}
if with.Adopted {
// First, before anything a module declares: what the mesh needs reachable, then its guard.
// The order a machine applies is the order written here.
@@ -1137,6 +1127,19 @@ type Contribution struct {
// requirement's name — everything providing `reverse-proxy` understands the same shape, which
// is what makes swapping one for another cost nothing.
Values map[string]any `json:"values"`
// Derived is what this provider's own definition said it derives for this consumer, already
// derived (novox/hq ADR 0188).
//
// **The provider is told, rather than recomputing it.** A served value may name the consumer's
// identity — a bucket named for who is asking, a database prefixed with it — and before this
// the rule lived twice: once in the provisioner's code, once transcribed into every consumer's
// definition. The mesh fills the provider's own statement here and delivers the same filled
// value to the consumer, so the two cannot disagree: there is no second computation to
// disagree with.
//
// Only the keys that are per-consumer. The rest of what the provider serves is the same for
// everyone and is in its own definition, where it already is.
Derived map[string]any `json:"derived,omitempty"`
}
// grantPath is where one consumer's sealed credential lands on the providing machine.
@@ -1228,12 +1231,17 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
// told about it and withdraws the login on its next pass.
continue
}
as := holderAs(ConsumerIdentity(g.Consumer, IdentitySource(g.Slug, g.From)), g.Local)
derived, err := r.derivedFor(g.Provision, as, settings)
if err != nil {
return nil, err
}
out[g.Provision] = append(out[g.Provision], Contribution{
From: g.From, Node: g.Consumer, At: g.At, Values: g.Values,
From: g.From, Node: g.Consumer, At: g.At, Values: g.Values, Derived: derived,
// One holder per local name: the identity the consumer is known by, and the local name
// after it where the module keeps several (ADR 0094). Not a login any backend checks —
// a secret is not a login — so the identity limit does not apply to the suffix.
As: holderAs(ConsumerIdentity(g.Consumer, IdentitySource(g.Slug, g.From)), g.Local),
As: as,
Secret: grantPath(directories[g.Provision], g.Consumer, holderAs(g.From, g.Local)),
})
if granted[g.Provision] == nil {
@@ -2088,87 +2096,3 @@ func portOfEndpoint(values map[string]any, ports map[string]int) {
values["port"] = port
}
}
// secretsReadBy is the file resources of this module's own secrets that a container or a daemon reads
// — named in its volumes, its environment or its env-files by the secret's placed path — as
// restart-on ids. Nothing for other shapes, and nothing for a scheduled or run-once process, which
// the host refuses a restart-on for (it runs again anyway, and reads the file afresh).
func secretsReadBy(resource map[string]any, m Manifest) []string {
kind := fmt.Sprint(resource["type"])
if kind != "container" && kind != "process" {
return nil
}
if resource["schedule"] != nil || resource["run-once"] == true {
return nil
}
var mentioned []string
for _, key := range []string{"volumes", "env", "env-file"} {
mentioned = append(mentioned, stringsIn(resource[key])...)
}
var out []string
for _, name := range sortedKeys(m.OwnSecrets) {
path := m.OwnSecrets[name].Path
if path == "" {
continue
}
for _, s := range mentioned {
// A volume is `source:destination[:mode]`; an env value or an env-file is the path itself.
if s == path || strings.HasPrefix(s, path+":") {
out = append(out, m.Module+"."+NeedID(name))
break
}
}
}
return out
}
// stringsIn is every string in a list or a map's values; nothing for anything else.
func stringsIn(v any) []string {
switch x := v.(type) {
case []any:
var out []string
for _, item := range x {
if s, ok := item.(string); ok {
out = append(out, s)
}
}
return out
case []string:
return x
case map[string]any:
var out []string
for _, k := range sortedKeys(x) {
if s, ok := x[k].(string); ok {
out = append(out, s)
}
}
return out
case map[string]string:
var out []string
for _, k := range sortedKeys(x) {
out = append(out, x[k])
}
return out
}
return nil
}
// withRestartOn is a resource's restart-on list with these ids added once each.
func withRestartOn(have any, add []string) []any {
var out []any
seen := map[string]bool{}
for _, id := range reflectsRenamed("", have) {
s := fmt.Sprint(id)
if !seen[s] {
seen[s] = true
out = append(out, s)
}
}
for _, id := range add {
if !seen[id] {
seen[id] = true
out = append(out, id)
}
}
return out
}
@@ -0,0 +1,297 @@
package catalogue
import (
"encoding/json"
"strings"
"testing"
)
// What a provider derives for each consumer, said once and delivered to both ends
// (novox/hq ADR 0188, issue 124).
//
// The failure these are written against: the object store's provisioner derived each consumer's
// bucket from the login the mesh minted, in its own code, and the mesh had no channel to tell the
// consumer which bucket that was — so all three consumers wrote the answer into their own
// definitions by hand. Two were right. One named a predecessor's bucket and would have
// authenticated successfully and been refused on every object. Each of them also named the
// machine the module happens to run on, which a definition may not do.
// store is an object store in the shape minio has: it serves a region and a port to everyone, and
// a bucket named for whoever is asking.
func store() Manifest {
return Manifest{
Module: "store", Version: "1",
Provides: FromAnywhere("s3-bucket"),
Listens: []Listening{{Port: 9000, Protocol: "tcp", From: FromMesh}},
Serves: map[string]map[string]any{"s3-bucket": {
"region": "eu-west",
"bucket": "${consumer:as:dns}",
}},
Receives: map[string]string{"s3-bucket": "/var/lib/store/grants/mesh.json"},
Grants: map[string]string{"s3-bucket": "/var/lib/store/grants"},
Resources: []map[string]any{{
"id": "server", "type": "container", "name": "store", "ports": []any{"9000"},
}},
}
}
// files is a consumer that writes the bucket into its own configuration — which is the thing it
// could not do before, and had to transcribe.
func files() Manifest {
return Manifest{
Module: "files", Version: "1", Slug: "files",
Requires: []string{"s3-bucket"},
Binds: map[string]string{"s3-bucket": "/var/lib/files/store.json"},
Secrets: map[string]string{"s3-bucket": "/var/lib/files/store.secret"},
Resources: []map[string]any{{
"id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600",
"content": "BUCKET=${bound:s3-bucket:bucket}\nREGION=${bound:s3-bucket:region}\n",
}},
}
}
// pics is a second consumer of the same provider on the same machine: two derivations, neither
// the other's.
func pics() Manifest {
return Manifest{
Module: "pics", Version: "1", Slug: "pics",
Requires: []string{"s3-bucket"},
Binds: map[string]string{"s3-bucket": "/var/lib/pics/store.json"},
Secrets: map[string]string{"s3-bucket": "/var/lib/pics/store.secret"},
Resources: []map[string]any{{
"id": "env", "type": "file", "path": "/var/lib/pics/env", "mode": "0600",
"content": "BUCKET=${bound:s3-bucket:bucket}\n",
}},
}
}
// The three places the derived value lands must agree, because agreeing is the whole point: the
// consumer's own file, the binding it reads as JSON, and the provider's contributions entry.
func TestADerivedValueReachesBothEndsAndAgrees(t *testing.T) {
r, err := Resolve(shelf(store(), files()), []string{"store", "files"}, reachable(), World{})
if err != nil {
t.Fatal(err)
}
out, err := r.Declaration(Rendering{Grants: []Grant{{
Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files",
Values: map[string]any{}, Sealed: "c2VhbGVk",
}}})
if err != nil {
t.Fatal(err)
}
// The mesh minted this identity for the consumer; the bucket is that identity as a DNS label.
// Derived here with the mesh's own function, so the test cannot agree with a wrong rule.
as := ConsumerIdentity("workstation", IdentitySource("files", "files"))
want := strings.ReplaceAll(as, "_", "-")
if want == as || !strings.Contains(as, "_") {
t.Fatalf("the mesh's identity %q has no separator to rewrite; this test proves nothing", as)
}
env := fileNamed(out, "files.env")
if env == nil {
t.Fatalf("the consumer was given no file: %v", out)
}
if got := env["content"].(string); !strings.Contains(got, "BUCKET="+want+"\n") {
t.Errorf("the consumer's own file was not told the bucket:\n%s\nwant BUCKET=%s", got, want)
}
binding := fileNamed(out, "files.bound-s3-bucket")
if binding == nil {
t.Fatalf("the consumer was given no binding: %v", out)
}
var said struct {
Serves map[string]any `json:"serves"`
}
if err := json.Unmarshal([]byte(binding["content"].(string)), &said); err != nil {
t.Fatal(err)
}
if said.Serves["bucket"] != want {
t.Errorf("the binding says the bucket is %q, want %q", said.Serves["bucket"], want)
}
// And what is the same for everybody is still the same for everybody.
if said.Serves["region"] != "eu-west" {
t.Errorf("the binding lost what the provider serves to all: %v", said.Serves)
}
given := storeGrants(t, out)
if len(given) != 1 {
t.Fatalf("the provider was told about %d consumer(s): %v", len(given), given)
}
if given[0].Derived["bucket"] != want {
t.Errorf("the provider was told the bucket is %v, and the consumer was told %q — "+
"the two ends disagree, which is the whole failure", given[0].Derived["bucket"], want)
}
// Only the per-consumer half. The region is the same for everyone and is already in the
// provider's own definition; repeating it here would be a copy to go stale.
if _, carried := given[0].Derived["region"]; carried {
t.Errorf("the provider was handed back what it already says for everyone: %v", given[0].Derived)
}
}
// Two consumers of one provider on one machine get two buckets, and neither gets the other's.
func TestTwoConsumersOfOneProviderGetTheirOwnDerivation(t *testing.T) {
r, err := Resolve(shelf(store(), files(), pics()),
[]string{"store", "files", "pics"}, reachable(), World{})
if err != nil {
t.Fatal(err)
}
out, err := r.Declaration(Rendering{Grants: []Grant{
{Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files",
Values: map[string]any{}, Sealed: "c2VhbGVk"},
{Provision: "s3-bucket", Consumer: "workstation", From: "pics", Slug: "pics",
Values: map[string]any{}, Sealed: "c2VhbGVk"},
}})
if err != nil {
t.Fatal(err)
}
forFiles := strings.ReplaceAll(ConsumerIdentity("workstation", IdentitySource("files", "files")), "_", "-")
forPics := strings.ReplaceAll(ConsumerIdentity("workstation", IdentitySource("pics", "pics")), "_", "-")
if forFiles == forPics {
t.Fatal("the two consumers were given the same identity; this test proves nothing")
}
if got := fileNamed(out, "files.env")["content"].(string); !strings.Contains(got, "BUCKET="+forFiles+"\n") {
t.Errorf("files was not given its own bucket:\n%s", got)
}
if got := fileNamed(out, "pics.env")["content"].(string); !strings.Contains(got, "BUCKET="+forPics+"\n") {
t.Errorf("pics was not given its own bucket:\n%s", got)
}
var buckets []any
for _, g := range storeGrants(t, out) {
buckets = append(buckets, g.Derived["bucket"])
}
if len(buckets) != 2 || buckets[0] == buckets[1] {
t.Errorf("the provider was told %v; it must be told one bucket per consumer", buckets)
}
}
// An operator may still set what the provider serves, and the mesh still derives the rest: the
// setting is laid on first, then the consumer's half is filled.
func TestASettingComposesWithADerivedValue(t *testing.T) {
r, err := Resolve(shelf(store(), files()), []string{"store", "files"}, reachable(), World{})
if err != nil {
t.Fatal(err)
}
out, err := r.Declaration(Rendering{
Settings: SettingsBy{"store": {{From: "the operator",
Values: map[string]any{"bucket": "team-${consumer:as:dns}"}}}},
Grants: []Grant{{Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files",
Values: map[string]any{}, Sealed: "c2VhbGVk"}},
})
if err != nil {
t.Fatal(err)
}
want := "team-" + strings.ReplaceAll(ConsumerIdentity("workstation", IdentitySource("files", "files")), "_", "-")
if got := fileNamed(out, "files.env")["content"].(string); !strings.Contains(got, "BUCKET="+want+"\n") {
t.Errorf("the operator's prefix did not survive the derivation:\n%s\nwant BUCKET=%s", got, want)
}
if given := storeGrants(t, out); given[0].Derived["bucket"] != want {
t.Errorf("the provider was told %v, the consumer %q", given[0].Derived["bucket"], want)
}
}
// A fact or an alphabet the mesh does not have is refused where the definition is, not where a
// consumer happens to be resolved — and the refusal says what may be said instead.
func TestAServedValueNamingSomethingTheMeshDoesNotHaveIsRefused(t *testing.T) {
for _, c := range []struct{ value, says string }{
{"${consumer:node}", "as"},
{"${consumer:as:punycode}", "dns"},
} {
m := store()
m.Serves["s3-bucket"]["bucket"] = c.value
raw, err := json.Marshal(m)
if err != nil {
t.Fatal(err)
}
_, err = ParseManifest(raw)
if err == nil {
t.Fatalf("%s was accepted", c.value)
}
if !strings.Contains(err.Error(), c.value) {
t.Errorf("the refusal of %s does not quote it: %v", c.value, err)
}
if !strings.Contains(err.Error(), c.says) {
t.Errorf("the refusal of %s does not say what may be said (%q): %v", c.value, c.says, err)
}
}
}
// `dns` is checked against an identity the mesh actually mints, not an invented string.
func TestTheDNSAlphabetIsTheMintedIdentityWithItsSeparatorRewritten(t *testing.T) {
as := ConsumerIdentity("anchor", IdentitySource("ncloud", "nextcloud"))
if err := CheckIdentity("anchor", IdentitySource("ncloud", "nextcloud")); err != nil {
t.Fatalf("the mesh would not mint this identity at all: %v", err)
}
label := asDNSLabel(as)
if strings.Contains(label, "_") {
t.Errorf("%q is not a DNS label", label)
}
if strings.ReplaceAll(label, "-", "_") != as {
t.Errorf("%q is not %q with its separator rewritten", label, as)
}
}
// The check that would have caught the one wrong instance: a consumer that writes the derived
// value into its own definition instead of asking for it is refused, whether it transcribed the
// right answer or a predecessor's.
func TestAConsumerThatTranscribesWhatItsProviderDerivesIsRefused(t *testing.T) {
as := ConsumerIdentity("workstation", IdentitySource("files", "files"))
transcribed := strings.ReplaceAll(as, "_", "-")
m := files()
m.Resources = []map[string]any{{
"id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600",
// Exactly what the provider will create — correct today, and a copy of a rule that is
// not this module's.
"content": "BUCKET=" + transcribed + "\n",
}}
r, err := Resolve(shelf(store(), m), []string{"store", "files"}, reachable(), World{})
if err != nil {
t.Fatal(err)
}
_, err = r.Declaration(Rendering{Grants: []Grant{{
Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files",
Values: map[string]any{}, Sealed: "c2VhbGVk",
}}})
if err == nil {
t.Fatal("a definition holding its own copy of the provider's naming rule was accepted")
}
if !strings.Contains(err.Error(), "${bound:s3-bucket:bucket}") {
t.Errorf("the refusal does not say what to write instead: %v", err)
}
// And a constant the provider serves to everyone is not a transcription: repeating it is
// redundant, not wrong, and refusing it would be the mesh policing style.
m.Resources = []map[string]any{{
"id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600",
"content": "REGION=eu-west\n",
}}
r, err = Resolve(shelf(store(), m), []string{"store", "files"}, reachable(), World{})
if err != nil {
t.Fatal(err)
}
if _, err := r.Declaration(Rendering{Grants: []Grant{{
Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files",
Values: map[string]any{}, Sealed: "c2VhbGVk",
}}}); err != nil {
t.Errorf("a value the provider serves to everyone was judged a transcription: %v", err)
}
}
func storeGrants(t *testing.T, out []map[string]any) []Contribution {
t.Helper()
for _, r := range out {
if r["path"] != "/var/lib/store/grants/mesh.json" {
continue
}
var parsed struct {
Given []Contribution `json:"given"`
}
if err := json.Unmarshal([]byte(r["content"].(string)), &parsed); err != nil {
t.Fatal(err)
}
return parsed.Given
}
t.Fatalf("the provider was given no contributions file: %v", out)
return nil
}
+2 -12
View File
@@ -1,8 +1,6 @@
package catalogue
import (
"crypto/sha256"
"encoding/hex"
"fmt"
"sort"
"strings"
@@ -45,16 +43,8 @@ func jailsInto(modules []Manifest, j *Jailing) []map[string]any {
out := make([]map[string]any, 0, len(jails)+1)
for _, d := range jails {
// **The filter's digest rides in the jail file.** fail2ban is restarted when this file
// changes, and the filter is a file of its own: a module that changed only what a failure
// looks like rewrote the filter on disk and left the running jail on the old pattern, with
// nothing said (novox/hq issue 191's rollout found it on gitea's sshd). Naming the filter's
// digest here makes a changed pattern a changed jail file, so the restart the service
// already takes on it covers the filter too.
sum := sha256.Sum256([]byte(d.jail.Failregex))
fmt.Fprintf(&composed, "\n# from %s, filter %s\n[%s]\nenabled = true\nfilter = %s\n%s\n",
d.module, hex.EncodeToString(sum[:])[:12], d.jail.Name, d.jail.Name,
strings.TrimRight(d.jail.Jail, "\n"))
fmt.Fprintf(&composed, "\n# from %s\n[%s]\nenabled = true\nfilter = %s\n%s\n",
d.module, d.jail.Name, d.jail.Name, strings.TrimRight(d.jail.Jail, "\n"))
// The filter is a file of its own, named as the jail's filter= references it.
out = append(out, map[string]any{
"id": "filter-" + d.jail.Name,
-26
View File
@@ -44,29 +44,3 @@ func TestTheComposedJailFileIsWrittenEvenWhenEmpty(t *testing.T) {
t.Fatalf("the empty composed jail file was not written alone: %v", files)
}
}
// A changed pattern restarts fail2ban (novox/hq issue 191's rollout): the service restarts when the
// composed jail file changes, and the filter is a file of its own, so the jail file names the
// filter's digest. Changing only the failregex must change the jail file; the same pattern must not.
func TestAChangedFilterChangesTheJailFile(t *testing.T) {
jailFile := func(failregex string) string {
modules := []Manifest{
{Module: "fail2ban", Jailing: &Jailing{Into: "/etc/fail2ban/jail.d/mesh.conf", FilterInto: "/etc/fail2ban/filter.d"}},
{Module: "gitea", Jails: []Jail{{Name: "gitea", Failregex: failregex, Jail: "port = 222"}}},
}
for _, f := range jailsInto(modules, modules[0].Jailing) {
if f["id"] == ComposedJailsID() {
return f["content"].(string)
}
}
t.Fatal("no composed jail file")
return ""
}
before := jailFile("web login failed from <HOST>")
if again := jailFile("web login failed from <HOST>"); again != before {
t.Errorf("the same pattern composed a different jail file, which would restart fail2ban for nothing")
}
if after := jailFile("web login failed from <HOST>\n Invalid user .* from <HOST>"); after == before {
t.Errorf("a changed pattern left the jail file as it was, so fail2ban keeps the old filter:\n%s", after)
}
}
+4 -49
View File
@@ -572,36 +572,6 @@ type Manifest struct {
// a module that could ask for it could read every credential on the bus — and the claim on
// `mesh-broker` is what authorises it, checked from this manifest alone.
BusUsers string `json:"bus-users,omitempty"`
// Bundles are this module's compiled bundles as the build produced them: what each is called,
// where it is, what it hashes to, what language it is in and which files a tool runtime loads
// from it (novox/hq ADR 0175, to-be 38).
//
// **Derived, never written.** The manifest in a repository says `build.artifacts`; the manifest
// the mesh holds says what came out, the way a resource naming an artifact comes to name a
// digest. Kept here because a tools bundle is referenced by no resource of the module's own —
// the node's runtime loads it, and the runtime is composed by the mesh — so without this the
// resolved manifest would carry no trace of the one artifact the runtime needs. A repository
// manifest that writes this beside a build is refused: it would be stating the build's output
// by hand.
Bundles []Bundle `json:"bundles,omitempty"`
}
// Bundle is one compiled bundle after it exists, as the resolved manifest carries it.
type Bundle struct {
Name string `json:"name"`
// Source is where a machine fetches it, kept without the store's address like every reference
// the mesh records (artifacts.go); Digest is what it must hash to.
Source string `json:"source"`
Digest string `json:"digest"`
// Language is what it was compiled from, which is what says how it is run.
Language string `json:"language,omitempty"`
// Entrypoints are the compiled files it was built around, relative to its root.
Entrypoints []string `json:"entrypoints,omitempty"`
// Loads are the entrypoints a node's tool runtime imports from it: what the artifact said, or
// every entrypoint for a module declaring tools that said nothing. Empty for a bundle that is
// run rather than loaded.
Loads []string `json:"loads,omitempty"`
}
// Build says how to produce this module's artifacts from its source.
@@ -738,16 +708,6 @@ type Artifact struct {
// somebody adds a helper. An empty list is a bundle that is run rather than loaded — a
// provisioner or a step, named by whatever runs it.
Entrypoints []string `json:"entrypoints,omitempty"`
// Loads are the entrypoints of this bundle the node's tool runtime loads (novox/hq ADR 0175,
// to-be 38): the module's tool code, each file registering its tools as it is imported. A
// subset of Entrypoints, for a bundle that also carries things that are RUN — a daemon, a
// step, a report — and must not have them imported into the runtime.
//
// Absent means every entrypoint, for a module that declares `tools`: a bundle holding the
// module's tools and nothing else is the ordinary case and should not have to say the same
// list twice. A module declaring no tools has nothing the runtime loads, whatever it compiles.
Loads []string `json:"loads,omitempty"`
}
// Kinds an artifact may be.
@@ -1373,15 +1333,6 @@ func ParseManifest(raw []byte) (Manifest, error) {
//
// Refused here because the alternative is a build that never returns, on a mesh new enough
// that nobody is watching it yet.
if m.Build != nil && len(m.Bundles) > 0 {
// The output of a build, written beside the build that produces it (ADR 0175). A resource
// naming a digest beside an `artifact` would be the same mistake, and is caught the same way:
// what the mesh derives, a repository does not state.
problems = append(problems, fmt.Sprintf(
"%s writes `bundles` beside its build. The mesh derives that from what the build "+
"produced; a manifest states `build.artifacts` and nothing about what came out",
m.Module))
}
if m.Build != nil && len(m.Build.Artifacts) > 0 {
for _, o := range m.Offers() {
if o != ArtifactStoreProvision {
@@ -1409,6 +1360,10 @@ func ParseManifest(raw []byte) (Manifest, error) {
"%s serves %q to whoever requires it, and does not provide it", m.Module, to))
}
}
// A served value may be derived for the consumer it is served to (novox/hq ADR 0188). Read
// here, where the definition is, rather than when somebody first requires it: a rule that
// would be refused at the first consumer is wrong from the moment it is written.
problems = append(problems, CheckServes(m)...)
for to, where := range m.Binds {
if !placedOrAbsolute(where) {
problems = append(problems, fmt.Sprintf(
@@ -1,42 +0,0 @@
package catalogue
import (
"strings"
"testing"
)
// A resolved manifest keeps a built artifact's reference in the store-relative form, never the
// address the builder reached the store by (novox/hq ADR 0155): on 2026-10-02 the first bundle
// resolved on the mesh carried the store's host in bundles[0].source and registration refused it
// as naming an installation. Archive resources are the same kind of reference and get the same.
func TestAResolvedReferenceIsKeptNotRouted(t *testing.T) {
m, err := ParseManifest([]byte(`{
"module": "sample", "version": "1",
"tools": ["one"],
"build": {"artifacts": [
{"name": "code", "kind": "bundle", "language": "typescript", "entrypoints": ["tools/index.js"]},
{"name": "files", "kind": "archive", "from": "files"}
]},
"resources": [{"id": "packed", "type": "archive", "path": "/opt/sample", "artifact": "files"}]
}`))
if err != nil {
t.Fatal(err)
}
digest := "sha256:" + strings.Repeat("ab", 32)
resolved, err := m.Resolve([]Built{
{Name: "code", Kind: ArtifactBundle, Reference: "http://store.example:5100/v2/sample/code/blobs/" + digest, Digest: digest},
{Name: "files", Kind: ArtifactArchive, Reference: "http://store.example:5100/v2/sample/files/blobs/" + digest, Digest: digest},
})
if err != nil {
t.Fatal(err)
}
if got, want := resolved.Bundles[0].Source, ArtifactStoreScheme+"sample/code/blobs/"+digest; got != want {
t.Errorf("bundle source %q, want the kept form %q", got, want)
}
if got, want := resolved.Resources[0]["source"], ArtifactStoreScheme+"sample/files/blobs/"+digest; got != want {
t.Errorf("archive source %q, want the kept form %q", got, want)
}
if problems := InstallationProblems(resolved); len(problems) != 0 {
t.Errorf("a resolved manifest names an installation: %v", problems)
}
}
-251
View File
@@ -1,251 +0,0 @@
package catalogue
import (
"fmt"
"sort"
"strings"
)
// The node's tool runtime, as the catalogue knows it (novox/hq ADR 0175, to-be 38).
//
// **One module is the runtime.** Where it is assigned, one process per machine serves every assigned
// module's tools and every held seat's verbs, on the host side, from the bundles each module's build
// produced — and no module needs a container to reach the bus with its tools. The name is a constant
// rather than a manifest field because a rule turns on it: the composer places the runtime's process
// where this module is, and registration refuses the old pattern once this module exists.
// RuntimeModule is the module that is the node's tool runtime. Mirrored in the broker package,
// which composes a principal of its own for it; the agreement test there holds the two to one string.
const RuntimeModule = "node-tools"
// BundleRoot is where a machine keeps the tools bundles the mesh delivers to it: under the mesh's
// own directory, beside the daemons the host unpacks there, and never where a package manager also
// writes. One directory per module, one per bundle beneath it, at a path that does not move with
// the version — so the runtime's process names each entrypoint once and is restarted, not
// recomposed, when a bundle changes.
const BundleRoot = "/var/lib/mesh/bundles"
// BundleID names the archive resource that delivers one of a module's bundles; prefixed with the
// module like every resource of its own.
func BundleID(bundle string) string { return "bundle-" + bundle }
// BundlePath is where one module's bundle is unpacked on a machine.
func BundlePath(module, bundle string) string { return BundleRoot + "/" + module + "/" + bundle }
// runtimeHere says whether this node's set includes the runtime module, which is what decides
// whether anything about tools changes on the machine (to-be 38 WP2): until the runtime is assigned,
// a node is sent exactly what it was sent before, bundles included, because a bundle nothing loads
// is bytes nobody reads.
func (r Resolution) runtimeHere() bool {
for _, m := range r.Modules {
if m.Module == RuntimeModule {
return true
}
}
return false
}
// bundleArchives is one archive per tools bundle of a module — a bundle the runtime LOADS something
// from — as the host fetches and unpacks any artifact (novox/hq ADR 0175 §3: a module brings its
// tools as a bundle, delivered by the host like any artifact, never an image). A bundle it loads
// nothing from is run rather than loaded: a daemon, a step, the runtime itself — delivered by the
// process that runs it, and not again here.
//
// The source is the kept reference; the per-resource pass that follows routes it through the
// artifact store as this network reaches it now, as it does every image and archive the mesh built.
func bundleArchives(m Manifest) []map[string]any {
var out []map[string]any
for _, b := range m.Bundles {
if len(b.Loads) == 0 {
continue
}
out = append(out, map[string]any{
"id": BundleID(b.Name), "type": "archive",
"source": b.Source, "digest": b.Digest,
"path": BundlePath(m.Module, b.Name),
})
}
return out
}
// RuntimeProcessID names the one process the mesh composes for a machine's runtime; prefixed with
// the runtime module like a resource of its own, because that module is what the host sees it as.
func RuntimeProcessID() string { return "runtime" }
// RuntimeToolModules is the variable the runtime reads the modules it serves from: one
// `<module>=<entrypoint>` per file it loads, comma-separated — several entries may name one module.
// RuntimeBrokerFile is where it reads the node's credential; RuntimeOperatorAccount and
// RuntimeOperatorHome are the machine's operator account and home, handed to every tool's
// environment (to-be 38 WP1), and absent on a machine with no account.
const (
RuntimeToolModules = "MESH_TOOL_MODULES"
RuntimeBrokerFile = "MESH_BROKER_FILE"
RuntimeOperatorAccount = "MESH_OPERATOR_ACCOUNT"
RuntimeOperatorHome = "MESH_OPERATOR_HOME"
)
// interpreterFor is how a bundle in a language is run: the program the host's unit starts, with the
// bundle's entrypoint after it. The one thing the composer takes from a language, and said here
// rather than in a manifest because the runtime's process is the mesh's to compose (to-be 38 WP3).
func interpreterFor(language string) (string, error) {
switch language {
case "typescript":
return "node", nil
}
return "", fmt.Errorf(
"%s is written in %q, and the mesh knows no interpreter to run a %q bundle with",
RuntimeModule, language, language)
}
// runtimeProcess is the one process a machine runs the node's tool runtime as (novox/hq ADR 0175,
// to-be 38 WP2.3): the runtime module's own bundle, run by its language's interpreter, told which
// modules it serves and from which files, where its credential is, and who the machine's operator
// is — and restarted when any bundle it loads or the credential it holds changes.
//
// Composed from the placed manifests, so the credential's path is where this node puts it. The
// runtime runs as the operator's account when the machine has one, which is what lets a tool that
// needs root escalate as the operator would (ADR 0175 §4); on a machine with no account it runs as
// root, and the two operator words are not set.
func (r Resolution) runtimeProcess(with Rendering) (map[string]any, error) {
var runtime *Manifest
for i := range r.Modules {
if r.Modules[i].Module == RuntimeModule {
runtime = &r.Modules[i]
}
}
if runtime == nil {
return nil, nil
}
if len(runtime.Bundles) != 1 {
return nil, fmt.Errorf(
"%s is assigned to %s and its build produced %d bundle(s); the runtime is one bundle "+
"the mesh runs, so the module declares exactly one (novox/hq to-be 38)",
RuntimeModule, r.Node, len(runtime.Bundles))
}
bundle := runtime.Bundles[0]
if len(bundle.Entrypoints) != 1 {
return nil, fmt.Errorf(
"%s's bundle %q names %d entrypoint(s); the runtime is run from one, so the module "+
"declares exactly one (novox/hq to-be 38)", RuntimeModule, bundle.Name, len(bundle.Entrypoints))
}
interpreter, err := interpreterFor(bundle.Language)
if err != nil {
return nil, err
}
credential, declared := runtime.OwnSecrets["broker"]
if !declared {
return nil, fmt.Errorf(
"%s declares no own secret named broker, and the node's credential is delivered there: "+
"a module that speaks on the bus declares \"own-secrets\": {\"broker\": <path>}",
RuntimeModule)
}
// What it serves, and from which files: every module on this machine that composes here, in
// name order, each bundle it loads from in the order the manifest gave. A module left out of
// the declaration — a filter on an adopted machine — is left out of this too, or the runtime
// would be told to load files that were never delivered.
var served []string
var restartOn []string
for _, m := range r.Modules {
if with.Adopted && m.Filtering != nil {
continue
}
for _, b := range m.Bundles {
if len(b.Loads) == 0 {
continue
}
for _, load := range b.Loads {
served = append(served, m.Module+"="+BundlePath(m.Module, b.Name)+"/"+load)
}
restartOn = append(restartOn, m.Module+"."+BundleID(b.Name))
}
}
sort.Strings(served)
restartOn = append(restartOn, RuntimeModule+"."+NeedID("broker"))
sort.Strings(restartOn)
env := map[string]string{
RuntimeToolModules: strings.Join(served, ","),
RuntimeBrokerFile: credential.Path,
}
process := map[string]any{
"id": RuntimeModule + "." + RuntimeProcessID(), "type": "process", "name": RuntimeModule,
"source": bundle.Source, "digest": bundle.Digest,
"run": []any{interpreter, bundle.Entrypoints[0]},
"env": env,
"restart-on": toAny(restartOn),
}
if r.Account != "" {
env[RuntimeOperatorAccount] = r.Account
env[RuntimeOperatorHome] = accountHomeOf(r.Account, r.AccountHome)
process["user"] = r.Account
}
// Routed through the artifact store as this network reaches it now, like everything the mesh
// built; refused with the same words when there is no store to route through.
if err := artifactsInto(process, RuntimeModule, with); err != nil {
return nil, err
}
return process, nil
}
func toAny(in []string) []any {
out := make([]any, 0, len(in))
for _, s := range in {
out = append(out, s)
}
return out
}
// RuntimeImageModule and RuntimeImageArtifact name the image every per-module tool container was
// built on: the tool runtime's own runtime image. With the runtime a module of its own, that image
// stays the way a module's SERVICE may be built and stops being the way tools reach a node (ADR 0175).
const (
RuntimeImageModule = "mesh-tools"
RuntimeImageArtifact = "runtime"
)
// ToolContainerOnTheRuntime says why a manifest is the pattern ADR 0175 retires — a module whose tools
// are served from a container built on the tool runtime's image — or nothing when it is not. Judged
// from the manifest's own `build.on` when it is a repository manifest, and from what its build stood
// on when it is a built one, because a resolved manifest carries no build. The gate itself is
// registration's (to-be 38 WP2.4): once the runtime module is in the catalogue, this is refused.
//
// Three things must hold, and each alone is fine: declaring tools (a bundle does that); a container
// (a module's service may well be one); building on the runtime's image (a service written against
// the SDK may). All three is a container whose purpose is tools, which the runtime now serves.
func ToolContainerOnTheRuntime(m Manifest, against []string) string {
if len(m.Tools) == 0 {
return ""
}
container := false
for _, r := range m.Resources {
if fmt.Sprint(r["type"]) == "container" {
container = true
}
}
if !container {
return ""
}
onTheRuntime := false
if m.Build != nil {
for _, on := range m.Build.On {
if on.Module == RuntimeImageModule && on.Artifact == RuntimeImageArtifact {
onTheRuntime = true
}
}
}
for _, ref := range against {
path, kept := InArtifactStore(Recorded(ref))
if kept && strings.HasPrefix(path, RuntimeImageModule+"/"+RuntimeImageArtifact+"@") {
onTheRuntime = true
}
}
if !onTheRuntime {
return ""
}
return fmt.Sprintf(
"%s declares tools and a container built on %s's %s image — a container whose purpose is "+
"serving tools. The node's tool runtime (%s) serves every module's tools from its bundle "+
"now (novox/hq ADR 0175, to-be 38); declare the tools as a bundle and drop the container",
m.Module, RuntimeImageModule, RuntimeImageArtifact, RuntimeModule)
}
-88
View File
@@ -1,88 +0,0 @@
package catalogue
import (
"strings"
"testing"
)
// The packet-filter manifest as it was the day the runtime was decided (novox/hq ADR 0175): tools,
// served from a container built on the tool runtime's image, with NET_ADMIN so the container could
// reach the filter. The exact pattern to-be 38 WP4 moves it off, and the one the gate refuses.
const thePacketFilterAsItWas = `{
"module": "nftables",
"version": "1",
"capabilities": ["firewall", "container-runtime"],
"claims": [{"name": "node-packet-filter", "scope": "node", "serves": ["rules", "reload", "remove"]}],
"filtering": {"into": "/etc/nftables.conf"},
"resources": [
{"id": "mesh-state", "type": "directory", "mode": "0700", "place": "mesh"},
{"id": "package", "type": "package", "package": "nftables"},
{"id": "unit", "type": "file", "path": "/etc/systemd/system/mesh-filter.service",
"content": "[Unit]\nDescription=The mesh's packet filter\n[Service]\nType=oneshot\nExecStart=nft -f /etc/nftables.conf\n", "mode": "0644"},
{"id": "load", "type": "service", "unit": "mesh-filter.service", "state": "running", "boot": "enabled",
"restart-on": ["unit"], "reload-on": ["filtering"]},
{"id": "runtime", "type": "container", "name": "mesh-nftables", "network": "host",
"capabilities": ["NET_ADMIN"],
"volumes": ["${dir:mesh-state}/broker:/run/secrets/broker:ro", "/etc/nftables.conf:/etc/nftables.conf:ro"],
"env": {"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_FILTER_FILE": "/etc/nftables.conf"},
"artifact": "runtime"}
],
"tools": ["firewall_rules"],
"own-secrets": {"broker": "${dir:mesh-state}/broker"},
"build": {
"on": [
{"arg": "BUILD_BASE", "module": "mesh-tools", "artifact": "build"},
{"arg": "RUNTIME_BASE", "module": "mesh-tools", "artifact": "runtime"}
],
"artifacts": [{"name": "runtime", "kind": "image", "from": "Dockerfile"}]
}
}`
func TestAToolContainerOnTheRuntimeImageIsNamedForWhatItIs(t *testing.T) {
m, err := ParseManifest([]byte(thePacketFilterAsItWas))
if err != nil {
t.Fatal(err)
}
// From the repository: the manifest says what it builds on.
why := ToolContainerOnTheRuntime(m, nil)
if why == "" {
t.Fatal("the packet filter's tool container was not recognised from its build")
}
for _, word := range []string{"nftables", "mesh-tools", "runtime", "ADR 0175", "bundle"} {
if !strings.Contains(why, word) {
t.Errorf("the refusal does not say %q: %s", word, why)
}
}
// Built: the manifest carries no build, and what it stood on says the same.
built, err := m.Resolve([]Built{{Name: "runtime", Kind: ArtifactImage,
Reference: ArtifactStoreScheme + "nftables/runtime@" + digest}})
if err != nil {
t.Fatal(err)
}
stoodOn := []string{"anchor.internal:5100/mesh-tools/build@" + digest, "anchor.internal:5100/mesh-tools/runtime@" + digest}
if ToolContainerOnTheRuntime(built, stoodOn) == "" {
t.Error("the packet filter's tool container was not recognised from what its build stood on")
}
if ToolContainerOnTheRuntime(built, nil) != "" {
t.Error("a built manifest with no record of its base was judged to be on the runtime")
}
// Each of the three alone is an ordinary module.
bundle := m
bundle.Resources = m.Resources[:len(m.Resources)-1]
if ToolContainerOnTheRuntime(bundle, nil) != "" {
t.Error("a module with tools and no container is the pattern the runtime serves, and was refused")
}
service := m
service.Tools = nil
if ToolContainerOnTheRuntime(service, nil) != "" {
t.Error("a service built against the SDK, declaring no tools, was refused")
}
elsewhere := m
elsewhere.Build = &Build{On: []BuildsOn{{Arg: "NODE_BASE", Image: "node@" + digest}},
Artifacts: m.Build.Artifacts}
if ToolContainerOnTheRuntime(elsewhere, nil) != "" {
t.Error("a tool container on a public base was refused as though it were on the runtime's")
}
}
-212
View File
@@ -1,212 +0,0 @@
package catalogue
import (
"fmt"
"strings"
"testing"
)
// The node's tool runtime (novox/hq ADR 0175, to-be 38): where the runtime module is assigned, a
// machine is sent every assigned module's tools bundle as an archive, and the runtime's own process
// loading them. Where it is not, the machine is sent exactly what it was sent before.
var bundleDigest = "sha256:" + strings.Repeat("b", 64)
// aToolsModule is a module whose tools come as a compiled bundle and nothing else — the shape every
// module takes once its tool container goes (to-be 38 WP4).
func aToolsModule(t *testing.T, name string, entrypoints ...string) Manifest {
t.Helper()
m := Manifest{Module: name, Version: "1", Tools: []string{"status"},
Build: &Build{Artifacts: []Artifact{
{Name: "tools", Kind: ArtifactBundle, Language: "typescript", Entrypoints: entrypoints},
}}}
resolved, err := m.Resolve([]Built{{Name: "tools", Kind: ArtifactBundle,
Reference: ArtifactStoreScheme + name + "/tools/blobs/" + bundleDigest, Digest: bundleDigest}})
if err != nil {
t.Fatal(err)
}
return resolved
}
// theRuntime is the runtime module as the catalogue holds it: its own bundle, run rather than
// loaded, and its broker secret to receive the node's credential in.
func theRuntime(t *testing.T) Manifest {
t.Helper()
m := Manifest{Module: RuntimeModule, Version: "1",
OwnSecrets: OwnSecrets{"broker": {Path: "/var/lib/mesh/" + RuntimeModule + "/broker"}},
Build: &Build{Artifacts: []Artifact{{Name: "runtime", Kind: ArtifactBundle, Language: "typescript",
Entrypoints: []string{"src/main.js"}}}}}
resolved, err := m.Resolve([]Built{{Name: "runtime", Kind: ArtifactBundle,
Reference: ArtifactStoreScheme + RuntimeModule + "/runtime/blobs/" + bundleDigest, Digest: bundleDigest}})
if err != nil {
t.Fatal(err)
}
return resolved
}
func TestABuildsBundlesAreCarriedOnTheResolvedManifest(t *testing.T) {
m := aToolsModule(t, "nftables", "tools/index.js")
if len(m.Bundles) != 1 {
t.Fatalf("the resolved manifest carries %d bundle(s), not the one the build made", len(m.Bundles))
}
b := m.Bundles[0]
if b.Name != "tools" || b.Digest != bundleDigest || b.Language != "typescript" ||
b.Source != ArtifactStoreScheme+"nftables/tools/blobs/"+bundleDigest ||
len(b.Entrypoints) != 1 || b.Entrypoints[0] != "tools/index.js" {
t.Errorf("the bundle is carried as %+v", b)
}
// A repository manifest may not write what the build derives.
raw := `{"module":"x","version":"1","build":{"artifacts":[{"name":"t","kind":"bundle","language":"typescript"}]},` +
`"bundles":[{"name":"t","source":"s","digest":"` + bundleDigest + `"}]}`
if _, err := ParseManifest([]byte(raw)); err == nil || !strings.Contains(err.Error(), "bundles") {
t.Errorf("a manifest stating its build's output by hand was accepted: %v", err)
}
}
func TestEveryToolsBundleIsDeliveredWhereTheRuntimeRuns(t *testing.T) {
store := Rendering{ArtifactStore: "anchor.internal:5101",
Needed: map[string]map[string]string{RuntimeModule: {"broker": "sealed-credential"}}}
nftables := aToolsModule(t, "nftables", "tools/index.js")
zsh := aToolsModule(t, "zsh", "tools/index.js", "tools/more.js")
t.Run("with the runtime, one archive per tools bundle", func(t *testing.T) {
r := Resolution{Node: "anchor", Modules: []Manifest{nftables, zsh, theRuntime(t)}}
out, err := r.Declaration(store)
if err != nil {
t.Fatal(err)
}
archive := fileNamed(out, "nftables."+BundleID("tools"))
if archive == nil {
t.Fatalf("nftables' tools bundle was not delivered: %v", ids(out))
}
if archive["type"] != "archive" || archive["digest"] != bundleDigest ||
archive["path"] != BundleRoot+"/nftables/tools" {
t.Errorf("delivered as %v", archive)
}
if archive["source"] != "http://anchor.internal:5101/v2/nftables/tools/blobs/"+bundleDigest {
t.Errorf("fetched from %v, not through the store as this network reaches it", archive["source"])
}
if fileNamed(out, "zsh."+BundleID("tools")) == nil {
t.Errorf("zsh's tools bundle was not delivered: %v", ids(out))
}
// The runtime's own bundle is run, not loaded: its process delivers it, not an archive.
if fileNamed(out, RuntimeModule+"."+BundleID("runtime")) != nil {
t.Error("the runtime's own bundle was delivered as an archive beside its process")
}
})
t.Run("without the runtime, nothing changes", func(t *testing.T) {
r := Resolution{Node: "anchor", Modules: []Manifest{nftables, zsh}}
out, err := r.Declaration(store)
if err != nil {
t.Fatal(err)
}
for _, id := range ids(out) {
if strings.Contains(id, BundleID("")) {
t.Errorf("%s was delivered to a machine running no runtime to load it", id)
}
}
})
}
func ids(out []map[string]any) []string {
var names []string
for _, r := range out {
names = append(names, r["id"].(string))
}
return names
}
// One process per machine runs the runtime from its own bundle, told what it serves and from where,
// where its credential is, and who the operator is — restarted when any of that changes.
func TestTheMachineRunsOneRuntimeLoadingEveryDeliveredBundle(t *testing.T) {
with := Rendering{ArtifactStore: "anchor.internal:5101",
Needed: map[string]map[string]string{RuntimeModule: {"broker": "sealed-credential"}}}
nftables := aToolsModule(t, "nftables", "tools/index.js")
// A bundle carrying a daemon beside its tools says which files the runtime loads.
showcase := Manifest{Module: "showcase", Version: "1", Tools: []string{"greet"},
Build: &Build{Artifacts: []Artifact{{Name: "code", Kind: ArtifactBundle, Language: "typescript",
Entrypoints: []string{"daemon/index.js", "tools/index.js"}, Loads: []string{"tools/index.js"}}}}}
showcase, err := showcase.Resolve([]Built{{Name: "code", Kind: ArtifactBundle,
Reference: ArtifactStoreScheme + "showcase/code/blobs/" + bundleDigest, Digest: bundleDigest}})
if err != nil {
t.Fatal(err)
}
r := Resolution{Node: "anchor", Account: "ops", Modules: []Manifest{nftables, showcase, theRuntime(t)}}
out, err := r.Declaration(with)
if err != nil {
t.Fatal(err)
}
process := fileNamed(out, RuntimeModule+"."+RuntimeProcessID())
if process == nil {
t.Fatalf("no runtime process was composed: %v", ids(out))
}
if process["type"] != "process" || process["name"] != RuntimeModule || process["digest"] != bundleDigest ||
process["source"] != "http://anchor.internal:5101/v2/"+RuntimeModule+"/runtime/blobs/"+bundleDigest {
t.Errorf("the runtime's process is %v", process)
}
if fmt.Sprint(process["run"]) != "[node src/main.js]" {
t.Errorf("the runtime is run as %v; its bundle's one entrypoint, by its language's interpreter", process["run"])
}
env := process["env"].(map[string]string)
if env[RuntimeToolModules] != "nftables="+BundleRoot+"/nftables/tools/tools/index.js,"+
"showcase="+BundleRoot+"/showcase/code/tools/index.js" {
t.Errorf("the runtime is told to serve %q: every loaded file, by module, and nothing a bundle runs", env[RuntimeToolModules])
}
if env[RuntimeBrokerFile] != "/var/lib/mesh/"+RuntimeModule+"/broker" {
t.Errorf("the runtime reads its credential at %q, not where the module's own secret is placed", env[RuntimeBrokerFile])
}
if env[RuntimeOperatorAccount] != "ops" || env[RuntimeOperatorHome] != "/home/ops" || process["user"] != "ops" {
t.Errorf("the operator is not handed to the runtime: %v as %v", env, process["user"])
}
// The credential the process reads belongs to the account it runs as, or it could not read it
// (to-be 38 WP3); other modules' secrets are left as their manifests say.
if credential := fileNamed(out, RuntimeModule+"."+NeedID("broker")); credential == nil || credential["owner"] != "ops" {
t.Errorf("the runtime's credential is not the account's to read: %v", credential)
}
restarts := fmt.Sprint(process["restart-on"])
for _, want := range []string{"nftables." + BundleID("tools"), "showcase." + BundleID("code"), RuntimeModule + "." + NeedID("broker")} {
if !strings.Contains(restarts, want) {
t.Errorf("the runtime is not restarted when %s changes: %s", want, restarts)
}
}
// After every bundle and the credential, so both exist before it starts.
names := ids(out)
if names[len(names)-1] != RuntimeModule+"."+RuntimeProcessID() {
t.Errorf("the runtime's process is not last: %v", names)
}
t.Run("a machine with no account runs it as root without the operator words", func(t *testing.T) {
out, err := Resolution{Node: "anchor", Modules: []Manifest{nftables, theRuntime(t)}}.Declaration(with)
if err != nil {
t.Fatal(err)
}
process := fileNamed(out, RuntimeModule+"."+RuntimeProcessID())
env := process["env"].(map[string]string)
if _, set := env[RuntimeOperatorAccount]; set {
t.Error("an operator account was named on a machine that has none")
}
if _, set := process["user"]; set {
t.Error("a user was set on a machine with no account")
}
if credential := fileNamed(out, RuntimeModule+"."+NeedID("broker")); credential == nil || credential["owner"] != nil {
t.Errorf("the runtime's credential was given an owner on a machine with no account: %v", credential)
}
})
t.Run("a runtime module built wrong is refused by name", func(t *testing.T) {
two := Manifest{Module: RuntimeModule, Version: "1", OwnSecrets: OwnSecrets{"broker": {Path: "/b"}},
Build: &Build{Artifacts: []Artifact{{Name: "runtime", Kind: ArtifactBundle, Language: "typescript",
Entrypoints: []string{"a.js", "b.js"}}}}}
resolved, err := two.Resolve([]Built{{Name: "runtime", Kind: ArtifactBundle,
Reference: ArtifactStoreScheme + "x/runtime/blobs/" + bundleDigest, Digest: bundleDigest}})
if err != nil {
t.Fatal(err)
}
_, err = Resolution{Node: "anchor", Modules: []Manifest{resolved}}.Declaration(with)
if err == nil || !strings.Contains(err.Error(), "entrypoint") {
t.Errorf("a runtime bundle with two entrypoints was composed: %v", err)
}
})
}
+1 -10
View File
@@ -96,17 +96,8 @@ var defaultSeats = []Seat{
// A build says what it does as it does it (novox/hq ADR 0157): `started` when work is taken,
// `log.<build id>` for every line, `built` for the outcome. The log's tail token is the build's
// id, so a reader follows one build by subject alone.
// **Node-scoped, and every holder takes from one queue** (novox/hq ADR 0190): a build is asked of
// the role, and whichever machine holding the seat is idle pulls it. One holder per machine is
// what the scope says; sharing the work is what a seat's queue has always done.
{Name: "node-build-agent", Scope: ScopeNode,
Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0190"},
// **Retired by ADR 0190, kept while a manifest still claims it.** The one build machine's seat.
// A claim to a seat the mesh no longer defines is refused, and the module holding this one is
// assigned on a live machine until build-agent replaces it — removing the row first would make
// that machine unresolvable in the meantime. Deleted once no registered manifest claims it.
{Name: "mesh-build-machine", Scope: ScopeMesh,
Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0190"},
Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0121"},
{Name: "node-dns-resolver", Scope: ScopeNode, Decision: "novox/hq ADR 0121"},
// The intrusion prevention's verbs (novox/hq ADR 0179): what a person asks a machine's ban list
// whatever keeps it — who is banned and why, ban one address, let one go. Every holder serves all
+3 -4
View File
@@ -44,10 +44,9 @@ func TestTheSeatsAreAClosedSetAndEachNamesItsDecision(t *testing.T) {
delivered[s.Delivers] = s.Name
}
}
// Seventeen since node-build-agent (novox/hq ADR 0190) — sixteen once the retired
// mesh-build-machine row goes, when no registered manifest claims it any more.
if len(Seats()) != 17 {
t.Errorf("the mesh defines %d seats rather than 17; the set is closed, so a change here is "+
// Sixteen since node-service-manager (novox/hq ADR 0177).
if len(Seats()) != 16 {
t.Errorf("the mesh defines %d seats rather than 16; the set is closed, so a change here is "+
"a decision (novox/hq ADR 0110): %s", len(Seats()), seatNames())
}
}
-52
View File
@@ -1,52 +0,0 @@
package catalogue
import (
"reflect"
"strings"
"testing"
)
// What reads one of a module's own secrets is restarted when the secret changes (novox/hq issue 203,
// issue 206): the build machine kept an hour-old credential open because its manifest restarted it
// on its environment file alone. Composed, so a manifest need not say it; a scheduled process is
// left alone, because the host refuses a restart-on for one and it reads the file afresh each run.
func TestAContainerReadingAnOwnSecretIsRestartedWhenItChanges(t *testing.T) {
m := Manifest{Module: "agent", Version: "1",
OwnSecrets: OwnSecrets{"broker": {Path: "/var/lib/mesh/agent/broker"}},
Resources: []map[string]any{
{"id": "mesh-state", "type": "directory", "path": "/var/lib/mesh/agent", "mode": "0700"},
{"id": "settings", "type": "file", "path": "/var/lib/mesh/agent/agent.env", "mode": "0600", "content": "A=1\n"},
{"id": "server", "type": "container", "name": "agent", "network": "host",
"image": "registry.example/agent@sha256:" + strings.Repeat("a", 64),
"volumes": []any{"/var/lib/mesh/agent:/run/mesh:ro", "/var/lib/mesh/agent/broker:/run/mesh/broker:ro"},
"env-file": []any{"/var/lib/mesh/agent/agent.env"},
"restart-on": []any{"settings"}},
{"id": "nightly", "type": "container", "name": "agent-nightly", "schedule": "0 3 * * *",
"image": "registry.example/agent@sha256:" + strings.Repeat("a", 64),
"volumes": []any{"/var/lib/mesh/agent/broker:/run/mesh/broker:ro"}},
{"id": "other", "type": "container", "name": "agent-other",
"image": "registry.example/agent@sha256:" + strings.Repeat("a", 64)},
}}
got, err := Resolve(shelf(m), []string{m.Module},
Node{Name: "anchor", At: "10.0.0.1", Capabilities: map[string]bool{"container-runtime": true}}, World{})
if err != nil {
t.Fatal(err)
}
out, err := got.Declaration(Rendering{Needed: map[string]map[string]string{"agent": {"broker": "SEALED"}}})
if err != nil {
t.Fatal(err)
}
by := map[string]map[string]any{}
for _, r := range out {
by[r["id"].(string)] = r
}
if want := []any{"agent.settings", "agent.needs-broker"}; !reflect.DeepEqual(by["agent.server"]["restart-on"], want) {
t.Fatalf("the server reads the credential and is not restarted on it: %v", by["agent.server"]["restart-on"])
}
if _, has := by["agent.nightly"]["restart-on"]; has {
t.Fatalf("a scheduled container was given a restart-on, which the host refuses: %v", by["agent.nightly"]["restart-on"])
}
if _, has := by["agent.other"]["restart-on"]; has {
t.Fatalf("a container that reads no secret was given one to restart on: %v", by["agent.other"]["restart-on"])
}
}
@@ -1,23 +0,0 @@
package catalogue
import (
"regexp"
"testing"
)
// The bus is never public (novox/hq ADR 0169). Its port is what the bus module declares, the mesh,
// and the control plane adds no opening of its own: a machine joins through the tunnel, so the
// broker's host is filtered like any other. Before this, the broker port was a foundation port and
// rendered from anywhere beside its from-the-mesh rule.
func TestTheBusPortIsReachedFromTheMeshAlone(t *testing.T) {
rules := []Rule{{Port: 4222, Protocol: "tcp", From: FromMesh, Because: []string{"nats"},
Why: []string{"the mesh bus"}}}
out := AsNftables(rules, []string{"10.10.0.1", "10.10.0.2"}, true, nil, []string{"eth0"}, "mesh0")
if !regexp.MustCompile(`ip saddr \{ 10\.10\.0\.1, 10\.10\.0\.2 \} tcp dport 4222 accept`).MatchString(out) {
t.Fatalf("the bus is not reachable from the mesh:\n%s", out)
}
if regexp.MustCompile(`(?m)^\s*tcp dport 4222 accept`).MatchString(out) {
t.Fatalf("the bus is reachable from anywhere:\n%s", out)
}
}
-3
View File
@@ -42,9 +42,6 @@ const (
BusModule = "module"
BusEnrolment = "enrolment"
BusPerson = "person"
// BusNodeTools is a machine's tool runtime (novox/hq ADR 0175): named like the module it
// stands for, recorded as what it is.
BusNodeTools = "node-tools"
)
// MintBusPassword makes a bus password and records its hash under a username, replacing whatever was
-61
View File
@@ -42,11 +42,6 @@ type Source struct {
// Seen is when the source was last looked at — by a build, by hand, or by the forge saying it
// moved. What a late report of an older move is judged against.
Seen time.Time
// Against is every artifact the build this manifest came from stood on, as recorded. Part of a
// module's provenance like the commit is, and what tells a built manifest's base when the manifest
// itself no longer carries its build (novox/hq to-be 38 WP2.4). Empty for a manifest handed over
// by hand, which carries its `build.on` itself.
Against []string
}
// Current reports whether what the mesh holds is what the source last had.
@@ -66,32 +61,6 @@ func (s Source) Current() bool {
// 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 {
// **Once the node's tool runtime is in the catalogue, the pattern it retires may not spread**
// (novox/hq ADR 0175, to-be 38 WP2.4): a module serving its tools from a container built on the
// runtime's image. Refused at registration, by name, for a module that is new to the catalogue
// or that was registered in another shape — the mechanism that keeps the old pattern from
// returning by habit. **Not refused for a module already registered in that shape**: the
// catalogue holds some thirty of them the day the runtime arrives, each moves to a bundle in
// its own change (to-be 38 WP4 onward), and a gate that refused every rebuild of every unmoved
// module in the meantime would stop the whole pipeline to make a point the record already makes.
// Before the runtime exists the pattern is accepted as it always was.
if m.Module != catalogue.RuntimeModule {
if why := catalogue.ToolContainerOnTheRuntime(m, from.Against); why != "" {
runtime, err := i.hasModule(ctx, catalogue.RuntimeModule)
if err != nil {
return err
}
if runtime {
already, err := i.registeredInThatShape(ctx, m.Module)
if err != nil {
return err
}
if !already {
return fmt.Errorf("%s is not registered: %s", m.Module, why)
}
}
}
}
raw, err := json.Marshal(m)
if err != nil {
return err
@@ -120,36 +89,6 @@ func (i *Inventory) RegisterModule(ctx context.Context, m catalogue.Manifest, fr
return err
}
// registeredInThatShape is whether the catalogue already holds this module as a tools container on
// the runtime's image — judged from the manifest it holds and what that module's newest build stood
// on, the same two things the gate judges a new registration by. False for a module the catalogue
// does not hold.
func (i *Inventory) registeredInThatShape(ctx context.Context, name string) (bool, error) {
held, err := i.Catalogue(ctx)
if err != nil {
return false, err
}
stored, has := held[name]
if !has {
return false, nil
}
against, err := i.BuiltAgainst(ctx)
if err != nil {
return false, err
}
return catalogue.ToolContainerOnTheRuntime(stored, against[name]) != "", nil
}
// hasModule is whether the catalogue holds a module of that name.
func (i *Inventory) hasModule(ctx context.Context, name string) (bool, error) {
var one int
err := i.store.Pool().QueryRow(ctx, `select 1 from module where name = $1`, name).Scan(&one)
if errors.Is(err, pgx.ErrNoRows) {
return false, nil
}
return err == nil, 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
-58
View File
@@ -685,61 +685,3 @@ func TestRegisteringWithoutProvenanceKeepsTheSeat(t *testing.T) {
t.Fatalf("a hand-registered manifest erased where the module comes from: %+v", got)
}
}
// Once the node's tool runtime is in the catalogue, a module serving its tools from a container
// built on the runtime's image is refused at registration, naming the record (novox/hq ADR 0175,
// to-be 38 WP2.4) — for a module new to the catalogue or one that had moved away from it; a module
// already standing in that shape is rebuilt as before, so the catalogue's pipeline keeps running
// while each moves (WP3's amendment). Before the runtime, it is accepted as it always was — so a
// mesh converts in the order the design says and nothing is refused before there is anything to
// move to.
func TestAToolContainerIsRefusedOnceTheRuntimeIsRegistered(t *testing.T) {
inv := fresh(t)
ctx := t.Context()
filter := catalogue.Manifest{Module: "nftables", Version: "1", Tools: []string{"firewall_rules"},
Resources: []map[string]any{{"id": "runtime", "type": "container", "name": "mesh-nftables"}}}
stoodOn := []string{catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:" + strings.Repeat("d", 64)}
// Before the runtime exists the old pattern is accepted as it always was — and built, which is
// how the catalogue comes to know what the module stood on.
if err := inv.RegisterModule(ctx, filter, Source{Repository: "/r", Against: stoodOn}); err != nil {
t.Fatalf("before the runtime exists the old pattern is accepted: %v", err)
}
built := aBuild("nf1", "nftables", "")
built.Against = stoodOn
if err := inv.RecordBuild(ctx, built); err != nil {
t.Fatal(err)
}
runtime := catalogue.Manifest{Module: catalogue.RuntimeModule, Version: "1"}
if err := inv.RegisterModule(ctx, runtime, Source{Repository: "/r"}); err != nil {
t.Fatal(err)
}
// **A module already registered in that shape is rebuilt without complaint** (to-be 38 WP2.4 as
// amended by WP3): some thirty of them stand the day the runtime arrives, and each moves in its
// own change. The gate is against the pattern spreading, not against the pipeline running.
if err := inv.RegisterModule(ctx, filter, Source{Repository: "/r", Against: stoodOn}); err != nil {
t.Fatalf("a rebuild of a module that already had the pattern was refused: %v", err)
}
// A module new to the catalogue in that shape is refused, naming the record.
newcomer := filter
newcomer.Module = "lamp"
err := inv.RegisterModule(ctx, newcomer, Source{Repository: "/r", Against: stoodOn})
if err == nil || !strings.Contains(err.Error(), "ADR 0175") {
t.Fatalf("a new module in the old pattern was registered beside the runtime: %v", err)
}
// And a module that had moved its tools to a bundle may not come back to a container.
moved := filter
moved.Resources = nil
if err := inv.RegisterModule(ctx, moved, Source{Repository: "/r", Against: stoodOn}); err != nil {
t.Fatalf("a module whose tools are a bundle was refused: %v", err)
}
unbuilt := aBuild("nf2", "nftables", "")
if err := inv.RecordBuild(ctx, unbuilt); err != nil {
t.Fatal(err)
}
err = inv.RegisterModule(ctx, filter, Source{Repository: "/r", Against: stoodOn})
if err == nil || !strings.Contains(err.Error(), "ADR 0175") {
t.Fatalf("a module that had moved returned to the old pattern unrefused: %v", err)
}
}
+1 -17
View File
@@ -18,17 +18,8 @@ const (
EdgeBuiltBy = "built-by"
// EdgeDeclared: the manifest's own `build.on`.
EdgeDeclared = "declared"
// EdgeWorkerOf: the module holds the build seat, whose worker the control plane defines
// (novox/hq issue 206). The one place a *running* order enters the graph: a build machine rolled
// before the controller that redefines its worker cannot bind it, and nothing can then build the
// controller that would end that — so the holder of the build seat follows the controller, and
// the controller is built by whichever build machine is running, as it always was.
EdgeWorkerOf = "worker-of"
)
// TheControlPlane is the module that defines every seat's worker on the bus.
const TheControlPlane = "mesh-controller"
// Edge is one dependency: From depends on To, in the way Kind says.
type Edge struct {
From string `json:"from"`
@@ -72,7 +63,7 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin
if r := repositoryKey(e.Source.Repository); r != "" {
byRepository[r] = append(byRepository[r], name)
}
if e.Manifest.ClaimsSeat("node-build-agent") || e.Manifest.ClaimsSeat("mesh-build-machine") {
if e.Manifest.ClaimsSeat("mesh-build-machine") {
builders = append(builders, name)
}
}
@@ -115,13 +106,6 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin
}
}
}
if known[TheControlPlane] {
for _, b := range builders {
if b != TheControlPlane {
add(b, TheControlPlane, EdgeWorkerOf)
}
}
}
sort.Slice(out, func(a, b int) bool {
if out[a].From != out[b].From {
return out[a].From < out[b].From
+1 -1
View File
@@ -12,7 +12,7 @@ func TestDependenciesAreOneRelationWithTheirKinds(t *testing.T) {
return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}}
}
builder := entry("builder", "http://forge/novox/mesh-catalog.git")
builder.Manifest.Claims = []catalogue.Claim{{Name: "node-build-agent", Scope: catalogue.ScopeNode}}
builder.Manifest.Claims = []catalogue.Claim{{Name: "mesh-build-machine", Scope: catalogue.ScopeMesh}}
plugin := entry("shop-plugin", "http://forge/novox/mesh-catalog.git")
plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}}
entries := []Entry{
-32
View File
@@ -1,32 +0,0 @@
package link
import "testing"
// A build machine serves the seat its credential claims (novox/hq ADR 0190 handover): the old
// `builder` keeps the old role, a `build-agent` takes the new, from one binary and no flag.
func TestABuildMachineServesTheSeatItsCredentialClaims(t *testing.T) {
if got := BuildSeatClaimed([]string{"mesh-build-machine"}); got != "mesh-build-machine" {
t.Errorf("a credential claiming the old role serves %q", got)
}
if got := BuildSeatClaimed([]string{"node-build-agent"}); got != TheBuildMachine {
t.Errorf("a credential claiming the new role serves %q", got)
}
if got := BuildSeatClaimed(nil); got != TheBuildMachine {
t.Errorf("a credential claiming nothing serves %q, want the current role", got)
}
if got := BuildSeatClaimed([]string{"", "node-build-agent"}); got != TheBuildMachine {
t.Errorf("an empty claim is skipped; got %q", got)
}
}
// What a machine says about a build is the event of the seat it took the build from, so an outcome
// is heard where the asker of that seat listens.
func TestABuildsEventsAreItsSeats(t *testing.T) {
if BuildOutcomeOf(TheBuildMachineBefore) != "mesh.seat.mesh-build-machine.event.built" {
t.Error(BuildOutcomeOf(TheBuildMachineBefore))
}
if BuildWorkOf(TheBuildMachine) != BuildWork() || BuildOutcomeOf(TheBuildMachine) != BuildOutcome() ||
BuildStartedOf(TheBuildMachine) != BuildStarted() || BuildLogOf(TheBuildMachine, "b1") != BuildLog("b1") {
t.Error("the no-argument forms must name the current role")
}
}
+7 -53
View File
@@ -19,57 +19,13 @@ import (
// act on or a declaration a node reconciles toward; a build is a request that takes minutes and has
// exactly one answer. Too long for request/reply, too particular to be an event.
// TheBuildMachine is the role a build is submitted to: node-scoped, held on every machine that
// builds, and the work shared among them (novox/hq ADR 0190). The name stays for every caller; what
// it names moved from the mesh's one build machine to whichever build agent is idle.
//
// **Switching a live mesh over, in order** — and why no step strands a build. The old seat's
// stream and worker (SEAT_MESH_BUILD_MACHINE, SEAT_MESH_BUILD_MACHINE_worker) stay on the bus until
// removed by hand, and the builder keeps draining them while it is assigned, because a machine
// serves the seat its credential claims (BuildSeatClaimed) and the controller asks the seat that
// has a holder (buildSeatAmong in the command) and hears both seats' outcomes:
//
// 1. Merge the controller and the host's first user list together; the new controller rolls and,
// seeing only the builder assigned, still asks mesh-build-machine — which the builder holds.
// 2. Merge the catalogue's build-agent; the builder builds it and the controller registers it.
// 3. On each machine that builds: `module issue build-agent --node <n>`, then `assign`, then
// `push`. The first holder appears, and from then on asks go to node-build-agent.
// 4. Unassign builder everywhere and `module forget` it.
// 5. By hand: delete SEAT_MESH_BUILD_MACHINE and its worker, drop the retired seat row and
// TheBuildMachineBefore with it, and the second entries in seatsTheControllerAsks and
// ControllerFollows.
const TheBuildMachine = "node-build-agent"
// TheBuildMachineBefore is the role a build was submitted to until ADR 0190: the mesh's one build
// machine, mesh-scoped. Kept named while the handover runs — a machine whose credential claims it
// still serves it, and the controller still hears its outcomes — and dropped with the retired seat
// row once nothing claims it.
const TheBuildMachineBefore = "mesh-build-machine"
// BuildSeatClaimed is the build role a machine serves: the first seat its credential claims, or the
// current role when the credential names none (a credential from before claims travelled in it, or
// one written by hand). **The credential decides, not the binary** (ADR 0190 handover): one build
// machine binary runs as the old `builder` on the old seat and as a `build-agent` on the new one,
// each taking the work the mesh issued it a credential for, so neither drains the other's queue
// and the switch needs no flag day.
func BuildSeatClaimed(claimed []string) string {
for _, seat := range claimed {
if seat != "" {
return seat
}
}
return TheBuildMachine
}
// TheBuildMachine is the role a build is submitted to.
const TheBuildMachine = "mesh-build-machine"
// BuildWork is where a build request lands, and BuildOutcome is where its result does. Derived from
// the seat, so both sides name the role and neither names the other. The no-argument forms name the
// current role; the `Of` forms take the seat, for the handover during which two roles exist.
func BuildWork() string { return BuildWorkOf(TheBuildMachine) }
func BuildOutcome() string { return BuildOutcomeOf(TheBuildMachine) }
func BuildWorkOf(seat string) string { return "mesh.seat." + seat + ".accept.build" }
func BuildOutcomeOf(seat string) string {
return "mesh.seat." + seat + ".event.built"
}
// the seat, so both sides name the role and neither names the other.
func BuildWork() string { return "mesh.seat." + TheBuildMachine + ".accept.build" }
func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.built" }
// BuildStarted is where a build machine says it has taken a build, and BuildLog is where it says
// what it is doing, one line per message, under the build's own id (novox/hq ADR 0157).
@@ -79,10 +35,8 @@ func BuildOutcomeOf(seat string) string {
// lived in one container's stderr on one machine. Every line is now an event of the role, retained
// with the rest of the mesh's events, so a reader follows a build live by subscribing its subject,
// or reads it back afterwards from the stream, and a viewer is a subscriber and nothing more.
func BuildStarted() string { return BuildStartedOf(TheBuildMachine) }
func BuildLog(id string) string { return BuildLogOf(TheBuildMachine, id) }
func BuildStartedOf(seat string) string { return "mesh.seat." + seat + ".event.started" }
func BuildLogOf(seat, id string) string { return "mesh.seat." + seat + ".event.log." + id }
func BuildStarted() string { return "mesh.seat." + TheBuildMachine + ".event.started" }
func BuildLog(id string) string { return "mesh.seat." + TheBuildMachine + ".event.log." + id }
// BuildStart is what a build machine says the moment it takes a build.
type BuildStart struct {
+1 -1
View File
@@ -26,7 +26,7 @@ func TestTheOldBusAnnouncesABuildUnderBothNames(t *testing.T) {
if KeyRoleBuilt != "built" {
t.Fatalf("the role's event is %q, and a holder emits its verbs bare", KeyRoleBuilt)
}
if TheBuildMachine != "node-build-agent" {
if TheBuildMachine != "mesh-build-machine" {
t.Fatalf("the role is %q", TheBuildMachine)
}
// The two must differ, or one publish would serve both and this doubling would be pointless.
+24 -67
View File
@@ -25,34 +25,16 @@ import (
type natsBuilds struct {
js *broker.JetStream
owned bool
// seat is the build role asked: the one that has a holder (ADR 0190 handover), chosen by the
// controller from what is assigned, so an ask lands where a machine is pulling.
seat string
}
// BuildsOverNATS is the asking side on the bus being built, asking the current build role. It dials,
// because the command that asks for a build is a one-shot and holds nothing else.
// BuildsOverNATS is the asking side on the bus being built. It dials, because the command that asks
// for a build is a one-shot and holds nothing else.
func BuildsOverNATS(address string) (Builders, error) {
return BuildsOverNATSOn(address, TheBuildMachine)
}
// BuildsOverNATSOn is the asking side for one named build role — during the handover from the one
// build machine to build agents, the role that has a holder (ADR 0190).
func BuildsOverNATSOn(address, seat string) (Builders, error) {
js, err := broker.Dial(address)
if err != nil {
return nil, fmt.Errorf("cannot reach the bus at %s to ask for a build: %w", address, err)
}
return &natsBuilds{js: js, owned: true, seat: seat}, nil
}
// role is the seat asked: what the asker was made for, or the current build role for one made
// without saying (a test building the struct by hand).
func (b *natsBuilds) role() string {
if b.seat == "" {
return TheBuildMachine
}
return b.seat
return &natsBuilds{js: js, owned: true}, nil
}
func (b *natsBuilds) Close() {
@@ -69,7 +51,7 @@ func (b *natsBuilds) Ask(ctx context.Context, request BuildRequest) error {
}
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil {
if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil {
return fmt.Errorf("cannot submit a build: %w", err)
}
return nil
@@ -81,7 +63,7 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest,
// Subscribed before the ask, so an outcome cannot arrive before there is anywhere for it to
// land. Core, not the stream: the asker is waiting now, and the durable copy of this outcome is
// the same event on EVENTS, which the controller records.
outcomes, err := b.js.Conn().SubscribeSync(BuildOutcomeOf(b.role()))
outcomes, err := b.js.Conn().SubscribeSync(BuildOutcome())
if err != nil {
return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err)
}
@@ -98,7 +80,7 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest,
// be assumed, because nothing else will ever say so.
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil {
if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil {
return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err)
}
@@ -133,16 +115,9 @@ type natsMachine struct {
sub *nats.Subscription
}
// MachineOverNATS takes build work from the current build role.
// MachineOverNATS takes build work from the role this machine holds.
func MachineOverNATS(js *broker.JetStream, on string) BuildMachine {
return MachineOverNATSOn(js, on, TheBuildMachine)
}
// MachineOverNATSOn takes build work from the role named — the one this machine's credential claims
// (ADR 0190 handover): its asks come from that seat's worker, and what it says about a build goes
// out as that seat's events, so an outcome is heard where the asker listens.
func MachineOverNATSOn(js *broker.JetStream, on, seat string) BuildMachine {
return &natsMachine{js: js, on: on, seat: seat}
return &natsMachine{js: js, on: on, seat: TheBuildMachine}
}
func (m *natsMachine) Close() {
@@ -151,31 +126,29 @@ func (m *natsMachine) Close() {
}
}
// Take binds to the role's worker and pulls one request at a time, handing each over.
// Take binds to the role's worker and hands each request over, one at a time.
//
// **Bound, never created.** The work queue and the worker on it are the controller's to define
// (design 25 §3), and a build machine reaches no part of the JetStream API — so a missing one is said
// as the mesh's to answer rather than quietly created with whatever this client defaults to.
//
// **Pulled, one at a time, by whichever holder is free** (novox/hq ADR 0190). Every machine holding
// the role binds this same worker; a machine asks for the next request only when it has finished
// the last, so a slow machine never holds an ask an idle one could take, and a machine that took
// five at once would run five container builds against one runtime and finish all of them slower
// than the first.
func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) error {
worker, found := broker.HolderConsumerFor(m.on, "build-agent",
worker, found := broker.HolderConsumerFor(m.on, "builder",
broker.DeclaredSeat{Name: m.seat, Accepts: []string{"build"}})
if !found {
return fmt.Errorf("%s accepts no work, so there is nothing for this machine to take", m.seat)
}
// One at a time, which the consumer's own ack-pending limit enforces rather than a prefetch
// setting: a machine that took five requests at once would run five container builds against one
// runtime and finish all of them slower than the first.
work := make(chan *nats.Msg, 1)
// **The consumer's own filter, not the one subject this machine cares about.** The client checks
// what is asked for against the consumer's filter and refuses anything that is not the same —
// "subject does not match consumer" — so subscribing `…accept.build` against a consumer filtered
// on `…accept.>` is rejected even though it is narrower. Learned twice now, on two different
// consumers, which is why it is written down here.
filter := worker.Filters[0]
sub, err := m.js.Context().PullSubscribe(filter, worker.Name,
sub, err := m.js.Context().ChanQueueSubscribe(filter, worker.Queue, work,
nats.Bind(worker.Stream, worker.Name), nats.ManualAck())
if err != nil {
return fmt.Errorf(
@@ -186,27 +159,13 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
m.sub = sub
for {
if ctx.Err() != nil {
select {
case <-ctx.Done():
return nil
case msg, ok := <-work:
if !ok {
return errors.New("the bus stopped delivering build work")
}
// One, and wait a while for it; an empty queue is a timeout, which is the normal state of a
// machine with nothing to build, and is asked again.
fetched, err := sub.Fetch(1, nats.Context(ctx))
switch {
case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded):
return nil
case errors.Is(err, nats.ErrTimeout):
continue
case err != nil:
if sub.IsValid() {
// A transient fault in asking — a reconnect, a slow server — is asked past rather
// than ending the machine; one that outlasts the ack wait redelivers nothing lost.
time.Sleep(time.Second)
continue
}
return fmt.Errorf("the bus stopped delivering build work: %w", err)
}
for _, msg := range fetched {
var request BuildRequest
if err := json.Unmarshal(msg.Data, &request); err != nil {
// Unreadable: terminated rather than retried, because the next attempt reads the same
@@ -219,7 +178,7 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
// the ask to a second machine nor counts the wait against its deliveries.
working := make(chan struct{})
go stillWorking(msg, working)
do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js, seat: m.seat})
do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js})
close(working)
}
}
@@ -230,8 +189,6 @@ type natsBuild struct {
msg *nats.Msg
on string
js *broker.JetStream
// seat is the role this build was taken from; what the machine says about it is that role's.
seat string
seq int
}
@@ -259,7 +216,7 @@ func (b *natsBuild) Announce(ctx context.Context, result BuildResult) error {
}
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
if _, err := b.js.Context().Publish(BuildOutcomeOf(b.seat), body, nats.Context(publish)); err != nil {
if _, err := b.js.Context().Publish(BuildOutcome(), body, nats.Context(publish)); err != nil {
return fmt.Errorf("cannot announce a build's outcome: %w", err)
}
return nil
@@ -279,7 +236,7 @@ func (b *natsBuild) Began(ctx context.Context) error {
}
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
if _, err := b.js.Context().Publish(BuildStartedOf(b.seat), body, nats.Context(publish)); err != nil {
if _, err := b.js.Context().Publish(BuildStarted(), body, nats.Context(publish)); err != nil {
return fmt.Errorf("cannot say a build started: %w", err)
}
return nil
@@ -297,7 +254,7 @@ func (b *natsBuild) Say(step, message string) {
if err != nil {
return
}
_ = b.js.Conn().Publish(BuildLogOf(b.seat, b.request.ID), body)
_ = b.js.Conn().Publish(BuildLog(b.request.ID), body)
}
func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) }
+3 -119
View File
@@ -48,7 +48,7 @@ func aBusWithTheBuildRole(t *testing.T) *broker.JetStream {
t.Fatal(err)
}
clean := func() {
_ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT")
_ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE")
for _, s := range broker.MeshStreams() {
_ = js.Context().PurgeStream(s.Name)
}
@@ -158,7 +158,7 @@ func TestNatsABuildIsTakenAndItsOutcomeReachesEverybody(t *testing.T) {
// And the work left the queue: a request a machine took and settled must not be given to another.
deadline := time.Now().Add(5 * time.Second)
for time.Now().Before(deadline) {
info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT")
info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE")
if err == nil && info.State.Msgs == 0 {
return
}
@@ -177,7 +177,7 @@ func TestNatsABuildWaitsForAMachineRatherThanFailing(t *testing.T) {
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
t.Fatal(err)
}
info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT")
info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE")
if err != nil || info.State.Msgs != 1 {
t.Fatalf("the work did not queue: %+v %v", info, err)
}
@@ -245,119 +245,3 @@ func TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue(t *testing.T) {
func quietLog() *log.Logger { return log.New(io.Discard, "", 0) }
var _ = quietLog
// Two machines holding the role share one queue (novox/hq ADR 0190): three asks, each machine takes
// one and the third waits until one of them is done; an ask is never handed to a machine that is
// busy; and a machine that stops mid-ask leaves its ask to the other.
func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testing.T) {
js := aBusWithTheBuildRole(t)
ctx, stop := context.WithCancel(context.Background())
defer stop()
for _, id := range []string{"w-1", "w-2", "w-3"} {
body, _ := json.Marshal(BuildRequest{ID: id, Repository: "/r"})
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
t.Fatal(err)
}
}
type taken struct{ machine, id string }
took := make(chan taken, 8)
release := map[string]chan struct{}{"anchor": make(chan struct{}), "laptop": make(chan struct{})}
machines := map[string]BuildMachine{}
for _, name := range []string{"anchor", "laptop"} {
name := name
m := MachineOverNATS(js, name)
machines[name] = m
defer m.Close()
go func() {
_ = m.Take(ctx, func(ctx context.Context, work Build) {
took <- taken{name, work.Request().ID}
<-release[name]
_ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: name})
_ = work.Done()
})
}()
}
// Each machine took exactly one, and they are different asks.
first := map[string]string{}
for i := 0; i < 2; i++ {
select {
case got := <-took:
if _, twice := first[got.machine]; twice {
t.Fatalf("%s was handed a second ask while busy with its first", got.machine)
}
first[got.machine] = got.id
case <-time.After(10 * time.Second):
t.Fatalf("only %d machine(s) took work; two idle holders should both have", len(first))
}
}
if first["anchor"] == first["laptop"] {
t.Fatalf("both machines took %q: the queue is not shared, it is copied", first["anchor"])
}
// The third waits: nobody is free.
select {
case got := <-took:
t.Fatalf("%s was handed %s while both machines were busy", got.machine, got.id)
case <-time.After(2 * time.Second):
}
// One finishes, and only then is the third taken — by that machine, the one that is free.
close(release["anchor"])
release["anchor"] = make(chan struct{})
select {
case got := <-took:
if got.machine != "anchor" {
t.Fatalf("the third ask went to %s, which is still busy", got.machine)
}
case <-time.After(10 * time.Second):
t.Fatal("the third ask was never taken after a machine became free")
}
// A machine that stops mid-ask leaves its ask unacknowledged, and the ack wait brings it round
// to whoever is left — the path TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue proves with
// an explicit hand-back, because the real wait is a minute. Here: the laptop goes, anchor
// finishes, and with nothing queued nothing more is taken by the machine that is left.
machines["laptop"].Close()
close(release["anchor"])
select {
case got := <-took:
t.Fatalf("%s took %s; the queue should be empty", got.machine, got.id)
case <-time.After(2 * time.Second):
}
}
// During the handover (ADR 0190) two build roles exist. A machine whose credential claims the retired
// one takes an ask published to that seat and answers as that seat; the asker of that seat hears it.
func TestNatsAMachineOnTheRetiredBuildRoleTakesThatRolesAsks(t *testing.T) {
js := aBusWithTheBuildRole(t)
seats := []broker.DeclaredSeat{{Name: TheBuildMachineBefore, Accepts: []string{"build"}, Emits: []string{"built"}}}
if err := broker.RaiseSeats(js, seats, map[string]broker.Holder{TheBuildMachineBefore: {Node: "anchor", Module: "builder"}}); err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") })
ctx, stop := context.WithCancel(context.Background())
defer stop()
machine := MachineOverNATSOn(js, "anchor", TheBuildMachineBefore)
defer machine.Close()
go func() {
_ = machine.Take(ctx, func(ctx context.Context, work Build) {
_ = work.Began(ctx)
_ = work.Announce(ctx, BuildResult{ID: work.Request().ID, Repository: work.Request().Repository, On: "anchor", Commit: "abc"})
_ = work.Done()
})
}()
asker, err := BuildsOverNATSOn(os.Getenv("MESH_TEST_NATS"), TheBuildMachineBefore)
if err != nil {
t.Fatal(err)
}
defer asker.Close()
result, err := asker.Submit(ctx, BuildRequest{ID: "build-old-seat", Repository: "r"}, 20*time.Second)
if err != nil {
t.Fatal(err)
}
if result.On != "anchor" || result.ID != "build-old-seat" {
t.Errorf("the retired role's holder did not answer: %+v", result)
}
}
+2 -3
View File
@@ -223,11 +223,10 @@ func kindOfSubject(subject string) (string, bool) {
return KindCatchUp, true
case broker.ControllerFollows[3]:
return KindSourceMoved, true
case BuildOutcome(), BuildOutcomeOf(TheBuildMachineBefore):
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
// with it did not change (novox/hq ADR 0121). From either build role while the handover
// runs (ADR 0190): the old builder still answers on the retired seat until it is unassigned.
// with it did not change (novox/hq ADR 0121).
return KindBuilt, true
}
return "", false