| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | package rethinkdb |
| 4 | |
| 5 | import ( |
| 6 | "encoding/json" |
| 7 | "errors" |
| 8 | "fmt" |
| 9 | |
| 10 | "github.com/netdata/netdata/go/plugins/pkg/stm" |
| 11 | ) |
| 12 | |
| 13 | type ( |
| 14 | // https://rethinkdb.com/docs/system-stats/ |
| 15 | serverStats struct { |
| 16 | ID []string `json:"id"` |
| 17 | Server string `json:"server"` |
| 18 | QueryEngine struct { |
| 19 | ClientConnections int64 `json:"client_connections" stm:"client_connections"` |
| 20 | ClientsActive int64 `json:"clients_active" stm:"clients_active"` |
| 21 | QueriesTotal int64 `json:"queries_total" stm:"queries_total"` |
| 22 | ReadDocsTotal int64 `json:"read_docs_total" stm:"read_docs_total"` |
| 23 | WrittenDocsTotal int64 `json:"written_docs_total" stm:"written_docs_total"` |
| 24 | } `json:"query_engine" stm:""` |
| 25 | |
| 26 | Error string `json:"error"` |
| 27 | } |
| 28 | ) |
| 29 | |
| 30 | func (c *Collector) collect() (map[string]int64, error) { |
| 31 | if c.rdb == nil { |
| 32 | conn, err := c.newConn(c.Config) |
| 33 | if err != nil { |
| 34 | return nil, err |
| 35 | } |
| 36 | c.rdb = conn |
| 37 | } |
| 38 | |
| 39 | mx := make(map[string]int64) |
| 40 | |
| 41 | if err := c.collectStats(mx); err != nil { |
| 42 | return nil, err |
| 43 | } |
| 44 | |
| 45 | return mx, nil |
| 46 | } |
| 47 | |
| 48 | func (c *Collector) collectStats(mx map[string]int64) error { |
| 49 | resp, err := c.rdb.stats() |
| 50 | if err != nil { |
| 51 | return err |
| 52 | } |
| 53 | |
| 54 | if len(resp) == 0 { |
| 55 | return errors.New("empty stats response from server") |
| 56 | } |
| 57 | |
| 58 | for _, v := range []string{ |
| 59 | "cluster_servers_stats_request_success", |
| 60 | "cluster_servers_stats_request_timeout", |
| 61 | "cluster_client_connections", |
| 62 | "cluster_clients_active", |
| 63 | "cluster_queries_total", |
| 64 | "cluster_read_docs_total", |
| 65 | "cluster_written_docs_total", |
| 66 | } { |
| 67 | mx[v] = 0 |
| 68 | } |
| 69 | |
| 70 | seen := make(map[string]bool) |
| 71 | |
| 72 | for _, bs := range resp[1:] { // skip cluster |
| 73 | var srv serverStats |
| 74 | |
| 75 | if err := json.Unmarshal(bs, &srv); err != nil { |
| 76 | return fmt.Errorf("invalid stats response: failed to unmarshal server data: %v", err) |
| 77 | } |
| 78 | if len(srv.ID[0]) == 0 { |
| 79 | return errors.New("invalid stats response: empty id") |
| 80 | } |
| 81 | if srv.ID[0] != "server" { |
| 82 | continue |
| 83 | } |
| 84 | if len(srv.ID) != 2 { |
| 85 | return fmt.Errorf("invalid stats response: unexpected server id: '%v'", srv.ID) |
| 86 | } |
| 87 | |
| 88 | srvUUID := srv.ID[1] |
| 89 | |
| 90 | seen[srvUUID] = true |
| 91 | |
| 92 | if !c.seenServers[srvUUID] { |
| 93 | c.seenServers[srvUUID] = true |
| 94 | c.addServerCharts(srvUUID, srv.Server) |
| 95 | } |
| 96 | |
| 97 | px := fmt.Sprintf("server_%s_", srv.ID[1]) // uuid |
| 98 | |
| 99 | mx[px+"stats_request_status_success"] = 0 |
| 100 | mx[px+"stats_request_status_timeout"] = 0 |
| 101 | if srv.Error != "" { |
| 102 | mx["cluster_servers_stats_request_timeout"]++ |
| 103 | mx[px+"stats_request_status_timeout"] = 1 |
| 104 | continue |
| 105 | } |
| 106 | mx["cluster_servers_stats_request_success"]++ |
| 107 | mx[px+"stats_request_status_success"] = 1 |
| 108 | |
| 109 | for k, v := range stm.ToMap(srv.QueryEngine) { |
| 110 | mx["cluster_"+k] += v |
| 111 | mx[px+k] = v |
| 112 | } |
| 113 | } |
| 114 | |
| 115 | for k := range c.seenServers { |
| 116 | if !seen[k] { |
| 117 | delete(c.seenServers, k) |
| 118 | c.removeServerCharts(k) |
| 119 | } |
| 120 | } |
| 121 | |
| 122 | return nil |
| 123 | } |