diff --git a/internal/analytics/postgres_analytics.go b/internal/analytics/postgres_analytics.go index 4b9968c6a..6af6cbbc8 100644 --- a/internal/analytics/postgres_analytics.go +++ b/internal/analytics/postgres_analytics.go @@ -345,13 +345,26 @@ func (s *PostgresAnalyticsStore) QueryMonthlyTotals(ctx context.Context, account func (s *PostgresAnalyticsStore) QueryByProvider(ctx context.Context, accountUUIDs []string, accountExternalIDsByProvider map[string][]string, startDate, endDate time.Time) ([]ProviderBreakdown, error) { accountClause, args := accountFilterClause(accountUUIDs, accountExternalIDsByProvider, []any{startDate, endDate}) + // Snapshot rows are run-rates written at (account, provider, service, + // region, commitment_type, timestamp) grain, so a (provider, service) + // bucket spans several rows at the SAME timestamp. Per H5 the inner query + // sums those rows into the bucket's instant total at each collection + // timestamp and the outer query averages the instant totals over time; a + // flat AVG would report the mean per-row run-rate instead of the bucket + // total (COR-02). // #nosec G201 — accountClause uses only internally-built placeholders. query := ` - SELECT provider, service, AVG(total_savings) as total_savings, AVG(coverage_percentage) as avg_coverage - FROM savings_snapshots - WHERE timestamp >= $1 - AND timestamp <= $2 - AND ` + accountClause + ` + SELECT provider, service, AVG(ts_savings) as total_savings, AVG(ts_coverage) as avg_coverage + FROM ( + SELECT provider, service, timestamp, + SUM(total_savings) as ts_savings, + AVG(coverage_percentage) as ts_coverage + FROM savings_snapshots + WHERE timestamp >= $1 + AND timestamp <= $2 + AND ` + accountClause + ` + GROUP BY provider, service, timestamp + ) per_ts GROUP BY provider, service ORDER BY total_savings DESC ` @@ -392,13 +405,22 @@ func (s *PostgresAnalyticsStore) QueryByService(ctx context.Context, accountUUID providerClause = fmt.Sprintf(" AND provider = $%d", len(args)) } + // Same H5 nested rollup as QueryByProvider: a (service, region) bucket + // spans accounts and commitment types at the same timestamp, so SUM the + // rows per timestamp first, then AVG the instant totals over time (COR-02). // #nosec G201 — accountClause / providerClause use only internally-built placeholders. query := fmt.Sprintf(` - SELECT service, region, AVG(total_savings) as total_savings, AVG(coverage_percentage) as avg_coverage - FROM savings_snapshots - WHERE timestamp >= $1 - AND timestamp <= $2 - AND %s%s + SELECT service, region, AVG(ts_savings) as total_savings, AVG(ts_coverage) as avg_coverage + FROM ( + SELECT service, region, timestamp, + SUM(total_savings) as ts_savings, + AVG(coverage_percentage) as ts_coverage + FROM savings_snapshots + WHERE timestamp >= $1 + AND timestamp <= $2 + AND %s%s + GROUP BY service, region, timestamp + ) per_ts GROUP BY service, region ORDER BY total_savings DESC `, accountClause, providerClause) diff --git a/internal/analytics/postgres_analytics_test.go b/internal/analytics/postgres_analytics_test.go index e2f651fdc..f55daf3a3 100644 --- a/internal/analytics/postgres_analytics_test.go +++ b/internal/analytics/postgres_analytics_test.go @@ -580,7 +580,7 @@ func TestQueryByProvider(t *testing.T) { AddRow("aws", "rds", 2500.0, f64ptr(85.0)). AddRow("aws", "elasticache", 1200.0, f64ptr(75.0)) - mock.ExpectQuery(`SELECT provider, service, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT provider, service, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY provider, service, timestamp`). WithArgs(startDate, now, []string{"account-123"}). WillReturnRows(rows) @@ -603,7 +603,7 @@ func TestQueryByProvider(t *testing.T) { now := time.Now().UTC() startDate := now.Add(-30 * 24 * time.Hour) - mock.ExpectQuery(`SELECT provider, service, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT provider, service, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY provider, service, timestamp`). WithArgs(startDate, now, []string{"account-123"}). WillReturnError(errors.New("database error")) @@ -632,7 +632,7 @@ func TestQueryByService(t *testing.T) { AddRow("rds", "us-east-1", 1800.0, f64ptr(90.0)). AddRow("rds", "us-west-2", 700.0, f64ptr(75.0)) - mock.ExpectQuery(`SELECT service, region, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT service, region, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY service, region, timestamp`). WithArgs(startDate, now, []string{"account-123"}, "aws"). WillReturnRows(rows) @@ -655,7 +655,7 @@ func TestQueryByService(t *testing.T) { now := time.Now().UTC() startDate := now.Add(-30 * 24 * time.Hour) - mock.ExpectQuery(`SELECT service, region, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT service, region, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY service, region, timestamp`). WithArgs(startDate, now, []string{"account-123"}, "aws"). WillReturnError(errors.New("database error")) @@ -1141,7 +1141,7 @@ func TestQueryByProviderRowScanError(t *testing.T) { "provider", // Missing other columns }).AddRow("aws").RowError(0, errors.New("scan error")) - mock.ExpectQuery(`SELECT provider, service, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT provider, service, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY provider, service, timestamp`). WithArgs(startDate, now, []string{"account-123"}). WillReturnRows(rows) @@ -1167,7 +1167,7 @@ func TestQueryByServiceRowScanError(t *testing.T) { "service", // Missing other columns }).AddRow("rds").RowError(0, errors.New("scan error")) - mock.ExpectQuery(`SELECT service, region, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT service, region, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY service, region, timestamp`). WithArgs(startDate, now, []string{"account-123"}, "aws"). WillReturnRows(rows) @@ -1258,7 +1258,7 @@ func TestQueryByProviderRowsErr(t *testing.T) { "provider", "service", "total_savings", "avg_coverage", }).AddRow("aws", "rds", 2500.0, f64ptr(85.0)) - mock.ExpectQuery(`SELECT provider, service, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT provider, service, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY provider, service, timestamp`). WithArgs(startDate, now, []string{"account-123"}). WillReturnRows(rows) @@ -1285,7 +1285,7 @@ func TestQueryByServiceRowsErr(t *testing.T) { "service", "region", "total_savings", "avg_coverage", }).AddRow("rds", "us-east-1", 1800.0, f64ptr(90.0)) - mock.ExpectQuery(`SELECT service, region, AVG\(total_savings\) as total_savings`). + mock.ExpectQuery(`(?s)SELECT service, region, AVG\(ts_savings\) as total_savings.*SUM\(total_savings\) as ts_savings.*GROUP BY service, region, timestamp`). WithArgs(startDate, now, []string{"account-123"}, "aws"). WillReturnRows(rows) diff --git a/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup.down.sql b/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup.down.sql new file mode 100644 index 000000000..26724229c --- /dev/null +++ b/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup.down.sql @@ -0,0 +1,24 @@ +-- 000078 down: restore the flat-AVG monthly_savings_summary definition from +-- 000067/000074 (which understates multi-row buckets, per COR-02), and +-- recreate the unique index using NULLS NOT DISTINCT (as 000076 left it). +DROP MATERIALIZED VIEW IF EXISTS monthly_savings_summary CASCADE; + +CREATE MATERIALIZED VIEW monthly_savings_summary AS +SELECT + DATE_TRUNC('month', timestamp) as month, + account_id, + cloud_account_id, + provider, + service, + AVG(total_savings) as total_savings, + AVG(coverage_percentage) as avg_coverage, + AVG(total_commitment) as total_commitment, + AVG(total_usage) as total_usage, + COUNT(*) as snapshot_count, + MAX(timestamp) as last_updated +FROM savings_snapshots +GROUP BY DATE_TRUNC('month', timestamp), account_id, cloud_account_id, provider, service; + +CREATE UNIQUE INDEX idx_monthly_savings_summary_unique + ON monthly_savings_summary (month, account_id, cloud_account_id, provider, service) + NULLS NOT DISTINCT; diff --git a/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup.up.sql b/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup.up.sql new file mode 100644 index 000000000..0f92eb44d --- /dev/null +++ b/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup.up.sql @@ -0,0 +1,62 @@ +-- 000078: monthly_savings_summary nested SUM-then-AVG rollup (COR-02). +-- +-- Snapshot rows are run-rates written at (account, provider, service, region, +-- commitment_type, timestamp) grain, so a (month, account, provider, service) +-- bucket contains several rows sharing the SAME timestamp (one per region / +-- commitment type). The 000067/000074 definition used a flat AVG across all +-- rows in the bucket, which returns the mean per-row run-rate instead of the +-- bucket's total: $100/mo in each of 5 regions reported total_savings=$100 +-- instead of $500. +-- +-- Apply the same H5 shape daily_savings_trend and provider_savings_summary +-- already use: the inner query sums the per-row run-rates into the bucket's +-- instant total at each collection timestamp, the outer query averages those +-- instant totals over the month so the result stays invariant to the +-- collection frequency. snapshot_count keeps its raw-row semantics via +-- SUM(ts_row_count) (cast back to BIGINT because SUM(bigint) yields NUMERIC). +-- +-- DROP + CREATE because a materialized view's defining query cannot be +-- altered in place. The unique index is recreated for the CONCURRENTLY +-- refresh path, and the view is created populated (no WITH NO DATA) so +-- refresh_savings_materialized_views() can keep using CONCURRENTLY. +DROP MATERIALIZED VIEW IF EXISTS monthly_savings_summary CASCADE; + +CREATE MATERIALIZED VIEW monthly_savings_summary AS +SELECT + month, + account_id, + cloud_account_id, + provider, + service, + AVG(ts_savings) as total_savings, + AVG(ts_coverage) as avg_coverage, + AVG(ts_commitment) as total_commitment, + AVG(ts_usage) as total_usage, + SUM(ts_row_count)::BIGINT as snapshot_count, + MAX(timestamp) as last_updated +FROM ( + SELECT + DATE_TRUNC('month', timestamp) as month, + timestamp, + account_id, + cloud_account_id, + provider, + service, + SUM(total_savings) as ts_savings, + AVG(coverage_percentage) as ts_coverage, + SUM(total_commitment) as ts_commitment, + SUM(total_usage) as ts_usage, + COUNT(*) as ts_row_count + FROM savings_snapshots + GROUP BY DATE_TRUNC('month', timestamp), timestamp, + account_id, cloud_account_id, provider, service +) per_ts +GROUP BY month, account_id, cloud_account_id, provider, service; + +-- Use plain columns with NULLS NOT DISTINCT (PostgreSQL 15+), consistent +-- with 000076 which switched from COALESCE to NULLS NOT DISTINCT so that +-- REFRESH MATERIALIZED VIEW CONCURRENTLY works (it requires plain-column +-- indexes with no expressions). +CREATE UNIQUE INDEX idx_monthly_savings_summary_unique + ON monthly_savings_summary (month, account_id, cloud_account_id, provider, service) + NULLS NOT DISTINCT; diff --git a/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup_test.go b/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup_test.go new file mode 100644 index 000000000..9efb42fa3 --- /dev/null +++ b/internal/database/postgres/migrations/000078_monthly_summary_nested_rollup_test.go @@ -0,0 +1,153 @@ +//go:build integration +// +build integration + +package migrations_test + +import ( + "context" + "testing" + "time" + + "github.com/LeanerCloud/CUDly/internal/analytics" + "github.com/LeanerCloud/CUDly/internal/database/postgres/migrations" + "github.com/LeanerCloud/CUDly/internal/database/postgres/testhelpers" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TestAnalyticsNestedRollup_COR02 replicates the COR-02 failing scenario: +// snapshot rows are run-rates written at (account, provider, service, region, +// commitment_type, timestamp) grain, so a breakdown bucket contains several +// rows sharing the SAME timestamp. The pre-fix flat AVG reported the mean +// per-row run-rate instead of the bucket's per-timestamp total, understating +// savings for any multi-region / multi-commitment-type estate. +// +// Fixture: provider aws, service rds, two collection timestamps T1 and T2 in +// the current month. +// +// T1: us-east-1 RI $100 + us-east-1 SavingsPlan $50 + us-west-2 RI $100 = $250 +// T2: us-east-1 RI $200 + us-east-1 SavingsPlan $100 + us-west-2 RI $100 = $400 +// +// Correct rollup (inner SUM per timestamp, outer AVG over timestamps, the H5 +// shape): provider/service total = AVG(250, 400) = 325. The pre-fix flat AVG +// returned 650/6 = 108.33. +func TestAnalyticsNestedRollup_COR02(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute) + defer cancel() + + container, err := testhelpers.SetupPostgresContainer(ctx, t) + if err != nil { + t.Skipf("Skipping test: could not setup postgres container: %v", err) + } + defer container.Cleanup(ctx) + + require.NoError(t, migrations.RunMigrations(ctx, container.DB.Pool(), getMigrationsPath(), "", "")) + + store := analytics.NewPostgresAnalyticsStore(container.DB) + + // Two timestamps inside the current month so the current-month partition + // (created by the migrations) holds the rows and QueryMonthlyTotals' + // current-month window matches them. + now := time.Now().UTC() + monthStart := time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, time.UTC) + t1 := monthStart.Add(12 * time.Hour) + t2 := monthStart.Add(13 * time.Hour) + + const account = "123456789012" + mkSnapshot := func(ts time.Time, region, commitmentType string, savings float64) *analytics.SavingsSnapshot { + return &analytics.SavingsSnapshot{ + AccountID: account, + Timestamp: ts, + Provider: "aws", + Service: "rds", + Region: region, + CommitmentType: commitmentType, + TotalCommitment: savings * 10, + TotalSavings: savings, + } + } + fixture := []*analytics.SavingsSnapshot{ + mkSnapshot(t1, "us-east-1", "RI", 100), + mkSnapshot(t1, "us-east-1", "SavingsPlan", 50), + mkSnapshot(t1, "us-west-2", "RI", 100), + mkSnapshot(t2, "us-east-1", "RI", 200), + mkSnapshot(t2, "us-east-1", "SavingsPlan", 100), + mkSnapshot(t2, "us-west-2", "RI", 100), + } + for _, snapshot := range fixture { + require.NoError(t, store.SaveSnapshot(ctx, snapshot)) + } + + accountFilter := map[string][]string{"aws": {account}} + start := monthStart + end := monthStart.Add(24 * time.Hour) + + t.Run("QueryByProvider sums rows per timestamp before averaging", func(t *testing.T) { + breakdowns, err := store.QueryByProvider(ctx, nil, accountFilter, start, end) + require.NoError(t, err) + require.Len(t, breakdowns, 1) + assert.Equal(t, "aws", breakdowns[0].Provider) + assert.Equal(t, "rds", breakdowns[0].Service) + // AVG(T1 total $250, T2 total $400) = $325; the pre-fix flat AVG + // across the 6 rows returned $108.33. + assert.InDelta(t, 325.0, breakdowns[0].TotalSavings, 0.01) + }) + + t.Run("QueryByService sums commitment types per timestamp before averaging", func(t *testing.T) { + breakdowns, err := store.QueryByService(ctx, nil, accountFilter, "aws", start, end) + require.NoError(t, err) + require.Len(t, breakdowns, 2) + + byRegion := map[string]float64{} + for _, b := range breakdowns { + assert.Equal(t, "rds", b.Service) + byRegion[b.Region] = b.TotalSavings + } + // us-east-1: AVG(T1 100+50, T2 200+100) = $225; pre-fix flat AVG + // across the 4 rows returned $112.50. + assert.InDelta(t, 225.0, byRegion["us-east-1"], 0.01) + // us-west-2 has a single row per timestamp, so both shapes agree. + assert.InDelta(t, 100.0, byRegion["us-west-2"], 0.01) + }) + + t.Run("monthly_savings_summary sums rows per timestamp before averaging", func(t *testing.T) { + _, err := container.DB.Exec(ctx, "REFRESH MATERIALIZED VIEW monthly_savings_summary") + require.NoError(t, err) + + summaries, err := store.QueryMonthlyTotals(ctx, nil, accountFilter, 1) + require.NoError(t, err) + require.Len(t, summaries, 1) + assert.Equal(t, "aws", summaries[0].Provider) + assert.Equal(t, "rds", summaries[0].Service) + // AVG(T1 total $250, T2 total $400) = $325; the pre-fix flat-AVG view + // definition reported $108.33. + assert.InDelta(t, 325.0, summaries[0].TotalSavings, 0.01) + // snapshot_count keeps its raw-row semantics across the rewrite. + assert.Equal(t, 6, summaries[0].SnapshotCount) + }) + + t.Run("down restores the flat-AVG view and up reapplies cleanly", func(t *testing.T) { + // Migrate down to version 77 (just below this migration) so 000078's + // down runs regardless of how many later migrations sit above it on + // main; a fixed-step RollbackMigrations would only undo the topmost + // migration and leave the nested-rollup view in place. + require.NoError(t, migrations.MigrateToVersion(ctx, container.DB.Pool(), getMigrationsPath(), 77)) + _, err := container.DB.Exec(ctx, "REFRESH MATERIALIZED VIEW monthly_savings_summary") + require.NoError(t, err) + + summaries, err := store.QueryMonthlyTotals(ctx, nil, accountFilter, 1) + require.NoError(t, err) + require.Len(t, summaries, 1) + // Rolled back to the 000067/000074 flat AVG: 650/6 = 108.33. + assert.InDelta(t, 650.0/6.0, summaries[0].TotalSavings, 0.01) + + require.NoError(t, migrations.RunMigrations(ctx, container.DB.Pool(), getMigrationsPath(), "", "")) + _, err = container.DB.Exec(ctx, "REFRESH MATERIALIZED VIEW monthly_savings_summary") + require.NoError(t, err) + + summaries, err = store.QueryMonthlyTotals(ctx, nil, accountFilter, 1) + require.NoError(t, err) + require.Len(t, summaries, 1) + assert.InDelta(t, 325.0, summaries[0].TotalSavings, 0.01) + }) +}