Discovery from what answers: grants, the controller announces its seat, JSON lists (hq ADR 0195, 0197) #241
@@ -147,6 +147,40 @@ func moduleCommand(ctx context.Context, args []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
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 {
|
if len(entries) == 0 {
|
||||||
fmt.Println("this mesh knows about no modules yet")
|
fmt.Println("this mesh knows about no modules yet")
|
||||||
return nil
|
return nil
|
||||||
@@ -733,3 +767,23 @@ func claimsFor(ctx context.Context, inv *inventory.Inventory, m catalogue.Manife
|
|||||||
}
|
}
|
||||||
return out, nil
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
@@ -44,6 +45,21 @@ func nodeCommand(ctx context.Context, args []string) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
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 {
|
if len(nodes) == 0 {
|
||||||
// Said rather than printed as nothing: an empty list and a failed read must never
|
// 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
|
// look the same, and this command answering "none" is only honest because getting
|
||||||
@@ -500,3 +516,13 @@ func orNotReported(s string) string {
|
|||||||
}
|
}
|
||||||
return s
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -151,6 +151,12 @@ func serve(ctx context.Context) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
defer stopServing()
|
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)
|
return server.Serve(ctx)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -6,8 +6,10 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"github.com/nats-io/nats.go/micro"
|
||||||
"os"
|
"os"
|
||||||
"os/exec"
|
"os/exec"
|
||||||
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/novox/mesh-controller/internal/catalogue"
|
"github.com/novox/mesh-controller/internal/catalogue"
|
||||||
@@ -66,14 +68,14 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
|
|||||||
case "status":
|
case "status":
|
||||||
return []string{"status", "--json"}, nil
|
return []string{"status", "--json"}, nil
|
||||||
case "nodes":
|
case "nodes":
|
||||||
return []string{"node", "list"}, nil
|
return []string{"node", "list", "--json"}, nil
|
||||||
case "node":
|
case "node":
|
||||||
if err := need("node"); err != nil {
|
if err := need("node"); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return []string{"node", "show", str("node")}, nil
|
return []string{"node", "show", str("node")}, nil
|
||||||
case "modules":
|
case "modules":
|
||||||
return []string{"module", "list"}, nil
|
return []string{"module", "list", "--json"}, nil
|
||||||
case "seats":
|
case "seats":
|
||||||
return []string{"seats", "--json"}, nil
|
return []string{"seats", "--json"}, nil
|
||||||
case "builds":
|
case "builds":
|
||||||
@@ -379,3 +381,43 @@ func splitCommandLine(line string) ([]string, error) {
|
|||||||
}
|
}
|
||||||
return words, nil
|
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,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"github.com/novox/mesh-controller/internal/link"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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
|
// which this package mirrors rather than reads, and a verb the seat does not declare is a
|
||||||
// subject nothing publishes.
|
// subject nothing publishes.
|
||||||
sub = append(sub, "mesh.seat."+ControllerSeat+".tool.>")
|
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
|
// 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
|
// (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
|
return Permissions{}, err
|
||||||
}
|
}
|
||||||
pub = append(pub, invoked...)
|
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:
|
case KindEnrolment:
|
||||||
// A leaked token is useless for anything but enrolling: it cannot read a declaration, hear
|
// 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
|
// away — no other principal may subscribe this namespace, and a caller's authority is
|
||||||
// still granted per tool, by name, on the publish side.
|
// still granted per tool, by name, on the publish side.
|
||||||
sub = append(sub, own+".tool.>")
|
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
|
// Its own membership (ADR 0160): the one subject a runtime derives for itself, read
|
||||||
// directly from the stream and followed live. Nothing else's.
|
// directly from the stream and followed live. Nothing else's.
|
||||||
sub = append(sub, MembershipSubject(p.Node, p.Module))
|
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
|
// 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
|
// 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.
|
// this node, as the holder's own principal would be granted them.
|
||||||
|
var serves []string
|
||||||
for _, d := range p.Carries {
|
for _, d := range p.Carries {
|
||||||
if !safeSubject.MatchString(d.Module) {
|
if !safeSubject.MatchString(d.Module) {
|
||||||
return Permissions{}, fmt.Errorf(
|
return Permissions{}, fmt.Errorf(
|
||||||
"%q cannot be part of a subject: a permission is a subject pattern, and this would widen it", d.Module)
|
"%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
|
own := "mesh.mod." + d.Module
|
||||||
sub = append(sub, own+".tool.>")
|
sub = append(sub, own+".tool.>")
|
||||||
// A tool that emits an event is the module's code and emits under the module's name
|
// 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
|
return Permissions{}, err
|
||||||
}
|
}
|
||||||
pub = append(pub, invoked...)
|
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
|
// 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.
|
// own long-lived process, which ADR 0175 leaves where it is; what moves here is tools.
|
||||||
sub = unique(sub)
|
sub = unique(sub)
|
||||||
@@ -765,3 +788,23 @@ func invokedSubjects(invokes []string) ([]string, error) {
|
|||||||
}
|
}
|
||||||
return out, nil
|
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.>"}
|
||||||
|
}
|
||||||
|
|||||||
@@ -233,7 +233,9 @@ func TestAPersonReachesNothingButTools(t *testing.T) {
|
|||||||
perms, _ := PermissionsFor(Principal{Kind: KindPerson, Module: "jo",
|
perms, _ := PermissionsFor(Principal{Kind: KindPerson, Module: "jo",
|
||||||
Invokes: []string{"*"}, PasswordHash: "x"})
|
Invokes: []string{"*"}, PasswordHash: "x"})
|
||||||
for _, p := range perms.Publish {
|
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)
|
t.Errorf("a person may publish %q, which is not a tool call", p)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -423,14 +425,18 @@ func TestTheRuntimeServesTheUnionAndConsumesNothing(t *testing.T) {
|
|||||||
if _, needed := ConsumerFor(p); needed {
|
if _, needed := ConsumerFor(p); needed {
|
||||||
t.Error("a consumer would be made for the runtime, which consumes nothing")
|
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.
|
// 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{}
|
seen := map[string]bool{}
|
||||||
for _, s := range append(append([]string{}, perms.Subscribe...), perms.Publish...) {
|
for _, s := range list {
|
||||||
if seen[s] {
|
if seen[s] {
|
||||||
t.Errorf("%s is granted twice", s)
|
t.Errorf("%s is granted twice", s)
|
||||||
}
|
}
|
||||||
seen[s] = true
|
seen[s] = true
|
||||||
}
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func contains(list []string, want string) bool {
|
func contains(list []string, want string) bool {
|
||||||
|
|||||||
+4
-4
@@ -25,7 +25,7 @@ accounts {
|
|||||||
users = [
|
users = [
|
||||||
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
|
{ 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.>"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
||||||
@@ -38,17 +38,17 @@ accounts {
|
|||||||
} }
|
} }
|
||||||
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
{ 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"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
{ 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"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
|
{ 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"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user