mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 04:38:48 +00:00
chore(metrics): move store regression proof into accepted proof file
The canonical completion guard requires the metrics-store hot-path proof in pkg/metrics/store_additional_test.go. Move the rollup-gap, upsert-spread, busy-retry and auto-vacuum regressions there and point the performance-and-scalability contract at that file. The runtime change and contract section are in the preceding fix(metrics) commit. Change-source: pulse-maintainer
This commit is contained in:
parent
e1d97b361e
commit
087bbf8c20
3 changed files with 180 additions and 181 deletions
|
|
@ -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`.
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue