Files
mesh-catalog/modules/nats/cmd/nats-tools/main.go
T
jochen f53fc0929b The bus's own tools: server, connections, subscriptions, streams, backlog, buckets, users, user_can
A Go bundle the runtime launches beside the nats module's server. It reads
the server's monitoring API and the composed user list — never a password
hash — and changes nothing. Reached directly when the endpoint is published,
through the container otherwise: its configuration binds monitoring to the
container's own loopback, so the published port answers nothing today.
2026-10-04 16:51:33 +02:00

133 lines
5.7 KiB
Go

// nats-tools: the mesh bus's own tools (novox/hq design 25), a Go bundle the node's runtime launches on the
// machine that runs the bus. Everything here reads — the server's monitoring API and the user list the
// controller composed — and nothing changes the bus. stdout is the MCP channel; this says nothing else.
package main
import (
"context"
"fmt"
"os"
"time"
stdio "git.novox.be/novox/mesh-sdk/go"
)
func str(description string) map[string]any {
return map[string]any{"type": "string", "description": description}
}
func num(description string) map[string]any {
return map[string]any{"type": "number", "description": description}
}
func arg(a map[string]any, k string) string { s, _ := a[k].(string); return s }
func limitOf(a map[string]any, fallback int) int {
if v, ok := a["limit"].(float64); ok && v > 0 {
return int(v)
}
return fallback
}
func tools(monitor Monitor, users func(context.Context) ([]byte, error)) []stdio.Tool {
get := func(path string) ([]byte, error) {
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
defer cancel()
return monitor(ctx, path)
}
conf := func() ([]byte, error) {
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
defer cancel()
return users(ctx)
}
jsz := "/jsz?accounts=1&streams=1&consumers=1&config=1"
return []stdio.Tool{
{Name: "nats_server",
Description: "The bus server as it runs: version, uptime, limits (max_payload included), connections, subscriptions, slow consumers, traffic, memory and CPU, JetStream totals, and its health.",
Run: func(map[string]any) (any, error) {
varz, err := get("/varz")
if err != nil {
return nil, err
}
health, _ := get("/healthz?js-enabled-only=true")
return Server(varz, health)
}},
{Name: "nats_connections",
Description: "Every client connected to the bus, named by the mesh user it is (a machine's host `node.<machine>`, a machine's runtime `<machine>.node-tools`, the controller, a person): address, uptime, round trip, traffic and pending bytes. Narrowed by a user or connection name, sorted by out_bytes (default), in_bytes, pending, msgs or subscriptions.",
Input: map[string]any{"user": str("part of a user or connection name, e.g. g14"), "sort": str("out_bytes, in_bytes, pending, msgs or subscriptions"), "limit": num("at most this many (default 50)")},
Run: func(a map[string]any) (any, error) {
connz, err := get("/connz?auth=1&limit=1024")
if err != nil {
return nil, err
}
return Connections(connz, arg(a, "user"), arg(a, "sort"), limitOf(a, 50))
}},
{Name: "nats_subscriptions",
Description: "Who listens on what: every subscription whose subject contains the text given, with its queue group, the messages it received and the user holding it, the busiest first.",
Input: map[string]any{"subject": str("part of a subject, e.g. mesh.seat.anthropic-licence-manager or $SRV"), "limit": num("at most this many (default 100)")},
Run: func(a map[string]any) (any, error) {
subsz, err := get("/subsz?subs=1&limit=100000")
if err != nil {
return nil, err
}
connz, _ := get("/connz?auth=1&limit=1024")
return Subscriptions(subsz, connz, arg(a, "subject"), limitOf(a, 100))
}},
{Name: "nats_streams",
Description: "JetStream: every stream with its subjects, retention, message count and size — or, given one stream, every consumer on it with what it has pending, unacknowledged and redelivered.",
Input: map[string]any{"stream": str("one stream, e.g. EVENTS, CONTROL, ASSIGNMENTS, SEAT_NODE_BUILD_AGENT or KV_<bucket>")},
Run: func(a map[string]any) (any, error) {
doc, err := get(jsz)
if err != nil {
return nil, err
}
return Streams(doc, arg(a, "stream"))
}},
{Name: "nats_backlog",
Description: "Every consumer that is behind — messages pending or delivered and not yet acknowledged — the furthest behind first. Empty when every module has handled what it was sent.",
Run: func(map[string]any) (any, error) {
doc, err := get(jsz)
if err != nil {
return nil, err
}
return Backlog(doc)
}},
{Name: "nats_buckets",
Description: "Every module's state on the bus (novox/hq ADR 0201): each key-value bucket with the module and state it belongs to, its keys, size, history kept, last change and live watchers.",
Run: func(map[string]any) (any, error) {
doc, err := get(jsz)
if err != nil {
return nil, err
}
return Buckets(doc)
}},
{Name: "nats_users",
Description: "The bus's users as the controller composed them — never a password hash. Without a user, each with how many grants it holds; with part of a user's name, its publish and subscribe grants in full.",
Input: map[string]any{"user": str("part of a user name, e.g. g14.node-tools or controller")},
Run: func(a map[string]any) (any, error) {
c, err := conf()
if err != nil {
return nil, err
}
return Users(c, arg(a, "user")), nil
}},
{Name: "nats_user_can",
Description: "Whether a bus user may publish or subscribe to a subject, and the grant that allows it — the question behind a request that timed out because the bus refused it.",
Input: map[string]any{"user": str("the user exactly, e.g. g14.node-tools"), "action": str("publish or subscribe"), "subject": str("a concrete subject, e.g. $KV.claude-code_servers.all.x")},
Run: func(a map[string]any) (any, error) {
c, err := conf()
if err != nil {
return nil, err
}
return UserCan(c, arg(a, "user"), arg(a, "action"), arg(a, "subject"))
}},
}
}
func main() {
if err := stdio.Serve("", tools(LiveMonitor(), ReadUsers)); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
}