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
This commit is contained in:
rcourtman 2026-09-24 08:45:38 +01:00
parent e2b41c379d
commit 2c8cb8435d
5 changed files with 359 additions and 87 deletions

View file

@ -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

View file

@ -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)
}
}
}

View file

@ -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) {

View file

@ -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 + ")"

View file

@ -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 (