Files
mesh-controller/internal/link/announce.go
T
jochen 58b4fcb8c8 A bundle stands on the toolchain it is compiled in (hq issue 211)
A manifest names its toolchain by language, not in build.on, so the planner did not know a bundle
depends on the module that publishes its toolchain and built the two in one tier: the bundle
against the old toolchain, recorded as built from the new commit. The edge is read from the
manifest, so it holds before any build recorded it, and a toolchain that moves rebuilds every
bundle compiled in it.
2026-10-03 22:18:20 +02:00

82 lines
2.7 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", "STATS"} {
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
}
// Statistics the protocol asks for; the controller keeps none per verb, so it answers its
// identity and its endpoints with nothing counted — an honest zero, not a refusal.
stats := micro.Stats{ServiceIdentity: info.ServiceIdentity, Type: micro.StatsResponseType}
for _, e := range info.Endpoints {
stats.Endpoints = append(stats.Endpoints, &micro.EndpointStats{Name: e.Name, Subject: e.Subject, QueueGroup: e.QueueGroup})
}
statsBody, err := json.Marshal(stats)
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
switch {
case len(subject) >= 9 && subject[:9] == "$SRV.PING":
body = pingBody
case len(subject) >= 10 && subject[:10] == "$SRV.STATS":
body = statsBody
}
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
}