Grants: a principal that serves tools subscribes $SRV.PING/$SRV.INFO and those questions under each name it serves — its own and no other's; the tool runtime and people may ask. The controller answers discovery for the mesh-controller seat in NATS's services format, one endpoint per verb it serves, with the seat's description and schema. module list --json says which modules declare tools, so the console expects an announcement only from those.
69 lines
2.1 KiB
Go
69 lines
2.1 KiB
Go
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
|
|
}
|