Both runtimes subscribed $SRV.<verb>.> as a wildcard; the grants allow the bare question and the service's own name and instance. The bus refused the wildcard, and the TypeScript runtime treats a refused subscription as fatal, so every per-module container crash-looped after the image rolled. They now subscribe exactly $SRV.<verb>, $SRV.<verb>.<name> and $SRV.<verb>.<name>.<id>.
212 lines
6.5 KiB
Go
212 lines
6.5 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"} {
|
|
// Exactly the questions asked of every service and of this one by name and instance — what the
|
|
// grants allow (novox/hq ADR 0197). A wildcard is refused by the bus.
|
|
for _, subject := range []string{"$SRV." + verb, "$SRV." + verb + "." + s.Name, "$SRV." + verb + "." + s.Name + "." + s.ID} {
|
|
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
|
|
}
|