master
go 280 lines 9.9 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package rethinkdb
4
5 import (
6 "context"
7 "fmt"
8 "sort"
9 "strings"
10
11 "github.com/netdata/netdata/go/plugins/pkg/funcapi"
12 "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/strmutil"
13 )
14
15 const (
16 runningQueriesMethodID = "running-queries"
17 rethinkMaxQueryTextLength = 4096
18 )
19
20 func runningQueriesMethodConfig() funcapi.MethodConfig {
21 return funcapi.MethodConfig{
22 ID: runningQueriesMethodID,
23 Name: "Running Queries",
24 UpdateEvery: 10,
25 Help: "Currently running queries from rethinkdb.jobs. WARNING: Query text may contain unmasked literals (potential PII).",
26 RequireCloud: true,
27 RequiredParams: []funcapi.ParamConfig{funcapi.BuildSortParam(rethinkRunningColumns)},
28 }
29 }
30
31 // Compile-time interface check.
32 var _ funcapi.MethodHandler = (*funcRunningQueries)(nil)
33
34 // funcRunningQueries handles the "running-queries" function for RethinkDB.
35 type funcRunningQueries struct {
36 router *funcRouter
37 }
38
39 func newFuncRunningQueries(r *funcRouter) *funcRunningQueries {
40 return &funcRunningQueries{router: r}
41 }
42
43 // MethodParams implements funcapi.MethodHandler.
44 func (f *funcRunningQueries) Cleanup(ctx context.Context) {}
45
46 func (f *funcRunningQueries) MethodParams(_ context.Context, method string) ([]funcapi.ParamConfig, error) {
47 if method != runningQueriesMethodID {
48 return nil, fmt.Errorf("unknown method: %s", method)
49 }
50 if f.router.collector.Functions.RunningQueries.Disabled {
51 return nil, fmt.Errorf("running-queries function disabled in configuration")
52 }
53 return []funcapi.ParamConfig{funcapi.BuildSortParam(rethinkRunningColumns)}, nil
54 }
55
56 // Handle implements funcapi.MethodHandler.
57 func (f *funcRunningQueries) Handle(ctx context.Context, method string, params funcapi.ResolvedParams) *funcapi.FunctionResponse {
58 if method != runningQueriesMethodID {
59 return funcapi.NotFoundResponse(method)
60 }
61 if f.router.collector.Functions.RunningQueries.Disabled {
62 return funcapi.UnavailableResponse("running-queries function has been disabled in configuration")
63 }
64
65 return f.collectRunningQueries(ctx, params.Column("__sort"))
66 }
67
68 func (f *funcRunningQueries) collectRunningQueries(ctx context.Context, sortColumn string) *funcapi.FunctionResponse {
69 c := f.router.collector
70
71 limit := c.runningQueriesLimit()
72 timeout := c.runningQueriesTimeout()
73
74 queryCtx, cancel := context.WithTimeout(ctx, timeout)
75 defer cancel()
76
77 if queryCtx.Err() == context.DeadlineExceeded {
78 return &funcapi.FunctionResponse{Status: 504, Message: "query timed out"}
79 }
80
81 rows, err := c.rdb.jobs(queryCtx)
82 if err != nil {
83 if queryCtx.Err() == context.DeadlineExceeded {
84 return &funcapi.FunctionResponse{Status: 504, Message: "query timed out"}
85 }
86 return &funcapi.FunctionResponse{Status: 500, Message: fmt.Sprintf("jobs query failed: %v", err)}
87 }
88
89 jobRows := make([]rethinkJobRow, 0, len(rows))
90 for _, row := range rows {
91 jobRows = append(jobRows, parseRethinkJob(row))
92 }
93
94 cs := rethinkColumnSet(rethinkRunningColumns)
95
96 if len(jobRows) == 0 {
97 return &funcapi.FunctionResponse{
98 Status: 200,
99 Message: "No running queries found.",
100 Help: "Currently running queries from rethinkdb.jobs",
101 Columns: cs.BuildColumns(),
102 Data: [][]any{},
103 DefaultSortColumn: mapRethinkSortColumn(sortColumn),
104 RequiredParams: []funcapi.ParamConfig{funcapi.BuildSortParam(rethinkRunningColumns)},
105 }
106 }
107
108 sortColumn = mapRethinkSortColumn(sortColumn)
109 sortRethinkRows(jobRows)
110 if len(jobRows) > limit {
111 jobRows = jobRows[:limit]
112 }
113
114 data := make([][]any, 0, len(jobRows))
115 for _, row := range jobRows {
116 out := make([]any, len(rethinkRunningColumns))
117 for i, col := range rethinkRunningColumns {
118 switch col.Name {
119 case "jobId":
120 out[i] = row.JobID
121 case "query":
122 out[i] = strmutil.TruncateText(row.Query, rethinkMaxQueryTextLength)
123 case "durationMs":
124 out[i] = row.DurationMs
125 case "type":
126 out[i] = row.Type
127 case "user":
128 out[i] = row.User
129 case "clientAddress":
130 out[i] = row.ClientAddress
131 case "clientPort":
132 out[i] = row.ClientPort
133 case "servers":
134 out[i] = row.Servers
135 default:
136 out[i] = nil
137 }
138 }
139 data = append(data, out)
140 }
141
142 return &funcapi.FunctionResponse{
143 Status: 200,
144 Help: "Currently running queries from rethinkdb.jobs. WARNING: Query text may contain unmasked literals (potential PII).",
145 Columns: cs.BuildColumns(),
146 Data: data,
147 DefaultSortColumn: sortColumn,
148 RequiredParams: []funcapi.ParamConfig{funcapi.BuildSortParam(rethinkRunningColumns)},
149 }
150 }
151
152 type rethinkColumn struct {
153 funcapi.ColumnMeta
154 sortOpt bool // whether this column appears as a sort option
155 sortLbl string // label for sort option dropdown
156 defaultSort bool // default sort column
157 }
158
159 // funcapi.SortableColumn interface implementation for rethinkColumn.
160 func (c rethinkColumn) IsSortOption() bool { return c.sortOpt }
161 func (c rethinkColumn) SortLabel() string { return c.sortLbl }
162 func (c rethinkColumn) IsDefaultSort() bool { return c.defaultSort }
163 func (c rethinkColumn) ColumnName() string { return c.Name }
164 func (c rethinkColumn) SortColumn() string { return "" }
165
166 func rethinkColumnSet(cols []rethinkColumn) funcapi.ColumnSet[rethinkColumn] {
167 return funcapi.Columns(cols, func(c rethinkColumn) funcapi.ColumnMeta { return c.ColumnMeta })
168 }
169
170 var rethinkRunningColumns = []rethinkColumn{
171 {ColumnMeta: funcapi.ColumnMeta{Name: "jobId", Tooltip: "Job ID", Type: funcapi.FieldTypeString, Visible: false, Sortable: true, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, UniqueKey: true, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount}},
172 {ColumnMeta: funcapi.ColumnMeta{Name: "query", Tooltip: "Query", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, Sticky: true, FullWidth: true, Wrap: true}},
173 {ColumnMeta: funcapi.ColumnMeta{Name: "durationMs", Tooltip: "Duration", Type: funcapi.FieldTypeDuration, Visible: true, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualBar, Transform: funcapi.FieldTransformDuration, Units: "milliseconds", DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax}, sortOpt: true, defaultSort: true, sortLbl: "Running queries by Duration"},
174 {ColumnMeta: funcapi.ColumnMeta{Name: "type", Tooltip: "Type", Type: funcapi.FieldTypeString, Visible: true, Sortable: true, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount}},
175 {ColumnMeta: funcapi.ColumnMeta{Name: "user", Tooltip: "User", Type: funcapi.FieldTypeString, Visible: true, Sortable: true, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount}},
176 {ColumnMeta: funcapi.ColumnMeta{Name: "clientAddress", Tooltip: "Client Address", Type: funcapi.FieldTypeString, Visible: false, Sortable: true, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount}},
177 {ColumnMeta: funcapi.ColumnMeta{Name: "clientPort", Tooltip: "Client Port", Type: funcapi.FieldTypeInteger, Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax}},
178 {ColumnMeta: funcapi.ColumnMeta{Name: "servers", Tooltip: "Servers", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}},
179 }
180
181 type rethinkJobRow struct {
182 JobID string
183 Query string
184 DurationMs float64
185 Type string
186 User string
187 ClientAddress string
188 ClientPort int64
189 Servers string
190 }
191
192 func parseRethinkJob(row map[string]any) rethinkJobRow {
193 info := mapStringAny(row["info"])
194 query := fmt.Sprint(info["query"])
195 user := fmt.Sprint(info["user"])
196 clientAddr := fmt.Sprint(info["client_address"])
197 clientPort := toInt64(info["client_port"])
198
199 servers := ""
200 if list, ok := row["servers"].([]any); ok {
201 ss := make([]string, 0, len(list))
202 for _, v := range list {
203 ss = append(ss, fmt.Sprint(v))
204 }
205 servers = strings.Join(ss, ",")
206 }
207
208 return rethinkJobRow{
209 JobID: fmt.Sprint(row["id"]),
210 Query: query,
211 DurationMs: toFloat64(row["duration_sec"]) * 1000,
212 Type: fmt.Sprint(row["type"]),
213 User: user,
214 ClientAddress: clientAddr,
215 ClientPort: clientPort,
216 Servers: servers,
217 }
218 }
219
220 func mapStringAny(v any) map[string]any {
221 if m, ok := v.(map[string]any); ok {
222 return m
223 }
224 return map[string]any{}
225 }
226
227 func toFloat64(v any) float64 {
228 switch t := v.(type) {
229 case float64:
230 return t
231 case float32:
232 return float64(t)
233 case int:
234 return float64(t)
235 case int64:
236 return float64(t)
237 case uint64:
238 return float64(t)
239 default:
240 return 0
241 }
242 }
243
244 func toInt64(v any) int64 {
245 switch t := v.(type) {
246 case int64:
247 return t
248 case int:
249 return int64(t)
250 case float64:
251 return int64(t)
252 case float32:
253 return int64(t)
254 default:
255 return 0
256 }
257 }
258
259 func mapRethinkSortColumn(input string) string {
260 for _, col := range rethinkRunningColumns {
261 if col.IsSortOption() && col.Name == input {
262 return col.Name
263 }
264 }
265 for _, col := range rethinkRunningColumns {
266 if col.IsDefaultSort() {
267 return col.Name
268 }
269 }
270 for _, col := range rethinkRunningColumns {
271 if col.IsSortOption() {
272 return col.Name
273 }
274 }
275 return ""
276 }
277
278 func sortRethinkRows(rows []rethinkJobRow) {
279 sort.Slice(rows, func(i, j int) bool { return rows[i].DurationMs > rows[j].DurationMs })
280 }