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