Files
mesh-catalog/modules/nats/cmd/nats-tools/tools.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

406 lines
14 KiB
Go

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)
}