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