diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 7675985..d61e53d 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -7,6 +7,7 @@ import ( "errors" "flag" "fmt" + "log" "os" "sort" "strings" @@ -124,6 +125,22 @@ func serve(ctx context.Context) error { return err } + // And the mesh's own verbs, as the seat this control plane holds (novox/hq ADR 0154). Served + // from the store's row, so what the seat declares is what is answered. + handlers, err := seatToolHandlers() + if err != nil { + return err + } + bus, isNATS := server.Bus().(link.OverNATS) + if !isNATS { + return errors.New("the mesh's verbs are served over the bus, and this control plane is not on it") + } + stopServing, err := bus.ServeSeatTools(catalogue.ControllerSeatName, handlers, log.New(os.Stdout, "", log.LstdFlags)) + if err != nil { + return err + } + defer stopServing() + return server.Serve(ctx) } diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go new file mode 100644 index 0000000..09170aa --- /dev/null +++ b/cmd/mesh-controller/seatverbs.go @@ -0,0 +1,190 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "os" + "os/exec" + "strings" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/link" +) + +// The mesh's own verbs, served as the mesh-controller seat's tools (novox/hq ADR 0154, design 33). +// +// **Each tool runs the command it names, in this same binary, and answers what it printed.** That is +// ADR 0035 taken literally: the logic lives once, in the command, and a surface is an adapter with no +// decisions in it. Running a fresh process rather than calling the function keeps two things true +// that calling it would not — every command opens and closes its own stores the way it does from a +// shell, and nothing a command prints to the process's standard output can leak into another call's +// answer. It also means a refusal is the same refusal in the same words, because it is the same +// output. + +// verbAnswer is what a verb answers: what the command printed, whether it succeeded, and — where the +// command speaks JSON — the same as data. +type verbAnswer struct { + Output string `json:"output"` + OK bool `json:"ok"` + Answer any `json:"answer,omitempty"` +} + +// argvFor is the command line a verb and its arguments become. Only the verbs the seat declares, and +// only the arguments each declares: a caller cannot reach a flag the schema did not name. +func argvFor(verb string, args map[string]any) ([]string, error) { + str := func(key string) string { + v, _ := args[key].(string) + return strings.TrimSpace(v) + } + need := func(keys ...string) error { + for _, k := range keys { + if str(k) == "" { + return fmt.Errorf("%s needs %q", verb, k) + } + } + return nil + } + switch verb { + case "status": + return []string{"status", "--json"}, nil + case "nodes": + return []string{"node", "list"}, nil + case "node": + if err := need("node"); err != nil { + return nil, err + } + return []string{"node", "show", str("node")}, nil + case "modules": + return []string{"module", "list"}, nil + case "seats": + return []string{"seats", "--json"}, nil + case "builds": + if m := str("module"); m != "" { + return []string{"builds", m}, nil + } + return []string{"builds"}, nil + case "plan": + if err := need("node"); err != nil { + return nil, err + } + return []string{"plan", str("node"), "--json"}, nil + case "assign", "unassign": + if err := need("node", "module"); err != nil { + return nil, err + } + return []string{verb, str("node"), str("module")}, 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 + // time out on every machine that takes a minute, and say nothing about the ones that did not. + if n := str("node"); n != "" { + return []string{"push", n, "--wait", "0"}, nil + } + return []string{"push", "--behind", "--wait", "0"}, nil + case "build": + if err := need("repository"); err != nil { + return nil, err + } + argv := []string{"build", str("repository"), "--wait", "0"} + if p := str("path"); p != "" { + argv = append(argv, "--path", p) + } + if r := str("ref"); r != "" { + argv = append(argv, "--ref", r) + } + return argv, nil + } + return nil, fmt.Errorf("%q is not a verb the %s seat serves", verb, catalogue.ControllerSeatName) +} + +// jsonVerbs are the verbs whose command speaks JSON, so the answer carries it as data as well. +var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true} + +// runVerb runs this binary with the given command line and gathers what it said. +func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { + self, err := os.Executable() + if err != nil { + return verbAnswer{}, err + } + cmd := exec.CommandContext(ctx, self, argv...) + // The same environment: the stores' credentials, the bus, the broker — everything a command run + // from a shell in this container would have, because it is that. + cmd.Env = os.Environ() + var out bytes.Buffer + cmd.Stdout = &out + cmd.Stderr = &out + runErr := cmd.Run() + answer := verbAnswer{Output: out.String(), OK: runErr == nil} + if jsonVerbs[argv[0]] && runErr == nil { + var parsed any + if json.Unmarshal(bytes.TrimSpace(out.Bytes()), &parsed) == nil { + answer.Answer = parsed + } + } + var exit *exec.ExitError + if runErr != nil && !errors.As(runErr, &exit) { + // Not the command refusing — the command not running at all, which is this process's fault. + return answer, fmt.Errorf("could not run %s: %w", strings.Join(argv, " "), runErr) + } + return answer, nil +} + +// seatToolHandlers are the handlers for every verb the mesh-controller seat declares, from the +// store's row, so a verb the row does not carry is not served and a verb it carries that this binary +// cannot run is said at start rather than at the first call. +func seatToolHandlers() (map[string]link.ToolHandler, error) { + seat, known := catalogue.SeatNamed(catalogue.ControllerSeatName) + if !known { + return nil, fmt.Errorf("this mesh defines no %s seat", catalogue.ControllerSeatName) + } + handlers := map[string]link.ToolHandler{} + for _, v := range seat.Serves { + verb := v.Name + if verb == "tools" { + handlers[verb] = func(ctx context.Context, _ json.RawMessage) (any, error) { + return seatTools(), nil + } + continue + } + if _, err := argvFor(verb, map[string]any{"node": "x", "module": "x", "repository": "x"}); err != nil { + return nil, fmt.Errorf("the %s seat's row declares %q, which this control plane cannot run: %w", + catalogue.ControllerSeatName, verb, err) + } + handlers[verb] = func(ctx context.Context, raw json.RawMessage) (any, error) { + args := map[string]any{} + if len(raw) > 0 { + if err := json.Unmarshal(raw, &args); err != nil { + return nil, fmt.Errorf("the arguments are not a JSON object: %w", err) + } + } + argv, err := argvFor(verb, args) + if err != nil { + return nil, err + } + return runVerb(ctx, argv) + } + } + return handlers, nil +} + +// seatTools is what `tools` answers: every seat with a protocol, and the tools each serves, from the +// mesh's own records — no holder in the path, so it is true while a holder restarts (design 33 §5). +func seatTools() map[string]any { + var seats []map[string]any + for _, s := range catalogue.SeatsWithAProtocol() { + if len(s.Serves) == 0 { + continue + } + var tools []map[string]any + for _, v := range s.Serves { + tools = append(tools, map[string]any{ + "name": v.Name, "description": v.Description, "input": v.Input, "output": v.Output, + }) + } + seats = append(seats, map[string]any{"seat": s.Name, "scope": s.Scope, "tools": tools}) + } + return map[string]any{"seats": seats} +} diff --git a/cmd/mesh-controller/seatverbs_test.go b/cmd/mesh-controller/seatverbs_test.go new file mode 100644 index 0000000..791ab9c --- /dev/null +++ b/cmd/mesh-controller/seatverbs_test.go @@ -0,0 +1,79 @@ +package main + +import ( + "strings" + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// Every verb the mesh-controller seat declares is one this binary can run, with the arguments the +// schema names and no other (novox/hq ADR 0154, ADR 0035). +func TestEveryDeclaredVerbHasACommandLine(t *testing.T) { + for _, v := range catalogue.ControllerVerbs { + if v.Name == "tools" { + continue + } + args := map[string]any{} + props, _ := v.Input["properties"].(map[string]any) + for name := range props { + args[name] = "x" + } + argv, err := argvFor(v.Name, args) + if err != nil { + t.Errorf("%s: %v", v.Name, err) + continue + } + if argv[0] == "" { + t.Errorf("%s: empty command", v.Name) + } + } +} + +// A required argument missing is refused in the verb's own words, before anything runs. +func TestAVerbMissingWhatItNeedsIsRefused(t *testing.T) { + if _, err := argvFor("node", map[string]any{}); err == nil || !strings.Contains(err.Error(), `node needs "node"`) { + t.Fatalf("node without a machine was accepted: %v", err) + } + if _, err := argvFor("upgrade", map[string]any{}); err == nil { + t.Fatal("a verb the seat does not serve was accepted") + } +} + +// A push and a build are sent, not waited for: the asker reads status for what happened. +func TestActsDoNotBlockTheCall(t *testing.T) { + argv, _ := argvFor("push", map[string]any{"node": "one"}) + if strings.Join(argv, " ") != "push one --wait 0" { + t.Fatalf("push waits: %v", argv) + } + argv, _ = argvFor("build", map[string]any{"repository": "novox/x", "path": "modules/x"}) + if strings.Join(argv, " ") != "build novox/x --wait 0 --path modules/x" { + t.Fatalf("build: %v", argv) + } +} + +// What `tools` answers is the seats' records, with each verb's schema. +func TestToolsAnswersTheSeatsRecords(t *testing.T) { + handlers, err := seatToolHandlers() + if err != nil { + t.Fatal(err) + } + if len(handlers) != len(catalogue.ControllerVerbs) { + t.Fatalf("%d handlers for %d verbs", len(handlers), len(catalogue.ControllerVerbs)) + } + answer := seatTools() + seats, _ := answer["seats"].([]map[string]any) + var found bool + for _, s := range seats { + if s["seat"] == catalogue.ControllerSeatName { + found = true + tools, _ := s["tools"].([]map[string]any) + if len(tools) != len(catalogue.ControllerVerbs) || tools[0]["input"] == nil { + t.Fatalf("the controller seat's tools are not listed in full: %v", tools) + } + } + } + if !found { + t.Fatal("the mesh-controller seat is not in the listing") + } +} diff --git a/internal/broker/invokes_test.go b/internal/broker/invokes_test.go index 58e227b..c2b76fd 100644 --- a/internal/broker/invokes_test.go +++ b/internal/broker/invokes_test.go @@ -42,7 +42,8 @@ func TestInvokingGrantsNothingButTheCall(t *testing.T) { if strings.Contains(p, ".event.") { t.Errorf("a module that only invokes may publish %q, an event it never declared", p) } - if strings.HasPrefix(p, "mesh.seat.") { + // A role's tools are tools (ADR 0132); a role's work queue and events are not. + if strings.HasPrefix(p, "mesh.seat.") && !strings.Contains(p, ".tool.") { t.Errorf("a module that only invokes may publish %q, a seat it neither holds nor uses", p) } } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 58a5e0b..e4b8d09 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -39,7 +39,11 @@ const ( // Seat is a role on the bus as a principal relates to it: the subjects it accepts, and those it // emits (novox/hq ADR 0118, design 29 §5). type Seat struct { - Name string + Name string + // Scope is where the seat has one holder. A node-scoped seat's tool carries the node in its + // subject, because one subject reaching six machines' holders is not an address + // (novox/hq ADR 0132, design 33 §4). Empty reads as mesh. + Scope string Accepts []string Emits []string Serves []string @@ -197,6 +201,13 @@ func PermissionsFor(p Principal) (Permissions, error) { // the new bus was refused the publish (2026-09-28). pub = append(pub, "mesh.mod.*.tool.>") + // **And the mesh's own verbs, as the seat it holds** (novox/hq ADR 0132, ADR 0154): + // `status`, `push`, `assign` are the mesh-controller seat's tools, served by its holder. The + // whole verb namespace of its own seat rather than a list: the list is the seat's protocol, + // which this package mirrors rather than reads, and a verb the seat does not declare is a + // subject nothing publishes. + sub = append(sub, "mesh.seat."+ControllerSeat+".tool.>") + // The two events it reacts to, and its ack subject on the stream they arrive from // (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the // controller a subscriber to every event in the mesh, and its permission list would stop @@ -343,7 +354,7 @@ func PermissionsFor(p Principal) (Permissions, error) { pub = append(pub, seatSubject(s, "event", e)) } for _, t := range s.Serves { - sub = append(sub, seatSubject(s, "tool", t)) + sub = append(sub, seatToolSubject(s, t, p.Node)) } } @@ -355,7 +366,7 @@ func PermissionsFor(p Principal) (Permissions, error) { pub = append(pub, seatSubject(s, "accept", a)) } for _, t := range s.Serves { - pub = append(pub, seatSubject(s, "tool", t)) + pub = append(pub, seatToolSubject(s, t, "*")) } } } @@ -405,6 +416,18 @@ func seatSubject(s Seat, kind, verb string) string { return "mesh.seat." + s.Name + "." + kind + "." + verb } +// seatToolSubject is where a role's tool is asked. Mesh-wide for a mesh-scoped seat; a node-scoped +// seat carries the node it is asked of, because a flat subject would reach every machine's holder +// and the queue group would silently pick a winner (novox/hq ADR 0132, design 33 §4). A holder +// subscribes its own node's; a user publishes any node's (`*`) and names the machine in the subject. +func seatToolSubject(s Seat, verb, node string) string { + base := seatSubject(s, "tool", verb) + if s.Scope == "node" && node != "" { + return base + "." + node + } + return base +} + // consumerStream and consumerDurable are the two halves of a consumer's identity, and they are // two functions because conflating them was a real bug. // @@ -615,13 +638,27 @@ func invokedSubjects(invokes []string) ([]string, error) { var out []string for _, t := range invokes { if t == "*" { - out = append(out, "mesh.mod.*.tool.>") + // Every module's tools and every role's (novox/hq ADR 0132): a role's verb is a tool + // like any other, addressed to the seat instead of a module. + out = append(out, "mesh.mod.*.tool.>", "mesh.seat.*.tool.>") + continue + } + if rest, isSeat := strings.CutPrefix(t, "seat:"); isSeat { + // A role's tool, `seat:.`. Both address shapes, because the grant is + // written without knowing the seat's scope: a mesh seat's verb is flat and a node + // seat's carries the machine (design 33 §4). + seat, verb, ok := strings.Cut(rest, ".") + if !ok || seat == "" || verb == "" { + return nil, fmt.Errorf( + "%q does not name a role's tool: one invokes seat:.", t) + } + out = append(out, "mesh.seat."+seat+".tool."+verb, "mesh.seat."+seat+".tool."+verb+".*") continue } module, tool, ok := strings.Cut(t, ".") if !ok || module == "" || tool == "" { return nil, fmt.Errorf( - "%q does not name a tool: one invokes ., or * for every one", t) + "%q does not name a tool: one invokes ., seat:., or * for every one", t) } out = append(out, "mesh.mod."+module+".tool."+tool) } diff --git a/internal/broker/seattools_test.go b/internal/broker/seattools_test.go new file mode 100644 index 0000000..3daff1a --- /dev/null +++ b/internal/broker/seattools_test.go @@ -0,0 +1,57 @@ +package broker + +import "testing" + +// A node-scoped seat's tool carries the node (novox/hq ADR 0132, design 33 §4): two nodes holding one +// node-scoped seat derive two addresses, and a user of the seat may publish any node's. +func TestTwoNodesHoldingOneNodeSeatDeriveTwoToolAddresses(t *testing.T) { + seat := Seat{Name: "node-dns-resolver", Scope: "node", Serves: []string{"lookup"}} + one, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "dnsmasq", Holds: []Seat{seat}, PasswordHash: "x"}) + two, _ := PermissionsFor(Principal{Kind: KindModule, Node: "two", Module: "dnsmasq", Holds: []Seat{seat}, PasswordHash: "x"}) + has(t, one.Subscribe, "mesh.seat.node-dns-resolver.tool.lookup.one") + has(t, two.Subscribe, "mesh.seat.node-dns-resolver.tool.lookup.two") + hasNot(t, one.Subscribe, "mesh.seat.node-dns-resolver.tool.lookup") + hasNot(t, one.Subscribe, "mesh.seat.node-dns-resolver.tool.lookup.two") + + user, _ := PermissionsFor(Principal{Kind: KindModule, Node: "three", Module: "asker", Uses: []Seat{seat}, PasswordHash: "x"}) + has(t, user.Publish, "mesh.seat.node-dns-resolver.tool.lookup.*") +} + +// A mesh-scoped seat's tool stays flat: nothing about it changes. +func TestAMeshSeatsToolIsAddressedToTheSeatAlone(t *testing.T) { + seat := Seat{Name: "git", Scope: "mesh", Serves: []string{"list_repos"}} + holder, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "gitea", Holds: []Seat{seat}, PasswordHash: "x"}) + has(t, holder.Subscribe, "mesh.seat.git.tool.list_repos") + user, _ := PermissionsFor(Principal{Kind: KindModule, Node: "two", Module: "asker", Uses: []Seat{seat}, PasswordHash: "x"}) + has(t, user.Publish, "mesh.seat.git.tool.list_repos") +} + +// The controller serves its own seat's verbs and may answer them (novox/hq ADR 0154). +func TestTheControllerServesItsSeatsToolsAndMayAnswer(t *testing.T) { + perms, err := PermissionsFor(Principal{Kind: KindController, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + has(t, perms.Subscribe, "mesh.seat.mesh-controller.tool.>") + if !perms.AllowResponses { + t.Fatal("the controller serves tools and may not answer one") + } +} + +// A grant to every tool reaches a role's tools too, and a role's tool is granted by name. +func TestAGrantReachesARolesTools(t *testing.T) { + all, _ := PermissionsFor(Principal{Kind: KindModule, Node: "desk", Module: "mesh-console", Invokes: []string{"*"}, PasswordHash: "x"}) + has(t, all.Publish, "mesh.seat.*.tool.>") + + one, err := PermissionsFor(Principal{Kind: KindPerson, Module: "jo", Invokes: []string{"seat:mesh-controller.status"}, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + has(t, one.Publish, "mesh.seat.mesh-controller.tool.status") + hasNot(t, one.Publish, "mesh.seat.mesh-controller.tool.push") + hasNot(t, one.Publish, "mesh.mod.*.tool.>") + + if _, err := PermissionsFor(Principal{Kind: KindPerson, Module: "jo", Invokes: []string{"seat:mesh-controller"}, PasswordHash: "x"}); err == nil { + t.Fatal("a role grant naming no verb was accepted") + } +} diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index faebd3b..e55b370 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -25,7 +25,7 @@ accounts { 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"] } - 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"] } + subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/catalogue/manifest.go b/internal/catalogue/manifest.go index 1d07488..db46878 100644 --- a/internal/catalogue/manifest.go +++ b/internal/catalogue/manifest.go @@ -252,8 +252,8 @@ type Manifest struct { // module claiming a seat answers what that seat's protocol promises (novox/hq ADR 0118). Tools []string `json:"tools,omitempty"` - // Invokes are the tools this module calls, each `.`, or the single entry `*` for - // every tool on the mesh (novox/hq ADR 0152). + // Invokes are the tools this module calls, each `.` or a role's `seat:.`, + // or the single entry `*` for every tool on the mesh (novox/hq ADR 0152, ADR 0154). // // **A grant, and only a grant.** The bus lets this module publish exactly those tool subjects // and nothing beside them — no event, no subscription, no seat. A module that declares none @@ -1795,6 +1795,7 @@ func invokeProblems(m Manifest) []string { if t == "*" { continue } + t = strings.TrimPrefix(t, "seat:") module, tool, named := strings.Cut(t, ".") if !named || !name.MatchString(module) || !toolName.MatchString(tool) { problems = append(problems, fmt.Sprintf( diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index b41b122..f521daa 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -37,7 +37,9 @@ type Seat struct { // today — which the bus refuses, because a namespace belongs to who it is named for. Accepts []string Emits []string - Serves []string + // Serves carries each verb in full — name, description, schema — because a role's tools are the + // mesh's to define and an agent's to call (novox/hq ADR 0132, design 33 §2). + Serves []Verb // Decision is the record that made it a seat. Decision string } @@ -51,8 +53,12 @@ var defaultSeats = []Seat{ // The control plane states what it did under the seat it holds (novox/hq ADR 0134): a role's // events belong to the role, so they keep their address while the holder is replaced. No accepts, // so no work queue is raised for it — only what its holder may say. - {Name: "mesh-controller", Scope: ScopeMesh, Decision: "novox/hq ADR 0079", - Emits: []string{"applied", "refused", "built-before"}}, + // And it serves the mesh's own verbs as the seat's tools (novox/hq ADR 0154): `status`, `push`, + // `assign` and the rest are a role's interface, not a container's, and stay addressable while + // the control plane is replaced. + {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"}, // **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 @@ -269,6 +275,14 @@ func CanHold(m Manifest, seat Seat) error { return fmt.Errorf("%s claims %s, whose holder answers for %q, and %s does not provide %q at %s scope", m.Module, seat.Name, seat.Delivers, m.Module, seat.Delivers, seat.Scope) } + // **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 { + 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", + m.Module, seat.Name, strings.Join(missing, ", ")) + } return nil } diff --git a/internal/catalogue/seats_declared.go b/internal/catalogue/seats_declared.go index bbc03eb..97a6198 100644 --- a/internal/catalogue/seats_declared.go +++ b/internal/catalogue/seats_declared.go @@ -38,8 +38,9 @@ type SeatDeclaration struct { Accepts []string `json:"accepts,omitempty"` // Emits are the verbs the holder publishes: 1:many, nobody obliged to act. Emits []string `json:"emits,omitempty"` - // Serves are the verbs the holder answers: request and reply, awaited. - Serves []string `json:"serves,omitempty"` + // Serves are the verbs the holder answers: request and reply, awaited. A bare name, or the + // verb in full with its schema (novox/hq ADR 0132). + Serves []Verb `json:"serves,omitempty"` // RetainSeconds is how long the inbound backlog survives with no holder, zero for the // mesh's default. Retention belongs to whoever owns the namespace (design 29 §3) — a seat @@ -62,7 +63,7 @@ func (s SeatDeclaration) At() string { func (s SeatDeclaration) verbs() []string { out := append([]string{}, s.Accepts...) out = append(out, s.Emits...) - return append(out, s.Serves...) + return append(out, VerbNames(s.Serves)...) } // declaredSeatProblems is what one manifest can be judged on alone. @@ -228,8 +229,8 @@ func unserved(m Manifest, s SeatDeclaration) []string { } var missing []string for _, t := range s.Serves { - if !has[t] { - missing = append(missing, t) + if !has[t.Name] { + missing = append(missing, t.Name) } } return missing diff --git a/internal/catalogue/seats_declared_test.go b/internal/catalogue/seats_declared_test.go index b9a0329..2bc6907 100644 --- a/internal/catalogue/seats_declared_test.go +++ b/internal/catalogue/seats_declared_test.go @@ -13,7 +13,7 @@ func problemsFor(t *testing.T, shelf Shelf) string { func telegram() Manifest { return Manifest{Module: "telegram", Tools: []string{"status"}, DefinesSeats: []SeatDeclaration{{ Name: "telegram-sender", Scope: ScopeMesh, - Accepts: []string{"send"}, Emits: []string{"delivered", "failed"}, Serves: []string{"status"}, + Accepts: []string{"send"}, Emits: []string{"delivered", "failed"}, Serves: []Verb{{Name: "status"}}, }}, Claims: []Claim{{Name: "telegram-sender", Scope: ScopeMesh}}} } diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go new file mode 100644 index 0000000..ef91cc2 --- /dev/null +++ b/internal/catalogue/verbs.go @@ -0,0 +1,133 @@ +package catalogue + +import ( + "bytes" + "encoding/json" + "fmt" +) + +// A Verb is one tool a role serves: its name, what it does, and the schema of its arguments and of +// its answer (novox/hq ADR 0132, design 33 §2). +// +// **A name alone is not callable by something that has never seen the mesh before**, which is the +// whole population a tool surface exists for. So a seat's protocol carries the definition, in the +// form an agent protocol already uses — a JSON schema for the input — so nothing translates between +// a seat's idea of an argument and the caller's. +// +// A manifest may still write a bare verb name (`"serves": ["price"]`); that is a Verb with only a +// name, and the module's runtime answers `tools` with the rest. The two forms read into one type so +// nothing downstream cares which was written. +type Verb struct { + Name string `json:"name"` + Description string `json:"description,omitempty"` + Input map[string]any `json:"input,omitempty"` + Output map[string]any `json:"output,omitempty"` +} + +func (v *Verb) UnmarshalJSON(raw []byte) error { + trimmed := bytes.TrimSpace(raw) + if len(trimmed) > 0 && trimmed[0] == '"' { + var name string + if err := json.Unmarshal(trimmed, &name); err != nil { + return err + } + *v = Verb{Name: name} + return nil + } + // Strictly, like the manifest around it: a misspelt key in a tool's definition would otherwise + // describe a tool nobody can call and refuse nothing. + type plain Verb + var p plain + decoder := json.NewDecoder(bytes.NewReader(trimmed)) + decoder.DisallowUnknownFields() + if err := decoder.Decode(&p); err != nil { + return fmt.Errorf("a served verb is a name or {name, description, input, output}: %w", err) + } + if p.Name == "" { + return fmt.Errorf("a served verb has no name: %s", trimmed) + } + *v = Verb(p) + return nil +} + +// VerbNames are the names alone, for the grants and the checks that care about nothing else. +func VerbNames(verbs []Verb) []string { + out := make([]string, 0, len(verbs)) + for _, v := range verbs { + out = append(out, v.Name) + } + return out +} + +// ControllerSeatName is the seat the control plane holds, whose tools are the mesh's own verbs. +const ControllerSeatName = "mesh-controller" + +// ControllerVerbs are the mesh's own verbs, as the `mesh-controller` seat's tools (novox/hq ADR 0154). +// +// **The same function the command line calls, and nothing the tool adds** (ADR 0035): each of these +// is a command the controller's binary already answers, run by the holder of the seat with the +// arguments below and answered with what the command printed. A verb here is a contract every future +// holder must implement, which is why the list is short and made of what an operator asks weekly. +// Additive within a version (design 33 §7); a verb that would break a caller takes a new version. +var ControllerVerbs = []Verb{ + {Name: "tools", Description: "Every seat's tools, from the mesh's own records: what each role " + + "answers, whether or not its holder is up. The mesh's own verbs are the mesh-controller seat's.", + Input: schema(nil, nil)}, + {Name: "status", Description: "What is wrong, what is quiet, what is out of date, and which " + + "machines are behind what the mesh would send them.", + Input: schema(nil, nil)}, + {Name: "nodes", Description: "Every machine the mesh knows, with whether it is converged or adopted.", + Input: schema(nil, nil)}, + {Name: "node", Description: "What one machine reported it can do, what it is assigned, and why.", + Input: schema(map[string]string{"node": "the machine's name"}, []string{"node"})}, + {Name: "modules", Description: "Every module the mesh holds: version, the commit it was built from, " + + "and which machines run it.", + Input: schema(nil, nil)}, + {Name: "seats", Description: "Every seat the mesh defines, what it delivers, and who holds it.", + Input: schema(nil, nil)}, + {Name: "builds", Description: "What has been built lately and what came of it, for every module or for one.", + Input: schema(map[string]string{"module": "one module's name; every module when absent"}, 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: "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: "build", Description: "Have the build machine build a repository and record what came out.", + Input: schema(map[string]string{ + "repository": "the repository's URL, or its path on the forge holding the git seat", + "path": "the module's directory inside it (optional)", + "ref": "the branch, tag or commit to build (optional)", + }, []string{"repository"})}, +} + +// schema is a JSON schema for an object of string properties, which is every argument the verbs +// above take. Kept small on purpose: a schema an agent cannot read is a tool it cannot call. +func schema(properties map[string]string, required []string) map[string]any { + props := map[string]any{} + for name, description := range properties { + props[name] = map[string]any{"type": "string", "description": description} + } + out := map[string]any{"type": "object", "properties": props} + if len(required) > 0 { + out["required"] = required + } + return out +} + +// unservedVerbs is what a seat promises and a claimant's `tools` does not answer. +func unservedVerbs(tools []string, promised []Verb) []string { + has := map[string]bool{} + for _, t := range tools { + has[t] = true + } + var missing []string + for _, v := range promised { + if !has[v.Name] { + missing = append(missing, v.Name) + } + } + return missing +} diff --git a/internal/catalogue/verbs_test.go b/internal/catalogue/verbs_test.go new file mode 100644 index 0000000..cda55ca --- /dev/null +++ b/internal/catalogue/verbs_test.go @@ -0,0 +1,72 @@ +package catalogue + +import ( + "strings" + "testing" +) + +// A served verb is written as a bare name or in full, and both read into one type (novox/hq ADR 0132). +func TestAServedVerbIsANameOrADefinition(t *testing.T) { + m, err := ParseManifest([]byte(`{"module":"till","version":"1","tools":["price","refund"],` + + `"seats":[{"name":"shop-till","serves":["price",{"name":"refund","description":"give it back",` + + `"input":{"type":"object","properties":{"order":{"type":"string"}}}}]}],` + + `"claims":[{"name":"shop-till","scope":"mesh"}]}`)) + if err != nil { + t.Fatal(err) + } + got := m.DefinesSeats[0].Serves + if len(got) != 2 || got[0].Name != "price" || got[1].Name != "refund" || got[1].Description != "give it back" { + t.Fatalf("verbs not read: %+v", got) + } + if got[1].Input["type"] != "object" { + t.Fatalf("the schema did not travel with the verb: %+v", got[1].Input) + } +} + +// A misspelt key inside a verb's definition is refused, like one anywhere else in the manifest. +func TestAVerbWithAnUnknownKeyIsRefused(t *testing.T) { + _, err := ParseManifest([]byte(`{"module":"till","version":"1",` + + `"seats":[{"name":"shop-till","serves":[{"name":"price","descripton":"typo"}]}]}`)) + if err == nil || !strings.Contains(err.Error(), "descripton") { + t.Fatalf("a verb with a misspelt key was accepted: %v", err) + } +} + +// Holding a mesh seat that serves verbs requires serving them, and the refusal names the verbs. +func TestHoldingAMeshSeatRequiresServingItsVerbs(t *testing.T) { + was := Seats() + t.Cleanup(func() { UseSeats(was) }) + UseSeats([]Seat{{Name: "mesh-controller", Scope: ScopeMesh, Decision: "test", + Serves: []Verb{{Name: "status"}, {Name: "push"}}}}) + seat, _ := SeatNamed("mesh-controller") + + partial := Manifest{Module: "a-controller", Tools: []string{"status"}, + Claims: []Claim{{Name: "mesh-controller", Scope: ScopeMesh}}} + err := CanHold(partial, seat) + if err == nil || !strings.Contains(err.Error(), "does not serve push") { + t.Fatalf("a holder missing a verb was not refused by name: %v", err) + } + whole := Manifest{Module: "a-controller", Tools: []string{"status", "push"}, + Claims: []Claim{{Name: "mesh-controller", Scope: ScopeMesh}}} + if err := CanHold(whole, seat); err != nil { + t.Fatalf("a holder serving every verb was refused: %v", err) + } +} + +// The mesh's own verbs are declared in full: an agent cannot call a name without a schema. +func TestEveryControllerVerbIsDescribedWithASchema(t *testing.T) { + seen := map[string]bool{} + for _, v := range ControllerVerbs { + if v.Description == "" || v.Input == nil || v.Input["type"] != "object" { + t.Errorf("%s: no description or no object schema", v.Name) + } + if seen[v.Name] { + t.Errorf("%s declared twice", v.Name) + } + seen[v.Name] = true + } + seat, _ := SeatNamed(ControllerSeatName) + if len(seat.Serves) != len(ControllerVerbs) { + t.Fatalf("the compiled mesh-controller seat serves %d verbs, the table has %d", len(seat.Serves), len(ControllerVerbs)) + } +} diff --git a/internal/inventory/busrecords.go b/internal/inventory/busrecords.go index 1de57e9..3efdc65 100644 --- a/internal/inventory/busrecords.go +++ b/internal/inventory/busrecords.go @@ -137,7 +137,8 @@ func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaratio } func asSeat(s catalogue.SeatDeclaration) broker.Seat { - return broker.Seat{Name: s.Name, Accepts: s.Accepts, Emits: s.Emits, Serves: s.Serves} + return broker.Seat{Name: s.Name, Scope: s.Scope, Accepts: s.Accepts, Emits: s.Emits, + Serves: catalogue.VerbNames(s.Serves)} } // MeshSeats are the mesh's own seats that carry a protocol, as the bus needs them: what to make a work @@ -146,7 +147,7 @@ func MeshSeats() []broker.DeclaredSeat { var out []broker.DeclaredSeat for _, s := range catalogue.SeatsWithAProtocol() { out = append(out, broker.DeclaredSeat{ - Name: s.Name, Accepts: s.Accepts, Emits: s.Emits, Serves: s.Serves, + Name: s.Name, Accepts: s.Accepts, Emits: s.Emits, Serves: catalogue.VerbNames(s.Serves), }) } return out diff --git a/internal/inventory/migrations/0047-a-seat-carries-its-tools.sql b/internal/inventory/migrations/0047-a-seat-carries-its-tools.sql new file mode 100644 index 0000000..b37f741 --- /dev/null +++ b/internal/inventory/migrations/0047-a-seat-carries-its-tools.sql @@ -0,0 +1,12 @@ +-- A seat's protocol lives in the store, not in the binary (novox/hq ADR 0129, ADR 0132, design 33 §2). +-- +-- ADR 0122 moved the seat set into this table with name, scope, delivers and decision, and the +-- protocol — what a role accepts, emits and serves — stayed compiled into the control plane and was +-- merged in as a row was read. Discovery that reads a binary disagrees with the mesh the moment the +-- two are on different versions, and a tool without a schema is not something an agent can call. So +-- the three halves become columns: accepts and emits as lists of verbs, serves as the verbs in full +-- ({name, description, input, output}). Seeded from the compiled defaults where a row has none, +-- additively thereafter (a verb a release adds joins the row; nothing is taken away). +alter table seat add column accepts jsonb not null default '[]'::jsonb; +alter table seat add column emits jsonb not null default '[]'::jsonb; +alter table seat add column serves jsonb not null default '[]'::jsonb; diff --git a/internal/inventory/seats.go b/internal/inventory/seats.go index b920f82..570bf49 100644 --- a/internal/inventory/seats.go +++ b/internal/inventory/seats.go @@ -2,6 +2,7 @@ package inventory import ( "context" + "encoding/json" "fmt" "github.com/novox/mesh-controller/internal/catalogue" @@ -17,7 +18,7 @@ import ( // Seats is every seat the mesh defines, read from the store. func (i *Inventory) Seats(ctx context.Context) ([]catalogue.Seat, error) { rows, err := i.store.Pool().Query(ctx, - `select name, scope, delivers, decided from seat order by name`) + `select name, scope, delivers, decided, accepts, emits, serves from seat order by name`) if err != nil { return nil, err } @@ -26,9 +27,21 @@ func (i *Inventory) Seats(ctx context.Context) ([]catalogue.Seat, error) { var seats []catalogue.Seat for rows.Next() { var s catalogue.Seat - if err := rows.Scan(&s.Name, &s.Scope, &s.Delivers, &s.Decision); err != nil { + var accepts, emits, serves []byte + if err := rows.Scan(&s.Name, &s.Scope, &s.Delivers, &s.Decision, &accepts, &emits, &serves); err != nil { return nil, err } + // The protocol, from the row (novox/hq ADR 0132). A row that predates the columns has empty + // lists, and UseSeats keeps the compiled protocol for it until the next seeding fills them. + if err := json.Unmarshal(accepts, &s.Accepts); err != nil { + return nil, fmt.Errorf("seat %s: accepts: %w", s.Name, err) + } + if err := json.Unmarshal(emits, &s.Emits); err != nil { + return nil, fmt.Errorf("seat %s: emits: %w", s.Name, err) + } + if err := json.Unmarshal(serves, &s.Serves); err != nil { + return nil, fmt.Errorf("seat %s: serves: %w", s.Name, err) + } seats = append(seats, s) } return seats, rows.Err() @@ -41,21 +54,116 @@ func (i *Inventory) Seats(ctx context.Context) ([]catalogue.Seat, error) { // exactly as it is, so an operator's rename in the table is not undone by the next deploy putting // the old name back. What a release removes from the defaults is not deleted here either; retiring a // seat is its own decision, not a silent consequence of it dropping out of the binary. +// +// **The protocol is seeded additively** (novox/hq ADR 0132, design 33 §7). A row that has none takes +// the compiled protocol whole — that is the compiled fallback becoming data, once. A row that has one +// gains any verb the defaults name and it lacks, and loses nothing: a seat's tools are an interface, +// additive within a version, and a verb an operator added to the row is theirs to keep. func (i *Inventory) SeedSeats(ctx context.Context, defaults []catalogue.Seat) (int, error) { var added int for _, s := range defaults { - tag, err := i.store.Pool().Exec(ctx, - `insert into seat (name, scope, delivers, decided) values ($1, $2, $3, $4) - on conflict (name) do nothing`, - s.Name, s.Scope, s.Delivers, s.Decision) + accepts, emits, serves, err := protocolJSON(s) if err != nil { return added, err } - added += int(tag.RowsAffected()) + tag, err := i.store.Pool().Exec(ctx, + `insert into seat (name, scope, delivers, decided, accepts, emits, serves) + values ($1, $2, $3, $4, $5, $6, $7) + on conflict (name) do nothing`, + s.Name, s.Scope, s.Delivers, s.Decision, accepts, emits, serves) + if err != nil { + return added, err + } + if n := int(tag.RowsAffected()); n > 0 { + added += n + continue + } + if err := i.widenProtocol(ctx, s); err != nil { + return added, err + } } return added, nil } +// widenProtocol adds to a seat's row whatever the defaults name and the row lacks, by verb name. +func (i *Inventory) widenProtocol(ctx context.Context, s catalogue.Seat) error { + var accepts, emits, serves []byte + if err := i.store.Pool().QueryRow(ctx, + `select accepts, emits, serves from seat where name = $1`, s.Name).Scan(&accepts, &emits, &serves); err != nil { + return err + } + var row catalogue.Seat + if err := json.Unmarshal(accepts, &row.Accepts); err != nil { + return err + } + if err := json.Unmarshal(emits, &row.Emits); err != nil { + return err + } + if err := json.Unmarshal(serves, &row.Serves); err != nil { + return err + } + changed := false + row.Accepts, changed = union(row.Accepts, s.Accepts, changed) + row.Emits, changed = union(row.Emits, s.Emits, changed) + have := map[string]bool{} + for _, v := range row.Serves { + have[v.Name] = true + } + for _, v := range s.Serves { + if !have[v.Name] { + row.Serves = append(row.Serves, v) + changed = true + } + } + if !changed { + return nil + } + a, e, sv, err := protocolJSON(row) + if err != nil { + return err + } + _, err = i.store.Pool().Exec(ctx, + `update seat set accepts = $2, emits = $3, serves = $4 where name = $1`, s.Name, a, e, sv) + return err +} + +func union(have, want []string, changed bool) ([]string, bool) { + seen := map[string]bool{} + for _, h := range have { + seen[h] = true + } + for _, w := range want { + if !seen[w] { + have = append(have, w) + seen[w] = true + changed = true + } + } + return have, changed +} + +func protocolJSON(s catalogue.Seat) (accepts, emits, serves []byte, err error) { + if accepts, err = json.Marshal(orEmpty(s.Accepts)); err != nil { + return + } + if emits, err = json.Marshal(orEmpty(s.Emits)); err != nil { + return + } + verbs := s.Serves + if verbs == nil { + verbs = []catalogue.Verb{} + } + serves, err = json.Marshal(verbs) + return +} + +func orEmpty(s []string) []string { + if s == nil { + return []string{} + } + return s +} + // Aliases is every former seat name and the seat it now resolves to (novox/hq ADR 0122). func (i *Inventory) Aliases(ctx context.Context) (map[string]string, error) { rows, err := i.store.Pool().Query(ctx, `select alias, seat from seat_alias`) diff --git a/internal/link/seattools.go b/internal/link/seattools.go new file mode 100644 index 0000000..11c7f75 --- /dev/null +++ b/internal/link/seattools.go @@ -0,0 +1,76 @@ +package link + +import ( + "context" + "encoding/json" + "fmt" + "log" + "time" + + "github.com/nats-io/nats.go" +) + +// A role's tools, served by its holder (novox/hq ADR 0132, ADR 0154). +// +// The mesh's own verbs — `status`, `push`, `assign` — are the mesh-controller seat's tools, and the +// control plane is that seat's holder. So it answers them here, on the seat's subjects, the way a +// module's runtime answers a module's: one request, one reply on the asker's own inbox, `{result}` or +// `{error}`. Nothing about the transport is the command's business; a handler is a function of its +// arguments and gets the same answer the command line prints. + +// ToolHandler answers one call of a role's tool. What it returns is marshalled as the result; an +// error is the tool answering with one, which is an answer and not a timeout. +type ToolHandler func(ctx context.Context, args json.RawMessage) (any, error) + +// SeatToolSubject is where a mesh-scoped seat's tool is asked (design 33 §4). +func SeatToolSubject(seat, verb string) string { return "mesh.seat." + seat + ".tool." + verb } + +// HandlerTimeout bounds one answer. A verb that runs a command — a push, a build with no wait — +// answers in seconds; anything that has not in this long is said to have not answered. +const HandlerTimeout = 5 * time.Minute + +// ServeSeatTools binds every handler on its seat's subject until stopped. A queue group per seat, so +// a second holder during a handover shares the calls rather than both answering one. +func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, logger *log.Logger) (func(), error) { + var subs []*nats.Subscription + stop := func() { + for _, s := range subs { + _ = s.Unsubscribe() + } + } + for verb, handle := range handlers { + verb, handle := verb, handle + subject := SeatToolSubject(seat, verb) + sub, err := b.Conn.QueueSubscribe(subject, "seat."+seat, func(msg *nats.Msg) { + // Its own goroutine per call: a slow `push` must not hold up a `status` asked beside it, + // and the library would otherwise run handlers one after another. + go func() { + ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout) + defer cancel() + args := json.RawMessage(msg.Data) + if len(args) == 0 { + args = json.RawMessage(`{}`) + } + var reply []byte + result, err := handle(ctx, args) + if err != nil { + reply, _ = json.Marshal(map[string]any{"error": err.Error()}) + } else if reply, err = json.Marshal(map[string]any{"result": result}); err != nil { + reply, _ = json.Marshal(map[string]any{"error": "the answer could not be written as JSON: " + err.Error()}) + } + if err := msg.Respond(reply); err != nil && logger != nil { + logger.Printf("%s: could not answer: %v", subject, err) + } + }() + }) + if err != nil { + stop() + return nil, fmt.Errorf("serving %s: %w", subject, err) + } + subs = append(subs, sub) + } + if logger != nil { + logger.Printf("serving %d tool(s) of the %s seat", len(handlers), seat) + } + return stop, nil +} diff --git a/module.json b/module.json index f96c596..05fab59 100644 --- a/module.json +++ b/module.json @@ -28,6 +28,20 @@ }, "secrets-owner": "65534:65534", "prepares": true, + "tools": [ + "tools", + "status", + "nodes", + "node", + "modules", + "seats", + "builds", + "plan", + "assign", + "unassign", + "push", + "build" + ], "resources": [ { "id": "mesh-state",