From bcccb4ceae82bf730492bbf6eca14b2ca84d6216 Mon Sep 17 00:00:00 2001 From: Matthew Staebler Date: Tue, 7 Jul 2026 22:48:30 -0400 Subject: [PATCH] TRT-2737: Rewrite test analysis queries to use partitioned tables Replace the legacy test_analysis_by_job_by_dates table and prow_test_analysis_by_variant_14d_view with direct queries against the new partitioned tables: - test_daily_totals for per-day test analysis (overall, by-job, by-variant views) - test_cumulative_summaries for range aggregation (QueryTestAnalysis) Variant filtering joins through prow_jobs.variant_combination_id instead of the dropped test_daily_summaries.variant_combination_id column. The old table and its BQ loader remain in place for safe rollback; removal is planned for a follow-up PR. Co-Authored-By: Claude Opus 4.6 TRT-2737: Use civil.Date for date boundaries and add SQL comments Normalize report windows to calendar-day boundaries using civil.Date instead of time.Time, eliminating time-component ambiguity when comparing against DATE columns. Add inline comments to the cumulative prefix-sum query explaining the self-join pattern. Co-Authored-By: Claude Opus 4.6 --- pkg/api/job_runs.go | 2 +- pkg/api/test_analysis.go | 239 ++++++++++++------------ pkg/api/test_analysis_test.go | 65 +++++++ pkg/db/query/test_queries.go | 25 ++- pkg/db/views.go | 48 +---- pkg/flags/postgres_benchmarking_test.go | 7 +- 6 files changed, 210 insertions(+), 176 deletions(-) create mode 100644 pkg/api/test_analysis_test.go diff --git a/pkg/api/job_runs.go b/pkg/api/job_runs.go index 2364b97b8e..5c749c26e8 100644 --- a/pkg/api/job_runs.go +++ b/pkg/api/job_runs.go @@ -492,7 +492,7 @@ func jobNamesTestResultFunc(dbc *db.DB, release string) testResultsByJobNameFunc analyzeSince := time.Now().Add(-14 * 24 * time.Hour) - q := dbc.DB.Raw(query.QueryTestAnalysis, analyzeSince, testName, jobNames, release) + q := dbc.DB.Raw(query.QueryTestAnalysis, analyzeSince, release, release, testName, jobNames) if q.Error != nil { return nil, q.Error } diff --git a/pkg/api/test_analysis.go b/pkg/api/test_analysis.go index cd8ed72a01..0a106980ec 100644 --- a/pkg/api/test_analysis.go +++ b/pkg/api/test_analysis.go @@ -1,71 +1,104 @@ package api import ( + "strings" "time" + "cloud.google.com/go/civil" log "github.com/sirupsen/logrus" + "gorm.io/gorm" "github.com/openshift/sippy/pkg/db" "github.com/openshift/sippy/pkg/filter" ) +const testAnalysisLookbackDays = 14 + type CountByDate struct { - Date string `json:"date"` - Group string `json:"group"` - PassPercentage float64 `json:"pass_percentage"` - FlakePercentage float64 `json:"flake_percentage"` - FailPercentage float64 `json:"fail_percentage"` - Runs int `json:"runs"` - Passes int `json:"passes"` - Flakes int `json:"flakes"` - Failures int `json:"failures"` + Date civil.Date `json:"date"` + Group string `json:"group"` + PassPercentage float64 `json:"pass_percentage"` + FlakePercentage float64 `json:"flake_percentage"` + FailPercentage float64 `json:"fail_percentage"` + Runs int `json:"runs"` + Passes int `json:"passes"` + Flakes int `json:"flakes"` + Failures int `json:"failures"` } -func GetTestAnalysisOverallFromDB(dbc *db.DB, filters *filter.Filter, release, testName string, reportEnd time.Time) (map[string][]CountByDate, error) { - var rows []CountByDate - jq := dbc.DB.Table("test_analysis_by_job_by_dates"). - Select(`test_id, - test_name, - to_date((date at time zone 'UTC')::text, 'YYYY-MM-DD'::text)::text as date, - 'overall' as group, - SUM(runs) as runs, - SUM(passes) as passes, - SUM(flakes) as flakes, - SUM(failures) as failures, - SUM(passes) * 100.0 / NULLIF(SUM(runs), 0) AS pass_percentage, - SUM(flakes) * 100.0 / NULLIF(SUM(runs), 0) AS flake_percentage, - SUM(failures) * 100.0 / NULLIF(SUM(runs), 0) AS fail_percentage`). - Joins("JOIN prow_jobs on prow_jobs.name = job_name"). - Where("test_analysis_by_job_by_dates.release = ?", release). - Where("test_name = ?", testName). - Where("date >= ?", time.Now().Add(-24*14*time.Hour)). - Order("date ASC"). - Group("date, test_id, test_name, test_analysis_by_job_by_dates.release") - - var allowedVariants, blockedVariants []string - if filters != nil { - for _, f := range filters.Items { - if f.Field == "variants" { - if f.Not { - blockedVariants = append(blockedVariants, f.Value) - } else { - allowedVariants = append(allowedVariants, f.Value) - } +func extractVariantFilters(filters *filter.Filter) (allowed, blocked []string) { + if filters == nil { + return nil, nil + } + for _, f := range filters.Items { + if f.Field == "variants" { + if f.Not { + blocked = append(blocked, f.Value) + } else { + allowed = append(allowed, f.Value) } } } + return allowed, blocked +} - for _, bv := range blockedVariants { - jq = jq.Where("NOT EXISTS (SELECT 1 FROM variant_combinations WHERE ? = any(variants) AND id = prow_jobs.variant_combination_id)", bv) +func applyVariantCombinationFilters(q *gorm.DB, allowed, blocked []string) *gorm.DB { + for _, bv := range blocked { + q = q.Where("NOT EXISTS (SELECT 1 FROM variant_combinations WHERE ? = ANY(variants) AND id = pj.variant_combination_id)", bv) } - - for _, av := range allowedVariants { - jq = jq.Where("prow_jobs.variant_combination_id IN (SELECT id FROM variant_combinations WHERE ? = any(variants))", av) + for _, av := range allowed { + q = q.Where("EXISTS (SELECT 1 FROM variant_combinations WHERE ? = ANY(variants) AND id = pj.variant_combination_id)", av) } + return q +} + +var testAnalysisAggColumns = []string{ + "SUM(tds.runs) AS runs", + "SUM(tds.successes) AS passes", + "SUM(tds.flakes) AS flakes", + "SUM(tds.failures) AS failures", + "SUM(tds.successes) * 100.0 / NULLIF(SUM(tds.runs), 0) AS pass_percentage", + "SUM(tds.flakes) * 100.0 / NULLIF(SUM(tds.runs), 0) AS flake_percentage", + "SUM(tds.failures) * 100.0 / NULLIF(SUM(tds.runs), 0) AS fail_percentage", +} + +// withAggColumns appends the shared test analysis aggregation columns +// to the supplied per-query columns and joins them into a SELECT string. +func withAggColumns(columns ...string) string { + return strings.Join(append(columns, testAnalysisAggColumns...), ", ") +} + +func testAnalysisBaseQuery(dbc *db.DB, filters *filter.Filter, release, testName string, since civil.Date) *gorm.DB { + q := dbc.DB.Table("test_daily_totals tds"). + Joins("JOIN tests t ON t.id = tds.test_id"). + Joins("JOIN prow_jobs pj ON pj.id = tds.prow_job_id"). + Where("tds.release = ?", release). + Where("t.name = ?", testName). + Where("tds.date >= ?", since). + Order("tds.date ASC") + + allowed, blocked := extractVariantFilters(filters) + return applyVariantCombinationFilters(q, allowed, blocked) +} + +func GetTestAnalysisOverallFromDB(dbc *db.DB, filters *filter.Filter, release, testName string, reportEnd time.Time) (map[string][]CountByDate, error) { + endDate := civil.DateOf(reportEnd.UTC()) + sinceDate := endDate.AddDays(-testAnalysisLookbackDays) + + var rows []CountByDate + jq := testAnalysisBaseQuery(dbc, filters, release, testName, sinceDate). + Select(withAggColumns( + "tds.test_id", + `t.name AS test_name`, + `tds.date AS date`, + `'overall' AS "group"`, + )). + Where("tds.date <= ?", endDate). + Group("tds.date, tds.test_id, t.name, tds.release") r := jq.Scan(&rows) if r.Error != nil { - log.WithError(r.Error).Error("error querying test analysis by job") + log.WithError(r.Error).Error("error querying test analysis overall") return nil, r.Error } @@ -86,48 +119,20 @@ func GetTestAnalysisByJobFromDB(dbc *db.DB, filters *filter.Filter, release, tes results["overall"] = overall } - jq := dbc.DB.Table("test_analysis_by_job_by_dates"). - Select(`test_id, - test_name, - to_date((date at time zone 'UTC')::text, 'YYYY-MM-DD'::text)::text as date, - prow_jobs.release, - job_name as group, - runs, - passes, - flakes, - failures, - variants, - passes * 100.0 / NULLIF(runs, 0) AS pass_percentage, - flakes * 100.0 / NULLIF(runs, 0) AS flake_percentage, - failures * 100.0 / NULLIF(runs, 0) AS fail_percentage`). - Joins("INNER JOIN prow_jobs on prow_jobs.name = job_name"). - Where("prow_jobs.release = ?", release). - Where("test_analysis_by_job_by_dates.release = ?", release). - Where("test_name = ?", testName). - Where("date <= ?", reportEnd). - Where("date >= ?", reportEnd.Add(-24*14*time.Hour)). - Order("date ASC") - - var allowedVariants, blockedVariants []string - if filters != nil { - for _, f := range filters.Items { - if f.Field == "variants" { - if f.Not { - blockedVariants = append(blockedVariants, f.Value) - } else { - allowedVariants = append(allowedVariants, f.Value) - } - } - } - } + endDate := civil.DateOf(reportEnd.UTC()) + sinceDate := endDate.AddDays(-testAnalysisLookbackDays) - for _, bv := range blockedVariants { - jq = jq.Where("NOT EXISTS (SELECT 1 FROM variant_combinations WHERE ? = any(variants) AND id = prow_jobs.variant_combination_id)", bv) - } - - for _, av := range allowedVariants { - jq = jq.Where("prow_jobs.variant_combination_id IN (SELECT id FROM variant_combinations WHERE ? = any(variants))", av) - } + jq := testAnalysisBaseQuery(dbc, filters, release, testName, sinceDate). + Select(withAggColumns( + "tds.test_id", + `t.name AS test_name`, + `tds.date AS date`, + "pj.release", + `pj.name AS "group"`, + "pj.variants", + )). + Where("tds.date <= ?", endDate). + Group("tds.date, tds.test_id, t.name, pj.release, pj.name, pj.variants") r := jq.Scan(&rows) if r.Error != nil { @@ -154,40 +159,34 @@ func GetTestAnalysisByVariantFromDB(dbc *db.DB, filters *filter.Filter, release, results["overall"] = overall } - vq := dbc.DB.Table("prow_test_analysis_by_variant_14d_view"). - Where("release = ?", release). - Where("test_name = ?", testName). - Where("date <= ?", reportEnd). - Select(`to_date((date at time zone 'UTC')::text, 'YYYY-MM-DD'::text)::text as date, - variant as group, - runs, - passes, - flakes, - failures, - passes * 100.0 / NULLIF(runs, 0) AS pass_percentage, - flakes * 100.0 / NULLIF(runs, 0) AS flake_percentage, - failures * 100.0 / NULLIF(runs, 0) AS fail_percentage`). - Order("date ASC") - - var allowedVariants, blockedVariants []string - if filters != nil { - for _, f := range filters.Items { - if f.Field == "variants" { - if f.Not { - blockedVariants = append(blockedVariants, f.Value) - } else { - allowedVariants = append(allowedVariants, f.Value) - } - } - } - - if len(blockedVariants) > 0 { - vq = vq.Where("variant NOT IN ?", blockedVariants) - } - - if len(allowedVariants) > 0 { - vq = vq.Where("variant IN ?", allowedVariants) - } + endDate := civil.DateOf(reportEnd.UTC()) + sinceDate := endDate.AddDays(-testAnalysisLookbackDays) + + inner := dbc.DB.Table("test_daily_totals tds"). + Select(withAggColumns( + "tds.test_id", + `t.name AS test_name`, + `tds.date AS date`, + `unnest(vc.variants) AS "group"`, + "tds.release", + )). + Joins("JOIN tests t ON t.id = tds.test_id"). + Joins("JOIN prow_jobs pj ON pj.id = tds.prow_job_id"). + Joins("JOIN variant_combinations vc ON vc.id = pj.variant_combination_id"). + Where("tds.release = ?", release). + Where("t.name = ?", testName). + Where("tds.date >= ?", sinceDate). + Where("tds.date <= ?", endDate). + Group("t.name, t.id, tds.test_id, tds.date, unnest(vc.variants), tds.release"). + Order("tds.date ASC") + + vq := dbc.DB.Table("(?) AS analysis", inner).Select("*") + allowed, blocked := extractVariantFilters(filters) + if len(blocked) > 0 { + vq = vq.Where(`"group" NOT IN ?`, blocked) + } + if len(allowed) > 0 { + vq = vq.Where(`"group" IN ?`, allowed) } r := vq.Scan(&rows) diff --git a/pkg/api/test_analysis_test.go b/pkg/api/test_analysis_test.go new file mode 100644 index 0000000000..909d3bb396 --- /dev/null +++ b/pkg/api/test_analysis_test.go @@ -0,0 +1,65 @@ +package api + +import ( + "testing" + + "github.com/stretchr/testify/assert" + + "github.com/openshift/sippy/pkg/filter" +) + +func TestExtractVariantFilters(t *testing.T) { + tests := []struct { + name string + filters *filter.Filter + expectedAllowed []string + expectedBlocked []string + }{ + { + name: "nil filter", + filters: nil, + expectedAllowed: nil, + expectedBlocked: nil, + }, + { + name: "no variant filters", + filters: &filter.Filter{Items: []filter.FilterItem{{Field: "name", Value: "test1"}}}, + expectedAllowed: nil, + expectedBlocked: nil, + }, + { + name: "allowed variants", + filters: &filter.Filter{Items: []filter.FilterItem{ + {Field: "variants", Value: "Platform:aws"}, + {Field: "variants", Value: "Network:ovn"}, + }}, + expectedAllowed: []string{"Platform:aws", "Network:ovn"}, + expectedBlocked: nil, + }, + { + name: "blocked variants", + filters: &filter.Filter{Items: []filter.FilterItem{ + {Field: "variants", Value: "Platform:aws", Not: true}, + }}, + expectedAllowed: nil, + expectedBlocked: []string{"Platform:aws"}, + }, + { + name: "mixed allowed and blocked", + filters: &filter.Filter{Items: []filter.FilterItem{ + {Field: "variants", Value: "Platform:aws"}, + {Field: "variants", Value: "Topology:single", Not: true}, + {Field: "name", Value: "ignored"}, + }}, + expectedAllowed: []string{"Platform:aws"}, + expectedBlocked: []string{"Topology:single"}, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + allowed, blocked := extractVariantFilters(tt.filters) + assert.Equal(t, tt.expectedAllowed, allowed) + assert.Equal(t, tt.expectedBlocked, blocked) + }) + } +} diff --git a/pkg/db/query/test_queries.go b/pkg/db/query/test_queries.go index c6e7b01012..afd11f1deb 100644 --- a/pkg/db/query/test_queries.go +++ b/pkg/db/query/test_queries.go @@ -66,12 +66,27 @@ const ( GROUP BY bug_tests.test_id` QueryTestAnalysis = ` - select current_successes, current_runs, + SELECT current_successes, current_runs, current_successes * 100.0 / NULLIF(current_runs, 0) AS current_pass_percentage - from ( - select sum(runs) as current_runs, sum(passes) as current_successes - from test_analysis_by_job_by_dates - where date >= ? AND test_name = ? AND job_name IN ? AND release = ? + FROM ( + SELECT + SUM(e.prefix_sum_successes - COALESCE(s.prefix_sum_successes, 0)) AS current_successes, + SUM(e.prefix_sum_runs - COALESCE(s.prefix_sum_runs, 0)) AS current_runs + -- e = latest cumulative snapshot for the release + FROM test_cumulative_summaries e + JOIN tests t ON t.id = e.test_id + JOIN prow_jobs pj ON pj.id = e.prow_job_id + -- s = prior-day baseline; the difference e - s gives the windowed totals + LEFT JOIN test_cumulative_summaries s + ON s.test_id = e.test_id + AND s.prow_job_id = e.prow_job_id + AND s.suite_id = e.suite_id + AND s.release = e.release + AND s.date = (?::date - INTERVAL '1 day')::date + WHERE e.date = (SELECT MAX(date) FROM test_cumulative_summaries WHERE release = ?) + AND e.release = ? + AND t.name = ? + AND pj.name IN ? ) t` ) diff --git a/pkg/db/views.go b/pkg/db/views.go index 51093d1537..763b7eef39 100644 --- a/pkg/db/views.go +++ b/pkg/db/views.go @@ -37,12 +37,7 @@ var PostgresMatViews = []PostgresView{ } // PostgresViews are regular, non-materialized views: -var PostgresViews = []PostgresView{ - { - Name: "prow_test_analysis_by_variant_14d_view", - Definition: testAnalysisByVariantView, - }, -} +var PostgresViews = []PostgresView{} type PostgresView struct { // Name is the name of the materialized view in postgres. @@ -234,47 +229,6 @@ FROM prow_job_runs JOIN prow_jobs ON prow_job_runs.prow_job_id = prow_jobs.id WHERE prow_job_runs."timestamp" >= |||TIMENOW||| - interval '90 days' ` -const testAnalysisByVariantView = ` -SELECT - byjob.test_id AS test_id, - byjob.test_name AS test_name, - byjob.date AS date, - unnest(prow_jobs.variants) AS variant, - prow_jobs.release, - SUM(runs) as runs, - SUM(passes) as passes, - SUM(flakes) as flakes, - SUM(failures) as failures -FROM - test_analysis_by_job_by_dates byjob - JOIN tests ON tests.id = byjob.test_id - JOIN prow_jobs ON prow_jobs.name = byjob.job_name -WHERE - byjob.date >= (|||TIMENOW||| - '15 days'::interval) -GROUP BY - tests.name, tests.id, byjob.test_id, byjob.test_name, date, unnest(prow_jobs.variants), prow_jobs.release -` - -const testAnalysisByJobMatView = ` -SELECT - tests.id AS test_id, - tests.name AS test_name, - date(prow_job_run_tests.prow_job_run_timestamp) AS date, - prow_job_run_tests.prow_job_run_release AS release, - prow_jobs.name AS job_name, - COUNT(*) FILTER (WHERE prow_job_run_tests.prow_job_run_timestamp >= (|||TIMENOW||| - '14 days'::interval) AND prow_job_run_tests.prow_job_run_timestamp <= |||TIMENOW|||) AS runs, - COUNT(*) FILTER (WHERE prow_job_run_tests.status = 1 AND prow_job_run_tests.prow_job_run_timestamp >= (|||TIMENOW||| - '14 days'::interval) AND prow_job_run_tests.prow_job_run_timestamp <= |||TIMENOW|||) AS passes, - COUNT(*) FILTER (WHERE prow_job_run_tests.status = 13 AND prow_job_run_tests.prow_job_run_timestamp >= (|||TIMENOW||| - '14 days'::interval) AND prow_job_run_tests.prow_job_run_timestamp <= |||TIMENOW|||) AS flakes, - COUNT(*) FILTER (WHERE prow_job_run_tests.status = 12 AND prow_job_run_tests.prow_job_run_timestamp >= (|||TIMENOW||| - '14 days'::interval) AND prow_job_run_tests.prow_job_run_timestamp <= |||TIMENOW|||) AS failures -FROM - prow_job_run_tests - JOIN tests ON tests.id = prow_job_run_tests.test_id - JOIN prow_jobs ON prow_jobs.id = prow_job_run_tests.prow_job_id -WHERE - prow_job_run_tests.prow_job_run_timestamp > (|||TIMENOW||| - '14 days'::interval) -GROUP BY - tests.name, tests.id, date(prow_job_run_tests.prow_job_run_timestamp), prow_job_run_tests.prow_job_run_release, prow_jobs.name -` // TODO: remove distinct once bug fixed re dupes in release_job_runs const payloadTestFailuresMatView = ` diff --git a/pkg/flags/postgres_benchmarking_test.go b/pkg/flags/postgres_benchmarking_test.go index acbcf5be64..315d4f3091 100644 --- a/pkg/flags/postgres_benchmarking_test.go +++ b/pkg/flags/postgres_benchmarking_test.go @@ -246,7 +246,7 @@ func getBenchmarkCases(asOf time.Time) []benchmarkCase { CurrentPassPercent float64 } var result testResult - res := dbc.DB.Raw(query.QueryTestAnalysis, analyzeSince, benchmarkTestName, []string{benchmarkJobName}, benchmarkRelease) + res := dbc.DB.Raw(query.QueryTestAnalysis, analyzeSince, benchmarkRelease, benchmarkRelease, benchmarkTestName, []string{benchmarkJobName}) if res.Error != nil { return res.Error } @@ -481,9 +481,10 @@ func getBenchmarkCases(asOf time.Time) []benchmarkCase { var result passRate res := dbc.DB.Raw(query.QueryTestAnalysis, time.Now().Add(-24*14*time.Hour), + benchmarkRelease, + benchmarkRelease, benchmarkTestName, - []string{benchmarkJobName}, - benchmarkRelease).Scan(&result) + []string{benchmarkJobName}).Scan(&result) if res.Error != nil { return res.Error }