master
go 335 lines 7.98 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package postgres
4
5 import (
6 "context"
7 "database/sql"
8 "fmt"
9 "regexp"
10 "strconv"
11 "time"
12
13 "github.com/jackc/pgx/v5"
14 "github.com/jackc/pgx/v5/stdlib"
15 )
16
17 const (
18 pgVersion94 = 9_04_00
19 pgVersion10 = 10_00_00
20 pgVersion11 = 11_00_00
21 pgVersion13 = 13_00_00
22 pgVersion14 = 14_00_00
23 pgVersion17 = 17_00_00
24 )
25
26 func (c *Collector) collect() (map[string]int64, error) {
27 if c.db == nil {
28 db, err := c.openPrimaryConnection()
29 if err != nil {
30 return nil, err
31 }
32 c.db = db
33 }
34
35 if c.pgVersion == 0 {
36 ver, err := c.doQueryServerVersion()
37 if err != nil {
38 return nil, fmt.Errorf("querying server version error: %v", err)
39 }
40 c.pgVersion = ver
41 c.Debugf("connected to PostgreSQL v%d", c.pgVersion)
42 }
43
44 if c.superUser == nil {
45 v, err := c.doQueryIsSuperUser()
46 if err != nil {
47 return nil, fmt.Errorf("querying is super user error: %v", err)
48 }
49 c.superUser = &v
50 c.Debugf("connected as super user: %v", *c.superUser)
51 }
52
53 if c.canExecutePgLsDir == nil && c.pgVersion >= pgVersion10 {
54 v, err := c.doQueryCanExecutePgLsDir()
55 if err != nil {
56 return nil, fmt.Errorf("querying can execute pg_ls_dir() error: %v", err)
57 }
58 c.canExecutePgLsDir = &v
59 c.Debugf("can execute pg_ls_dir(): %v", *c.canExecutePgLsDir)
60 }
61
62 if c.pgIsInRecovery == nil {
63 v, err := c.doQueryPGIsInRecovery()
64 if err != nil {
65 return nil, fmt.Errorf("querying recovery status error: %v", err)
66 }
67 c.pgIsInRecovery = &v
68 c.Debugf("the instance is in recovery mode: %v", *c.pgIsInRecovery)
69 }
70
71 now := time.Now()
72
73 if now.Sub(c.recheckSettingsTime) > c.recheckSettingsEvery {
74 c.recheckSettingsTime = now
75 maxConn, err := c.doQuerySettingsMaxConnections()
76 if err != nil {
77 return nil, fmt.Errorf("querying settings max connections error: %v", err)
78 }
79 c.mx.maxConnections = maxConn
80
81 maxLocks, err := c.doQuerySettingsMaxLocksHeld()
82 if err != nil {
83 return nil, fmt.Errorf("querying settings max locks held error: %v", err)
84 }
85 c.mx.maxLocksHeld = maxLocks
86 }
87
88 c.resetMetrics()
89
90 if c.pgVersion >= pgVersion10 {
91 // need 'backend_type' in pg_stat_activity
92 c.addXactQueryRunningTimeChartsOnce.Do(func() {
93 c.addTransactionsRunTimeHistogramChart()
94 c.addQueriesRunTimeHistogramChart()
95 })
96 }
97 if c.isSuperUser() {
98 c.addWALFilesChartsOnce.Do(c.addWALFilesCharts)
99 }
100
101 if err := c.doQueryGlobalMetrics(); err != nil {
102 return nil, err
103 }
104 if err := c.doQueryReplicationMetrics(); err != nil {
105 return nil, err
106 }
107 if err := c.doQueryDatabasesMetrics(); err != nil {
108 return nil, err
109 }
110 if c.dbSr != nil {
111 if err := c.doQueryQueryableDatabases(); err != nil {
112 return nil, err
113 }
114 }
115 if err := c.doQueryTablesMetrics(); err != nil {
116 return nil, err
117 }
118 if err := c.doQueryIndexesMetrics(); err != nil {
119 return nil, err
120 }
121
122 if now.Sub(c.doSlowTime) > c.doSlowEvery {
123 c.doSlowTime = now
124 if err := c.doQueryBloat(); err != nil {
125 return nil, err
126 }
127 if err := c.doQueryColumns(); err != nil {
128 return nil, err
129 }
130 }
131
132 mx := make(map[string]int64)
133 c.collectMetrics(mx)
134
135 return mx, nil
136 }
137
138 func (c *Collector) openPrimaryConnection() (*sql.DB, error) {
139 if c.CloudAuth.IsEnabled() {
140 cfg, err := pgx.ParseConfig(c.DSN)
141 if err != nil {
142 return nil, fmt.Errorf("error on parsing DSN [%s]: %v", c.DSN, err)
143 }
144 return c.openAzureADConnection(cfg, "Postgres database")
145 }
146
147 db, err := sql.Open("pgx", c.DSN)
148 if err != nil {
149 return nil, fmt.Errorf("error on opening a connection with the Postgres database [%s]: %v", c.DSN, err)
150 }
151
152 db.SetMaxOpenConns(1)
153 db.SetMaxIdleConns(1)
154 db.SetConnMaxLifetime(10 * time.Minute)
155
156 ctx, cancel := context.WithTimeout(context.Background(), c.Timeout.Duration())
157 defer cancel()
158
159 if err := db.PingContext(ctx); err != nil {
160 _ = db.Close()
161 return nil, fmt.Errorf("error on pinging the Postgres database [%s]: %v", c.DSN, err)
162 }
163
164 return db, nil
165 }
166
167 func (c *Collector) openSecondaryConnection(dbname string) (*sql.DB, string, error) {
168 cfg, err := pgx.ParseConfig(c.DSN)
169 if err != nil {
170 return nil, "", fmt.Errorf("error on parsing DSN [%s]: %v", c.DSN, err)
171 }
172
173 cfg.Database = dbname
174
175 if c.CloudAuth.IsEnabled() {
176 db, err := c.openAzureADConnection(cfg, fmt.Sprintf("secondary Postgres database [%s]", dbname))
177 return db, "", err
178 }
179
180 connStr := stdlib.RegisterConnConfig(cfg)
181
182 db, err := sql.Open("pgx", connStr)
183 if err != nil {
184 stdlib.UnregisterConnConfig(connStr)
185 return nil, "", fmt.Errorf("error on opening a secondary connection with the Postgres database [%s]: %v", dbname, err)
186 }
187
188 db.SetMaxOpenConns(1)
189 db.SetMaxIdleConns(1)
190 db.SetConnMaxLifetime(10 * time.Minute)
191
192 ctx, cancel := context.WithTimeout(context.Background(), c.Timeout.Duration())
193 defer cancel()
194
195 if err := db.PingContext(ctx); err != nil {
196 stdlib.UnregisterConnConfig(connStr)
197 _ = db.Close()
198 return nil, "", fmt.Errorf("error on pinging the secondary Postgres database [%s]: %v", dbname, err)
199 }
200
201 return db, connStr, nil
202 }
203
204 func (c *Collector) openAzureADConnection(cfg *pgx.ConnConfig, target string) (*sql.DB, error) {
205 if c.azureTokenProvider == nil {
206 return nil, fmt.Errorf("cloud auth token provider is not initialized for %s", target)
207 }
208
209 db := stdlib.OpenDB(*cfg, stdlib.OptionBeforeConnect(c.azureADBeforeConnect))
210
211 db.SetMaxOpenConns(1)
212 db.SetMaxIdleConns(1)
213 db.SetConnMaxLifetime(10 * time.Minute)
214
215 ctx, cancel := context.WithTimeout(context.Background(), c.Timeout.Duration())
216 defer cancel()
217
218 if err := db.PingContext(ctx); err != nil {
219 _ = db.Close()
220 return nil, fmt.Errorf("error on pinging the %s: %v", target, err)
221 }
222
223 return db, nil
224 }
225
226 func (c *Collector) azureADBeforeConnect(ctx context.Context, cfg *pgx.ConnConfig) error {
227 token, _, err := c.azureTokenProvider.Token(ctx)
228 if err != nil {
229 return err
230 }
231 cfg.Password = token
232 return nil
233 }
234
235 func (c *Collector) isSuperUser() bool { return c.superUser != nil && *c.superUser }
236
237 func (c *Collector) canQueryReplicationSlotFiles() bool {
238 return c.isSuperUser() || (c.canExecutePgLsDir != nil && *c.canExecutePgLsDir)
239 }
240
241 func (c *Collector) isPGInRecovery() bool { return c.pgIsInRecovery != nil && *c.pgIsInRecovery }
242
243 func (c *Collector) getDBMetrics(name string) *dbMetrics {
244 db, ok := c.mx.dbs[name]
245 if !ok {
246 db = &dbMetrics{name: name}
247 c.mx.dbs[name] = db
248 }
249 return db
250 }
251
252 func (c *Collector) getTableMetrics(name, db, schema string) *tableMetrics {
253 key := name + "_" + db + "_" + schema
254 m, ok := c.mx.tables[key]
255 if !ok {
256 m = &tableMetrics{db: db, schema: schema, name: name}
257 c.mx.tables[key] = m
258 }
259 return m
260 }
261
262 func (c *Collector) hasTableMetrics(name, db, schema string) bool {
263 key := name + "_" + db + "_" + schema
264 _, ok := c.mx.tables[key]
265 return ok
266 }
267
268 func (c *Collector) getIndexMetrics(name, table, db, schema string) *indexMetrics {
269 key := name + "_" + table + "_" + db + "_" + schema
270 m, ok := c.mx.indexes[key]
271 if !ok {
272 m = &indexMetrics{name: name, db: db, schema: schema, table: table}
273 c.mx.indexes[key] = m
274 }
275 return m
276 }
277
278 func (c *Collector) hasIndexMetrics(name, table, db, schema string) bool {
279 key := name + "_" + table + "_" + db + "_" + schema
280 _, ok := c.mx.indexes[key]
281 return ok
282 }
283
284 func (c *Collector) getReplAppMetrics(name string) *replStandbyAppMetrics {
285 app, ok := c.mx.replApps[name]
286 if !ok {
287 app = &replStandbyAppMetrics{name: name}
288 c.mx.replApps[name] = app
289 }
290 return app
291 }
292
293 func (c *Collector) getReplSlotMetrics(name string) *replSlotMetrics {
294 slot, ok := c.mx.replSlots[name]
295 if !ok {
296 slot = &replSlotMetrics{name: name}
297 c.mx.replSlots[name] = slot
298 }
299 return slot
300 }
301
302 func parseInt(s string) int64 {
303 v, _ := strconv.ParseInt(s, 10, 64)
304 return v
305 }
306
307 func parseFloat(s string) int64 {
308 v, _ := strconv.ParseFloat(s, 64)
309 return int64(v)
310 }
311
312 //go:fix inline
313 func newInt(v int64) *int64 {
314 return new(v)
315 }
316
317 func calcPercentage(value, total int64) (v int64) {
318 if total == 0 {
319 return 0
320 }
321 if v = value * 100 / total; v < 0 {
322 v = -v
323 }
324 return v
325 }
326
327 func calcDeltaPercentage(a, b incDelta) int64 {
328 return calcPercentage(a.delta(), a.delta()+b.delta())
329 }
330
331 func removeSpaces(s string) string {
332 return reSpace.ReplaceAllString(s, "_")
333 }
334
335 var reSpace = regexp.MustCompile(`\s+`)