// 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 at the least: every instance answers one // request, and how many will is what is being found out. var Window = 750 * time.Millisecond // Patience is how long the console waits, at the most, for the full answer of a service that has // said it is there. **A fixed window lost the largest runtime** (2026-10-05): the laptop's answer — // 48 modules, 341 endpoints, 164 kB — went up to the broker on another machine and back, arrived // last every time (a median of 365 ms, one in 25 after 813 ms on a quiet link, later still while the // runtime re-served), and when it missed the 750 ms window the console said its modules ran nowhere. // A PING answer is a hundred bytes and arrives at once, so who is there is known early; what each // serves is waited for until it arrives or this passes, and a service that never sends it is named. var Patience = 5 * time.Second // Shortened is the service metadata key a runtime sets when its announcement was cut to fit the bus. const Shortened = "shortened" // 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), } // **A seat's verb carries its schema; a module's tool does not.** Every tool's schema, once per // subject it answers on, made a runtime's one answer outgrow the bus's largest message (1 MiB) once // machines served a couple of hundred tools — and the answer that cannot be sent is silence: the // machine vanished from discovery (2026-10-04). A module's tool's schema is asked of the runtime that // serves it when a person describes it (its `tools` verb); a seat's verbs are few, and a seat has no // `tools` verb of its own to ask. if e.Kind == KindSeat { 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 } // Too large for the bus is no answer at all, and was silent: shorten every description to its // first line and say so, rather than vanish; still too large, leave descriptions out. if limit := conn.MaxPayload(); limit > 0 && int64(len(body)) > limit { if info, ok := v.(micro.Info); ok { body = shorten(conn.Logf, info, body, limit) } } 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 } // shorten is an info_response cut to fit the bus: descriptions to their first line, then none, the // service's metadata saying which (so discovery can tell a person what it was given). func shorten(logf func(string, ...any), info micro.Info, body []byte, limit int64) []byte { md := map[string]string{} for k, v := range info.Metadata { md[k] = v } info.Metadata = md for _, step := range []string{"descriptions cut to their first line", "descriptions left out"} { for i := range info.Endpoints { e := info.Endpoints[i].Metadata if step == "descriptions left out" { delete(e, "description") } else { e["description"] = firstLine(e["description"]) } } md[Shortened] = step shorter, err := json.Marshal(info) if err != nil { return body } logf("[mesh-tools] what this runtime serves is %d bytes, beyond the bus's %d; announced with %s (%d bytes)", len(body), limit, step, len(shorter)) if int64(len(shorter)) <= limit { return shorter } body = shorter } return body } // Heard is one service instance's answer to discovery: what it serves, how long it took to arrive, // and how large it was. type Heard struct { Info micro.Info Took time.Duration Bytes int Shortened string // why its descriptions are not whole, or "" } // Machine is the machine an instance runs on: its metadata's, or the one its endpoints name. func (h Heard) Machine() string { if n := h.Info.Metadata["node"]; n != "" { return n } for _, e := range h.Info.Endpoints { if n := e.Metadata["node"]; n != "" { return n } } return "" } // Instance names one service instance. type Instance struct { Name string ID string Machine string } // Discovery is what one round of discovery heard, and who said it was there but did not say what it // serves in time. type Discovery struct { At time.Time Waited time.Duration Heard []Heard Silent []Instance } // Gather asks every service on the bus who it is and what it serves, and waits at least the window // and at most the patience: until every instance that answered PING has answered INFO. An answer // that is not the protocol's is skipped. func Gather(conn *bus.Conn) (Discovery, error) { start := time.Now() d := Discovery{At: start.UTC()} pings, stopPings, err := conn.Collect("$SRV.PING", nil) if err != nil { return d, err } defer stopPings() infos, stopInfos, err := conn.Collect("$SRV.INFO", nil) if err != nil { return d, err } defer stopInfos() type key struct{ name, id string } pinged := map[key]int{} answered := map[key]int{} where := map[key]string{} var order []key complete := func() bool { for k, n := range pinged { if answered[k] < n { return false } } return true } least := time.NewTimer(Window) defer least.Stop() most := time.NewTimer(Patience) defer most.Stop() windowOver := false for { select { case a := <-pings: var p micro.Ping if json.Unmarshal(a.Data, &p) != nil || p.Type != micro.PingResponseType { continue } k := key{p.Name, p.ID} if _, seen := pinged[k]; !seen { order = append(order, k) } pinged[k]++ if where[k] == "" { where[k] = p.Metadata["node"] } case a := <-infos: var i micro.Info if json.Unmarshal(a.Data, &i) != nil || i.Type != micro.InfoResponseType { continue } answered[key{i.Name, i.ID}]++ d.Heard = append(d.Heard, Heard{Info: i, Took: a.At.Sub(start), Bytes: len(a.Data), Shortened: i.Metadata[Shortened]}) case <-least.C: windowOver = true case <-most.C: d.Waited = time.Since(start) for _, k := range order { if answered[k] < pinged[k] { d.Silent = append(d.Silent, Instance{Name: k.name, ID: k.id, Machine: where[k]}) } } return d, nil } if windowOver && complete() { d.Waited = time.Since(start) return d, 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 } // firstLine is a description's first sentence or line, at most 160 characters. func firstLine(s string) string { if i := strings.IndexAny(s, "\n"); i >= 0 { s = s[:i] } if i := strings.Index(s, ". "); i >= 0 { s = s[:i+1] } if len(s) > 160 { s = s[:157] + "..." } return s }