diff --git a/docs/release-control/v6/internal/subsystems/performance-and-scalability.md b/docs/release-control/v6/internal/subsystems/performance-and-scalability.md index 29fb7354e..a50d82d4e 100644 --- a/docs/release-control/v6/internal/subsystems/performance-and-scalability.md +++ b/docs/release-control/v6/internal/subsystems/performance-and-scalability.md @@ -3190,7 +3190,7 @@ had its rows stranded below the upper tier's checkpoint and eventually purged by retention. `rollupTier` now stops at the gap, leaving the checkpoint at the gap start; the next invocation skips the gap once it is confirmed empty. The `nextSourceRollupBucket` prefix skip still runs only before the first source -window of an invocation, preserving bounded catch-up. `pkg/metrics/store_test.go` +window of an invocation, preserving bounded catch-up. `pkg/metrics/store_additional_test.go` pins the interspersed-gap backfill in `TestStoreRollupTierBackfillsInterspersedGap` and keeps `TestStoreRollupTierEmptyWindowPreservesCheckpoint`. @@ -3204,7 +3204,7 @@ reports `database is locked (5) (SQLITE_BUSY)`, which an equality check against retrying. `isRetryableWriteError` matches the BUSY/LOCKED code family with a substring fallback for closed-pool errors. A direct aggregate-tier write also preserves an existing rollup row's `min_value`/`max_value` via `COALESCE` rather -than nulling the stored spread. `pkg/metrics/store_test.go` pins both in +than nulling the stored spread. `pkg/metrics/store_additional_test.go` pins both in `TestIsRetryableWriteErrorMatchesDriverBusyCode` and `TestStoreWriteBatchPreservesRollupSpread`. @@ -3216,5 +3216,5 @@ per-connection setting until that connection's `VACUUM` rewrites the file header. Running them on separate pool connections could leave the file at NONE and repeat the full-file VACUUM on every restart. The migration now pins a connection, verifies the persisted mode, and leaves the conversion for a later -restart if verification fails. `pkg/metrics/store_test.go` pins the persisted +restart if verification fails. `pkg/metrics/store_additional_test.go` pins the persisted mode in `TestStoreAutoVacuumPersistsAcrossRestart`. diff --git a/pkg/metrics/store_additional_test.go b/pkg/metrics/store_additional_test.go index 83eca705b..44b981d8e 100644 --- a/pkg/metrics/store_additional_test.go +++ b/pkg/metrics/store_additional_test.go @@ -1542,3 +1542,180 @@ func TestConcurrentWritesDuringShutdownDoNotPanic(t *testing.T) { close(stop) wg.Wait() } + +// TestStoreRollupTierBackfillsInterspersedGap verifies that a rollup checkpoint +// does not leap over a gap between already-processed source windows. A late +// backfill into that gap must still be aggregated on a later run instead of +// being stranded below the checkpoint and eventually purged. +func TestStoreRollupTierBackfillsInterspersedGap(t *testing.T) { + dir := t.TempDir() + cfg := DefaultConfig(dir) + cfg.DBPath = filepath.Join(dir, "metrics-rollup-gap.db") + cfg.FlushInterval = time.Hour + + store, err := NewStore(cfg) + if err != nil { + t.Fatalf("NewStore returned error: %v", err) + } + defer store.Close() + + base := time.Now().UTC().Add(-4 * time.Hour).Truncate(time.Hour) + firstHour := base.Unix() + thirdHour := base.Add(2 * time.Hour).Unix() + gapHour := base.Add(time.Hour).Unix() + + if _, err := store.db.Exec( + `INSERT INTO metrics (resource_type, resource_id, metric_type, value, min_value, max_value, timestamp, tier) VALUES + ('vm','vm-gap','cpu',1.0,1.0,1.0,?, 'minute'), + ('vm','vm-gap','cpu',3.0,3.0,3.0,?, 'minute')`, firstHour, thirdHour); err != nil { + t.Fatalf("insert source rows: %v", err) + } + + store.rollupTier(TierMinute, TierHourly, time.Hour, time.Hour) + + // Backfill the interspersed gap after the first rollup ran. + if _, err := store.db.Exec( + `INSERT INTO metrics (resource_type, resource_id, metric_type, value, min_value, max_value, timestamp, tier) VALUES + ('vm','vm-gap','cpu',2.0,2.0,2.0,?, 'minute')`, gapHour); err != nil { + t.Fatalf("backfill gap row: %v", err) + } + + store.rollupTier(TierMinute, TierHourly, time.Hour, time.Hour) + + var count int + if err := store.db.QueryRow( + `SELECT COUNT(*) FROM metrics WHERE tier='hourly' AND resource_id='vm-gap' AND timestamp=?`, gapHour).Scan(&count); err != nil { + t.Fatalf("query hourly gap: %v", err) + } + if count != 1 { + t.Fatalf("expected backfilled gap hour %d to roll up to hourly, got %d rows", gapHour, count) + } +} + +// TestStoreWriteBatchPreservesRollupSpread verifies that a direct write to an +// aggregate tier does not erase an existing rollup row's min/max spread. The +// write statement omits min/max, so excluded.* are NULL and a naive upsert +// would null the stored spread. +func TestStoreWriteBatchPreservesRollupSpread(t *testing.T) { + dir := t.TempDir() + cfg := DefaultConfig(dir) + cfg.DBPath = filepath.Join(dir, "metrics-spread.db") + cfg.FlushInterval = time.Hour + + store, err := NewStore(cfg) + if err != nil { + t.Fatalf("NewStore returned error: %v", err) + } + defer store.Close() + + ts := time.Now().UTC().Add(-2 * time.Hour).Truncate(time.Minute) + if _, err := store.db.Exec( + `INSERT INTO metrics (resource_type, resource_id, metric_type, value, min_value, max_value, timestamp, tier) VALUES + ('vm','vm-spread','cpu',5.0,1.0,9.0,?, 'minute')`, ts.Unix()); err != nil { + t.Fatalf("insert rollup row: %v", err) + } + + store.writeBatch([]bufferedMetric{ + {resourceType: "vm", resourceID: "vm-spread", metricType: "cpu", value: 7.0, timestamp: ts, tier: TierMinute}, + }) + + var minValue, maxValue sql.NullFloat64 + if err := store.db.QueryRow( + `SELECT min_value, max_value FROM metrics WHERE tier='minute' AND resource_id='vm-spread' AND timestamp=?`, ts.Unix()).Scan(&minValue, &maxValue); err != nil { + t.Fatalf("query spread: %v", err) + } + if !minValue.Valid || !maxValue.Valid { + t.Fatalf("direct write erased rollup spread: min=%+v max=%+v", minValue, maxValue) + } + if minValue.Float64 != 1.0 || maxValue.Float64 != 9.0 { + t.Fatalf("unexpected spread after direct write: min=%v max=%v", minValue.Float64, maxValue.Float64) + } +} + +func TestIsRetryableWriteError(t *testing.T) { + cases := []struct { + name string + err error + want bool + }{ + {"nil", nil, false}, + {"driver busy", fmt.Errorf("database is locked (5) (SQLITE_BUSY)"), true}, + {"legacy busy", fmt.Errorf("database is locked"), true}, + {"closed pool", fmt.Errorf("sql: database is closed"), true}, + {"unrelated", fmt.Errorf("no such table: metrics"), false}, + } + for _, tc := range cases { + if got := isRetryableWriteError(tc.err); got != tc.want { + t.Fatalf("%s: isRetryableWriteError(%v) = %v, want %v", tc.name, tc.err, got, tc.want) + } + } +} + +// TestIsRetryableWriteErrorMatchesDriverBusyCode proves the classifier matches +// a real SQLITE_BUSY error returned by the modernc driver, not just its message. +func TestIsRetryableWriteErrorMatchesDriverBusyCode(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "busy.db") + + holder, err := sql.Open("sqlite", "file:"+path+"?_pragma=busy_timeout(0)") + if err != nil { + t.Fatalf("open holder: %v", err) + } + defer holder.Close() + holder.SetMaxOpenConns(1) + if _, err := holder.Exec("CREATE TABLE t (x INTEGER)"); err != nil { + t.Fatalf("create table: %v", err) + } + if _, err := holder.Exec("BEGIN IMMEDIATE"); err != nil { + t.Fatalf("begin immediate: %v", err) + } + defer holder.Exec("ROLLBACK") + + waiter, err := sql.Open("sqlite", "file:"+path+"?_pragma=busy_timeout(0)") + if err != nil { + t.Fatalf("open waiter: %v", err) + } + defer waiter.Close() + + _, err = waiter.Exec("BEGIN IMMEDIATE") + if err == nil { + waiter.Exec("ROLLBACK") + t.Skip("driver did not report SQLITE_BUSY on a contended write lock") + } + if !isRetryableWriteError(err) { + t.Fatalf("real driver busy error not classified retryable: %v", err) + } +} + +// TestStoreAutoVacuumPersistsAcrossRestart verifies that the one-time +// auto-vacuum conversion is committed to the database file and remains in +// effect after reopening the store, rather than being re-run every restart. +func TestStoreAutoVacuumPersistsAcrossRestart(t *testing.T) { + dir := t.TempDir() + cfg := DefaultConfig(dir) + cfg.DBPath = filepath.Join(dir, "metrics-autovacuum.db") + cfg.FlushInterval = time.Hour + + for i := 0; i < 2; i++ { + store, err := NewStore(cfg) + if err != nil { + t.Fatalf("open %d: %v", i, err) + } + if err := store.WaitForMaintenance(10 * time.Second); err != nil { + store.Close() + t.Fatalf("maintenance %d: %v", i, err) + } + var mode int + if err := store.db.QueryRow("PRAGMA auto_vacuum").Scan(&mode); err != nil { + store.Close() + t.Fatalf("pragma %d: %v", i, err) + } + if mode != 2 { + store.Close() + t.Fatalf("open %d: auto_vacuum = %d, want 2 (INCREMENTAL)", i, mode) + } + if err := store.Close(); err != nil { + t.Fatalf("close %d: %v", i, err) + } + } +} diff --git a/pkg/metrics/store_test.go b/pkg/metrics/store_test.go index 7bb6a8e9d..552f16e0f 100644 --- a/pkg/metrics/store_test.go +++ b/pkg/metrics/store_test.go @@ -1,7 +1,6 @@ package metrics import ( - "database/sql" "fmt" "path/filepath" "testing" @@ -795,180 +794,3 @@ func TestQueryAllBatch(t *testing.T) { } }) } - -// TestStoreRollupTierBackfillsInterspersedGap verifies that a rollup checkpoint -// does not leap over a gap between already-processed source windows. A late -// backfill into that gap must still be aggregated on a later run instead of -// being stranded below the checkpoint and eventually purged. -func TestStoreRollupTierBackfillsInterspersedGap(t *testing.T) { - dir := t.TempDir() - cfg := DefaultConfig(dir) - cfg.DBPath = filepath.Join(dir, "metrics-rollup-gap.db") - cfg.FlushInterval = time.Hour - - store, err := NewStore(cfg) - if err != nil { - t.Fatalf("NewStore returned error: %v", err) - } - defer store.Close() - - base := time.Now().UTC().Add(-4 * time.Hour).Truncate(time.Hour) - firstHour := base.Unix() - thirdHour := base.Add(2 * time.Hour).Unix() - gapHour := base.Add(time.Hour).Unix() - - if _, err := store.db.Exec( - `INSERT INTO metrics (resource_type, resource_id, metric_type, value, min_value, max_value, timestamp, tier) VALUES - ('vm','vm-gap','cpu',1.0,1.0,1.0,?, 'minute'), - ('vm','vm-gap','cpu',3.0,3.0,3.0,?, 'minute')`, firstHour, thirdHour); err != nil { - t.Fatalf("insert source rows: %v", err) - } - - store.rollupTier(TierMinute, TierHourly, time.Hour, time.Hour) - - // Backfill the interspersed gap after the first rollup ran. - if _, err := store.db.Exec( - `INSERT INTO metrics (resource_type, resource_id, metric_type, value, min_value, max_value, timestamp, tier) VALUES - ('vm','vm-gap','cpu',2.0,2.0,2.0,?, 'minute')`, gapHour); err != nil { - t.Fatalf("backfill gap row: %v", err) - } - - store.rollupTier(TierMinute, TierHourly, time.Hour, time.Hour) - - var count int - if err := store.db.QueryRow( - `SELECT COUNT(*) FROM metrics WHERE tier='hourly' AND resource_id='vm-gap' AND timestamp=?`, gapHour).Scan(&count); err != nil { - t.Fatalf("query hourly gap: %v", err) - } - if count != 1 { - t.Fatalf("expected backfilled gap hour %d to roll up to hourly, got %d rows", gapHour, count) - } -} - -// TestStoreWriteBatchPreservesRollupSpread verifies that a direct write to an -// aggregate tier does not erase an existing rollup row's min/max spread. The -// write statement omits min/max, so excluded.* are NULL and a naive upsert -// would null the stored spread. -func TestStoreWriteBatchPreservesRollupSpread(t *testing.T) { - dir := t.TempDir() - cfg := DefaultConfig(dir) - cfg.DBPath = filepath.Join(dir, "metrics-spread.db") - cfg.FlushInterval = time.Hour - - store, err := NewStore(cfg) - if err != nil { - t.Fatalf("NewStore returned error: %v", err) - } - defer store.Close() - - ts := time.Now().UTC().Add(-2 * time.Hour).Truncate(time.Minute) - if _, err := store.db.Exec( - `INSERT INTO metrics (resource_type, resource_id, metric_type, value, min_value, max_value, timestamp, tier) VALUES - ('vm','vm-spread','cpu',5.0,1.0,9.0,?, 'minute')`, ts.Unix()); err != nil { - t.Fatalf("insert rollup row: %v", err) - } - - store.writeBatch([]bufferedMetric{ - {resourceType: "vm", resourceID: "vm-spread", metricType: "cpu", value: 7.0, timestamp: ts, tier: TierMinute}, - }) - - var minValue, maxValue sql.NullFloat64 - if err := store.db.QueryRow( - `SELECT min_value, max_value FROM metrics WHERE tier='minute' AND resource_id='vm-spread' AND timestamp=?`, ts.Unix()).Scan(&minValue, &maxValue); err != nil { - t.Fatalf("query spread: %v", err) - } - if !minValue.Valid || !maxValue.Valid { - t.Fatalf("direct write erased rollup spread: min=%+v max=%+v", minValue, maxValue) - } - if minValue.Float64 != 1.0 || maxValue.Float64 != 9.0 { - t.Fatalf("unexpected spread after direct write: min=%v max=%v", minValue.Float64, maxValue.Float64) - } -} - -func TestIsRetryableWriteError(t *testing.T) { - cases := []struct { - name string - err error - want bool - }{ - {"nil", nil, false}, - {"driver busy", fmt.Errorf("database is locked (5) (SQLITE_BUSY)"), true}, - {"legacy busy", fmt.Errorf("database is locked"), true}, - {"closed pool", fmt.Errorf("sql: database is closed"), true}, - {"unrelated", fmt.Errorf("no such table: metrics"), false}, - } - for _, tc := range cases { - if got := isRetryableWriteError(tc.err); got != tc.want { - t.Fatalf("%s: isRetryableWriteError(%v) = %v, want %v", tc.name, tc.err, got, tc.want) - } - } -} - -// TestIsRetryableWriteErrorMatchesDriverBusyCode proves the classifier matches -// a real SQLITE_BUSY error returned by the modernc driver, not just its message. -func TestIsRetryableWriteErrorMatchesDriverBusyCode(t *testing.T) { - dir := t.TempDir() - path := filepath.Join(dir, "busy.db") - - holder, err := sql.Open("sqlite", "file:"+path+"?_pragma=busy_timeout(0)") - if err != nil { - t.Fatalf("open holder: %v", err) - } - defer holder.Close() - holder.SetMaxOpenConns(1) - if _, err := holder.Exec("CREATE TABLE t (x INTEGER)"); err != nil { - t.Fatalf("create table: %v", err) - } - if _, err := holder.Exec("BEGIN IMMEDIATE"); err != nil { - t.Fatalf("begin immediate: %v", err) - } - defer holder.Exec("ROLLBACK") - - waiter, err := sql.Open("sqlite", "file:"+path+"?_pragma=busy_timeout(0)") - if err != nil { - t.Fatalf("open waiter: %v", err) - } - defer waiter.Close() - - _, err = waiter.Exec("BEGIN IMMEDIATE") - if err == nil { - waiter.Exec("ROLLBACK") - t.Skip("driver did not report SQLITE_BUSY on a contended write lock") - } - if !isRetryableWriteError(err) { - t.Fatalf("real driver busy error not classified retryable: %v", err) - } -} - -// TestStoreAutoVacuumPersistsAcrossRestart verifies that the one-time -// auto-vacuum conversion is committed to the database file and remains in -// effect after reopening the store, rather than being re-run every restart. -func TestStoreAutoVacuumPersistsAcrossRestart(t *testing.T) { - dir := t.TempDir() - cfg := DefaultConfig(dir) - cfg.DBPath = filepath.Join(dir, "metrics-autovacuum.db") - cfg.FlushInterval = time.Hour - - for i := 0; i < 2; i++ { - store, err := NewStore(cfg) - if err != nil { - t.Fatalf("open %d: %v", i, err) - } - if err := store.WaitForMaintenance(10 * time.Second); err != nil { - store.Close() - t.Fatalf("maintenance %d: %v", i, err) - } - var mode int - if err := store.db.QueryRow("PRAGMA auto_vacuum").Scan(&mode); err != nil { - store.Close() - t.Fatalf("pragma %d: %v", i, err) - } - if mode != 2 { - store.Close() - t.Fatalf("open %d: auto_vacuum = %d, want 2 (INCREMENTAL)", i, mode) - } - if err := store.Close(); err != nil { - t.Fatalf("close %d: %v", i, err) - } - } -}