Merge pull request 'The bus's own tools (nats): server, connections, subscriptions, streams, backlog, buckets, users, user_can' (#280) from feat/nats-tools into main
This commit is contained in:
@@ -0,0 +1,132 @@
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
package main
|
||||
|
||||
// Reaching the bus server's own monitoring API (the HTTP endpoints nats-server serves: varz, connz,
|
||||
// subsz, jsz, healthz). Read-only by the server's own design: nothing here can change the bus.
|
||||
//
|
||||
// **Two ways in, the first that answers.** The module's configuration binds the endpoint inside its
|
||||
// container; published on the machine's loopback, a host-side client reaches it directly — and where the
|
||||
// binding is the container's own loopback (as the module shipped until this bundle), only a process inside
|
||||
// the container can. So the endpoint is asked directly first, and through the container second; the day
|
||||
// the binding is moved, the second way simply stops being used.
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Monitor fetches a monitoring path, e.g. "/varz".
|
||||
type Monitor func(ctx context.Context, path string) ([]byte, error)
|
||||
|
||||
func env(key, fallback string) string {
|
||||
if v := strings.TrimSpace(os.Getenv(key)); v != "" {
|
||||
return v
|
||||
}
|
||||
return fallback
|
||||
}
|
||||
|
||||
// LiveMonitor is the server as this machine reaches it.
|
||||
func LiveMonitor() Monitor {
|
||||
base := env("MESH_NATS_MONITOR", "http://127.0.0.1:8222")
|
||||
container := env("MESH_NATS_CONTAINER", "mesh-broker-nats")
|
||||
client := &http.Client{Timeout: 5 * time.Second}
|
||||
return func(ctx context.Context, path string) ([]byte, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, base+path, nil)
|
||||
if err == nil {
|
||||
if resp, err := client.Do(req); err == nil {
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode == http.StatusOK {
|
||||
return io.ReadAll(io.LimitReader(resp.Body, 64<<20))
|
||||
}
|
||||
}
|
||||
}
|
||||
// Through the container: its own loopback is where the endpoint listens.
|
||||
out, err := exec.CommandContext(ctx, "docker", "exec", container, "wget", "-qO-", "http://127.0.0.1:8222"+path).Output()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("the bus's monitoring endpoint answered neither at %s nor inside %s: %w", base, container, err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
}
|
||||
|
||||
// ReadUsers reads the composed user list the server includes: in the container, where it is mounted
|
||||
// read-only; the machine's copy is readable by root alone.
|
||||
func ReadUsers(ctx context.Context) ([]byte, error) {
|
||||
container := env("MESH_NATS_CONTAINER", "mesh-broker-nats")
|
||||
path := env("MESH_NATS_USERS_FILE", "/etc/nats/accounts.conf")
|
||||
out, err := exec.CommandContext(ctx, "docker", "exec", container, "cat", path).Output()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("the bus's user list could not be read inside %s: %w", container, err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
@@ -0,0 +1,405 @@
|
||||
package main
|
||||
|
||||
// What the bus's tools answer, shaped from the server's own monitoring documents. Pure: each takes the
|
||||
// raw document and the caller's arguments, so every answer is tested against a captured document.
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// ---- the server ----------------------------------------------------------------------------------
|
||||
|
||||
// Server is what a person asks first: which server, how loaded, what it allows.
|
||||
func Server(varz, healthz []byte) (map[string]any, error) {
|
||||
var v map[string]any
|
||||
if err := json.Unmarshal(varz, &v); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := map[string]any{}
|
||||
for _, k := range []string{"version", "go", "start", "uptime", "max_payload", "max_connections", "max_pending",
|
||||
"connections", "total_connections", "subscriptions", "slow_consumers", "in_msgs", "out_msgs", "in_bytes",
|
||||
"out_bytes", "mem", "cpu", "cores"} {
|
||||
if val, ok := v[k]; ok {
|
||||
out[k] = val
|
||||
}
|
||||
}
|
||||
if js, ok := v["jetstream"].(map[string]any); ok {
|
||||
if stats, ok := js["stats"].(map[string]any); ok {
|
||||
out["jetstream"] = stats
|
||||
}
|
||||
}
|
||||
var h map[string]any
|
||||
if json.Unmarshal(healthz, &h) == nil {
|
||||
out["health"] = h
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// ---- connections ---------------------------------------------------------------------------------
|
||||
|
||||
// Connection is one client, named by the mesh user it authenticated as.
|
||||
type Connection struct {
|
||||
CID int64 `json:"cid"`
|
||||
User string `json:"user"`
|
||||
Name string `json:"name"`
|
||||
IP string `json:"ip"`
|
||||
Uptime string `json:"uptime"`
|
||||
Idle string `json:"idle"`
|
||||
RTT string `json:"rtt"`
|
||||
InMsgs int64 `json:"in_msgs"`
|
||||
OutMsgs int64 `json:"out_msgs"`
|
||||
InBytes int64 `json:"in_bytes"`
|
||||
OutBytes int64 `json:"out_bytes"`
|
||||
PendingBytes int64 `json:"pending_bytes"`
|
||||
Subscriptions int64 `json:"subscriptions"`
|
||||
}
|
||||
|
||||
func connectionsOf(connz []byte) ([]Connection, error) {
|
||||
var doc struct {
|
||||
Connections []struct {
|
||||
Connection
|
||||
AuthorizedUser string `json:"authorized_user"`
|
||||
} `json:"connections"`
|
||||
}
|
||||
if err := json.Unmarshal(connz, &doc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make([]Connection, 0, len(doc.Connections))
|
||||
for _, c := range doc.Connections {
|
||||
conn := c.Connection
|
||||
conn.User = c.AuthorizedUser
|
||||
out = append(out, conn)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// Connections is every client, narrowed to users containing `user`, sorted by `by` (out_bytes,
|
||||
// in_bytes, pending, msgs, subscriptions), at most `limit`.
|
||||
func Connections(connz []byte, user, by string, limit int) (map[string]any, error) {
|
||||
all, err := connectionsOf(connz)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out []Connection
|
||||
for _, c := range all {
|
||||
if user == "" || strings.Contains(c.User, user) || strings.Contains(c.Name, user) {
|
||||
out = append(out, c)
|
||||
}
|
||||
}
|
||||
key := map[string]func(Connection) int64{
|
||||
"in_bytes": func(c Connection) int64 { return c.InBytes }, "pending": func(c Connection) int64 { return c.PendingBytes },
|
||||
"msgs": func(c Connection) int64 { return c.InMsgs + c.OutMsgs },
|
||||
"subscriptions": func(c Connection) int64 { return c.Subscriptions },
|
||||
}[by]
|
||||
if key == nil {
|
||||
key = func(c Connection) int64 { return c.OutBytes }
|
||||
}
|
||||
sort.SliceStable(out, func(i, j int) bool { return key(out[i]) > key(out[j]) })
|
||||
total := len(out)
|
||||
if limit > 0 && len(out) > limit {
|
||||
out = out[:limit]
|
||||
}
|
||||
return map[string]any{"total": total, "connections": out}, nil
|
||||
}
|
||||
|
||||
// ---- subscriptions -------------------------------------------------------------------------------
|
||||
|
||||
// Subscriptions is who listens on what: subjects containing `subject`, with the user listening when
|
||||
// connz names it, the busiest first.
|
||||
func Subscriptions(subsz, connz []byte, subject string, limit int) (map[string]any, error) {
|
||||
var doc struct {
|
||||
Total int `json:"num_subscriptions"`
|
||||
List []struct {
|
||||
Subject string `json:"subject"`
|
||||
Queue string `json:"qgroup"`
|
||||
Msgs int64 `json:"msgs"`
|
||||
CID int64 `json:"cid"`
|
||||
} `json:"subscriptions_list"`
|
||||
}
|
||||
if err := json.Unmarshal(subsz, &doc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
users := map[int64]string{}
|
||||
if conns, err := connectionsOf(connz); err == nil {
|
||||
for _, c := range conns {
|
||||
users[c.CID] = c.User
|
||||
}
|
||||
}
|
||||
type sub struct {
|
||||
Subject string `json:"subject"`
|
||||
Queue string `json:"queue,omitempty"`
|
||||
Msgs int64 `json:"msgs"`
|
||||
User string `json:"user,omitempty"`
|
||||
}
|
||||
var out []sub
|
||||
for _, s := range doc.List {
|
||||
if subject != "" && !strings.Contains(s.Subject, subject) {
|
||||
continue
|
||||
}
|
||||
out = append(out, sub{Subject: s.Subject, Queue: s.Queue, Msgs: s.Msgs, User: users[s.CID]})
|
||||
}
|
||||
sort.SliceStable(out, func(i, j int) bool { return out[i].Msgs > out[j].Msgs })
|
||||
matched := len(out)
|
||||
if limit > 0 && len(out) > limit {
|
||||
out = out[:limit]
|
||||
}
|
||||
return map[string]any{"subscriptions_on_server": doc.Total, "matched": matched, "subscriptions": out}, nil
|
||||
}
|
||||
|
||||
// ---- JetStream -----------------------------------------------------------------------------------
|
||||
|
||||
type jszDoc struct {
|
||||
Streams int `json:"streams"`
|
||||
Consumers int `json:"consumers"`
|
||||
Messages int64 `json:"messages"`
|
||||
Bytes int64 `json:"bytes"`
|
||||
Accounts []struct {
|
||||
Streams []struct {
|
||||
Name string `json:"name"`
|
||||
Config struct {
|
||||
Subjects []string `json:"subjects"`
|
||||
Retention string `json:"retention"`
|
||||
MaxAge int64 `json:"max_age"`
|
||||
MaxMsgsPerSubject int64 `json:"max_msgs_per_subject"`
|
||||
MaxBytes int64 `json:"max_bytes"`
|
||||
Description string `json:"description"`
|
||||
} `json:"config"`
|
||||
State struct {
|
||||
Messages int64 `json:"messages"`
|
||||
Bytes int64 `json:"bytes"`
|
||||
FirstSeq int64 `json:"first_seq"`
|
||||
LastSeq int64 `json:"last_seq"`
|
||||
LastTS string `json:"last_ts"`
|
||||
NumSubjects int64 `json:"num_subjects"`
|
||||
ConsumerCount int64 `json:"consumer_count"`
|
||||
} `json:"state"`
|
||||
Consumers []struct {
|
||||
Name string `json:"name"`
|
||||
Config struct {
|
||||
FilterSubject string `json:"filter_subject"`
|
||||
FilterSubjects []string `json:"filter_subjects"`
|
||||
DeliverSubject string `json:"deliver_subject"`
|
||||
MaxDeliver int64 `json:"max_deliver"`
|
||||
} `json:"config"`
|
||||
NumPending int64 `json:"num_pending"`
|
||||
NumAckPending int64 `json:"num_ack_pending"`
|
||||
NumRedelivered int64 `json:"num_redelivered"`
|
||||
NumWaiting int64 `json:"num_waiting"`
|
||||
Delivered struct {
|
||||
LastActive string `json:"last_active"`
|
||||
} `json:"delivered"`
|
||||
} `json:"consumer_detail"`
|
||||
} `json:"stream_detail"`
|
||||
} `json:"account_details"`
|
||||
}
|
||||
|
||||
// Streams is every stream with its retention and size — or one stream with every consumer's backlog.
|
||||
func Streams(jsz []byte, stream string) (map[string]any, error) {
|
||||
var doc jszDoc
|
||||
if err := json.Unmarshal(jsz, &doc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var rows []map[string]any
|
||||
for _, a := range doc.Accounts {
|
||||
for _, s := range a.Streams {
|
||||
if stream != "" && s.Name != stream {
|
||||
continue
|
||||
}
|
||||
row := map[string]any{"name": s.Name, "subjects": s.Config.Subjects, "retention": s.Config.Retention,
|
||||
"messages": s.State.Messages, "bytes": s.State.Bytes, "subjects_held": s.State.NumSubjects,
|
||||
"consumers": s.State.ConsumerCount, "last": s.State.LastTS}
|
||||
if s.Config.MaxAge > 0 {
|
||||
row["max_age_hours"] = s.Config.MaxAge / 3_600_000_000_000
|
||||
}
|
||||
if s.Config.MaxMsgsPerSubject > 0 {
|
||||
row["max_per_subject"] = s.Config.MaxMsgsPerSubject
|
||||
}
|
||||
if stream != "" {
|
||||
row["description"] = s.Config.Description
|
||||
var consumers []map[string]any
|
||||
for _, c := range s.Consumers {
|
||||
filters := c.Config.FilterSubjects
|
||||
if c.Config.FilterSubject != "" {
|
||||
filters = append(filters, c.Config.FilterSubject)
|
||||
}
|
||||
kind := "pull"
|
||||
if c.Config.DeliverSubject != "" {
|
||||
kind = "push"
|
||||
}
|
||||
consumers = append(consumers, map[string]any{"name": c.Name, "kind": kind, "filters": filters,
|
||||
"pending": c.NumPending, "unacknowledged": c.NumAckPending, "redelivered": c.NumRedelivered,
|
||||
"waiting_pulls": c.NumWaiting, "last_delivered": c.Delivered.LastActive, "max_deliver": c.Config.MaxDeliver})
|
||||
}
|
||||
sort.Slice(consumers, func(i, j int) bool { return consumers[i]["name"].(string) < consumers[j]["name"].(string) })
|
||||
row["consumer_detail"] = consumers
|
||||
}
|
||||
rows = append(rows, row)
|
||||
}
|
||||
}
|
||||
sort.Slice(rows, func(i, j int) bool { return rows[i]["name"].(string) < rows[j]["name"].(string) })
|
||||
if stream != "" && len(rows) == 0 {
|
||||
return nil, fmt.Errorf("no stream called %s", stream)
|
||||
}
|
||||
return map[string]any{"streams": doc.Streams, "consumers": doc.Consumers, "messages": doc.Messages,
|
||||
"bytes": doc.Bytes, "detail": rows}, nil
|
||||
}
|
||||
|
||||
// Backlog is every consumer with work it has not finished — pending or unacknowledged — the busiest first.
|
||||
func Backlog(jsz []byte) (map[string]any, error) {
|
||||
var doc jszDoc
|
||||
if err := json.Unmarshal(jsz, &doc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var rows []map[string]any
|
||||
for _, a := range doc.Accounts {
|
||||
for _, s := range a.Streams {
|
||||
for _, c := range s.Consumers {
|
||||
if c.NumPending == 0 && c.NumAckPending == 0 {
|
||||
continue
|
||||
}
|
||||
rows = append(rows, map[string]any{"stream": s.Name, "consumer": c.Name, "pending": c.NumPending,
|
||||
"unacknowledged": c.NumAckPending, "redelivered": c.NumRedelivered, "last_delivered": c.Delivered.LastActive})
|
||||
}
|
||||
}
|
||||
}
|
||||
sort.Slice(rows, func(i, j int) bool {
|
||||
return rows[i]["pending"].(int64)+rows[i]["unacknowledged"].(int64) > rows[j]["pending"].(int64)+rows[j]["unacknowledged"].(int64)
|
||||
})
|
||||
return map[string]any{"consumers_behind": len(rows), "backlog": rows}, nil
|
||||
}
|
||||
|
||||
// Buckets is every module's state (novox/hq ADR 0201): the key-value buckets, with their keys and size.
|
||||
func Buckets(jsz []byte) (map[string]any, error) {
|
||||
var doc jszDoc
|
||||
if err := json.Unmarshal(jsz, &doc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var rows []map[string]any
|
||||
for _, a := range doc.Accounts {
|
||||
for _, s := range a.Streams {
|
||||
name, ok := strings.CutPrefix(s.Name, "KV_")
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
module, state, _ := strings.Cut(name, "_")
|
||||
rows = append(rows, map[string]any{"bucket": name, "module": module, "state": state,
|
||||
"keys_including_deleted": s.State.NumSubjects, "values_kept": s.State.Messages, "bytes": s.State.Bytes,
|
||||
"history": s.Config.MaxMsgsPerSubject, "last_change": s.State.LastTS, "watchers": s.State.ConsumerCount})
|
||||
}
|
||||
}
|
||||
sort.Slice(rows, func(i, j int) bool { return rows[i]["bucket"].(string) < rows[j]["bucket"].(string) })
|
||||
return map[string]any{"buckets": rows}, nil
|
||||
}
|
||||
|
||||
// ---- users and grants ----------------------------------------------------------------------------
|
||||
|
||||
// User is one bus user's grants, as the controller composed them — never its password hash.
|
||||
type User struct {
|
||||
User string `json:"user"`
|
||||
Publish []string `json:"publish"`
|
||||
Subscribe []string `json:"subscribe"`
|
||||
AllowResponses bool `json:"allow_responses"`
|
||||
}
|
||||
|
||||
var (
|
||||
userBlock = regexp.MustCompile(`\{\s*user:\s*"([^"]+)"`)
|
||||
allowList = regexp.MustCompile(`(publish|subscribe):\s*\{\s*allow:\s*\[([^\]]*)\]`)
|
||||
quoted = regexp.MustCompile(`"((?:[^"\\]|\\.)*)"`)
|
||||
)
|
||||
|
||||
// ParseUsers reads the composed user list (the controller's broker.ComposeAccounts shape).
|
||||
func ParseUsers(conf []byte) []User {
|
||||
text := string(conf)
|
||||
idx := userBlock.FindAllStringSubmatchIndex(text, -1)
|
||||
var out []User
|
||||
for i, m := range idx {
|
||||
end := len(text)
|
||||
if i+1 < len(idx) {
|
||||
end = idx[i+1][0]
|
||||
}
|
||||
block := text[m[0]:end]
|
||||
u := User{User: text[m[2]:m[3]], Publish: []string{}, Subscribe: []string{}, AllowResponses: strings.Contains(block, "allow_responses")}
|
||||
for _, a := range allowList.FindAllStringSubmatch(block, -1) {
|
||||
var subjects []string
|
||||
for _, q := range quoted.FindAllStringSubmatch(a[2], -1) {
|
||||
subjects = append(subjects, q[1])
|
||||
}
|
||||
if a[1] == "publish" {
|
||||
u.Publish = subjects
|
||||
} else {
|
||||
u.Subscribe = subjects
|
||||
}
|
||||
}
|
||||
out = append(out, u)
|
||||
}
|
||||
sort.Slice(out, func(i, j int) bool { return out[i].User < out[j].User })
|
||||
return out
|
||||
}
|
||||
|
||||
// Users is every user's grants, or those of users containing `user`; with `grants`, the subjects too.
|
||||
func Users(conf []byte, user string) map[string]any {
|
||||
all := ParseUsers(conf)
|
||||
if user == "" {
|
||||
type summary struct {
|
||||
User string `json:"user"`
|
||||
Publish int `json:"publish_grants"`
|
||||
Subscribe int `json:"subscribe_grants"`
|
||||
}
|
||||
var out []summary
|
||||
for _, u := range all {
|
||||
out = append(out, summary{u.User, len(u.Publish), len(u.Subscribe)})
|
||||
}
|
||||
return map[string]any{"users": out}
|
||||
}
|
||||
var out []User
|
||||
for _, u := range all {
|
||||
if strings.Contains(u.User, user) {
|
||||
out = append(out, u)
|
||||
}
|
||||
}
|
||||
return map[string]any{"users": out}
|
||||
}
|
||||
|
||||
// SubjectMatches says whether a concrete subject falls under a grant's pattern: `*` one token, `>` the rest.
|
||||
func SubjectMatches(pattern, subject string) bool {
|
||||
p, s := strings.Split(pattern, "."), strings.Split(subject, ".")
|
||||
for i, tok := range p {
|
||||
if tok == ">" {
|
||||
return len(s) > i
|
||||
}
|
||||
if i >= len(s) || (tok != "*" && tok != s[i]) {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return len(p) == len(s)
|
||||
}
|
||||
|
||||
// UserCan answers whether a user may publish or subscribe to a subject, and which grant allows it.
|
||||
func UserCan(conf []byte, user, action, subject string) (map[string]any, error) {
|
||||
for _, u := range ParseUsers(conf) {
|
||||
if u.User != user {
|
||||
continue
|
||||
}
|
||||
grants := u.Publish
|
||||
if action == "subscribe" {
|
||||
grants = u.Subscribe
|
||||
} else if action != "publish" {
|
||||
return nil, fmt.Errorf("action is publish or subscribe, not %q", action)
|
||||
}
|
||||
for _, g := range grants {
|
||||
if SubjectMatches(g, subject) {
|
||||
return map[string]any{"user": user, "action": action, "subject": subject, "allowed": true, "by": g}, nil
|
||||
}
|
||||
}
|
||||
answer := map[string]any{"user": user, "action": action, "subject": subject, "allowed": false}
|
||||
if action == "publish" && u.AllowResponses {
|
||||
answer["note"] = "a reply to a request this user received is allowed regardless (allow_responses)"
|
||||
}
|
||||
return answer, nil
|
||||
}
|
||||
return nil, fmt.Errorf("the bus has no user %q; nats_users lists them", user)
|
||||
}
|
||||
@@ -0,0 +1,127 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
const connz = `{"connections":[
|
||||
{"cid":1,"ip":"192.0.2.1","name":"mesh-host/anchor","authorized_user":"node.anchor","out_bytes":10,"in_bytes":5,"pending_bytes":0,"subscriptions":2,"in_msgs":3,"out_msgs":4},
|
||||
{"cid":2,"ip":"192.0.2.2","name":"laptop.node-tools","authorized_user":"laptop.node-tools","out_bytes":900,"in_bytes":50,"pending_bytes":7,"subscriptions":300,"in_msgs":30,"out_msgs":40}]}`
|
||||
|
||||
const subsz = `{"num_subscriptions":3,"subscriptions_list":[
|
||||
{"subject":"$SRV.INFO","msgs":4,"cid":2},
|
||||
{"subject":"mesh.seat.x.tool.current","qgroup":"","msgs":9,"cid":2},
|
||||
{"subject":"mesh.node.anchor.declare","msgs":1,"cid":1}]}`
|
||||
|
||||
const jsz = `{"streams":2,"consumers":2,"messages":12,"bytes":300,"account_details":[{"stream_detail":[
|
||||
{"name":"EVENTS","config":{"subjects":["mesh.mod.*.event.>"],"retention":"limits","max_age":604800000000000,"max_msgs_per_subject":10000},
|
||||
"state":{"messages":10,"bytes":200,"num_subjects":3,"consumer_count":2,"last_ts":"2026-10-04T10:00:00Z"},
|
||||
"consumer_detail":[
|
||||
{"name":"anchor_audit","config":{"filter_subject":"mesh.mod.*.event.>"},"num_pending":5,"num_ack_pending":1,"num_redelivered":2,"delivered":{"last_active":"2026-10-04T10:00:00Z"}},
|
||||
{"name":"anchor_quiet","config":{"filter_subjects":["mesh.mod.a.event.b"]},"num_pending":0,"num_ack_pending":0}]},
|
||||
{"name":"KV_claude-code_servers","config":{"subjects":["$KV.claude-code_servers.>"],"retention":"limits","max_msgs_per_subject":1},
|
||||
"state":{"messages":2,"bytes":100,"num_subjects":2,"consumer_count":1,"last_ts":"2026-10-04T11:00:00Z"}}]}]}`
|
||||
|
||||
const users = `accounts {
|
||||
MESH {
|
||||
jetstream: enabled
|
||||
users = [
|
||||
{ user: "controller", password: "$2a$10$secret-hash-one", permissions: {
|
||||
publish: { allow: ["mesh.control.>", "$JS.API.>"] }
|
||||
subscribe: { allow: ["mesh.control.>", "_INBOX.controller.>"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
{ user: "laptop.node-tools", password: "$2a$10$secret-hash-two", permissions: {
|
||||
publish: { allow: ["$KV.claude-code_servers.>", "mesh.mod.*.tool.>"] }
|
||||
subscribe: { allow: ["mesh.mod.claude-code.tool.>"] }
|
||||
} }
|
||||
]
|
||||
}
|
||||
}`
|
||||
|
||||
func TestConnectionsAreNamedByTheirUserAndSorted(t *testing.T) {
|
||||
out, err := Connections([]byte(connz), "", "", 10)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
cs := out["connections"].([]Connection)
|
||||
if len(cs) != 2 || cs[0].User != "laptop.node-tools" || cs[1].User != "node.anchor" {
|
||||
t.Fatalf("%+v", cs)
|
||||
}
|
||||
narrowed, _ := Connections([]byte(connz), "anchor", "pending", 10)
|
||||
if n := narrowed["total"].(int); n != 1 {
|
||||
t.Fatalf("narrowed to %d", n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSubscriptionsNameTheirUser(t *testing.T) {
|
||||
out, err := Subscriptions([]byte(subsz), []byte(connz), "SRV", 10)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
raw, _ := json.Marshal(out)
|
||||
if !strings.Contains(string(raw), `"user":"laptop.node-tools"`) || out["matched"].(int) != 1 {
|
||||
t.Fatalf("%s", raw)
|
||||
}
|
||||
}
|
||||
|
||||
func TestStreamsBacklogAndBuckets(t *testing.T) {
|
||||
all, err := Streams([]byte(jsz), "")
|
||||
if err != nil || len(all["detail"].([]map[string]any)) != 2 {
|
||||
t.Fatalf("%v %v", all, err)
|
||||
}
|
||||
one, _ := Streams([]byte(jsz), "EVENTS")
|
||||
row := one["detail"].([]map[string]any)[0]
|
||||
if row["max_age_hours"].(int64) != 168 || len(row["consumer_detail"].([]map[string]any)) != 2 {
|
||||
t.Fatalf("%v", row)
|
||||
}
|
||||
if _, err := Streams([]byte(jsz), "NOPE"); err == nil {
|
||||
t.Fatal("an unknown stream answered")
|
||||
}
|
||||
b, _ := Backlog([]byte(jsz))
|
||||
if b["consumers_behind"].(int) != 1 {
|
||||
t.Fatalf("%v", b)
|
||||
}
|
||||
k, _ := Buckets([]byte(jsz))
|
||||
rows := k["buckets"].([]map[string]any)
|
||||
if len(rows) != 1 || rows[0]["module"] != "claude-code" || rows[0]["state"] != "servers" {
|
||||
t.Fatalf("%v", rows)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUsersNeverCarryAPasswordHash(t *testing.T) {
|
||||
for _, out := range []map[string]any{Users([]byte(users), ""), Users([]byte(users), "laptop")} {
|
||||
raw, _ := json.Marshal(out)
|
||||
if strings.Contains(string(raw), "secret-hash") || strings.Contains(string(raw), "$2a$") {
|
||||
t.Fatalf("a hash in %s", raw)
|
||||
}
|
||||
}
|
||||
us := ParseUsers([]byte(users))
|
||||
if len(us) != 2 || len(us[0].Publish) != 2 || !us[0].AllowResponses || us[1].AllowResponses {
|
||||
t.Fatalf("%+v", us)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWhetherAUserMayDoSomething(t *testing.T) {
|
||||
ok, err := UserCan([]byte(users), "laptop.node-tools", "publish", "$KV.claude-code_servers.all.x")
|
||||
if err != nil || ok["allowed"] != true || ok["by"] != "$KV.claude-code_servers.>" {
|
||||
t.Fatalf("%v %v", ok, err)
|
||||
}
|
||||
no, _ := UserCan([]byte(users), "laptop.node-tools", "publish", "$KV.claude-code_holdings.laptop")
|
||||
if no["allowed"] != false {
|
||||
t.Fatalf("%v", no)
|
||||
}
|
||||
if _, err := UserCan([]byte(users), "nobody", "publish", "x"); err == nil {
|
||||
t.Fatal("an unknown user answered")
|
||||
}
|
||||
for _, c := range []struct {
|
||||
p, s string
|
||||
want bool
|
||||
}{{"a.*.c", "a.b.c", true}, {"a.>", "a.b.c", true}, {"a.>", "a", false}, {"a.b", "a.b.c", false}, {"a.*", "a.b.c", false}} {
|
||||
if SubjectMatches(c.p, c.s) != c.want {
|
||||
t.Errorf("%s ~ %s", c.p, c.s)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
module nats-tools
|
||||
|
||||
go 1.25.0
|
||||
|
||||
require git.novox.be/novox/mesh-sdk/go v0.1.7
|
||||
@@ -0,0 +1,2 @@
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.7 h1:C0sTQmtTiyYH7bnqZb7PusXnqA37gKuT7Nqjn9gG47w=
|
||||
git.novox.be/novox/mesh-sdk/go v0.1.7/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY=
|
||||
@@ -25,9 +25,19 @@
|
||||
"port": 4222,
|
||||
"protocol": "tcp",
|
||||
"from": "mesh",
|
||||
"why": "the mesh bus \u2014 every link the mesh has, over TLS, reached across the overlay"
|
||||
"why": "the mesh bus — every link the mesh has, over TLS, reached across the overlay"
|
||||
}
|
||||
],
|
||||
"tools": [
|
||||
"nats_server",
|
||||
"nats_connections",
|
||||
"nats_subscriptions",
|
||||
"nats_streams",
|
||||
"nats_backlog",
|
||||
"nats_buckets",
|
||||
"nats_users",
|
||||
"nats_user_can"
|
||||
],
|
||||
"guards": [
|
||||
8222
|
||||
],
|
||||
@@ -48,7 +58,7 @@
|
||||
"id": "server-conf",
|
||||
"type": "file",
|
||||
"path": "/var/lib/nats-module/conf/nats.conf",
|
||||
"content": "# The nats module's own server settings. Declared by the module, because a port, a TLS path\n# and a store directory are properties of the container this module raises: they live in its\n# image and its mounts and change when it does.\n#\n# The mesh writes accounts.conf beside this one and nothing else. A controller that wrote the\n# whole file would have to be kept in step with a Dockerfile it never sees.\n\nport: 4222\nhttp: 127.0.0.1:8222\n\n# The mesh's own broker certificate \u2014 the one every machine already pins by fingerprint and the\n# controller already trusts (MESH_BROKER_CERTIFICATE). Serving the new bus with it means no\n# machine's pin changes when it moves, and no second certificate exists to be wrong about.\ntls {\n cert_file: \"/tls/tls.crt\"\n key_file: \"/tls/tls.key\"\n}\n\n# **No `verify`, deliberately, and it was `verify: true` until a probe ran this image.** That\n# setting makes the server demand a *client* certificate, and nothing in the mesh presents one: a\n# host pins this server's exact certificate and authenticates with the password the mesh minted\n# (novox/hq ADR 0004, design 25 \u00a74), and so does a module's runtime. With it on, every connection\n# in the mesh is refused at the TLS handshake, before any password is looked at \u2014 and the error is\n# \"client didn't provide a certificate\", which reads as a client fault.\n#\n# TLS is still required: a tls block is what makes it required, and verify only decides whether\n# client certificates are checked. What is given up is a second factor the mesh has no machinery\n# to issue or rotate \u2014 a certificate per module per node \u2014 and what is kept is stronger than a\n# name check in both directions: an exact pin outward, a per-user password inward.\n\njetstream {\n store_dir: \"/data\"\n}\n\n# Every user of the mesh, composed by the controller and rewritten whenever a module is\n# assigned, a node enrols or a person's access changes.\n#\n# **Relative, and in this same directory, because it has to be.** An absolute include path is\n# resolved relative to the including file's directory, not from the root: nats-server given\n# `include /etc/nats/accounts.conf` from /etc/nats-server/nats.conf looks for\n# /etc/nats-server/etc/nats/accounts.conf and refuses to start. Verified against the server.\ninclude accounts.conf\n",
|
||||
"content": "# The nats module's own server settings. Declared by the module, because a port, a TLS path\n# and a store directory are properties of the container this module raises: they live in its\n# image and its mounts and change when it does.\n#\n# The mesh writes accounts.conf beside this one and nothing else. A controller that wrote the\n# whole file would have to be kept in step with a Dockerfile it never sees.\n\nport: 4222\nhttp: 127.0.0.1:8222\n\n# The mesh's own broker certificate — the one every machine already pins by fingerprint and the\n# controller already trusts (MESH_BROKER_CERTIFICATE). Serving the new bus with it means no\n# machine's pin changes when it moves, and no second certificate exists to be wrong about.\ntls {\n cert_file: \"/tls/tls.crt\"\n key_file: \"/tls/tls.key\"\n}\n\n# **No `verify`, deliberately, and it was `verify: true` until a probe ran this image.** That\n# setting makes the server demand a *client* certificate, and nothing in the mesh presents one: a\n# host pins this server's exact certificate and authenticates with the password the mesh minted\n# (novox/hq ADR 0004, design 25 §4), and so does a module's runtime. With it on, every connection\n# in the mesh is refused at the TLS handshake, before any password is looked at — and the error is\n# \"client didn't provide a certificate\", which reads as a client fault.\n#\n# TLS is still required: a tls block is what makes it required, and verify only decides whether\n# client certificates are checked. What is given up is a second factor the mesh has no machinery\n# to issue or rotate — a certificate per module per node — and what is kept is stronger than a\n# name check in both directions: an exact pin outward, a per-user password inward.\n\njetstream {\n store_dir: \"/data\"\n}\n\n# Every user of the mesh, composed by the controller and rewritten whenever a module is\n# assigned, a node enrols or a person's access changes.\n#\n# **Relative, and in this same directory, because it has to be.** An absolute include path is\n# resolved relative to the including file's directory, not from the root: nats-server given\n# `include /etc/nats/accounts.conf` from /etc/nats-server/nats.conf looks for\n# /etc/nats-server/etc/nats/accounts.conf and refuses to start. Verified against the server.\ninclude accounts.conf\n",
|
||||
"mode": "0644"
|
||||
},
|
||||
{
|
||||
@@ -86,6 +96,22 @@
|
||||
"name": "server",
|
||||
"kind": "image",
|
||||
"from": "Dockerfile"
|
||||
},
|
||||
{
|
||||
"name": "tools",
|
||||
"kind": "bundle",
|
||||
"language": "go",
|
||||
"system": "arch",
|
||||
"from": "cmd/nats-tools",
|
||||
"binary": "nats-tools",
|
||||
"loads": [
|
||||
"nats-tools"
|
||||
],
|
||||
"env": {
|
||||
"MESH_NATS_MONITOR": "http://127.0.0.1:8222",
|
||||
"MESH_NATS_CONTAINER": "mesh-broker-nats",
|
||||
"MESH_NATS_USERS_FILE": "/etc/nats/accounts.conf"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user