...

Source file src/github.com/cybertec-postgresql/pgwatch/v6/internal/reaper/metric_test.go

Documentation: github.com/cybertec-postgresql/pgwatch/v6/internal/reaper

     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  // mockDefinerWriter implements sinks.Writer and sinks.MetricsDefiner for testing.
    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  	// Regression test for https://github.com/cybertec-postgresql/pgwatch/issues/1091
   245  	// When a source config change (e.g. custom_tags) triggers a gatherer restart AND the preset
   246  	// interval is updated simultaneously, the new interval must be picked up after reload.
   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  		// Attach a mock connection so CloseResourcesForRemovedMonitoredDBs doesn't panic
   278  		// when the custom_tags change triggers a full source restart.
   279  		mockConn, err := pgxmock.NewPool()
   280  		require.NoError(t, err)
   281  		mockConn.ExpectClose()
   282  		r.monitoredSources[0].(*sources.DbConn).Conn = mockConn
   283  
   284  		// Simulate what happens between two Reap iterations:
   285  		// 1. source custom_tags change triggers restart
   286  		// 2. preset interval also changes
   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