1 package reaper
2
3 import (
4 "context"
5 "testing"
6
7 "github.com/cybertec-postgresql/pgwatch/v6/internal/cmdopts"
8 "github.com/cybertec-postgresql/pgwatch/v6/internal/log"
9 "github.com/cybertec-postgresql/pgwatch/v6/internal/metrics"
10 "github.com/cybertec-postgresql/pgwatch/v6/internal/sinks"
11 "github.com/cybertec-postgresql/pgwatch/v6/internal/sources"
12 "github.com/cybertec-postgresql/pgwatch/v6/internal/testutil"
13 "github.com/pashagolub/pgxmock/v5"
14 "github.com/stretchr/testify/assert"
15 "github.com/stretchr/testify/require"
16 )
17
18
19 type mockDefinerWriter struct {
20 defineErr error
21 defineCalled bool
22 receivedDefs *metrics.Metrics
23 }
24
25 func (m *mockDefinerWriter) SyncMetric(string, string, sinks.SyncOp) error { return nil }
26 func (m *mockDefinerWriter) Write(metrics.MeasurementEnvelope) error { return nil }
27 func (m *mockDefinerWriter) DefineMetrics(defs *metrics.Metrics) error {
28 m.defineCalled = true
29 m.receivedDefs = defs
30 return m.defineErr
31 }
32
33 var (
34 initialMetricDefs = metrics.MetricDefs{
35 "metric1": metrics.Metric{Description: "metric1"},
36 }
37 initialPresetDefs = metrics.PresetDefs{
38 "preset1": metrics.Preset{Description: "preset1", Metrics: metrics.MetricIntervals{"metric1": 1.0}},
39 }
40
41 newMetricDefs = metrics.MetricDefs{
42 "metric2": metrics.Metric{Description: "metric2"},
43 }
44 newPresetDefs = metrics.PresetDefs{
45 "preset2": metrics.Preset{Description: "preset2", Metrics: metrics.MetricIntervals{"metric2": 2.0}},
46 }
47 )
48
49 func TestReaper_FetchStatsDirectlyFromOS(t *testing.T) {
50 a := assert.New(t)
51 r := &reaper{Options: &cmdopts.Options{}}
52 t.Run("metrics directly fetchable when on same host", func(*testing.T) {
53 conn, _ := pgxmock.NewPool(pgxmock.QueryMatcherOption(pgxmock.QueryMatcherEqual))
54 expq := conn.ExpectQuery("SELECT COALESCE(inet_client_addr(), inet_server_addr()) IS NULL")
55 expq.Times(uint(len(directlyFetchableOSMetrics)))
56 md := &sources.DbConn{Conn: conn}
57 for _, m := range directlyFetchableOSMetrics {
58 expq.WillReturnRows(pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(true))
59 a.True(IsDirectlyFetchableMetric(md, m), "Expected %s to be directly fetchable", m)
60 a.NotPanics(func() {
61 _, _ = r.FetchStatsDirectlyFromOS(context.Background(), md, m)
62 })
63 }
64 })
65
66 t.Run("cpu_load not directly fetchable when not on same host", func(*testing.T) {
67 remoteConn, _ := pgxmock.NewPool(pgxmock.QueryMatcherOption(pgxmock.QueryMatcherEqual))
68 remoteConn.ExpectQuery("SELECT COALESCE(inet_client_addr(), inet_server_addr()) IS NULL").
69 WillReturnRows(pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(false))
70 remoteMd := &sources.DbConn{Conn: remoteConn}
71 a.False(IsDirectlyFetchableMetric(remoteMd, metricCPULoad),
72 "cpu_load should not be directly fetchable when pgwatch is not on the same host as PostgreSQL")
73 })
74 }
75
76 func TestConcurrentMetricDefs_Assign(t *testing.T) {
77 concurrentDefs := NewConcurrentMetricDefs()
78 concurrentDefs.Assign(&metrics.Metrics{
79 MetricDefs: initialMetricDefs,
80 PresetDefs: initialPresetDefs,
81 })
82
83 concurrentDefs.Assign(&metrics.Metrics{
84 MetricDefs: newMetricDefs,
85 PresetDefs: newPresetDefs,
86 })
87
88 assert.Equal(t, newMetricDefs, concurrentDefs.MetricDefs, "MetricDefs should be updated")
89 assert.Equal(t, newPresetDefs, concurrentDefs.PresetDefs, "PresetDefs should be updated")
90 }
91
92 func TestConcurrentMetricDefs_RandomAccess(t *testing.T) {
93 a := assert.New(t)
94
95 concurrentDefs := NewConcurrentMetricDefs()
96 concurrentDefs.Assign(&metrics.Metrics{
97 MetricDefs: initialMetricDefs,
98 PresetDefs: initialPresetDefs,
99 })
100
101 go a.NotPanics(func() {
102 for range 1000 {
103 _, ok1 := concurrentDefs.GetMetricDef("metric1")
104 _, ok2 := concurrentDefs.GetMetricDef("metric2")
105 a.True(ok1 || ok2, "Expected metric1 or metric3 to exist at any time")
106 _, ok1 = concurrentDefs.GetPresetDef("preset1")
107 _, ok2 = concurrentDefs.GetPresetDef("preset2")
108 a.True(ok1 || ok2, "Expected preset1 or preset2 to exist at any time")
109 m1 := concurrentDefs.GetPresetMetrics("preset1")
110 m2 := concurrentDefs.GetPresetMetrics("preset2")
111 a.True(m1 != nil || m2 != nil, "Expected preset1 or preset2 metrics to be non-empty")
112 }
113 })
114
115 go a.NotPanics(func() {
116 for range 1000 {
117 concurrentDefs.Assign(&metrics.Metrics{
118 MetricDefs: newMetricDefs,
119 PresetDefs: newPresetDefs,
120 })
121 }
122 })
123 }
124
125 func TestReaper_LoadMetrics(t *testing.T) {
126 ctx := log.WithLogger(t.Context(), log.NewNoopLogger())
127
128 t.Run("returns error from GetMetrics", func(t *testing.T) {
129 r := newReaper(ctx, &cmdopts.Options{
130 MetricsReaderWriter: &testutil.MockMetricsReaderWriter{
131 GetMetricsFunc: func() (*metrics.Metrics, error) {
132 return nil, assert.AnError
133 },
134 },
135 })
136 assert.ErrorIs(t, r.LoadMetrics(), assert.AnError)
137 })
138
139 t.Run("updates metricDefs on success", func(t *testing.T) {
140 defs := &metrics.Metrics{
141 MetricDefs: metrics.MetricDefs{"m1": {Description: "M1"}},
142 PresetDefs: metrics.PresetDefs{"p1": {Description: "P1", Metrics: metrics.MetricIntervals{"m1": 1.0}}},
143 }
144 r := newReaper(ctx, &cmdopts.Options{
145 MetricsReaderWriter: &testutil.MockMetricsReaderWriter{
146 GetMetricsFunc: func() (*metrics.Metrics, error) { return defs, nil },
147 },
148 })
149 assert.NoError(t, r.LoadMetrics())
150
151 m, ok := metricDefs.GetMetricDef("m1")
152 assert.True(t, ok)
153 assert.Equal(t, "M1", m.Description)
154
155 p, ok := metricDefs.GetPresetDef("p1")
156 assert.True(t, ok)
157 assert.Equal(t, "P1", p.Description)
158 })
159
160 t.Run("calls DefineMetrics on MetricsDefiner sink", func(t *testing.T) {
161 defs := &metrics.Metrics{
162 MetricDefs: metrics.MetricDefs{"m2": {Description: "M2"}},
163 PresetDefs: metrics.PresetDefs{},
164 }
165 mock := &mockDefinerWriter{}
166 r := newReaper(ctx, &cmdopts.Options{
167 MetricsReaderWriter: &testutil.MockMetricsReaderWriter{
168 GetMetricsFunc: func() (*metrics.Metrics, error) { return defs, nil },
169 },
170 SinksWriter: mock,
171 })
172 assert.NoError(t, r.LoadMetrics())
173 assert.True(t, mock.defineCalled, "DefineMetrics should be called on MetricsDefiner sink")
174 assert.Equal(t, defs, mock.receivedDefs)
175 })
176
177 t.Run("DefineMetrics error is logged not returned", func(t *testing.T) {
178 defs := &metrics.Metrics{
179 MetricDefs: metrics.MetricDefs{"m3": {Description: "M3"}},
180 PresetDefs: metrics.PresetDefs{},
181 }
182 mock := &mockDefinerWriter{defineErr: assert.AnError}
183 r := newReaper(ctx, &cmdopts.Options{
184 MetricsReaderWriter: &testutil.MockMetricsReaderWriter{
185 GetMetricsFunc: func() (*metrics.Metrics, error) { return defs, nil },
186 },
187 SinksWriter: mock,
188 })
189 assert.NoError(t, r.LoadMetrics(), "DefineMetrics error should not propagate")
190 assert.True(t, mock.defineCalled)
191 })
192
193 t.Run("resolves preset metrics for monitored sources", func(t *testing.T) {
194 defs := &metrics.Metrics{
195 MetricDefs: metrics.MetricDefs{
196 "m1": {Description: "M1"},
197 "m2": {Description: "M2"},
198 },
199 PresetDefs: metrics.PresetDefs{
200 "preset1": {Metrics: metrics.MetricIntervals{"m1": 10.0}},
201 "standby1": {Metrics: metrics.MetricIntervals{"m2": 20.0}},
202 },
203 }
204 r := newReaper(ctx, &cmdopts.Options{
205 MetricsReaderWriter: &testutil.MockMetricsReaderWriter{
206 GetMetricsFunc: func() (*metrics.Metrics, error) { return defs, nil },
207 },
208 })
209 r.monitoredSources = sources.SourceConns{
210 sources.NewDbConn(sources.Source{
211 Name: "src1",
212 PresetMetrics: "preset1",
213 PresetMetricsStandby: "standby1",
214 }),
215 }
216 assert.NoError(t, r.LoadMetrics())
217
218 sc := r.monitoredSources[0].(*sources.DbConn)
219 assert.Equal(t, metrics.MetricIntervals{"m1": 10.0}, sc.Metrics)
220 assert.Equal(t, metrics.MetricIntervals{"m2": 20.0}, sc.MetricsStandby)
221 })
222
223 t.Run("skips preset resolution for sources without presets", func(t *testing.T) {
224 customMetrics := metrics.MetricIntervals{"cpu": 5.0}
225 defs := &metrics.Metrics{
226 MetricDefs: metrics.MetricDefs{"cpu": {Description: "CPU"}},
227 PresetDefs: metrics.PresetDefs{},
228 }
229 r := newReaper(ctx, &cmdopts.Options{
230 MetricsReaderWriter: &testutil.MockMetricsReaderWriter{
231 GetMetricsFunc: func() (*metrics.Metrics, error) { return defs, nil },
232 },
233 })
234 r.monitoredSources = sources.SourceConns{
235 sources.NewDbConn(sources.Source{
236 Name: "src2",
237 Metrics: customMetrics,
238 }),
239 }
240 assert.NoError(t, r.LoadMetrics())
241 assert.Equal(t, customMetrics, r.monitoredSources[0].(*sources.DbConn).Metrics, "custom metrics should be unchanged")
242 })
243
244
245
246
247 t.Run("preset interval update is applied after source config change", func(t *testing.T) {
248 ctx := log.WithLogger(t.Context(), log.NewNoopLogger())
249 initialDefs := &metrics.Metrics{
250 MetricDefs: metrics.MetricDefs{"test_metric": {}},
251 PresetDefs: metrics.PresetDefs{
252 "test_preset": {Metrics: metrics.MetricIntervals{"test_metric": 1}},
253 },
254 }
255 src := sources.Source{
256 Name: "test_source",
257 IsEnabled: true,
258 Kind: sources.SourcePostgres,
259 ConnStr: "postgres://localhost:5432/testdb",
260 CustomTags: map[string]string{"version": "1.0"},
261 PresetMetrics: "test_preset",
262 }
263 getMetricsFn := func() (*metrics.Metrics, error) { return initialDefs, nil }
264 r := newReaper(ctx, &cmdopts.Options{
265 SourcesReaderWriter: &testutil.MockSourcesReaderWriter{
266 GetSourcesFunc: func() (sources.Sources, error) { return sources.Sources{src}, nil },
267 },
268 MetricsReaderWriter: &testutil.MockMetricsReaderWriter{
269 GetMetricsFunc: func() (*metrics.Metrics, error) { return getMetricsFn() },
270 },
271 SinksWriter: &sinks.MultiWriter{},
272 })
273 require.NoError(t, r.LoadSources(ctx))
274 require.NoError(t, r.LoadMetrics())
275 assert.Equal(t, metrics.MetricIntervals{"test_metric": 1}, r.monitoredSources[0].(*sources.DbConn).Metrics)
276
277
278
279 mockConn, err := pgxmock.NewPool()
280 require.NoError(t, err)
281 mockConn.ExpectClose()
282 r.monitoredSources[0].(*sources.DbConn).Conn = mockConn
283
284
285
286
287 src.CustomTags = map[string]string{"version": "2.0"}
288 updatedDefs := &metrics.Metrics{
289 MetricDefs: metrics.MetricDefs{"test_metric": {}},
290 PresetDefs: metrics.PresetDefs{
291 "test_preset": {Metrics: metrics.MetricIntervals{"test_metric": 2}},
292 },
293 }
294 getMetricsFn = func() (*metrics.Metrics, error) { return updatedDefs, nil }
295
296 require.NoError(t, r.LoadSources(ctx))
297 require.NoError(t, r.LoadMetrics())
298 assert.Equal(t, metrics.MetricIntervals{"test_metric": 2}, r.monitoredSources[0].(*sources.DbConn).Metrics,
299 "preset interval should be updated after source config change triggered a restart")
300 })
301 }
302
303 func TestChangeDetectionResults(t *testing.T) {
304 a := assert.New(t)
305
306 cdr := &ChangeDetectionResults{
307 Target: "test",
308 Created: 2,
309 Altered: 3,
310 Dropped: 1,
311 }
312
313 a.Equal(6, cdr.Total(), "Total should be sum of Created, Altered, and Dropped")
314 a.Equal("test: 2/3/1", cdr.String(), "String representation should match expected format")
315 }
316