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.
133 lines
5.7 KiB
Go
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)
|
|
}
|
|
}
|