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