Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 32 additions & 10 deletions internal/analytics/postgres_analytics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
`
Expand Down Expand Up @@ -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)
Expand Down
16 changes: 8 additions & 8 deletions internal/analytics/postgres_analytics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand All @@ -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"))

Expand Down Expand Up @@ -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)

Expand All @@ -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"))

Expand Down Expand Up @@ -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)

Expand All @@ -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)

Expand Down Expand Up @@ -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)

Expand All @@ -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)

Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Original file line number Diff line number Diff line change
@@ -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;
Original file line number Diff line number Diff line change
@@ -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)
})
}
Loading