diff --git a/cmd/mesh-controller/modules.go b/cmd/mesh-controller/modules.go index 2a91408..4b9fdf3 100644 --- a/cmd/mesh-controller/modules.go +++ b/cmd/mesh-controller/modules.go @@ -147,6 +147,40 @@ func moduleCommand(ctx context.Context, args []string) error { if err != nil { return err } + // **The same list, for something other than a person** (novox/hq ADR 0195): what each module + // is, where it runs, whether it is current, and what it says of itself. + if len(args) > 1 && args[1] == "--json" { + type listed struct { + Module string `json:"module"` + Version string `json:"version"` + Built string `json:"built,omitempty"` + Head string `json:"head,omitempty"` + Current bool `json:"current"` + Provided bool `json:"provided,omitempty"` + // Tools says whether the module answers tools anywhere it runs: a list of its own, + // a bundle the runtime serves, or a seat's verbs it claims (novox/hq ADR 0197) — + // what the console checks the bus's answers against. + Tools bool `json:"tools"` + On []string `json:"on"` + Provides []string `json:"provides,omitempty"` + Requires []string `json:"requires,omitempty"` + Claims []string `json:"claims,omitempty"` + Capabilities []string `json:"capabilities,omitempty"` + } + out := make([]listed, 0, len(entries)) + for _, e := range entries { + m := e.Manifest + l := listed{Module: m.Module, Version: m.Version, Built: e.Source.BuiltFrom, Head: e.Source.Head, + Current: e.Provided || e.Source.Repository == "" || e.Source.Current(), Provided: e.Provided, + On: append([]string{}, e.On...), Provides: m.Offers(), Requires: m.Requires, + Capabilities: m.Capabilities, Tools: declaresTools(m)} + for _, c := range m.Claims { + l.Claims = append(l.Claims, c.At()+"/"+c.Name) + } + out = append(out, l) + } + return printJSON(out) + } if len(entries) == 0 { fmt.Println("this mesh knows about no modules yet") return nil @@ -733,3 +767,23 @@ func claimsFor(ctx context.Context, inv *inventory.Inventory, m catalogue.Manife } return out, nil } + +// declaresTools is whether a module answers tools wherever it runs (novox/hq ADR 0197): it names +// tools of its own, its build delivers a bundle the node's runtime serves, or it claims a seat +// whose verbs it serves. A module with none is never expected to announce anything. +func declaresTools(m catalogue.Manifest) bool { + if len(m.Tools) > 0 { + return true + } + for _, b := range m.Bundles { + if len(b.Loads) > 0 { + return true + } + } + for _, c := range m.Claims { + if len(c.Serves) > 0 { + return true + } + } + return false +} diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index fe96b48..04a38e4 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -2,6 +2,7 @@ package main import ( "context" + "encoding/json" "errors" "flag" "fmt" @@ -44,6 +45,21 @@ func nodeCommand(ctx context.Context, args []string) error { if err != nil { return err } + // **The same list, for something other than a person** — the console's discovery reads it + // (novox/hq ADR 0195), and a reader that parses a printed column breaks when it is reworded. + if len(args) > 1 && args[1] == "--json" { + type listed struct { + Name string `json:"name"` + Heard string `json:"heard"` + Mode string `json:"mode"` + ID string `json:"id"` + } + out := make([]listed, 0, len(nodes)) + for _, n := range nodes { + out = append(out, listed{Name: n.Name, Heard: heardFrom(n), Mode: modeOf(n), ID: n.ID}) + } + return printJSON(out) + } if len(nodes) == 0 { // Said rather than printed as nothing: an empty list and a failed read must never // look the same, and this command answering "none" is only honest because getting @@ -500,3 +516,13 @@ func orNotReported(s string) string { } return s } + +// printJSON prints a value as indented JSON, the shape every `--json` answers in. +func printJSON(v any) error { + body, err := json.MarshalIndent(v, "", " ") + if err != nil { + return err + } + fmt.Println(string(body)) + return nil +} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 001f235..3f5f73f 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -151,6 +151,12 @@ func serve(ctx context.Context) error { return err } defer stopServing() + // And says so on the bus (novox/hq ADR 0197): what it serves, as the NATS services protocol asks. + stopAnnouncing, err := bus.Announce(seatAnnouncement(handlers), log.New(os.Stdout, "", log.LstdFlags)) + if err != nil { + return err + } + defer stopAnnouncing() return server.Serve(ctx) } diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 85d3fe3..523ab08 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -6,8 +6,10 @@ import ( "encoding/json" "errors" "fmt" + "github.com/nats-io/nats.go/micro" "os" "os/exec" + "sort" "strings" "github.com/novox/mesh-controller/internal/catalogue" @@ -66,14 +68,14 @@ func argvFor(verb string, args map[string]any) ([]string, error) { case "status": return []string{"status", "--json"}, nil case "nodes": - return []string{"node", "list"}, nil + return []string{"node", "list", "--json"}, 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 + return []string{"module", "list", "--json"}, nil case "seats": return []string{"seats", "--json"}, nil case "builds": @@ -379,3 +381,43 @@ func splitCommandLine(line string) ([]string, error) { } return words, nil } + +// seatAnnouncement is what the controller says it serves on the bus (novox/hq ADR 0197): the +// mesh-controller seat, one endpoint per verb it answers, each with the seat's own description and +// argument schema — the same facts `tools` answers from the records, as NATS's services format. +func seatAnnouncement(handlers map[string]link.ToolHandler) micro.Info { + about := map[string]catalogue.Verb{} + for _, s := range catalogue.SeatsWithAProtocol() { + if s.Name == catalogue.ControllerSeatName { + for _, v := range s.Serves { + about[v.Name] = v + } + } + } + verbs := make([]string, 0, len(handlers)) + for verb := range handlers { + verbs = append(verbs, verb) + } + sort.Strings(verbs) + var endpoints []micro.EndpointInfo + for _, verb := range verbs { + schema, _ := json.Marshal(about[verb].Input) + endpoints = append(endpoints, micro.EndpointInfo{ + Name: verb, + Subject: link.SeatToolSubject(catalogue.ControllerSeatName, verb), + QueueGroup: "seat." + catalogue.ControllerSeatName, + Metadata: map[string]string{ + "description": about[verb].Description, "schema": string(schema), + "seat": catalogue.ControllerSeatName, "scope": "mesh", + }, + }) + } + return micro.Info{ + ServiceIdentity: micro.ServiceIdentity{ + Name: catalogue.ControllerSeatName, ID: "controller", Version: "1.0.0", + Metadata: map[string]string{"seat": catalogue.ControllerSeatName, "scope": "mesh"}, + }, + Description: "the mesh's own verbs, answered by the holder of the mesh-controller seat", + Endpoints: endpoints, + } +} diff --git a/cmd/mesh-controller/seatverbs_test.go b/cmd/mesh-controller/seatverbs_test.go index 95cb8ba..df3177f 100644 --- a/cmd/mesh-controller/seatverbs_test.go +++ b/cmd/mesh-controller/seatverbs_test.go @@ -2,6 +2,8 @@ package main import ( "context" + "fmt" + "github.com/novox/mesh-controller/internal/link" "strings" "testing" @@ -249,3 +251,44 @@ func TestARowAheadOfThisBuildIsServedAnyway(t *testing.T) { } } } + +// novox/hq ADR 0195: the console's discovery reads the machines and the modules; they answer as JSON, +// as status and seats do, so nothing parses a printed column. +func TestTheNodesAndModulesVerbsAnswerAsJSON(t *testing.T) { + for verb, want := range map[string]string{"nodes": "[node list --json]", "modules": "[module list --json]"} { + argv, err := argvFor(verb, map[string]any{}) + if err != nil { + t.Fatal(err) + } + if fmt.Sprint(argv) != want { + t.Errorf("%s runs %v, want %s", verb, argv, want) + } + } +} + +// novox/hq ADR 0197: the controller announces exactly the verbs it serves, each on the subject and +// queue it serves it on, with the seat's own description and schema, in NATS's services format. +func TestTheControllerAnnouncesTheVerbsItServes(t *testing.T) { + handlers, _, err := seatToolHandlers() + if err != nil { + t.Fatal(err) + } + info := seatAnnouncement(handlers) + if info.Name != catalogue.ControllerSeatName || info.ID == "" || info.Version == "" { + t.Fatalf("the service is not named for the seat: %+v", info.ServiceIdentity) + } + if len(info.Endpoints) != len(handlers) { + t.Fatalf("%d endpoints announced for %d verbs served", len(info.Endpoints), len(handlers)) + } + for _, e := range info.Endpoints { + if _, served := handlers[e.Name]; !served { + t.Errorf("%s is announced and not served", e.Name) + } + if e.Subject != link.SeatToolSubject(catalogue.ControllerSeatName, e.Name) || e.QueueGroup != "seat."+catalogue.ControllerSeatName { + t.Errorf("%s is announced on %s/%s, not where it is served", e.Name, e.Subject, e.QueueGroup) + } + if e.Metadata["description"] == "" || e.Metadata["schema"] == "" || e.Metadata["scope"] != "mesh" { + t.Errorf("%s is announced without its description, schema or scope: %v", e.Name, e.Metadata) + } + } +} diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 025bcd5..f547d03 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -236,6 +236,8 @@ func PermissionsFor(p Principal) (Permissions, error) { // 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.>") + // And says so (novox/hq ADR 0197): it answers discovery for the seat it serves. + sub = append(sub, announcing(ControllerSeat)...) // 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 @@ -271,6 +273,9 @@ func PermissionsFor(p Principal) (Permissions, error) { return Permissions{}, err } pub = append(pub, invoked...) + // And may ask what answers (novox/hq ADR 0197): a question every service answers about + // itself, its replies to the asker's own inbox. + pub = append(pub, discovering()...) case KindEnrolment: // A leaked token is useless for anything but enrolling: it cannot read a declaration, hear @@ -327,6 +332,15 @@ func PermissionsFor(p Principal) (Permissions, error) { // away — no other principal may subscribe this namespace, and a caller's authority is // still granted per tool, by name, on the publish side. sub = append(sub, own+".tool.>") + // It says what it serves (novox/hq ADR 0197): discovery for its own name and every seat it + // holds a verb of, answered by the runtime that serves them. + announced := []string{p.Module} + for _, s := range p.Holds { + if len(s.Serves) > 0 { + announced = append(announced, s.Name) + } + } + sub = append(sub, announcing(announced...)...) // Its own membership (ADR 0160): the one subject a runtime derives for itself, read // directly from the stream and followed live. Nothing else's. sub = append(sub, MembershipSubject(p.Node, p.Module)) @@ -413,11 +427,16 @@ func PermissionsFor(p Principal) (Permissions, error) { // module's own principal has, for the same reason: the tools a module serves are what its // code answers, and a list here would be a second copy of it. Each held seat's verbs on // this node, as the holder's own principal would be granted them. + var serves []string for _, d := range p.Carries { if !safeSubject.MatchString(d.Module) { return Permissions{}, fmt.Errorf( "%q cannot be part of a subject: a permission is a subject pattern, and this would widen it", d.Module) } + serves = append(serves, d.Module) + for _, s := range d.Holds { + serves = append(serves, s.Name) + } own := "mesh.mod." + d.Module sub = append(sub, own+".tool.>") // A tool that emits an event is the module's code and emits under the module's name @@ -444,6 +463,10 @@ func PermissionsFor(p Principal) (Permissions, error) { return Permissions{}, err } pub = append(pub, invoked...) + // It says what it serves and may ask what answers (novox/hq ADR 0197): the runtime answers + // discovery for each module and seat it carries, and the console it is asks the bus. + sub = append(sub, announcing(serves...)...) + pub = append(pub, discovering()...) // Nothing about consumers: it consumes nothing. A module's reactions to events are its // own long-lived process, which ADR 0175 leaves where it is; what moves here is tools. sub = unique(sub) @@ -765,3 +788,23 @@ func invokedSubjects(invokes []string) ([]string, error) { } return out, nil } + +// announcing is what a principal that serves tools subscribes to answer the NATS services +// protocol's discovery (novox/hq ADR 0197): the questions asked of every service, and those asked of +// each name it serves — its own and no other's, so it cannot answer for a service it is not. +func announcing(names ...string) []string { + out := []string{"$SRV.PING", "$SRV.INFO"} + for _, n := range names { + if !safeSubject.MatchString(n) { + continue + } + out = append(out, "$SRV.PING."+n, "$SRV.PING."+n+".>", "$SRV.INFO."+n, "$SRV.INFO."+n+".>") + } + return out +} + +// discovering is what a principal publishes to ask what answers (novox/hq ADR 0197): the services +// protocol's discovery requests, whose replies come to its own inbox. +func discovering() []string { + return []string{"$SRV.PING", "$SRV.PING.>", "$SRV.INFO", "$SRV.INFO.>"} +} diff --git a/internal/broker/nats_test.go b/internal/broker/nats_test.go index 1a60c2e..dad5ed9 100644 --- a/internal/broker/nats_test.go +++ b/internal/broker/nats_test.go @@ -233,7 +233,9 @@ func TestAPersonReachesNothingButTools(t *testing.T) { perms, _ := PermissionsFor(Principal{Kind: KindPerson, Module: "jo", Invokes: []string{"*"}, PasswordHash: "x"}) for _, p := range perms.Publish { - if !strings.Contains(p, ".tool.") { + // A tool call, or asking what answers (novox/hq ADR 0197) — a question every service + // answers about itself, which claims nothing and controls nothing. + if !strings.Contains(p, ".tool.") && !strings.HasPrefix(p, "$SRV.") { t.Errorf("a person may publish %q, which is not a tool call", p) } } @@ -423,13 +425,17 @@ func TestTheRuntimeServesTheUnionAndConsumesNothing(t *testing.T) { if _, needed := ConsumerFor(p); needed { t.Error("a consumer would be made for the runtime, which consumes nothing") } - // Each subject once: the file is read as the mesh's authority model. - seen := map[string]bool{} - for _, s := range append(append([]string{}, perms.Subscribe...), perms.Publish...) { - if seen[s] { - t.Errorf("%s is granted twice", s) + // Each subject once in each list: the file is read as the mesh's authority model. One subject may + // stand in both — the runtime answers discovery on `$SRV.INFO` and, as the console, asks it + // (novox/hq ADR 0197) — because subscribing and publishing are two different grants. + for _, list := range [][]string{perms.Subscribe, perms.Publish} { + seen := map[string]bool{} + for _, s := range list { + if seen[s] { + t.Errorf("%s is granted twice", s) + } + seen[s] = true } - seen[s] = true } } diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index ba18138..45d607f 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.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] } - subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } + subscribe: { allow: ["$JS.API.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { @@ -38,17 +38,17 @@ accounts { } } { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } - subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } + subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.telegram", "$SRV.INFO.telegram.>", "$SRV.PING", "$SRV.PING.telegram", "$SRV.PING.telegram.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.audit"] } - subscribe: { allow: ["_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] } + subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.audit", "$SRV.INFO.audit.>", "$SRV.PING", "$SRV.PING.audit", "$SRV.PING.audit.>", "_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: { publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] } - subscribe: { allow: ["_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] } + subscribe: { allow: ["$SRV.INFO", "$SRV.INFO.shop", "$SRV.INFO.shop.>", "$SRV.PING", "$SRV.PING.shop", "$SRV.PING.shop.>", "_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] } allow_responses: { max: 1, ttl: "1m" } } } ] diff --git a/internal/link/announce.go b/internal/link/announce.go new file mode 100644 index 0000000..f822100 --- /dev/null +++ b/internal/link/announce.go @@ -0,0 +1,68 @@ +package link + +import ( + "encoding/json" + "fmt" + "log" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/micro" +) + +// What answers announces itself (novox/hq ADR 0197). A holder that serves a seat's verbs answers the +// NATS services protocol's discovery — `$SRV.PING` and `$SRV.INFO`, and each by its service's name +// and instance — with exactly what it serves, in NATS's own format, so the console and the standard +// `nats micro` commands learn what exists from what answers rather than from a roster. + +// DiscoverySubjects are where one service instance is asked to say what it is. +func DiscoverySubjects(name, id string) []string { + var out []string + for _, verb := range []string{"PING", "INFO"} { + out = append(out, "$SRV."+verb, "$SRV."+verb+"."+name, "$SRV."+verb+"."+name+"."+id) + } + return out +} + +// Announce answers discovery for one service until stopped. The answer is fixed at the call: a holder +// whose verbs change announces again. Every instance answers, so there is no queue group. +func (b OverNATS) Announce(info micro.Info, logger *log.Logger) (func(), error) { + info.Type = micro.InfoResponseType + infoBody, err := json.Marshal(info) + if err != nil { + return nil, err + } + pingBody, err := json.Marshal(micro.Ping{ServiceIdentity: info.ServiceIdentity, Type: micro.PingResponseType}) + if err != nil { + return nil, err + } + var subs []*nats.Subscription + done := make(chan struct{}) + stop := func() { + close(done) + for _, s := range subs { + _ = s.Unsubscribe() + } + } + for _, subject := range DiscoverySubjects(info.Name, info.ID) { + subject := subject + body := infoBody + if len(subject) >= 9 && subject[:9] == "$SRV.PING" { + body = pingBody + } + bind := func() (*nats.Subscription, error) { + return b.Conn.Subscribe(subject, func(msg *nats.Msg) { + if err := msg.Respond(body); err != nil && logger != nil { + logger.Printf("%s: could not answer: %v", subject, err) + } + }) + } + sub, err := bind() + if err != nil { + stop() + return nil, fmt.Errorf("announcing %s on %s: %w", info.Name, subject, err) + } + subs = append(subs, sub) + go keepBound(sub, bind, subject, done, logger) + } + return stop, nil +}