master
go 820 lines 43.1 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package postgres
4
5 import (
6 "context"
7 "fmt"
8 "strings"
9
10 "github.com/netdata/netdata/go/plugins/pkg/funcapi"
11 "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/sqlquery"
12 "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/strmutil"
13 )
14
15 const (
16 topQueriesMethodID = "top-queries"
17 maxQueryTextLength = 4096
18 paramSort = "__sort"
19 )
20
21 // queryStatsSourceName is the type for query stats source
22 type queryStatsSourceName string
23
24 const (
25 queryStatsSourcePgStatMonitor queryStatsSourceName = "pg_stat_monitor"
26 queryStatsSourcePgStatStatements queryStatsSourceName = "pg_stat_statements"
27 queryStatsSourceNone queryStatsSourceName = ""
28 )
29
30 // getQueryStatsSource detects and returns the best available query stats source.
31 // Prefers pg_stat_monitor if available, falls back to pg_stat_statements.
32 // Result is cached after first detection.
33 func (f *funcTopQueries) getQueryStatsSource(ctx context.Context) (queryStatsSourceName, error) {
34 c := f.router.collector
35
36 // Fast path: return cached result
37 c.pgStatStatementsMu.RLock()
38 source := c.queryStatsSource
39 c.pgStatStatementsMu.RUnlock()
40 if source != "" {
41 return queryStatsSourceName(source), nil
42 }
43
44 // Slow path: detect and cache
45 c.pgStatStatementsMu.Lock()
46 defer c.pgStatStatementsMu.Unlock()
47
48 // Double-check after acquiring write lock
49 if c.queryStatsSource != "" {
50 return queryStatsSourceName(c.queryStatsSource), nil
51 }
52
53 // Check pg_stat_monitor first (preferred)
54 var hasPgStatMonitor bool
55 query := `SELECT EXISTS(SELECT 1 FROM pg_extension WHERE extname = 'pg_stat_monitor')`
56 if err := c.db.QueryRowContext(ctx, query).Scan(&hasPgStatMonitor); err != nil {
57 return queryStatsSourceNone, fmt.Errorf("failed to check pg_stat_monitor: %v", err)
58 }
59 if hasPgStatMonitor {
60 c.queryStatsSource = string(queryStatsSourcePgStatMonitor)
61 c.pgStatMonitorAvail = true
62 return queryStatsSourcePgStatMonitor, nil
63 }
64
65 // Fall back to pg_stat_statements
66 var hasPgStatStatements bool
67 query = `SELECT EXISTS(SELECT 1 FROM pg_extension WHERE extname = 'pg_stat_statements')`
68 if err := c.db.QueryRowContext(ctx, query).Scan(&hasPgStatStatements); err != nil {
69 return queryStatsSourceNone, fmt.Errorf("failed to check pg_stat_statements: %v", err)
70 }
71 if hasPgStatStatements {
72 c.queryStatsSource = string(queryStatsSourcePgStatStatements)
73 c.pgStatStatementsAvail = true
74 return queryStatsSourcePgStatStatements, nil
75 }
76
77 return queryStatsSourceNone, nil
78 }
79
80 // detectPgStatStatementsColumns queries the database to find available columns
81 func (f *funcTopQueries) detectPgStatStatementsColumns(ctx context.Context) (map[string]bool, error) {
82 c := f.router.collector
83
84 // Fast path: return cached result
85 c.pgStatStatementsMu.RLock()
86 if c.pgStatStatementsColumns != nil {
87 cols := c.pgStatStatementsColumns
88 c.pgStatStatementsMu.RUnlock()
89 return cols, nil
90 }
91 c.pgStatStatementsMu.RUnlock()
92
93 // Slow path: query and cache
94 c.pgStatStatementsMu.Lock()
95 defer c.pgStatStatementsMu.Unlock()
96
97 // Double-check after acquiring write lock
98 if c.pgStatStatementsColumns != nil {
99 return c.pgStatStatementsColumns, nil
100 }
101
102 cols, err := sqlquery.FetchTableColumns(
103 ctx,
104 c.db,
105 "public",
106 "pg_stat_statements",
107 sqlquery.PlaceholderDollar,
108 nil,
109 )
110 if err != nil {
111 return nil, fmt.Errorf("failed to query columns: %v", err)
112 }
113
114 // Cache the result
115 c.pgStatStatementsColumns = cols
116 return cols, nil
117 }
118
119 // detectPgStatMonitorColumns queries the database to find available columns
120 func (f *funcTopQueries) detectPgStatMonitorColumns(ctx context.Context) (map[string]bool, error) {
121 c := f.router.collector
122
123 // Fast path: return cached result
124 c.pgStatStatementsMu.RLock()
125 if c.pgStatMonitorColumns != nil {
126 cols := c.pgStatMonitorColumns
127 c.pgStatStatementsMu.RUnlock()
128 return cols, nil
129 }
130 c.pgStatStatementsMu.RUnlock()
131
132 // Slow path: query and cache
133 c.pgStatStatementsMu.Lock()
134 defer c.pgStatStatementsMu.Unlock()
135
136 // Double-check after acquiring write lock
137 if c.pgStatMonitorColumns != nil {
138 return c.pgStatMonitorColumns, nil
139 }
140
141 cols, err := sqlquery.FetchTableColumns(
142 ctx,
143 c.db,
144 "public",
145 "pg_stat_monitor",
146 sqlquery.PlaceholderDollar,
147 nil,
148 )
149 if err != nil {
150 return nil, fmt.Errorf("failed to query columns: %v", err)
151 }
152
153 // Cache the result
154 c.pgStatMonitorColumns = cols
155 return cols, nil
156 }
157
158 // pgColumn defines metadata for a pg_stat_statements/pg_stat_monitor column.
159 // Embeds funcapi.ColumnMeta for UI rendering and adds PG-specific fields.
160 type pgColumn struct {
161 funcapi.ColumnMeta
162
163 // DBColumn is the database column expression (e.g., "s.queryid::text", "d.datname")
164 DBColumn string
165 // IsSortOption indicates whether this column appears in the sort dropdown
166 IsSortOption bool
167 // SortLabel is the label shown in the sort dropdown (if IsSortOption)
168 SortLabel string
169 // IsDefaultSort indicates whether this is the default sort column
170 IsDefaultSort bool
171 // OnlyPgStatMonitor indicates this column only exists in pg_stat_monitor
172 OnlyPgStatMonitor bool
173 }
174
175 // pgColumnSet creates a ColumnSet from a slice of pgColumn.
176 func pgColumnSet(cols []pgColumn) funcapi.ColumnSet[pgColumn] {
177 return funcapi.Columns(cols, func(c pgColumn) funcapi.ColumnMeta { return c.ColumnMeta })
178 }
179
180 // pgAllColumns defines ALL possible columns from pg_stat_statements.
181 // Order matters - this determines column index in the response.
182 var pgAllColumns = []pgColumn{
183 // Core identification columns (always present)
184 {ColumnMeta: funcapi.ColumnMeta{Name: "queryid", Tooltip: "Query ID", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount, Filter: funcapi.FieldFilterMultiselect, UniqueKey: true, Sortable: true}, DBColumn: "s.queryid::text"},
185 {ColumnMeta: funcapi.ColumnMeta{Name: "query", Tooltip: "Query", Type: funcapi.FieldTypeString, Visible: true, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount, Filter: funcapi.FieldFilterMultiselect, Sticky: true, FullWidth: true, Sortable: true}, DBColumn: "s.query"},
186 {ColumnMeta: funcapi.ColumnMeta{Name: "database", Tooltip: "Database", Type: funcapi.FieldTypeString, Visible: true, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount, Filter: funcapi.FieldFilterMultiselect, Sortable: true}, DBColumn: "d.datname"},
187 {ColumnMeta: funcapi.ColumnMeta{Name: "user", Tooltip: "User", Type: funcapi.FieldTypeString, Visible: true, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount, Filter: funcapi.FieldFilterMultiselect, Sortable: true}, DBColumn: "u.usename"},
188
189 // Execution count (always present)
190 {ColumnMeta: funcapi.ColumnMeta{Name: "calls", Tooltip: "Calls", Type: funcapi.FieldTypeInteger, Visible: true, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.calls", IsSortOption: true, SortLabel: "Number of Calls"},
191
192 // Execution time columns (names vary by version - detected dynamically)
193 // PG <13: total_time, mean_time, min_time, max_time, stddev_time
194 // PG 13+: total_exec_time, mean_exec_time, min_exec_time, max_exec_time, stddev_exec_time
195 {ColumnMeta: funcapi.ColumnMeta{Name: "totalTime", Tooltip: "Total Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: true, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "total_time", IsSortOption: true, SortLabel: "Total Execution Time", IsDefaultSort: true},
196 {ColumnMeta: funcapi.ColumnMeta{Name: "meanTime", Tooltip: "Mean Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: true, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "mean_time", IsSortOption: true, SortLabel: "Average Execution Time"},
197 {ColumnMeta: funcapi.ColumnMeta{Name: "minTime", Tooltip: "Min Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMin, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "min_time"},
198 {ColumnMeta: funcapi.ColumnMeta{Name: "maxTime", Tooltip: "Max Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "max_time"},
199 {ColumnMeta: funcapi.ColumnMeta{Name: "stddevTime", Tooltip: "Stddev Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "stddev_time"},
200
201 // Planning time columns (PG 13+ only)
202 {ColumnMeta: funcapi.ColumnMeta{Name: "plans", Tooltip: "Plans", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.plans"},
203 {ColumnMeta: funcapi.ColumnMeta{Name: "totalPlanTime", Tooltip: "Total Plan Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "total_plan_time"},
204 {ColumnMeta: funcapi.ColumnMeta{Name: "meanPlanTime", Tooltip: "Mean Plan Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "mean_plan_time"},
205 {ColumnMeta: funcapi.ColumnMeta{Name: "minPlanTime", Tooltip: "Min Plan Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMin, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "min_plan_time"},
206 {ColumnMeta: funcapi.ColumnMeta{Name: "maxPlanTime", Tooltip: "Max Plan Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "max_plan_time"},
207 {ColumnMeta: funcapi.ColumnMeta{Name: "stddevPlanTime", Tooltip: "Stddev Plan Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "stddev_plan_time"},
208
209 // Row count (always present)
210 {ColumnMeta: funcapi.ColumnMeta{Name: "rows", Tooltip: "Rows", Type: funcapi.FieldTypeInteger, Visible: true, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.rows", IsSortOption: true, SortLabel: "Rows Returned"},
211
212 // Shared buffer statistics (always present)
213 {ColumnMeta: funcapi.ColumnMeta{Name: "sharedBlksHit", Tooltip: "Shared Blocks Hit", Type: funcapi.FieldTypeInteger, Visible: true, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.shared_blks_hit", IsSortOption: true, SortLabel: "Shared Blocks Hit (Cache)"},
214 {ColumnMeta: funcapi.ColumnMeta{Name: "sharedBlksRead", Tooltip: "Shared Blocks Read", Type: funcapi.FieldTypeInteger, Visible: true, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.shared_blks_read", IsSortOption: true, SortLabel: "Shared Blocks Read (Disk I/O)"},
215 {ColumnMeta: funcapi.ColumnMeta{Name: "sharedBlksDirtied", Tooltip: "Shared Blocks Dirtied", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.shared_blks_dirtied"},
216 {ColumnMeta: funcapi.ColumnMeta{Name: "sharedBlksWritten", Tooltip: "Shared Blocks Written", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.shared_blks_written"},
217
218 // Local buffer statistics (always present)
219 {ColumnMeta: funcapi.ColumnMeta{Name: "localBlksHit", Tooltip: "Local Blocks Hit", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.local_blks_hit"},
220 {ColumnMeta: funcapi.ColumnMeta{Name: "localBlksRead", Tooltip: "Local Blocks Read", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.local_blks_read"},
221 {ColumnMeta: funcapi.ColumnMeta{Name: "localBlksDirtied", Tooltip: "Local Blocks Dirtied", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.local_blks_dirtied"},
222 {ColumnMeta: funcapi.ColumnMeta{Name: "localBlksWritten", Tooltip: "Local Blocks Written", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.local_blks_written"},
223
224 // Temp buffer statistics (always present)
225 {ColumnMeta: funcapi.ColumnMeta{Name: "tempBlksRead", Tooltip: "Temp Blocks Read", Type: funcapi.FieldTypeInteger, Visible: true, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.temp_blks_read"},
226 {ColumnMeta: funcapi.ColumnMeta{Name: "tempBlksWritten", Tooltip: "Temp Blocks Written", Type: funcapi.FieldTypeInteger, Visible: true, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.temp_blks_written", IsSortOption: true, SortLabel: "Temp Blocks Written"},
227
228 // I/O timing (requires track_io_timing, always present but may be 0)
229 {ColumnMeta: funcapi.ColumnMeta{Name: "blkReadTime", Tooltip: "Block Read Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: true, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.blk_read_time"},
230 {ColumnMeta: funcapi.ColumnMeta{Name: "blkWriteTime", Tooltip: "Block Write Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: true, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.blk_write_time"},
231
232 // WAL statistics (PG 13+ only)
233 {ColumnMeta: funcapi.ColumnMeta{Name: "walRecords", Tooltip: "WAL Records", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.wal_records"},
234 {ColumnMeta: funcapi.ColumnMeta{Name: "walFpi", Tooltip: "WAL Full Page Images", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.wal_fpi"},
235 {ColumnMeta: funcapi.ColumnMeta{Name: "walBytes", Tooltip: "WAL Bytes", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.wal_bytes"},
236
237 // JIT statistics (PG 15+ only)
238 {ColumnMeta: funcapi.ColumnMeta{Name: "jitFunctions", Tooltip: "JIT Functions", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_functions"},
239 {ColumnMeta: funcapi.ColumnMeta{Name: "jitGenerationTime", Tooltip: "JIT Generation Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_generation_time"},
240 {ColumnMeta: funcapi.ColumnMeta{Name: "jitInliningCount", Tooltip: "JIT Inlining Count", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_inlining_count"},
241 {ColumnMeta: funcapi.ColumnMeta{Name: "jitInliningTime", Tooltip: "JIT Inlining Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_inlining_time"},
242 {ColumnMeta: funcapi.ColumnMeta{Name: "jitOptimizationCount", Tooltip: "JIT Optimization Count", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_optimization_count"},
243 {ColumnMeta: funcapi.ColumnMeta{Name: "jitOptimizationTime", Tooltip: "JIT Optimization Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_optimization_time"},
244 {ColumnMeta: funcapi.ColumnMeta{Name: "jitEmissionCount", Tooltip: "JIT Emission Count", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_emission_count"},
245 {ColumnMeta: funcapi.ColumnMeta{Name: "jitEmissionTime", Tooltip: "JIT Emission Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.jit_emission_time"},
246
247 // Temp file statistics (PG 15+ only)
248 {ColumnMeta: funcapi.ColumnMeta{Name: "tempBlkReadTime", Tooltip: "Temp Block Read Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.temp_blk_read_time"},
249 {ColumnMeta: funcapi.ColumnMeta{Name: "tempBlkWriteTime", Tooltip: "Temp Block Write Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.temp_blk_write_time"},
250
251 // pg_stat_monitor-specific columns (only available with pg_stat_monitor extension)
252 {ColumnMeta: funcapi.ColumnMeta{Name: "applicationName", Tooltip: "Application Name", Type: funcapi.FieldTypeString, Visible: true, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount, Filter: funcapi.FieldFilterMultiselect, Sortable: true}, DBColumn: "s.application_name", OnlyPgStatMonitor: true},
253 {ColumnMeta: funcapi.ColumnMeta{Name: "clientIp", Tooltip: "Client IP Address", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount, Filter: funcapi.FieldFilterMultiselect, Sortable: true}, DBColumn: "s.client_ip::text", OnlyPgStatMonitor: true},
254 {ColumnMeta: funcapi.ColumnMeta{Name: "cmdType", Tooltip: "Query Type", Type: funcapi.FieldTypeString, Visible: true, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Summary: funcapi.FieldSummaryCount, Filter: funcapi.FieldFilterMultiselect, Visualization: funcapi.FieldVisualPill, Sortable: true}, DBColumn: "s.cmd_type_text", OnlyPgStatMonitor: true},
255 {ColumnMeta: funcapi.ColumnMeta{Name: "comments", Tooltip: "Query Comments", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Sortable: true}, DBColumn: "s.comments", OnlyPgStatMonitor: true},
256 {ColumnMeta: funcapi.ColumnMeta{Name: "relations", Tooltip: "Involved Tables", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Sortable: true}, DBColumn: "array_to_string(s.relations, ', ')", OnlyPgStatMonitor: true},
257 {ColumnMeta: funcapi.ColumnMeta{Name: "cpuUserTime", Tooltip: "User CPU Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.cpu_user_time", IsSortOption: true, SortLabel: "User CPU Time", OnlyPgStatMonitor: true},
258 {ColumnMeta: funcapi.ColumnMeta{Name: "cpuSysTime", Tooltip: "System CPU Time", Type: funcapi.FieldTypeDuration, Units: "milliseconds", Visible: false, Transform: funcapi.FieldTransformDuration, DecimalPoints: 2, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummarySum, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.cpu_sys_time", IsSortOption: true, SortLabel: "System CPU Time", OnlyPgStatMonitor: true},
259 {ColumnMeta: funcapi.ColumnMeta{Name: "elevel", Tooltip: "Error Level", Type: funcapi.FieldTypeInteger, Visible: false, Transform: funcapi.FieldTransformNumber, Sort: funcapi.FieldSortDescending, Summary: funcapi.FieldSummaryMax, Filter: funcapi.FieldFilterRange, Sortable: true}, DBColumn: "s.elevel", OnlyPgStatMonitor: true},
260 {ColumnMeta: funcapi.ColumnMeta{Name: "sqlcode", Tooltip: "SQL Error Code", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Filter: funcapi.FieldFilterMultiselect, Sortable: true}, DBColumn: "s.sqlcode", OnlyPgStatMonitor: true},
261 {ColumnMeta: funcapi.ColumnMeta{Name: "message", Tooltip: "Error Message", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, FullWidth: true, Sortable: true}, DBColumn: "s.message", OnlyPgStatMonitor: true},
262 {ColumnMeta: funcapi.ColumnMeta{Name: "toplevel", Tooltip: "Top-level Statement", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortAscending, Filter: funcapi.FieldFilterMultiselect, Sortable: true}, DBColumn: "s.toplevel::text", OnlyPgStatMonitor: true},
263 {ColumnMeta: funcapi.ColumnMeta{Name: "bucketStartTime", Tooltip: "Bucket Start Time", Type: funcapi.FieldTypeString, Visible: false, Transform: funcapi.FieldTransformNone, Sort: funcapi.FieldSortDescending, Sortable: true}, DBColumn: "s.bucket_start_time::text", OnlyPgStatMonitor: true},
264 }
265
266 // pgChartGroupDefs defines chart groupings for columns. These are applied at runtime via decoratePgColumns.
267 var pgChartGroupDefs = []struct {
268 key string
269 title string
270 columns []string
271 defaultChart bool
272 }{
273 {key: "Calls", title: "Number of Calls", columns: []string{"calls"}, defaultChart: true},
274 {key: "Time", title: "Execution Time", columns: []string{"totalTime", "meanTime", "minTime", "maxTime", "stddevTime"}, defaultChart: true},
275 {key: "PlanTime", title: "Planning Time", columns: []string{"totalPlanTime", "meanPlanTime", "minPlanTime", "maxPlanTime", "stddevPlanTime"}},
276 {key: "Plans", title: "Plans", columns: []string{"plans"}},
277 {key: "Rows", title: "Rows Returned", columns: []string{"rows"}},
278 {key: "SharedBlocks", title: "Shared Blocks", columns: []string{"sharedBlksHit", "sharedBlksRead", "sharedBlksDirtied", "sharedBlksWritten"}},
279 {key: "LocalBlocks", title: "Local Blocks", columns: []string{"localBlksHit", "localBlksRead", "localBlksDirtied", "localBlksWritten"}},
280 {key: "TempBlocks", title: "Temp Blocks", columns: []string{"tempBlksRead", "tempBlksWritten"}},
281 {key: "IOTime", title: "Block I/O Time", columns: []string{"blkReadTime", "blkWriteTime"}},
282 {key: "WALRecords", title: "WAL Records", columns: []string{"walRecords", "walFpi"}},
283 {key: "WALBytes", title: "WAL Bytes", columns: []string{"walBytes"}},
284 {key: "JITCounts", title: "JIT Counts", columns: []string{"jitFunctions", "jitInliningCount", "jitOptimizationCount", "jitEmissionCount"}},
285 {key: "JITTime", title: "JIT Time", columns: []string{"jitGenerationTime", "jitInliningTime", "jitOptimizationTime", "jitEmissionTime"}},
286 {key: "TempIOTime", title: "Temp Block I/O Time", columns: []string{"tempBlkReadTime", "tempBlkWriteTime"}},
287 // pg_stat_monitor-specific chart groups
288 {key: "CPUTime", title: "CPU Time", columns: []string{"cpuUserTime", "cpuSysTime"}},
289 {key: "Errors", title: "Error Info", columns: []string{"elevel", "sqlcode", "message"}},
290 }
291
292 // pgLabelColumnIDs defines which columns are available for group-by.
293 var pgLabelColumnIDs = map[string]bool{
294 "database": true,
295 "user": true,
296 "applicationName": true, // pg_stat_monitor only
297 "cmdType": true, // pg_stat_monitor only
298 }
299
300 const pgPrimaryLabelID = "database"
301
302 func topQueriesMethodConfig() funcapi.MethodConfig {
303 return funcapi.MethodConfig{
304 ID: topQueriesMethodID,
305 Name: "Top Queries",
306 UpdateEvery: 10,
307 Help: "Top SQL queries from pg_stat_statements",
308 RequireCloud: true,
309 RequiredParams: []funcapi.ParamConfig{
310 {
311 ID: paramSort,
312 Name: "Filter By",
313 Help: "Select the primary sort column",
314 Selection: funcapi.ParamSelect,
315 Options: buildPgSortOptions(),
316 UniqueView: true,
317 },
318 },
319 }
320 }
321
322 // Compile-time interface check.
323 var _ funcapi.MethodHandler = (*funcTopQueries)(nil)
324
325 // funcTopQueries handles the "top-queries" function for PostgreSQL.
326 type funcTopQueries struct {
327 router *funcRouter
328 }
329
330 func newFuncTopQueries(r *funcRouter) *funcTopQueries {
331 return &funcTopQueries{router: r}
332 }
333
334 // MethodParams implements funcapi.MethodHandler.
335 func (f *funcTopQueries) MethodParams(ctx context.Context, method string) ([]funcapi.ParamConfig, error) {
336 if f.router.collector.Functions.TopQueries.Disabled {
337 return nil, fmt.Errorf("top-queries function disabled in configuration")
338 }
339 if f.router.collector.db == nil {
340 return nil, fmt.Errorf("collector is still initializing")
341 }
342 return f.topQueriesParams(ctx)
343 }
344
345 // Handle implements funcapi.MethodHandler.
346 func (f *funcTopQueries) Handle(ctx context.Context, method string, params funcapi.ResolvedParams) *funcapi.FunctionResponse {
347 if f.router.collector.Functions.TopQueries.Disabled {
348 return funcapi.UnavailableResponse("top-queries function has been disabled in configuration")
349 }
350 if f.router.collector.db == nil {
351 return funcapi.UnavailableResponse("collector is still initializing, please retry in a few seconds")
352 }
353 queryCtx, cancel := context.WithTimeout(ctx, f.router.collector.topQueriesTimeout())
354 defer cancel()
355 return f.collectTopQueries(queryCtx, params.Column(paramSort))
356 }
357
358 // buildPgSortOptions builds sort options from pgAllColumns.
359 func buildPgSortOptions() []funcapi.ParamOption {
360 var opts []funcapi.ParamOption
361 sortDir := funcapi.FieldSortDescending
362 for _, col := range pgAllColumns {
363 if col.IsSortOption {
364 opts = append(opts, funcapi.ParamOption{
365 ID: col.Name,
366 Column: col.Name,
367 Name: "Top queries by " + col.SortLabel,
368 Default: col.IsDefaultSort,
369 Sort: &sortDir,
370 })
371 }
372 }
373 return opts
374 }
375
376 // collectTopQueries queries pg_stat_statements or pg_stat_monitor for top queries.
377 // It auto-detects pg_stat_monitor and uses it when available, falling back to pg_stat_statements.
378 func (f *funcTopQueries) collectTopQueries(ctx context.Context, sortColumn string) *funcapi.FunctionResponse {
379 c := f.router.collector
380
381 // Auto-detect best available query stats source
382 source, err := f.getQueryStatsSource(ctx)
383 if err != nil {
384 return funcapi.InternalErrorResponse("failed to detect query stats source: %v", err)
385 }
386 if source == queryStatsSourceNone {
387 return funcapi.UnavailableResponse("No query statistics extension is installed in this database. " +
388 "Install pg_stat_monitor (recommended) or pg_stat_statements:\n\n" +
389 "For pg_stat_monitor:\n" +
390 " ALTER SYSTEM SET shared_preload_libraries = 'pg_stat_monitor';\n" +
391 " -- restart PostgreSQL\n" +
392 " CREATE EXTENSION pg_stat_monitor;\n\n" +
393 "For pg_stat_statements:\n" +
394 " ALTER SYSTEM SET shared_preload_libraries = 'pg_stat_statements';\n" +
395 " -- restart PostgreSQL\n" +
396 " CREATE EXTENSION pg_stat_statements;")
397 }
398
399 // Detect available columns based on source
400 var availableCols map[string]bool
401 if source == queryStatsSourcePgStatMonitor {
402 availableCols, err = f.detectPgStatMonitorColumns(ctx)
403 } else {
404 availableCols, err = f.detectPgStatStatementsColumns(ctx)
405 }
406 if err != nil {
407 return funcapi.InternalErrorResponse("failed to detect available columns: %v", err)
408 }
409
410 // Build list of columns to query based on what's available and source
411 queryCols := f.buildAvailableColumns(availableCols, source)
412 if len(queryCols) == 0 {
413 return funcapi.InternalErrorResponse("no queryable columns found in %s", source)
414 }
415
416 // Map and validate sort column
417 actualSortCol := f.mapAndValidateSortColumn(sortColumn, availableCols, source)
418
419 // Get query limit (default 500)
420 limit := c.topQueriesLimit()
421
422 // Build and execute query
423 query := f.buildDynamicSQL(queryCols, actualSortCol, limit, source)
424 rows, err := c.db.QueryContext(ctx, query)
425 if err != nil {
426 if ctx.Err() == context.DeadlineExceeded {
427 return funcapi.ErrorResponse(504, "query timed out")
428 }
429 return funcapi.InternalErrorResponse("query failed: %v", err)
430 }
431 defer rows.Close()
432
433 // Process rows and build response
434 data, err := f.scanDynamicRows(rows, queryCols)
435 if err != nil {
436 return funcapi.InternalErrorResponse("%s", err)
437 }
438
439 if err := rows.Err(); err != nil {
440 return funcapi.InternalErrorResponse("rows iteration error: %v", err)
441 }
442
443 // Build dynamic sort options from available columns (only those actually detected)
444 sortParam, sortOptions := f.topQueriesSortParam(queryCols)
445
446 // Find default sort column from metadata
447 defaultSort := ""
448 for _, col := range queryCols {
449 if col.IsDefaultSort && col.IsSortOption {
450 defaultSort = col.Name
451 break
452 }
453 }
454 // Fallback to first sort option if no default
455 if defaultSort == "" && len(sortOptions) > 0 {
456 defaultSort = sortOptions[0].ID
457 }
458
459 // Decorate columns with chart/label metadata and build using ColumnSet
460 annotatedCols := decoratePgColumns(queryCols)
461 cs := pgColumnSet(annotatedCols)
462
463 // Build help message based on source
464 helpMsg := "Top SQL queries from pg_stat_statements"
465 if source == queryStatsSourcePgStatMonitor {
466 helpMsg = "Top SQL queries from pg_stat_monitor (includes application, client IP, CPU time, and error info)"
467 }
468
469 return &funcapi.FunctionResponse{
470 Status: 200,
471 Help: helpMsg,
472 Columns: cs.BuildColumns(),
473 Data: data,
474 DefaultSortColumn: defaultSort,
475 RequiredParams: []funcapi.ParamConfig{sortParam},
476 ChartingConfig: cs.BuildCharting(),
477 }
478 }
479
480 // decoratePgColumns adds label and chart metadata to columns for ColumnSet builders.
481 func decoratePgColumns(cols []pgColumn) []pgColumn {
482 out := make([]pgColumn, len(cols))
483 index := make(map[string]int, len(cols))
484 for i, col := range cols {
485 out[i] = col
486 index[col.Name] = i
487 }
488
489 // Mark groupby columns
490 for i := range out {
491 if pgLabelColumnIDs[out[i].Name] {
492 out[i].GroupBy = &funcapi.GroupByOptions{
493 IsDefault: out[i].Name == pgPrimaryLabelID,
494 }
495 }
496 }
497
498 // Mark chart columns
499 for _, group := range pgChartGroupDefs {
500 for _, key := range group.columns {
501 idx, ok := index[key]
502 if !ok {
503 continue
504 }
505 out[idx].Chart = &funcapi.ChartOptions{
506 Group: group.key,
507 Title: group.title,
508 IsDefault: group.defaultChart,
509 }
510 }
511 }
512
513 return out
514 }
515
516 // buildAvailableColumns returns column metadata for columns that exist in this PG version and source.
517 func (f *funcTopQueries) buildAvailableColumns(availableCols map[string]bool, source queryStatsSourceName) []pgColumn {
518 c := f.router.collector
519 var result []pgColumn
520 isPgStatMonitor := source == queryStatsSourcePgStatMonitor
521
522 for _, col := range pgAllColumns {
523 // Skip pg_stat_monitor-only columns when using pg_stat_statements
524 if col.OnlyPgStatMonitor && !isPgStatMonitor {
525 continue
526 }
527
528 // Extract the actual column name for availability check
529 colName := col.DBColumn
530
531 // Strip array_to_string wrapper FIRST (before table prefix removal)
532 // e.g., "array_to_string(s.relations, ', ')" -> "s.relations"
533 if after, ok := strings.CutPrefix(colName, "array_to_string("); ok {
534 colName = after
535 if idx := strings.Index(colName, ","); idx != -1 {
536 colName = colName[:idx]
537 }
538 }
539
540 // Remove table prefix (e.g., "s.relations" -> "relations")
541 if idx := strings.LastIndex(colName, "."); idx != -1 {
542 colName = colName[idx+1:]
543 }
544
545 // Remove PostgreSQL type cast suffix (e.g., "queryid::text" -> "queryid")
546 if idx := strings.Index(colName, "::"); idx != -1 {
547 colName = colName[:idx]
548 }
549
550 // Handle version-specific column names for time columns.
551 // pg_stat_statements PG 13+ and pg_stat_monitor both use: total_exec_time, mean_exec_time, etc.
552 // pg_stat_statements < PG 13 uses: total_time, mean_time, etc.
553 actualColName := colName
554 if isPgStatMonitor || c.pgVersion >= pgVersion13 {
555 switch colName {
556 case "total_time":
557 actualColName = "total_exec_time"
558 case "mean_time":
559 actualColName = "mean_exec_time"
560 case "min_time":
561 actualColName = "min_exec_time"
562 case "max_time":
563 actualColName = "max_exec_time"
564 case "stddev_time":
565 actualColName = "stddev_exec_time"
566 }
567 }
568
569 // Check if column exists (either directly or via join)
570 // Join columns (database, user) come from other tables (d.datname, u.usename)
571 // pg_stat_monitor has datname directly, pg_stat_statements needs join
572 isJoinCol := col.Name == "database" || col.Name == "user"
573 if isPgStatMonitor && col.Name == "database" {
574 // pg_stat_monitor has datname directly
575 isJoinCol = false
576 }
577 if isJoinCol || availableCols[actualColName] {
578 // Create a copy with the actual column name for this version
579 colCopy := col
580 if actualColName != colName {
581 // Update DBColumn to use the version-specific name with alias
582 if strings.HasPrefix(col.DBColumn, "s.") {
583 colCopy.DBColumn = "s." + actualColName
584 }
585 }
586 // For pg_stat_monitor, database comes from s.datname directly
587 if isPgStatMonitor && col.Name == "database" {
588 colCopy.DBColumn = "s.datname"
589 }
590 result = append(result, colCopy)
591 }
592 }
593
594 return result
595 }
596
597 // mapAndValidateSortColumn maps the semantic sort column to actual SQL column.
598 func (f *funcTopQueries) mapAndValidateSortColumn(sortColumn string, availableCols map[string]bool, source queryStatsSourceName) string {
599 c := f.router.collector
600 isPgStatMonitor := source == queryStatsSourcePgStatMonitor
601
602 // Map column ID back to DBColumn
603 for _, col := range pgAllColumns {
604 if col.Name == sortColumn || col.DBColumn == sortColumn {
605 // Get actual column name (strip table prefix and type cast)
606 colName := col.DBColumn
607 if idx := strings.LastIndex(colName, "."); idx != -1 {
608 colName = colName[idx+1:]
609 }
610 if idx := strings.Index(colName, "::"); idx != -1 {
611 colName = colName[:idx]
612 }
613
614 // Handle version-specific mapping for time columns.
615 // pg_stat_statements PG 13+ and pg_stat_monitor both use: total_exec_time, mean_exec_time, etc.
616 // pg_stat_statements < PG 13 uses: total_time, mean_time, etc.
617 if isPgStatMonitor || c.pgVersion >= pgVersion13 {
618 switch colName {
619 case "total_time":
620 colName = "total_exec_time"
621 case "mean_time":
622 colName = "mean_exec_time"
623 case "min_time":
624 colName = "min_exec_time"
625 case "max_time":
626 colName = "max_exec_time"
627 case "stddev_time":
628 colName = "stddev_exec_time"
629 }
630 }
631
632 // Validate column exists
633 if availableCols[colName] {
634 return colName
635 }
636 }
637 }
638
639 // Default fallback
640 if isPgStatMonitor || c.pgVersion >= pgVersion13 {
641 return "total_exec_time"
642 }
643 return "total_time"
644 }
645
646 // buildDynamicSQL builds the SQL query with only available columns.
647 func (f *funcTopQueries) buildDynamicSQL(cols []pgColumn, sortColumn string, limit int, source queryStatsSourceName) string {
648 c := f.router.collector
649 var selectCols []string
650
651 for _, col := range cols {
652 colExpr := col.DBColumn
653
654 // Handle version-specific column names for time columns.
655 // pg_stat_statements PG 13+ and pg_stat_monitor both use: total_exec_time, mean_exec_time, etc.
656 // pg_stat_statements < PG 13 uses: total_time, mean_time, etc.
657 if source == queryStatsSourcePgStatMonitor || c.pgVersion >= pgVersion13 {
658 switch {
659 case strings.HasSuffix(colExpr, ".total_time"):
660 colExpr = strings.Replace(colExpr, ".total_time", ".total_exec_time", 1)
661 case strings.HasSuffix(colExpr, ".mean_time"):
662 colExpr = strings.Replace(colExpr, ".mean_time", ".mean_exec_time", 1)
663 case strings.HasSuffix(colExpr, ".min_time"):
664 colExpr = strings.Replace(colExpr, ".min_time", ".min_exec_time", 1)
665 case strings.HasSuffix(colExpr, ".max_time"):
666 colExpr = strings.Replace(colExpr, ".max_time", ".max_exec_time", 1)
667 case strings.HasSuffix(colExpr, ".stddev_time"):
668 colExpr = strings.Replace(colExpr, ".stddev_time", ".stddev_exec_time", 1)
669 case colExpr == "total_time":
670 colExpr = "total_exec_time"
671 case colExpr == "mean_time":
672 colExpr = "mean_exec_time"
673 case colExpr == "min_time":
674 colExpr = "min_exec_time"
675 case colExpr == "max_time":
676 colExpr = "max_exec_time"
677 case colExpr == "stddev_time":
678 colExpr = "stddev_exec_time"
679 }
680 }
681
682 // Use column ID as the SQL alias for consistent naming
683 // Use double quotes to handle reserved keywords like "database", "user"
684 selectCols = append(selectCols, fmt.Sprintf("%s AS \"%s\"", colExpr, col.Name))
685 }
686
687 // Build query based on source
688 if source == queryStatsSourcePgStatMonitor {
689 // pg_stat_monitor has datname and username columns directly
690 return fmt.Sprintf(`
691 SELECT %s
692 FROM pg_stat_monitor s
693 JOIN pg_user u ON s.userid = u.usesysid
694 ORDER BY "%s" DESC
695 LIMIT %d
696 `, strings.Join(selectCols, ", "), sortColumn, limit)
697 }
698
699 // pg_stat_statements needs joins for database and user names
700 return fmt.Sprintf(`
701 SELECT %s
702 FROM pg_stat_statements s
703 JOIN pg_database d ON s.dbid = d.oid
704 JOIN pg_user u ON s.userid = u.usesysid
705 ORDER BY "%s" DESC
706 LIMIT %d
707 `, strings.Join(selectCols, ", "), sortColumn, limit)
708 }
709
710 // scanDynamicRows scans rows into the data array based on column types.
711 // Uses sql.Null* types to handle NULL values safely.
712 func (f *funcTopQueries) scanDynamicRows(rows dbRows, cols []pgColumn) ([][]any, error) {
713 specs := make([]sqlquery.ScanColumnSpec, len(cols))
714 for i, col := range cols {
715 specs[i] = pgTopQueriesScanSpec(col)
716 }
717
718 data, err := sqlquery.ScanTypedRows(rows, specs)
719 if err != nil {
720 return nil, fmt.Errorf("row scan failed: %v", err)
721 }
722 return data, nil
723 }
724
725 func pgTopQueriesScanSpec(col pgColumn) sqlquery.ScanColumnSpec {
726 spec := sqlquery.ScanColumnSpec{}
727 switch col.Type {
728 case funcapi.FieldTypeString:
729 spec.Type = sqlquery.ScanValueString
730 if col.Name == "query" {
731 spec.Transform = func(v any) any {
732 s, _ := v.(string)
733 return strmutil.TruncateText(s, maxQueryTextLength)
734 }
735 }
736 case funcapi.FieldTypeInteger:
737 spec.Type = sqlquery.ScanValueInteger
738 case funcapi.FieldTypeFloat, funcapi.FieldTypeDuration:
739 spec.Type = sqlquery.ScanValueFloat
740 default:
741 spec.Type = sqlquery.ScanValueString
742 }
743 return spec
744 }
745
746 // buildDynamicSortOptions builds sort options from available columns.
747 // Returns only sort options for columns that actually exist in the database.
748 func (f *funcTopQueries) buildDynamicSortOptions(cols []pgColumn) []funcapi.ParamOption {
749 var sortOpts []funcapi.ParamOption
750 seen := make(map[string]bool)
751 sortDir := funcapi.FieldSortDescending
752
753 for _, col := range cols {
754 if col.IsSortOption && !seen[col.Name] {
755 seen[col.Name] = true
756 sortOpts = append(sortOpts, funcapi.ParamOption{
757 ID: col.Name,
758 Column: col.Name,
759 Name: col.SortLabel,
760 Default: col.IsDefaultSort,
761 Sort: &sortDir,
762 })
763 }
764 }
765 return sortOpts
766 }
767
768 func (f *funcTopQueries) topQueriesSortParam(queryCols []pgColumn) (funcapi.ParamConfig, []funcapi.ParamOption) {
769 sortOptions := f.buildDynamicSortOptions(queryCols)
770 sortParam := funcapi.ParamConfig{
771 ID: paramSort,
772 Name: "Filter By",
773 Help: "Select the primary sort column",
774 Selection: funcapi.ParamSelect,
775 Options: sortOptions,
776 UniqueView: true,
777 }
778 return sortParam, sortOptions
779 }
780
781 func (f *funcTopQueries) topQueriesParams(ctx context.Context) ([]funcapi.ParamConfig, error) {
782 // Auto-detect best available query stats source
783 source, err := f.getQueryStatsSource(ctx)
784 if err != nil {
785 return nil, err
786 }
787 if source == queryStatsSourceNone {
788 return nil, fmt.Errorf("no query statistics extension is installed (pg_stat_monitor or pg_stat_statements)")
789 }
790
791 // Detect available columns based on source
792 var availableCols map[string]bool
793 if source == queryStatsSourcePgStatMonitor {
794 availableCols, err = f.detectPgStatMonitorColumns(ctx)
795 } else {
796 availableCols, err = f.detectPgStatStatementsColumns(ctx)
797 }
798 if err != nil {
799 return nil, err
800 }
801
802 queryCols := f.buildAvailableColumns(availableCols, source)
803 if len(queryCols) == 0 {
804 return nil, fmt.Errorf("no queryable columns found in %s", source)
805 }
806
807 sortParam, _ := f.topQueriesSortParam(queryCols)
808
809 return []funcapi.ParamConfig{sortParam}, nil
810 }
811
812 // Cleanup implements funcapi.MethodHandler.
813 func (f *funcTopQueries) Cleanup(ctx context.Context) {}
814
815 // dbRows interface for testing
816 type dbRows interface {
817 Next() bool
818 Scan(dest ...any) error
819 Err() error
820 }