1 package reaper
2
3 import (
4 "context"
5 "testing"
6 "time"
7
8 "github.com/cybertec-postgresql/pgwatch/v6/internal/cmdopts"
9 "github.com/cybertec-postgresql/pgwatch/v6/internal/db"
10 "github.com/cybertec-postgresql/pgwatch/v6/internal/log"
11 "github.com/cybertec-postgresql/pgwatch/v6/internal/metrics"
12 "github.com/cybertec-postgresql/pgwatch/v6/internal/sinks"
13 "github.com/cybertec-postgresql/pgwatch/v6/internal/sources"
14 "github.com/cybertec-postgresql/pgwatch/v6/internal/testutil"
15 "github.com/stretchr/testify/assert"
16 "github.com/stretchr/testify/require"
17 )
18
19
20
21 func setupIntegrationDB(t *testing.T) (*sources.DbConn, func()) {
22 t.Helper()
23 if testing.Short() {
24 t.Skip("skipping integration test in short mode")
25 }
26
27 pgContainer, tearDown, err := testutil.SetupPostgresContainer()
28 require.NoError(t, err, "failed to start postgres container")
29
30 connStr, err := pgContainer.ConnectionString(testutil.TestContext, "sslmode=disable")
31 require.NoError(t, err, "failed to get connection string")
32
33 pool, err := db.New(testutil.TestContext, connStr)
34 require.NoError(t, err, "failed to create connection pool")
35
36 md := sources.NewDbConn(sources.Source{
37 Name: "integration_test",
38 Kind: sources.SourcePostgres,
39 })
40 md.Conn = pool
41 err = md.FetchRuntimeInfo(testutil.TestContext, true)
42 require.NoError(t, err, "failed to fetch runtime info")
43
44 return md, func() {
45 pool.Close()
46 tearDown()
47 }
48 }
49
50
51
52
53 func TestIntegration_ExecuteBatch(t *testing.T) {
54 md, tearDown := setupIntegrationDB(t)
55 defer tearDown()
56
57 ctx := log.WithLogger(context.Background(), log.NewNoopLogger())
58
59 metricDefs.MetricDefs["integ_version"] = metrics.Metric{
60 SQLs: metrics.SQLs{0: "SELECT version() AS pg_version"},
61 }
62 metricDefs.MetricDefs["integ_uptime"] = metrics.Metric{
63 SQLs: metrics.SQLs{0: "SELECT extract(epoch from now() - pg_postmaster_start_time())::int8 AS uptime_seconds"},
64 }
65 defer func() {
66 delete(metricDefs.MetricDefs, "integ_version")
67 delete(metricDefs.MetricDefs, "integ_uptime")
68 }()
69
70 md.Metrics = metrics.MetricIntervals{
71 "integ_version": 30,
72 "integ_uptime": 60,
73 }
74
75 r := &reaper{
76 Options: &cmdopts.Options{
77 Metrics: metrics.CmdOpts{},
78 Sinks: sinks.CmdOpts{},
79 },
80 measurementCh: make(chan metrics.MeasurementEnvelope, 10),
81 measurementCache: NewInstanceMetricCache(),
82 }
83 sr := NewDbConnReaper(r, md)
84
85 _ = sr.executeBatch(ctx, []batchEntry{
86 {metricName: "integ_version", metric: metricDefs.MetricDefs["integ_version"], sql: "SELECT version() AS pg_version"},
87 {metricName: "integ_uptime", metric: metricDefs.MetricDefs["integ_uptime"], sql: "SELECT extract(epoch from now() - pg_postmaster_start_time())::int8 AS uptime_seconds"},
88 })
89
90 received := make(map[string]metrics.MeasurementEnvelope)
91 for range 10 {
92 select {
93 case msg := <-r.measurementCh:
94 received[msg.MetricName] = msg
95 default:
96 }
97 }
98
99 assert.Contains(t, received, "integ_version")
100 assert.Contains(t, received, "integ_uptime")
101
102 if msg, ok := received["integ_version"]; ok {
103 assert.Equal(t, "integration_test", msg.DBName)
104 assert.NotEmpty(t, msg.Data)
105 assert.Contains(t, msg.Data[0]["pg_version"], "PostgreSQL")
106 }
107
108 if msg, ok := received["integ_uptime"]; ok {
109 assert.Equal(t, "integration_test", msg.DBName)
110 assert.NotEmpty(t, msg.Data)
111 assert.NotNil(t, msg.Data[0]["uptime_seconds"])
112 }
113 }
114
115
116
117
118 func TestIntegration_SourceReaper_RunCollectsMetrics(t *testing.T) {
119 md, tearDown := setupIntegrationDB(t)
120 defer tearDown()
121
122 metricDefs.MetricDefs["integ_run_version"] = metrics.Metric{
123 SQLs: metrics.SQLs{0: "SELECT version() AS pg_version"},
124 }
125 metricDefs.MetricDefs["integ_run_size"] = metrics.Metric{
126 SQLs: metrics.SQLs{0: "SELECT pg_database_size(current_database()) AS db_size_bytes"},
127 }
128 defer func() {
129 delete(metricDefs.MetricDefs, "integ_run_version")
130 delete(metricDefs.MetricDefs, "integ_run_size")
131 }()
132
133 md.Metrics = metrics.MetricIntervals{
134 "integ_run_version": 5,
135 "integ_run_size": 5,
136 }
137
138 r := &reaper{
139 Options: &cmdopts.Options{
140 Metrics: metrics.CmdOpts{},
141 Sinks: sinks.CmdOpts{},
142 },
143 measurementCh: make(chan metrics.MeasurementEnvelope, 20),
144 measurementCache: NewInstanceMetricCache(),
145 }
146 sr := NewDbConnReaper(r, md)
147
148 ctx, cancel := context.WithCancel(log.WithLogger(context.Background(), log.NewNoopLogger()))
149
150 done := make(chan struct{})
151 go func() {
152 sr.Reap(ctx)
153 close(done)
154 }()
155
156 received := make(map[string]metrics.MeasurementEnvelope)
157 deadline := time.After(15 * time.Second)
158 for len(received) < 2 {
159 select {
160 case msg := <-r.measurementCh:
161 received[msg.MetricName] = msg
162 case <-deadline:
163 t.Fatal("timed out waiting for measurements")
164 }
165 }
166
167 cancel()
168 <-done
169
170 assert.Contains(t, received, "integ_run_version")
171 assert.Contains(t, received, "integ_run_size")
172
173 vMsg := received["integ_run_version"]
174 assert.Equal(t, "integration_test", vMsg.DBName)
175 assert.NotEmpty(t, vMsg.Data)
176 assert.Contains(t, vMsg.Data[0]["pg_version"], "PostgreSQL")
177
178 sMsg := received["integ_run_size"]
179 assert.Equal(t, "integration_test", sMsg.DBName)
180 assert.NotEmpty(t, sMsg.Data)
181 }
182
183 func TestIntegration_SourceReaper_RunExcludesMetricsByNodeStatus(t *testing.T) {
184 md, tearDown := setupIntegrationDB(t)
185 defer tearDown()
186
187 helperSetNodeStatus := func(status string) {
188 metricDefs.MetricDefs["test_metric"] = metrics.Metric{
189 SQLs: metrics.SQLs{0: "SELECT 1 AS value"},
190 NodeStatus: status,
191 }
192 metricDefs.MetricDefs["server_log_event_counts"] = metrics.Metric{
193 SQLs: metrics.SQLs{0: "SELECT 1 AS value"},
194 NodeStatus: status,
195 }
196 metricDefs.MetricDefs["psutil_cpu"] = metrics.Metric{
197 SQLs: metrics.SQLs{0: "SELECT 1 AS value"},
198 NodeStatus: status,
199 }
200 metricDefs.MetricDefs[specialMetricInstanceUp] = metrics.Metric{
201 SQLs: metrics.SQLs{0: "SELECT 1 AS value"},
202 NodeStatus: status,
203 }
204 }
205
206 r := &reaper{
207 Options: &cmdopts.Options{
208 Metrics: metrics.CmdOpts{},
209 Sinks: sinks.CmdOpts{},
210 },
211 measurementCh: make(chan metrics.MeasurementEnvelope, 10),
212 measurementCache: NewInstanceMetricCache(),
213 }
214
215
216
217 md.Metrics = metrics.MetricIntervals{
218 "test_metric": 5,
219 "server_log_event_counts": 5,
220 "psutil_cpu": 5,
221 specialMetricInstanceUp: 5,
222 }
223
224 t.Run("primary-only/standby-only metrics get excluded when node is standby/primary", func(t *testing.T) {
225 states := []string{"primary", "standby"}
226 for _, state := range states {
227 ctx, cancel := context.WithCancel(log.WithLogger(context.Background(), log.NewNoopLogger()))
228
229 md.Lock()
230 md.IsInRecovery = true
231 if state == "standby" {
232 md.IsInRecovery = false
233 }
234 md.Unlock()
235
236 helperSetNodeStatus(state)
237
238 sr := NewDbConnReaper(r, md)
239 go func() {
240 sr.Reap(ctx)
241 }()
242
243 select {
244 case msg := <-r.measurementCh:
245 t.Errorf("Expected no measurement for primary-only metrics on standby, but got: %s", msg.MetricName)
246 case <-time.After(2 * time.Second):
247 }
248
249 cancel()
250 }
251 })
252
253 t.Run("primary-only/standby-only metrics get executed when node is primary/standby", func(t *testing.T) {
254 states := []string{"primary", "standby", ""}
255 for _, state := range states {
256 ctx, cancel := context.WithCancel(log.WithLogger(context.Background(), log.NewNoopLogger()))
257
258 md.Lock()
259 md.IsInRecovery = false
260 if state == "standby" {
261 md.IsInRecovery = true
262 }
263 md.Unlock()
264
265 helperSetNodeStatus(state)
266
267 sr := NewDbConnReaper(r, md)
268 go func() {
269 sr.Reap(ctx)
270 }()
271
272 time.Sleep(2 * time.Second)
273 assert.GreaterOrEqual(t, len(r.measurementCh), 3)
274 cancel()
275
276 for range len(r.measurementCh) {
277
278 select {
279 case <-r.measurementCh:
280 default:
281 }
282 }
283 }
284 })
285 }
286