master
go 619 lines 28 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package mongo
4
5 import (
6 "context"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "sort"
11 "time"
12
13 "github.com/netdata/netdata/go/plugins/pkg/funcapi"
14 "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/strmutil"
15
16 "go.mongodb.org/mongo-driver/bson"
17 "go.mongodb.org/mongo-driver/mongo/options"
18 )
19
20 const (
21 topQueriesMethodID = "top-queries"
22 topQueriesMaxTextLength = 4096
23 topQueriesDefaultLimit = 500
24 topQueriesHelpText = "Top queries from MongoDB Profiler (system.profile). " +
25 "WARNING: Query text may contain unmasked literals (potential PII). " +
26 "Requires profiling enabled on target databases (db.setProfilingLevel)."
27 )
28
29 func topQueriesMethodConfig() funcapi.MethodConfig {
30 return funcapi.MethodConfig{
31 ID: topQueriesMethodID,
32 Name: "Top Queries",
33 UpdateEvery: 10,
34 Help: topQueriesHelpText,
35 RequireCloud: true,
36 RequiredParams: []funcapi.ParamConfig{funcapi.BuildSortParam(topQueriesColumns)},
37 }
38 }
39
40 const topQueriesParamSort = "__sort"
41
42 // topQueriesColumn defines metadata for a MongoDB profile column.
43 // Embeds funcapi.ColumnMeta for UI rendering and adds MongoDB-specific fields.
44 type topQueriesColumn struct {
45 funcapi.ColumnMeta
46
47 // DBField is the MongoDB document field name (e.g., "millis")
48 DBField string
49 // sortOpt indicates whether this column can be used for sorting in params
50 sortOpt bool
51 // sortLbl is the label for the sort option dropdown
52 sortLbl string
53 // defaultSort indicates whether this is the default sort column
54 defaultSort bool
55 }
56
57 // funcapi.SortableColumn interface implementation for topQueriesColumn.
58 func (c topQueriesColumn) IsSortOption() bool { return c.sortOpt }
59 func (c topQueriesColumn) SortLabel() string { return fmt.Sprintf("Top queries by %s", c.sortLbl) }
60 func (c topQueriesColumn) IsDefaultSort() bool { return c.defaultSort }
61 func (c topQueriesColumn) ColumnName() string { return c.Name }
62 func (c topQueriesColumn) SortColumn() string { return c.DBField }
63
64 // topQueriesColumns defines all available columns from system.profile.
65 // Ordered by display priority (index).
66 var topQueriesColumns = []topQueriesColumn{
67 // Core fields (visible by default)
68 {ColumnMeta: funcapi.ColumnMeta{Name: "timestamp", Tooltip: "Timestamp", Type: funcapi.FieldTypeTimestamp, Visible: true, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummaryMax, Transform: funcapi.FieldTransformDatetime, UniqueKey: true}, DBField: "ts", sortOpt: true, sortLbl: "Timestamp"},
69 {ColumnMeta: funcapi.ColumnMeta{Name: "namespace", Tooltip: "Namespace", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Sticky: true, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, ExpandFilter: true, GroupBy: &funcapi.GroupByOptions{IsDefault: true}}, DBField: "ns"},
70 {ColumnMeta: funcapi.ColumnMeta{Name: "operation", Tooltip: "Operation", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{}}, DBField: "op"},
71 {ColumnMeta: funcapi.ColumnMeta{Name: "query", Tooltip: "Query", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, FullWidth: true, Wrap: true, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "command"},
72 {ColumnMeta: funcapi.ColumnMeta{Name: "execution_time", Tooltip: "Execution Time", Type: funcapi.FieldTypeDuration, Units: "seconds", Visible: true, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualBar, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformDuration, DecimalPoints: 3, Chart: &funcapi.ChartOptions{Group: "Time", Title: "Execution Time", IsDefault: true}}, DBField: "millis", sortOpt: true, defaultSort: true, sortLbl: "Execution Time"},
73 {ColumnMeta: funcapi.ColumnMeta{Name: "docs_examined", Tooltip: "Docs Examined", Type: funcapi.FieldTypeInteger, Visible: true, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Docs", Title: "Documents"}}, DBField: "docsExamined", sortOpt: true, sortLbl: "Docs Examined"},
74 {ColumnMeta: funcapi.ColumnMeta{Name: "keys_examined", Tooltip: "Keys Examined", Type: funcapi.FieldTypeInteger, Visible: true, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Docs", Title: "Documents"}}, DBField: "keysExamined", sortOpt: true, sortLbl: "Keys Examined"},
75 {ColumnMeta: funcapi.ColumnMeta{Name: "docs_returned", Tooltip: "Docs Returned", Type: funcapi.FieldTypeInteger, Visible: true, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Docs", Title: "Documents"}}, DBField: "nreturned", sortOpt: true, sortLbl: "Docs Returned"},
76 {ColumnMeta: funcapi.ColumnMeta{Name: "plan_summary", Tooltip: "Plan Summary", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "planSummary"},
77
78 // Secondary fields
79 {ColumnMeta: funcapi.ColumnMeta{Name: "client", Tooltip: "Client", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{}}, DBField: "client"},
80 {ColumnMeta: funcapi.ColumnMeta{Name: "user", Tooltip: "User", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{}}, DBField: "user"},
81 {ColumnMeta: funcapi.ColumnMeta{Name: "docs_deleted", Tooltip: "Docs Deleted", Type: funcapi.FieldTypeInteger, Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Docs", Title: "Documents"}}, DBField: "ndeleted", sortOpt: true, sortLbl: "Docs Deleted"},
82 {ColumnMeta: funcapi.ColumnMeta{Name: "docs_inserted", Tooltip: "Docs Inserted", Type: funcapi.FieldTypeInteger, Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Docs", Title: "Documents"}}, DBField: "ninserted", sortOpt: true, sortLbl: "Docs Inserted"},
83 {ColumnMeta: funcapi.ColumnMeta{Name: "docs_modified", Tooltip: "Docs Modified", Type: funcapi.FieldTypeInteger, Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Docs", Title: "Documents"}}, DBField: "nModified", sortOpt: true, sortLbl: "Docs Modified"},
84 {ColumnMeta: funcapi.ColumnMeta{Name: "response_length", Tooltip: "Response Length", Type: funcapi.FieldTypeInteger, Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Response", Title: "Response Size"}}, DBField: "responseLength", sortOpt: true, sortLbl: "Response Length"},
85 {ColumnMeta: funcapi.ColumnMeta{Name: "num_yield", Tooltip: "Num Yield", Type: funcapi.FieldTypeInteger, Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualValue, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformNumber, Chart: &funcapi.ChartOptions{Group: "Yield", Title: "Yield"}}, DBField: "numYield", sortOpt: true, sortLbl: "Num Yield"},
86 {ColumnMeta: funcapi.ColumnMeta{Name: "app_name", Tooltip: "App Name", Type: funcapi.FieldTypeString, Visible: true, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText, GroupBy: &funcapi.GroupByOptions{}}, DBField: "appName"},
87 {ColumnMeta: funcapi.ColumnMeta{Name: "cursor_exhausted", Tooltip: "Cursor Exhausted", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "cursorExhausted"},
88 {ColumnMeta: funcapi.ColumnMeta{Name: "has_sort_stage", Tooltip: "Has Sort Stage", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "hasSortStage"},
89 {ColumnMeta: funcapi.ColumnMeta{Name: "uses_disk", Tooltip: "Uses Disk", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "usedDisk"},
90 {ColumnMeta: funcapi.ColumnMeta{Name: "from_multi_planner", Tooltip: "From Multi Planner", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "fromMultiPlanner"},
91 {ColumnMeta: funcapi.ColumnMeta{Name: "replanned", Tooltip: "Replanned", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "replanned"},
92
93 // Version-specific fields (hidden by default)
94 {ColumnMeta: funcapi.ColumnMeta{Name: "query_hash", Tooltip: "Query Hash", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "queryHash"}, // 4.2+
95 {ColumnMeta: funcapi.ColumnMeta{Name: "plan_cache_key", Tooltip: "Plan Cache Key", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "planCacheKey"}, // 4.2+
96 {ColumnMeta: funcapi.ColumnMeta{Name: "planning_time", Tooltip: "Planning Time", Type: funcapi.FieldTypeDuration, Units: "seconds", Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualBar, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformDuration, DecimalPoints: 3, Chart: &funcapi.ChartOptions{Group: "Time", Title: "Execution Time"}}, DBField: "planningTimeMicros", sortOpt: true, sortLbl: "Planning Time"},
97 {ColumnMeta: funcapi.ColumnMeta{Name: "cpu_time", Tooltip: "CPU Time", Type: funcapi.FieldTypeDuration, Units: "seconds", Visible: false, Sortable: true, Filter: funcapi.FieldFilterRange, Visualization: funcapi.FieldVisualBar, Summary: funcapi.FieldSummarySum, Transform: funcapi.FieldTransformDuration, DecimalPoints: 3, Chart: &funcapi.ChartOptions{Group: "Time", Title: "Execution Time"}}, DBField: "cpuNanos", sortOpt: true, sortLbl: "CPU Time"},
98 {ColumnMeta: funcapi.ColumnMeta{Name: "query_framework", Tooltip: "Query Framework", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "queryFramework"}, // 7.0+
99 {ColumnMeta: funcapi.ColumnMeta{Name: "query_shape_hash", Tooltip: "Query Shape Hash", Type: funcapi.FieldTypeString, Visible: false, Sortable: false, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualValue, Transform: funcapi.FieldTransformText}, DBField: "queryShapeHash"}, // 8.0+
100 }
101
102 // topQueriesProfileDocument represents a document from system.profile
103 type topQueriesProfileDocument struct {
104 // Core fields (always present)
105 Timestamp time.Time `bson:"ts"`
106 Op string `bson:"op"`
107 Ns string `bson:"ns"`
108 Command bson.M `bson:"command"`
109 Millis int64 `bson:"millis"`
110 PlanSummary string `bson:"planSummary"`
111
112 // Common fields
113 DocsExamined int64 `bson:"docsExamined"`
114 KeysExamined int64 `bson:"keysExamined"`
115 Nreturned int64 `bson:"nreturned"`
116 Client string `bson:"client"`
117 User string `bson:"user"`
118 Ndeleted int64 `bson:"ndeleted"`
119 Ninserted int64 `bson:"ninserted"`
120 NModified int64 `bson:"nModified"`
121 ResponseLength int64 `bson:"responseLength"`
122 NumYield int64 `bson:"numYield"`
123 AppName string `bson:"appName"`
124
125 // Boolean fields (pointers for nil detection)
126 CursorExhausted *bool `bson:"cursorExhausted"`
127 HasSortStage *bool `bson:"hasSortStage"`
128 UsedDisk *bool `bson:"usedDisk"`
129 FromMultiPlanner *bool `bson:"fromMultiPlanner"`
130 Replanned *bool `bson:"replanned"`
131
132 // Version-specific fields
133 QueryHash string `bson:"queryHash"` // 4.2+
134 PlanCacheKey string `bson:"planCacheKey"` // 4.2+
135 PlanningTimeMicros *int64 `bson:"planningTimeMicros"` // 6.2+
136 CpuNanos *int64 `bson:"cpuNanos"` // 6.3+ Linux only
137 QueryFramework string `bson:"queryFramework"` // 7.0+
138 QueryShapeHash string `bson:"queryShapeHash"` // 8.0+
139 }
140
141 // funcTopQueries implements funcapi.MethodHandler for MongoDB top-queries.
142 // All function-related logic is encapsulated here, keeping Collector focused on metrics collection.
143 type funcTopQueries struct {
144 router *funcRouter
145 }
146
147 func newFuncTopQueries(r *funcRouter) *funcTopQueries {
148 return &funcTopQueries{router: r}
149 }
150
151 // Compile-time interface check.
152 var _ funcapi.MethodHandler = (*funcTopQueries)(nil)
153
154 func (f *funcTopQueries) Cleanup(ctx context.Context) {}
155
156 // MethodParams implements funcapi.MethodHandler.
157 func (f *funcTopQueries) MethodParams(ctx context.Context, method string) ([]funcapi.ParamConfig, error) {
158 if f.router.collector.conn == nil {
159 return nil, fmt.Errorf("collector is still initializing")
160 }
161 switch method {
162 case topQueriesMethodID:
163 if f.router.collector.Functions.TopQueries.Disabled {
164 return nil, fmt.Errorf("top-queries function disabled in configuration")
165 }
166 return f.methodParams(ctx)
167 default:
168 return nil, fmt.Errorf("unknown method: %s", method)
169 }
170 }
171
172 // Handle implements funcapi.MethodHandler.
173 func (f *funcTopQueries) Handle(ctx context.Context, method string, params funcapi.ResolvedParams) *funcapi.FunctionResponse {
174 if f.router.collector.conn == nil {
175 return funcapi.UnavailableResponse("collector is still initializing, please retry in a few seconds")
176 }
177
178 switch method {
179 case topQueriesMethodID:
180 if f.router.collector.Functions.TopQueries.Disabled {
181 return funcapi.UnavailableResponse("top-queries function has been disabled in configuration")
182 }
183 queryCtx, cancel := context.WithTimeout(ctx, f.router.collector.topQueriesTimeout())
184 defer cancel()
185 return f.collectData(queryCtx, params.Column(topQueriesParamSort))
186 default:
187 return funcapi.NotFoundResponse(method)
188 }
189 }
190
191 func (f *funcTopQueries) methodParams(ctx context.Context) ([]funcapi.ParamConfig, error) {
192
193 databases, err := f.getDatabases()
194 if err != nil {
195 return nil, err
196 }
197
198 availableFields, err := f.detectProfileFields(ctx, databases)
199 if err != nil {
200 return nil, err
201 }
202
203 availableCols := f.buildAvailableColumns(availableFields)
204 sortParam := funcapi.BuildSortParam(availableCols)
205 return []funcapi.ParamConfig{sortParam}, nil
206 }
207
208 func (f *funcTopQueries) collectData(ctx context.Context, sortColumn string) *funcapi.FunctionResponse {
209 limit := f.router.collector.topQueriesLimit()
210
211 // Build valid sort columns map from metadata
212 validSortCols := make(map[string]bool)
213 for _, col := range topQueriesColumns {
214 if col.IsSortOption() {
215 validSortCols[col.DBField] = true
216 }
217 }
218 if !validSortCols[sortColumn] {
219 sortColumn = "millis" // safe default
220 }
221
222 databases, err := f.getDatabases()
223 if err != nil {
224 return &funcapi.FunctionResponse{
225 Status: 500,
226 Message: fmt.Sprintf("failed to list databases: %v", err),
227 }
228 }
229
230 // Detect available fields (with caching)
231 availableFields, err := f.detectProfileFields(ctx, databases)
232 if err != nil {
233 f.router.collector.Debugf("failed to detect profile fields: %v", err)
234 }
235 availableCols := f.buildAvailableColumns(availableFields)
236 cs := f.columnSet(availableCols)
237 sortParam := funcapi.BuildSortParam(availableCols)
238
239 // Query system.profile from each database
240 var allDocs []topQueriesProfileDocument
241 var profilingDisabledDBs []string
242 var failedDBs []string
243 var successfulDBs int
244
245 for _, dbName := range databases {
246 docs, enabled, err := f.querySystemProfile(ctx, dbName, sortColumn, limit)
247 if err != nil {
248 // Check for timeout (parent or child context)
249 if ctx.Err() == context.DeadlineExceeded || errors.Is(err, context.DeadlineExceeded) {
250 return &funcapi.FunctionResponse{Status: 504, Message: "query timed out"}
251 }
252 f.router.collector.Debugf("failed to query system.profile in %s: %v", dbName, err)
253 failedDBs = append(failedDBs, dbName)
254 continue
255 }
256
257 if !enabled {
258 profilingDisabledDBs = append(profilingDisabledDBs, dbName)
259 continue
260 }
261
262 successfulDBs++
263 allDocs = append(allDocs, docs...)
264 }
265
266 // Check if all databases failed with errors (not just profiling disabled)
267 if successfulDBs == 0 && len(failedDBs) > 0 && len(profilingDisabledDBs) == 0 {
268 return &funcapi.FunctionResponse{
269 Status: 500,
270 Message: fmt.Sprintf("failed to query all databases: %v", failedDBs),
271 }
272 }
273
274 // Check if profiling is disabled everywhere (no successful queries, no errors, only disabled)
275 if len(allDocs) == 0 && len(profilingDisabledDBs) > 0 && len(failedDBs) == 0 {
276 return &funcapi.FunctionResponse{
277 Status: 503,
278 Message: fmt.Sprintf(
279 "Database profiling is disabled. Enable it with: db.setProfilingLevel(1, {slowms: 100}). "+
280 "Disabled databases: %v", profilingDisabledDBs),
281 }
282 }
283
284 // Check if we have a mix of failures and disabled profiling (no successful queries at all)
285 if successfulDBs == 0 && (len(failedDBs) > 0 || len(profilingDisabledDBs) > 0) {
286 msg := "No databases could be queried successfully."
287 if len(failedDBs) > 0 {
288 msg += fmt.Sprintf(" Failed: %v.", failedDBs)
289 }
290 if len(profilingDisabledDBs) > 0 {
291 msg += fmt.Sprintf(" Profiling disabled: %v.", profilingDisabledDBs)
292 }
293 return &funcapi.FunctionResponse{
294 Status: 503,
295 Message: msg,
296 }
297 }
298
299 // Build empty response structure
300 emptyResponse := &funcapi.FunctionResponse{
301 Status: 200,
302 Message: "No slow queries found. Profiling may be disabled or no queries exceeded the slowms threshold.",
303 Help: topQueriesHelpText,
304 Columns: cs.BuildColumns(),
305 Data: [][]any{},
306 DefaultSortColumn: "execution_time",
307 RequiredParams: []funcapi.ParamConfig{sortParam},
308 ChartingConfig: cs.BuildCharting(),
309 }
310
311 if len(allDocs) == 0 {
312 return emptyResponse
313 }
314
315 // Sort all documents by the requested column
316 f.sortDocuments(allDocs, sortColumn)
317
318 // Apply limit
319 if len(allDocs) > limit {
320 allDocs = allDocs[:limit]
321 }
322
323 // Convert to response format: [][]any (array of arrays, ordered by column index)
324 data := make([][]any, 0, len(allDocs))
325 for _, doc := range allDocs {
326 row := make([]any, len(availableCols))
327
328 // Fill each column based on available columns
329 for i, col := range availableCols {
330 switch col.Name {
331 case "timestamp":
332 row[i] = doc.Timestamp.Format(time.RFC3339Nano)
333 case "namespace":
334 row[i] = doc.Ns
335 case "operation":
336 row[i] = doc.Op
337 case "query":
338 cmdJSON, err := json.Marshal(doc.Command)
339 if err != nil {
340 cmdJSON = []byte("{}")
341 }
342 row[i] = strmutil.TruncateText(string(cmdJSON), topQueriesMaxTextLength)
343 case "execution_time":
344 row[i] = float64(doc.Millis) / 1000.0 // ms to seconds
345 case "docs_examined":
346 row[i] = doc.DocsExamined
347 case "keys_examined":
348 row[i] = doc.KeysExamined
349 case "docs_returned":
350 row[i] = doc.Nreturned
351 case "plan_summary":
352 row[i] = doc.PlanSummary
353 case "client":
354 row[i] = doc.Client
355 case "user":
356 row[i] = doc.User
357 case "docs_deleted":
358 row[i] = doc.Ndeleted
359 case "docs_inserted":
360 row[i] = doc.Ninserted
361 case "docs_modified":
362 row[i] = doc.NModified
363 case "response_length":
364 row[i] = doc.ResponseLength
365 case "num_yield":
366 row[i] = doc.NumYield
367 case "app_name":
368 row[i] = doc.AppName
369 case "cursor_exhausted":
370 row[i] = optionalBool(doc.CursorExhausted)
371 case "has_sort_stage":
372 row[i] = optionalBool(doc.HasSortStage)
373 case "uses_disk":
374 row[i] = optionalBool(doc.UsedDisk)
375 case "from_multi_planner":
376 row[i] = optionalBool(doc.FromMultiPlanner)
377 case "replanned":
378 row[i] = optionalBool(doc.Replanned)
379 case "query_hash":
380 row[i] = doc.QueryHash
381 case "plan_cache_key":
382 row[i] = doc.PlanCacheKey
383 case "planning_time":
384 row[i] = optionalDuration(doc.PlanningTimeMicros, 1000000.0) // us to seconds
385 case "cpu_time":
386 row[i] = optionalDuration(doc.CpuNanos, 1000000000.0) // ns to seconds
387 case "query_framework":
388 row[i] = doc.QueryFramework
389 case "query_shape_hash":
390 row[i] = doc.QueryShapeHash
391 default:
392 row[i] = nil
393 }
394 }
395
396 data = append(data, row)
397 }
398
399 return &funcapi.FunctionResponse{
400 Status: 200,
401 Help: topQueriesHelpText,
402 Columns: cs.BuildColumns(),
403 Data: data,
404 DefaultSortColumn: "execution_time",
405 RequiredParams: []funcapi.ParamConfig{sortParam},
406 ChartingConfig: cs.BuildCharting(),
407 }
408 }
409
410 func (f *funcTopQueries) columnSet(cols []topQueriesColumn) funcapi.ColumnSet[topQueriesColumn] {
411 return funcapi.Columns(cols, func(c topQueriesColumn) funcapi.ColumnMeta { return c.ColumnMeta })
412 }
413
414 func (f *funcTopQueries) getDatabases() ([]string, error) {
415 databases, err := f.router.collector.conn.listDatabaseNames()
416 if err != nil {
417 return nil, err
418 }
419
420 var filteredDBs []string
421 for _, dbName := range databases {
422 if dbName == "admin" || dbName == "local" || dbName == "config" {
423 continue
424 }
425 if f.router.collector.dbSelector != nil && !f.router.collector.dbSelector.MatchString(dbName) {
426 continue
427 }
428 filteredDBs = append(filteredDBs, dbName)
429 }
430
431 return filteredDBs, nil
432 }
433
434 // detectProfileFields detects available fields in system.profile using double-checked locking
435 func (f *funcTopQueries) detectProfileFields(ctx context.Context, databases []string) (map[string]bool, error) {
436 // Fast path: return cached
437 f.router.collector.topQueriesColsMu.RLock()
438 if f.router.collector.topQueriesCols != nil {
439 cols := f.router.collector.topQueriesCols
440 f.router.collector.topQueriesColsMu.RUnlock()
441 return cols, nil
442 }
443 f.router.collector.topQueriesColsMu.RUnlock()
444
445 // Slow path: detect and cache
446 f.router.collector.topQueriesColsMu.Lock()
447 defer f.router.collector.topQueriesColsMu.Unlock()
448
449 // Double-check after acquiring write lock
450 if f.router.collector.topQueriesCols != nil {
451 return f.router.collector.topQueriesCols, nil
452 }
453
454 client, ok := f.router.collector.conn.(*mongoClient)
455 if !ok || client == nil || client.client == nil {
456 return nil, fmt.Errorf("client not initialized")
457 }
458
459 available := make(map[string]bool)
460
461 // Always include core fields that are guaranteed to exist
462 coreFields := []string{"ts", "op", "ns", "command", "millis"}
463 for _, fld := range coreFields {
464 available[fld] = true
465 }
466
467 // Sample documents from system.profile to detect available fields
468 for _, dbName := range databases {
469 queryCtx, cancel := context.WithTimeout(ctx, f.router.collector.topQueriesTimeout())
470
471 collection := client.client.Database(dbName).Collection("system.profile")
472 opts := options.FindOne().SetSort(bson.D{{Key: "$natural", Value: -1}})
473
474 var doc bson.M
475 err := collection.FindOne(queryCtx, bson.M{}, opts).Decode(&doc)
476 cancel()
477
478 if err != nil {
479 continue // No documents or profiling disabled
480 }
481
482 // Add all fields found in this document
483 for field := range doc {
484 available[field] = true
485 }
486 }
487
488 f.router.collector.topQueriesCols = available
489 return available, nil
490 }
491
492 // buildAvailableColumns returns columns that are available based on detected fields.
493 func (f *funcTopQueries) buildAvailableColumns(available map[string]bool) []topQueriesColumn {
494 var result []topQueriesColumn
495 for _, col := range topQueriesColumns {
496 // command field maps to query column, always include
497 if col.DBField == "command" || available[col.DBField] {
498 result = append(result, col)
499 }
500 }
501 return result
502 }
503
504 // querySystemProfile queries the system.profile collection for a specific database
505 func (f *funcTopQueries) querySystemProfile(ctx context.Context, dbName, sortColumn string, limit int) ([]topQueriesProfileDocument, bool, error) {
506 client, ok := f.router.collector.conn.(*mongoClient)
507 if !ok || client == nil || client.client == nil {
508 return nil, false, fmt.Errorf("client not initialized")
509 }
510
511 queryCtx, cancel := context.WithTimeout(ctx, f.router.collector.topQueriesTimeout())
512 defer cancel()
513
514 // Check if profiling is enabled for this database
515 var profilingStatus struct {
516 Was int `bson:"was"`
517 }
518 err := client.client.Database(dbName).RunCommand(queryCtx, bson.D{{Key: "profile", Value: -1}}).Decode(&profilingStatus)
519 if err != nil {
520 return nil, false, fmt.Errorf("failed to check profiling status: %w", err)
521 }
522
523 if profilingStatus.Was == 0 {
524 return nil, false, nil // Profiling disabled
525 }
526
527 // Query system.profile
528 collection := client.client.Database(dbName).Collection("system.profile")
529
530 // Build sort order (descending for all except timestamp which can be either)
531 sortOrder := -1 // descending by default
532 sortField := sortColumn
533
534 findOpts := options.Find().
535 SetSort(bson.D{{Key: sortField, Value: sortOrder}}).
536 SetLimit(int64(limit))
537
538 cursor, err := collection.Find(queryCtx, bson.M{}, findOpts)
539 if err != nil {
540 return nil, true, fmt.Errorf("find failed: %w", err)
541 }
542 defer cursor.Close(queryCtx)
543
544 var docs []topQueriesProfileDocument
545 if err := cursor.All(queryCtx, &docs); err != nil {
546 return nil, true, fmt.Errorf("cursor.All failed: %w", err)
547 }
548
549 return docs, true, nil
550 }
551
552 // sortDocuments sorts documents in place by the specified column (descending)
553 func (f *funcTopQueries) sortDocuments(docs []topQueriesProfileDocument, sortColumn string) {
554 sort.Slice(docs, func(i, j int) bool {
555 switch sortColumn {
556 case "millis":
557 return docs[i].Millis > docs[j].Millis
558 case "docsExamined":
559 return docs[i].DocsExamined > docs[j].DocsExamined
560 case "keysExamined":
561 return docs[i].KeysExamined > docs[j].KeysExamined
562 case "nreturned":
563 return docs[i].Nreturned > docs[j].Nreturned
564 case "ts":
565 return docs[i].Timestamp.After(docs[j].Timestamp)
566 case "ndeleted":
567 return docs[i].Ndeleted > docs[j].Ndeleted
568 case "ninserted":
569 return docs[i].Ninserted > docs[j].Ninserted
570 case "nModified":
571 return docs[i].NModified > docs[j].NModified
572 case "responseLength":
573 return docs[i].ResponseLength > docs[j].ResponseLength
574 case "numYield":
575 return docs[i].NumYield > docs[j].NumYield
576 case "planningTimeMicros":
577 vi := int64(0)
578 vj := int64(0)
579 if docs[i].PlanningTimeMicros != nil {
580 vi = *docs[i].PlanningTimeMicros
581 }
582 if docs[j].PlanningTimeMicros != nil {
583 vj = *docs[j].PlanningTimeMicros
584 }
585 return vi > vj
586 case "cpuNanos":
587 vi := int64(0)
588 vj := int64(0)
589 if docs[i].CpuNanos != nil {
590 vi = *docs[i].CpuNanos
591 }
592 if docs[j].CpuNanos != nil {
593 vj = *docs[j].CpuNanos
594 }
595 return vi > vj
596 default:
597 return docs[i].Millis > docs[j].Millis
598 }
599 })
600 }
601
602 // optionalDuration converts an optional int64 pointer to float64 seconds, returning nil if nil.
603 func optionalDuration(v *int64, divisor float64) any {
604 if v == nil {
605 return nil
606 }
607 return float64(*v) / divisor
608 }
609
610 // optionalBool converts a bool pointer to a display value.
611 func optionalBool(v *bool) any {
612 if v == nil {
613 return nil
614 }
615 if *v {
616 return "Yes"
617 }
618 return "No"
619 }