...

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

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

     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  // setupIntegrationDB starts a real Postgres container and returns a SourceConn
    20  // with a live pgxpool connection. The caller must call tearDown when done.
    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  // TestIntegration_ExecuteBatch verifies the full executeBatch path against a real
    51  // Postgres instance: builds a pgx.Batch from metric definitions, sends it, and
    52  // receives MeasurementEnvelopes on the measurement channel.
    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  // TestIntegration_SourceReaper_RunCollectsMetrics verifies the full Run loop:
   116  // creates a SourceReaper with two SQL metrics, lets it run one tick against
   117  // a real Postgres container, and verifies that both metric envelopes arrive.
   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  	// using psutil_*, server_log_event_counts, instance_up
   216  	// to ensure specially-handled metrics have the same behaviour
   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", ""} // "" => should fetch all as well
   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  				// empty channel to ensure correctness in subsequent runs
   278  				select {
   279  				case <-r.measurementCh:
   280  				default:
   281  				}
   282  			}
   283  		}
   284  	})
   285  }
   286