Every serving principal may answer the services discovery for what it serves; the controller announces its seat (hq ADR 0197)
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.
This commit is contained in:
@@ -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