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.
406 lines
14 KiB
Go
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)
|
|
}
|