From 2c8cb8435df51cef3134b2b085acaa21afa797c5 Mon Sep 17 00:00:00 2001 From: rcourtman <8825017+rcourtman@users.noreply.github.com> Date: Thu, 24 Sep 2026 08:45:38 +0100 Subject: [PATCH] Make the time-major metrics index the unique identity index Metric identity lived in idx_metrics_lookup, ordered (type, id, metric, tier, timestamp). Once a series retains history, that order gives every series its own insertion point, so each poll's commit rewrites one index leaf page per series through the WAL and again at checkpoint. On a three-node homelab running v6.4.5-rc.2 that came to 330-607 KB/s of disk writes, about 10.6 KB per ~60-byte sample and 27-50 GB a day (#1966). idx_metrics_query_all already holds the same five columns time-major, so one poll's samples for a resource share pages. It now carries uniqueness and serves every read, and the metric-major lookup index is dropped. With an hour of retained history for 480 series, 30 polls write 6,022 WAL frames instead of 25,310 (4.2x fewer); the empty-table issue-1124 estate drops from 36,308 frames to 20,223, and its ceiling tightens from 40,000 to 23,000. The tier-reconciliation overlap probe now names the identity index too; left to the planner it seeks on (tier, timestamp) alone once the metric-major tree is gone, which TestRetainedQueryPlansUseIndexes caught. Metric-filtered reads now scan a resource's other metrics in the window: a single-metric 500-node query went from 45 to 59 us and the 24-hour 50-disk smart_temp batch from about 0.25 to 0.49 s. All-metric dashboard reads are unchanged. Continuous write volume is the user-visible harm, so the contract records that trade explicitly. The migration reuses the crash-safe swap: databases whose identity is already enforced by a unique index defer the rebuild to startup maintenance and swap in one transaction, and a database with no unique index migrates synchronously. Keeping the name idx_metrics_query_all means a downgrade rebuilds only the lookup index and leaves no orphan. Refs #1966 --- .../subsystems/performance-and-scalability.md | 53 +++-- pkg/metrics/identity_index_test.go | 188 ++++++++++++++++++ .../issue1124_write_amplification_test.go | 60 +++--- pkg/metrics/store.go | 132 ++++++++---- pkg/metrics/store_query_plan_test.go | 13 +- 5 files changed, 359 insertions(+), 87 deletions(-) create mode 100644 pkg/metrics/identity_index_test.go 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 7a8e38fe0..6d27232b2 100644 --- a/docs/release-control/v6/internal/subsystems/performance-and-scalability.md +++ b/docs/release-control/v6/internal/subsystems/performance-and-scalability.md @@ -1963,26 +1963,41 @@ WAL checkpoints more aggressive again, reopens duplicate-write failures, removes the bounded rollup-window checkpointing, or removes the metrics DB path/cadence controls must re-prove the metrics-store hot path with the owned store tests rather than assuming the earlier vacuum fixes are sufficient. -The metrics identity and single-series lookup must also share one unique +Metric identity and every range read must also share one unique, time-major +`idx_metrics_query_all(resource_type, resource_id, tier, timestamp, +metric_type)` B-tree. Reintroducing a second full-width identity index is a +write-amplification regression even when retained `metrics.db` size looks small: +every sample would dirty both trees in the WAL and every checkpoint would copy +both sets of pages back into the main file. The column order is part of the +invariant. A metric-major tree such as the retired `idx_metrics_lookup(resource_type, resource_id, metric_type, tier, timestamp)` -B-tree. Reintroducing a second full-width identity index is a write-amplification -regression even when retained `metrics.db` size looks small: every sample would -dirty both trees in the WAL and every checkpoint would copy both sets of pages -back into the main file. Legacy databases may keep their existing unique index -authoritative while the replacement index is built, but that O(rows) migration -must run on deferred startup maintenance and swap indexes in one SQLite -transaction so constructor latency, concurrent reads, uniqueness, and crash -recovery remain intact. The deterministic issue-1124 profile must continue to -attribute writes by logical payload, table/index `dbstat` pages, page-cache -writes/spills, WAL frames, checkpoint frames, process write bytes where the -host exposes them, and traced sync calls rather than treating final file size -as a proxy. Its checked invariant persists 157,452 samples across 2,197 mixed -provider resources and caps the no-auto-checkpoint workload at 40,000 WAL -frames; the former four-index schema produces 50,516 frames while the -consolidated schema produces 35,030. Checkpoint-threshold changes must also -prove physical writes, WAL allocation, write latency, restart, and the -four-connection concurrent-read case; a no-reader byte reduction alone is not -sufficient to widen the production WAL bound. +gives every series its own insertion point once it retains history, so each +commit rewrites one leaf page per series; the time-major order lets one poll's +samples for a resource share pages. With one hour of retained history for 480 +series, 30 polls write 25,310 WAL frames under the metric-major layout and 6,022 +under the time-major one (`TestMetricsIdentityIndexSteadyStateWrites` requires +at most 35%), consistent with a production install measuring about 10.6 KB of +WAL and checkpoint writes per ~60-byte sample before the change (#1966). The +accepted cost is on metric-filtered reads, which now scan a resource's other +metrics in the window: a single-metric 500-node query moved from 45 to 59 +microseconds and the 24-hour 50-disk `smart_temp` batch from about 0.25 to 0.49 +s, while all-metric dashboard reads are unchanged. Legacy databases keep their +existing unique index authoritative while the replacement index is built, but +that O(rows) migration must run on deferred startup maintenance and swap indexes +in one SQLite transaction so constructor latency, concurrent reads, uniqueness, +and crash recovery remain intact; a database with no unique identity index at +all migrates synchronously so upserts never run without a conflict target. The +deterministic issue-1124 profile must continue to attribute writes by logical +payload, table/index `dbstat` pages, page-cache writes/spills, WAL frames, +checkpoint frames, process write bytes where the host exposes them, and traced +sync calls rather than treating final file size as a proxy. Its checked +invariant persists 157,452 samples across 2,197 mixed provider resources and +caps the no-auto-checkpoint workload at 23,000 WAL frames; the v6.1.1 four-index +schema produces 50,516 frames, the metric-major lookup schema 36,308, and the +time-major identity schema 20,223. Checkpoint-threshold changes must also prove +physical writes, WAL allocation, write latency, restart, and the four-connection +concurrent-read case; a no-reader byte reduction alone is not sufficient to +widen the production WAL bound. Retention must also return freed SQLite pages to the OS proportionally to the current freelist, bounded per cycle, and on every retention cycle rather than only when that cycle deleted rows: a fixed small diff --git a/pkg/metrics/identity_index_test.go b/pkg/metrics/identity_index_test.go new file mode 100644 index 000000000..d02d3d623 --- /dev/null +++ b/pkg/metrics/identity_index_test.go @@ -0,0 +1,188 @@ +package metrics + +import ( + "fmt" + "path/filepath" + "testing" + "time" + + "github.com/rs/zerolog" + "github.com/rs/zerolog/log" +) + +// The issue-1124 invariant starts from an empty table, where each series' +// few samples share leaf pages with its neighbours whichever column order the +// identity tree uses. A running install retains hours of history per series, +// so a metric-major identity tree gives every series its own insertion page +// and each commit rewrites one page per series. These tests seed that retained +// history before measuring. + +const ( + steadyStateResources = 60 + steadyStateHistory = 360 // one hour of 10-second raw samples per series + steadyStateTicks = 30 +) + +var steadyStateMetricTypes = []string{"cpu", "memory", "disk", "netin", "netout", "diskread", "diskwrite", "temperature"} + +var steadyStateBase = time.Unix(1_700_000_000, 0).UTC() + +// newSteadyStateStore returns a store holding one hour of raw history for +// every series. layout "metric-major" reinstalls the retired unique +// idx_metrics_lookup beside a non-unique time-major index; any other value +// keeps the store's own schema. +func newSteadyStateStore(tb testing.TB, layout string) (*Store, string) { + tb.Helper() + dir := tb.TempDir() + cfg := DefaultConfig(dir) + cfg.DBPath = filepath.Join(dir, "metrics.db") + cfg.FlushInterval = time.Hour + cfg.RollupInterval = time.Hour + cfg.RetentionRaw = 10 * 365 * 24 * time.Hour + store, err := NewStore(cfg) + if err != nil { + tb.Fatalf("NewStore: %v", err) + } + tb.Cleanup(func() { _ = store.Close() }) + if err := store.WaitForMaintenance(30 * time.Second); err != nil { + tb.Fatalf("startup maintenance: %v", err) + } + if layout == "metric-major" { + if _, err := store.db.Exec(` + DROP INDEX idx_metrics_query_all; + CREATE INDEX idx_metrics_query_all + ON metrics(resource_type, resource_id, tier, timestamp, metric_type); + CREATE UNIQUE INDEX idx_metrics_lookup + ON metrics(resource_type, resource_id, metric_type, tier, timestamp); + `); err != nil { + tb.Fatalf("install metric-major identity index: %v", err) + } + } + + store.db.SetMaxOpenConns(1) + store.db.SetMaxIdleConns(1) + tx, err := store.db.Begin() + if err != nil { + tb.Fatalf("begin history seed: %v", err) + } + stmt, err := tx.Prepare(` + INSERT INTO metrics(resource_type, resource_id, metric_type, value, timestamp, tier) + VALUES('vm', ?, ?, ?, ?, 'raw') + `) + if err != nil { + _ = tx.Rollback() + tb.Fatalf("prepare history seed: %v", err) + } + for sample := 0; sample < steadyStateHistory; sample++ { + ts := steadyStateBase.Add(time.Duration(sample) * 10 * time.Second).Unix() + for resource := 0; resource < steadyStateResources; resource++ { + for i, metric := range steadyStateMetricTypes { + if _, err := stmt.Exec(steadyStateResourceID(resource), metric, float64((sample+i)%100), ts); err != nil { + _ = stmt.Close() + _ = tx.Rollback() + tb.Fatalf("seed history: %v", err) + } + } + } + } + if err := stmt.Close(); err != nil { + _ = tx.Rollback() + tb.Fatalf("close history seed: %v", err) + } + if err := tx.Commit(); err != nil { + tb.Fatalf("commit history seed: %v", err) + } + return store, cfg.DBPath +} + +func steadyStateResourceID(resource int) string { + return fmt.Sprintf("cluster-a:pve-%d:%d", resource%4, 100+resource) +} + +// steadyStateWALFrames writes steadyStateTicks polls, one commit each, through +// the production ingestion path and returns the WAL frames they produced. +func steadyStateWALFrames(t *testing.T, layout string) int64 { + t.Helper() + store, dbPath := newSteadyStateStore(t, layout) + if _, err := store.db.Exec(`PRAGMA wal_autocheckpoint=0`); err != nil { + t.Fatalf("disable auto-checkpoint: %v", err) + } + if _, err := store.db.Exec(`PRAGMA wal_checkpoint(TRUNCATE)`); err != nil { + t.Fatalf("reset WAL: %v", err) + } + + for tick := 0; tick < steadyStateTicks; tick++ { + ts := steadyStateBase.Add(time.Duration(steadyStateHistory+tick) * 10 * time.Second) + batch := make([]WriteMetric, 0, steadyStateResources*len(steadyStateMetricTypes)) + for resource := 0; resource < steadyStateResources; resource++ { + for i, metric := range steadyStateMetricTypes { + batch = append(batch, WriteMetric{ + ResourceType: "vm", + ResourceID: steadyStateResourceID(resource), + MetricType: metric, + Value: float64((tick + i) % 100), + Timestamp: ts, + Tier: TierRaw, + }) + } + } + store.WriteBatchSync(batch) + } + + var persisted int + if err := store.db.QueryRow(`SELECT COUNT(*) FROM metrics`).Scan(&persisted); err != nil { + t.Fatalf("count persisted metrics: %v", err) + } + series := steadyStateResources * len(steadyStateMetricTypes) + if want := series * (steadyStateHistory + steadyStateTicks); persisted != want { + t.Fatalf("persisted rows=%d, want %d", persisted, want) + } + pageSize := issue1124PragmaInt(t, store, "page_size") + return (issue1124FileSize(t, dbPath+"-wal") - 32) / (pageSize + 24) +} + +// TestMetricsIdentityIndexSteadyStateWrites pins the steady-state write cost of +// the time-major identity index against the retired metric-major one. Measured +// on 2026-09-24: 25,310 WAL frames for the metric-major layout and 6,022 for +// the time-major layout over 30 polls of 480 series. +func TestMetricsIdentityIndexSteadyStateWrites(t *testing.T) { + suppressTestLogs(t) + metricMajor := steadyStateWALFrames(t, "metric-major") + timeMajor := steadyStateWALFrames(t, "") + t.Logf("WAL frames over %d polls: metric-major=%d time-major=%d", steadyStateTicks, metricMajor, timeMajor) + if timeMajor*100 > metricMajor*35 { + t.Fatalf( + "time-major identity index wrote %d WAL frames against %d for the metric-major layout; want at most 35%%", + timeMajor, + metricMajor, + ) + } +} + +// BenchmarkQuerySingleMetricOfWideResource reads one metric of a resource that +// carries eight. A time-major identity index reads the resource's other +// metrics in the window as well, so this is the read the index order can slow. +func BenchmarkQuerySingleMetricOfWideResource(b *testing.B) { + logger := log.Logger + log.Logger = zerolog.Nop() + b.Cleanup(func() { log.Logger = logger }) + store, _ := newSteadyStateStore(b, "") + store.db.SetMaxOpenConns(4) + store.db.SetMaxIdleConns(4) + start := steadyStateBase.Add(-time.Second) + end := steadyStateBase.Add(time.Duration(steadyStateHistory) * 10 * time.Second) + points, err := store.Query("vm", steadyStateResourceID(7), "memory", start, end, 0) + if err != nil { + b.Fatalf("Query: %v", err) + } + if len(points) != steadyStateHistory { + b.Fatalf("points=%d, want %d", len(points), steadyStateHistory) + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + if _, err := store.Query("vm", steadyStateResourceID(7), "memory", start, end, 0); err != nil { + b.Fatalf("Query: %v", err) + } + } +} diff --git a/pkg/metrics/issue1124_write_amplification_test.go b/pkg/metrics/issue1124_write_amplification_test.go index 016227deb..6959f93db 100644 --- a/pkg/metrics/issue1124_write_amplification_test.go +++ b/pkg/metrics/issue1124_write_amplification_test.go @@ -399,8 +399,10 @@ func TestStoreIdentityMigrationKeepsReadersAvailable(t *testing.T) { // TestMetricsWriteAmplificationInvariant persists 157,452 deterministic // multi-provider samples, enough to force repeated B-tree splits and expose an // accidentally restored identity index. The frame ceiling includes 14% margin -// over the consolidated schema's measured 35,030 frames; the v6.1.1 schema -// produces 50,516 frames and fails this invariant. +// over the time-major identity schema's measured 20,223 frames; the retired +// metric-major lookup schema produces 36,308 frames and the v6.1.1 schema +// 50,516, and both fail this invariant. It starts from an empty table, so +// TestMetricsIdentityIndexSteadyStateWrites covers retained history. func TestMetricsWriteAmplificationInvariant(t *testing.T) { suppressTestLogs(t) dir := t.TempDir() @@ -463,7 +465,7 @@ func TestMetricsWriteAmplificationInvariant(t *testing.T) { pageSize := issue1124PragmaInt(t, store, "page_size") walBytes := issue1124FileSize(t, dbPath+"-wal") walFrames := (walBytes - 32) / (pageSize + 24) - const maxWALFrames = 40_000 + const maxWALFrames = 23_000 if walFrames > maxWALFrames { t.Fatalf( "WAL frames=%d for %d samples, limit=%d; redundant indexes or transaction churn regressed write amplification", @@ -493,8 +495,10 @@ func TestMetricsWriteAmplificationInvariant(t *testing.T) { // multi-provider estate. It reports WAL frames and per-B-tree dbstat pages // separately, so retained database size is never used as a write-volume proxy. // -// Run with PULSE_METRICS_WRITE_PROFILE=full (v6.1.1 layout) or consolidated -// and optionally PULSE_METRICS_WRITE_PROFILE_TICKS=N. The default 30 ticks +// Run with PULSE_METRICS_WRITE_PROFILE=full (v6.1.1 layout), lookup (the +// v6.1.2 to v6.4 metric-major unique index) or consolidated (the time-major +// identity index), each optionally -batched, and optionally +// PULSE_METRICS_WRITE_PROFILE_TICKS=N. The default 30 ticks // persist 393,630 samples across 2,197 resources. Set // PULSE_METRICS_WRITE_PROFILE_AUTOCHECKPOINT=1 for production checkpoint/fsync // tracing, PULSE_METRICS_WRITE_PROFILE_CHECKPOINT_PAGES=N for a bounded @@ -506,9 +510,9 @@ func TestIssue1124WriteAmplificationProfile(t *testing.T) { if mode == "" { t.Skip("set " + issue1124ProfileEnv + " to run the write-amplification profile") } - if mode != "full" && mode != "consolidated" && mode != "full-batched" && - mode != "consolidated-batched" && mode != "without-rowid" { - t.Fatalf("%s must be full, consolidated, full-batched, consolidated-batched, or without-rowid, got %q", issue1124ProfileEnv, mode) + if mode != "full" && mode != "lookup" && mode != "consolidated" && mode != "full-batched" && + mode != "lookup-batched" && mode != "consolidated-batched" && mode != "without-rowid" { + t.Fatalf("%s must be full, lookup, consolidated, full-batched, lookup-batched, consolidated-batched, or without-rowid, got %q", issue1124ProfileEnv, mode) } ticks := 30 @@ -562,7 +566,9 @@ func TestIssue1124WriteAmplificationProfile(t *testing.T) { // Reconstruct the v6.1.1 schema so released and consolidated layouts // remain comparable after the production migration lands. if _, err := store.db.Exec(` - DROP INDEX idx_metrics_lookup; + DROP INDEX idx_metrics_query_all; + CREATE INDEX idx_metrics_query_all + ON metrics(resource_type, resource_id, tier, timestamp, metric_type); CREATE INDEX idx_metrics_lookup ON metrics(resource_type, resource_id, metric_type, tier, timestamp); CREATE UNIQUE INDEX idx_metrics_unique @@ -572,6 +578,20 @@ func TestIssue1124WriteAmplificationProfile(t *testing.T) { t.Fatalf("install v6.1.1 index layout: %v", err) } } + if mode == "lookup" || mode == "lookup-batched" { + // Reconstruct the v6.1.2 to v6.4 schema, whose unique identity tree + // was metric-major beside a non-unique time-major read index. + if _, err := store.db.Exec(` + DROP INDEX idx_metrics_query_all; + CREATE INDEX idx_metrics_query_all + ON metrics(resource_type, resource_id, tier, timestamp, metric_type); + CREATE UNIQUE INDEX idx_metrics_lookup + ON metrics(resource_type, resource_id, metric_type, tier, timestamp); + `); err != nil { + _ = store.Close() + t.Fatalf("install metric-major lookup index layout: %v", err) + } + } if mode == "without-rowid" { if _, err := store.db.Exec(` DROP TABLE metrics; @@ -671,7 +691,7 @@ func TestIssue1124WriteAmplificationProfile(t *testing.T) { batch[i] = metric logicalBytes += issue1124LogicalBytes(metric) } - if mode == "full-batched" || mode == "consolidated-batched" { + if mode == "full-batched" || mode == "lookup-batched" || mode == "consolidated-batched" { tickBatch = append(tickBatch, batch...) } else { writeStarted := time.Now() @@ -850,7 +870,7 @@ func TestIssue1124WriteAmplificationProfile(t *testing.T) { RetentionFreelist: retentionFreelist, DBStats: dbStats, } - if mode == "full-batched" || mode == "consolidated-batched" { + if mode == "full-batched" || mode == "lookup-batched" || mode == "consolidated-batched" { report.WriteCalls = ticks } encoded, err := json.Marshal(report) @@ -1191,19 +1211,12 @@ func issue1124CreateLegacyStore(t *testing.T, dbPath string, rows int, duplicate func issue1124AssertConsolidatedIndexes(t *testing.T, store *Store) { t.Helper() - matches, err := store.metricsIndexMatches("idx_metrics_lookup", true, metricsIdentityColumns) + matches, err := store.metricsIndexMatches(metricsIdentityIndex, true, metricsIdentityColumns) if err != nil { - t.Fatalf("inspect consolidated lookup index: %v", err) + t.Fatalf("inspect consolidated identity index: %v", err) } if !matches { - t.Fatal("idx_metrics_lookup is not the expected unique identity/range index") - } - exists, err := store.metricsIndexExists("idx_metrics_unique") - if err != nil { - t.Fatalf("inspect obsolete unique index: %v", err) - } - if exists { - t.Fatal("obsolete idx_metrics_unique still exists") + t.Fatalf("%s is not the expected unique time-major identity index", metricsIdentityIndex) } rows, err := store.db.Query(`PRAGMA index_list(metrics)`) @@ -1228,9 +1241,10 @@ func issue1124AssertConsolidatedIndexes(t *testing.T, store *Store) { if err := rows.Err(); err != nil { t.Fatalf("consolidated index rows: %v", err) } + // Retired metric-major trees (idx_metrics_lookup, idx_metrics_unique) + // must be gone: each would be dirtied by every insert. want := map[string]bool{ - "idx_metrics_lookup": true, - "idx_metrics_query_all": false, + "idx_metrics_query_all": true, "idx_metrics_tier_time": false, } if len(indexes) != len(want) { diff --git a/pkg/metrics/store.go b/pkg/metrics/store.go index 79e6af59a..fdfeb29fb 100644 --- a/pkg/metrics/store.go +++ b/pkg/metrics/store.go @@ -10,6 +10,7 @@ import ( "net/url" "os" "path/filepath" + "slices" "strconv" "strings" "sync" @@ -467,7 +468,9 @@ func (s *Store) initSchema() error { CREATE INDEX IF NOT EXISTS idx_metrics_tier_time ON metrics(tier, timestamp); - -- Covering index for Unified History (QueryAll) performance + -- Identity and range-read index. Created non-unique here so reads + -- can name it on every schema generation; ensureMetricsIdentityIndex + -- makes it the unique identity index. CREATE INDEX IF NOT EXISTS idx_metrics_query_all ON metrics(resource_type, resource_id, tier, timestamp, metric_type); @@ -494,41 +497,83 @@ func (s *Store) initSchema() error { return nil } -var metricsIdentityColumns = []string{"resource_type", "resource_id", "metric_type", "tier", "timestamp"} +// metricsIdentityIndex is the one B-tree that enforces sample identity and +// serves every range read. Its columns are time-major within a resource, so +// the samples one poll writes for a resource share leaf pages. The retired +// metric-major order gave every series its own insertion point, and each +// commit then rewrote one page per series through the WAL and checkpoint. +// Upsert conflict targets match it by column set, not order. +const metricsIdentityIndex = "idx_metrics_query_all" + +var metricsIdentityColumns = []string{"resource_type", "resource_id", "tier", "timestamp", "metric_type"} + +// retiredMetricsIdentityIndexes are the metric-major trees earlier schemas +// kept over the same five columns: idx_metrics_unique up to v6.1.1 and the +// unique idx_metrics_lookup that replaced it. Either is authoritative for +// identity until the replacement commits. +var retiredMetricsIdentityIndexes = []string{"idx_metrics_lookup", "idx_metrics_unique"} // ensureMetricsIdentityIndex keeps one B-tree for both metric identity and -// single-series range queries. Older databases have two indexes containing the -// same five columns in different orders: idx_metrics_lookup serves reads while -// idx_metrics_unique enforces identity. Every insert dirties both trees, which -// materially amplifies WAL and checkpoint writes. +// range queries. Every additional tree over the same columns is dirtied by +// every insert, which materially amplifies WAL and checkpoint writes. func (s *Store) ensureMetricsIdentityIndex() error { - lookupCurrent, err := s.metricsIndexMatches("idx_metrics_lookup", true, metricsIdentityColumns) + current, err := s.metricsIndexMatches(metricsIdentityIndex, true, metricsIdentityColumns) if err != nil { return fmt.Errorf("inspect metrics identity index: %w", err) } - legacyUniqueExists, err := s.metricsIndexExists("idx_metrics_unique") + retired, retiredUnique, err := s.existingRetiredMetricsIdentityIndexes() if err != nil { - return fmt.Errorf("inspect legacy metrics unique index: %w", err) + return fmt.Errorf("inspect retired metrics identity indexes: %w", err) } - - if legacyUniqueExists { - // v6.1.1 and earlier already have an authoritative unique index. Keep - // serving with that crash-safe schema and defer the O(rows) rebuild to - // the startup-maintenance worker so NewStore latency stays bounded. - s.identityMigrationPending.Store(true) - log.Info().Msg("Scheduled metrics identity index consolidation") + if current && len(retired) == 0 { return nil } - if lookupCurrent { + if current || retiredUnique { + // Identity is already enforced by a unique index. Keep serving with + // that crash-safe schema and defer the O(rows) rebuild to the + // startup-maintenance worker so NewStore latency stays bounded. + s.identityMigrationPending.Store(true) + log.Info().Strs("retired_indexes", retired).Msg("Scheduled metrics identity index consolidation") return nil } return s.migrateMetricsIdentityIndex() } +// existingRetiredMetricsIdentityIndexes reports which retired trees remain +// and whether any of them still enforces identity. +func (s *Store) existingRetiredMetricsIdentityIndexes() (existing []string, anyUnique bool, err error) { + rows, err := s.db.Query(`PRAGMA index_list(metrics)`) + if err != nil { + return nil, false, err + } + defer rows.Close() + + for rows.Next() { + var ( + sequence int + index string + unique int + origin string + partial int + ) + if err := rows.Scan(&sequence, &index, &unique, &origin, &partial); err != nil { + return nil, false, err + } + if !slices.Contains(retiredMetricsIdentityIndexes, index) { + continue + } + existing = append(existing, index) + if unique == 1 && partial == 0 { + anyUnique = true + } + } + return existing, anyUnique, rows.Err() +} + func (s *Store) migrateMetricsIdentityIndex() error { if err := s.replaceMetricsIdentityIndex(); err == nil { - log.Info().Msg("Consolidated metrics lookup and identity indexes") + log.Info().Msg("Consolidated metrics identity index") return nil } else if !isUniqueConstraintError(err) { return fmt.Errorf("consolidate metrics identity index: %w", err) @@ -546,27 +591,41 @@ func (s *Store) migrateMetricsIdentityIndex() error { return nil } +// replaceMetricsIdentityIndex rebuilds the identity index as unique when it is +// not already, then drops the retired trees, all in one transaction so a crash +// leaves the previous authoritative index in place. func (s *Store) replaceMetricsIdentityIndex() error { + current, err := s.metricsIndexMatches(metricsIdentityIndex, true, metricsIdentityColumns) + if err != nil { + return fmt.Errorf("inspect metrics identity index: %w", err) + } + tx, err := s.db.Begin() if err != nil { return fmt.Errorf("begin identity index migration: %w", err) } defer func() { _ = tx.Rollback() }() - if _, err := tx.Exec(`DROP INDEX IF EXISTS idx_metrics_lookup`); err != nil { - return fmt.Errorf("drop old lookup index: %w", err) - } - if metricsIdentityMigrationHook != nil { + if !current { + if _, err := tx.Exec(`DROP INDEX IF EXISTS ` + metricsIdentityIndex); err != nil { + return fmt.Errorf("drop non-unique identity index: %w", err) + } + if metricsIdentityMigrationHook != nil { + metricsIdentityMigrationHook() + } + if _, err := tx.Exec(` + CREATE UNIQUE INDEX ` + metricsIdentityIndex + ` + ON metrics(` + strings.Join(metricsIdentityColumns, ", ") + `) + `); err != nil { + return fmt.Errorf("create unique identity index: %w", err) + } + } else if metricsIdentityMigrationHook != nil { metricsIdentityMigrationHook() } - if _, err := tx.Exec(` - CREATE UNIQUE INDEX idx_metrics_lookup - ON metrics(resource_type, resource_id, metric_type, tier, timestamp) - `); err != nil { - return fmt.Errorf("create consolidated lookup index: %w", err) - } - if _, err := tx.Exec(`DROP INDEX IF EXISTS idx_metrics_unique`); err != nil { - return fmt.Errorf("drop old unique index: %w", err) + for _, name := range retiredMetricsIdentityIndexes { + if _, err := tx.Exec(`DROP INDEX IF EXISTS ` + name); err != nil { + return fmt.Errorf("drop retired identity index %s: %w", name, err) + } } if err := tx.Commit(); err != nil { return fmt.Errorf("commit identity index migration: %w", err) @@ -1590,10 +1649,7 @@ func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, st endParam := startParam + 1 tierParam := endParam + 1 stepParam := tierParam + len(tiers) - index := "idx_metrics_query_all" - if len(metricTypes) > 0 { - index = "idx_metrics_lookup" - } + index := metricsIdentityIndex scope := func(alias string, tierIndex int) string { clause := alias + ".resource_type = :p1 AND " + alias + ".resource_id IN (" + idSlots + ")" if len(metricTypes) > 0 { @@ -1618,8 +1674,11 @@ func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, st // fallback observation. Both checks use this statement's snapshot. branch += " AND (NOT EXISTS (SELECT 1 FROM metrics AS coverage INDEXED BY " + index + " WHERE " + scope("coverage", j) + ") OR NOT EXISTS (" bucket := max(tierBucketSeconds(tier), tierBucketSeconds(preferred)) + // Name the identity index here too: with a time-major identity + // tree the planner otherwise seeks this correlated probe on + // (tier, timestamp) alone and walks every resource in the bucket. branch += fmt.Sprintf(` - SELECT 1 FROM metrics AS h + SELECT 1 FROM metrics AS h INDEXED BY `+index+` WHERE h.resource_type = m.resource_type AND h.resource_id = m.resource_id AND h.metric_type = m.metric_type AND h.tier = :p%d AND h.timestamp >= MAX(:p%d, (m.timestamp / %d) * %d) @@ -1698,10 +1757,9 @@ func (s *Store) queryRetainedChunk(resourceType string, resourceIDs []string, me // order the value projection for each presence probe. idSlots := strings.TrimSuffix(strings.Repeat("?,", len(resourceIDs)), ",") presenceScope := "resource_type = ? AND resource_id IN (" + idSlots + ")" - index := "idx_metrics_query_all" + index := metricsIdentityIndex if len(metricTypes) > 0 { presenceScope += " AND metric_type IN (" + strings.TrimSuffix(strings.Repeat("?,", len(metricTypes)), ",") + ")" - index = "idx_metrics_lookup" } presenceScope += " AND tier = ? AND timestamp >= ? AND timestamp <= ?" check := "EXISTS (SELECT 1 FROM metrics INDEXED BY " + index + " WHERE " + presenceScope + ")" diff --git a/pkg/metrics/store_query_plan_test.go b/pkg/metrics/store_query_plan_test.go index d894a45c6..4df466fe9 100644 --- a/pkg/metrics/store_query_plan_test.go +++ b/pkg/metrics/store_query_plan_test.go @@ -83,7 +83,7 @@ func TestQueryPlansUseIndexes(t *testing.T) { AND tier = ? AND timestamp >= ? AND timestamp < ? GROUP BY resource_type, resource_id, metric_type, bucket_ts`, args: []any{int64(60), int64(60), "minute", "vm", "vm-1", "cpu", "raw", int64(0), farFuture}, - wantIndex: "idx_metrics_lookup", + wantIndex: metricsIdentityIndex, }, { name: "batched rollup aggregation insert", @@ -102,8 +102,8 @@ func TestQueryPlansUseIndexes(t *testing.T) { GROUP BY resource_type, resource_id, metric_type, bucket_ts`, args: []any{int64(60), int64(60), "minute", "raw", int64(0), farFuture}, // The SELECT filters on (tier, timestamp) without resource/metric - // columns. SQLite uses an index-ordered scan (SCAN USING INDEX) - // on idx_metrics_lookup. A TEMP B-TREE may still be used for + // columns. SQLite may use an index-ordered scan (SCAN USING INDEX) + // on the identity index. A TEMP B-TREE may still be used for // GROUP BY. The INSERT side uses the same unique index for ON CONFLICT. allowIndexScan: true, }, @@ -359,7 +359,7 @@ func containsCoveringIndexScan(plan string) bool { } // containsIndexScan returns true if any plan line shows an index-ordered scan -// on the metrics table (e.g., "SCAN metrics USING INDEX idx_metrics_lookup"). +// on the metrics table (e.g., "SCAN metrics USING INDEX idx_metrics_query_all"). // This differs from a covering-index scan in that it reads table rows via the // index, but still avoids a full table scan. func containsIndexScan(plan string) bool { @@ -405,13 +405,10 @@ func newPlanTestDB(t *testing.T) *sql.DB { tier TEXT NOT NULL DEFAULT 'raw' ); - CREATE UNIQUE INDEX IF NOT EXISTS idx_metrics_lookup - ON metrics(resource_type, resource_id, metric_type, tier, timestamp); - CREATE INDEX IF NOT EXISTS idx_metrics_tier_time ON metrics(tier, timestamp); - CREATE INDEX IF NOT EXISTS idx_metrics_query_all + CREATE UNIQUE INDEX IF NOT EXISTS idx_metrics_query_all ON metrics(resource_type, resource_id, tier, timestamp, metric_type); CREATE TABLE IF NOT EXISTS metrics_meta (