Merge pull request 'The mesh's own verbs are the mesh-controller seat's tools' (#166) from feat/the-mesh-answers-for-itself into main
Reviewed-on: #166
This commit was merged in pull request #166.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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}
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
+41
-4
@@ -40,6 +40,10 @@ const (
|
||||
// emits (novox/hq ADR 0118, design 29 §5).
|
||||
type Seat struct {
|
||||
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:<seat>.<verb>`. 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:<seat>.<verb>", 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 <module>.<tool>, or * for every one", t)
|
||||
"%q does not name a tool: one invokes <module>.<tool>, seat:<seat>.<verb>, or * for every one", t)
|
||||
}
|
||||
out = append(out, "mesh.mod."+module+".tool."+tool)
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -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: {
|
||||
|
||||
@@ -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 `<module>.<tool>`, or the single entry `*` for
|
||||
// every tool on the mesh (novox/hq ADR 0152).
|
||||
// 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).
|
||||
//
|
||||
// **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(
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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}}}
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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))
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
+115
-7
@@ -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`)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
+14
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user