Compare commits

..
Author SHA1 Message Date
mesh-admin ba6a3ea829 Merge pull request 'The rest of the mesh resolves each machine with its pins and names what it leaves out; a built-by edge never widens a plan; plans stop' (#199) from fix/a-machine-not-on-the-network-is-said into main 2026-10-01 16:25:09 +00:00
jschoubben dcf6278523 The rest of the mesh resolves each machine with its own pins, and says which machine it leaves out
The second pass of theRestOfTheMesh resolved every machine without its pins. Since a machine with two
providers of one provision is refused unless a pin names one (195/196), the control node was refused
there and vanished: every seat it holds read as unheld, the build machine refused what needs the git
seat, the roll-out was refused — and nothing said why (hq issue 188). Each machine is now resolved as
its plan resolves it, with its pins; a machine left out is named with the resolver's words.
2026-10-01 18:22:17 +02:00
jschoubben d94f9e7f9f A built-by edge orders and gates a plan; it never widens it — and plans stop
The first live plan took seventy-five modules along for a controller change: the builder packages the
controller's source, everything is built by the builder, so everything was reachable. A module built by
the build machine is not changed by a new build machine. Reachability now follows the code and build
edges only; built-by still orders a tier after the build machine and gates it on the machine's roll-out.
plans stop <id> ends a plan by hand: what was asked still builds and registers, nothing further is asked.
2026-10-01 18:16:47 +02:00
jschoubben 882687afd4 Merge remote-tracking branch 'origin/main' into fix/a-machine-not-on-the-network-is-said 2026-10-01 18:16:14 +02:00
jschoubben 884c088949 A machine the network filter drops is said, not skipped in silence (hq issue 188)
onTheNetwork resolves every machine unchecked and skipped one whose resolution refused. A machine
skipped there has no address, so its own plan fails on the first placeholder that needs one, in another
module's words, every seat held on it reads as unheld, and what is built from it cannot be built — four
symptoms, none naming the refusal (2026-10-01, the control node, forty minutes). The refusal is now said
where it happens, in the resolver's own words.
2026-10-01 18:14:34 +02:00
mesh-admin 76f27562a1 Merge pull request 'The plans verb: what the last merges produced and where each stands (hq ADR 0162)' (#198) from feat/the-plans-verb into main 2026-10-01 16:14:09 +00:00
jschoubben 89d43e0dad gofmt 2026-10-01 18:04:58 +02:00
jschoubben 09b636e628 The plans verb: what the last merges produced and where each stands (hq ADR 0162)
The mesh's seat answers plans — the recent plans with their tier, what each waits for and since when,
or one plan whole given its id — so the console reads a merge's progress where it reads everything
else, instead of a person reading the daemon's log.
2026-10-01 17:58:49 +02:00
mesh-admin d2ed6d50c9 Merge pull request 'A merge produces a tiered plan the mesh keeps (hq ADR 0162)' (#197) from feat/a-merge-produces-a-tiered-plan into main 2026-10-01 15:57:10 +00:00
jschoubben fdf338dbd1 A plan is kept, advanced and resumed from the store: the test 2026-10-01 17:56:06 +02:00
jschoubben a69a832248 A merge produces a tiered plan the mesh keeps (hq ADR 0162)
A module's dependencies are one relation in the catalogue — stands-on, packages, built-by, declared —
answered by one call. A merge takes what moved and everything reachable from it, sorts the set into
tiers (a code dependency in the same tier, a build dependency after its base is built, a runtime
dependency after the build machine is built and running; the build machine's own base comes first,
built by the one that runs), writes the plan to the store, asks the first tier and returns. Every
outcome advances the plan; a ticker advances what outcomes cannot; a controller replaced mid-plan
resumes it. status lists open plans and names one that has waited too long.
2026-10-01 17:53:55 +02:00
mesh-admin bdcbda801e Merge pull request 'The first pass does not refuse two providers beside the consumer (hotfix for #195)' (#196) from fix/first-pass-does-not-refuse-two-providers-beside into main 2026-10-01 15:34:23 +00:00
jschoubben 5a7ed56f61 The first pass does not refuse two providers beside the consumer
#195 refused two modules beside a consumer that both answer a bound
provision — in every pass. The first pass exists only to answer what a node
offers, and a refusal there makes the machine vanish from every other node's
world (resolve.go says so for its sibling case): with novox refused for its
own acme-ca, ace's plan lost the vault and the identity provider and read
"nothing in this mesh provides secret". Seen live within minutes of the
rollout. The second pass refuses it, where it is asked, as before.
2026-10-01 17:34:17 +02:00
mesh-admin d54ecb3bf2 Merge pull request 'A pin names the module as well as the node; a node that answers twice is refused' (#195) from feat/pin-names-the-provider into main 2026-10-01 15:26:02 +00:00
jschoubben cf84117638 A pin names the module as well as the node; a node that answers twice is refused
A provider is a (node, module) pair (design 23), and the pin — the one way a
consumer names its provider — named only the node. Two modules on one node
can both answer a provision (public-acme and step-ca both offer acme-ca on
novox), and then the resolver, given a pin naming that node, took the last
provider listed: a coin flip. The same ambiguity beside the consumer was
settled by a map walk — random per plan — which is how novox's own
route-proxy got its issuer (novox/hq #258).

- `pin <node> <provision> <from-node> <module>`: both halves, always. The
  console gains `pin` and `unpin`. The provider may be on the consumer's own
  node, since two modules beside it can both answer.
- The resolver refuses ambiguity instead of picking, across machines and
  beside the consumer alike, naming every candidate as node/module and the
  form of the pin that settles it. A plain capability that grants nothing
  and serves nothing (three shells beside an editor) is not a choice to put
  to anybody and stays as it was.
- provision_pin gains a nullable module (0050); records made before are
  completed where the node they name answers once, and left for a person
  where it answers twice (0051).
- The provider of something already satisfied is looked for among what was
  assigned, not only what the walk has reached — a consumer reached before
  the provider beside it no longer loses its binding.
- The start-time check that every declared verb is runnable samples each
  verb's required arguments from its schema instead of three guessed keys.

Live consequence: a node that has two providers of one bound provision
assigned (novox: acme-ca) resolves only once pinned —
`pin novox acme-ca novox public-acme`.
2026-10-01 17:23:31 +02:00
mesh-admin 8d4e940866 Merge pull request 'The build machine takes one ask at a time, and says so while it builds' (#194) from fix/the-build-machine-takes-one-ask-at-a-time into main 2026-10-01 14:19:18 +00:00
jschoubben 11b20b10ff The build machine takes one ask at a time, and says so while it builds
With the worker consumer's default of many deliveries in flight, every ask behind the one being
built was delivered at once, left unacknowledged for the length of the build, redelivered after the
ack wait and dropped after the fifth time: on 2026-10-01 twenty-six of forty-three builds asked in two
minutes were never built and the queue read as empty (hq issue 186). The holder's worker now has one
in flight, and a running build tells the bus it is still working, as the controller's long handlers
do, so a build longer than the ack wait is neither redelivered nor counted out.
2026-10-01 16:16:13 +02:00
mesh-admin 7e701c0db2 Merge pull request 'A merge of a base rebuilds what stands on it' (#193) from fix/a-merge-of-a-base-rebuilds-what-stands-on-it into main 2026-10-01 14:12:32 +00:00
jschoubben 853c63b181 A merge of a base rebuilds what stands on it
The controller knew which modules were built against which base artifacts and used it only when asked
(build --on). A merge that rebuilt the runtime image left forty-two modules on the old image until
somebody asked, twice, by hand (hq issue 186). The merge now takes every module standing on what moved,
through every layer, into the same rebuild, in base order — the same rebuild the flag does, asked by
the merge that made it necessary.
2026-10-01 16:09:26 +02:00
mesh-admin 84ac840ff4 Merge pull request 'The vault's seat, and a report that carries the machine's profile (hq ADR 0161)' (#192) from feat/the-vaults-seat-and-the-machines-profile into main 2026-10-01 14:02:10 +00:00
jschoubben e8e502343f The vault's seat, and a report that carries the machine's profile (hq ADR 0161)
mesh-vault joins the mesh's own set — mesh-scoped, delivering secret — because the controller seals
every minted credential with it, which is the test for a seat of the mesh's own; a second provider is
a second claimant, refused by name (issue 106). A report may carry the machine's profile, detected
again by the apply that reports, and the latest replaces the enrolled one: a machine that switched
its network manager is a machine whose uplink holder lacks a capability at its next push (issue 138).
2026-10-01 15:58:04 +02:00
mesh-admin 1edc44b25f Merge pull request 'The store seat promises two verbs, databases and query (ADR 0159)' (#186) from feat/the-store-seat-has-verbs into main 2026-10-01 13:45:29 +00:00
mesh-admin 5a1b37e477 Merge pull request 'A push issues the memberships of the machines it told' (#191) from fix/a-push-issues-memberships into main 2026-10-01 13:45:20 +00:00
jschoubben 4a6a4eadeb A push issues the memberships of the machines it told
sendTo — the path a roll-out, a rotation and a secret change take — issued memberships after its
sends (ADR 0160); the push command, which sends the same declarations through its own loop, did
not, so the one command operators run issued none. Said once in the same words after the sends.
2026-10-01 15:41:09 +02:00
mesh-admin 69f559ab4c Merge pull request 'A refused membership does not stop the controller' (#190) from fix/a-refused-membership-does-not-stop-the-controller into main 2026-10-01 13:31:27 +00:00
jschoubben 05b90f966a A refused membership does not stop the controller
A stream publish waits for its acknowledgement as long as its context lives, and the server never
acknowledges a publish it refuses. Issuing memberships after a push used the daemon's own context, so
the one refused membership of 2026-10-01 (hq issue 183) held the controller's receive loop for good:
no report, no build outcome, no merge was heard until a restart (hq issue 185). Issuing one
membership is now bounded to ten seconds, and a push says how many could not be issued and stands —
the machines keep the shape they derive until the next push.
2026-10-01 15:30:04 +02:00
mesh-admin 184913b620 Merge pull request 'The controller may publish the memberships it issues' (#189) from fix/the-controller-may-publish-memberships into main 2026-10-01 13:26:48 +00:00
jschoubben de7aed5016 The controller may publish the memberships it issues
The server refused every membership the controller published after a push (2026-10-01, Permissions
Violation for Publish to mesh.assignment.<node>.<module>): the controller's own grant named the
control, node and JetStream subjects and not the assignments it alone issues (hq ADR 0160). Broker
golden regenerated; one line differs.
2026-10-01 15:20:32 +02:00
jschoubben 18958154f0 A claim may name the verbs it serves for its seat (hq ADR 0160)
The store's databases is not postgres's postgres_list_databases, and a holder may serve both. A
claim's serves names the seat's verbs the module implements for the role; absent, the module's own
tools must list every verb the seat promises, which is how a module named like its seat says they
are one and the same. Registration refuses a claim naming a verb the seat never promised, and the
credential's claims carry the claim's own verbs to the runtime.
2026-10-01 15:18:34 +02:00
jschoubben a2b1f9e936 Merge remote-tracking branch 'origin/main' into feat/the-store-seat-has-verbs 2026-10-01 15:16:52 +02:00
mesh-admin c9a0b1f9f4 Merge pull request 'The mesh issues an assignment's subjects: a membership per module on a machine (ADR 0160)' (#188) from feat/the-mesh-issues-an-assignments-subjects into main 2026-10-01 12:45:45 +00:00
jschoubben f450303e8e A person invoking one tool holds it both ways it is addressed (the test, after #185) 2026-10-01 14:45:31 +02:00
jschoubben 603ad61142 The mesh issues an assignment's subjects: a membership per module on a machine (hq ADR 0160)
For every module on every machine the controller composes what that instance serves — its machine's
address always, the module's plain address in a queue when it is alone or its definition says its
instances are interchangeable — the verbs of the seats it holds at the seats' subjects, where its
events land, and what it may reach, resolved the same way for the modules it invokes. Published
beside the node's declaration on `mesh.assignment.<node>.<module>`, last per subject in a stream
that allows direct reads, and the account may read exactly its own. Composed from the same records
the bus's accounts are, so what a runtime serves and what its account may are one composition.
`instances: interchangeable` is the one fact a definition states for it.

The shape issued is the shape the mesh already had, so nothing moves when the membership arrives;
the runtime that reads it instead of deriving it is the next piece.
2026-10-01 14:43:26 +02:00
jschoubben e7cff3d38e The store seat promises two verbs, databases and query (hq ADR 0159)
Seeded additively into the seat's row at the controller's next start; a holder must list tools of
these names (design 33 §3), so this merges after the database engine's definition does.
2026-10-01 14:02:46 +02:00
53 changed files with 2326 additions and 172 deletions
+2
View File
@@ -581,6 +581,8 @@ type answers struct {
// pair that answers "has it caught up", which waiting alone cannot (the sent digest is
// recorded at send, not at apply).
reported []inventory.Reported
// plans is what the last merges produced and where each stands (novox/hq ADR 0162).
plans []inventory.Plan
// refused is why a machine cannot be worked out at all, by name. A different thing from every
// other answer here: those are about a machine that was told something, and this is about one
// that cannot be told anything — it never reaches waiting, because nothing was computed for it
+52
View File
@@ -0,0 +1,52 @@
package main
import (
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// A merge that rebuilds a base rebuilds what stands on it, through every layer, and nothing else
// (novox/hq issue 186): the runtime image moving means every module built on it moves too, and a
// module built on one of those moves as well.
func TestAMergeOfABaseTakesWhatStandsOnItAlong(t *testing.T) {
entry := func(name string) inventory.Entry {
return inventory.Entry{Manifest: catalogue.Manifest{Module: name}}
}
entries := []inventory.Entry{entry("mesh-tools"), entry("shop"), entry("shop-plugin"), entry("postgres"), entry("unrelated")}
against := map[string][]string{
"shop": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
"shop-plugin": {catalogue.ArtifactStoreScheme + "shop/runtime@sha256:b"},
"postgres": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
"unrelated": {catalogue.ArtifactStoreScheme + "alpine/base@sha256:c"},
}
got := dependentsOf([]inventory.Entry{entry("mesh-tools")}, entries, against)
var names []string
for _, e := range got {
names = append(names, e.Manifest.Module)
}
want := map[string]bool{"shop": true, "shop-plugin": true, "postgres": true}
if len(names) != len(want) {
t.Fatalf("rebuilt %v; wanted exactly the three that stand on the runtime, directly or through shop", names)
}
for _, n := range names {
if !want[n] {
t.Fatalf("%s was rebuilt and stands on nothing that moved (%v)", n, names)
}
}
// The dependents come in base order when the merge orders them: the runtime, then shop, then
// the plugin that stands on shop.
ordered := orderByBases(append([]inventory.Entry{entry("mesh-tools")}, got...), against)
pos := map[string]int{}
for i, e := range ordered {
pos[e.Manifest.Module] = i
}
if !(pos["mesh-tools"] < pos["shop"] && pos["shop"] < pos["shop-plugin"]) {
t.Fatalf("not in base order: %v", ordered)
}
// Nothing moved: nothing follows.
if more := dependentsOf(nil, entries, against); len(more) != 0 {
t.Fatalf("with nothing moved, %d module(s) were rebuilt", len(more))
}
}
+13 -1
View File
@@ -72,6 +72,8 @@ func run() error {
return askCommand(ctx, args[1:])
case "builds":
return buildsCommand(ctx, args[1:])
case "plans":
return plansCommand(ctx, args[1:])
case "pin":
return pinCommand(ctx, args[1:], true)
case "unpin":
@@ -203,7 +205,8 @@ func usage() {
licence refresh <name> mint a new access token and seal it to every holder
rotate <provision> [--consumer <n>] a new credential for every holder, both ends at once
ask <module> <tool> [json] call one of a module's tools over the broker, and print its answer
pin <node> <provision> <from> which node this one gets a provision from
pin <node> <provision> <from-node> <module>
which provider this one gets a provision from: the module, and its node
unpin <node> <provision> put that question back
plan <node> [--files|--json] what that node would run, and why
push [<node>] [--behind] send a node everything it should be, or only those that need it
@@ -252,12 +255,21 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error {
switch {
case err != nil && result.Failed != "":
fmt.Printf("%s: %v\n", result.ID, err)
if result.Module != "" {
planBuilt(ctx, b.inv, result.Module, result.Commit, result.Failed)
} else {
planFailedBuild(ctx, b.inv, result)
}
return nil
case err != nil:
fmt.Printf("%s: heard and recorded, and not registered: %v\n", result.ID, err)
if manifest.Module != "" {
planBuilt(ctx, b.inv, manifest.Module, result.Commit, err.Error())
}
return nil
}
fmt.Printf("%s: %s %s registered, built on %s from %s\n",
result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit))
planBuilt(ctx, b.inv, manifest.Module, result.Commit, "")
return nil
}
+9 -12
View File
@@ -411,8 +411,8 @@ func describeOffers(offers []catalogue.Offer) string {
// database should not change where an existing machine gets its data the day a second one
// arrives.
func pinCommand(ctx context.Context, args []string, setting bool) error {
if setting && len(args) != 3 {
return errors.New("pin <node> <provision> <from-node>")
if setting && len(args) != 4 {
return errors.New("pin <node> <provision> <from-node> <module>")
}
if !setting && len(args) != 2 {
return errors.New("unpin <node> <provision>")
@@ -431,17 +431,12 @@ func pinCommand(ctx context.Context, args []string, setting bool) error {
fmt.Printf("%s is no longer told where to get %s from\n", args[0], args[1])
return nil
}
if args[0] == args[2] {
// Allowed by nothing here, and worth saying rather than resolving into a confusing
// refusal later: a node providing something to itself is a node-scoped provision, and
// this field is for the other kind.
return fmt.Errorf("%s cannot get %s from itself; that would be a provision this machine "+
"provides, which does not need saying", args[0], args[1])
}
if err := inv.PinProvision(ctx, args[0], args[1], args[2]); err != nil {
// The provider's node may be this same machine: two modules beside the consumer can both
// answer a provision, and then the module is the whole question (novox/hq #258).
if err := inv.PinProvision(ctx, args[0], args[1], args[2], args[3]); err != nil {
return err
}
fmt.Printf("%s gets %s from %s\n", args[0], args[1], args[2])
fmt.Printf("%s gets %s from %s/%s\n", args[0], args[1], args[2], args[3])
fmt.Printf(" run `push %s` to send it\n", args[0])
return nil
}
@@ -716,7 +711,9 @@ func claimsFor(ctx context.Context, inv *inventory.Inventory, m catalogue.Manife
claimed := seatClaimed{Seat: c.Name, Scope: c.At()}
if s, known := byName[c.Name]; known {
claimed.Scope = s.Scope
claimed.Serves = catalogue.VerbNames(s.Serves)
// The verbs the runtime serves for the seat: the claim's own when it names them
// (ADR 0160), else every verb the seat promises, which its tools then answer.
claimed.Serves = c.ServesFor(catalogue.Manifest{Tools: catalogue.VerbNames(s.Serves)})
}
out = append(out, claimed)
}
+3
View File
@@ -567,6 +567,9 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: caps},
catalogue.World{Unchecked: true, Holdings: holdings})
if err != nil {
// Said, not skipped in silence: a machine dropped here loses its address, and every
// plan that names it fails in another module's words (novox/hq issue 188).
fmt.Fprintf(os.Stderr, "%s is not counted as on the network: it does not resolve: %v\n", p.Name, err)
continue
}
for _, m := range got.Modules {
+14 -3
View File
@@ -7,6 +7,7 @@ import (
"errors"
"flag"
"fmt"
"os"
"sort"
"strings"
@@ -281,8 +282,9 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
for _, o := range others {
got, err := catalogue.Resolve(shelf, o.assigned, o.node, catalogue.World{Unchecked: true, Holdings: holdings})
if err != nil {
// Their set does not resolve for some other reason. Not this node's problem to
// report, and nothing of theirs is running, so it offers nothing.
// Said, not skipped: a machine dropped here offers nothing and holds nothing as far
// as every other machine's plan can tell (novox/hq issue 188).
fmt.Fprintf(os.Stderr, "%s is left out of the rest of the mesh: it does not resolve: %v\n", o.node.Name, err)
continue
}
firstHeld = append(firstHeld, got.Claims...)
@@ -313,8 +315,17 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
world := catalogue.World{Offered: offered, Held: firstHeld, Holdings: holdings}
var held []catalogue.Held
for _, o := range others {
got, err := catalogue.Resolve(shelf, o.assigned, o.node, world)
// Each machine is resolved with its own pins, as its plan is: a machine that needs one to
// settle two providers would otherwise be refused here and vanish from the mesh — every
// seat it holds unheld, every build that needs one refused (2026-10-01, the control node;
// novox/hq issue 188).
theirs := world
if pins, err := inv.PinsFor(ctx, o.node.Name); err == nil {
theirs.Pinned = pins
}
got, err := catalogue.Resolve(shelf, o.assigned, o.node, theirs)
if err != nil {
fmt.Fprintf(os.Stderr, "%s is left out of the rest of the mesh: it does not resolve: %v\n", o.node.Name, err)
continue
}
held = append(held, got.Claims...)
+58
View File
@@ -4,6 +4,7 @@ import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"flag"
"fmt"
@@ -113,6 +114,9 @@ func serve(ctx context.Context) error {
// And build results nobody was waiting for. A build triggered any other way than `build`
// would otherwise be reported into the void, which is the same as not reporting it.
server.Records(builds{inv})
// Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for
// machines to report moves when they have, and a plan left by a replaced controller resumes.
go planTicker(ctx, inv)
// And what the catalogue decided a build meant. The builder's own result is already handled
// above; this is the other half — the control plane is the only one of the three that knows
// which machines run the thing, so it is the one that acts (novox/hq ADR 0072).
@@ -383,6 +387,15 @@ func pushCommand(ctx context.Context, args []string) error {
}
release()
fmt.Printf("\n%d node(s) told\n", len(sending))
// And each machine's memberships, as every other send does (ADR 0160): a push is the one most
// operators run, and on 2026-10-01 it was the one path that issued none.
var told []string
for _, s := range sending {
told = append(told, s.node)
}
if err := issueMemberships(ctx, open, server, told); err != nil {
return err
}
// **A named push leaves the mesh consistent, not just the machine it named** (novox/hq
// issue 057, ADR 0083). Assigning a cross-node consumer mints a provision, and the PROVIDER's
@@ -690,6 +703,51 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
}
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
}
// And every assignment on those machines its membership (novox/hq ADR 0160): composed from the
// same records the bus's accounts are, so what a runtime serves and what its account may are one
// composition. Issued after the declaration, because the runtime it is for arrives with it.
return issueMemberships(ctx, open, server, names)
}
// issueMemberships publishes the membership of every module on the named machines.
func issueMemberships(ctx context.Context, open *stores, server *link.Server, names []string) error {
records, err := open.inventory.BusRecords(ctx)
if err != nil {
return err
}
where := broker.PlacementsOf(records, records.Interchangeable)
bus, ok := server.Bus().(link.OverNATS)
if !ok {
return nil
}
// The declarations are sent and recorded by now; a membership that cannot be issued is said
// and does not unsay them. Every runtime without one serves the shape it derives (ADR 0160), so
// the push stands, the first failure is named once, and the next push tries again.
issued, failed := 0, 0
var first error
for _, node := range names {
for _, d := range records.Assigned[node] {
body, err := json.Marshal(broker.MembershipFor(node, d, where))
if err != nil {
return err
}
if err := bus.PublishMembership(ctx, node, d.Module, body); err != nil {
if first == nil {
first = err
}
failed++
continue
}
issued++
}
}
if issued > 0 {
fmt.Printf(" issued %d membership(s)\n", issued)
}
if failed > 0 {
fmt.Printf(" %d membership(s) could not be issued; the first: %v — the machines keep what "+
"they derive until the next push\n", failed, first)
}
return nil
}
+4 -1
View File
@@ -39,6 +39,9 @@ type meshStatus struct {
// whose is older is still working — and Waiting cannot tell those apart, because the sent
// digest is recorded at send, not at apply.
Reported []machineReported `json:"reported"`
// Plans is what the last merges produced and where each stands (novox/hq ADR 0162): the
// open ones first, each saying its tier, what it waits for, and whether it has waited too long.
Plans []planStatus `json:"plans"`
// Unresolved is every machine that cannot be worked out at all, with what the mesh said when
// it tried. **A machine here is in none of the lists above**: nothing was computed for it, so
// there is nothing to compare it against and nothing it can be behind — which is why a
@@ -152,7 +155,7 @@ func statusAsJSON(asked answers) ([]byte, error) {
out := meshStatus{Machines: len(nodes), Wrong: []machineDoing{},
Quiet: []machineQuiet{}, Behind: []moduleBehind{}, Waiting: []machineWaiting{},
Reported: []machineReported{}, Unresolved: []machineUnresolved{},
Network: asked.network, Adopted: adoptedNodes(nodes)}
Network: asked.network, Adopted: adoptedNodes(nodes), Plans: planStatuses(asked.plans, time.Now())}
// In a stated order, so two readings of an unchanged mesh are the same document.
untakenNodes := make([]string, 0, len(asked.untaken))
for name := range asked.untaken {
+608
View File
@@ -0,0 +1,608 @@
package main
import (
"context"
"flag"
"fmt"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A merge produces a tiered plan the mesh keeps (novox/hq ADR 0162).
//
// The handler that hears the merge computes the plan from the catalogue's one dependency relation,
// writes it to the store, asks the first tier and returns — the receive loop is never held by a
// build. Every outcome taken in advances the plan it belongs to; a ticker advances what outcomes
// alone cannot (a tier waiting for machines to report); a controller replaced mid-plan finds the
// plan where it left it.
// planWaitBound is how long a plan may wait on one thing before `status` names it red.
const planWaitBound = 30 * time.Minute
// tiersOf sorts a set of modules into tiers along the ordering edges among them: tier 0 depends
// on nothing else in the set, tier 1 only on tier 0, and so on. An edge to a module outside the set says
// nothing about the order inside it. A cycle — which the catalogue should never produce — puts
// what remains in one last tier rather than losing it, and is said by the caller.
func tiersOf(set []string, edges []inventory.Edge) [][]string {
in := map[string]bool{}
for _, m := range set {
in[m] = true
}
deps := map[string]map[string]bool{}
for _, m := range set {
deps[m] = map[string]bool{}
}
for _, e := range edges {
// A code dependency — B packages A's source — rebuilds B with A, in the same tier: B's
// build needs nothing of A's first. The other kinds order: stands-on and declared after
// the base is built, built-by after the build machine is built and running — except for
// what the build machine itself stands on. The runtime image is built by the builder and
// the builder is built on the runtime image; the image comes first, built by the builder
// that is running, which is the only one there could be.
if !in[e.From] || !in[e.To] || e.From == e.To || e.Kind == inventory.EdgePackages {
continue
}
if e.Kind == inventory.EdgeBuiltBy && isBaseOf(e.From, e.To, edges, in) {
continue
}
deps[e.From][e.To] = true
}
placed := map[string]bool{}
var tiers [][]string
for len(placed) < len(set) {
var tier []string
for _, m := range set {
if placed[m] {
continue
}
free := true
for d := range deps[m] {
if !placed[d] {
free = false
break
}
}
if free {
tier = append(tier, m)
}
}
if len(tier) == 0 {
// A cycle: everything left, together, and the caller says so.
for _, m := range set {
if !placed[m] {
tier = append(tier, m)
}
}
}
sort.Strings(tier)
for _, m := range tier {
placed[m] = true
}
tiers = append(tiers, tier)
}
return tiers
}
// isBaseOf says whether `to` stands on `base`, directly or through other bases in the set, along
// the build edges alone.
func isBaseOf(base, to string, edges []inventory.Edge, in map[string]bool) bool {
seen := map[string]bool{}
var walk func(string) bool
walk = func(m string) bool {
if m == base {
return true
}
if seen[m] {
return false
}
seen[m] = true
for _, e := range edges {
if e.From == m && in[e.To] && (e.Kind == inventory.EdgeStandsOn || e.Kind == inventory.EdgeDeclared) && walk(e.To) {
return true
}
}
return false
}
return walk(to)
}
// reachableFrom is the moved modules plus everything that depends on them, through every layer:
// what a merge rebuilds. Along the code and build edges only: a module *built by* the build machine
// is not changed by a new build machine, so a built-by edge orders and gates a plan and never
// widens it — the first plan of 2026-10-01 took the whole catalogue along for a controller change.
func reachableFrom(moved []string, edges []inventory.Edge) []string {
in := map[string]bool{}
for _, m := range moved {
in[m] = true
}
for grew := true; grew; {
grew = false
for _, e := range edges {
if e.Kind == inventory.EdgeBuiltBy {
continue
}
if in[e.To] && !in[e.From] {
in[e.From] = true
grew = true
}
}
}
out := make([]string, 0, len(in))
for m := range in {
out = append(out, m)
}
sort.Strings(out)
return out
}
// hasCycle says whether the tiers' last tier holds modules that still depend on each other.
func hasCycle(tiers [][]string, edges []inventory.Edge) bool {
if len(tiers) == 0 {
return false
}
last := map[string]bool{}
for _, m := range tiers[len(tiers)-1] {
last[m] = true
}
for _, e := range edges {
if last[e.From] && last[e.To] {
return true
}
}
return false
}
// planFor is the plan a merge produces: the moved modules and everything reachable from them,
// tiered, with the merge it answers.
func planOfMerge(m link.SourceMoved, moved []string, edges []inventory.Edge) inventory.Plan {
set := reachableFrom(moved, edges)
tiers := tiersOf(set, edges)
modules := map[string]*inventory.PlanModule{}
for _, name := range set {
modules[name] = &inventory.PlanModule{}
}
return inventory.Plan{
ID: fmt.Sprintf("plan-%d", time.Now().UnixNano()),
Repository: m.Owner + "/" + m.Repo,
Commit: m.Commit,
Created: time.Now().UTC(),
State: inventory.PlanBuilding,
Tiers: tiers,
Modules: modules,
}
}
// gates is what the next tier needs running from this one: a module of the tier that a later
// tier is built by — the runtime dependency — and whose policy rolls it out, must be applied by
// the machines running it before the next tier is asked. A base an image stands on need only be
// built; a source another module packages need not even be that.
func gates(p inventory.Plan, edges []inventory.Edge, rollsOut func(string) bool) []string {
if p.Tier >= len(p.Tiers) {
return nil
}
inTier := map[string]bool{}
all := map[string]bool{}
for _, tier := range p.Tiers {
for _, m := range tier {
all[m] = true
}
}
for _, m := range p.Tiers[p.Tier] {
inTier[m] = true
}
later := map[string]bool{}
for _, tier := range p.Tiers[p.Tier+1:] {
for _, m := range tier {
later[m] = true
}
}
seen := map[string]bool{}
var out []string
for _, e := range edges {
if later[e.From] && inTier[e.To] && e.Kind == inventory.EdgeBuiltBy && !seen[e.To] && rollsOut(e.To) &&
!isBaseOf(e.From, e.To, edges, all) {
seen[e.To] = true
out = append(out, e.To)
}
}
sort.Strings(out)
return out
}
// applied says whether every machine running the module has reported since the module was built.
func applied(module string, builtAt time.Time, running []string, reports []inventory.Reported) (bool, []string) {
at := map[string]*time.Time{}
for _, r := range reports {
at[r.Node] = r.At
}
var waiting []string
for _, n := range running {
if t := at[n]; t == nil || t.Before(builtAt) {
waiting = append(waiting, n)
}
}
return len(waiting) == 0, waiting
}
// askTier asks the build machine for every module of the tier, and marks each asked. A module
// the catalogue no longer holds, or whose ask could not be made, is a failure of the plan: a tier
// half asked is a tier that will never complete.
func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) error {
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
byName := map[string]inventory.Entry{}
for _, e := range entries {
byName[e.Manifest.Module] = e
}
now := time.Now().UTC()
for _, name := range p.Tiers[p.Tier] {
state := p.Modules[name]
if state == nil {
state = &inventory.PlanModule{}
p.Modules[name] = state
}
e, known := byName[name]
if !known {
state.State = "failed"
state.Why = "no longer in the catalogue"
p.State = inventory.PlanFailed
p.Note = name + " is no longer in the catalogue"
continue
}
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
fmt.Printf(" tier %d: ", p.Tier)
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 0); err != nil {
state.State = "failed"
state.Why = err.Error()
p.State = inventory.PlanFailed
p.Note = fmt.Sprintf("%s could not be asked for: %v", name, err)
continue
}
state.State = "asked"
state.AskedAt = &now
}
return nil
}
// planBuilt marks a module built (or failed) in every open plan whose current tier holds it, and
// advances what that completes. Called from the daemon's take-in of every outcome.
func planBuilt(ctx context.Context, inv *inventory.Inventory, module, commit, failed string) {
plans, err := inv.OpenPlans(ctx)
if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err)
return
}
now := time.Now().UTC()
for i := range plans {
p := &plans[i]
if p.Tier >= len(p.Tiers) {
continue
}
inTier := false
for _, m := range p.Tiers[p.Tier] {
if m == module {
inTier = true
}
}
if !inTier {
continue
}
state := p.Modules[module]
if state == nil {
state = &inventory.PlanModule{}
p.Modules[module] = state
}
if failed != "" {
state.State = "failed"
state.Why = failed
p.State = inventory.PlanFailed
p.Note = fmt.Sprintf("%s failed to build in tier %d", module, p.Tier)
} else {
state.State = "built"
state.BuiltAt = &now
state.Commit = commit
}
if err := inv.SavePlan(ctx, *p); err != nil {
fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err)
continue
}
if p.State == inventory.PlanFailed {
fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note)
}
}
advancePlans(ctx, inv)
}
// advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built
// and whose gates are applied gives way to the next; the last tier done is the plan done. Called
// after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes.
func advancePlans(ctx context.Context, inv *inventory.Inventory) {
plans, err := inv.OpenPlans(ctx)
if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err)
return
}
if len(plans) == 0 {
return
}
edges, err := inv.Dependencies(ctx)
if err != nil {
fmt.Printf("plans: cannot read the dependencies: %v\n", err)
return
}
rollsOut := func(module string) bool {
u, err := inv.UpgradeOf(ctx, module)
return err == nil && u.RollOut
}
for i := range plans {
p := &plans[i]
for p.Open() {
moved, err := advanceOnce(ctx, inv, p, edges, rollsOut)
if err != nil {
fmt.Printf("%s: %v\n", p.ID, err)
break
}
if err := inv.SavePlan(ctx, *p); err != nil {
fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err)
break
}
if !moved {
break
}
}
}
}
// advanceOnce takes one step of one plan and says whether anything changed.
func advanceOnce(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan,
edges []inventory.Edge, rollsOut func(string) bool) (bool, error) {
if p.Tier >= len(p.Tiers) {
p.State = inventory.PlanDone
fmt.Printf("%s: done — %s at %s, %d tier(s)\n", p.ID, p.Repository, short(p.Commit), len(p.Tiers))
return true, nil
}
tier := p.Tiers[p.Tier]
// Not yet asked: ask.
unasked := 0
for _, m := range tier {
if s := p.Modules[m]; s == nil || s.State == "" {
unasked++
}
}
if unasked == len(tier) {
if err := askTier(ctx, inv, p); err != nil {
return false, err
}
return true, nil
}
// Asked: wait for every build.
var latest time.Time
for _, m := range tier {
s := p.Modules[m]
if s == nil || s.State != "built" {
return false, nil
}
if s.BuiltAt != nil && s.BuiltAt.After(latest) {
latest = *s.BuiltAt
}
}
// Built: wait for what the next tier needs running.
needed := gates(*p, edges, rollsOut)
if len(needed) > 0 {
reports, err := inv.LastReports(ctx)
if err != nil {
return false, err
}
var waiting []string
for _, m := range needed {
running, err := inv.Running(ctx, m)
if err != nil {
return false, err
}
builtAt := latest
if s := p.Modules[m]; s != nil && s.BuiltAt != nil {
builtAt = *s.BuiltAt
}
if ok, on := applied(m, builtAt, running, reports); !ok {
waiting = append(waiting, fmt.Sprintf("%s on %s", m, strings.Join(on, ", ")))
}
}
if len(waiting) > 0 {
note := "tier " + fmt.Sprint(p.Tier) + " built; waiting for " + strings.Join(waiting, "; ") + " to be applied"
changed := p.State != inventory.PlanRolling || p.Note != note
p.State = inventory.PlanRolling
p.Note = note
return changed, nil
}
}
p.Tier++
p.State = inventory.PlanBuilding
p.Note = ""
if p.Tier < len(p.Tiers) {
fmt.Printf("%s: tier %d done; asking tier %d: %s\n", p.ID, p.Tier-1, p.Tier, strings.Join(p.Tiers[p.Tier], ", "))
}
return true, nil
}
// planTicker advances open plans on a timer, for the steps outcomes alone cannot take.
func planTicker(ctx context.Context, inv *inventory.Inventory) {
advancePlans(ctx, inv)
tick := time.NewTicker(30 * time.Second)
defer tick.Stop()
for {
select {
case <-ctx.Done():
return
case <-tick.C:
advancePlans(ctx, inv)
}
}
}
// planLine is one plan as `status` says it.
func planLine(p inventory.Plan, now time.Time) string {
where := fmt.Sprintf("tier %d of %d", min(p.Tier+1, len(p.Tiers)), len(p.Tiers))
switch p.State {
case inventory.PlanDone:
return fmt.Sprintf("%s %s done, %d tier(s)", p.Repository, short(p.Commit), len(p.Tiers))
case inventory.PlanFailed:
return fmt.Sprintf("%s %s FAILED at %s: %s", p.Repository, short(p.Commit), where, p.Note)
}
since := now.Sub(p.Updated).Round(time.Minute)
late := ""
if since > planWaitBound {
late = " — LATE"
}
what := "building"
if p.State == inventory.PlanRolling {
what = p.Note
}
return fmt.Sprintf("%s %s %s, %s for %s%s", p.Repository, short(p.Commit), where, what, since, late)
}
// planFailedBuild marks the module a failed build was for when the result names no module: by the
// repository and path the plan's modules were asked at.
func planFailedBuild(ctx context.Context, inv *inventory.Inventory, result link.BuildResult) {
entries, err := inv.Catalogued(ctx)
if err != nil {
return
}
for _, e := range entries {
if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path {
planBuilt(ctx, inv, e.Manifest.Module, result.Commit, result.Failed)
return
}
}
}
func repositoryMatches(a, b string) bool {
trim := func(s string) string { return strings.ToLower(strings.TrimSuffix(s, ".git")) }
return trim(a) == trim(b) || strings.HasSuffix(trim(a), "/"+trim(b)) || strings.HasSuffix(trim(b), "/"+trim(a))
}
// planStatus is one plan as `status --json` says it.
type planStatus struct {
ID string `json:"id"`
Repository string `json:"repository"`
Commit string `json:"commit"`
State string `json:"state"`
Tier int `json:"tier"`
Tiers int `json:"tiers"`
Waiting string `json:"waiting,omitempty"`
Since time.Time `json:"since"`
Late bool `json:"late"`
}
func planStatuses(plans []inventory.Plan, now time.Time) []planStatus {
out := make([]planStatus, 0, len(plans))
for _, p := range plans {
ps := planStatus{ID: p.ID, Repository: p.Repository, Commit: p.Commit, State: p.State,
Tier: p.Tier, Tiers: len(p.Tiers), Since: p.Updated}
if p.Open() {
ps.Waiting = p.Note
if ps.Waiting == "" {
ps.Waiting = "builds of tier " + fmt.Sprint(p.Tier)
}
ps.Late = now.Sub(p.Updated) > planWaitBound
}
out = append(out, ps)
}
return out
}
// openPlans is the open plans among the recent ones, and how many have waited past the bound.
func openPlans(plans []inventory.Plan) ([]inventory.Plan, int) {
var open []inventory.Plan
late := 0
for _, p := range plans {
if p.Open() {
open = append(open, p)
if time.Since(p.Updated) > planWaitBound {
late++
}
}
}
return open, late
}
// plansCommand says what the last merges produced and where each stands; given an id, one plan
// tier by tier with every module's state.
func plansCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("plans", flag.ContinueOnError)
limit := set.Int("n", 10, "how many to show")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
now := time.Now()
if len(positionals) == 1 {
p, err := inv.PlanByID(ctx, positionals[0])
if err != nil {
return err
}
fmt.Printf("%s — %s\n", p.ID, planLine(p, now))
for i, tier := range p.Tiers {
marker := " "
if i == p.Tier && p.Open() {
marker = ">"
}
fmt.Printf("%s tier %d\n", marker, i)
for _, m := range tier {
s := p.Modules[m]
state := "not yet asked"
if s != nil && s.State != "" {
state = s.State
if s.Commit != "" {
state += " from " + short(s.Commit)
}
if s.Why != "" {
state += ": " + s.Why
}
}
fmt.Printf(" %-22s %s\n", m, state)
}
}
return nil
}
if len(positionals) == 2 && positionals[0] == "stop" {
p, err := inv.PlanByID(ctx, positionals[1])
if err != nil {
return err
}
if !p.Open() {
return fmt.Errorf("%s is already %s", p.ID, p.State)
}
p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
if err := inv.SavePlan(ctx, p); err != nil {
return err
}
fmt.Printf("%s stopped at tier %d of %d; what was asked still builds and registers, nothing further is asked\n",
p.ID, p.Tier, len(p.Tiers))
return nil
}
plans, err := inv.RecentPlans(ctx, *limit)
if err != nil {
return err
}
if len(plans) == 0 {
fmt.Println("no merge has produced a plan yet")
return nil
}
for _, p := range plans {
fmt.Printf("%-28s %s\n", p.ID, planLine(p, now))
}
return nil
}
+117
View File
@@ -0,0 +1,117 @@
package main
import (
"testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A merge produces a tiered plan (novox/hq ADR 0162): what moved and everything reachable from it,
// sorted so a tier depends only on earlier ones — with the three kinds of dependency told apart.
func TestAMergeIsPlannedInTiersAlongTheThreeKindsOfDependency(t *testing.T) {
edges := []inventory.Edge{
// build dependencies: images on the runtime, a plugin on one of them
{From: "shop", To: "mesh-tools", Kind: inventory.EdgeStandsOn},
{From: "postgres", To: "mesh-tools", Kind: inventory.EdgeStandsOn},
{From: "shop-plugin", To: "shop", Kind: inventory.EdgeDeclared},
// a code dependency: the proxy packages the controller's source — same tier
{From: "route-proxy", To: "mesh-controller", Kind: inventory.EdgePackages},
{From: "builder", To: "mesh-controller", Kind: inventory.EdgePackages},
// runtime dependencies: everything source-built is built by the builder
{From: "shop", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "postgres", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "shop-plugin", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "route-proxy", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "mesh-controller", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "mesh-tools", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "builder", To: "mesh-tools", Kind: inventory.EdgeStandsOn},
{From: "unrelated", To: "alpine", Kind: inventory.EdgeStandsOn},
}
// The runtime image moved: everything on it, and what is built by what is on it.
set := reachableFrom([]string{"mesh-tools"}, edges)
// What stands on the runtime, and the builder that stands on it; not the controller, which the
// builder merely builds, nor the proxy that packages the controller.
want := []string{"builder", "mesh-tools", "postgres", "shop", "shop-plugin"}
if len(set) != len(want) {
t.Fatalf("reachable from the runtime: %v, want %v", set, want)
}
tiers := tiersOf(set, edges)
pos := map[string]int{}
for i, tier := range tiers {
for _, m := range tier {
pos[m] = i
}
}
if pos["mesh-tools"] != 0 || pos["builder"] != 1 {
t.Fatalf("the runtime then the builder: %v", tiers)
}
if !(pos["shop"] > pos["builder"] && pos["postgres"] > pos["builder"]) {
t.Fatalf("what the builder builds comes after the builder: %v", tiers)
}
if pos["shop-plugin"] <= pos["shop"] {
t.Fatalf("a plugin after what it is declared on: %v", tiers)
}
if hasCycle(tiers, edges) {
t.Fatalf("no cycle here: %v", tiers)
}
// The controller alone moved: the proxy with it, nothing else.
small := reachableFrom([]string{"mesh-controller"}, edges)
if len(small) != 3 {
t.Fatalf("a controller merge rebuilds the controller and what packages it: %v", small)
}
// The builder packages the controller's source (same tier by that edge) and the controller is
// built by the builder (next tier by that one): the builder first, then the controller and the
// proxy together — a code dependency in one tier, a runtime dependency across tiers.
smallTiers := tiersOf(small, edges)
if len(smallTiers) != 2 || smallTiers[0][0] != "builder" || len(smallTiers[1]) != 2 {
t.Fatalf("the builder, then the controller and the proxy together: %v", smallTiers)
}
// The builder alone moved: the builder, and nothing it builds.
if only := reachableFrom([]string{"builder"}, edges); len(only) != 1 {
t.Fatalf("a build machine change rebuilds the build machine alone: %v", only)
}
// Only a runtime dependency gates on deployment, and only when the module rolls out.
p := planOfMerge(link.SourceMoved{Owner: "novox", Repo: "mesh-tools", Commit: "abc"}, []string{"mesh-tools"}, edges)
p.Tier = pos["builder"]
rollsOut := func(m string) bool { return m == "builder" }
if g := gates(p, edges, rollsOut); len(g) != 1 || g[0] != "builder" {
t.Fatalf("the builder gates the tier after it: %v", g)
}
p.Tier = 0
if g := gates(p, edges, rollsOut); len(g) != 0 {
t.Fatalf("the runtime image is a build dependency and gates nothing: %v", g)
}
if g := gates(p, edges, func(string) bool { return false }); len(g) != 0 {
t.Fatalf("a module that only records its upgrade gates nothing: %v", g)
}
}
// A gate is open once every machine running the module has reported after it was built.
func TestAGateOpensWhenTheMachinesHaveReportedSinceTheBuild(t *testing.T) {
built := time.Date(2026, 10, 1, 15, 0, 0, 0, time.UTC)
before, after := built.Add(-time.Minute), built.Add(time.Minute)
reports := []inventory.Reported{{Node: "anchor", At: &after}, {Node: "home-server", At: &before}}
ok, waiting := applied("builder", built, []string{"anchor", "home-server"}, reports)
if ok || len(waiting) != 1 || waiting[0] != "home-server" {
t.Fatalf("one machine has not reported since the build: ok=%v waiting=%v", ok, waiting)
}
if ok, _ := applied("builder", built, []string{"anchor"}, reports); !ok {
t.Fatal("the machine that reported after the build holds the gate open")
}
if ok, _ := applied("builder", built, nil, reports); !ok {
t.Fatal("a module running nowhere gates nothing")
}
}
// A cycle is not lost: what remains is one last tier, and the caller says so.
func TestACycleIsOneLastTierAndSaidSo(t *testing.T) {
edges := []inventory.Edge{{From: "a", To: "b", Kind: inventory.EdgeStandsOn}, {From: "b", To: "a", Kind: inventory.EdgeStandsOn}}
tiers := tiersOf([]string{"a", "b"}, edges)
if len(tiers) != 1 || len(tiers[0]) != 2 || !hasCycle(tiers, edges) {
t.Fatalf("a cycle should be one tier of two, said: %v", tiers)
}
}
+38 -10
View File
@@ -69,6 +69,14 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
return []string{"builds", m}, nil
}
return []string{"builds"}, nil
case "plans":
if id := str("stop"); id != "" {
return []string{"plans", "stop", id}, nil
}
if id := str("id"); id != "" {
return []string{"plans", id}, nil
}
return []string{"plans"}, nil
case "plan":
if err := need("node"); err != nil {
return nil, err
@@ -79,6 +87,16 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
return nil, err
}
return []string{verb, str("node"), str("module")}, nil
case "pin":
if err := need("node", "provision", "from", "module"); err != nil {
return nil, err
}
return []string{"pin", str("node"), str("provision"), str("from"), str("module")}, nil
case "unpin":
if err := need("node", "provision"); err != nil {
return nil, err
}
return []string{"unpin", str("node"), str("provision")}, nil
case "push":
// Sent and not waited for: the asker reads `status` for what the machine did, which is
// what a person at a shell does too. A tool call that blocked for a push's whole apply would
@@ -102,15 +120,6 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
// the answer the caller needs.
return []string{"rotate"}, nil
case "build":
// Three shapes, as the command has them: a repository, a base every module built on it
// is rebuilt from (`--on`), or everything behind its source (`--behind`). Asked, not
// waited for, the same as a single build.
if on := str("on"); on != "" {
return []string{"build", "--on", on, "--wait", "0"}, nil
}
if b := str("behind"); b != "" && b != "no" && b != "false" {
return []string{"build", "--behind", "--wait", "0"}, nil
}
if err := need("repository"); err != nil {
return nil, err
}
@@ -185,7 +194,7 @@ func seatToolHandlers() (map[string]link.ToolHandler, error) {
}
continue
}
if _, err := argvFor(verb, map[string]any{"node": "x", "module": "x", "repository": "x"}); err != nil {
if _, err := argvFor(verb, sampleArguments(v)); err != nil {
return nil, fmt.Errorf("the %s seat's row declares %q, which this control plane cannot run: %w",
catalogue.ControllerSeatName, verb, err)
}
@@ -224,3 +233,22 @@ func seatTools() map[string]any {
}
return map[string]any{"seats": seats}
}
// sampleArguments is one of every argument a verb's schema requires, so the check at start proves the
// verb runnable rather than that it happens to want the arguments the check guessed.
func sampleArguments(v catalogue.Verb) map[string]any {
sample := map[string]any{"node": "x", "module": "x", "repository": "x"}
switch required := v.Input["required"].(type) {
case []string:
for _, k := range required {
sample[k] = "x"
}
case []any:
for _, k := range required {
if name, ok := k.(string); ok {
sample[name] = "x"
}
}
}
return sample
}
-16
View File
@@ -139,19 +139,3 @@ func TestAJSONVerbsAnswerIsItsStandardOutput(t *testing.T) {
t.Fatalf("stderr and stdout are both what the command said: %s", answer.Output)
}
}
// The build tool has the command's three shapes (ADR 0157's follow-up, 2026-10-01): a repository, a
// base whose dependents are rebuilt, or everything behind its source — each asked, not waited for.
func TestTheBuildToolRebuildsWhatStandsOnABase(t *testing.T) {
argv, _ := argvFor("build", map[string]any{"on": "mesh-tools"})
if strings.Join(argv, " ") != "build --on mesh-tools --wait 0" {
t.Fatalf("a base: %v", argv)
}
argv, _ = argvFor("build", map[string]any{"behind": "yes"})
if strings.Join(argv, " ") != "build --behind --wait 0" {
t.Fatalf("behind: %v", argv)
}
if _, err := argvFor("build", map[string]any{}); err == nil {
t.Fatal("a build naming nothing was accepted")
}
}
+16
View File
@@ -138,6 +138,18 @@ func printStatus(asked answers) error {
len(quiet), strings.Join(said, "\n "))
}
if open, late := openPlans(asked.plans); len(open) > 0 {
fmt.Printf("%d plan(s) open", len(open))
if late > 0 {
fmt.Printf(", %d waiting past %s", late, planWaitBound)
}
fmt.Println(":")
for _, p := range open {
fmt.Printf(" %s\n", planLine(p, time.Now()))
}
fmt.Println()
}
if len(behind) > 0 {
var names []string
for m := range behind {
@@ -342,6 +354,10 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) {
if err != nil {
return answers{}, err
}
out.plans, err = inv.RecentPlans(ctx, 5)
if err != nil {
return answers{}, err
}
// And which machines are not running what the mesh would send them. The same question as a
// module being behind its source, one level down: that one says the catalogue is out of date,
+59 -22
View File
@@ -295,17 +295,31 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
"built from\n", m.Owner, m.Repo, m.Base, m.Commit)
return nil
}
against, err := inv.BuiltAgainst(ctx)
// A merge produces a plan the mesh keeps (novox/hq ADR 0162): what moved and everything that
// depends on it, along the catalogue's one dependency relation, sorted into tiers. The plan is
// written before any build is asked; the first tier is asked; this returns. Outcomes advance it.
edges, err := inv.Dependencies(ctx)
if err != nil {
return notNow(err)
}
ordered := orderByBases(moved, against)
names := make([]string, 0, len(ordered))
for _, e := range ordered {
names = append(names, e.Manifest.Module)
var movedNames []string
for _, e := range moved {
movedNames = append(movedNames, e.Manifest.Module)
}
fmt.Printf("%s/%s merged into %s (%.8s); building %s\n",
m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", "))
plan := planOfMerge(m, movedNames, edges)
if hasCycle(plan.Tiers, edges) {
fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n",
strings.Join(plan.Tiers[len(plan.Tiers)-1], ", "))
}
if err := inv.SavePlan(ctx, plan); err != nil {
return notNow(err)
}
var tiers []string
for i, t := range plan.Tiers {
tiers = append(tiers, fmt.Sprintf("%d: %s", i, strings.Join(t, ", ")))
}
fmt.Printf("%s/%s merged into %s (%.8s); plan %s, %d module(s) in %d tier(s)\n %s\n",
m.Owner, m.Repo, m.Base, m.Commit, plan.ID, len(plan.Modules), len(plan.Tiers), strings.Join(tiers, "\n "))
if len(packaging) > 0 {
var also []string
for _, e := range packaging {
@@ -314,22 +328,11 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
fmt.Printf(" %s package source from it, so they are rebuilt and their own source record "+
"is left where it is\n", strings.Join(also, ", "))
}
var failed []string
for _, e := range ordered {
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 20*time.Minute); err != nil {
fmt.Printf(" %s: %v\n", e.Manifest.Module, err)
failed = append(failed, e.Manifest.Module)
// A base that failed is a reason to stop: what stands on it would be built against
// the old one, and report success (novox/hq 04-ISSUES/131).
if standsOn(ordered, e.Manifest.Module, against) {
fmt.Printf(" stopping: %s is a base of what was still to build\n", e.Manifest.Module)
break
}
}
if err := askTier(ctx, inv, &plan); err != nil {
return notNow(err)
}
if len(failed) > 0 {
fmt.Printf("%d of %d not built: %s\n", len(failed), len(ordered), strings.Join(failed, ", "))
if err := inv.SavePlan(ctx, plan); err != nil {
return notNow(err)
}
return nil
}
@@ -522,3 +525,37 @@ func isHistory(mergedAt string, seen time.Time) bool {
}
return at.Before(seen)
}
// dependentsOf is every catalogued module that stands on one of the moved modules, directly or
// through another dependent, and is not itself among them — in the catalogue's order, so the
// answer is the same each time. A module standing on nothing that moved is left alone: a merge
// rebuilds what it changed and what is built on top of that, not the catalogue.
func dependentsOf(moved, entries []inventory.Entry, against map[string][]string) []inventory.Entry {
bases := map[string]bool{}
for _, e := range moved {
bases[e.Manifest.Module] = true
}
var out []inventory.Entry
taken := map[string]bool{}
for grew := true; grew; {
grew = false
for _, e := range entries {
name := e.Manifest.Module
if bases[name] || taken[name] {
continue
}
for base := range bases {
if standsOnModule(e, base, against) {
taken[name] = true
out = append(out, e)
grew = true
break
}
}
}
for _, e := range out {
bases[e.Manifest.Module] = true
}
}
return out
}
+8 -1
View File
@@ -161,8 +161,15 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
Queue: "holders",
AckWaitSeconds: 60,
MaxDeliver: 5,
// **One in flight.** A holder works one ask at a time, so the server hands it one at a
// time: with the default of many, every ask behind the one being worked was delivered,
// left unacknowledged for the length of the work, redelivered after the ack wait, and
// after the fifth time dropped — on 2026-10-01 twenty-six of forty-three builds asked in
// two minutes were never built, and the queue read as empty (novox/hq issue 186).
MaxAckPending: 1,
Why: fmt.Sprintf("%s on %s holds %s; it acknowledges after the work is done, so a "+
"crash mid-work redelivers rather than loses", module, node, seat.Name),
"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
}
+13
View File
@@ -153,3 +153,16 @@ func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) {
has(t, perms.Publish, "$JS.ACK.NODES."+c.Name+".>")
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):
// asks queued behind the one being worked wait in the stream rather than being delivered,
// left to expire and dropped after the fifth redelivery.
func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) {
c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}})
if !found {
t.Fatal("a seat that accepts work has no worker")
}
if c.MaxAckPending != 1 {
t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending)
}
}
+1
View File
@@ -144,6 +144,7 @@ func (j *JetStream) EnsureStream(s Stream) error {
MaxMsgsPerSubject: int64(s.MaxMsgsPerSubject),
Description: s.Why,
}
want.AllowDirect = s.Direct
if s.Retention == RetentionLastPerSubject {
// Last-per-subject is a limits stream with one message kept per subject, not a
// retention policy of its own — the state shape, spelled the way the server spells it.
+131
View File
@@ -0,0 +1,131 @@
package broker
import (
"sort"
"strings"
)
// What the mesh issues an assignment to serve and to reach (novox/hq ADR 0160).
//
// A module's code names its tools and its events; **where they land is the mesh's to decide**, and
// it decided it twice — once in the runtime, once here, by one rule compiled into both. Now the
// controller composes a membership for every module on every machine and publishes it to a subject
// only that assignment reads; the runtime serves exactly what the membership says, and the account's
// grant is the same composition read the other way. The shape issued today is the shape the mesh
// already had, so nothing moves when a membership first arrives; only who decides it moves.
// Membership is one assignment's subjects: what this instance of a module on this machine serves,
// and what it may reach.
type Membership struct {
Node string `json:"node"`
Module string `json:"module"`
// Serves is every address a tool of this instance answers on. `{tool}` stands for the tool's
// own name, which the module knows and the mesh does not need to: the mesh issues the address,
// the runtime fills the name. An address with a queue is shared with the module's other
// instances, and the bus hands each call to one of them; an address without is this instance's.
Serves []Served `json:"serves"`
// Seats is every verb of a seat this instance holds, at the subject the seat's callers use.
Seats []SeatServed `json:"seats,omitempty"`
// Emits is where an event of this module lands; `{event}` stands for the event's name.
Emits string `json:"emits"`
// Reaches is each tool this module may call, `<module>.<tool>`, to the subjects that reach it:
// the first is whichever instance answers, when the mesh issued one; the rest name a machine.
Reaches map[string][]string `json:"reaches,omitempty"`
// Tools is where this instance answers what it serves — the runtime's one verb of its own.
Tools string `json:"tools"`
}
// Served is one address a tool is answered on.
type Served struct {
Subject string `json:"subject"`
Queue string `json:"queue,omitempty"`
}
// SeatServed is one verb of a held seat, where its callers ask.
type SeatServed struct {
Seat string `json:"seat"`
Verb string `json:"verb"`
Subject string `json:"subject"`
}
// MembershipSubject is the one address a runtime derives for itself: where its own membership is
// published, from the two names its credential carries. Everything else is in the membership.
func MembershipSubject(node, module string) string {
return "mesh.assignment." + node + "." + module
}
// Placements is where every module runs, for deciding which instance answers for the module.
type Placements struct {
// Nodes is each module's machines.
Nodes map[string][]string
// Interchangeable is each module whose definition says its instances are the same anywhere,
// so the module's plain subject is issued to all of them in one queue.
Interchangeable map[string]bool
}
// AnswersForTheModule says whether an instance of a module on one machine is issued the module's
// plain subject: when it is the only instance, or when the definition says instances are
// interchangeable. A stateful module on two machines gets only its machines' subjects, so a call
// that names none reaches nothing rather than the wrong store.
func (p Placements) AnswersForTheModule(module string) bool {
return len(p.Nodes[module]) <= 1 || p.Interchangeable[module]
}
// MembershipFor composes one assignment's membership from what it declared and where everything
// runs. The subjects are the ones PermissionsFor grants, derived here once more only until the
// grant itself is read from the membership — which is the next step, not this one.
func MembershipFor(node string, d Declared, where Placements) Membership {
own := "mesh.mod." + d.Module
m := Membership{
Node: node, Module: d.Module,
Emits: own + ".event.{event}",
Tools: own + ".tool.tools",
}
// This machine's address always; the module's when this instance answers for the module.
m.Serves = append(m.Serves, Served{Subject: own + ".tool.{tool}." + node})
if where.AnswersForTheModule(d.Module) {
m.Serves = append(m.Serves, Served{Subject: own + ".tool.{tool}", Queue: "serve." + d.Module})
}
for _, s := range d.Holds {
for _, verb := range s.Serves {
m.Seats = append(m.Seats, SeatServed{Seat: s.Name, Verb: verb, Subject: seatToolSubject(s, verb, node)})
}
}
if len(d.Invokes) > 0 {
m.Reaches = map[string][]string{}
for _, t := range d.Invokes {
if t == "*" || strings.HasPrefix(t, "seat:") {
continue // every tool, or a role's: addressed by name, not resolved per instance
}
module, tool, ok := strings.Cut(t, ".")
if !ok {
continue
}
var reach []string
if where.AnswersForTheModule(module) {
reach = append(reach, "mesh.mod."+module+".tool."+tool)
}
nodes := append([]string{}, where.Nodes[module]...)
sort.Strings(nodes)
for _, n := range nodes {
reach = append(reach, "mesh.mod."+module+".tool."+tool+"."+n)
}
m.Reaches[t] = reach
}
}
return m
}
// PlacementsOf reads where everything runs from the records the bus's accounts are composed from.
func PlacementsOf(r Records, interchangeable map[string]bool) Placements {
p := Placements{Nodes: map[string][]string{}, Interchangeable: interchangeable}
for node, declared := range r.Assigned {
for _, d := range declared {
p.Nodes[d.Module] = append(p.Nodes[d.Module], node)
}
}
for _, nodes := range p.Nodes {
sort.Strings(nodes)
}
return p
}
+75
View File
@@ -0,0 +1,75 @@
package broker
import (
"reflect"
"testing"
)
// The mesh issues an assignment's subjects (novox/hq ADR 0160): a module alone on one machine
// answers for the module and for its machine; a stateful module on two machines answers only for
// each machine; one that says its instances are interchangeable answers for the module everywhere;
// a holder serves its seat's verbs; and what a module may reach is resolved the same way.
func TestAMembershipIsIssuedFromWhereEverythingRuns(t *testing.T) {
records := Records{Assigned: map[string][]Declared{
"anchor": {
{Module: "postgres", Serves: []string{"query"}, Holds: []Seat{{Name: "mesh-store", Scope: "mesh", Serves: []string{"databases", "query"}}}},
{Module: "catalog", Invokes: []string{"postgres.query", "search.find"}},
},
"home-server": {
{Module: "postgres"},
{Module: "search"},
{Module: "dashboard", Invokes: []string{"postgres.query"}},
},
"laptop": {{Module: "search"}},
}, Interchangeable: map[string]bool{"search": true}}
where := PlacementsOf(records, records.Interchangeable)
pg := MembershipFor("anchor", records.Assigned["anchor"][0], where)
if !reflect.DeepEqual(pg.Serves, []Served{{Subject: "mesh.mod.postgres.tool.{tool}.anchor"}}) {
t.Fatalf("a stateful module on two machines answers only for its machine: %+v", pg.Serves)
}
if len(pg.Seats) != 2 || pg.Seats[0].Subject != "mesh.seat.mesh-store.tool.databases" {
t.Fatalf("the holder serves the seat's verbs at the seat's subjects: %+v", pg.Seats)
}
if pg.Emits != "mesh.mod.postgres.event.{event}" || pg.Tools != "mesh.mod.postgres.tool.tools" {
t.Fatalf("events and the tools verb: %+v", pg)
}
search := MembershipFor("laptop", records.Assigned["laptop"][0], where)
if !reflect.DeepEqual(search.Serves, []Served{
{Subject: "mesh.mod.search.tool.{tool}.laptop"},
{Subject: "mesh.mod.search.tool.{tool}", Queue: "serve.search"},
}) {
t.Fatalf("an interchangeable module answers for the module in the queue too: %+v", search.Serves)
}
dashboard := MembershipFor("home-server", records.Assigned["home-server"][2], where)
if !reflect.DeepEqual(dashboard.Serves, []Served{
{Subject: "mesh.mod.dashboard.tool.{tool}.home-server"},
{Subject: "mesh.mod.dashboard.tool.{tool}", Queue: "serve.dashboard"},
}) {
t.Fatalf("a module alone on one machine answers for the module: %+v", dashboard.Serves)
}
if !reflect.DeepEqual(dashboard.Reaches["postgres.query"],
[]string{"mesh.mod.postgres.tool.query.anchor", "mesh.mod.postgres.tool.query.home-server"}) {
t.Fatalf("reaching a stateful module names each machine and no plain subject: %v", dashboard.Reaches)
}
catalog := MembershipFor("anchor", records.Assigned["anchor"][1], where)
if !reflect.DeepEqual(catalog.Reaches["search.find"],
[]string{"mesh.mod.search.tool.find", "mesh.mod.search.tool.find.home-server", "mesh.mod.search.tool.find.laptop"}) {
t.Fatalf("reaching an interchangeable module offers the plain subject first: %v", catalog.Reaches)
}
if MembershipSubject("anchor", "postgres") != "mesh.assignment.anchor.postgres" {
t.Fatal("the one subject a runtime derives for itself")
}
}
func TestAnAccountMayReadItsOwnMembershipAndNoOthers(t *testing.T) {
perms, err := PermissionsFor(Principal{Kind: KindModule, Node: "anchor", Module: "postgres", PasswordHash: "x"})
if err != nil {
t.Fatal(err)
}
has(t, perms.Subscribe, "mesh.assignment.anchor.postgres")
has(t, perms.Publish, "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.anchor.postgres")
hasNot(t, perms.Subscribe, "mesh.assignment.>")
}
+8 -2
View File
@@ -173,8 +173,10 @@ func PermissionsFor(p Principal) (Permissions, error) {
switch p.Kind {
case KindController:
// The controller owns the mesh's own traffic and the streams. It is the only writer of
// stream definitions (design 25 §3), so it alone reaches the JetStream API.
pub = []string{"mesh.control.>", "mesh.node.>", "$JS.API.>"}
// stream definitions (design 25 §3), so it alone reaches the JetStream API — and it alone
// issues memberships (novox/hq ADR 0160), which it publishes into the assignments stream
// after each push; refused by the server on 2026-10-01 until this line named them.
pub = []string{"mesh.control.>", "mesh.node.>", "mesh.assignment.>", "$JS.API.>"}
// **And where its consumers deliver.** A push consumer delivers on `_DELIVER.<its name>`,
// and a client bound to it subscribes exactly that; the server refused it for every
// principal the first time one bound a consumer (2026-09-28). Each kind below is granted
@@ -298,6 +300,10 @@ func PermissionsFor(p Principal) (Permissions, error) {
// away — no other principal may subscribe this namespace, and a caller's authority is
// still granted per tool, by name, on the publish side.
sub = append(sub, own+".tool.>")
// Its own membership (ADR 0160): the one subject a runtime derives for itself, read
// directly from the stream and followed live. Nothing else's.
sub = append(sub, MembershipSubject(p.Node, p.Module))
pub = append(pub, "$JS.API.DIRECT.GET."+AssignmentsStream+"."+MembershipSubject(p.Node, p.Module))
// 1b. The tools it calls, if its manifest says it calls any (novox/hq ADR 0152). The same
// grant a person gets and derived the same way, so "what may this module ask" is
+14
View File
@@ -47,8 +47,14 @@ type Stream struct {
// Why is carried into the assertion so an operator reading the server's own state finds the
// reason there, rather than only in a repository they may not have.
Why string
// Direct lets a client read a subject's last message without a consumer, which is how a
// runtime reads its own membership with no JetStream API beyond one request (ADR 0160).
Direct bool
}
// AssignmentsStream holds every assignment's membership, the newest per subject.
const AssignmentsStream = "ASSIGNMENTS"
// MeshStreams is the foundation set, in the order a person reads it.
//
// **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live
@@ -78,6 +84,14 @@ func MeshStreams() []Stream {
Why: "one declaration per node, always the newest; a node that sees sequence n refuses " +
"n-1 by construction (issue 107)",
},
{
Name: AssignmentsStream,
Subjects: []string{"mesh.assignment.*.*"},
Retention: RetentionLastPerSubject,
Direct: true,
Why: "one membership per assignment, always the newest: what the mesh issued this module " +
"on this machine to serve and to reach (ADR 0160); read directly by the runtime it is for",
},
{
Name: EventsStream,
// A seat's own events ride here too: they are 1:many like any event, and the
+4 -3
View File
@@ -123,9 +123,10 @@ func subjectMatches(filter, subject string) bool {
// Each relationship's retention is the thing that makes it what it is (design 29 §4).
func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) {
want := map[string]Retention{
"CONTROL": RetentionWorkQueue,
"NODES": RetentionLastPerSubject,
"EVENTS": RetentionLimits,
"CONTROL": RetentionWorkQueue,
"NODES": RetentionLastPerSubject,
"EVENTS": RetentionLimits,
"ASSIGNMENTS": RetentionLastPerSubject,
}
got := map[string]Retention{}
for _, s := range MeshStreams() {
+7 -7
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled
users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.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"] }
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>"] }
allow_responses: { max: 1, ttl: "1m" }
} }
@@ -37,18 +37,18 @@ accounts {
subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] }
} }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "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.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
allow_responses: { max: 1, ttl: "1m" }
} }
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit"] }
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.audit"] }
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
allow_responses: { max: 1, ttl: "1m" }
} }
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["_INBOX.two.shop.>", "mesh.mod.shop.tool.>"] }
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] }
allow_responses: { max: 1, ttl: "1m" }
} }
]
+3
View File
@@ -47,6 +47,9 @@ type Records struct {
Enrolling []string
// People is each person's name against the tools they may invoke, `*` for an administrator.
People map[string][]string
// Interchangeable is each module whose definition says its instances are the same anywhere
// (ADR 0160), which decides whether the module's plain subject is issued to every instance.
Interchangeable map[string]bool
}
// Users is every user the composed file should contain, in the order it will be written.
+5 -5
View File
@@ -25,7 +25,7 @@ func reachable() Node {
func onNetwork(nodes ...string) map[string][]Provider {
out := make([]Provider, 0, len(nodes))
for _, n := range nodes {
out = append(out, Provider{Node: n, At: n + ".internal"})
out = append(out, Provider{Node: n, At: n + ".internal", Module: "postgres"})
}
return map[string][]Provider{"postgres-database": out}
}
@@ -99,7 +99,7 @@ func TestSayingWhichOneSettlesIt(t *testing.T) {
got, err := Resolve(brokeredShelf(), []string{"meshboard"}, reachable(),
World{
Offered: onNetwork("anchor", "archive"),
Pinned: map[string]string{"postgres-database": "archive"},
Pinned: map[string]Chosen{"postgres-database": {Node: "archive", Module: "postgres"}},
})
if err != nil {
t.Fatal(err)
@@ -115,7 +115,7 @@ func TestBeingPointedAtAMachineThatDoesNotProvideItIsRefused(t *testing.T) {
_, err := Resolve(brokeredShelf(), []string{"meshboard"}, reachable(),
World{
Offered: onNetwork("anchor", "archive"),
Pinned: map[string]string{"postgres-database": "somewhere-else"},
Pinned: map[string]Chosen{"postgres-database": {Node: "somewhere-else", Module: "postgres"}},
})
if err == nil {
t.Fatal("a machine was silently given a different database from the one chosen")
@@ -131,12 +131,12 @@ func TestOneProviderDoesNotOverruleAChoice(t *testing.T) {
_, err := Resolve(brokeredShelf(), []string{"meshboard"}, reachable(),
World{
Offered: onNetwork("anchor"),
Pinned: map[string]string{"postgres-database": "archive"},
Pinned: map[string]Chosen{"postgres-database": {Node: "archive", Module: "postgres"}},
})
if err == nil {
t.Fatal("the only database was used although another was chosen")
}
if !strings.Contains(err.Error(), "only anchor provides it") {
if !strings.Contains(err.Error(), "only anchor/postgres provides it") {
t.Fatalf("the refusal does not say what is available: %v", err)
}
}
+100
View File
@@ -0,0 +1,100 @@
package catalogue
import (
"sort"
)
// Chosen is the provider somebody named for a provision: the module, and the node it runs on. Both,
// always (novox/hq #258) — a provision comes from a module, and the same module on two machines is
// two answers, so neither half alone says which. Module is empty only on a record made before this
// was asked, and such a record is honoured exactly as long as it is unambiguous.
type Chosen struct {
Node string
Module string
}
func (c Chosen) String() string {
if c.Module == "" {
return c.Node
}
return c.Node + "/" + c.Module
}
// matches is whether this provider is the one chosen.
func (c Chosen) matches(p Provider) bool {
return p.Node == c.Node && (c.Module == "" || p.Module == c.Module)
}
// among is every offered provider the choice names — one, when the choice is whole.
func (c Chosen) among(where []Provider) []Provider {
var out []Provider
for _, p := range where {
if c.matches(p) {
out = append(out, p)
}
}
return out
}
// nameOf is how a refusal names a provider: the node and the module on it.
func nameOf(p Provider) string {
return Chosen{Node: p.Node, Module: p.Module}.String()
}
// providerNames is every provider named, sorted, for a refusal to list.
func providerNames(where []Provider) []string {
out := make([]string, 0, len(where))
for _, p := range where {
out = append(out, nameOf(p))
}
sort.Strings(out)
return out
}
// providersHere is which modules in this node's own set offer a provision, sorted.
func providersHere(catalogue map[string]Manifest, here func(string) bool, want string) []string {
var out []string
for name, m := range catalogue {
if !here(name) {
continue
}
for _, o := range m.Offers() {
if o == want {
out = append(out, name)
break
}
}
}
sort.Strings(out)
return out
}
// servedByOne is what one provider beside the consumer says a consumer needs to know, or nothing.
//
// Serving is *whether* a need is created at all when the provider is on this same machine (novox/hq
// 04-ISSUES/038's sibling): a need never created is a binding the consumer never gets. The manifest
// alone answers that; the values are settled later, with the node's settings.
func servedByOne(m Manifest, want string) map[string]any {
if _, ok := m.Serves[want]; ok {
return ServedOn(m, want, nil)
}
return nil
}
// sharedByOne is the own secret that provider names as its credential (ADR 0158), or "" when it
// gives each consumer its own.
func sharedByOne(m Manifest, want string) string {
if own, shared := m.SharedCredentialOf(want); shared {
return own
}
return ""
}
func oneOf(list []string, s string) bool {
for _, x := range list {
if x == s {
return true
}
}
return false
}
+21
View File
@@ -107,6 +107,27 @@ func TestCanHoldJudgesClaimScopeAndWhatTheSeatDelivers(t *testing.T) {
if err := CanHold(cannotAnswer, seat); err == nil || !strings.Contains(err.Error(), `does not provide "mesh-bus"`) {
t.Fatalf("a holder that cannot answer for the seat was allowed: %v", err)
}
// A holder's own tools need not be the seat's verbs: the claim may name what it serves for the
// role (ADR 0160), and then only those count — and only the seat's verbs may be named.
promising := seat
promising.Serves = []Verb{{Name: "databases"}, {Name: "query"}}
engine := newBroker()
engine.Tools = []string{"engine_list_databases", "engine_query"}
if err := CanHold(engine, promising); err == nil || !strings.Contains(err.Error(), "does not serve databases, query") {
t.Fatalf("a holder whose tools are not the seat's verbs was allowed without saying what it serves: %v", err)
}
engine.Claims[0].Serves = []string{"databases", "query"}
if err := CanHold(engine, promising); err != nil {
t.Fatalf("a claim naming the seat's verbs was refused: %v", err)
}
engine.Claims[0].Serves = []string{"databases"}
if err := CanHold(engine, promising); err == nil || !strings.Contains(err.Error(), "does not serve query") {
t.Fatalf("a claim naming half the verbs was allowed: %v", err)
}
engine.Claims[0].Serves = []string{"databases", "query", "engine_query"}
if err := CanHold(engine, promising); err == nil || !strings.Contains(err.Error(), "does not promise") {
t.Fatalf("a claim naming a verb the seat never promised was allowed: %v", err)
}
// And the judgement follows the store's row, not a compiled copy.
busSeatDelivering(t, "amqp")
seat, _ = SeatNamed("mesh-broker")
+31
View File
@@ -52,6 +52,22 @@ type Claim struct {
Name string `json:"name"`
// Scope defaults to the node, which is where nearly everything singular is singular.
Scope string `json:"scope,omitempty"`
// Serves names the seat's verbs this module implements for the role, when its own tools are
// not the seat's (novox/hq ADR 0159, 0160): the store's `databases` is not postgres's
// `postgres_list_databases`, and a holder may well serve both. The runtime serves an
// implementation registered under the seat's name on the seat's subjects. Absent, the
// module's own `tools` must list every verb the seat promises, which is how a module named
// like its seat — the catalogue, the records — says they are one and the same.
Serves []string `json:"serves,omitempty"`
}
// ServesFor is what this claim offers a seat's protocol: the verbs it names, else the module's
// own tools.
func (c Claim) ServesFor(m Manifest) []string {
if len(c.Serves) > 0 {
return c.Serves
}
return m.Tools
}
// At is this claim's scope, with the default applied.
@@ -293,6 +309,12 @@ type Manifest struct {
// module claiming a seat answers what that seat's protocol promises (novox/hq ADR 0118).
Tools []string `json:"tools,omitempty"`
// Instances says whether this module's instances are the same anywhere — `interchangeable` —
// so a call that names no machine may be answered by any of them (novox/hq ADR 0160). A fact
// about the software, not about the bus: a stateless web tool says it; a database does not,
// and its instances are then each addressed by machine, never confused for one another.
Instances string `json:"instances,omitempty"`
// Invokes are the tools this module calls, each `<module>.<tool>` or a role's `seat:<seat>.<verb>`,
// or the single entry `*` for every tool on the mesh (novox/hq ADR 0152, ADR 0154).
//
@@ -1173,6 +1195,11 @@ func ParseManifest(raw []byte) (Manifest, error) {
m.Module, r))
}
}
if m.Instances != "" && m.Instances != InstancesInterchangeable {
problems = append(problems, fmt.Sprintf(
"%s says its instances are %q; the one word is %q, for a module that is the same on every machine",
m.Module, m.Instances, InstancesInterchangeable))
}
for _, offer := range m.Provides {
p := offer.Name
if !name.MatchString(p) {
@@ -1960,3 +1987,7 @@ func (o OwnSecrets) Paths() map[string]string {
}
return out
}
// 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.
const InstancesInterchangeable = "interchangeable"
+115
View File
@@ -0,0 +1,115 @@
package catalogue
import (
"strings"
"testing"
)
// Two modules on one node both provide acme-ca — public-acme (Let's Encrypt) and step-ca (the
// mesh's own authority) on novox — and a route-proxy elsewhere must get the public one (novox/hq
// #258). A pin names the module as well as the node, so that it can say which.
func issuerShelf() map[string]Manifest {
return shelf(
Manifest{Module: "public-acme", Version: "1", Provides: FromAnywhere("acme-ca"),
Serves: map[string]map[string]any{"acme-ca": {"at": "acme-v02.api.letsencrypt.org"}}},
Manifest{Module: "step-ca", Version: "1", Provides: FromAnywhere("acme-ca"),
Serves: map[string]map[string]any{"acme-ca": {"at": "novox.internal"}}},
Manifest{Module: "route-proxy", Version: "1", Requires: []string{"acme-ca"}},
)
}
func twoIssuersOnOneNode() map[string][]Provider {
return map[string][]Provider{"acme-ca": {
{Node: "novox", At: "novox.internal", Module: "public-acme", Serves: map[string]any{"at": "acme-v02.api.letsencrypt.org"}},
{Node: "novox", At: "novox.internal", Module: "step-ca", Serves: map[string]any{"at": "novox.internal"}},
}}
}
func TestTwoProvidersOnOneNodeAreRefusedWithBothNamed(t *testing.T) {
// The refusal must name the pair, because a node alone cannot tell them apart.
_, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode()})
if err == nil {
t.Fatal("two providers on one node were resolved by picking")
}
for _, want := range []string{"novox/public-acme", "novox/step-ca", "<node> <module>"} {
if !strings.Contains(err.Error(), want) {
t.Fatalf("the refusal does not say %q: %v", want, err)
}
}
}
func TestAPinNamesTheModule(t *testing.T) {
got, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode(),
Pinned: map[string]Chosen{"acme-ca": {Node: "novox", Module: "public-acme"}}})
if err != nil {
t.Fatal(err)
}
if len(got.Needs) != 1 || got.Needs[0].From != "novox" || got.Needs[0].Serves["at"] != "acme-v02.api.letsencrypt.org" {
t.Fatalf("the named module was not the one taken: %+v", got.Needs)
}
}
func TestARecordNamingOnlyTheNodeIsRefusedWhenThatNodeAnswersTwice(t *testing.T) {
// A pin from before the module was asked for. It once took the last one listed — a coin flip.
_, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode(), Pinned: map[string]Chosen{"acme-ca": {Node: "novox"}}})
if err == nil {
t.Fatal("a node that answers twice was resolved by picking")
}
if !strings.Contains(err.Error(), "provides it 2 times") || !strings.Contains(err.Error(), "pin workstation acme-ca novox <module>") {
t.Fatalf("the refusal does not ask for the module: %v", err)
}
}
func TestAPinNamingAModuleThatDoesNotProvideItIsRefused(t *testing.T) {
_, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode(), Pinned: map[string]Chosen{"acme-ca": {Node: "novox", Module: "gitea"}}})
if err == nil || !strings.Contains(err.Error(), "novox/gitea does not provide it") {
t.Fatalf("a module that does not provide it was not refused by name: %v", err)
}
}
func TestTwoProvidersBesideTheConsumerAreRefusedUntilOneIsNamed(t *testing.T) {
// The same ambiguity on the consumer's own machine. This was settled by a map walk — random,
// per plan — which is how novox's own proxy got its issuer.
_, err := Resolve(issuerShelf(), []string{"route-proxy", "public-acme", "step-ca"}, reachable(), World{})
if err == nil {
t.Fatal("two providers beside the consumer were resolved by picking")
}
if !strings.Contains(err.Error(), "workstation provides \"acme-ca\" 2 times") || !strings.Contains(err.Error(), "public-acme, step-ca") {
t.Fatalf("the refusal does not list them: %v", err)
}
}
func TestAPinSettlesTwoProvidersBesideTheConsumer(t *testing.T) {
got, err := Resolve(issuerShelf(), []string{"route-proxy", "public-acme", "step-ca"}, reachable(),
World{Pinned: map[string]Chosen{"acme-ca": {Node: "workstation", Module: "public-acme"}}})
if err != nil {
t.Fatal(err)
}
var found bool
for _, n := range got.Needs {
if n.Name == "acme-ca" {
found = true
if n.Serves["at"] != "acme-v02.api.letsencrypt.org" {
t.Fatalf("the named module was not the one taken: %+v", n)
}
}
}
if !found {
t.Fatalf("no need for acme-ca was created: %+v", got.Needs)
}
}
func TestTheFirstPassDoesNotRefuseTwoProvidersBesideTheConsumer(t *testing.T) {
// The first pass answers only what a node offers. Refused there, the node vanishes from every
// other node's world — and the whole mesh loses its vault for an ambiguity one machine has to
// settle. The second pass is where it is refused, and the test above proves it is.
if _, err := Resolve(issuerShelf(), []string{"route-proxy", "public-acme", "step-ca"}, reachable(),
World{Unchecked: true}); err != nil {
t.Fatalf("the first pass refused what only the second may: %v", err)
}
}
+80 -53
View File
@@ -48,11 +48,12 @@ type World struct {
Holdings []Held
// Offered is what other nodes provide at mesh scope, and everything needed to use it.
Offered map[string][]Provider
// Pinned is which node this machine was told to get a provision from, by name. Only consulted
// when more than one node could answer -- a choice recorded before it was needed should not
// start meaning something the day a second provider appears, and one recorded and then made
// unnecessary should not quietly stop applying either.
Pinned map[string]string
// Pinned is which provider this machine was told to get a provision from, by name: a module and
// the node it runs on, both (novox/hq #258). Only consulted when more than one could answer -- a
// choice recorded before it was needed should not start meaning something the day a second
// provider appears, and one recorded and then made unnecessary should not quietly stop applying
// either.
Pinned map[string]Chosen
// Licences is every provision answered by a **record rather than a node**, by provision name.
//
// novox/hq ADR 0024: a hosted model is on nobody's machine and is reached over the public
@@ -260,6 +261,10 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
// still "choose one" after somebody has chosen one. That makes the remedy useless, and it is
// how this read when first used.
satisfied := map[string]bool{}
// Which modules a person assigned here, hostable. The walk marks a module chosen only when it
// reaches it, and a consumer may be reached before the provider beside it — so the provider of
// something already satisfied is looked for among these as well as among the chosen.
assignedHere := map[string]bool{}
// Everything a person assigned goes in first, except what this machine cannot run. Those are
// choices already made, and a requirement one of them answers is not a choice to put back to
@@ -284,6 +289,7 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
for _, o := range m.Offers() {
satisfied[o] = true
}
assignedHere[a] = true
}
because[a] = "assigned"
queue = append(queue, a)
@@ -311,6 +317,51 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
// commonest arrangement of all — a service and its database on one node — the weakest
// handling, silently.
if satisfied[want] && !isModule(catalogue, want) {
here := func(name string) bool { return chosen[name] || assignedHere[name] }
local := providersHere(catalogue, here, want)
// Which of them it matters to choose between. A plain capability — a shell, a display
// server — asks nothing of whoever answers it, and three shells beside an editor are
// not a choice to put to anybody. One that grants a credential, or serves a fact the
// consumer cannot guess, becomes a binding, and a binding is to one provider.
matter := local
if !brokered[want] {
matter = nil
for _, name := range local {
if _, ok := catalogue[name].Serves[want]; ok {
matter = append(matter, name)
}
}
}
var by Manifest
switch len(matter) {
case 0:
// Nothing to bind to; satisfied by its presence, as it was.
case 1:
by = catalogue[matter[0]]
default:
// Two modules on this machine answer it. Taking whichever a map walk met first
// was the rule until novox/hq #258 — random, per plan — and the same stance as
// across machines applies: ambiguity is refused, never resolved by picking.
//
// **Not in the first pass.** That pass exists only to answer *what does this node
// offer*, and refusing there makes the machine vanish rather than report a problem
// (the sibling case below says why): every other node then loses what this one
// provides — the vault, the identity provider — and refuses for a fault that is
// this node's to settle. The second pass refuses it properly, where it is asked.
if world.Unchecked {
continue
}
c, pinned := world.Pinned[want]
if !pinned || c.Node != node.Name || !oneOf(matter, c.Module) {
reported[want] = true
problems = append(problems, fmt.Sprintf(
"%s provides %q %d times, wanted by %s — say which with `pin %s %s %s <module>`: %s",
node.Name, want, len(matter), because[want], node.Name, want, node.Name,
strings.Join(matter, ", ")))
continue
}
by = catalogue[c.Module]
}
if brokered[want] {
// Answered here, and still a need: the provider is this node.
//
@@ -326,9 +377,9 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
}
needs = append(needs, Needed{
Name: want, From: node.Name, At: at,
Serves: servedHere(catalogue, chosen, want), For: because[want],
SharedOwn: sharedHere(catalogue, chosen, want)})
} else if served := servedHere(catalogue, chosen, want); len(served) > 0 {
Serves: servedByOne(by, want), For: because[want],
SharedOwn: sharedByOne(by, want)})
} else if served := servedByOne(by, want); len(served) > 0 {
// Answered here with no credential to mint, but the provider serves facts the
// consumer cannot guess — a port, a model name — and so still needs a binding.
// **The reachability rule does not apply**: both ends are on this same machine, so
@@ -358,11 +409,7 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
if brokered[want] {
reported[want] = true
where := world.Offered[want]
names := make([]string, 0, len(where))
for _, p := range where {
names = append(names, p.Node)
}
sort.Strings(names)
names := providerNames(where)
take := func(p Provider) {
if node.At != "" && p.At == "" || node.At == "" && p.At != "" || node.At == "" && p.At == "" {
// One of them is not on the private network, so there is no path between
@@ -396,17 +443,17 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
"nothing in this mesh provides %q, wanted by %s %s",
want, because[want], remedy))
case len(where) == 1:
if chosenNode, pinned := world.Pinned[want]; pinned && chosenNode != where[0].Node {
if c, pinned := world.Pinned[want]; pinned && !c.matches(where[0]) {
// One provider, and it is not the one this machine was told to use. Silently
// using the other would be the mesh overruling a choice somebody made.
problems = append(problems, fmt.Sprintf(
"%s was told to get %q from %s, and only %s provides it",
node.Name, want, chosenNode, where[0].Node))
node.Name, want, c, nameOf(where[0])))
break
}
take(where[0])
default:
chosenNode, pinned := world.Pinned[want]
c, pinned := world.Pinned[want]
if !pinned {
// **The seat's holder answers, when a seat delivers this** (novox/hq ADR 0110).
// Not a guess, which ADR 0009 refuses: the choice was made once, mesh-wide, by
@@ -418,27 +465,32 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
break
}
problems = append(problems, fmt.Sprintf(
"%d nodes provide %q, wanted by %s — say which with `pin %s %s <node>`: %s",
"%d providers of %q, wanted by %s — say which with `pin %s %s <node> <module>`: %s",
len(where), want, because[want], node.Name, want,
strings.Join(names, ", ")))
break
}
var chosen *Provider
for i, w := range where {
if w.Node == chosenNode {
chosen = &where[i]
}
}
if chosen == nil {
// Pointed at a machine that does not answer this. Refused rather than
matching := c.among(where)
switch len(matching) {
case 0:
// Pointed at a provider that does not answer this. Refused rather than
// falling back to another: a fallback would quietly move somebody's data to
// a machine they did not choose, which is the whole reason this is asked.
problems = append(problems, fmt.Sprintf(
"%s was told to get %q from %s, and %s does not provide it — these do: %s",
node.Name, want, chosenNode, chosenNode, strings.Join(names, ", ")))
break
node.Name, want, c, c, strings.Join(names, ", ")))
case 1:
take(matching[0])
default:
// A record naming only the node, from before a pin named the module, and that
// node answers twice. This once took the last one listed (novox/hq #258): a
// coin flip, handed to whoever reads the certificate it chose.
problems = append(problems, fmt.Sprintf(
"%s was told to get %q from %s, and %s provides it %d times — say which with "+
"`pin %s %s %s <module>`: %s",
node.Name, want, c, c.Node, len(matching), node.Name, want, c.Node,
strings.Join(providerNames(matching), ", ")))
}
take(*chosen)
}
continue
}
@@ -597,31 +649,6 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
// need that is never created is a binding the consumer never gets. It is right about that from the
// manifest alone, which is why walking the catalogue mid-resolution is enough here and is not
// enough for the values.
// sharedHere is the own secret the provider of a provision on this same machine names as its
// credential (ADR 0158), or "" when the provider gives each consumer its own.
func sharedHere(catalogue map[string]Manifest, chosen map[string]bool, want string) string {
for name, m := range catalogue {
if !chosen[name] {
continue
}
if own, shared := m.SharedCredentialOf(want); shared {
return own
}
}
return ""
}
func servedHere(catalogue map[string]Manifest, chosen map[string]bool, want string) map[string]any {
for name, m := range catalogue {
if !chosen[name] {
continue
}
if _, ok := m.Serves[want]; ok {
return ServedOn(m, want, nil)
}
}
return nil
}
// isModule reports whether a name is a module in its own right rather than only something
// modules provide.
+24 -3
View File
@@ -59,13 +59,26 @@ var defaultSeats = []Seat{
{Name: ControllerSeatName, Scope: ScopeMesh, Decision: "novox/hq ADR 0079",
Emits: []string{"applied", "refused", "built-before"},
Serves: ControllerVerbs},
{Name: "mesh-store", Scope: ScopeMesh, Delivers: "postgres-database", Decision: "novox/hq ADR 0079"},
// The store's first verbs (novox/hq ADR 0159): the smallest set that makes the store askable,
// served by whichever module holds the seat with tools of these names.
{Name: "mesh-store", Scope: ScopeMesh, Delivers: "postgres-database", Decision: "novox/hq ADR 0079",
Serves: []Verb{
{Name: "databases", Description: "Every database the store holds, with its on-disk size.",
Input: schema(map[string]string{}, nil)},
{Name: "query", Description: "One read-only statement against one database the store holds.",
Input: schema(map[string]string{"database": "the database to query", "sql": "the read-only statement"},
[]string{"database", "sql"})},
}},
// **Delivers the mesh's own bus, not `amqp`.** Those were the same word until
// ADR 0127 separated them: `amqp` is a backing service a module may require, and this seat is
// the mesh's own transport. ADR 0128 then made that connection something a module requires
// rather than receives ambiently — 23 of the catalogue's modules never speak, and an ambient
// connection would mint a credential for each.
{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "novox/hq ADR 0079"},
// The vault: the controller seals every minted credential with what it provides, which is the
// test for a seat of the mesh's own (novox/hq ADR 0161) — a second provider of `secret` is a
// second claimant, refused by name, rather than a candidate for a pin.
{Name: "mesh-vault", Scope: ScopeMesh, Delivers: "secret", Decision: "novox/hq ADR 0161"},
// Named for its scope since 2026-09-30 (novox/hq ADR 0156); `the-artifact-store` resolves to it as
// an alias on a mesh that predates the rename. It serves artifacts of every kind a build makes —
// images and archives, by digest — which is why the provision is the artifact store and not an
@@ -286,11 +299,19 @@ func CanHold(m Manifest, seat Seat) error {
// **Serving the seat's tools is a condition of holding it** (novox/hq ADR 0132). A holder that
// does not answer what the role promises is every caller's timeout, found at registration and
// at handover instead, naming the verbs rather than the fact that something is missing.
if missing := unservedVerbs(m.Tools, seat.Serves); len(missing) > 0 {
if missing := unservedVerbs(claimed.ServesFor(m), seat.Serves); len(missing) > 0 {
return fmt.Errorf("%s claims %s but does not serve %s, which that seat's protocol promises "+
"(novox/hq ADR 0132) — a holder lists every verb its seat declares under tools",
"(novox/hq ADR 0132) — a holder names every verb its seat declares, under the claim's "+
"serves or among its own tools",
m.Module, seat.Name, strings.Join(missing, ", "))
}
// And nothing the seat does not promise: a verb named here that the protocol lacks is served
// to nobody, which is a typo the holder would otherwise discover as a caller's timeout.
if extra := unpromised(claimed.Serves, seat.Serves); len(extra) > 0 {
return fmt.Errorf("%s claims %s and says it serves %s, which that seat's protocol does not "+
"promise — a claim's serves names the seat's verbs and nothing else",
m.Module, seat.Name, strings.Join(extra, ", "))
}
return nil
}
+4 -4
View File
@@ -208,7 +208,7 @@ func CatalogueProblems(shelf Shelf) []string {
}
// A holder that does not answer what the seat promises is a caller's timeout, found
// at assignment instead.
if missing := unserved(m, s); len(missing) > 0 {
if missing := unserved(m, c, s); len(missing) > 0 {
problems = append(problems, fmt.Sprintf(
"%s claims %s but does not serve %s, which that seat's protocol promises",
module, c.Name, strings.Join(missing, ", ")))
@@ -221,10 +221,10 @@ func CatalogueProblems(shelf Shelf) []string {
// unserved is what a seat's protocol promises and the claimant does not answer. Only the tools
// are checked: `accepts` and `emits` are wired by the runtime from the declaration, while a tool
// is code the module either has or has not written.
func unserved(m Manifest, s SeatDeclaration) []string {
// is code the module either has or has not written — under the claim's serves, or among its own.
func unserved(m Manifest, c Claim, s SeatDeclaration) []string {
has := map[string]bool{}
for _, t := range m.Tools {
for _, t := range c.ServesFor(m) {
has[t] = true
}
var missing []string
+16
View File
@@ -70,6 +70,22 @@ func TestAHolderMustServeWhatItsSeatPromises(t *testing.T) {
}
}
// The claim may say what it serves for the role instead, when the module's own tools are not the
// seat's verbs (ADR 0160).
func TestAClaimMayNameWhatItServesForTheSeat(t *testing.T) {
m := telegram()
m.Tools = []string{"telegram_send"}
m.Claims = append([]Claim(nil), m.Claims...)
for i := range m.Claims {
if m.Claims[i].Name == "telegram-sender" {
m.Claims[i].Serves = []string{"status"}
}
}
if got := problemsFor(t, Shelf{"telegram": m}); strings.Contains(got, "does not serve") {
t.Fatalf("a claim naming the seat's verb was refused: %s", got)
}
}
// A seat with no protocol is a marker: which module is this node's showcase, or its packet filter.
// Most node-scoped seats are markers, so refusing one would refuse the majority of the set.
func TestASeatWithoutAProtocolIsAMarkerNotAMistake(t *testing.T) {
+2 -2
View File
@@ -44,7 +44,7 @@ func TestTheSeatsAreAClosedSetAndEachNamesItsDecision(t *testing.T) {
delivered[s.Delivers] = s.Name
}
}
if len(Seats()) != 14 {
if len(Seats()) != 15 {
t.Errorf("the mesh defines %d seats rather than 14; the set is closed, so a change here is "+
"a decision (novox/hq ADR 0110): %s", len(Seats()), seatNames())
}
@@ -249,7 +249,7 @@ func TestAPinStillWinsOverTheSeat(t *testing.T) {
// A consumer coupled to one provider's contents has said so, and the seat does not overrule it.
got, err := Resolve(registryShelf(), []string{"builder"}, reachable(),
World{Offered: twoRegistries(), Held: giteaHoldsTheSeat(),
Pinned: map[string]string{"npm-package-registry": "archive"}})
Pinned: map[string]Chosen{"npm-package-registry": {Node: "archive", Module: "verdaccio"}}})
if err != nil {
t.Fatal(err)
}
+25
View File
@@ -0,0 +1,25 @@
package catalogue
import "testing"
// The vault's provision is one the controller itself dereferences — every minted credential is
// sealed with it — so it is delivered by a seat of the mesh's own, and a second provider is a second
// claimant refused by name rather than a candidate for a pin (novox/hq ADR 0161, issue 106).
func TestTheVaultsSeatDeliversSecret(t *testing.T) {
seat, known := SeatNamed("mesh-vault")
if !known {
t.Fatal("mesh-vault is not in the mesh's own set")
}
if seat.Scope != ScopeMesh || seat.Delivers != "secret" {
t.Fatalf("mesh-vault is %s-scoped and delivers %q; one per mesh, delivering secret", seat.Scope, seat.Delivers)
}
vault := Manifest{Module: "mesh-vault", Provides: []Offer{{Name: "secret", Scope: ScopeMesh}},
Claims: []Claim{{Name: "mesh-vault", Scope: ScopeMesh}}}
if err := CanHold(vault, seat); err != nil {
t.Fatalf("the vault, claiming its seat and providing secret, was refused: %v", err)
}
another := Manifest{Module: "other-vault", Provides: []Offer{{Name: "secret", Scope: ScopeMesh}}}
if err := CanHold(another, seat); err == nil {
t.Fatal("a provider of secret that does not claim the seat was allowed to hold it")
}
}
+33 -4
View File
@@ -91,12 +91,28 @@ var ControllerVerbs = []Verb{
"module": "one module's name; every module when absent",
"log": "a build's id (as `builds` lists it): print what the build machine said, line by line",
}, nil)},
{Name: "plans", Description: "What the last merges produced and where each stands (novox/hq ADR 0162): " +
"the tiers, the tier a plan is at, what it waits for and since when; one plan whole, given its id.",
Input: schema(map[string]string{
"id": "a plan's id (as `plans` lists them): that plan, tier by tier",
"stop": "a plan's id: stop it — what was asked still builds, nothing further is asked",
}, nil)},
{Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it.",
Input: schema(map[string]string{"node": "the machine's name"}, []string{"node"})},
{Name: "assign", Description: "Put a module on a machine. Refused with the mesh's own words when it cannot resolve there.",
Input: schema(map[string]string{"node": "the machine's name", "module": "the module's name"}, []string{"node", "module"})},
{Name: "unassign", Description: "Take a module off a machine.",
Input: schema(map[string]string{"node": "the machine's name", "module": "the module's name"}, []string{"node", "module"})},
{Name: "pin", Description: "Tell a machine which provider answers a provision for it — the module, and the node " +
"it runs on, both. Asked for when more than one could answer; the refusal lists them.",
Input: schema(map[string]string{
"node": "the machine's name",
"provision": "the provision, as the consumer requires it",
"from": "the node the chosen provider runs on",
"module": "the module providing it there",
}, []string{"node", "provision", "from", "module"})},
{Name: "unpin", Description: "Take that choice back, putting the question to the mesh again.",
Input: schema(map[string]string{"node": "the machine's name", "provision": "the provision"}, []string{"node", "provision"})},
{Name: "push", Description: "Send a machine everything it should be — or every machine that is behind, when no machine is named.",
Input: schema(map[string]string{"node": "the machine's name; every machine behind when absent"}, nil)},
{Name: "rotate", Description: "Replace a credential. A pair credential, by provision (and a consuming machine, " +
@@ -113,12 +129,10 @@ var ControllerVerbs = []Verb{
{Name: "build", Description: "Have the build machine build a repository. Answers at once with the build's id: " +
"`builds` with that id follows it line by line, and the module is registered when the outcome comes.",
Input: schema(map[string]string{
"on": "instead of a repository: a module whose artifacts others stand on; every module built on it is rebuilt (the rebuild a changed base needs)",
"behind": "instead of a repository: \"yes\" rebuilds every module the mesh holds older than its source has",
"repository": "the repository's URL, or its path on the forge holding the git seat (owner/name)",
"path": "the module's directory inside it (optional)",
"ref": "the branch, tag or commit to build (optional)",
}, nil)},
}, []string{"repository"})},
}
// schema is a JSON schema for an object of string properties, which is every argument the verbs
@@ -135,7 +149,22 @@ func schema(properties map[string]string, required []string) map[string]any {
return out
}
// unservedVerbs is what a seat promises and a claimant's `tools` does not answer.
// unpromised is what a claim says it serves and the seat's protocol never promised.
func unpromised(serves []string, promised []Verb) []string {
has := map[string]bool{}
for _, v := range promised {
has[v.Name] = true
}
var extra []string
for _, s := range serves {
if !has[s] {
extra = append(extra, s)
}
}
return extra
}
// unservedVerbs is what a seat promises and a claimant's offer for it does not answer.
func unservedVerbs(tools []string, promised []Verb) []string {
has := map[string]bool{}
for _, t := range tools {
+5 -1
View File
@@ -49,7 +49,8 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
}
}
out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{}}
out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{},
Interchangeable: map[string]bool{}}
for _, n := range nodes {
out.Nodes = append(out.Nodes, n.Name)
modules, err := i.Assigned(ctx, n.Name)
@@ -72,6 +73,9 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
"be derived", module, n.Name)
}
out.Assigned[n.Name] = append(out.Assigned[n.Name], declaredFor(m, seats))
if m.Instances == catalogue.InstancesInterchangeable {
out.Interchangeable[m.Module] = true
}
}
}
+21 -16
View File
@@ -841,24 +841,28 @@ func (i *Inventory) SettingsFor(ctx context.Context, nodeName, module string) ([
return layers, rows.Err()
}
// PinProvision records which node a machine gets a provision from.
// PinProvision records which provider a machine gets a provision from: a module, and the node it
// runs on — both, always (novox/hq #258). A provision comes from a module, and the same module on
// two machines is two answers, so neither half alone says which.
//
// Only needed when more than one node could answer. Recordable before that, because a mesh with
// one database should not change where an existing machine gets its data the day a second
// arrives.
func (i *Inventory) PinProvision(ctx context.Context, nodeName, provision, provider string) error {
// Only needed when more than one could answer. Recordable before that, because a mesh with one
// database should not change where an existing machine gets its data the day a second arrives.
func (i *Inventory) PinProvision(ctx context.Context, nodeName, provision, providerNode, module string) error {
if strings.TrimSpace(module) == "" {
return fmt.Errorf("a pin names the module providing %q as well as the node it runs on", provision)
}
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
from, err := i.NodeByName(ctx, provider)
from, err := i.NodeByName(ctx, providerNode)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into provision_pin (node, name, provider) values ($1, $2, $3)
on conflict (node, name) do update set provider = excluded.provider, pinned_at = now()`,
node.ID, provision, from.ID)
`insert into provision_pin (node, name, provider, module) values ($1, $2, $3, $4)
on conflict (node, name) do update set provider = excluded.provider, module = excluded.module, pinned_at = now()`,
node.ID, provision, from.ID, module)
return err
}
@@ -879,27 +883,28 @@ func (i *Inventory) UnpinProvision(ctx context.Context, nodeName, provision stri
return nil
}
// PinsFor is what a node was told about where its provisions come from.
func (i *Inventory) PinsFor(ctx context.Context, nodeName string) (map[string]string, error) {
// PinsFor is what a node was told about where its provisions come from. A record from before a pin
// named the module carries the node alone; the resolver honours it while it is unambiguous.
func (i *Inventory) PinsFor(ctx context.Context, nodeName string) (map[string]catalogue.Chosen, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
`select p.name, n.name from provision_pin p join node n on n.id = p.provider
`select p.name, n.name, coalesce(p.module, '') from provision_pin p join node n on n.id = p.provider
where p.node = $1`, node.ID)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]string{}
out := map[string]catalogue.Chosen{}
for rows.Next() {
var name, provider string
if err := rows.Scan(&name, &provider); err != nil {
var name, provider, module string
if err := rows.Scan(&name, &provider, &module); err != nil {
return nil, err
}
out[name] = provider
out[name] = catalogue.Chosen{Node: provider, Module: module}
}
return out, rows.Err()
}
+4 -4
View File
@@ -385,19 +385,19 @@ func TestAPinSurvivesAndCanBeChanged(t *testing.T) {
t.Fatal(err)
}
}
if err := inv.PinProvision(ctx, "user", "postgres-database", "first"); err != nil {
if err := inv.PinProvision(ctx, "user", "postgres-database", "first", "postgres"); err != nil {
t.Fatal(err)
}
// Changing the answer replaces it rather than adding a second, or a machine would be told to
// use two databases and nothing would say which.
if err := inv.PinProvision(ctx, "user", "postgres-database", "second"); err != nil {
if err := inv.PinProvision(ctx, "user", "postgres-database", "second", "postgres"); err != nil {
t.Fatal(err)
}
pins, err := inv.PinsFor(ctx, "user")
if err != nil {
t.Fatal(err)
}
if len(pins) != 1 || pins["postgres-database"] != "second" {
if len(pins) != 1 || pins["postgres-database"].Node != "second" || pins["postgres-database"].Module != "postgres" {
t.Fatalf("got %v", pins)
}
if err := inv.UnpinProvision(ctx, "user", "postgres-database"); err != nil {
@@ -422,7 +422,7 @@ func TestAPinGoesWhenTheProviderLeavesTheMesh(t *testing.T) {
t.Fatal(err)
}
}
if err := inv.PinProvision(ctx, "consumer", "postgres-database", "provider"); err != nil {
if err := inv.PinProvision(ctx, "consumer", "postgres-database", "provider", "postgres"); err != nil {
t.Fatal(err)
}
if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'provider'`); err != nil {
+124
View File
@@ -0,0 +1,124 @@
package inventory
import (
"context"
"sort"
"strings"
"github.com/novox/mesh-controller/internal/catalogue"
)
// The kinds of edge in the catalogue's one dependency relation (novox/hq ADR 0162).
const (
// EdgeStandsOn: the module's artifact is built on the other's.
EdgeStandsOn = "stands-on"
// EdgePackages: the module's build reads the other's repository.
EdgePackages = "packages"
// EdgeBuiltBy: the module is built by the holder of the build-machine seat.
EdgeBuiltBy = "built-by"
// EdgeDeclared: the manifest's own `build.on`.
EdgeDeclared = "declared"
)
// Edge is one dependency: From depends on To, in the way Kind says.
type Edge struct {
From string `json:"from"`
To string `json:"to"`
Kind string `json:"kind"`
}
// Dependencies is the catalogue's dependency relation, whole: every module the mesh holds, with
// an edge to each module it depends on and the kind of dependency on the edge. One answer, so
// nothing else computes an edge (novox/hq ADR 0162) — the merge handler, `build --on` and the
// overview all read this.
//
// Four sources, one relation: a manifest's `build.on`; the artifacts the latest build was made
// against (an `artifact-store://<module>/…` reference is an edge to that module); the repositories
// the latest build read (an edge to the module whose source that is); and the build machine, which
// every source-built module is built by.
func (i *Inventory) Dependencies(ctx context.Context) ([]Edge, error) {
entries, err := i.Catalogued(ctx)
if err != nil {
return nil, err
}
against, err := i.BuiltAgainst(ctx)
if err != nil {
return nil, err
}
read, err := i.ReadRepositories(ctx)
if err != nil {
return nil, err
}
return dependenciesOf(entries, against, read), nil
}
// dependenciesOf is Dependencies over what was read, so a test can hand it a catalogue.
func dependenciesOf(entries []Entry, against map[string][]string, read map[string][]ReadRepository) []Edge {
known := map[string]bool{}
byRepository := map[string][]string{}
var builders []string
for _, e := range entries {
name := e.Manifest.Module
known[name] = true
if r := repositoryKey(e.Source.Repository); r != "" {
byRepository[r] = append(byRepository[r], name)
}
if e.Manifest.ClaimsSeat("mesh-build-machine") {
builders = append(builders, name)
}
}
seen := map[Edge]bool{}
var out []Edge
add := func(from, to, kind string) {
if from == to || !known[to] {
return
}
e := Edge{From: from, To: to, Kind: kind}
if !seen[e] {
seen[e] = true
out = append(out, e)
}
}
for _, e := range entries {
name := e.Manifest.Module
if e.Manifest.Build != nil {
for _, on := range e.Manifest.Build.On {
if on.Module != "" {
add(name, on.Module, EdgeDeclared)
}
}
}
for _, ref := range against[name] {
if rest, ok := strings.CutPrefix(ref, catalogue.ArtifactStoreScheme); ok {
if base, _, found := strings.Cut(rest, "/"); found {
add(name, base, EdgeStandsOn)
}
}
}
for _, r := range read[name] {
for _, other := range byRepository[repositoryKey(r.Repository)] {
add(name, other, EdgePackages)
}
}
if e.Source.Repository != "" {
for _, b := range builders {
add(name, b, EdgeBuiltBy)
}
}
}
sort.Slice(out, func(a, b int) bool {
if out[a].From != out[b].From {
return out[a].From < out[b].From
}
if out[a].To != out[b].To {
return out[a].To < out[b].To
}
return out[a].Kind < out[b].Kind
})
return out
}
// repositoryKey is a repository as compared: lower-cased, without a trailing `.git`.
func repositoryKey(repository string) string {
return strings.ToLower(strings.TrimSuffix(strings.TrimSpace(repository), ".git"))
}
+64
View File
@@ -0,0 +1,64 @@
package inventory
import (
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// The catalogue's one dependency relation (novox/hq ADR 0162): four kinds of edge from one call.
func TestDependenciesAreOneRelationWithTheirKinds(t *testing.T) {
entry := func(name, repository string) Entry {
return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}}
}
builder := entry("builder", "http://forge/novox/mesh-catalog.git")
builder.Manifest.Claims = []catalogue.Claim{{Name: "mesh-build-machine", Scope: catalogue.ScopeMesh}}
plugin := entry("shop-plugin", "http://forge/novox/mesh-catalog.git")
plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}}
entries := []Entry{
entry("mesh-tools", "http://forge/novox/mesh-tools.git"),
entry("shop", "http://forge/novox/mesh-catalog.git"),
plugin,
builder,
entry("mesh-controller", "http://forge/novox/mesh-controller.git"),
entry("route-proxy", "http://forge/novox/mesh-catalog.git"),
{Manifest: catalogue.Manifest{Module: "hand-made"}},
}
against := map[string][]string{
"shop": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
"builder": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
}
read := map[string][]ReadRepository{
"route-proxy": {{Repository: "http://forge/novox/mesh-controller.git", Ref: "main"}},
}
got := dependenciesOf(entries, against, read)
has := func(from, to, kind string) bool {
for _, e := range got {
if e == (Edge{From: from, To: to, Kind: kind}) {
return true
}
}
return false
}
for _, want := range []Edge{
{"shop", "mesh-tools", EdgeStandsOn},
{"builder", "mesh-tools", EdgeStandsOn},
{"shop-plugin", "shop", EdgeDeclared},
{"route-proxy", "mesh-controller", EdgePackages},
{"shop", "builder", EdgeBuiltBy},
{"mesh-controller", "builder", EdgeBuiltBy},
{"mesh-tools", "builder", EdgeBuiltBy},
} {
if !has(want.From, want.To, want.Kind) {
t.Errorf("missing %+v in %+v", want, got)
}
}
if has("builder", "builder", EdgeBuiltBy) {
t.Error("the builder is not built by itself")
}
for _, e := range got {
if e.From == "hand-made" {
t.Errorf("a module with no source depends on nothing: %+v", e)
}
}
}
@@ -0,0 +1,12 @@
-- A pin names the module as well as the node (novox/hq #258).
--
-- 0008 said "not a module: the same module on two machines is two answers, and which machine is the
-- whole question". Half right. Two modules on one machine can both answer a provision — public-acme
-- and step-ca both offer acme-ca on novox — and then which *module* is the whole question, and a
-- node alone cannot ask it. The resolver, given a node that answered twice, took the last one listed.
--
-- A provider is a (node, module) pair (design 23), and a pin names the pair. Nullable, so a record
-- made before this was asked keeps meaning what it meant: honoured while that node answers once,
-- refused with the module asked for when it answers twice.
alter table provision_pin add column module text;
@@ -0,0 +1,19 @@
-- The records already made are completed where the mesh can tell: a pin naming a node on which
-- exactly one assigned module offers the provision (or is the module itself, for a requirement that
-- names a module) gets that module. A node that answers twice is left to say which — the resolver
-- refuses it with the module asked for, rather than this guessing on its behalf.
update provision_pin p
set module = sub.module
from (
select p2.node, p2.name, min(a.module) as module, count(distinct a.module) as answers
from provision_pin p2
join assignment a on a.node = p2.provider
join module m on m.name = a.module
where p2.module is null
and (a.module = p2.name
or exists (select 1
from jsonb_array_elements(coalesce(m.manifest -> 'provides', '[]'::jsonb)) e
where (case when jsonb_typeof(e) = 'string' then e #>> '{}' else e ->> 'name' end) = p2.name))
group by p2.node, p2.name
) sub
where sub.node = p.node and sub.name = p.name and sub.answers = 1;
@@ -0,0 +1,19 @@
-- A merge produces a tiered plan the mesh keeps (novox/hq ADR 0162): what the merge changed and
-- everything standing on it, sorted into tiers, each module's state, and the tier the plan is at.
-- Kept so a controller replaced mid-plan resumes it, and so `status` can say what a merge still
-- waits for.
create table release_plan (
id text primary key,
repository text not null,
commit_hash text not null,
created timestamptz not null default now(),
updated timestamptz not null default now(),
-- building: a tier's builds are asked; rolling: the tier is built and the machines are applying
-- what a later tier needs running; done; failed.
state text not null,
tier int not null default 0,
tiers jsonb not null,
modules jsonb not null,
note text not null default ''
);
create index release_plan_open on release_plan (created) where state in ('building', 'rolling');
+5 -2
View File
@@ -42,8 +42,11 @@ func TestAPersonMayCallToolsAndNothingElse(t *testing.T) {
if err != nil {
t.Fatal(err)
}
if len(perms.Publish) != 1 || perms.Publish[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" {
t.Errorf("ada may publish %v, which should be the one tool and nothing else", perms.Publish)
// The one tool, both ways it is addressed (novox/hq ADR 0159): to whichever instance
// answers, and to the instance on one machine. Nothing else.
if len(perms.Publish) != 2 || perms.Publish[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" ||
perms.Publish[1] != "mesh.mod.mesh-catalog.tool.catalog_tools.*" {
t.Errorf("ada may publish %v, which should be the one tool, both ways addressed, and nothing else", perms.Publish)
}
for _, s := range perms.Publish {
if strings.HasPrefix(s, "mesh.control") || strings.HasPrefix(s, "mesh.node") ||
+86
View File
@@ -0,0 +1,86 @@
package inventory
import (
"context"
"os"
"testing"
)
// A pin made before it named the module (migration 0051, novox/hq #258): completed where the node it
// names answers once, left for a person where it answers twice.
func legacyPin(t *testing.T, inv *Inventory, node, provision, provider string) {
t.Helper()
_, err := inv.store.Pool().Exec(context.Background(),
`insert into provision_pin (node, name, provider)
select u.id, $2, p.id from node u, node p where u.name = $1 and p.name = $3`,
node, provision, provider)
if err != nil {
t.Fatal(err)
}
}
func completeEarlierPins(t *testing.T, inv *Inventory) {
t.Helper()
sql, err := os.ReadFile("migrations/0051-a-pin-made-before-is-completed.sql")
if err != nil {
t.Fatal(err)
}
if _, err := inv.store.Pool().Exec(context.Background(), string(sql)); err != nil {
t.Fatal(err)
}
}
func TestAPinMadeBeforeIsCompletedWhenTheNodeAnswersOnce(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
for _, n := range []string{"user", "provider"} {
if _, err := inv.AddNode(ctx, n); err != nil {
t.Fatal(err)
}
}
for _, m := range []string{"postgres", "redis"} {
if err := inv.RegisterModule(ctx, manifest(m, []string{m + "-database"}, nil), Source{}); err != nil {
t.Fatal(err)
}
if _, err := inv.Assign(ctx, "provider", m); err != nil {
t.Fatal(err)
}
}
legacyPin(t, inv, "user", "postgres-database", "provider")
completeEarlierPins(t, inv)
pins, err := inv.PinsFor(ctx, "user")
if err != nil {
t.Fatal(err)
}
if got := pins["postgres-database"]; got.Node != "provider" || got.Module != "postgres" {
t.Fatalf("the record was not completed with the one module that answers: %+v", got)
}
}
func TestAPinMadeBeforeIsLeftOpenWhenTheNodeAnswersTwice(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
for _, n := range []string{"user", "provider"} {
if _, err := inv.AddNode(ctx, n); err != nil {
t.Fatal(err)
}
}
for _, m := range []string{"public-acme", "step-ca"} {
if err := inv.RegisterModule(ctx, manifest(m, []string{"acme-ca"}, nil), Source{}); err != nil {
t.Fatal(err)
}
if _, err := inv.Assign(ctx, "provider", m); err != nil {
t.Fatal(err)
}
}
legacyPin(t, inv, "user", "acme-ca", "provider")
completeEarlierPins(t, inv)
pins, err := inv.PinsFor(ctx, "user")
if err != nil {
t.Fatal(err)
}
if got := pins["acme-ca"]; got.Node != "provider" || got.Module != "" {
t.Fatalf("a node that answers twice was guessed for: %+v", got)
}
}
+123
View File
@@ -0,0 +1,123 @@
package inventory
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
)
// A Plan is what a merge produces (novox/hq ADR 0162): the modules it changed and everything
// standing on them, sorted into tiers, each module's state, and the tier the plan is at. Kept in
// the store so a controller replaced mid-plan resumes it, and so `status` can say what a merge
// still waits for.
type Plan struct {
ID string `json:"id"`
Repository string `json:"repository"`
Commit string `json:"commit"`
Created time.Time `json:"created"`
Updated time.Time `json:"updated"`
State string `json:"state"`
Tier int `json:"tier"`
Tiers [][]string `json:"tiers"`
Modules map[string]*PlanModule `json:"modules"`
Note string `json:"note,omitempty"`
}
// PlanModule is one module's state within a plan.
type PlanModule struct {
// State: asked, built, failed; empty for a module whose tier has not been asked yet.
State string `json:"state,omitempty"`
AskedAt *time.Time `json:"asked_at,omitempty"`
BuiltAt *time.Time `json:"built_at,omitempty"`
Commit string `json:"commit,omitempty"`
Why string `json:"why,omitempty"`
}
// The states a plan passes through.
const (
PlanBuilding = "building"
PlanRolling = "rolling"
PlanDone = "done"
PlanFailed = "failed"
)
// Open says whether the plan is still being worked.
func (p Plan) Open() bool { return p.State == PlanBuilding || p.State == PlanRolling }
// SavePlan writes a plan, new or changed, whole: the plan is small and read as one thing.
func (i *Inventory) SavePlan(ctx context.Context, p Plan) error {
tiers, err := json.Marshal(p.Tiers)
if err != nil {
return err
}
modules, err := json.Marshal(p.Modules)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note)
values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9)
on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier,
tiers = excluded.tiers, modules = excluded.modules, note = excluded.note`,
p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note)
return err
}
// OpenPlans is every plan still being worked, oldest first.
func (i *Inventory) OpenPlans(ctx context.Context) ([]Plan, error) {
return i.plans(ctx, `where state in ('building', 'rolling') order by created`)
}
// RecentPlans is the last few plans, newest first, open or not — what the overview shows.
func (i *Inventory) RecentPlans(ctx context.Context, limit int) ([]Plan, error) {
return i.plans(ctx, fmt.Sprintf(`order by created desc limit %d`, limit))
}
// PlanByID is one plan.
func (i *Inventory) PlanByID(ctx context.Context, id string) (Plan, error) {
plans, err := i.plans(ctx, `where id = '`+id+`'`)
if err != nil {
return Plan{}, err
}
if len(plans) == 0 {
return Plan{}, fmt.Errorf("no plan %s", id)
}
return plans[0], nil
}
func (i *Inventory) plans(ctx context.Context, tail string) ([]Plan, error) {
rows, err := i.store.Pool().Query(ctx,
`select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note
from release_plan `+tail)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Plan
for rows.Next() {
var p Plan
var tiers, modules []byte
if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State,
&p.Tier, &tiers, &modules, &p.Note); err != nil {
return nil, err
}
if err := json.Unmarshal(tiers, &p.Tiers); err != nil {
return nil, err
}
if err := json.Unmarshal(modules, &p.Modules); err != nil {
return nil, err
}
if p.Modules == nil {
p.Modules = map[string]*PlanModule{}
}
out = append(out, p)
}
if errors.Is(rows.Err(), pgx.ErrNoRows) {
return nil, nil
}
return out, rows.Err()
}
+48
View File
@@ -0,0 +1,48 @@
package inventory
import (
"testing"
"time"
)
// A plan is a record the mesh keeps and resumes (novox/hq ADR 0162): written whole, read back open,
// advanced, and gone from the open ones when done.
func TestAPlanIsKeptAdvancedAndResumedFromTheStore(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
p := Plan{ID: "plan-1", Repository: "novox/mesh-tools", Commit: "abc", Created: time.Now().UTC(),
State: PlanBuilding, Tiers: [][]string{{"mesh-tools"}, {"builder"}, {"shop"}},
Modules: map[string]*PlanModule{"mesh-tools": {}, "builder": {}, "shop": {}}}
if err := inv.SavePlan(ctx, p); err != nil {
t.Fatal(err)
}
open, err := inv.OpenPlans(ctx)
if err != nil || len(open) != 1 || open[0].ID != "plan-1" || len(open[0].Tiers) != 3 {
t.Fatalf("the plan was not kept whole: %v %+v", err, open)
}
// Another controller picks it up where it was left: a tier advanced and a module built.
now := time.Now().UTC()
resumed := open[0]
resumed.Tier = 1
resumed.Modules["mesh-tools"].State = "built"
resumed.Modules["mesh-tools"].BuiltAt = &now
resumed.State = PlanRolling
resumed.Note = "tier 0 built; waiting for builder on anchor to be applied"
if err := inv.SavePlan(ctx, resumed); err != nil {
t.Fatal(err)
}
again, err := inv.PlanByID(ctx, "plan-1")
if err != nil || again.Tier != 1 || again.Modules["mesh-tools"].State != "built" || again.State != PlanRolling {
t.Fatalf("the advanced plan did not come back as left: %v %+v", err, again)
}
again.State = PlanDone
if err := inv.SavePlan(ctx, again); err != nil {
t.Fatal(err)
}
if open, _ = inv.OpenPlans(ctx); len(open) != 0 {
t.Fatalf("a done plan is not open: %+v", open)
}
if recent, _ := inv.RecentPlans(ctx, 5); len(recent) != 1 || recent[0].State != PlanDone {
t.Fatalf("a done plan is still among the recent ones: %+v", recent)
}
}
+6
View File
@@ -173,7 +173,13 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
_ = msg.Term()
continue
}
// A build outlives the acknowledgement window many times over; said while it runs,
// as the controller says it for its own long handlers, so the server neither hands
// the ask to a second machine nor counts the wait against its deliveries.
working := make(chan struct{})
go stillWorking(msg, working)
do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js})
close(working)
}
}
}
+19
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"github.com/novox/mesh-controller/internal/broker"
"time"
"github.com/nats-io/nats.go"
@@ -145,6 +146,24 @@ func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body
return nil
}
// PublishMembership issues one assignment what it serves and reaches (novox/hq ADR 0160), last per
// subject, so the runtime that connects later reads the current one and one that is running follows.
// MembershipWait bounds how long issuing one membership may take. A publish the server refuses is
// never acknowledged, and a stream publish waits for its acknowledgement for as long as its
// context lives: on 2026-10-01 the daemon's own context was that long, and one refused membership
// held the controller's receive loop for good (novox/hq issue 185).
const MembershipWait = 10 * time.Second
func (b OverNATS) PublishMembership(ctx context.Context, node, module string, body []byte) error {
ctx, cancel := context.WithTimeout(ctx, MembershipWait)
defer cancel()
_, err := b.JS.Publish(broker.MembershipSubject(node, module), body, nats.Context(ctx))
if err != nil {
return fmt.Errorf("issuing %s on %s its membership: %w", module, node, err)
}
return nil
}
func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error {
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
if err != nil {
+7
View File
@@ -320,6 +320,13 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err
return false, err
}
}
if len(report.Profile) > 0 {
// The latest wins, as at enrolment: a capability the machine lost is one the plan must
// stop counting on (novox/hq ADR 0161).
if err := e.Inventory.RecordProfile(ctx, node.ID, report.Profile); err != nil {
return false, err
}
}
// What it says about the tunnel it carried (novox/hq ADR 0105), whenever it says it.
if report.Tunnel != nil {
if err := e.Inventory.RecordCarriedTunnel(ctx, node.ID, inventory.Carried{
+45
View File
@@ -0,0 +1,45 @@
package link_test
import (
"testing"
"github.com/novox/mesh-controller/internal/link"
)
// A report may carry the machine's profile, detected again by the apply that reports, and the latest
// replaces what enrolment recorded (novox/hq ADR 0161): a machine that switched its network manager
// is a machine whose uplink holder lacks a capability at its next push, not at its next enrolment.
func TestAReportsProfileReplacesTheEnrolledOne(t *testing.T) {
e, _, _ := anEnrolledHub(t)
ctx := t.Context()
first := map[string]any{"capabilities": []any{map[string]any{"name": "uplink-networkmanager", "present": true}}}
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Profile: first}); err != nil {
t.Fatal(err)
}
got, err := e.Inventory.Profile(ctx, "anchor")
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Name != "uplink-networkmanager" || !got[0].Present {
t.Fatalf("the report's profile was not kept: %+v", got)
}
// The machine switched managers; the next report says so and the old fact is gone.
second := map[string]any{"capabilities": []any{map[string]any{"name": "uplink-systemd-networkd", "present": true}}}
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Profile: second}); err != nil {
t.Fatal(err)
}
got, err = e.Inventory.Profile(ctx, "anchor")
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Name != "uplink-systemd-networkd" {
t.Fatalf("the latest profile did not replace the earlier one: %+v", got)
}
// A report with no profile leaves the last one standing.
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Host: "1"}); err != nil {
t.Fatal(err)
}
if got, _ = e.Inventory.Profile(ctx, "anchor"); len(got) != 1 {
t.Fatalf("a report without a profile erased it: %+v", got)
}
}
+6
View File
@@ -193,6 +193,12 @@ type Report struct {
// refuses it whole — which is right, and makes every new field a flag day that the mesh could
// not see coming.
Host string `json:"host,omitempty"`
// Profile is what the machine can do, detected again by this apply (novox/hq ADR 0161): the
// same shape enrolment sends, so a machine that gained or lost a capability — switched its
// network manager — is known at its next push and not at its next enrolment. Absent from a host
// older than this, and then the enrolment's profile stands.
Profile map[string]any `json:"profile,omitempty"`
// Reachable is what can be reached on the machine now: every listening socket and every
// published container port. Only an adopted node reports it; it is what converging previews.
Reachable []Reach `json:"reachable,omitempty"`