The Go runtime answers $SRV.PING, $SRV.INFO and $SRV.STATS (and per name and id) in the io.nats.micro.v1 format with what it serves at the moment it is asked: one service per runtime process, since the bus admits one reply per request from each responder, and one endpoint per tool per subject, its metadata saying module, seat, scope, machine, description, schema and whether the module is interchangeable. Serving is unchanged. The console gathers one $SRV.INFO request's answers instead of asking the catalogue's roster and each module's tools, and reads the controller's records as JSON for what should have answered: an assignment with tools that did not announce is named, a module without tools never is. The text parsers of node list and module list are gone. Packages share the test bus: go test -p 1.
210 lines
6.3 KiB
Go
210 lines
6.3 KiB
Go
// Package announce answers the NATS services protocol's discovery for what a runtime serves, and
|
|
// gathers the answers (novox/hq ADR 0197).
|
|
//
|
|
// A runtime does not re-serve its tools through a services library: serving is unchanged. It answers
|
|
// `$SRV.PING`, `$SRV.INFO` and `$SRV.STATS` — and the same followed by its service name, and by its
|
|
// name and id — in the format NATS's own tools read, with what it is serving at the moment it is
|
|
// asked. One service per runtime process: the bus admits one reply per request from each responder,
|
|
// so a runtime serving many modules and seats answers once, one endpoint per tool per subject, and
|
|
// says in each endpoint's metadata which module, seat, scope and machine it is.
|
|
package announce
|
|
|
|
import (
|
|
"encoding/json"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go/micro"
|
|
|
|
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
|
)
|
|
|
|
// Version is the announced service version (semver, as the protocol requires).
|
|
const Version = "0.1.0"
|
|
|
|
// Window is how long the console gathers discovery answers: every instance answers one request, and
|
|
// how many will is what is being found out.
|
|
var Window = 750 * time.Millisecond
|
|
|
|
// Kinds of endpoint.
|
|
const (
|
|
KindTool = "tool" // a module's own tool
|
|
KindSeat = "seat" // a seat's verb, served by the module holding it
|
|
)
|
|
|
|
// Endpoint is one tool served on one subject, as it is announced.
|
|
type Endpoint struct {
|
|
Kind string
|
|
Module string // the module whose code answers
|
|
Tool string // the tool's or the verb's name
|
|
Seat string // for a seat's verb
|
|
Scope string // "mesh" or "node", for a seat's verb
|
|
Node string
|
|
Description string
|
|
Schema json.RawMessage
|
|
Interchangeable bool
|
|
Subject string
|
|
Queue string
|
|
}
|
|
|
|
// Name is the endpoint's name as the protocol allows it — letters, digits, `-` and `_` — the
|
|
// prefix and the tool joined by `__`; the metadata, not the name, is what identifies it.
|
|
func (e Endpoint) Name() string {
|
|
prefix := e.Module
|
|
if e.Kind == KindSeat {
|
|
prefix = e.Seat
|
|
}
|
|
return clean(prefix) + "__" + clean(e.Tool)
|
|
}
|
|
|
|
func clean(s string) string {
|
|
var b strings.Builder
|
|
for _, r := range s {
|
|
if r == '-' || r == '_' || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') {
|
|
b.WriteRune(r)
|
|
} else {
|
|
b.WriteRune('_')
|
|
}
|
|
}
|
|
return b.String()
|
|
}
|
|
|
|
func (e Endpoint) info() micro.EndpointInfo {
|
|
md := map[string]string{
|
|
"kind": e.Kind, "module": e.Module, "tool": e.Tool, "node": e.Node,
|
|
"description": e.Description, "interchangeable": boolWord(e.Interchangeable),
|
|
}
|
|
schema := strings.TrimSpace(string(e.Schema))
|
|
if schema == "" || schema == "null" {
|
|
schema = "{}"
|
|
}
|
|
md["schema"] = schema
|
|
if e.Kind == KindSeat {
|
|
md["seat"] = e.Seat
|
|
md["scope"] = e.Scope
|
|
}
|
|
return micro.EndpointInfo{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue, Metadata: md}
|
|
}
|
|
|
|
func boolWord(b bool) string {
|
|
if b {
|
|
return "true"
|
|
}
|
|
return "false"
|
|
}
|
|
|
|
// Service is who answers: the runtime's name, its instance, and a word about it.
|
|
type Service struct {
|
|
Name string
|
|
ID string
|
|
Description string
|
|
Metadata map[string]string
|
|
}
|
|
|
|
func (s Service) identity() micro.ServiceIdentity {
|
|
md := s.Metadata
|
|
if md == nil {
|
|
md = map[string]string{}
|
|
}
|
|
return micro.ServiceIdentity{Name: s.Name, ID: s.ID, Version: Version, Metadata: md}
|
|
}
|
|
|
|
// Info is the info_response for these endpoints.
|
|
func Info(s Service, endpoints []Endpoint) micro.Info {
|
|
out := micro.Info{ServiceIdentity: s.identity(), Type: micro.InfoResponseType,
|
|
Description: s.Description, Endpoints: []micro.EndpointInfo{}}
|
|
for _, e := range endpoints {
|
|
out.Endpoints = append(out.Endpoints, e.info())
|
|
}
|
|
return out
|
|
}
|
|
|
|
// Serve answers discovery for one service until stopped, asking `current` for its endpoints each
|
|
// time — so what is announced is what is served now, re-served memberships included.
|
|
func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error) {
|
|
started := time.Now().UTC()
|
|
answer := func(subject string, _ []byte) []byte {
|
|
parts := strings.Split(subject, ".")
|
|
if len(parts) < 2 || parts[0] != "$SRV" {
|
|
return nil
|
|
}
|
|
if len(parts) >= 3 && parts[2] != s.Name {
|
|
return nil // another service's
|
|
}
|
|
if len(parts) >= 4 && parts[3] != s.ID {
|
|
return nil // another instance's
|
|
}
|
|
var v any
|
|
switch parts[1] {
|
|
case "PING":
|
|
v = micro.Ping{ServiceIdentity: s.identity(), Type: micro.PingResponseType}
|
|
case "INFO":
|
|
v = Info(s, current())
|
|
case "STATS":
|
|
st := micro.Stats{ServiceIdentity: s.identity(), Type: micro.StatsResponseType, Started: started,
|
|
Endpoints: []*micro.EndpointStats{}}
|
|
for _, e := range current() {
|
|
st.Endpoints = append(st.Endpoints, µ.EndpointStats{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue})
|
|
}
|
|
v = st
|
|
default:
|
|
return nil
|
|
}
|
|
body, err := json.Marshal(v)
|
|
if err != nil {
|
|
return nil
|
|
}
|
|
return body
|
|
}
|
|
var stops []func()
|
|
for _, verb := range []string{"PING", "INFO", "STATS"} {
|
|
for _, subject := range []string{"$SRV." + verb, "$SRV." + verb + ".>"} {
|
|
stop, err := conn.Raw(subject, answer)
|
|
if err != nil {
|
|
for _, st := range stops {
|
|
st()
|
|
}
|
|
return func() {}, err
|
|
}
|
|
stops = append(stops, stop)
|
|
}
|
|
}
|
|
return func() {
|
|
for _, st := range stops {
|
|
st()
|
|
}
|
|
}, nil
|
|
}
|
|
|
|
// Gather asks every service on the bus what it serves and answers what came back within the window.
|
|
// An answer that is not an info_response is skipped.
|
|
func Gather(conn *bus.Conn) ([]micro.Info, error) {
|
|
raw, err := conn.Gather("$SRV.INFO", nil, Window)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var out []micro.Info
|
|
for _, b := range raw {
|
|
var i micro.Info
|
|
if json.Unmarshal(b, &i) != nil || i.Type != micro.InfoResponseType {
|
|
continue
|
|
}
|
|
out = append(out, i)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// Endpoints reads an info_response's endpoints back into what they announce.
|
|
func Endpoints(i micro.Info) []Endpoint {
|
|
var out []Endpoint
|
|
for _, e := range i.Endpoints {
|
|
md := e.Metadata
|
|
out = append(out, Endpoint{
|
|
Kind: md["kind"], Module: md["module"], Tool: md["tool"], Seat: md["seat"], Scope: md["scope"],
|
|
Node: md["node"], Description: md["description"], Schema: json.RawMessage(md["schema"]),
|
|
Interchangeable: md["interchangeable"] == "true", Subject: e.Subject, Queue: e.QueueGroup,
|
|
})
|
|
}
|
|
return out
|
|
}
|