diff --git a/modules/nats/cmd/nats-tools/main.go b/modules/nats/cmd/nats-tools/main.go new file mode 100644 index 0000000..80e99d5 --- /dev/null +++ b/modules/nats/cmd/nats-tools/main.go @@ -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.`, a machine's runtime `.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_")}, + 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) + } +} diff --git a/modules/nats/cmd/nats-tools/monitor.go b/modules/nats/cmd/nats-tools/monitor.go new file mode 100644 index 0000000..a1b0707 --- /dev/null +++ b/modules/nats/cmd/nats-tools/monitor.go @@ -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 +} diff --git a/modules/nats/cmd/nats-tools/tools.go b/modules/nats/cmd/nats-tools/tools.go new file mode 100644 index 0000000..15f0648 --- /dev/null +++ b/modules/nats/cmd/nats-tools/tools.go @@ -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) +} diff --git a/modules/nats/cmd/nats-tools/tools_test.go b/modules/nats/cmd/nats-tools/tools_test.go new file mode 100644 index 0000000..677510b --- /dev/null +++ b/modules/nats/cmd/nats-tools/tools_test.go @@ -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) + } + } +} diff --git a/modules/nats/go.mod b/modules/nats/go.mod new file mode 100644 index 0000000..a5442db --- /dev/null +++ b/modules/nats/go.mod @@ -0,0 +1,5 @@ +module nats-tools + +go 1.25.0 + +require git.novox.be/novox/mesh-sdk/go v0.1.7 diff --git a/modules/nats/go.sum b/modules/nats/go.sum new file mode 100644 index 0000000..b474419 --- /dev/null +++ b/modules/nats/go.sum @@ -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= diff --git a/modules/nats/module.json b/modules/nats/module.json index 4da1e3a..3cafb56 100644 --- a/modules/nats/module.json +++ b/modules/nats/module.json @@ -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" + } } ] }