Compare commits
23
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bdcbda801e | ||
|
|
5a7ed56f61 | ||
|
|
d54ecb3bf2 | ||
|
|
cf84117638 | ||
|
|
8d4e940866 | ||
|
|
11b20b10ff | ||
|
|
7e701c0db2 | ||
|
|
853c63b181 | ||
|
|
84ac840ff4 | ||
|
|
e8e502343f | ||
|
|
1edc44b25f | ||
|
|
5a1b37e477 | ||
|
|
4a6a4eadeb | ||
|
|
69f559ab4c | ||
|
|
05b90f966a | ||
|
|
184913b620 | ||
|
|
de7aed5016 | ||
|
|
18958154f0 | ||
|
|
a2b1f9e936 | ||
|
|
c9a0b1f9f4 | ||
|
|
f450303e8e | ||
|
|
603ad61142 | ||
|
|
e7cff3d38e |
@@ -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))
|
||||
}
|
||||
}
|
||||
@@ -203,7 +203,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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
@@ -383,6 +384,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 +700,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
|
||||
}
|
||||
|
||||
|
||||
@@ -79,6 +79,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
|
||||
@@ -176,7 +186,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)
|
||||
}
|
||||
@@ -215,3 +225,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
|
||||
}
|
||||
|
||||
@@ -299,6 +299,20 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
|
||||
if err != nil {
|
||||
return notNow(err)
|
||||
}
|
||||
// And whatever stands on what moved. A base rebuilt without its dependents is a mesh half on
|
||||
// the old image until somebody remembers to ask — on 2026-10-01 forty-two modules, twice, by
|
||||
// hand (novox/hq issue 186). The relation is the one `build --on` reads; this is the same
|
||||
// rebuild, asked by the merge that made it necessary, in base order.
|
||||
standing := dependentsOf(moved, entries, against)
|
||||
if len(standing) > 0 {
|
||||
var on []string
|
||||
for _, e := range standing {
|
||||
on = append(on, e.Manifest.Module)
|
||||
}
|
||||
fmt.Printf(" %d module(s) stand on what moved and are rebuilt with it: %s\n",
|
||||
len(standing), strings.Join(on, ", "))
|
||||
moved = append(moved, standing...)
|
||||
}
|
||||
ordered := orderByBases(moved, against)
|
||||
names := make([]string, 0, len(ordered))
|
||||
for _, e := range ordered {
|
||||
@@ -522,3 +536,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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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.>")
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
@@ -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" }
|
||||
} }
|
||||
]
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -97,6 +97,16 @@ var ControllerVerbs = []Verb{
|
||||
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, " +
|
||||
@@ -133,7 +143,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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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;
|
||||
@@ -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") ||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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"`
|
||||
|
||||
Reference in New Issue
Block a user