| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | package postgres |
| 4 | |
| 5 | import ( |
| 6 | "fmt" |
| 7 | ) |
| 8 | |
| 9 | func (c *Collector) doQueryReplicationMetrics() error { |
| 10 | if err := c.doQueryReplStandbyAppWALDelta(); err != nil { |
| 11 | return fmt.Errorf("querying replication standby app wal delta error: %v", err) |
| 12 | } |
| 13 | |
| 14 | if c.pgVersion >= pgVersion10 { |
| 15 | if err := c.doQueryReplStandbyAppWALLag(); err != nil { |
| 16 | return fmt.Errorf("querying replication standby app wal lag error: %v", err) |
| 17 | } |
| 18 | } |
| 19 | |
| 20 | if c.pgVersion >= pgVersion10 && c.canQueryReplicationSlotFiles() { |
| 21 | if err := c.doQueryReplSlotFiles(); err != nil { |
| 22 | return fmt.Errorf("querying replication slot files error: %v", err) |
| 23 | } |
| 24 | } |
| 25 | |
| 26 | return nil |
| 27 | } |
| 28 | |
| 29 | func (c *Collector) doQueryReplStandbyAppWALDelta() error { |
| 30 | q := queryReplicationStandbyAppDelta(c.pgVersion) |
| 31 | |
| 32 | var app string |
| 33 | return c.doQuery(q, func(column, value string, _ bool) { |
| 34 | switch column { |
| 35 | case "application_name": |
| 36 | app = value |
| 37 | c.getReplAppMetrics(app).updated = true |
| 38 | default: |
| 39 | // TODO: delta calculation was changed in https://github.com/netdata/netdata/go/plugins/plugin/go.d/pull/1039 |
| 40 | // - 'replay_delta' (probably other deltas too?) can be negative |
| 41 | // - Also, WAL delta != WAL lag after that PR |
| 42 | v := max(parseInt(value), 0) |
| 43 | switch column { |
| 44 | case "sent_delta": |
| 45 | c.getReplAppMetrics(app).walSentDelta += v |
| 46 | case "write_delta": |
| 47 | c.getReplAppMetrics(app).walWriteDelta += v |
| 48 | case "flush_delta": |
| 49 | c.getReplAppMetrics(app).walFlushDelta += v |
| 50 | case "replay_delta": |
| 51 | c.getReplAppMetrics(app).walReplayDelta += v |
| 52 | } |
| 53 | } |
| 54 | }) |
| 55 | } |
| 56 | |
| 57 | func (c *Collector) doQueryReplStandbyAppWALLag() error { |
| 58 | q := queryReplicationStandbyAppLag() |
| 59 | |
| 60 | var app string |
| 61 | return c.doQuery(q, func(column, value string, _ bool) { |
| 62 | switch column { |
| 63 | case "application_name": |
| 64 | app = value |
| 65 | c.getReplAppMetrics(app).updated = true |
| 66 | case "write_lag": |
| 67 | c.getReplAppMetrics(app).walWriteLag += parseInt(value) |
| 68 | case "flush_lag": |
| 69 | c.getReplAppMetrics(app).walFlushLag += parseInt(value) |
| 70 | case "replay_lag": |
| 71 | c.getReplAppMetrics(app).walReplayLag += parseInt(value) |
| 72 | } |
| 73 | }) |
| 74 | } |
| 75 | |
| 76 | func (c *Collector) doQueryReplSlotFiles() error { |
| 77 | q := queryReplicationSlotFiles(c.pgVersion) |
| 78 | |
| 79 | var slot string |
| 80 | return c.doQuery(q, func(column, value string, _ bool) { |
| 81 | switch column { |
| 82 | case "slot_name": |
| 83 | slot = value |
| 84 | c.getReplSlotMetrics(slot).updated = true |
| 85 | case "replslot_wal_keep": |
| 86 | c.getReplSlotMetrics(slot).walKeep += parseInt(value) |
| 87 | case "replslot_files": |
| 88 | c.getReplSlotMetrics(slot).files += parseInt(value) |
| 89 | } |
| 90 | }) |
| 91 | } |