mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 04:38:48 +00:00
perf(metrics): reuse retained query plans and output series
Canonical tier reconciliation rebuilt SQL and probed absent preferred tiers for each fallback point, regressing batch reads and allocation costs. Reuse bounded query shapes with current bindings and snapshot-scoped absence checks, then append consecutive points directly to their output series. Preserve coverage and ordering semantics and verify fresh bindings after new preferred observations arrive. Integrate current main test additions.
This commit is contained in:
commit
f26668aa6d
9 changed files with 489 additions and 71 deletions
|
|
@ -62,7 +62,7 @@ reproduction evidence, not a representative customer success rate.
|
|||
| Step | Work | Acceptance | Current state |
|
||||
|---|---|---|---|
|
||||
| 1. Product contract and baseline | Map the current loop and sources of judgment. Record telemetry populations and gaps. | Every identified decision has an owner. Activity is not labelled usefulness. | Complete for this redesign scope. Contract, ownership decisions and baseline limits are recorded. |
|
||||
| 2. Shared evidence | Preserve canonical risk reasons and SMART counters, source/time semantics and history across tools/turns. | Regression tests preserve unknown versus zero and all canonical evidence. Real responses can inspect the same facts as the product. | Implemented and qualified for the named shared-evidence defects. Canonical disk detail, risk and cadence pass real data-path proof. Affected package and concurrency checks pass. Commit f01db995ed corrects the PR benchmark regressions. Exact-base worker comparisons and full metrics/database and focused race checks pass. Final landing CI remains open. Real-model interpretation failures remain tracked in step 5. |
|
||||
| 2. Shared evidence | Preserve canonical risk reasons and SMART counters, source/time semantics and history across tools/turns. | Regression tests preserve unknown versus zero and all canonical evidence. Real responses can inspect the same facts as the product. | Implemented and qualified for the named shared-evidence defects. Canonical disk detail, risk and cadence pass real data-path proof. Affected package and concurrency checks pass. Integrated CI later exposed remaining query and allocation regressions. The final bounded query-reuse correction passes complete selected exact-base worker comparisons and full metrics/database and focused race checks. Final landing CI remains open. Real-model interpretation failures remain tracked in step 5. |
|
||||
| 3. Diagnostic orchestration | Correct proposal-as-proof. Audit triage budgets, unmatched-signal evaluation, assessment completion and investigation cutoffs. | No code-written causal conclusion. No quality inferred from tool, flag or finding counts. Each retained pass has an objective reason. Safety boundaries and incomplete outcomes remain explicit. | Proposal promotion and capture inference were removed in c5d2f56dda. Commit 668af3fe6b removes investigation success-call floors, checkpoint instructions and generic call-count wrap-up rules. The detection slice removes contextless follow-up passes, flag/report-count policy and first-finding completion modes. Full chat and AI suites, focused API and conversation race tests pass. Real-model/action outcome qualification remains open. |
|
||||
| 4. Issue through verified outcome | Follow existing issue/investigation/action records into Assistant, approval, execution and independent readback. | Accepted proposal is visibly distinct from execution and verification. Rejected or unsupported actions do not become success. Uncertainty can survive an action proposal. | Existing foundation, full journey qualification pending. |
|
||||
| 5. Ground-truth qualification and landing | Extend existing qualification tooling only where necessary. Exercise healthy/unhealthy, dependency, missing-access, storage/backup and approved/rejected action cases. Inspect the final browser journey at desktop and narrow widths. | Record exact source/model/permissions, evidence, decisions, faults/misses, latency and verification. Fix in-scope failures, pass appropriate proofs and land scoped commits. | Pending. |
|
||||
|
|
@ -910,3 +910,52 @@ corresponding independent lab probe. The ordinary homelab file-read failure
|
|||
proof remains useful but does not close that autonomous qualification gap.
|
||||
Those cases require suitable independent ground truth through the existing
|
||||
qualification framework before the overall goal can complete.
|
||||
|
||||
## Retained-read performance correction, 2026-09-06
|
||||
|
||||
Build and Test run `34008823529` for integration `f779bf064ab4` failed its
|
||||
benchmark comparison against exact main `3f74c0c27304`. All other jobs passed,
|
||||
and Core E2E run `34008823534` passed all eight browser shards. The benchmark
|
||||
failure comprised fifteen time/allocation comparisons, including batch reads
|
||||
up to 59% slower and small-read bytes about 198% higher. Earlier narrow worker
|
||||
comparisons omitted the actual QueryAllBatch cases and did not establish this
|
||||
integrated performance result. The initial status report calling that CI job
|
||||
passed was incorrect and was explicitly corrected.
|
||||
|
||||
The canonical reader now reuses bounded SQL templates and numbered bindings,
|
||||
checks absent preferred tiers once within the current statement, and appends
|
||||
consecutive results directly to their series. It retains tier reconciliation,
|
||||
current snapshots, scope isolation and output semantics. No benchmark threshold
|
||||
was relaxed. An initial correction still regressed the plain raw-read case by
|
||||
11% and was revised before landing.
|
||||
|
||||
Final production `pkg/metrics/store.go` SHA-256:
|
||||
`9a66d1acea82c17ca540abe9a9ee66e089a36420620e2ebb8ca7e6d102191a8d`.
|
||||
Ten alternating 100ms samples on pulse-dev, Go 1.26.8 and GOMAXPROCS=4,
|
||||
compare the final implementation with exact base `3f74c0c27304`. The complete
|
||||
selected Query, QueryAllBatch, RollupCandidate, fleet dashboard, history API,
|
||||
chart batch and NormalizeRoute benchmark families have no statistically
|
||||
significant greater-than-10% regression in time, bytes or allocations. Plain
|
||||
raw reads show no significant time change. Batch reads improve 16–42%, with
|
||||
bytes reduced 28–34%. Single-metric downsampling is 5.66% slower and the
|
||||
1,000-point rollup candidate is 6.04% slower, both inside the unchanged gate.
|
||||
The bounded history API shows no significant time change. These are selected
|
||||
worker comparisons, not a claim that final remote CI has passed.
|
||||
|
||||
Private raw samples and benchstat outputs are under
|
||||
`tmp/patrol-f779-v2-selected/` and `tmp/patrol-f779-v2-normalize/` at the
|
||||
workspace root. Focused tests cover newly appearing preferred buckets after a
|
||||
cached absent-tier read, current resource family, metric, window and display
|
||||
step bindings, per-series overlap and batch parity. Integration with main
|
||||
`6c000837e27b` brings three test-only changes and no additional runtime changes.
|
||||
Full metrics and database packages pass on pulse-dev in 77.265s and 0.272s.
|
||||
Focused retained coverage, fresh bindings and batch identity race proof passes
|
||||
in 5.349s. The two latency-based concurrent SLO tests run in the ordinary suite
|
||||
but explicitly skip under the race detector. A separate canonical hot-path
|
||||
regression exercises eight concurrent query scopes and mixed display steps
|
||||
through a one-connection pool without latency assertions. It passes normally
|
||||
(0.043s) and under race (2.106s), with no cross-query binding leakage or lock
|
||||
inversion. Its worker log is `patrol-retained-concurrent-bindings.log`.
|
||||
Private full/race logs are under
|
||||
`tmp/patrol-f779-final-proof/`. Exact staged hook and remote landing remain
|
||||
pending. This correction does not close the outstanding real-model/action outcome gap.
|
||||
|
|
|
|||
File diff suppressed because one or more lines are too long
|
|
@ -45,8 +45,18 @@ Display-aggregated reads use one transaction. Indexed existence checks identify
|
|||
which retention tiers have observations in that snapshot. Every present tier
|
||||
participates in the shared indexed overlap query. Presence never stands in for
|
||||
per-series coverage. The store retains at most 32 compiled read-statement shapes
|
||||
and evaluates them again inside each current read snapshot.
|
||||
It never caches tier presence, query results or timestamp windows. Less common
|
||||
and 32 SQL templates keyed only by parameter counts, tier order and aggregation
|
||||
shape. Numbered SQLite bindings supply each current identity, window and display
|
||||
step once across all branches. An uncorrelated existence guard in the same
|
||||
statement avoids per-observation overlap probes when a preferred tier is absent.
|
||||
If any preferred observation exists, the correlated same-series check still
|
||||
owns coverage. Both checks are reevaluated in the current snapshot.
|
||||
The template lock performs no database I/O and is separate from statement
|
||||
preparation. The store never caches tier presence, query results or timestamp
|
||||
windows. Consecutive output points append directly to the current series slice,
|
||||
flushing it to the result map on series change and completion. Interleaved
|
||||
metrics resume their existing slices. A single chunk returns its result directly
|
||||
without copying the outer map. Less common
|
||||
shapes run uncached after the bound is reached. Preparation occurs before
|
||||
acquiring the transaction, including with a single-connection pool. The shared
|
||||
database instrumentation preserves timing for transaction-bound prepared
|
||||
|
|
|
|||
92
internal/alerts/active_checkpoint_failure_test.go
Normal file
92
internal/alerts/active_checkpoint_failure_test.go
Normal file
|
|
@ -0,0 +1,92 @@
|
|||
package alerts
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A recovery-file failure must be reported without discarding a successful
|
||||
// SQLite checkpoint, including an empty checkpoint after resolution.
|
||||
func TestSQLiteCheckpointSurvivesRecoveryMirrorRenameFailure(t *testing.T) {
|
||||
for _, resolved := range []bool{false, true} {
|
||||
name := "fired"
|
||||
if resolved {
|
||||
name = "resolved"
|
||||
}
|
||||
t.Run(name, func(t *testing.T) {
|
||||
dataDir := t.TempDir()
|
||||
alert := durableRestoreAlert(time.Now().Add(-72 * time.Hour).UTC())
|
||||
initial := []*Alert{}
|
||||
if resolved {
|
||||
initial = append(initial, alert)
|
||||
}
|
||||
writeActiveRecoveryFixture(t, dataDir, initial)
|
||||
m := NewManagerWithDataDir(dataDir, WithDurableAlertStore())
|
||||
t.Cleanup(m.Stop)
|
||||
if !m.activeStateAuthoritative.Load() {
|
||||
t.Fatal("durable constructor did not establish SQLite authority")
|
||||
}
|
||||
|
||||
// A non-empty directory forces rename failure even when tests run as
|
||||
// root; chmod-based fault injection would not reliably do so.
|
||||
mirror := filepath.Join(dataDir, "alerts", "active-alerts.json")
|
||||
if err := os.Remove(mirror); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.Mkdir(mirror, alertsDirPerm); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sentinel := filepath.Join(mirror, "preserve")
|
||||
if err := os.WriteFile(sentinel, []byte("unchanged"), alertsFilePerm); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Change only memory: SaveActiveAlerts, not a lifecycle append,
|
||||
// must be responsible for the durable state checked after restart.
|
||||
m.mu.Lock()
|
||||
if resolved {
|
||||
m.removeActiveAlertNoLock(alert.ID)
|
||||
} else {
|
||||
m.setActiveAlertNoLock(alert.ID, alert)
|
||||
}
|
||||
m.mu.Unlock()
|
||||
err := m.SaveActiveAlerts()
|
||||
if err == nil || !strings.Contains(err.Error(), "SQLite active alert checkpoint succeeded but recovery persistence failed") || !strings.Contains(err.Error(), "failed to rename") {
|
||||
t.Fatalf("checkpoint error = %v; want reported mirror rename failure after SQLite success", err)
|
||||
}
|
||||
// Close SQLite before shutdown so its final save cannot repair a
|
||||
// missing checkpoint and conceal a failure of the explicit save.
|
||||
m.SetEventLog(nil)
|
||||
m.Stop()
|
||||
remaining, err := filepath.Glob(filepath.Join(dataDir, "alerts", "active-alerts-*.json.tmp"))
|
||||
if err != nil || len(remaining) != 0 {
|
||||
t.Fatalf("temporary mirrors left after failed writes = %v, error = %v", remaining, err)
|
||||
}
|
||||
if got, err := os.ReadFile(sentinel); err != nil || string(got) != "unchanged" {
|
||||
t.Fatalf("rename failure damaged destination: %q, %v", got, err)
|
||||
}
|
||||
if err := os.Remove(sentinel); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.Remove(mirror); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
restarted := NewManagerWithDataDir(dataDir, WithDurableAlertStore())
|
||||
t.Cleanup(restarted.Stop)
|
||||
if !restarted.activeStateAuthoritative.Load() {
|
||||
t.Fatal("restart did not establish SQLite authority")
|
||||
}
|
||||
alerts := restarted.snapshotActiveAlerts()
|
||||
if resolved {
|
||||
if len(alerts) != 0 {
|
||||
t.Fatalf("resolved alert resurrected: %+v", alerts)
|
||||
}
|
||||
} else if len(alerts) != 1 || alerts[0].ID != alert.ID || !alerts[0].StartTime.Equal(alert.StartTime) || !alerts[0].Acknowledged || alerts[0].AckUser != alert.AckUser {
|
||||
t.Fatalf("checkpoint lost incident identity or acknowledgement: %+v", alerts)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
102
internal/alerts/pbs_datastore_restart_test.go
Normal file
102
internal/alerts/pbs_datastore_restart_test.go
Normal file
|
|
@ -0,0 +1,102 @@
|
|||
package alerts
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/alerts/eventlog"
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/models"
|
||||
)
|
||||
|
||||
// Missing capacity and a measured empty datastore both arrive with Usage=0.
|
||||
// Only the latter may recover an incident, including after SQLite restore.
|
||||
func TestPBSDatastoreMissingCapacityRestartCallbacks(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
fired, resolved := make(chan string, 16), make(chan string, 16)
|
||||
start := func() *Manager {
|
||||
m := NewManagerWithDataDir(dir)
|
||||
t.Cleanup(m.Stop)
|
||||
m.EnableEventLog()
|
||||
if !m.activeStateAuthoritative.Load() {
|
||||
t.Fatal("SQLite active state is not authoritative")
|
||||
}
|
||||
m.UpdateConfig(AlertConfig{Enabled: true, ActivationState: ActivationActive,
|
||||
StorageDefault: HysteresisThreshold{Trigger: 95, Clear: 90},
|
||||
Overrides: map[string]ThresholdConfig{"pbs-primary/backups": {Usage: &HysteresisThreshold{Trigger: 80, Clear: 70}}},
|
||||
})
|
||||
disableTestTimeThresholds(m)
|
||||
m.SetAlertCallback(func(a *Alert) { fired <- a.ID })
|
||||
m.SetResolvedCallback(func(id string) { resolved <- id })
|
||||
return m
|
||||
}
|
||||
storage := models.Storage{ID: "pbs-primary-backups", AliasIDs: []string{"pbs-primary/backups"}, Name: "backups", Instance: "pbs-primary", Type: "pbs", Status: "online", Total: 1000, Used: 850, Free: 150, Usage: 85}
|
||||
id := canonicalMetricStateID(storage.ID, "usage")
|
||||
observe := func(m *Manager, s models.Storage) {
|
||||
for range 5 {
|
||||
m.CheckStorage(s)
|
||||
}
|
||||
}
|
||||
requireCallback := func(ch <-chan string) {
|
||||
t.Helper()
|
||||
select {
|
||||
case got := <-ch:
|
||||
if got != id {
|
||||
t.Fatalf("callback ID = %q, want %q", got, id)
|
||||
}
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("missing lifecycle callback")
|
||||
}
|
||||
}
|
||||
requireQuiet := func() {
|
||||
t.Helper()
|
||||
select {
|
||||
case got := <-fired:
|
||||
t.Fatalf("unexpected firing callback %q", got)
|
||||
case got := <-resolved:
|
||||
t.Fatalf("unexpected recovery callback %q", got)
|
||||
case <-time.After(50 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
m := start()
|
||||
observe(m, storage)
|
||||
original := *testRequireActiveAlert(t, m, id)
|
||||
requireCallback(fired)
|
||||
missing := storage
|
||||
missing.Usage, missing.Total, missing.Used, missing.Free = 0, 0, 0, 0
|
||||
for _, restart := range []bool{false, true} {
|
||||
if restart {
|
||||
m.Stop()
|
||||
m = start()
|
||||
}
|
||||
observe(m, missing)
|
||||
if got := testRequireActiveAlert(t, m, id); !got.StartTime.Equal(original.StartTime) {
|
||||
t.Fatal("missing capacity replaced incident identity")
|
||||
}
|
||||
requireQuiet()
|
||||
}
|
||||
empty := storage
|
||||
empty.Usage, empty.Used, empty.Free = 0, 0, empty.Total
|
||||
observe(m, empty)
|
||||
if testHasActiveAlert(t, m, id) {
|
||||
t.Fatal("confirmed empty datastore did not recover")
|
||||
}
|
||||
requireCallback(resolved)
|
||||
m.Stop()
|
||||
m = start()
|
||||
observe(m, missing)
|
||||
if testHasActiveAlert(t, m, id) {
|
||||
t.Fatal("missing capacity resurrected resolved incident")
|
||||
}
|
||||
requireQuiet()
|
||||
observe(m, storage)
|
||||
if got := testRequireActiveAlert(t, m, id); !got.StartTime.After(original.StartTime) {
|
||||
t.Fatal("refire reused original incident")
|
||||
}
|
||||
requireCallback(fired)
|
||||
requireQuiet()
|
||||
for kind, want := range map[string]int{eventlog.TypeFired: 2, eventlog.TypeResolved: 1} {
|
||||
if got := len(queryAlertEvents(t, m, eventlog.Filter{Types: []string{kind}})); got != want {
|
||||
t.Fatalf("%s events = %d, want %d", kind, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -201,6 +201,12 @@ type Store struct {
|
|||
readMu sync.Mutex
|
||||
readStatements map[string]*pdb.InstrumentedStmt
|
||||
|
||||
// SQL templates contain only parameter positions, never query values.
|
||||
// Keep this lock separate from statement preparation, which may wait for
|
||||
// a database connection while another caller owns a read transaction.
|
||||
queryMu sync.Mutex
|
||||
readQueryShapes map[retainedQueryShape]string
|
||||
|
||||
// Write buffer
|
||||
bufferMu sync.Mutex
|
||||
buffer []bufferedMetric
|
||||
|
|
@ -1296,6 +1302,9 @@ func (s *Store) queryBatch(
|
|||
if len(tiers) == 0 {
|
||||
return map[string]map[string][]MetricPoint{}, nil
|
||||
}
|
||||
if len(unique) <= queryAllBatchChunkSize {
|
||||
return s.queryRetainedChunk(resourceType, unique, normalizedMetricTypes, start, end, stepSecs, tiers)
|
||||
}
|
||||
|
||||
result := make(map[string]map[string][]MetricPoint, len(unique))
|
||||
// Reconcile every series in one snapshot per chunk. A resource with one
|
||||
|
|
@ -1338,13 +1347,64 @@ func normalizeMetricTypes(metricTypes []string) []string {
|
|||
return normalized
|
||||
}
|
||||
|
||||
type retainedQueryShape struct {
|
||||
resources, metrics int
|
||||
tiers [4]Tier
|
||||
aggregate, groupSeries bool
|
||||
}
|
||||
|
||||
func retainedQueryParameters(resourceType string, resourceIDs, metricTypes []string, start, end time.Time, stepSecs int64, tiers []Tier) []interface{} {
|
||||
params := make([]interface{}, 0, len(resourceIDs)+len(metricTypes)+len(tiers)+4)
|
||||
params = append(params, resourceType)
|
||||
for _, id := range resourceIDs {
|
||||
params = append(params, id)
|
||||
}
|
||||
for _, metric := range metricTypes {
|
||||
params = append(params, metric)
|
||||
}
|
||||
params = append(params, start.Unix(), end.Unix())
|
||||
for _, tier := range tiers {
|
||||
params = append(params, string(tier))
|
||||
}
|
||||
if stepSecs > 1 {
|
||||
params = append(params, stepSecs)
|
||||
}
|
||||
return params
|
||||
}
|
||||
|
||||
func (s *Store) retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, start, end time.Time, stepSecs int64, tiers []Tier, groupSeries bool) (string, []interface{}) {
|
||||
shape := retainedQueryShape{resources: len(resourceIDs), metrics: len(metricTypes), aggregate: stepSecs > 1, groupSeries: groupSeries}
|
||||
cacheable := len(tiers) <= len(shape.tiers)
|
||||
copy(shape.tiers[:], tiers)
|
||||
if cacheable {
|
||||
s.queryMu.Lock()
|
||||
query := s.readQueryShapes[shape]
|
||||
s.queryMu.Unlock()
|
||||
if query != "" {
|
||||
return query, retainedQueryParameters(resourceType, resourceIDs, metricTypes, start, end, stepSecs, tiers)
|
||||
}
|
||||
}
|
||||
query, params := retainedQuerySQL(resourceType, resourceIDs, metricTypes, start, end, stepSecs, tiers, groupSeries)
|
||||
if cacheable {
|
||||
s.queryMu.Lock()
|
||||
if len(s.readQueryShapes) < maxRetainedReadStatements {
|
||||
if s.readQueryShapes == nil {
|
||||
s.readQueryShapes = make(map[retainedQueryShape]string)
|
||||
}
|
||||
s.readQueryShapes[shape] = query
|
||||
}
|
||||
s.queryMu.Unlock()
|
||||
}
|
||||
return query, params
|
||||
}
|
||||
|
||||
// retainedQuerySQL reconciles overlapping storage buckets before any display
|
||||
// aggregation. The preferred tier owns its bucket, lower tiers fill uncovered
|
||||
// buckets, and a coarser fallback is omitted if a preferred point overlaps it.
|
||||
// A bucket's presence does not establish continuous underlying collection.
|
||||
// All probes are scoped to the same series and requested timestamp window.
|
||||
// Presence never establishes continuous collection or per-series coverage.
|
||||
// Every probe is restricted to the requested identities and timestamp window.
|
||||
func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, start, end time.Time, stepSecs int64, tiers []Tier, groupSeries bool) (string, []interface{}) {
|
||||
// Query dimensions fixed by the caller need not be decoded per point.
|
||||
params := retainedQueryParameters(resourceType, resourceIDs, metricTypes, start, end, stepSecs, tiers)
|
||||
identityColumns := ""
|
||||
if len(resourceIDs) != 1 {
|
||||
identityColumns += "resource_id, "
|
||||
|
|
@ -1352,59 +1412,56 @@ func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, st
|
|||
if len(metricTypes) != 1 {
|
||||
identityColumns += "metric_type, "
|
||||
}
|
||||
idSlots := strings.TrimSuffix(strings.Repeat("?,", len(resourceIDs)), ",")
|
||||
metricClause := ""
|
||||
if len(metricTypes) > 0 {
|
||||
metricClause = " AND m.metric_type IN (" + strings.TrimSuffix(strings.Repeat("?,", len(metricTypes)), ",") + ")"
|
||||
// Reuse SQLite's numbered bindings across every branch and overlap probe.
|
||||
// Each identity, timestamp and display step is bound only once per read.
|
||||
slots := func(first, count int) string {
|
||||
values := make([]string, count)
|
||||
for i := range values {
|
||||
values[i] = fmt.Sprintf("?%d", first+i)
|
||||
}
|
||||
return strings.Join(values, ",")
|
||||
}
|
||||
idSlots := slots(2, len(resourceIDs))
|
||||
metricSlots := slots(2+len(resourceIDs), len(metricTypes))
|
||||
startParam := 2 + len(resourceIDs) + len(metricTypes)
|
||||
endParam := startParam + 1
|
||||
tierParam := endParam + 1
|
||||
stepParam := tierParam + len(tiers)
|
||||
index := "idx_metrics_query_all"
|
||||
if len(metricTypes) > 0 {
|
||||
index = "idx_metrics_lookup"
|
||||
}
|
||||
// Keep tier and timestamp constraints in the index even when an ordering
|
||||
// index could avoid a sort by walking unrelated retention tiers.
|
||||
scope := `m.resource_type = ? AND m.resource_id IN (` + idSlots + `)` + metricClause + `
|
||||
AND m.tier = ? AND m.timestamp >= ? AND m.timestamp <= ?`
|
||||
branches := make([]string, 0, len(tiers))
|
||||
var params []interface{}
|
||||
appendScopeParams := func(tier Tier) {
|
||||
params = append(params, resourceType)
|
||||
for _, id := range resourceIDs {
|
||||
params = append(params, id)
|
||||
scope := func(alias string, tierIndex int) string {
|
||||
clause := alias + ".resource_type = ?1 AND " + alias + ".resource_id IN (" + idSlots + ")"
|
||||
if len(metricTypes) > 0 {
|
||||
clause += " AND " + alias + ".metric_type IN (" + metricSlots + ")"
|
||||
}
|
||||
for _, metric := range metricTypes {
|
||||
params = append(params, metric)
|
||||
}
|
||||
params = append(params, string(tier), start.Unix(), end.Unix())
|
||||
return clause + fmt.Sprintf(" AND %s.tier = ?%d AND %s.timestamp >= ?%d AND %s.timestamp <= ?%d", alias, tierParam+tierIndex, alias, startParam, alias, endParam)
|
||||
}
|
||||
projection := identityColumns + `m.timestamp, m.value,
|
||||
COALESCE(m.min_value, m.value) AS min_value, COALESCE(m.max_value, m.value) AS max_value`
|
||||
directAggregate := stepSecs > 1 && len(tiers) == 1
|
||||
bucketExpression := fmt.Sprintf("(timestamp / ?%d) * ?%d + (?%d / 2)", stepParam, stepParam, stepParam)
|
||||
if directAggregate {
|
||||
// A snapshot containing one tier needs no intermediate projection.
|
||||
// Aggregate its values directly, avoiding expression materialization
|
||||
// for every input row while keeping the same bounded output.
|
||||
projection = identityColumns + `(m.timestamp / ?) * ? + (? / 2) AS bucket_ts,
|
||||
projection = identityColumns + bucketExpression + ` AS bucket_ts,
|
||||
AVG(m.value), MIN(COALESCE(m.min_value, m.value)), MAX(COALESCE(m.max_value, m.value))`
|
||||
params = append(params, stepSecs, stepSecs, stepSecs)
|
||||
}
|
||||
branches := make([]string, 0, len(tiers))
|
||||
for i, tier := range tiers {
|
||||
branch := `SELECT ` + projection + `
|
||||
FROM metrics AS m INDEXED BY ` + index + ` WHERE ` + scope
|
||||
appendScopeParams(tier)
|
||||
for _, preferred := range tiers[:i] {
|
||||
branch += ` AND NOT EXISTS (`
|
||||
// Retention intervals nest on UTC minute/hour/day boundaries. Using the
|
||||
// larger interval makes both finer and coarser overlap probes indexable.
|
||||
branch := "SELECT " + projection + " FROM metrics AS m INDEXED BY " + index + " WHERE " + scope("m", i)
|
||||
for j, preferred := range tiers[:i] {
|
||||
// SQLite evaluates the uncorrelated existence check once. An empty
|
||||
// preferred tier must not incur a correlated index probe for every
|
||||
// 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))
|
||||
branch += fmt.Sprintf(`
|
||||
SELECT 1 FROM metrics AS h
|
||||
WHERE h.resource_type = m.resource_type AND h.resource_id = m.resource_id
|
||||
AND h.metric_type = m.metric_type AND h.tier = ?
|
||||
AND h.timestamp >= MAX(?, (m.timestamp / %d) * %d)
|
||||
AND h.timestamp <= ? AND h.timestamp < (m.timestamp / %d) * %d + %d
|
||||
)`, bucket, bucket, bucket, bucket, bucket)
|
||||
params = append(params, string(preferred), start.Unix(), end.Unix())
|
||||
AND h.metric_type = m.metric_type AND h.tier = ?%d
|
||||
AND h.timestamp >= MAX(?%d, (m.timestamp / %d) * %d)
|
||||
AND h.timestamp <= ?%d AND h.timestamp < (m.timestamp / %d) * %d + %d
|
||||
))`, tierParam+j, startParam, bucket, bucket, endParam, bucket, bucket, bucket)
|
||||
}
|
||||
branches = append(branches, branch)
|
||||
}
|
||||
|
|
@ -1412,22 +1469,12 @@ func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, st
|
|||
if directAggregate {
|
||||
query += " GROUP BY " + identityColumns + "bucket_ts ORDER BY " + identityColumns + "bucket_ts ASC"
|
||||
} else if stepSecs > 1 {
|
||||
// Reconcile first, then aggregate inside SQLite so a bounded chart does
|
||||
// not allocate and scan every retained observation in Go.
|
||||
query = `SELECT ` + identityColumns + `
|
||||
(timestamp / ?) * ? + (? / 2) AS bucket_ts,
|
||||
AVG(value), MIN(min_value), MAX(max_value)
|
||||
FROM (` + query + `)
|
||||
GROUP BY ` + identityColumns + `bucket_ts
|
||||
ORDER BY ` + identityColumns + `bucket_ts ASC`
|
||||
params = append([]interface{}{stepSecs, stepSecs, stepSecs}, params...)
|
||||
query = "SELECT " + identityColumns + bucketExpression + ` AS bucket_ts,
|
||||
AVG(value), MIN(min_value), MAX(max_value)
|
||||
FROM (` + query + ") GROUP BY " + identityColumns + "bucket_ts ORDER BY " + identityColumns + "bucket_ts ASC"
|
||||
} else {
|
||||
orderColumns := identityColumns
|
||||
if !groupSeries && len(metricTypes) == 0 {
|
||||
// The all-metric index orders each resource by time. Interleaved
|
||||
// metrics still append in timestamp order within each output series,
|
||||
// so plain reads need no extra sort by metric. Streaming display
|
||||
// aggregation explicitly requests contiguous series instead.
|
||||
orderColumns = ""
|
||||
if len(resourceIDs) != 1 {
|
||||
orderColumns = "resource_id, "
|
||||
|
|
@ -1553,7 +1600,7 @@ func (s *Store) queryRetainedChunk(resourceType string, resourceIDs []string, me
|
|||
if streamBuckets {
|
||||
queryStep = 0
|
||||
}
|
||||
sqlQuery, params := retainedQuerySQL(resourceType, resourceIDs, metricTypes, start, end, queryStep, tiers, streamBuckets)
|
||||
sqlQuery, params := s.retainedQuerySQL(resourceType, resourceIDs, metricTypes, start, end, queryStep, tiers, streamBuckets)
|
||||
|
||||
queryRows := func() (*sql.Rows, error) { return tx.Query(sqlQuery, params...) }
|
||||
if tx == nil {
|
||||
|
|
@ -1589,14 +1636,35 @@ func (s *Store) queryRetainedChunk(resourceType string, resourceIDs []string, me
|
|||
result := make(map[string]map[string][]MetricPoint, len(resourceIDs))
|
||||
seriesCapacity := estimateQueryAllBatchSeriesCapacity(start, end, stepSecs)
|
||||
|
||||
// Consecutive points in one series append directly to its slice. Flush on
|
||||
// a series change so interleaved metrics can resume their existing slices
|
||||
// without performing repeated nested-map lookups for every observation.
|
||||
var outputResource, outputMetric string
|
||||
var outputMetrics map[string][]MetricPoint
|
||||
var outputPoints []MetricPoint
|
||||
flushSeries := func() {
|
||||
if outputMetrics != nil {
|
||||
outputMetrics[outputMetric] = outputPoints
|
||||
}
|
||||
}
|
||||
appendPoint := func(resourceID, metricType string, point MetricPoint) {
|
||||
if result[resourceID] == nil {
|
||||
result[resourceID] = make(map[string][]MetricPoint, 8)
|
||||
if outputMetrics == nil || outputResource != resourceID || outputMetric != metricType {
|
||||
flushSeries()
|
||||
if outputMetrics == nil || outputResource != resourceID {
|
||||
outputResource = resourceID
|
||||
outputMetrics = result[resourceID]
|
||||
if outputMetrics == nil {
|
||||
outputMetrics = make(map[string][]MetricPoint, 8)
|
||||
result[resourceID] = outputMetrics
|
||||
}
|
||||
}
|
||||
outputMetric = metricType
|
||||
outputPoints = outputMetrics[metricType]
|
||||
if outputPoints == nil {
|
||||
outputPoints = make([]MetricPoint, 0, seriesCapacity)
|
||||
}
|
||||
}
|
||||
if result[resourceID][metricType] == nil {
|
||||
result[resourceID][metricType] = make([]MetricPoint, 0, seriesCapacity)
|
||||
}
|
||||
result[resourceID][metricType] = append(result[resourceID][metricType], point)
|
||||
outputPoints = append(outputPoints, point)
|
||||
}
|
||||
var bucketResource, bucketMetric string
|
||||
var bucketStart int64
|
||||
|
|
@ -1656,6 +1724,7 @@ func (s *Store) queryRetainedChunk(resourceType string, resourceIDs []string, me
|
|||
}
|
||||
|
||||
flushBucket()
|
||||
flushSeries()
|
||||
return result, nil
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ import (
|
|||
"testing"
|
||||
"time"
|
||||
|
||||
pdb "github.com/rcourtman/pulse-go-rewrite/pkg/db"
|
||||
"github.com/rs/zerolog"
|
||||
"github.com/rs/zerolog/log"
|
||||
)
|
||||
|
|
@ -1250,3 +1251,53 @@ func TestCommercialHistoryRetentionNeverExpandsOperatorPolicy(t *testing.T) {
|
|||
t.Fatalf("commercial ceiling expanded shorter operator retention: %v", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Exercise shared query reuse under race detection without latency assertions.
|
||||
// A one-connection pool also exposes preparation/read-transaction lock inversion.
|
||||
func TestStoreRetainedConcurrentQueryBindings(t *testing.T) {
|
||||
db := newPlanTestDB(t)
|
||||
db.SetMaxOpenConns(1)
|
||||
store := &Store{db: pdb.Wrap(db, "concurrent-retained-bindings")}
|
||||
base := time.Date(2026, 9, 6, 0, 0, 0, 0, time.UTC)
|
||||
type scope struct {
|
||||
family, id, metric string
|
||||
value float64
|
||||
}
|
||||
var scopes []scope
|
||||
for i := 0; i < 8; i++ {
|
||||
item := scope{family: "node", id: strconv.Itoa(i / 2), metric: "cpu", value: float64(i + 1)}
|
||||
if i%2 != 0 {
|
||||
item.family = "vm"
|
||||
}
|
||||
if i%4 >= 2 {
|
||||
item.metric = "memory"
|
||||
}
|
||||
if _, err := db.Exec(`INSERT INTO metrics(resource_type,resource_id,metric_type,tier,timestamp,value) VALUES (?,?,?,'raw',?,?)`, item.family, item.id, item.metric, base.Add(70*time.Second).Unix(), item.value); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
scopes = append(scopes, item)
|
||||
}
|
||||
start := make(chan struct{})
|
||||
var workers sync.WaitGroup
|
||||
for _, item := range scopes {
|
||||
workers.Add(1)
|
||||
go func(item scope) {
|
||||
defer workers.Done()
|
||||
<-start
|
||||
for i := 0; i < 12; i++ {
|
||||
step := []int64{0, 60, 120}[i%3]
|
||||
points, err := store.Query(item.family, item.id, item.metric, base.Add(time.Duration(i)*time.Second), base.Add(3*time.Minute), step)
|
||||
if err != nil {
|
||||
t.Errorf("scope %+v step %d: %v", item, step, err)
|
||||
return
|
||||
}
|
||||
if len(points) != 1 || points[0].Value != item.value {
|
||||
t.Errorf("scope %+v step %d received another query's values: %+v", item, step, points)
|
||||
return
|
||||
}
|
||||
}
|
||||
}(item)
|
||||
}
|
||||
close(start)
|
||||
workers.Wait()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -179,7 +179,7 @@ func TestRetainedReadStatementsReadCurrentSnapshot(t *testing.T) {
|
|||
db.SetMaxOpenConns(1)
|
||||
store := &Store{db: pdb.Wrap(db, "presence-snapshot")}
|
||||
end := time.Unix(2000000040, 0)
|
||||
start := end.Add(-time.Hour)
|
||||
start := end.Add(-3 * time.Hour)
|
||||
query := func(id string) []MetricPoint {
|
||||
t.Helper()
|
||||
points, err := store.Query("node", id, "cpu", start, end, step)
|
||||
|
|
@ -191,6 +191,15 @@ func TestRetainedReadStatementsReadCurrentSnapshot(t *testing.T) {
|
|||
if points := query("a"); len(points) != 0 {
|
||||
t.Fatalf("empty inventory: %+v", points)
|
||||
}
|
||||
// Warm the fallback query while the preferred minute tier is absent.
|
||||
// A later preferred bucket must replace its overlapping raw point,
|
||||
// even when the same SQL template and statement are reused.
|
||||
if _, err := db.Exec(`INSERT INTO metrics(resource_type,resource_id,metric_type,tier,timestamp,value) VALUES ('node','a','cpu','raw',?,23)`, end.Add(-40*time.Second).Unix()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if points := query("a"); len(points) != 1 || points[0].Value != 23 {
|
||||
t.Fatalf("raw fallback was hidden: %+v", points)
|
||||
}
|
||||
if _, err := db.Exec(`INSERT INTO metrics(resource_type,resource_id,metric_type,tier,timestamp,value) VALUES ('node','a','cpu','minute',?,17)`, end.Add(-time.Minute).Unix()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
|
@ -221,6 +230,43 @@ func TestRetainedReadStatementsReadCurrentSnapshot(t *testing.T) {
|
|||
if len(store.readStatements) > maxRetainedReadStatements {
|
||||
t.Fatalf("unbounded compiled SQL retention: %d", len(store.readStatements))
|
||||
}
|
||||
if len(store.readQueryShapes) > maxRetainedReadStatements {
|
||||
t.Fatalf("unbounded SQL template retention: %d", len(store.readQueryShapes))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetainedQueryTemplatesBindCurrentWindowAndStep(t *testing.T) {
|
||||
db := newPlanTestDB(t)
|
||||
db.SetMaxOpenConns(1)
|
||||
store := &Store{db: pdb.Wrap(db, "query-template-bindings")}
|
||||
base := time.Date(2026, 9, 6, 0, 0, 0, 0, time.UTC)
|
||||
for _, row := range []struct {
|
||||
family, metric string
|
||||
offset time.Duration
|
||||
value float64
|
||||
}{{"node", "cpu", 10 * time.Second, 5}, {"node", "cpu", 70 * time.Second, 15}, {"vm", "cpu", 70 * time.Second, 91}, {"node", "memory", 70 * time.Second, 99}} {
|
||||
if _, err := db.Exec(`INSERT INTO metrics(resource_type,resource_id,metric_type,tier,timestamp,value) VALUES (?,'a',?,'raw',?,?)`, row.family, row.metric, base.Add(row.offset).Unix(), row.value); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
for _, tc := range []struct {
|
||||
family, metric string
|
||||
start time.Duration
|
||||
step int64
|
||||
values []float64
|
||||
}{{"node", "cpu", 0, 60, []float64{5, 15}}, {"node", "cpu", 0, 120, []float64{10}}, {"node", "cpu", time.Minute, 120, []float64{15}}, {"vm", "cpu", 0, 120, []float64{91}}, {"node", "memory", 0, 120, []float64{99}}} {
|
||||
points, err := store.Query(tc.family, "a", tc.metric, base.Add(tc.start), base.Add(3*time.Minute), tc.step)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var values []float64
|
||||
for _, point := range points {
|
||||
values = append(values, point.Value)
|
||||
}
|
||||
if !reflect.DeepEqual(values, tc.values) {
|
||||
t.Fatalf("%+v returned stale bindings: %+v", tc, points)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -225,14 +225,13 @@ test.describe("TrueNAS alert resource links", () => {
|
|||
await expect(
|
||||
page.getByRole("heading", { name: "Resource incidents" }),
|
||||
).toBeVisible();
|
||||
const incidentsPanel = page
|
||||
.locator("div")
|
||||
.filter({
|
||||
has: page.getByRole("heading", { name: "Resource incidents" }),
|
||||
})
|
||||
.last();
|
||||
await expect(page.getByText("TrueNAS Main").first()).toBeVisible();
|
||||
await expect(page.getByText("· 1 incident")).toBeVisible();
|
||||
// Mobile keeps the underlying history mounted behind the investigation
|
||||
// sheet, so a page-wide count can match both surfaces.
|
||||
const investigation = isMobile
|
||||
? page.getByTestId("mobile-alert-investigation-scroll")
|
||||
: page.locator("#main");
|
||||
await expect(investigation.getByText("TrueNAS Main").first()).toBeVisible();
|
||||
await expect(investigation.getByText("· 1 incident", { exact: true })).toBeVisible();
|
||||
|
||||
await page.screenshot({ path: SCREENSHOT_PATH, fullPage: true });
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue