| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | package elasticsearch |
| 4 | |
| 5 | import ( |
| 6 | "context" |
| 7 | "fmt" |
| 8 | "net/url" |
| 9 | "sort" |
| 10 | "strconv" |
| 11 | "strings" |
| 12 | "time" |
| 13 | |
| 14 | "github.com/netdata/netdata/go/plugins/pkg/funcapi" |
| 15 | "github.com/netdata/netdata/go/plugins/pkg/web" |
| 16 | "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/strmutil" |
| 17 | ) |
| 18 | |
| 19 | const ( |
| 20 | topQueriesMethodID = "top-queries" |
| 21 | topQueriesMaxTextLength = 4096 |
| 22 | ) |
| 23 | |
| 24 | func topQueriesMethodConfig() funcapi.MethodConfig { |
| 25 | return funcapi.MethodConfig{ |
| 26 | ID: topQueriesMethodID, |
| 27 | Name: "Top Queries", |
| 28 | UpdateEvery: 10, |
| 29 | Help: "Running queries from Elasticsearch Tasks API", |
| 30 | RequireCloud: true, |
| 31 | RequiredParams: []funcapi.ParamConfig{funcapi.BuildSortParam(topQueriesColumns)}, |
| 32 | } |
| 33 | } |
| 34 | |
| 35 | type topQueriesColumn struct { |
| 36 | funcapi.ColumnMeta |
| 37 | sortOpt bool // whether this column appears as a sort option |
| 38 | sortLbl string // label for sort option dropdown |
| 39 | defaultSort bool // default sort column |
| 40 | } |
| 41 | |
| 42 | // funcapi.SortableColumn interface implementation for topQueriesColumn. |
| 43 | func (c topQueriesColumn) IsSortOption() bool { return c.sortOpt } |
| 44 | func (c topQueriesColumn) SortLabel() string { return c.sortLbl } |
| 45 | func (c topQueriesColumn) IsDefaultSort() bool { return c.defaultSort } |
| 46 | func (c topQueriesColumn) ColumnName() string { return c.Name } |
| 47 | func (c topQueriesColumn) SortColumn() string { return "" } |
| 48 | |
| 49 | var topQueriesColumns = []topQueriesColumn{ |
| 50 | {ColumnMeta: funcapi.ColumnMeta{Name: "taskId", Tooltip: "Task ID", Type: funcapi.FieldTypeString, Visible: false, Sortable: true, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, UniqueKey: true, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryCount}, sortOpt: true, sortLbl: "Top queries by Task ID"}, |
| 51 | {ColumnMeta: funcapi.ColumnMeta{Name: "node", Tooltip: "Node ID", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{}}}, |
| 52 | {ColumnMeta: funcapi.ColumnMeta{Name: "nodeName", Tooltip: "Node Name", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{IsDefault: true}}}, |
| 53 | {ColumnMeta: funcapi.ColumnMeta{Name: "action", Tooltip: "Action", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{}}}, |
| 54 | {ColumnMeta: funcapi.ColumnMeta{Name: "type", Tooltip: "Type", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{}}}, |
| 55 | {ColumnMeta: funcapi.ColumnMeta{Name: "description", Tooltip: "Description", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, Sticky: true, FullWidth: true, Wrap: true}}, |
| 56 | {ColumnMeta: funcapi.ColumnMeta{Name: "startTime", Tooltip: "Start Time", Type: funcapi.FieldTypeTimestamp, Visible: true, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformDatetime, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax}, sortOpt: true, sortLbl: "Top queries by Start Time"}, |
| 57 | {ColumnMeta: funcapi.ColumnMeta{Name: "runningTime", Tooltip: "Running Time", 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.FieldSummarySum, Chart: &funcapi.ChartOptions{Group: "RunningTime", Title: "Running Time", IsDefault: true}}, sortOpt: true, defaultSort: true, sortLbl: "Top queries by Running Time"}, |
| 58 | {ColumnMeta: funcapi.ColumnMeta{Name: "cancellable", Tooltip: "Cancellable", Type: funcapi.FieldTypeBoolean, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformNone}}, |
| 59 | {ColumnMeta: funcapi.ColumnMeta{Name: "cancelled", Tooltip: "Cancelled", Type: funcapi.FieldTypeBoolean, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformNone}}, |
| 60 | } |
| 61 | |
| 62 | type topQueriesResponse struct { |
| 63 | Nodes map[string]struct { |
| 64 | Name string `json:"name"` |
| 65 | Tasks map[string]topQueriesTask `json:"tasks"` |
| 66 | } `json:"nodes"` |
| 67 | } |
| 68 | |
| 69 | type topQueriesTask struct { |
| 70 | ID int64 `json:"id"` |
| 71 | Action string `json:"action"` |
| 72 | Type string `json:"type"` |
| 73 | Description string `json:"description"` |
| 74 | StartTimeInMillis int64 `json:"start_time_in_millis"` |
| 75 | RunningTimeInNanos int64 `json:"running_time_in_nanos"` |
| 76 | Cancellable bool `json:"cancellable"` |
| 77 | Cancelled bool `json:"cancelled"` |
| 78 | } |
| 79 | |
| 80 | type topQueriesRow struct { |
| 81 | TaskID string |
| 82 | NodeID string |
| 83 | NodeName string |
| 84 | Action string |
| 85 | Type string |
| 86 | Description string |
| 87 | StartTime time.Time |
| 88 | RunningTime time.Duration |
| 89 | Cancellable bool |
| 90 | Cancelled bool |
| 91 | } |
| 92 | |
| 93 | // funcTopQueries implements funcapi.MethodHandler for Elasticsearch top-queries. |
| 94 | type funcTopQueries struct { |
| 95 | router *funcRouter |
| 96 | } |
| 97 | |
| 98 | func newFuncTopQueries(r *funcRouter) *funcTopQueries { |
| 99 | return &funcTopQueries{router: r} |
| 100 | } |
| 101 | |
| 102 | // Compile-time interface check. |
| 103 | var _ funcapi.MethodHandler = (*funcTopQueries)(nil) |
| 104 | |
| 105 | func (f *funcTopQueries) Cleanup(ctx context.Context) {} |
| 106 | |
| 107 | // MethodParams implements funcapi.MethodHandler. |
| 108 | func (f *funcTopQueries) MethodParams(_ context.Context, method string) ([]funcapi.ParamConfig, error) { |
| 109 | switch method { |
| 110 | case topQueriesMethodID: |
| 111 | if f.router.collector.Functions.TopQueries.Disabled { |
| 112 | return nil, fmt.Errorf("top-queries function disabled in configuration") |
| 113 | } |
| 114 | return []funcapi.ParamConfig{funcapi.BuildSortParam(topQueriesColumns)}, nil |
| 115 | default: |
| 116 | return nil, fmt.Errorf("unknown method: %s", method) |
| 117 | } |
| 118 | } |
| 119 | |
| 120 | // Handle implements funcapi.MethodHandler. |
| 121 | func (f *funcTopQueries) Handle(ctx context.Context, method string, params funcapi.ResolvedParams) *funcapi.FunctionResponse { |
| 122 | if f.router.collector.httpClient == nil { |
| 123 | return funcapi.UnavailableResponse("collector is still initializing, please retry in a few seconds") |
| 124 | } |
| 125 | |
| 126 | switch method { |
| 127 | case topQueriesMethodID: |
| 128 | if f.router.collector.Functions.TopQueries.Disabled { |
| 129 | return funcapi.UnavailableResponse("top-queries function has been disabled in configuration") |
| 130 | } |
| 131 | queryCtx, cancel := context.WithTimeout(ctx, f.router.collector.topQueriesTimeout()) |
| 132 | defer cancel() |
| 133 | return f.collectData(queryCtx, params.Column("__sort")) |
| 134 | default: |
| 135 | return funcapi.NotFoundResponse(method) |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | func (f *funcTopQueries) collectData(ctx context.Context, sortColumn string) *funcapi.FunctionResponse { |
| 140 | limit := f.router.collector.topQueriesLimit() |
| 141 | |
| 142 | req, err := web.NewHTTPRequestWithPath(f.router.collector.RequestConfig, "/_tasks") |
| 143 | if err != nil { |
| 144 | return &funcapi.FunctionResponse{Status: 500, Message: err.Error()} |
| 145 | } |
| 146 | req = req.WithContext(ctx) |
| 147 | q := url.Values{} |
| 148 | q.Set("actions", "*search") |
| 149 | q.Set("detailed", "true") |
| 150 | req.URL.RawQuery = q.Encode() |
| 151 | |
| 152 | var resp topQueriesResponse |
| 153 | if err := web.DoHTTP(f.router.collector.httpClient).RequestJSON(req, &resp); err != nil { |
| 154 | if ctx.Err() == context.DeadlineExceeded { |
| 155 | return &funcapi.FunctionResponse{Status: 504, Message: "query timed out"} |
| 156 | } |
| 157 | return &funcapi.FunctionResponse{Status: 500, Message: fmt.Sprintf("tasks query failed: %v", err)} |
| 158 | } |
| 159 | |
| 160 | rows := make([]topQueriesRow, 0, 100) |
| 161 | for nodeID, node := range resp.Nodes { |
| 162 | for taskID, task := range node.Tasks { |
| 163 | rows = append(rows, topQueriesRow{ |
| 164 | TaskID: taskID, |
| 165 | NodeID: nodeID, |
| 166 | NodeName: node.Name, |
| 167 | Action: task.Action, |
| 168 | Type: task.Type, |
| 169 | Description: task.Description, |
| 170 | StartTime: time.UnixMilli(task.StartTimeInMillis), |
| 171 | RunningTime: time.Duration(task.RunningTimeInNanos), |
| 172 | Cancellable: task.Cancellable, |
| 173 | Cancelled: task.Cancelled, |
| 174 | }) |
| 175 | } |
| 176 | } |
| 177 | |
| 178 | cs := f.columnSet(topQueriesColumns) |
| 179 | sortParam := funcapi.BuildSortParam(topQueriesColumns) |
| 180 | |
| 181 | if len(rows) == 0 { |
| 182 | return &funcapi.FunctionResponse{ |
| 183 | Status: 200, |
| 184 | Message: "No running search tasks found.", |
| 185 | Help: "Running queries from Elasticsearch Tasks API", |
| 186 | Columns: cs.BuildColumns(), |
| 187 | Data: [][]any{}, |
| 188 | DefaultSortColumn: "runningTime", |
| 189 | RequiredParams: []funcapi.ParamConfig{sortParam}, |
| 190 | ChartingConfig: cs.BuildCharting(), |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | sortColumn = f.mapSortColumn(sortColumn) |
| 195 | f.sortRows(rows, sortColumn) |
| 196 | |
| 197 | if len(rows) > limit { |
| 198 | rows = rows[:limit] |
| 199 | } |
| 200 | |
| 201 | data := make([][]any, 0, len(rows)) |
| 202 | for _, row := range rows { |
| 203 | out := make([]any, len(topQueriesColumns)) |
| 204 | for i, col := range topQueriesColumns { |
| 205 | switch col.Name { |
| 206 | case "taskId": |
| 207 | out[i] = row.TaskID |
| 208 | case "node": |
| 209 | out[i] = row.NodeID |
| 210 | case "nodeName": |
| 211 | out[i] = row.NodeName |
| 212 | case "action": |
| 213 | out[i] = row.Action |
| 214 | case "type": |
| 215 | out[i] = row.Type |
| 216 | case "description": |
| 217 | out[i] = strmutil.TruncateText(row.Description, topQueriesMaxTextLength) |
| 218 | case "startTime": |
| 219 | out[i] = row.StartTime.Format(time.RFC3339Nano) |
| 220 | case "runningTime": |
| 221 | out[i] = float64(row.RunningTime) / float64(time.Millisecond) |
| 222 | case "cancellable": |
| 223 | out[i] = row.Cancellable |
| 224 | case "cancelled": |
| 225 | out[i] = row.Cancelled |
| 226 | default: |
| 227 | out[i] = nil |
| 228 | } |
| 229 | } |
| 230 | data = append(data, out) |
| 231 | } |
| 232 | |
| 233 | return &funcapi.FunctionResponse{ |
| 234 | Status: 200, |
| 235 | Help: "Running queries from Elasticsearch Tasks API", |
| 236 | Columns: cs.BuildColumns(), |
| 237 | Data: data, |
| 238 | DefaultSortColumn: "runningTime", |
| 239 | RequiredParams: []funcapi.ParamConfig{sortParam}, |
| 240 | ChartingConfig: cs.BuildCharting(), |
| 241 | } |
| 242 | } |
| 243 | |
| 244 | func (f *funcTopQueries) columnSet(cols []topQueriesColumn) funcapi.ColumnSet[topQueriesColumn] { |
| 245 | return funcapi.Columns(cols, func(c topQueriesColumn) funcapi.ColumnMeta { return c.ColumnMeta }) |
| 246 | } |
| 247 | |
| 248 | func (f *funcTopQueries) mapSortColumn(col string) string { |
| 249 | switch col { |
| 250 | case "runningTime", "startTime", "taskId": |
| 251 | return col |
| 252 | default: |
| 253 | return "runningTime" |
| 254 | } |
| 255 | } |
| 256 | |
| 257 | func (f *funcTopQueries) sortRows(rows []topQueriesRow, sortColumn string) { |
| 258 | switch sortColumn { |
| 259 | case "startTime": |
| 260 | sort.Slice(rows, func(i, j int) bool { |
| 261 | return rows[i].StartTime.After(rows[j].StartTime) |
| 262 | }) |
| 263 | case "taskId": |
| 264 | sort.Slice(rows, func(i, j int) bool { |
| 265 | left, lok := f.parseTaskID(rows[i].TaskID) |
| 266 | right, rok := f.parseTaskID(rows[j].TaskID) |
| 267 | if lok && rok { |
| 268 | return left > right |
| 269 | } |
| 270 | return rows[i].TaskID > rows[j].TaskID |
| 271 | }) |
| 272 | default: |
| 273 | sort.Slice(rows, func(i, j int) bool { |
| 274 | return rows[i].RunningTime > rows[j].RunningTime |
| 275 | }) |
| 276 | } |
| 277 | } |
| 278 | |
| 279 | func (f *funcTopQueries) parseTaskID(taskID string) (int64, bool) { |
| 280 | id := taskID |
| 281 | if idx := strings.LastIndex(id, ":"); idx != -1 { |
| 282 | id = id[idx+1:] |
| 283 | } |
| 284 | val, err := strconv.ParseInt(id, 10, 64) |
| 285 | if err != nil { |
| 286 | return 0, false |
| 287 | } |
| 288 | return val, true |
| 289 | } |