Compare commits

..
Author SHA1 Message Date
mesh-admin 78de54381b Merge pull request 'A seat's work is shared by its holders: node-build-agent, pulled one ask at a time (hq ADR 0190)' (#228) from feat/a-seats-work-is-shared-by-its-holders into main 2026-10-03 00:53:42 +00:00
jochen ff5ef0ab60 The controller asks the build role that has a holder, and hears both roles' outcomes (hq ADR 0190, the handover)
A controller that asked node-build-agent from its first run would queue every build where nothing
pulls, and the build that registers build-agent — the first holder — would be among them. So the
role is chosen at ask time from the catalogue: the current role when any assigned module claims it,
the retired one while only the builder does, the current one when neither. Outcomes are followed on
both seats, the controller may publish to both, and a build's log is read under whichever role did
it; a machine on the retired role is proven on the bus to take that role's asks. The switch order
is written where the role is named, and the retired half is marked for removal with the seat row.
2026-10-03 02:51:03 +02:00
jochen a5d6a1187c A build machine serves the seat its credential claims (hq ADR 0190, the handover)
After the build role moved to node-build-agent, nothing would hold it until build-agent is
registered — and registering build-agent needs a build outcome that only the running builder
could produce, bound as it was to the old seat by name. One binary, two roles: the seat a machine
serves is the first its credential claims, as the mesh writes the claims beside the credential it
issues (ADR 0159); the old builder keeps draining mesh-build-machine, a build-agent takes
node-build-agent, and what each says about a build goes out as that seat's events, so an outcome
is heard where the asker of that seat listens. A credential naming no claim serves the current role.
2026-10-03 02:47:49 +02:00
mesh-admin 06d0a4bfe8 Merge pull request 'A changed jail filter restarts fail2ban' (#231) from jschoubben/jail-filter-restart into main 2026-10-02 21:28:32 +00:00
jschoubben d293a0deaf A changed jail filter restarts fail2ban
fail2ban restarts when the composed jail file changes, and each filter is
a file of its own, so a module that changed only its failregex left the
running jail on the old pattern. The jail file now names each filter's
digest.
2026-10-02 23:28:25 +02:00
mesh-admin 11a44e1ec4 Merge pull request 'A resolved manifest keeps a built reference in the store's own form, never the address it was reached by (hq ADR 0155)' (#230) from fix/a-resolved-reference-is-kept-not-routed into main 2026-10-02 21:14:06 +00:00
jochen 28c7f05b39 A resolved manifest keeps a built reference in the store's own form, never the address it was reached by (hq ADR 0155)
The first bundle resolved on the mesh carried the store's host in bundles[0].source, and registration
refused node-tools as naming an installation — rightly. The build record already keeps the
store-relative form; the resolved manifest now keeps the same for bundles and archive resources,
and composition routes it through the store a machine reaches, as it already did for a kept one.
2026-10-02 23:13:22 +02:00
mesh-admin 93a0c6202e Merge pull request 'The bus is never public: the broker port is no longer a foundation opening (hq ADR 0169)' (#229) from jschoubben/the-bus-is-never-public into main 2026-10-02 20:52:28 +00:00
jschoubben 7d82751862 The bus is never public: the broker port is no longer a foundation opening
The controller widened the bus's from-mesh port to from-anywhere on the
broker's host so a machine could enrol before it had a tunnel. ADR 0169
has machines join through the tunnel and decides the bus is never public;
every live bus connection already comes from the mesh.
2026-10-02 22:52:22 +02:00
jochen 905f3363c9 Two machines holding the build role share one queue, and neither is handed an ask while busy (hq ADR 0190)
Against a real bus: three asks, two machines; each takes one, the third waits until one is free
and then goes to that one; a machine that stops leaves nothing taken twice. The redelivery of an
ask a dead machine held is the ack wait's, proven by the hand-back test beside this one.

And the order a live mesh switches over in, written where the role is named: queued builds first,
then this controller, then build-agent assigned where machines build, then the builder and the old
seat's stream forgotten.
2026-10-02 22:35:39 +02:00
jochen 9f9d9b3b25 A seat's holders pull one ask at a time from one shared worker (hq ADR 0190, issue 186)
The worker a holder bound was a push consumer in a queue group with one ask in flight: right for
one holder, and with two it would still be a queue of one — the server hands a pushed ask to
whichever subscriber it picks, busy or not, and the in-flight cap is per consumer, not per holder.
Now the worker is pulled: every machine holding the seat binds the same durable and fetches one
ask when it has finished the last, so an idle machine is the one that takes the next, the asks in
flight are bounded by the holders working, and nothing is delivered that nobody asked for — which
is also what ended the race issue 186 describes. A holder's grants trade the delivery subject for
MSG.NEXT on the worker; the ack grant and the heartbeat that keeps a long build alive stay.

Proven against a real bus: the build round trip, a backlog taken by a machine that arrives later,
and work handed back by one machine coming round again.
2026-10-02 22:34:44 +02:00
jochen bde4b61b3b The build role is the node-scoped seat node-build-agent, and its work is shared by every holder (hq ADR 0190)
One build machine built everything, in a queue of one, because the seat was mesh-scoped and a
mesh seat has one holder. ADR 0190 makes building a node role: node-build-agent, held on every
machine that builds, with the work asked of the role and taken by whichever holder is idle. The
work subject of a node-scoped seat carries no node — that token is for a seat's tools, asked of
one machine (design 33 §4) — so holders on several machines read one queue; a test now says so.

The retired mesh-build-machine row stays while the builder module's registered manifest claims
it: a claim to a seat the mesh no longer defines is refused, and the machine holding it would be
unresolvable until build-agent replaces it. Removed once no manifest claims it.

The installer's genesis template (in the host's repository) still grants the controller the old
seat's subjects; its test here says so until that template names node-build-agent.
2026-10-02 22:31:55 +02:00
mesh-admin d07018f3c5 Merge pull request 'WP3: a TypeScript bundle carries what it runs with, the runtime's credential belongs to its account, and the gate refuses spreading not standing (hq to-be 38)' (#226) from feat/wp3-node-tools into main 2026-10-02 19:57:21 +00:00
jochen 729a5f9e6c The gate refuses the old tool-container pattern spreading, not a rebuild of what already stands (hq to-be 38 WP2.4, amended by WP3)
Once the runtime module is registered, a module serving tools from a container built on the
runtime's image is refused at registration — as written, including every rebuild of the thirty-odd
modules already in that shape, from the day the runtime arrives until WP4 onward moves each one.
That would stop the catalogue's whole pipeline to make a point the record already makes. Now a
module already registered in that shape — judged from the manifest the catalogue holds and what
its newest build stood on, the same two things a new registration is judged by — is rebuilt as
before; a module new to the catalogue in that shape, or one that had moved to a bundle and comes
back, is refused naming the record.
2026-10-02 21:44:44 +02:00
jochen 773b561f5e The runtime's credential belongs to the account it runs as (hq to-be 38 WP3)
The runtime's process is composed `user: <account>` where the node has one, and its broker file was
root's at 0600: a credential the process could not read. Composed in the declaration 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 there the runtime runs as root and the file stays root's.
2026-10-02 21:44:44 +02:00
jochen ca7e81e964 A TypeScript bundle carries what it runs with: the toolchain's runtime directory is copied into it (hq to-be 38 WP3, ADR 0188)
A bundle that compiled was not yet a bundle that ran. The compiler resolved `import "nats"` from
the toolchain image's own node_modules and the pack took only what the compiler wrote, so what a
machine unpacked could not find a single dependency — and Node would have read the bare `.js` as
CommonJS besides. No TypeScript bundle had run live to show it; the runtime's own is the first that
must. A toolchain now names a Dependencies directory in its image, copied whole into the output's
root after the compile by a second run in the same image: for TypeScript /app/runtime, which the
runtime's image puts a `"type": "module"` package.json and its pruned node_modules at. An older
image without it fails the build by name rather than packing a bundle that starts nowhere. The
SDK's and the runtime's dependencies, nothing module-specific yet: a skeleton, by ADR 0188 §5.
2026-10-02 21:44:44 +02:00
mesh-admin 08bb56f1d3 Merge pull request 'The controller composes one tool runtime per node: its principal, the bundles, its process, and the gate (hq ADR 0175, to-be 38 WP2)' (#223) from feat/the-operators-machine into main 2026-10-02 19:27:52 +00:00
jochen 8eb6c3e94a A tool container on the runtime's image is refused once the runtime is registered (hq ADR 0175, to-be 38 WP2.4)
A module declaring tools, a container, and a build on mesh-tools' runtime image is a container whose
purpose is serving tools — the pattern the node's tool runtime retires. Once node-tools is in the
catalogue, registering one is refused by name, with the record that says why; before, it is accepted
as it always was, so a mesh converts in the design's order and nothing is refused before there is
anything to move to. This is the mechanism that keeps the old pattern from returning by habit.

Judged from a repository manifest's own build.on, and for a built manifest — which carries no build
— from what its build stood on, now recorded beside the commit as part of a module's provenance.
2026-10-02 18:30:14 +02:00
jochen 9d4d847dc4 The machine runs one tool runtime, loading every delivered bundle (hq ADR 0175, to-be 38 WP2.3)
Where node-tools is in a node's set, the declaration ends with one process: the runtime module's own
bundle, run from its one entrypoint by its language's interpreter, told in MESH_TOOL_MODULES every
<module>=<file> the machine's bundles load, where its credential is (the module's own broker secret
as this node places it), and — on a machine with an operator account — who the operator is, running
as that account so a tool that needs root can escalate as the operator would. Restarted when any
bundle it loads or the credential changes. A machine with no account runs it as root without the two
operator words; a machine without the runtime is sent nothing new.

A bundle says which of its entrypoints the runtime LOADS (`loads`), because one bundle may carry a
daemon beside its tools and importing the daemon into the runtime would start it there; absent, a
module declaring tools has every entrypoint loaded. And the TypeScript toolchain is rooted at the
module, so an entrypoint lands at the path it is named by — the runtime loading bundles by their
declared paths is what made the compiler's common-directory default visible.
2026-10-02 18:26:54 +02:00
jochen b1df688c62 A module's tools bundle reaches the machine as an archive where the runtime runs (hq ADR 0175, to-be 38 WP2.2)
The resolved manifest now carries what the build compiled — each bundle's source, digest, language
and entrypoints — because a tools bundle is named by no resource of the module's own: the node's
runtime loads it, and until this the mesh held no trace of the one artifact that runtime needs. A
repository manifest that writes `bundles` beside its build is refused: the mesh derives it.

Where the runtime module is in a node's set, every assigned module's bundle with entrypoints is
composed as an archive under the mesh's own directory, routed through the artifact store like any
image or archive the mesh built. A bundle without entrypoints is run rather than loaded and is
delivered by the process that runs it. A node without the runtime is sent exactly what it was.
2026-10-02 18:19:10 +02:00
jochen 1f86ed4135 One runtime principal per node carries every assigned module's tools (hq ADR 0175, to-be 38 WP2.1)
Where the node-tools module is assigned, the machine's bus user list gains one principal of kind
node-tools in place of that module's own: it may subscribe every carried module's tool namespace
and every held seat's verbs on its node, read and follow every membership on its node, call any
tool anywhere, answer what it is asked — and consume nothing, because tools are what it runs.
Every other module keeps its own principal, so a module still serving tools from its container
holds its own credential until it moves.

Named exactly as the module it stands for, so `module issue` and `rollout mint` deliver its
credential through the path a module's already takes, into node-tools' own `broker` secret. The
runtime module's name is one constant in each of the broker and catalogue packages, held to one
string by the agreement test, because a rule turns on it.
2026-10-02 18:16:14 +02:00
55 changed files with 1818 additions and 1686 deletions
+21 -1
View File
@@ -137,7 +137,12 @@ func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
if err != nil { if err != nil {
return nil, err return nil, err
} }
return link.MachineOverNATS(js, on), nil // **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
} }
// answer does one build and says what happened, whichever way it went. // answer does one build and says what happened, whichever way it went.
@@ -453,6 +458,21 @@ type Credential struct {
// as two fields and this machine joins them once, here, to dial. // as two fields and this machine joins them once, here, to dial.
User string `json:"user,omitempty"` User string `json:"user,omitempty"`
Password string `json:"password,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 // onTheNewBus is whether a credential is for the bus being built: its address says so, and the
+27
View File
@@ -0,0 +1,27 @@
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)
}
}
+57 -10
View File
@@ -430,11 +430,13 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
} }
fmt.Println() fmt.Println()
ask, err := askOver(server) seat := buildSeatHeld(ctx)
ask, err := askOverOn(seat)
if err != nil { if err != nil {
return err return err
} }
defer ask.Close() defer ask.Close()
fmt.Printf(" of %s\n", seat)
if wait == 0 { if wait == 0 {
// Asked and not waited for (novox/hq issue 176): the outcome is the role's event, and the // Asked and not waited for (novox/hq issue 176): the outcome is the role's event, and the
@@ -522,6 +524,8 @@ func takeIn(ctx context.Context, inv *inventory.Inventory, result link.BuildResu
recorded := inventory.Source{ recorded := inventory.Source{
Repository: result.Repository, Path: result.Path, Ref: result.Ref, Repository: result.Repository, Path: result.Path, Ref: result.Ref,
BuiltFrom: result.Commit, Head: result.Commit, 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 != "" { if result.Source != nil && result.Source.Seat != "" {
recorded.Repository, recorded.Seat = result.Source.Repository, result.Source.Seat recorded.Repository, recorded.Seat = result.Source.Repository, result.Source.Seat
@@ -533,10 +537,6 @@ func takeIn(ctx context.Context, inv *inventory.Inventory, result link.BuildResu
if err := inv.RegisterModule(ctx, manifest, recorded); err != nil { if err := inv.RegisterModule(ctx, manifest, recorded); err != nil {
return manifest, kept, err return manifest, kept, err
} }
// The keep set just moved, and new bytes just landed (novox/hq ADR 0189). Asked here rather
// than on a timer of its own: this is the only moment either is true. Never fatal — the build
// worked and the module is registered.
collect(ctx, inv)
return manifest, kept, nil return manifest, kept, nil
} }
@@ -557,7 +557,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
} }
defer server.Close() defer server.Close()
ask, err := askOver(server) ask, err := askOverOn(buildSeatHeld(ctx))
if err != nil { if err != nil {
return err return err
} }
@@ -670,12 +670,56 @@ 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 // **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 // 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. // being built it dials, because a build request is a one-shot and holds nothing else.
func askOver(_ *link.Server) (link.Builders, error) { func askOverOn(seat string) (link.Builders, error) {
address, err := broker.BusAddress() address, err := broker.BusAddress()
if err != nil { if err != nil {
return nil, err return nil, err
} }
return link.BuildsOverNATS(address) 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
} }
// buildLog prints everything a build machine said about one build, read back from the bus. // buildLog prints everything a build machine said about one build, read back from the bus.
@@ -695,10 +739,13 @@ func buildLog(ctx context.Context, id string) error {
} }
defer js.Close() defer js.Close()
sub, err := js.Context().PullSubscribe(link.BuildLog(id), "", // 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, "",
nats.BindStream(broker.EventsStream), nats.DeliverAll(), nats.AckNone()) nats.BindStream(broker.EventsStream), nats.DeliverAll(), nats.AckNone())
if err != nil { if err != nil {
return fmt.Errorf("cannot read %s from the bus: %w", link.BuildLog(id), err) return fmt.Errorf("cannot read %s from the bus: %w", lines, err)
} }
defer func() { _ = sub.Unsubscribe() }() defer func() { _ = sub.Unsubscribe() }()
+43
View File
@@ -0,0 +1,43 @@
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)
}
}
-84
View File
@@ -1,84 +0,0 @@
package main
import (
"context"
"errors"
"fmt"
"os"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/inventory"
)
// Letting the artifact store go of what the mesh no longer keeps (novox/hq ADR 0189, issue 108).
//
// **Run where the records change.** A build is the moment new bytes landed in the store and the
// moment the keep set moved, so it is the moment to say what may go — and it needs no timer of
// its own. Reclaiming the bytes is the store's own nightly step; this only decides.
//
// Never fatal to a build. The build succeeded, the module is registered, and a store that could
// not be reached is a thing to say rather than a reason to undo any of that. The next build asks
// again, and the references it could not collect are still uncollected, so nothing is lost by
// having failed.
// collect asks the store to let go of everything the mesh made and no longer keeps, and records
// what it let go of. Says what it did and what it could not; returns nothing, because nothing
// upstream should branch on it.
func collect(ctx context.Context, inv *inventory.Inventory) {
references, err := inv.ToCollect(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "could not work out what the artifact store may let go of: %v\n", err)
return
}
if len(references) == 0 {
return
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "could not read the catalogue to find the artifact store: %v\n", err)
return
}
// As the mesh reaches it from the network. Empty means the store is not on the network — on a
// mesh being raised it is not yet, and there the store holds one build of anything and has
// nothing to collect.
address, err := artifactStoreAddress(ctx, inv, shelf, "")
if err != nil || address == "" {
if err != nil {
fmt.Fprintf(os.Stderr, "could not find the artifact store to collect from: %v\n", err)
}
return
}
store := artifacts.Store{Address: address}
var done []string
var refused int
for _, reference := range references {
switch err := store.LetGo(ctx, reference); {
case err == nil, errors.Is(err, artifacts.Gone):
// Gone is the outcome wanted, already true. Recorded so the next sweep does not ask
// again for ever.
done = append(done, reference)
default:
refused++
if refused == 1 {
// Once per sweep. A store that refuses one refuses all of them, and a hundred
// identical lines would bury the reason.
fmt.Fprintf(os.Stderr, "the artifact store kept %s: %v\n", reference, err)
}
}
}
if len(done) > 0 {
if err := inv.MarkCollected(ctx, done); err != nil {
// Said, and that is all: the artifacts are gone either way, and the only cost of an
// unrecorded collection is that the next sweep asks about them again.
fmt.Fprintf(os.Stderr, "the store let go of %d artifact(s) and the record of it did not keep: %v\n",
len(done), err)
return
}
fmt.Fprintf(os.Stderr, "the artifact store let go of %d artifact(s) the mesh no longer keeps\n",
len(done))
}
if refused > 0 {
fmt.Fprintf(os.Stderr, "%d artifact(s) were not collected; the next build asks again\n", refused)
}
}
@@ -1,33 +0,0 @@
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)
}
}
+11 -1
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() user := broker.Principal{Kind: broker.KindModule, Node: node, Module: m.Module}.Username()
password, err := inv.MintBusPassword(ctx, inventory.BusUser{ password, err := inv.MintBusPassword(ctx, inventory.BusUser{
Username: user, Kind: inventory.BusModule, Node: node, Module: m.Module, Username: user, Kind: busKindOf(m.Module), Node: node, Module: m.Module,
}) })
if err != nil { if err != nil {
return err return err
@@ -577,6 +577,16 @@ func issueOnTheNewBus(ctx context.Context, inv *inventory.Inventory, m catalogue
return issueWith(ctx, inv, m, node, busAddress, known, reachable, user, password) 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 // 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 // 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 // so the move can issue every module against a bus whose address it worked out itself
+2 -1
View File
@@ -1,6 +1,7 @@
package main package main
import ( import (
"reflect"
"strings" "strings"
"testing" "testing"
"time" "time"
@@ -191,7 +192,7 @@ func TestWhatAHandedOverModuleRecordsAboutItsSource(t *testing.T) {
t.Fatalf("the source records as %+v", from) t.Fatalf("the source records as %+v", from)
} }
// A manifest with no provenance at all is legitimate: fixing something in a hurry. // A manifest with no provenance at all is legitimate: fixing something in a hurry.
if from, err := whereItComesFrom("", "", "", "", false); err != nil || from != (inventory.Source{}) { if from, err := whereItComesFrom("", "", "", "", false); err != nil || !reflect.DeepEqual(from, inventory.Source{}) {
t.Fatalf("a manifest handed over with no provenance was refused: %+v, %v", from, err) t.Fatalf("a manifest handed over with no provenance was refused: %+v, %v", from, err)
} }
for _, c := range []struct { for _, c := range []struct {
+5 -37
View File
@@ -16,8 +16,6 @@ import (
"github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/licences" "github.com/novox/mesh-controller/internal/licences"
"github.com/novox/mesh-controller/internal/overlay" "github.com/novox/mesh-controller/internal/overlay"
"net"
"strconv"
) )
// working out what one machine should be. // working out what one machine should be.
@@ -648,25 +646,12 @@ func renderingFor(ctx context.Context, open *stores, node string,
names[name] = at names[name] = at
} }
// The ports the mesh itself needs open, which no module declares. Read from the broker this // **The bus is never public** (novox/hq ADR 0169). It was a foundation port — widened from the
// control plane was told about rather than written down twice: the address a node is handed in // broker's own `from: mesh` to from-anywhere on the broker's host, so a machine could enrol
// its token and the port its machine must accept on are the same fact. // 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.
// **Only on the node that listens on it** (novox/hq issue: the broker opening leaked onto // Nothing the mesh itself needs is opened beyond what a module declares.
// 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 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 // 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 // copy, outside the store (novox/hq ADR 0085, amended). Read only; nothing here mints. The
@@ -1380,23 +1365,6 @@ func composeBusUsers(ctx context.Context, inv *inventory.Inventory,
return broker.ComposeAccounts(filled) 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 // 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. // 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 { func providerModuleOf(resolved catalogue.Resolution, open *stores, ctx context.Context, n catalogue.Needed) string {
+4 -2
View File
@@ -346,7 +346,9 @@ func rolloutMint(ctx context.Context, again bool) error {
} }
machines++ machines++
case broker.KindModule: 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.
if p.Module == "mesh-controller" { if p.Module == "mesh-controller" {
// The control plane is a module too, and its `broker` secret is the old bus's // 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 // credential it is still using while this runs. Writing the new bus's blob there
@@ -365,7 +367,7 @@ func rolloutMint(ctx context.Context, again bool) error {
skipped++ skipped++
continue continue
} }
password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: inventory.BusModule, Node: p.Node, Module: p.Module}) password, err := inv.MintBusPassword(ctx, inventory.BusUser{Username: p.Username(), Kind: busKindOf(p.Module), Node: p.Node, Module: p.Module})
if err != nil { if err != nil {
return err return err
} }
-99
View File
@@ -1,99 +0,0 @@
// Package artifacts speaks to the mesh's artifact store over its own door.
//
// Only what the mesh needs that nothing else does: letting go of something it put there
// (novox/hq ADR 0189, issue 108). Pushing is the builder's, through the container runtime; reading
// is every machine's, through its runtime. This is the one operation that belongs to the thing
// holding the records, because it is the only one that is a decision rather than a transfer.
package artifacts
import (
"context"
"fmt"
"net/http"
"strings"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
)
// Store is the artifact store at an address, as this machine reaches it.
type Store struct {
// Address is `host:port` — the store as the caller reaches it now, composed and never
// recorded (novox/hq 04-ISSUES/102).
Address string
// HTTP is the client used; nil is a client with a modest timeout.
HTTP *http.Client
}
// Gone is the answer when the store does not hold it: the outcome wanted, already true.
var Gone = fmt.Errorf("the store does not hold it")
// LetGo asks the store to drop one artifact the mesh recorded making.
//
// Takes a reference as the mesh records it — `artifact-store://<module>/<artifact>@sha256:…` for
// an image, `…/blobs/sha256:…` for an archive — because that is the identity every record uses,
// and composes the address here at the moment of use.
//
// Returns Gone when the store answers that it does not have it. That is not a failure: the sweep
// wants the artifact absent, and it is. It is distinguished from success only so a caller can say
// which of the two happened.
func (s Store) LetGo(ctx context.Context, reference string) error {
path, kept := catalogue.InArtifactStore(reference)
if !kept {
// Nothing the mesh put in its own store. Refused rather than attempted: composing a
// delete for a reference of unknown shape is how a sweep reaches something that is not
// the mesh's.
return fmt.Errorf("%s is not a reference into the mesh's artifact store", reference)
}
if s.Address == "" {
return fmt.Errorf("this mesh has no artifact store on its network to ask about %s", reference)
}
repository, kind, digest, err := split(path)
if err != nil {
return err
}
url := "http://" + s.Address + "/v2/" + repository + "/" + kind + "/" + digest
request, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil)
if err != nil {
return err
}
client := s.HTTP
if client == nil {
client = &http.Client{Timeout: 30 * time.Second}
}
response, err := client.Do(request)
if err != nil {
return err
}
defer response.Body.Close()
switch response.StatusCode {
case http.StatusAccepted, http.StatusOK, http.StatusNoContent:
return nil
case http.StatusNotFound:
return Gone
case http.StatusMethodNotAllowed:
// The registry was started without deletion enabled. Said plainly, because the remedy is
// a setting on the store's module and not anything about this artifact.
return fmt.Errorf(
"the artifact store refuses deletion: its server was started without it enabled "+
"(REGISTRY_STORAGE_DELETE_ENABLED), so nothing can be collected until the store "+
"module is applied again (novox/hq ADR 0189). Asking about %s", reference)
default:
return fmt.Errorf("the artifact store answered %s for %s", response.Status, reference)
}
}
// split reads a recorded path into the repository, which endpoint names the thing, and the digest.
//
// Two shapes, which are the two the mesh records: `<repository>@sha256:<hex>` is a manifest, and
// `<repository>/blobs/sha256:<hex>` is a blob.
func split(path string) (repository, kind, digest string, err error) {
if before, after, ok := strings.Cut(path, "@sha256:"); ok {
return before, "manifests", "sha256:" + after, nil
}
if before, after, ok := strings.Cut(path, "/blobs/sha256:"); ok {
return before, "blobs", "sha256:" + after, nil
}
return "", "", "", fmt.Errorf("%q names nothing the store holds by digest", path)
}
-94
View File
@@ -1,94 +0,0 @@
package artifacts
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// Asking the store to let go of what the mesh no longer keeps (novox/hq ADR 0189, issue 108).
//
// A fake store records what it was asked to delete, so what is asserted is the mesh's decision
// and the shape of the request — not the registry's behaviour, which is the registry's to test.
func fakeStore(t *testing.T, answer int) (Store, *[]string) {
t.Helper()
var asked []string
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodDelete {
t.Errorf("the store was asked %s %s; collecting is a delete", r.Method, r.URL.Path)
}
asked = append(asked, r.URL.Path)
w.WriteHeader(answer)
}))
t.Cleanup(server.Close)
return Store{Address: strings.TrimPrefix(server.URL, "http://")}, &asked
}
func TestAnImageAndAnArchiveAreAskedForAtTheirOwnEndpoints(t *testing.T) {
// The two shapes the mesh records: a manifest by digest, and a blob by digest. They are
// different endpoints, and asking at the wrong one answers 404 — which this would then
// record as collected, leaving the bytes on disk for ever while the record says otherwise.
store, asked := fakeStore(t, http.StatusAccepted)
ctx := context.Background()
image := catalogue.ArtifactStoreScheme + "web/app@sha256:abc123"
archive := catalogue.ArtifactStoreScheme + "web/config/blobs/sha256:def456"
if err := store.LetGo(ctx, image); err != nil {
t.Fatal(err)
}
if err := store.LetGo(ctx, archive); err != nil {
t.Fatal(err)
}
want := []string{"/v2/web/app/manifests/sha256:abc123", "/v2/web/config/blobs/sha256:def456"}
if len(*asked) != 2 || (*asked)[0] != want[0] || (*asked)[1] != want[1] {
t.Fatalf("the store was asked %v; want %v", *asked, want)
}
}
func TestAStoreThatDoesNotHaveItAnswersGone(t *testing.T) {
// The outcome wanted, already true. Told apart from success only so the sweep can say which
// happened; both are recorded, because retrying for ever is the thing to avoid.
store, _ := fakeStore(t, http.StatusNotFound)
err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/app@sha256:abc123")
if !errors.Is(err, Gone) {
t.Fatalf("a store that does not hold it answered %v, want Gone", err)
}
}
func TestAStoreWithDeletionOffSaysSoAndNamesTheRemedy(t *testing.T) {
// The registry answers 405 when it was started without deletion enabled. The remedy is a
// setting on the store's module, and saying "405" would send somebody to the wrong place.
store, _ := fakeStore(t, http.StatusMethodNotAllowed)
err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/app@sha256:abc123")
if err == nil {
t.Fatal("a store that refuses deletion was read as success")
}
if !strings.Contains(err.Error(), "REGISTRY_STORAGE_DELETE_ENABLED") {
t.Fatalf("the refusal does not name the remedy: %v", err)
}
}
func TestAReferenceThatIsNotTheMeshsOwnIsNeverAsked(t *testing.T) {
// The whole safety of the sweep is that it names only what the mesh recorded putting there.
// A reference of another shape — a vendor's image, a package version — is refused rather
// than composed into a delete somewhere that is not the mesh's store.
store, asked := fakeStore(t, http.StatusAccepted)
for _, reference := range []string{
"docker.io/library/registry@sha256:abc123",
"registry@sha256:abc123",
"1.4.2",
} {
if err := store.LetGo(context.Background(), reference); err == nil {
t.Errorf("%s was asked about; it is not a reference into the mesh's store", reference)
}
}
if len(*asked) != 0 {
t.Fatalf("the store was asked about %v", *asked)
}
}
@@ -234,3 +234,13 @@ func admitsSubject(pattern, subject []string) bool {
} }
return len(pattern) == len(subject) 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)
}
}
+16 -15
View File
@@ -144,12 +144,18 @@ func ConsumerFor(p Principal) (Consumer, bool) {
}, true }, true
} }
// HolderConsumerFor is the worker a seat's holder gets on that seat's work queue. // HolderConsumerFor is the worker a seat's holders share on that seat's work queue.
// //
// **A queue group even though the seat guarantees one holder.** The seat is *authority* — who may // **One worker for every holder, and each holder pulls one ask when it is idle** (novox/hq ADR
// be the telegram sender — and the queue group is *delivery*. Tie delivery to the seat and the // 0190). The seat is *authority* — who may be the telegram sender — and the worker is *delivery*,
// day somebody allows two holders for throughput, every message is processed twice with nothing // kept separate so that relaxing one changes nothing about the other: a node-scoped seat has a
// reporting it. Kept separate, relaxing one changes nothing about the other. // 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.
func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) { func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) {
if len(seat.Accepts) == 0 { if len(seat.Accepts) == 0 {
return Consumer{}, false return Consumer{}, false
@@ -158,18 +164,13 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
Name: "SEAT_" + upperSnake(seat.Name) + "_worker", Name: "SEAT_" + upperSnake(seat.Name) + "_worker",
Stream: seatStreamName(seat.Name), Stream: seatStreamName(seat.Name),
Filters: []string{"mesh.seat." + seat.Name + ".accept.>"}, Filters: []string{"mesh.seat." + seat.Name + ".accept.>"},
Queue: "holders",
AckWaitSeconds: 60, AckWaitSeconds: 60,
MaxDeliver: 5, MaxDeliver: 5,
// **One in flight.** A holder works one ask at a time, so the server hands it one at a // As many in flight as there are holders working, which pulling bounds by itself: a holder
// time: with the default of many, every ask behind the one being worked was delivered, // fetches one and fetches again only after it acknowledged. The server's default stands.
// left unacknowledged for the length of the work, redelivered after the ack wait, and Why: fmt.Sprintf("%s on %s holds %s; every holder pulls one ask at a time from this worker "+
// after the fifth time dropped — on 2026-10-01 twenty-six of forty-three builds asked in "and acknowledges after the work is done, so a crash mid-work redelivers rather than "+
// two minutes were never built, and the queue read as empty (novox/hq issue 186). "loses and an idle holder is the one that takes the next ask", module, node, seat.Name),
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 }, true
} }
+25 -12
View File
@@ -88,15 +88,20 @@ func TestAModuleThatConsumesNothingGetsNoConsumer(t *testing.T) {
} }
} }
// The seat is authority and the queue group is delivery. Tie them together and the day somebody // The seat is authority and the worker is delivery (novox/hq ADR 0190): one worker per seat, shared
// allows two holders, every message is processed twice with nothing reporting it. // by every holder and pulled from, so a second holder takes the next ask rather than a copy of the
func TestAHoldersWorkerUsesAQueueGroupAnyway(t *testing.T) { // same one — which is what a queue group used to guard, and what pulling one durable gives outright.
func TestAHoldersWorkerIsOneSharedByItsHolders(t *testing.T) {
c, ok := HolderConsumerFor("one", "telegram", telegramSeat()) c, ok := HolderConsumerFor("one", "telegram", telegramSeat())
if !ok { if !ok {
t.Fatal("the holder of a seat with inbound work got no worker") t.Fatal("the holder of a seat with inbound work got no worker")
} }
if c.Queue == "" { two, _ := HolderConsumerFor("two", "telegram", telegramSeat())
t.Fatal("the worker is not in a queue group, so a second holder would double-process") 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.Stream != "SEAT_TELEGRAM_SENDER" { if c.Stream != "SEAT_TELEGRAM_SENDER" {
t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream) t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream)
@@ -154,15 +159,23 @@ func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) {
has(t, perms.Subscribe, c.Filters[0]) has(t, perms.Subscribe, c.Filters[0])
} }
// A holder works one ask at a time, so the server hands it one at a time (novox/hq issue 186): // Every holder of a seat shares one worker and pulls from it (novox/hq ADR 0190): no queue group
// asks queued behind the one being worked wait in the stream rather than being delivered, // and no delivery subject, because a push consumer hands the next ask to whichever subscriber the
// left to expire and dropped after the fifth redelivery. // server picks, busy or not; and no cap of one in flight, because pulling bounds the asks in flight
func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) { // by the holders that are free — which is what ended the race of issue 186, where asks delivered
c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}}) // 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"}})
if !found { if !found {
t.Fatal("a seat that accepts work has no worker") t.Fatal("a seat that accepts work has no worker")
} }
if c.MaxAckPending != 1 { if c.Queue != "" || c.Push {
t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending) 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)
} }
} }
+34
View File
@@ -73,3 +73,37 @@ func TestAnAccountMayReadItsOwnMembershipAndNoOthers(t *testing.T) {
has(t, perms.Publish, "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.anchor.postgres") has(t, perms.Publish, "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.anchor.postgres")
hasNot(t, perms.Subscribe, "mesh.assignment.>") 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)
}
}
+102 -10
View File
@@ -34,8 +34,20 @@ const (
// authority is a list of tools and nothing else — not control, not declarations, not builds, // 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. // and no ability to answer anything, because a person asks.
KindPerson Kind = "person" 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 // 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). // emits (novox/hq ADR 0118, design 29 §5).
type Seat struct { type Seat struct {
@@ -74,6 +86,13 @@ type Principal struct {
// a namespace no such module owns. Every service started and the graph stayed empty. // a namespace no such module owns. Every service started and the graph stayed empty.
Watches []Seat 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 // 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 // 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. // (novox/hq ADR 0152) — the console's does, and nothing else's.
@@ -91,10 +110,13 @@ type Principal struct {
PasswordHash string PasswordHash string
} }
// meshSeatsTheControllerUses are the roles the mesh's own flows submit work to. Named rather than // seatsTheControllerAsks 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 // 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. // seat has to be stated, and a list is what makes "which roles does the mesh itself talk to" answerable.
var meshSeatsTheControllerUses = []string{"mesh-build-machine"} // 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"}
// enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the // 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. // controller may answer an enrolment is derived from the same constant the user is named from.
@@ -112,7 +134,10 @@ func (p Principal) Username() string {
switch p.Kind { switch p.Kind {
case KindPerson: case KindPerson:
return "person." + p.Module return "person." + p.Module
case KindModule: 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.
return p.Node + "." + p.Module return p.Node + "." + p.Module
case KindNode: case KindNode:
return "node." + p.Node return "node." + p.Node
@@ -186,7 +211,9 @@ 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 // 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 // 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. // like the catalogue does — which is why no holder needs to publish into anybody's inbox.
for _, seat := range meshSeatsTheControllerUses { // 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 {
pub = append(pub, "mesh.seat."+seat+".accept.>") pub = append(pub, "mesh.seat."+seat+".accept.>")
} }
// **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own // **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own
@@ -346,13 +373,17 @@ func PermissionsFor(p Principal) (Permissions, error) {
// 3. Seats it holds: full participation. // 3. Seats it holds: full participation.
for _, s := range p.Holds { for _, s := range p.Holds {
// Taking work from the role's queue: the worker consumer it binds (asked about, // Taking work from the role's queue: the worker consumer every holder shares (asked
// delivered on, acknowledged), each on the seat's own stream. The first machine to // about, pulled from, acknowledged), on the seat's own stream (novox/hq ADR 0190). A
// take work over the new bus was refused the asking (2026-09-28). // 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).
worker := "SEAT_" + upperSnake(s.Name) + "_worker" worker := "SEAT_" + upperSnake(s.Name) + "_worker"
stream := seatStreamName(s.Name) stream := seatStreamName(s.Name)
sub = append(sub, "_DELIVER."+worker, "_DELIVER."+worker+".>") pub = append(pub,
pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>") "$JS.API.CONSUMER.INFO."+stream+"."+worker,
"$JS.API.CONSUMER.MSG.NEXT."+stream+"."+worker,
"$JS.ACK."+stream+"."+worker+".>")
for _, a := range s.Accepts { for _, a := range s.Accepts {
sub = append(sub, seatSubject(s, "accept", a)) sub = append(sub, seatSubject(s, "accept", a))
} }
@@ -375,6 +406,48 @@ func PermissionsFor(p Principal) (Permissions, error) {
pub = append(pub, seatToolSubject(s, t, "*")) 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 { if p.Kind == KindPerson {
@@ -382,6 +455,11 @@ func PermissionsFor(p Principal) (Permissions, error) {
// consumer, because nothing is delivered to a person — they ask and are answered. // consumer, because nothing is delivered to a person — they ask and are answered.
sub = append(sub, p.inbox()) 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 { if p.Kind == KindModule || p.Kind == KindNode || p.Kind == KindController {
// Its own reply space, and nothing wider. // Its own reply space, and nothing wider.
@@ -403,7 +481,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 // 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 // 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. // person are never asked anything, and are granted nothing here.
AllowResponses: p.Kind == KindModule || p.Kind == KindController, AllowResponses: p.Kind == KindModule || p.Kind == KindController || p.Kind == KindNodeTools,
}, nil }, nil
} }
@@ -625,6 +703,20 @@ func ComposeAccounts(principals []Principal) (string, error) {
return b.String(), nil 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 { func quoted(values []string) string {
if len(values) == 0 { if len(values) == 0 {
return "" return ""
+91
View File
@@ -371,3 +371,94 @@ 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.>")
}
+6 -1
View File
@@ -217,10 +217,15 @@ var ControllerFollows = []string{
// A build's outcome, which is the build-machine role's own event now (ADR 0121) rather than a // 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 // message on the control branch. Same three audiences, one publish: whoever asked, this, and the
// catalogue. // catalogue.
seatEventSubject("mesh-build-machine", "built"), seatEventSubject("node-build-agent", "built"),
// The forge's merges: what moved a source, so the mesh builds what that source produces // 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. // without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name.
moduleEventSubject("gitea", "pull.merged"), 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 // 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 jetstream: enabled
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { 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"] } 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.>"] } 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"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { { 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"] } subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] }
} } } }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { { 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.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } 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: ["_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"] } subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
+20
View File
@@ -62,13 +62,33 @@ func Users(r Records) ([]Principal, error) {
for _, node := range sortedCopy(r.Nodes) { for _, node := range sortedCopy(r.Nodes) {
out = append(out, Principal{Kind: KindNode, Node: node}) 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] { 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{ out = append(out, Principal{
Kind: KindModule, Node: node, Module: d.Module, Kind: KindModule, Node: node, Module: d.Module,
Emits: d.Emits, Consumes: d.Consumes, Serves: d.Serves, Emits: d.Emits, Consumes: d.Consumes, Serves: d.Serves,
Holds: d.Holds, Uses: d.Uses, Watches: d.Watches, Invokes: d.Invokes, 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) { for _, node := range sortedCopy(r.Enrolling) {
out = append(out, Principal{Kind: KindEnrolment, Node: node}) out = append(out, Principal{Kind: KindEnrolment, Node: node})
+51
View File
@@ -245,3 +245,54 @@ func TestAUserListIsComposedBeforeAnythingMovesOntoTheBus(t *testing.T) {
t.Errorf("the composed list does not contain the machine running the bus") 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")
}
}
+20
View File
@@ -968,6 +968,26 @@ func compile(ctx context.Context, run Runner, tree string, chain Toolchain,
if _, err := run(ctx, tree, "docker", invocation...); err != nil { if _, err := run(ctx, tree, "docker", invocation...); err != nil {
return "", err 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 return filepath.Join(tree, out), nil
} }
+40 -1
View File
@@ -82,6 +82,41 @@ func TestABundleIsCompiledAndPackedWithNoDockerfile(t *testing.T) {
if !strings.HasPrefix(digest, "sha256:") { if !strings.HasPrefix(digest, "sha256:") {
t.Fatalf("the bundle was not pinned: %v", got.Manifest.Resources[0]) 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 // **Refused before anything is built, naming what to build first.** A base the mesh has not built
@@ -145,9 +180,13 @@ func TestTwoBundlesInOneModuleArePackedSeparately(t *testing.T) {
t.Fatalf("a module with two bundles did not build: %v", err) t.Fatalf("a module with two bundles did not build: %v", err)
} }
// Compiled into two different places. // Compiled into two different places. Only the compile lines: the copy of each bundle's
// dependencies names the same directory again, deliberately.
var outputs []string var outputs []string
for _, line := range r.ran { for _, line := range r.ran {
if !strings.Contains(line, "--outDir") {
continue
}
for _, part := range strings.Fields(line) { for _, part := range strings.Fields(line) {
if strings.HasPrefix(part, ".mesh-build/") { if strings.HasPrefix(part, ".mesh-build/") {
outputs = append(outputs, part) outputs = append(outputs, part)
+28 -4
View File
@@ -57,6 +57,23 @@ type Toolchain struct {
// carrying its debug info. The mistake was believing a comment rather than reading the file it // carrying its debug info. The mistake was believing a comment rather than reading the file it
// produced (novox/hq 04-ISSUES/161). // produced (novox/hq 04-ISSUES/161).
LinkerFlags []string 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, // 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). // for a language whose binaries are pinned to one at link time (novox/hq ADR 0005).
// //
@@ -107,14 +124,21 @@ var toolchains = []Toolchain{
// symlinks to a launcher that requires its library relatively — and the base image's own // 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. // 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. // 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{ Compile: []string{
"node", "/app/node_modules/typescript/bin/tsc", "node", "/app/node_modules/typescript/bin/tsc",
"--module", "NodeNext", "--moduleResolution", "NodeNext", "--module", "NodeNext", "--moduleResolution", "NodeNext",
"--target", "ES2022", "--target", "ES2022", "--rootDir", ".",
}, },
OutputFlag: "--outDir", OutputFlag: "--outDir",
Unit: UnitSources, Unit: UnitSources,
SourceExt: ".ts", SourceExt: ".ts",
Dependencies: "/app/runtime",
}, },
{ {
Language: "go", Language: "go",
+4 -12
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 // 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 // module may not reach a binding it does not have — the same boundary as a secret, for the same
// reason. // reason.
func knownFor(m Manifest, needs []Needed, node string) (map[string]map[string]string, error) { func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]string {
out := map[string]map[string]string{} out := map[string]map[string]string{}
for _, want := range m.Wants() { for _, want := range m.Wants() {
for i := range needs { for i := range needs {
@@ -54,20 +54,12 @@ func knownFor(m Manifest, needs []Needed, node string) (map[string]map[string]st
if n.Name != want || n.For != m.Module { if n.Name != want || n.For != m.Module {
continue continue
} }
as := ConsumerIdentity(node, IdentitySource(m.Slug, m.Module))
values := map[string]string{ values := map[string]string{
"at": n.At, "at": n.At,
"from": n.From, "from": n.From,
"as": as, "as": ConsumerIdentity(node, IdentitySource(m.Slug, m.Module)),
} }
// What the provider derives for this consumer rather than for all of them for key, value := range n.Serves {
// (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, // 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. // which is what a float would write and what a connection string would refuse.
values[key] = plainly(value) values[key] = plainly(value)
@@ -75,7 +67,7 @@ func knownFor(m Manifest, needs []Needed, node string) (map[string]map[string]st
out[want] = values out[want] = values
} }
} }
return out, nil return out
} }
// withOwnNames adds a module's own composed names to what it may name from one binding: // withOwnNames adds a module's own composed names to what it may name from one binding:
+46 -1
View File
@@ -67,6 +67,36 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
out := m out := m
out.Build = nil out.Build = nil
out.Resources = 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 { for _, r := range m.Resources {
named, _ := r["artifact"].(string) named, _ := r["artifact"].(string)
if named == "" { if named == "" {
@@ -110,7 +140,8 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) {
// The same on the wire: both are bytes fetched by digest and unpacked. They differ in // 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 // how they were made — one packed as it stood, the other compiled first — and a
// machine has no reason to care which. // machine has no reason to care which.
filled["source"] = artifact.Reference // Kept, not routed, for the reason the bundles above are (ADR 0155).
filled["source"] = Recorded(artifact.Reference)
filled["digest"] = artifact.Digest filled["digest"] = artifact.Digest
// **And `${version}`, so a resource can name a place that is this build's alone** // **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 // (novox/hq ADR 0141, 04-ISSUES/142). A component is unpacked into a directory named
@@ -191,6 +222,20 @@ func (b *Build) problems(module string) []string {
"%s: %q is a bundle and says no language, so nothing can choose a compiler "+ "%s: %q is a bundle and says no language, so nothing can choose a compiler "+
"for it", module, a.Name)) "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 // **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 // 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 // it was not built for (novox/hq ADR 0005); an artifact that says nothing would be
-290
View File
@@ -1,290 +0,0 @@
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
}
+34 -51
View File
@@ -508,10 +508,27 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
return nil, fmt.Errorf( return nil, fmt.Errorf(
"%s needs a secret called %q and none was made for it", m.Module, name) "%s needs a secret called %q and none was made for it", m.Module, name)
} }
first = append(first, ownedBy(m.SecretsOwner, map[string]any{ // 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{
"id": NeedID(name), "type": "file", "path": m.OwnSecrets[name].Path, "sealed": sealed, "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 // 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 // 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 // find each present — refusing clearly if the operator has not provided it — before it
@@ -628,16 +645,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
if err != nil { if err != nil {
return nil, err return nil, err
} }
as := ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module)) file, err := boundFile(*found, m.Binds[to], ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module)), own)
// 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 { if err != nil {
return nil, err return nil, err
} }
@@ -703,10 +711,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
return nil, err return nil, err
} }
// And what its bindings say, for the half of a connection that is not secret. // And what its bindings say, for the half of a connection that is not secret.
known, err := knownFor(m, r.Needs, r.Node) known := 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 // 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 // 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. // beside it said the same facts. Filled from the same answer, so the two cannot disagree.
@@ -723,11 +728,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
} }
local := *answered local := *answered
local.For = m.Module local.For = m.Module
here, err := knownFor(m, []Needed{local}, r.Node) for provision, values := range knownFor(m, []Needed{local}, r.Node) {
if err != nil {
return nil, err
}
for provision, values := range here {
known[provision] = values known[provision] = values
} }
} }
@@ -752,17 +753,6 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
// And the machine underneath, which no binding of its own can tell it. // And the machine underneath, which no binding of its own can tell it.
thisMachine := machineFacts(r, with.Names, with.MeshRange) 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 // 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). // one of them as its environment without saying so (ADR 0086, issue 041).
secretFiles := secretFilesOf(resources) secretFiles := secretFilesOf(resources)
@@ -912,6 +902,17 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
out = append(out, fact) 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 { if with.Adopted {
// First, before anything a module declares: what the mesh needs reachable, then its guard. // First, before anything a module declares: what the mesh needs reachable, then its guard.
// The order a machine applies is the order written here. // The order a machine applies is the order written here.
@@ -1127,19 +1128,6 @@ type Contribution struct {
// requirement's name — everything providing `reverse-proxy` understands the same shape, which // requirement's name — everything providing `reverse-proxy` understands the same shape, which
// is what makes swapping one for another cost nothing. // is what makes swapping one for another cost nothing.
Values map[string]any `json:"values"` 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. // grantPath is where one consumer's sealed credential lands on the providing machine.
@@ -1231,17 +1219,12 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
// told about it and withdraws the login on its next pass. // told about it and withdraws the login on its next pass.
continue 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{ out[g.Provision] = append(out[g.Provision], Contribution{
From: g.From, Node: g.Consumer, At: g.At, Values: g.Values, Derived: derived, From: g.From, Node: g.Consumer, At: g.At, Values: g.Values,
// One holder per local name: the identity the consumer is known by, and the local name // 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 — // 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. // a secret is not a login — so the identity limit does not apply to the suffix.
As: as, As: holderAs(ConsumerIdentity(g.Consumer, IdentitySource(g.Slug, g.From)), g.Local),
Secret: grantPath(directories[g.Provision], g.Consumer, holderAs(g.From, g.Local)), Secret: grantPath(directories[g.Provision], g.Consumer, holderAs(g.From, g.Local)),
}) })
if granted[g.Provision] == nil { if granted[g.Provision] == nil {
@@ -1,297 +0,0 @@
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
}
+12 -2
View File
@@ -1,6 +1,8 @@
package catalogue package catalogue
import ( import (
"crypto/sha256"
"encoding/hex"
"fmt" "fmt"
"sort" "sort"
"strings" "strings"
@@ -43,8 +45,16 @@ func jailsInto(modules []Manifest, j *Jailing) []map[string]any {
out := make([]map[string]any, 0, len(jails)+1) out := make([]map[string]any, 0, len(jails)+1)
for _, d := range jails { for _, d := range jails {
fmt.Fprintf(&composed, "\n# from %s\n[%s]\nenabled = true\nfilter = %s\n%s\n", // **The filter's digest rides in the jail file.** fail2ban is restarted when this file
d.module, d.jail.Name, d.jail.Name, strings.TrimRight(d.jail.Jail, "\n")) // 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"))
// The filter is a file of its own, named as the jail's filter= references it. // The filter is a file of its own, named as the jail's filter= references it.
out = append(out, map[string]any{ out = append(out, map[string]any{
"id": "filter-" + d.jail.Name, "id": "filter-" + d.jail.Name,
+26
View File
@@ -44,3 +44,29 @@ func TestTheComposedJailFileIsWrittenEvenWhenEmpty(t *testing.T) {
t.Fatalf("the empty composed jail file was not written alone: %v", files) 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)
}
}
+49 -79
View File
@@ -572,6 +572,36 @@ type Manifest struct {
// a module that could ask for it could read every credential on the bus — and the claim on // 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. // `mesh-broker` is what authorises it, checked from this manifest alone.
BusUsers string `json:"bus-users,omitempty"` 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. // Build says how to produce this module's artifacts from its source.
@@ -708,6 +738,16 @@ type Artifact struct {
// somebody adds a helper. An empty list is a bundle that is run rather than loaded — a // 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. // provisioner or a step, named by whatever runs it.
Entrypoints []string `json:"entrypoints,omitempty"` 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. // Kinds an artifact may be.
@@ -1333,6 +1373,15 @@ func ParseManifest(raw []byte) (Manifest, error) {
// //
// Refused here because the alternative is a build that never returns, on a mesh new enough // Refused here because the alternative is a build that never returns, on a mesh new enough
// that nobody is watching it yet. // 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 { if m.Build != nil && len(m.Build.Artifacts) > 0 {
for _, o := range m.Offers() { for _, o := range m.Offers() {
if o != ArtifactStoreProvision { if o != ArtifactStoreProvision {
@@ -1360,10 +1409,6 @@ func ParseManifest(raw []byte) (Manifest, error) {
"%s serves %q to whoever requires it, and does not provide it", m.Module, to)) "%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 { for to, where := range m.Binds {
if !placedOrAbsolute(where) { if !placedOrAbsolute(where) {
problems = append(problems, fmt.Sprintf( problems = append(problems, fmt.Sprintf(
@@ -1511,13 +1556,6 @@ func ParseManifest(raw []byte) (Manifest, error) {
} }
} }
} }
// **A scheduled step may hold this module's own containers still while it runs**
// (novox/hq ADR 0189). What the host judges is the declaration it receives — whether each
// id is a container placed on that machine; what belongs here is what only the definition
// shows: that the ids are this module's, that they are containers, and that the step is
// scheduled. A module naming a neighbour's container would be a module that can stop the
// mesh, and the manifest is where that is visible.
problems = append(problems, whileStoppedProblems(m, r, hasSchedule(r))...)
} }
for name, own := range m.OwnSecrets { for name, own := range m.OwnSecrets {
if !placedOrAbsolute(own.Path) { if !placedOrAbsolute(own.Path) {
@@ -2049,71 +2087,3 @@ func (o OwnSecrets) Paths() map[string]string {
// InstancesInterchangeable is the one value of a definition's `instances`: the module is the same // InstancesInterchangeable is the one value of a definition's `instances`: the module is the same
// on every machine, so any instance may answer for the module. // on every machine, so any instance may answer for the module.
const InstancesInterchangeable = "interchangeable" const InstancesInterchangeable = "interchangeable"
// WhileStopped is the resource key naming the containers a scheduled step holds still while it
// runs (novox/hq ADR 0189). Carried to the host unchanged, like `schedule`.
const WhileStopped = "while-stopped"
// hasSchedule is whether a resource declares a cadence, as a string.
func hasSchedule(r map[string]any) bool {
s, _ := r["schedule"].(string)
return s != ""
}
// whileStoppedProblems judges one container's maintenance window against its own definition
// (novox/hq ADR 0189).
//
// Three things the manifest is the only place to see: that the step is scheduled (a one-time
// offline job says *before* rather than *instead of* — at apply the host already has a window,
// because the declaration is applied in order and a run-once step gates what follows); that every
// id it names is **this module's own** container; and that it does not name itself.
//
// The host checks the fourth — that the container is actually placed on that machine — because
// that is a fact about the declaration and not about the definition.
func whileStoppedProblems(m Manifest, r map[string]any, scheduled bool) []string {
raw, present := r[WhileStopped]
if !present {
return nil
}
ids, ok := raw.([]any)
if !ok {
return []string{fmt.Sprintf(
"%s declares %s on %v as a %T; it is a list of this module's container ids",
m.Module, WhileStopped, r["id"], raw)}
}
var problems []string
if len(ids) > 0 && !scheduled {
problems = append(problems, fmt.Sprintf(
"%s declares %s on %v, which has no schedule. A maintenance window is for a recurring "+
"step: at apply the mesh already has one, because a run-once step gates what is "+
"declared after it (novox/hq ADR 0189)", m.Module, WhileStopped, r["id"]))
}
containers := map[string]bool{}
for _, own := range m.Resources {
if fmt.Sprint(own["type"]) == "container" {
containers[fmt.Sprint(own["id"])] = true
}
}
for _, each := range ids {
id, ok := each.(string)
if !ok {
problems = append(problems, fmt.Sprintf(
"%s declares %s on %v naming a %T; each entry is a container's id",
m.Module, WhileStopped, r["id"], each))
continue
}
if id == fmt.Sprint(r["id"]) {
problems = append(problems, fmt.Sprintf(
"%s declares %s on %v naming itself", m.Module, WhileStopped, r["id"]))
continue
}
if !containers[id] {
problems = append(problems, fmt.Sprintf(
"%s declares %s on %v naming %q, which is not a container this module declares. "+
"A step may hold still its own module's containers and nobody else's — one "+
"that could quiesce a neighbour could stop the mesh",
m.Module, WhileStopped, r["id"], id))
}
}
return problems
}
@@ -0,0 +1,42 @@
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
@@ -0,0 +1,251 @@
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
@@ -0,0 +1,88 @@
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
@@ -0,0 +1,212 @@
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)
}
})
}
+10 -1
View File
@@ -96,8 +96,17 @@ var defaultSeats = []Seat{
// A build says what it does as it does it (novox/hq ADR 0157): `started` when work is taken, // 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 // `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. // 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, {Name: "mesh-build-machine", Scope: ScopeMesh,
Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0121"}, Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0190"},
{Name: "node-dns-resolver", Scope: ScopeNode, 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 // 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 // whatever keeps it — who is banned and why, ban one address, let one go. Every holder serves all
+4 -3
View File
@@ -44,9 +44,10 @@ func TestTheSeatsAreAClosedSetAndEachNamesItsDecision(t *testing.T) {
delivered[s.Delivers] = s.Name delivered[s.Delivers] = s.Name
} }
} }
// Sixteen since node-service-manager (novox/hq ADR 0177). // Seventeen since node-build-agent (novox/hq ADR 0190) — sixteen once the retired
if len(Seats()) != 16 { // mesh-build-machine row goes, when no registered manifest claims it any more.
t.Errorf("the mesh defines %d seats rather than 16; the set is closed, so a change here is "+ if len(Seats()) != 17 {
t.Errorf("the mesh defines %d seats rather than 17; the set is closed, so a change here is "+
"a decision (novox/hq ADR 0110): %s", len(Seats()), seatNames()) "a decision (novox/hq ADR 0110): %s", len(Seats()), seatNames())
} }
} }
@@ -0,0 +1,23 @@
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)
}
}
-89
View File
@@ -1,89 +0,0 @@
package catalogue
import (
"encoding/json"
"strings"
"testing"
)
// A scheduled step may hold its module's own containers still while it runs (novox/hq ADR 0189).
//
// The host judges what it receives — whether each id is a container on that machine. What the
// definition is the only place to see is judged here, near whoever wrote it.
func aStoreManifest(step map[string]any) []byte {
m := map[string]any{
"module": "distribution", "version": "1",
"resources": []any{
map[string]any{"id": "store", "type": "container", "name": "mesh-registry",
"image": "registry@sha256:" + strings.Repeat("a", 64)},
step,
},
}
raw, _ := json.Marshal(m)
return raw
}
func TestAMaintenanceWindowOnItsOwnModulesContainerIsAccepted(t *testing.T) {
raw := aStoreManifest(map[string]any{
"id": "collect", "type": "container", "name": "mesh-registry-collect",
"image": "registry@sha256:" + strings.Repeat("a", 64),
"schedule": "30 3 * * *", "while-stopped": []any{"store"},
})
if _, err := ParseManifest(raw); err != nil {
t.Fatalf("a step holding its own module's container still was refused: %v", err)
}
}
func TestAMaintenanceWindowIsRefusedWhereTheDefinitionShowsItCannotMean(t *testing.T) {
for _, c := range []struct {
name string
step map[string]any
says string
}{
{
"on a step with no schedule",
map[string]any{"id": "collect", "type": "container", "name": "c",
"image": "registry@sha256:" + strings.Repeat("a", 64),
"while-stopped": []any{"store"}},
"gates what is declared after it",
},
{
"on a run-once step, which already has order",
map[string]any{"id": "collect", "type": "container", "name": "c",
"image": "registry@sha256:" + strings.Repeat("a", 64),
"run-once": true, "while-stopped": []any{"store"}},
"A maintenance window is for a recurring step",
},
{
"naming a container this module does not declare",
map[string]any{"id": "collect", "type": "container", "name": "c",
"image": "registry@sha256:" + strings.Repeat("a", 64),
"schedule": "30 3 * * *", "while-stopped": []any{"the-broker"}},
"could quiesce a neighbour could stop the mesh",
},
{
"naming itself",
map[string]any{"id": "collect", "type": "container", "name": "c",
"image": "registry@sha256:" + strings.Repeat("a", 64),
"schedule": "30 3 * * *", "while-stopped": []any{"collect"}},
"naming itself",
},
{
"written as something that is not a list",
map[string]any{"id": "collect", "type": "container", "name": "c",
"image": "registry@sha256:" + strings.Repeat("a", 64),
"schedule": "30 3 * * *", "while-stopped": "store"},
"a list of this module's container ids",
},
} {
_, err := ParseManifest(aStoreManifest(c.step))
if err == nil {
t.Errorf("%s was accepted", c.name)
continue
}
if !strings.Contains(err.Error(), c.says) {
t.Errorf("%s: the refusal does not say %q:\n%v", c.name, c.says, err)
}
}
}
+3
View File
@@ -42,6 +42,9 @@ const (
BusModule = "module" BusModule = "module"
BusEnrolment = "enrolment" BusEnrolment = "enrolment"
BusPerson = "person" 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 // MintBusPassword makes a bus password and records its hash under a username, replacing whatever was
+61
View File
@@ -42,6 +42,11 @@ type Source struct {
// Seen is when the source was last looked at — by a build, by hand, or by the forge saying it // 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. // moved. What a late report of an older move is judged against.
Seen time.Time 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. // Current reports whether what the mesh holds is what the source last had.
@@ -61,6 +66,32 @@ func (s Source) Current() bool {
// gains a requirement, a claim, a resource. What matters is that the change is visible the next // 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. // time a node is resolved, which it is.
func (i *Inventory) RegisterModule(ctx context.Context, m catalogue.Manifest, from Source) error { 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) raw, err := json.Marshal(m)
if err != nil { if err != nil {
return err return err
@@ -89,6 +120,36 @@ func (i *Inventory) RegisterModule(ctx context.Context, m catalogue.Manifest, fr
return err 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. // 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 // This is the whole of noticing. Nothing here builds anything — it writes down that the two
+58
View File
@@ -685,3 +685,61 @@ func TestRegisteringWithoutProvenanceKeepsTheSeat(t *testing.T) {
t.Fatalf("a hand-registered manifest erased where the module comes from: %+v", got) 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)
}
}
-242
View File
@@ -1,242 +0,0 @@
package inventory
import (
"context"
"encoding/json"
"strings"
)
// What the artifact store keeps, and what it may let go (novox/hq ADR 0189, issue 108).
//
// The store has never collected anything: every build pushes another layer set and nothing has
// ever removed one. The registry's own answer — collect what no tag names — is wrong here, because
// the mesh pushes each artifact under one moving tag and pins machines by digest, so every build
// but the newest is untagged and some machine may still be running it.
//
// **So the mesh decides, from its own records, and it never has to look in the store to do it.**
// It has never put anything there it did not record, which means every digest it could remove is
// already in a build row. A digest the mesh did not record making is therefore never named here —
// not as a safety margin but as the rule restated, and it is what keeps the sweep away from the
// images genesis pushed before any record existed (04-ISSUES/102, F4).
// KeptBuilds is how many successful builds of each module keep their artifacts, counting the
// newest. The newest is what the mesh hands a machine now; the four behind it are how far back a
// release that turns out wrong can be taken.
const KeptBuilds = 5
// ToCollect is every artifact the mesh made, no longer keeps, and has not already collected.
//
// Three reasons an artifact stays, and nothing else is a reason:
//
// - **a definition names it** — the reference appears in a module's recorded manifest, which is
// what the mesh would hand a machine now. No age limit: this is the floor;
// - **the mesh can still go back to it** — it is an artifact of one of the KeptBuilds most
// recent successful builds of its module;
// - it was already collected, in which case there is nothing left to do.
//
// Returned in a stated order so two runs over the same records ask for the same things in the
// same sequence, which is what makes a failed sweep safe to simply run again.
func (i *Inventory) ToCollect(ctx context.Context) ([]string, error) {
keep, err := i.keptReferences(ctx)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
// Every artifact of every successful build, oldest first, minus what has already been
// collected. A failed build published nothing, so it names nothing to remove.
`select b.made
from build b
where b.failed = '' and b.module is not null and b.module <> ''
order by b.at asc, b.id asc`)
if err != nil {
return nil, err
}
defer rows.Close()
collected, err := i.alreadyCollected(ctx)
if err != nil {
return nil, err
}
seen := map[string]bool{}
var out []string
for rows.Next() {
var raw []byte
if err := rows.Scan(&raw); err != nil {
return nil, err
}
var made []Artifact
if err := json.Unmarshal(raw, &made); err != nil {
// One unreadable record must not stop the rest being collected — and an artifact this
// row named is simply not offered, which errs toward keeping.
continue
}
for _, a := range made {
if a.Reference == "" || keep[a.Reference] || collected[a.Reference] || seen[a.Reference] {
continue
}
seen[a.Reference] = true
out = append(out, a.Reference)
}
}
return out, rows.Err()
}
// keptReferences is every artifact reference the mesh still keeps, for either of the two reasons.
func (i *Inventory) keptReferences(ctx context.Context) (map[string]bool, error) {
keep := map[string]bool{}
// **Whatever a definition the mesh holds names.** Read as text rather than by walking the
// resource shapes: a reference may be a container's image, a bundle's source, or a field some
// later kind of resource grows, and what matters is only whether the mesh could hand this
// string to a machine. A manifest that mentions it is a manifest that might.
manifests, err := i.store.Pool().Query(ctx, `select manifest::text from module where manifest is not null`)
if err != nil {
return nil, err
}
defer manifests.Close()
var named []string
for manifests.Next() {
var text string
if err := manifests.Scan(&text); err != nil {
return nil, err
}
named = append(named, text)
}
if err := manifests.Err(); err != nil {
return nil, err
}
// The KeptBuilds most recent successful builds of each module, whole.
recent, err := i.store.Pool().Query(ctx,
`select made from (
select made, row_number() over (partition by module order by at desc, id desc) as back
from build
where failed = '' and module is not null and module <> ''
) ranked where back <= $1`, KeptBuilds)
if err != nil {
return nil, err
}
defer recent.Close()
for recent.Next() {
var raw []byte
if err := recent.Scan(&raw); err != nil {
return nil, err
}
var made []Artifact
if err := json.Unmarshal(raw, &made); err != nil {
continue
}
for _, a := range made {
if a.Reference != "" {
keep[a.Reference] = true
}
}
}
if err := recent.Err(); err != nil {
return nil, err
}
// And anything a manifest mentions. Done after the recent set so the scan runs over the
// candidates rather than over every reference ever recorded: a manifest holds a reference
// composed with the store's address or kept bare, so the search is for the digest within it.
if len(named) > 0 {
all, err := i.everyReferenceMade(ctx)
if err != nil {
return nil, err
}
for _, reference := range all {
if keep[reference] {
continue
}
digest := digestIn(reference)
if digest == "" {
// Not something the store holds by digest; nothing here can speak for it, so it
// is kept rather than guessed about.
keep[reference] = true
continue
}
for _, text := range named {
if strings.Contains(text, digest) {
keep[reference] = true
break
}
}
}
}
return keep, nil
}
// everyReferenceMade is every artifact reference any successful build recorded.
func (i *Inventory) everyReferenceMade(ctx context.Context) ([]string, error) {
rows, err := i.store.Pool().Query(ctx,
`select made from build where failed = '' and module is not null and module <> ''`)
if err != nil {
return nil, err
}
defer rows.Close()
seen := map[string]bool{}
var out []string
for rows.Next() {
var raw []byte
if err := rows.Scan(&raw); err != nil {
return nil, err
}
var made []Artifact
if err := json.Unmarshal(raw, &made); err != nil {
continue
}
for _, a := range made {
if a.Reference == "" || seen[a.Reference] {
continue
}
seen[a.Reference] = true
out = append(out, a.Reference)
}
}
return out, rows.Err()
}
// digestIn is the `sha256:<hex>` a reference names, empty when it names none.
func digestIn(reference string) string {
for _, marker := range []string{"@sha256:", "/sha256:"} {
if _, after, ok := strings.Cut(reference, marker); ok {
return "sha256:" + after
}
}
return ""
}
// alreadyCollected is what the store has already been asked to let go.
func (i *Inventory) alreadyCollected(ctx context.Context) (map[string]bool, error) {
rows, err := i.store.Pool().Query(ctx, `select reference from artifact_collected`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]bool{}
for rows.Next() {
var reference string
if err := rows.Scan(&reference); err != nil {
return nil, err
}
out[reference] = true
}
return out, rows.Err()
}
// MarkCollected records that the store no longer holds these.
//
// **A store that answered "not found" is recorded too.** The outcome wanted is that the artifact
// is gone, and it is; retrying it every sweep for ever is the failure this table exists to
// prevent. Only a store that could not be reached, or refused, leaves a reference unmarked — and
// then the next sweep asks again, which is what should happen.
func (i *Inventory) MarkCollected(ctx context.Context, references []string) error {
for _, reference := range references {
if _, err := i.store.Pool().Exec(ctx,
`insert into artifact_collected (reference) values ($1) on conflict (reference) do nothing`,
reference); err != nil {
return err
}
}
return nil
}
-147
View File
@@ -1,147 +0,0 @@
package inventory
import (
"context"
"fmt"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// What the store keeps, and what it may let go (novox/hq ADR 0189, issue 108).
//
// The store has collected nothing since it was raised, and the registry's own answer — collect
// what no tag names — would delete images machines are running, because the mesh pushes under one
// moving tag and pins by digest. So the rule is the mesh's, read from its own records, and these
// are the three reasons an artifact stays and the one reason it goes.
// ref is an artifact reference as the mesh records one.
func ref(module, artifact string, n int) string {
return fmt.Sprintf("%s%s/%s@sha256:%064x", catalogue.ArtifactStoreScheme, module, artifact, n)
}
// built records one successful build of a module publishing one image.
func built(t *testing.T, inv *Inventory, id, module string, n int) string {
t.Helper()
reference := ref(module, "app", n)
b := aBuild(id, module, "")
b.Made = []Artifact{{Name: "app", Kind: "image", Reference: reference}}
if err := inv.RecordBuild(context.Background(), b); err != nil {
t.Fatal(err)
}
return reference
}
func TestTheStoreKeepsTheRecentBuildsAndLetsGoOfTheRest(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
// Eight builds of one module, oldest first. Five are kept — the newest, and the four a
// release that turns out wrong can be taken back to.
var made []string
for i := 1; i <= 8; i++ {
made = append(made, built(t, inv, fmt.Sprintf("b%02d", i), "web", i))
}
go_, err := inv.ToCollect(ctx)
if err != nil {
t.Fatal(err)
}
want := made[:3] // the three oldest
if len(go_) != len(want) {
t.Fatalf("offered %v to collect; want the %d oldest of %d", go_, len(want), len(made))
}
for i := range want {
if go_[i] != want[i] {
t.Fatalf("offered %v; want %v — and in that order, so a failed sweep is safe to run again",
go_, want)
}
}
}
func TestADefinitionNamingAnArtifactKeepsItHoweverOldItIs(t *testing.T) {
// The floor: no age limit. A module recorded at an older commit still names what the mesh
// would hand a machine now, and that is what must not be collected out from under it.
inv := fresh(t)
ctx := context.Background()
var made []string
for i := 1; i <= 8; i++ {
made = append(made, built(t, inv, fmt.Sprintf("b%02d", i), "web", i))
}
oldest := made[0]
// A definition the mesh holds, whose container runs that oldest image.
m := catalogue.Manifest{Module: "web", Version: "1", Resources: []map[string]any{{
"id": "app", "type": "container", "name": "web", "image": oldest,
}}}
if err := inv.RegisterModule(ctx, m, Source{Repository: "https://forge.invalid/web.git"}); err != nil {
t.Fatal(err)
}
go_, err := inv.ToCollect(ctx)
if err != nil {
t.Fatal(err)
}
for _, reference := range go_ {
if reference == oldest {
t.Fatalf("the mesh offered to collect %s, which a definition it holds names", oldest)
}
}
if len(go_) != 2 {
t.Fatalf("offered %v; want the two oldest that nothing names", go_)
}
}
func TestWhatHasBeenCollectedIsNotOfferedAgain(t *testing.T) {
// Without this the sweep reissues a delete for every artifact it has ever collected, every
// time it runs, for ever — a number of requests that grows with the mesh's whole history.
inv := fresh(t)
ctx := context.Background()
for i := 1; i <= 7; i++ {
built(t, inv, fmt.Sprintf("b%02d", i), "web", i)
}
first, err := inv.ToCollect(ctx)
if err != nil {
t.Fatal(err)
}
if len(first) != 2 {
t.Fatalf("offered %v, want two", first)
}
if err := inv.MarkCollected(ctx, first); err != nil {
t.Fatal(err)
}
again, err := inv.ToCollect(ctx)
if err != nil {
t.Fatal(err)
}
if len(again) != 0 {
t.Fatalf("offered %v again after collecting it", again)
}
}
func TestAFailedBuildNamesNothingToCollectAndEachModuleIsCountedOnItsOwn(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
// A failed build published nothing, so it is neither kept nor collected — and it must not
// count against the module's five.
for i := 1; i <= 6; i++ {
built(t, inv, fmt.Sprintf("w%02d", i), "web", i)
}
if err := inv.RecordBuild(ctx, aBuild("w99", "web", "the recipe would not build")); err != nil {
t.Fatal(err)
}
// And a second module with three builds keeps all three: five each, not five between them.
for i := 1; i <= 3; i++ {
built(t, inv, fmt.Sprintf("d%02d", i), "db", 100+i)
}
go_, err := inv.ToCollect(ctx)
if err != nil {
t.Fatal(err)
}
if len(go_) != 1 || go_[0] != ref("web", "app", 1) {
t.Fatalf("offered %v; want only web's oldest — db's three are all within its five", go_)
}
}
+1 -1
View File
@@ -63,7 +63,7 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin
if r := repositoryKey(e.Source.Repository); r != "" { if r := repositoryKey(e.Source.Repository); r != "" {
byRepository[r] = append(byRepository[r], name) byRepository[r] = append(byRepository[r], name)
} }
if e.Manifest.ClaimsSeat("mesh-build-machine") { if e.Manifest.ClaimsSeat("node-build-agent") || e.Manifest.ClaimsSeat("mesh-build-machine") {
builders = append(builders, name) builders = append(builders, name)
} }
} }
+1 -1
View File
@@ -12,7 +12,7 @@ func TestDependenciesAreOneRelationWithTheirKinds(t *testing.T) {
return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}} return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}}
} }
builder := entry("builder", "http://forge/novox/mesh-catalog.git") builder := entry("builder", "http://forge/novox/mesh-catalog.git")
builder.Manifest.Claims = []catalogue.Claim{{Name: "mesh-build-machine", Scope: catalogue.ScopeMesh}} builder.Manifest.Claims = []catalogue.Claim{{Name: "node-build-agent", Scope: catalogue.ScopeNode}}
plugin := entry("shop-plugin", "http://forge/novox/mesh-catalog.git") plugin := entry("shop-plugin", "http://forge/novox/mesh-catalog.git")
plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}} plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}}
entries := []Entry{ entries := []Entry{
@@ -1,23 +0,0 @@
-- What the artifact store no longer keeps (novox/hq ADR 0189, issue 108).
--
-- The mesh removes from its store only what it put there and can account for: every digest it
-- could remove is already in a build record, so the sweep reads its own records rather than
-- enumerating the store. What it does not get from those records is whether it has already
-- removed something -- `build.made` says what that build published, for ever, which is history
-- and not an index of what is on disk.
--
-- Without this the sweep would reissue a delete for every artifact it has ever collected, every
-- time it runs, and each one would answer 404 -- a number of requests that grows with the mesh's
-- whole history and never shrinks.
--
-- Keyed by the reference as the mesh records it (`artifact-store://<module>/<artifact>@sha256:…`),
-- because that is the identity the record uses everywhere else. Not a foreign key to build: two
-- builds can publish the same digest (the same source built twice produces the same bytes), and
-- what is collected is the artifact, not the attempt that made it.
create table artifact_collected (
reference text primary key,
-- When the store answered. Kept so a reader of an old build record can tell "this artifact is
-- gone" from "this artifact was never there", which are different kinds of surprise.
at timestamptz not null default now()
);
+32
View File
@@ -0,0 +1,32 @@
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")
}
}
+53 -7
View File
@@ -19,13 +19,57 @@ import (
// act on or a declaration a node reconciles toward; a build is a request that takes minutes and has // 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. // exactly one answer. Too long for request/reply, too particular to be an event.
// TheBuildMachine is the role a build is submitted to. // TheBuildMachine is the role a build is submitted to: node-scoped, held on every machine that
const TheBuildMachine = "mesh-build-machine" // 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
}
// BuildWork is where a build request lands, and BuildOutcome is where its result does. Derived from // 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 seat, so both sides name the role and neither names the other. The no-argument forms name the
func BuildWork() string { return "mesh.seat." + TheBuildMachine + ".accept.build" } // current role; the `Of` forms take the seat, for the handover during which two roles exist.
func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.built" } 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"
}
// BuildStarted is where a build machine says it has taken a build, and BuildLog is where it says // 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). // what it is doing, one line per message, under the build's own id (novox/hq ADR 0157).
@@ -35,8 +79,10 @@ func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.bui
// lived in one container's stderr on one machine. Every line is now an event of the role, retained // 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, // 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. // or reads it back afterwards from the stream, and a viewer is a subscriber and nothing more.
func BuildStarted() string { return "mesh.seat." + TheBuildMachine + ".event.started" } func BuildStarted() string { return BuildStartedOf(TheBuildMachine) }
func BuildLog(id string) string { return "mesh.seat." + TheBuildMachine + ".event.log." + id } 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 }
// BuildStart is what a build machine says the moment it takes a build. // BuildStart is what a build machine says the moment it takes a build.
type BuildStart struct { type BuildStart struct {
+1 -1
View File
@@ -26,7 +26,7 @@ func TestTheOldBusAnnouncesABuildUnderBothNames(t *testing.T) {
if KeyRoleBuilt != "built" { if KeyRoleBuilt != "built" {
t.Fatalf("the role's event is %q, and a holder emits its verbs bare", KeyRoleBuilt) t.Fatalf("the role's event is %q, and a holder emits its verbs bare", KeyRoleBuilt)
} }
if TheBuildMachine != "mesh-build-machine" { if TheBuildMachine != "node-build-agent" {
t.Fatalf("the role is %q", TheBuildMachine) t.Fatalf("the role is %q", TheBuildMachine)
} }
// The two must differ, or one publish would serve both and this doubling would be pointless. // The two must differ, or one publish would serve both and this doubling would be pointless.
+68 -25
View File
@@ -25,16 +25,34 @@ import (
type natsBuilds struct { type natsBuilds struct {
js *broker.JetStream js *broker.JetStream
owned bool 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. It dials, because the command that asks // BuildsOverNATS is the asking side on the bus being built, asking the current build role. It dials,
// for a build is a one-shot and holds nothing else. // because the command that asks for a build is a one-shot and holds nothing else.
func BuildsOverNATS(address string) (Builders, error) { 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) js, err := broker.Dial(address)
if err != nil { if err != nil {
return nil, fmt.Errorf("cannot reach the bus at %s to ask for a build: %w", address, err) return nil, fmt.Errorf("cannot reach the bus at %s to ask for a build: %w", address, err)
} }
return &natsBuilds{js: js, owned: true}, nil 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
} }
func (b *natsBuilds) Close() { func (b *natsBuilds) Close() {
@@ -51,7 +69,7 @@ func (b *natsBuilds) Ask(ctx context.Context, request BuildRequest) error {
} }
publish, cancel := context.WithTimeout(ctx, 30*time.Second) publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel() defer cancel()
if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil { if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil {
return fmt.Errorf("cannot submit a build: %w", err) return fmt.Errorf("cannot submit a build: %w", err)
} }
return nil return nil
@@ -63,7 +81,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 // 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 // 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. // the same event on EVENTS, which the controller records.
outcomes, err := b.js.Conn().SubscribeSync(BuildOutcome()) outcomes, err := b.js.Conn().SubscribeSync(BuildOutcomeOf(b.role()))
if err != nil { if err != nil {
return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err) return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err)
} }
@@ -80,7 +98,7 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest,
// be assumed, because nothing else will ever say so. // be assumed, because nothing else will ever say so.
publish, cancel := context.WithTimeout(ctx, 30*time.Second) publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel() defer cancel()
if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil { if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil {
return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err) return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err)
} }
@@ -115,9 +133,16 @@ type natsMachine struct {
sub *nats.Subscription sub *nats.Subscription
} }
// MachineOverNATS takes build work from the role this machine holds. // MachineOverNATS takes build work from the current build role.
func MachineOverNATS(js *broker.JetStream, on string) BuildMachine { func MachineOverNATS(js *broker.JetStream, on string) BuildMachine {
return &natsMachine{js: js, on: on, seat: TheBuildMachine} 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}
} }
func (m *natsMachine) Close() { func (m *natsMachine) Close() {
@@ -126,29 +151,31 @@ func (m *natsMachine) Close() {
} }
} }
// Take binds to the role's worker and hands each request over, one at a time. // Take binds to the role's worker and pulls one request at a time, handing each over.
// //
// **Bound, never created.** The work queue and the worker on it are the controller's to define // **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 // (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. // 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 { func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) error {
worker, found := broker.HolderConsumerFor(m.on, "builder", worker, found := broker.HolderConsumerFor(m.on, "build-agent",
broker.DeclaredSeat{Name: m.seat, Accepts: []string{"build"}}) broker.DeclaredSeat{Name: m.seat, Accepts: []string{"build"}})
if !found { if !found {
return fmt.Errorf("%s accepts no work, so there is nothing for this machine to take", m.seat) 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 // **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 — // 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 // "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 // on `…accept.>` is rejected even though it is narrower. Learned twice now, on two different
// consumers, which is why it is written down here. // consumers, which is why it is written down here.
filter := worker.Filters[0] filter := worker.Filters[0]
sub, err := m.js.Context().ChanQueueSubscribe(filter, worker.Queue, work, sub, err := m.js.Context().PullSubscribe(filter, worker.Name,
nats.Bind(worker.Stream, worker.Name), nats.ManualAck()) nats.Bind(worker.Stream, worker.Name), nats.ManualAck())
if err != nil { if err != nil {
return fmt.Errorf( return fmt.Errorf(
@@ -159,13 +186,27 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
m.sub = sub m.sub = sub
for { for {
select { if ctx.Err() != nil {
case <-ctx.Done():
return nil return nil
case msg, ok := <-work: }
if !ok { // One, and wait a while for it; an empty queue is a timeout, which is the normal state of a
return errors.New("the bus stopped delivering build work") // 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 var request BuildRequest
if err := json.Unmarshal(msg.Data, &request); err != nil { if err := json.Unmarshal(msg.Data, &request); err != nil {
// Unreadable: terminated rather than retried, because the next attempt reads the same // Unreadable: terminated rather than retried, because the next attempt reads the same
@@ -178,7 +219,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. // the ask to a second machine nor counts the wait against its deliveries.
working := make(chan struct{}) working := make(chan struct{})
go stillWorking(msg, working) go stillWorking(msg, working)
do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js}) do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js, seat: m.seat})
close(working) close(working)
} }
} }
@@ -189,7 +230,9 @@ type natsBuild struct {
msg *nats.Msg msg *nats.Msg
on string on string
js *broker.JetStream js *broker.JetStream
seq int // seat is the role this build was taken from; what the machine says about it is that role's.
seat string
seq int
} }
func (b *natsBuild) Request() BuildRequest { return b.request } func (b *natsBuild) Request() BuildRequest { return b.request }
@@ -216,7 +259,7 @@ func (b *natsBuild) Announce(ctx context.Context, result BuildResult) error {
} }
publish, cancel := context.WithTimeout(ctx, 30*time.Second) publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel() defer cancel()
if _, err := b.js.Context().Publish(BuildOutcome(), body, nats.Context(publish)); err != nil { if _, err := b.js.Context().Publish(BuildOutcomeOf(b.seat), body, nats.Context(publish)); err != nil {
return fmt.Errorf("cannot announce a build's outcome: %w", err) return fmt.Errorf("cannot announce a build's outcome: %w", err)
} }
return nil return nil
@@ -236,7 +279,7 @@ func (b *natsBuild) Began(ctx context.Context) error {
} }
publish, cancel := context.WithTimeout(ctx, 30*time.Second) publish, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel() defer cancel()
if _, err := b.js.Context().Publish(BuildStarted(), body, nats.Context(publish)); err != nil { if _, err := b.js.Context().Publish(BuildStartedOf(b.seat), body, nats.Context(publish)); err != nil {
return fmt.Errorf("cannot say a build started: %w", err) return fmt.Errorf("cannot say a build started: %w", err)
} }
return nil return nil
@@ -254,7 +297,7 @@ func (b *natsBuild) Say(step, message string) {
if err != nil { if err != nil {
return return
} }
_ = b.js.Conn().Publish(BuildLog(b.request.ID), body) _ = b.js.Conn().Publish(BuildLogOf(b.seat, b.request.ID), body)
} }
func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) } func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) }
+119 -3
View File
@@ -48,7 +48,7 @@ func aBusWithTheBuildRole(t *testing.T) *broker.JetStream {
t.Fatal(err) t.Fatal(err)
} }
clean := func() { clean := func() {
_ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") _ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT")
for _, s := range broker.MeshStreams() { for _, s := range broker.MeshStreams() {
_ = js.Context().PurgeStream(s.Name) _ = 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. // 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) deadline := time.Now().Add(5 * time.Second)
for time.Now().Before(deadline) { for time.Now().Before(deadline) {
info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE") info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT")
if err == nil && info.State.Msgs == 0 { if err == nil && info.State.Msgs == 0 {
return return
} }
@@ -177,7 +177,7 @@ func TestNatsABuildWaitsForAMachineRatherThanFailing(t *testing.T) {
if _, err := js.Context().Publish(BuildWork(), body); err != nil { if _, err := js.Context().Publish(BuildWork(), body); err != nil {
t.Fatal(err) t.Fatal(err)
} }
info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE") info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT")
if err != nil || info.State.Msgs != 1 { if err != nil || info.State.Msgs != 1 {
t.Fatalf("the work did not queue: %+v %v", info, err) t.Fatalf("the work did not queue: %+v %v", info, err)
} }
@@ -245,3 +245,119 @@ func TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue(t *testing.T) {
func quietLog() *log.Logger { return log.New(io.Discard, "", 0) } func quietLog() *log.Logger { return log.New(io.Discard, "", 0) }
var _ = quietLog 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)
}
}
+3 -2
View File
@@ -223,10 +223,11 @@ func kindOfSubject(subject string) (string, bool) {
return KindCatchUp, true return KindCatchUp, true
case broker.ControllerFollows[3]: case broker.ControllerFollows[3]:
return KindSourceMoved, true return KindSourceMoved, true
case BuildOutcome(): case BuildOutcome(), BuildOutcomeOf(TheBuildMachineBefore):
// A build's outcome is the role's event now, so it arrives on the events stream rather than // 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 // 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). // 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.
return KindBuilt, true return KindBuilt, true
} }
return "", false return "", false