diff --git a/cmd/sippy/main.go b/cmd/sippy/main.go index a0715123a2..37cdeb38d0 100644 --- a/cmd/sippy/main.go +++ b/cmd/sippy/main.go @@ -57,6 +57,7 @@ func main() { NewVersionCommand(), NewAnnotateJobRunsCommand(), NewSeedDataCommand(), + NewVerifyCommand(), ) rootCmd.PersistentFlags().StringVar(&logLevel, "log-level", "info", diff --git a/cmd/sippy/verify.go b/cmd/sippy/verify.go new file mode 100644 index 0000000000..287c1bfb1e --- /dev/null +++ b/cmd/sippy/verify.go @@ -0,0 +1,140 @@ +package main + +import ( + "context" + "errors" + "fmt" + "time" + + "cloud.google.com/go/civil" + log "github.com/sirupsen/logrus" + "github.com/spf13/cobra" + "github.com/spf13/pflag" + + bqcachedclient "github.com/openshift/sippy/pkg/bigquery" + "github.com/openshift/sippy/pkg/db/verify" + "github.com/openshift/sippy/pkg/flags" + "github.com/openshift/sippy/pkg/flags/configflags" + "github.com/openshift/sippy/pkg/variantregistry" +) + +type VerifyFlags struct { + Date string + Checks []string + Release string + DBFlags *flags.PostgresFlags + BigQueryFlags *flags.BigQueryFlags + GoogleCloudFlags *flags.GoogleCloudFlags + ConfigFlags *configflags.ConfigFlags +} + +func NewVerifyFlags(now time.Time) *VerifyFlags { + return &VerifyFlags{ + Date: civil.DateOf(now.UTC()).AddDays(-2).String(), + DBFlags: flags.NewPostgresDatabaseFlags(), + BigQueryFlags: flags.NewBigQueryFlags(), + GoogleCloudFlags: flags.NewGoogleCloudFlags(), + ConfigFlags: configflags.NewConfigFlags(), + } +} + +func (f *VerifyFlags) BindFlags(fs *pflag.FlagSet) { + f.DBFlags.BindFlags(fs) + f.BigQueryFlags.BindFlags(fs) + f.GoogleCloudFlags.BindFlags(fs) + f.ConfigFlags.BindFlags(fs) + fs.StringVar(&f.Date, "date", f.Date, "UTC calendar date to verify (YYYY-MM-DD; defaults to the day before yesterday)") + fs.StringArrayVar(&f.Checks, "check", nil, "Check to run; repeat for multiple checks (bq-completeness, daily-totals, cumulative-summaries; defaults to all)") + fs.StringVar(&f.Release, "release", "", "Verify only this release (defaults to every configured and discovered release)") +} + +type verifyCommandDependencies struct { + now time.Time + run func(context.Context, *VerifyFlags, civil.Date, []verify.Check) (verify.Result, error) +} + +func NewVerifyCommand() *cobra.Command { + return newVerifyCommandWithDependencies(verifyCommandDependencies{now: time.Now(), run: runVerify}) +} + +func newVerifyCommandWithDependencies(dependencies verifyCommandDependencies) *cobra.Command { + f := NewVerifyFlags(dependencies.now) + cmd := &cobra.Command{ + Use: "verify", + Short: "Verify daily data integrity without modifying storage", + Args: cobra.NoArgs, + SilenceUsage: true, + RunE: func(cmd *cobra.Command, args []string) error { + date, err := civil.ParseDate(f.Date) + if err != nil { + return fmt.Errorf("invalid --date %q: expected YYYY-MM-DD: %w", f.Date, err) + } + checks, err := verify.ParseChecks(f.Checks) + if err != nil { + return err + } + result, runErr := dependencies.run(cmd.Context(), f, date, checks) + if runErr != nil && len(result.Summaries) == 0 { + for _, check := range checks { + result.Summaries = append(result.Summaries, verify.Summary{ + Check: check, Release: f.Release, Date: date, Passed: false, Error: runErr.Error(), + }) + } + } + result.Sort() + result.Log(log.StandardLogger()) + if runErr != nil { + return runErr + } + if !result.Passed() { + return fmt.Errorf("one or more verification checks failed") + } + return nil + }, + } + f.BindFlags(cmd.Flags()) + return cmd +} + +func runVerify(ctx context.Context, verifyFlags *VerifyFlags, date civil.Date, checks []verify.Check) (verify.Result, error) { + dbc, err := verifyFlags.DBFlags.GetDBClient() + if err != nil { + return verify.Result{}, fmt.Errorf("getting PostgreSQL client: %w", err) + } + + runner := verify.Runner{PostgreSQL: verify.NewPostgreSQL(dbc)} + var bqClient *bqcachedclient.Client + if verify.ContainsCheck(checks, verify.CheckBQCompleteness) { + var initializationErrors []error + config, configErr := verifyFlags.ConfigFlags.GetConfig() + if configErr != nil { + initializationErrors = append(initializationErrors, fmt.Errorf("loading Sippy config: %w", configErr)) + } else { + runner.Config = config + overrides, overrideErr := variantregistry.BuildSyntheticReleaseJobOverrides(config.Releases) + if overrideErr != nil { + initializationErrors = append(initializationErrors, fmt.Errorf("building synthetic release overrides: %w", overrideErr)) + } else { + runner.SyntheticReleaseOverrides = overrides + } + } + + opCtx, queryCtx := bqcachedclient.OpCtxForCronEnv(ctx, "verify") + bqClient, err = verifyFlags.BigQueryFlags.GetBigQueryClient( + queryCtx, opCtx, nil, verifyFlags.GoogleCloudFlags.ServiceAccountCredentialFile, + ) + if err != nil { + initializationErrors = append(initializationErrors, fmt.Errorf("initializing BigQuery client: %w", err)) + } else { + runner.BigQuery = verify.NewBigQuery(bqClient) + defer func() { + if closeErr := bqClient.BQ.Close(); closeErr != nil { + log.WithError(closeErr).Warn("closing BigQuery client") + } + }() + } + runner.BigQueryInitializationError = errors.Join(initializationErrors...) + } + + return runner.Run(ctx, verify.Options{Date: date, Checks: checks, Release: verifyFlags.Release}), nil +} diff --git a/cmd/sippy/verify_test.go b/cmd/sippy/verify_test.go new file mode 100644 index 0000000000..ff5016e91b --- /dev/null +++ b/cmd/sippy/verify_test.go @@ -0,0 +1,131 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "testing" + "time" + + "cloud.google.com/go/civil" + log "github.com/sirupsen/logrus" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/openshift/sippy/pkg/db/verify" +) + +func TestVerifyCommandDefaultsAndSelection(t *testing.T) { + now := time.Date(2026, 8, 28, 1, 30, 0, 0, time.FixedZone("west", -7*60*60)) + tests := []struct { + name string + args []string + wantDate civil.Date + wantChecks []verify.Check + wantRel string + }{ + { + name: "UTC day before yesterday and all checks", + wantDate: civil.Date{Year: 2026, Month: 8, Day: 26}, + wantChecks: verify.AllChecks, + }, + { + name: "explicit repeatable selection", + args: []string{"--date=2024-02-29", "--check=daily-totals", "--check=bq-completeness", "--release=4.20"}, + wantDate: civil.Date{Year: 2024, Month: 2, Day: 29}, + wantChecks: []verify.Check{verify.CheckBQCompleteness, verify.CheckDailyTotals}, + wantRel: "4.20", + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + called := false + cmd := newVerifyCommandWithDependencies(verifyCommandDependencies{ + now: now, + run: func(_ context.Context, flags *VerifyFlags, date civil.Date, checks []verify.Check) (verify.Result, error) { + called = true + assert.Equal(t, tt.wantDate, date) + assert.Equal(t, tt.wantChecks, checks) + assert.Equal(t, tt.wantRel, flags.Release) + return verify.Result{Summaries: []verify.Summary{{Check: checks[0], Date: date, Passed: true}}}, nil + }, + }) + cmd.SetArgs(tt.args) + require.NoError(t, cmd.Execute()) + assert.True(t, called) + }) + } +} + +func TestVerifyCommandValidationAndExit(t *testing.T) { + tests := []struct { + name string + args []string + result verify.Result + wantErr string + called bool + }{ + {name: "invalid date", args: []string{"--date=nope"}, wantErr: "invalid --date", called: false}, + {name: "invalid check", args: []string{"--check=nope"}, wantErr: "invalid --check", called: false}, + {name: "mismatch returns failure", result: verify.Result{Summaries: []verify.Summary{{Check: verify.CheckDailyTotals, Date: civil.Date{Year: 2026, Month: 1, Day: 1}, Passed: false}}}, wantErr: "one or more verification checks failed", called: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + called := false + cmd := newVerifyCommandWithDependencies(verifyCommandDependencies{ + now: time.Date(2026, 8, 27, 0, 0, 0, 0, time.UTC), + run: func(context.Context, *VerifyFlags, civil.Date, []verify.Check) (verify.Result, error) { + called = true + return tt.result, nil + }, + }) + cmd.SetArgs(tt.args) + err := cmd.Execute() + require.ErrorContains(t, err, tt.wantErr) + assert.Equal(t, tt.called, called) + }) + } +} + +func TestVerifyCommandHasNoFixFlag(t *testing.T) { + cmd := NewVerifyCommand() + assert.Nil(t, cmd.Flags().Lookup("fix")) + assert.Nil(t, cmd.PersistentFlags().Lookup("fix")) +} + +func TestVerifyCommandLogsSummaryOnRunError(t *testing.T) { + logger := log.StandardLogger() + var output bytes.Buffer + previousOutput := logger.Out + previousFormatter := logger.Formatter + previousLevel := logger.Level + logger.SetOutput(&output) + logger.SetFormatter(&log.JSONFormatter{}) + logger.SetLevel(log.InfoLevel) + t.Cleanup(func() { + logger.SetOutput(previousOutput) + logger.SetFormatter(previousFormatter) + logger.SetLevel(previousLevel) + }) + + runErr := errors.New("database unavailable") + cmd := newVerifyCommandWithDependencies(verifyCommandDependencies{ + now: time.Date(2026, 8, 27, 0, 0, 0, 0, time.UTC), + run: func(context.Context, *VerifyFlags, civil.Date, []verify.Check) (verify.Result, error) { + return verify.Result{}, runErr + }, + }) + cmd.SetArgs([]string{"--check=daily-totals", "--release=4.20"}) + require.ErrorIs(t, cmd.Execute(), runErr) + + var record map[string]any + require.NoError(t, json.Unmarshal(output.Bytes(), &record)) + assert.Equal(t, "verification summary", record["msg"]) + assert.Equal(t, "error", record["level"]) + assert.Equal(t, "daily-totals", record["check"]) + assert.Equal(t, "4.20", record["release"]) + assert.Equal(t, "2026-08-25", record["date"]) + assert.Equal(t, false, record["passed"]) + assert.Equal(t, runErr.Error(), record["error"]) +} diff --git a/docs/features/daily-data-integrity-verification.md b/docs/features/daily-data-integrity-verification.md new file mode 100644 index 0000000000..c618a89aab --- /dev/null +++ b/docs/features/daily-data-integrity-verification.md @@ -0,0 +1,53 @@ +# Daily data integrity verification + +The `sippy verify` command performs read-only checks of one UTC calendar day in +the Prow data pipeline. It reports discrepancies but never repairs data. + +## Usage + +```console +sippy verify [--date YYYY-MM-DD] [--check CHECK]... [--release RELEASE] +``` + +`--date` defaults to the UTC calendar day before yesterday. `--check` may be +repeated and accepts `bq-completeness`, `daily-totals`, and +`cumulative-summaries`; omitting it runs all three. `--release` limits every +selected check to one release. Without it, the command checks every release +definition and every non-empty, non-deleted historical release discovered in +`prow_jobs`. There is intentionally no active-release filter, so this can +include pseudo-releases with no data on the selected day. + +## Checks + +- `bq-completeness` compares deduplicated numeric Prow build IDs attributed to + each release in BigQuery with `prow_job_runs`. Both sources use the Prow + start-time half-open interval `[date 00:00:00Z, next date 00:00:00Z)`. + BigQuery retains the loader's terminal-state and non-null URL filters. + Malformed BigQuery build IDs are failures. +- `daily-totals` recomputes counts from `prow_job_run_tests` and compares them + in both directions with `test_daily_totals`. It uses the production composite + run join, normalizes a null suite to ID 0, separates lifecycle values, and + excludes runs labeled `InfraFailure`. Only successes, failures, flakes, and + runs are compared; timestamps are not compared. +- `cumulative-summaries` checks that each target-day cumulative row equals the + previous day's prefix counters plus the target day's daily counters. Keys + without daily data must carry forward, and first-day keys equal their daily + counters. The four prefix counters checked are successes, failures, flakes, + and runs. + +## Credentials and output + +PostgreSQL uses the standard Sippy database flags. BigQuery and Google +credential flags are only used when `bq-completeness` is selected. Selecting +that check without usable service-account credentials is a failed check; it is +not silently skipped. PostgreSQL-only selections require no Google +credentials. + +The command emits one bounded structured summary record for every applicable +`(check, release, date)` and separate deterministically ordered discrepancy +records. It runs all selected checks before returning final status. + +Exit status is 0 only when every selected check passes. Mismatches, malformed +IDs, missing selected BigQuery credentials, and operational errors return exit +status 1. The command has no `--fix` mode and performs no writes, migrations, +remediation, alerting, or `prow_job_run_test_outputs` verification. diff --git a/pkg/bigquery/bqlabel/labels.go b/pkg/bigquery/bqlabel/labels.go index eb35f29690..479c4072b8 100644 --- a/pkg/bigquery/bqlabel/labels.go +++ b/pkg/bigquery/bqlabel/labels.go @@ -102,6 +102,7 @@ const ( CacheLookup QueryValue = "cache-lookup" GATestStatusLoader QueryValue = "ga-test-status-loader" BackendDisruptionByRun QueryValue = "backend-disruption-by-run" + VerifyProwJobs QueryValue = "verify-prow-jobs" ) // sanitizeLabelValue sanitizes a label value to meet BigQuery requirements: diff --git a/pkg/dataloader/prowloader/prow.go b/pkg/dataloader/prowloader/prow.go index 13809c1125..c14f37843b 100644 --- a/pkg/dataloader/prowloader/prow.go +++ b/pkg/dataloader/prowloader/prow.go @@ -56,26 +56,24 @@ import ( var gcsPathStrip = regexp.MustCompile(`.*/gs/[^/]+/`) type ProwLoader struct { - ctx context.Context - dbc *db.DB - errors []error - githubClient *github.Client - bigQueryClient *bqcachedclient.Client - maxConcurrency int - prowJobCache map[string]*models.ProwJob - variantManager testidentification.VariantManager - syntheticTestManager synthetictests.SyntheticTestManager - syntheticReleaseJobOverrides *releaseoverride.SyntheticReleaseOverrides - releases []string - releaseSet sets.Set[string] - releaseRegexps map[string][]*regexp.Regexp - config *v1config.SippyConfig - ghCommenter *commenter.GitHubCommenter - gcsClient *storage.Client - promPusher *push.Pusher - loadSince *time.Time - labelsCache map[string]pq.StringArray - currentDate civil.Date + ctx context.Context + dbc *db.DB + errors []error + githubClient *github.Client + bigQueryClient *bqcachedclient.Client + maxConcurrency int + prowJobCache map[string]*models.ProwJob + variantManager testidentification.VariantManager + syntheticTestManager synthetictests.SyntheticTestManager + releases []string + releaseAttributor *ReleaseAttributor + config *v1config.SippyConfig + ghCommenter *commenter.GitHubCommenter + gcsClient *storage.Client + promPusher *push.Pusher + loadSince *time.Time + labelsCache map[string]pq.StringArray + currentDate civil.Date } func New( @@ -93,41 +91,30 @@ func New( loadSince *time.Time, syntheticReleaseJobOverrides *releaseoverride.SyntheticReleaseOverrides) *ProwLoader { - compiledRegexps := make(map[string][]*regexp.Regexp, len(releases)) - for _, release := range releases { - if cfg, ok := config.Releases[release]; ok { - for _, expr := range cfg.Regexp { - re, err := regexp.Compile(expr) - if err != nil { - log.WithError(err).WithField("release", release).WithField("regex", expr).Error("invalid regex in configuration") - continue - } - compiledRegexps[release] = append(compiledRegexps[release], re) - } - } - } - + releaseAttributor := NewReleaseAttributor(releases, config, syntheticReleaseJobOverrides) return &ProwLoader{ - ctx: ctx, - dbc: dbc, - gcsClient: gcsClient, - githubClient: githubClient, - bigQueryClient: bigQueryClient, - maxConcurrency: 50, - syntheticTestManager: syntheticTestManager, - syntheticReleaseJobOverrides: syntheticReleaseJobOverrides, - variantManager: variantManager, - releases: releases, - releaseSet: sets.New[string](releases...), - releaseRegexps: compiledRegexps, - config: config, - ghCommenter: ghCommenter, - promPusher: promPusher, - loadSince: loadSince, - currentDate: civil.DateOf(time.Now().UTC()), + ctx: ctx, + dbc: dbc, + gcsClient: gcsClient, + githubClient: githubClient, + bigQueryClient: bigQueryClient, + maxConcurrency: 50, + syntheticTestManager: syntheticTestManager, + variantManager: variantManager, + releases: releases, + releaseAttributor: releaseAttributor, + config: config, + ghCommenter: ghCommenter, + promPusher: promPusher, + loadSince: loadSince, + currentDate: civil.DateOf(time.Now().UTC()), } } +func (pl *ProwLoader) matchRelease(pj *prow.ProwJob) string { + return pl.releaseAttributor.Match(pj) +} + const DefaultLookbackDays = 14 func resolveFrom(since *time.Time, to time.Time) time.Time { @@ -358,43 +345,6 @@ func (pl *ProwLoader) Load() { } } -// matchRelease returns the release a prow job belongs to, or "" if it -// doesn't match any configured release. For /payload sub-jobs (identified -// by releaseJobName annotation + PR refs), it returns the Presubmits -// pseudo-release. -func (pl *ProwLoader) matchRelease(pj *prow.ProwJob) string { - if _, ok := pj.Annotations["releaseJobName"]; ok && pj.Spec.Refs != nil { - if pl.releaseSet.Has(models.ReleasePresubmits) { - return models.ReleasePresubmits - } - return "" - } - - jobName := pj.Spec.Job - if release, ok := pl.syntheticReleaseJobOverrides.Lookup(jobName); ok { - if pl.releaseSet.Has(release) { - return release - } - return "" - } - - for _, release := range pl.releases { - cfg, ok := pl.config.Releases[release] - if !ok { - continue - } - if val, ok := cfg.Jobs[jobName]; val && ok { - return release - } - for _, re := range pl.releaseRegexps[release] { - if re.MatchString(jobName) { - return release - } - } - } - return "" -} - // isPayloadPresubmit returns true if the prow job is a /payload sub-job. func isPayloadPresubmit(pj *prow.ProwJob) bool { _, hasAnnotation := pj.Annotations["releaseJobName"] diff --git a/pkg/dataloader/prowloader/prow_test.go b/pkg/dataloader/prowloader/prow_test.go index 2ee5d50c39..a72cfaedaf 100644 --- a/pkg/dataloader/prowloader/prow_test.go +++ b/pkg/dataloader/prowloader/prow_test.go @@ -1,12 +1,10 @@ package prowloader import ( - "regexp" "testing" "time" "github.com/stretchr/testify/assert" - "k8s.io/apimachinery/pkg/util/sets" v1config "github.com/openshift/sippy/pkg/apis/config/v1" "github.com/openshift/sippy/pkg/apis/prow" @@ -171,11 +169,11 @@ func TestMatchRelease(t *testing.T) { { name: "payload presubmit with Presubmits in release set", pl: &ProwLoader{ - releaseSet: sets.New[string](models.ReleasePresubmits), - releases: []string{models.ReleasePresubmits}, - config: &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{}}, - releaseRegexps: map[string][]*regexp.Regexp{}, - syntheticReleaseJobOverrides: releaseoverride.New(), + releaseAttributor: NewReleaseAttributor( + []string{models.ReleasePresubmits}, + &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{}}, + releaseoverride.New(), + ), }, pj: &prow.ProwJob{ Annotations: map[string]string{"releaseJobName": "periodic-ci-openshift-release-master-nightly-4.18-e2e-aws-ovn"}, @@ -186,11 +184,11 @@ func TestMatchRelease(t *testing.T) { { name: "payload presubmit without Presubmits in release set", pl: &ProwLoader{ - releaseSet: sets.New[string]("4.18"), - releases: []string{"4.18"}, - config: &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{}}, - releaseRegexps: map[string][]*regexp.Regexp{}, - syntheticReleaseJobOverrides: releaseoverride.New(), + releaseAttributor: NewReleaseAttributor( + []string{"4.18"}, + &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{}}, + releaseoverride.New(), + ), }, pj: &prow.ProwJob{ Annotations: map[string]string{"releaseJobName": "some-job"}, @@ -201,13 +199,13 @@ func TestMatchRelease(t *testing.T) { { name: "regular job matching configured release by regex", pl: &ProwLoader{ - releaseSet: sets.New[string]("4.18"), - releases: []string{"4.18"}, - config: &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{ - "4.18": {}, - }}, - releaseRegexps: map[string][]*regexp.Regexp{"4.18": {regexp.MustCompile(`-4\.18-`)}}, - syntheticReleaseJobOverrides: releaseoverride.New(), + releaseAttributor: NewReleaseAttributor( + []string{"4.18"}, + &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{ + "4.18": {Regexp: []string{`-4\.18-`}}, + }}, + releaseoverride.New(), + ), }, pj: &prow.ProwJob{ Annotations: map[string]string{}, @@ -218,13 +216,13 @@ func TestMatchRelease(t *testing.T) { { name: "regular job matching no release", pl: &ProwLoader{ - releaseSet: sets.New[string]("4.18"), - releases: []string{"4.18"}, - config: &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{ - "4.18": {}, - }}, - releaseRegexps: map[string][]*regexp.Regexp{"4.18": {regexp.MustCompile(`-4\.18-`)}}, - syntheticReleaseJobOverrides: releaseoverride.New(), + releaseAttributor: NewReleaseAttributor( + []string{"4.18"}, + &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{ + "4.18": {Regexp: []string{`-4\.18-`}}, + }}, + releaseoverride.New(), + ), }, pj: &prow.ProwJob{ Annotations: map[string]string{}, diff --git a/pkg/dataloader/prowloader/release_attribution.go b/pkg/dataloader/prowloader/release_attribution.go new file mode 100644 index 0000000000..bf91dd6f27 --- /dev/null +++ b/pkg/dataloader/prowloader/release_attribution.go @@ -0,0 +1,88 @@ +package prowloader + +import ( + "regexp" + + log "github.com/sirupsen/logrus" + "k8s.io/apimachinery/pkg/util/sets" + + v1config "github.com/openshift/sippy/pkg/apis/config/v1" + "github.com/openshift/sippy/pkg/apis/prow" + "github.com/openshift/sippy/pkg/db/models" + "github.com/openshift/sippy/pkg/releaseoverride" +) + +// ReleaseAttributor applies the same release matching rules to Prow jobs no +// matter whether they are being loaded or independently verified. +type ReleaseAttributor struct { + releases []string + releaseSet sets.Set[string] + releaseRegexps map[string][]*regexp.Regexp + config *v1config.SippyConfig + syntheticReleaseJobOverrides *releaseoverride.SyntheticReleaseOverrides +} + +func NewReleaseAttributor(releases []string, config *v1config.SippyConfig, syntheticReleaseJobOverrides *releaseoverride.SyntheticReleaseOverrides) *ReleaseAttributor { + if config == nil { + config = &v1config.SippyConfig{} + } + compiledRegexps := make(map[string][]*regexp.Regexp, len(releases)) + for _, release := range releases { + if cfg, ok := config.Releases[release]; ok { + for _, expr := range cfg.Regexp { + re, err := regexp.Compile(expr) + if err != nil { + log.WithError(err).WithFields(log.Fields{"release": release, "regex": expr}).Error("invalid regex in configuration") + continue + } + compiledRegexps[release] = append(compiledRegexps[release], re) + } + } + } + + return &ReleaseAttributor{ + releases: releases, + releaseSet: sets.New[string](releases...), + releaseRegexps: compiledRegexps, + config: config, + syntheticReleaseJobOverrides: syntheticReleaseJobOverrides, + } +} + +// Match returns the release a Prow job belongs to, or an empty string when it +// does not match one of the releases supplied to the attributor. +func (a *ReleaseAttributor) Match(pj *prow.ProwJob) string { + if a == nil || pj == nil { + return "" + } + if _, ok := pj.Annotations["releaseJobName"]; ok && pj.Spec.Refs != nil { + if a.releaseSet.Has(models.ReleasePresubmits) { + return models.ReleasePresubmits + } + return "" + } + + jobName := pj.Spec.Job + if release, ok := a.syntheticReleaseJobOverrides.Lookup(jobName); ok { + if a.releaseSet.Has(release) { + return release + } + return "" + } + + for _, release := range a.releases { + cfg, ok := a.config.Releases[release] + if !ok { + continue + } + if enabled, ok := cfg.Jobs[jobName]; ok && enabled { + return release + } + for _, re := range a.releaseRegexps[release] { + if re.MatchString(jobName) { + return release + } + } + } + return "" +} diff --git a/pkg/dataloader/prowloader/release_attribution_test.go b/pkg/dataloader/prowloader/release_attribution_test.go new file mode 100644 index 0000000000..242ebc1d65 --- /dev/null +++ b/pkg/dataloader/prowloader/release_attribution_test.go @@ -0,0 +1,54 @@ +package prowloader + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + v1config "github.com/openshift/sippy/pkg/apis/config/v1" + "github.com/openshift/sippy/pkg/apis/prow" + "github.com/openshift/sippy/pkg/db/models" + "github.com/openshift/sippy/pkg/releaseoverride" +) + +func TestReleaseAttributor(t *testing.T) { + overrides := releaseoverride.New() + require.NoError(t, overrides.AddExact("synthetic-exact", "rosa-stage")) + require.NoError(t, overrides.AddRegexp(`^synthetic-regexp-`, "rosa-stage")) + config := &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{ + "4.20": {Jobs: map[string]bool{"configured": true, "disabled": false}, Regexp: []string{`-4\.20-`}}, + "rosa-stage": {Synthetic: true}, + }} + attributor := NewReleaseAttributor([]string{"4.20", "rosa-stage", models.ReleasePresubmits}, config, overrides) + + tests := []struct { + name string + job *prow.ProwJob + want string + }{ + {name: "configured exact", job: jobForAttribution("configured", nil, false), want: "4.20"}, + {name: "configured disabled", job: jobForAttribution("disabled", nil, false), want: ""}, + {name: "configured regexp", job: jobForAttribution("periodic-4.20-e2e", nil, false), want: "4.20"}, + {name: "synthetic exact has priority", job: jobForAttribution("synthetic-exact", nil, false), want: "rosa-stage"}, + {name: "synthetic regexp", job: jobForAttribution("synthetic-regexp-job", nil, false), want: "rosa-stage"}, + {name: "payload presubmit", job: jobForAttribution("generated-name", map[string]string{"releaseJobName": "canonical"}, true), want: models.ReleasePresubmits}, + {name: "payload annotation without refs", job: jobForAttribution("generated-name", map[string]string{"releaseJobName": "canonical"}, false), want: ""}, + {name: "unmatched", job: jobForAttribution("other", nil, false), want: ""}, + {name: "nil job", job: nil, want: ""}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, attributor.Match(tt.job)) + }) + } + assert.Equal(t, "", (*ReleaseAttributor)(nil).Match(jobForAttribution("configured", nil, false))) +} + +func jobForAttribution(name string, annotations map[string]string, hasRefs bool) *prow.ProwJob { + job := &prow.ProwJob{Spec: prow.ProwJobSpec{Job: name}, Annotations: annotations} + if hasRefs { + job.Spec.Refs = &prow.Refs{} + } + return job +} diff --git a/pkg/db/verify/bq_completeness.go b/pkg/db/verify/bq_completeness.go new file mode 100644 index 0000000000..8617e537cb --- /dev/null +++ b/pkg/db/verify/bq_completeness.go @@ -0,0 +1,246 @@ +package verify + +import ( + "context" + "fmt" + "regexp" + "sort" + "strconv" + "strings" + "time" + + "cloud.google.com/go/bigquery" + "cloud.google.com/go/civil" + "google.golang.org/api/iterator" + "k8s.io/apimachinery/pkg/util/sets" + + v1config "github.com/openshift/sippy/pkg/apis/config/v1" + "github.com/openshift/sippy/pkg/apis/prow" + bqcachedclient "github.com/openshift/sippy/pkg/bigquery" + "github.com/openshift/sippy/pkg/bigquery/bqlabel" + "github.com/openshift/sippy/pkg/dataloader/prowloader" + "github.com/openshift/sippy/pkg/releaseoverride" +) + +type BQCompletenessVerifier struct { + PostgreSQL ProwJobRunReader + BigQuery BigQueryReader + BigQueryInitializationError error + Config *v1config.SippyConfig + SyntheticReleaseOverrides *releaseoverride.SyntheticReleaseOverrides +} + +func (v *BQCompletenessVerifier) Verify(ctx context.Context, scope Scope) Result { + result := Result{} + if v.BigQueryInitializationError != nil { + for _, release := range scope.Releases { + result.Summaries = append(result.Summaries, operationalSummary(CheckBQCompleteness, release, scope.Date, v.BigQueryInitializationError)) + } + return result + } + if v.BigQuery == nil { + err := fmt.Errorf("BigQuery client is not initialized") + for _, release := range scope.Releases { + result.Summaries = append(result.Summaries, operationalSummary(CheckBQCompleteness, release, scope.Date, err)) + } + return result + } + + start := scope.Date.In(time.UTC) + end := scope.Date.AddDays(1).In(time.UTC) + jobs, err := v.BigQuery.ProwJobs(ctx, start, end) + if err != nil { + for _, release := range scope.Releases { + result.Summaries = append(result.Summaries, operationalSummary(CheckBQCompleteness, release, scope.Date, err)) + } + return result + } + postgresIDs, err := v.PostgreSQL.ProwJobRunIDs(ctx, start, end) + if err != nil { + for _, release := range scope.Releases { + result.Summaries = append(result.Summaries, operationalSummary(CheckBQCompleteness, release, scope.Date, err)) + } + return result + } + + attributor := prowloader.NewReleaseAttributor(scope.Releases, v.Config, v.SyntheticReleaseOverrides) + bqIDs := make(map[string]map[BuildID]struct{}, len(scope.Releases)) + malformedSets := make(map[string]sets.Set[string], len(scope.Releases)) + for _, job := range jobs { + pj := &prow.ProwJob{ + Annotations: job.Annotations, + Spec: prow.ProwJobSpec{Job: job.JobName}, + } + if job.HasRefs { + pj.Spec.Refs = &prow.Refs{} + } + release := attributor.Match(pj) + if release == "" { + continue + } + value := strings.TrimSpace(job.BuildID) + id, err := strconv.ParseUint(value, 10, 64) + if err != nil { + if malformedSets[release] == nil { + malformedSets[release] = sets.New[string]() + } + malformedSets[release].Insert(value) + continue + } + if bqIDs[release] == nil { + bqIDs[release] = make(map[BuildID]struct{}) + } + bqIDs[release][BuildID(id)] = struct{}{} + } + for _, release := range scope.Releases { + malformed := sets.List(malformedSets[release]) + summary, discrepancies := CompareBuildIDs(release, scope.Date, bqIDs[release], postgresIDs[release], malformed) + result.Summaries = append(result.Summaries, summary) + result.Discrepancies = append(result.Discrepancies, discrepancies...) + } + return result +} + +func (p *PostgreSQL) ProwJobRunIDs(ctx context.Context, start, end time.Time) (map[string]map[BuildID]struct{}, error) { + if err := p.validate(); err != nil { + return nil, err + } + type row struct { + Release string + ID uint64 + } + var rows []row + err := p.dbc.DB.WithContext(ctx).Raw(` + SELECT prow_job_release AS release, id + FROM prow_job_runs + WHERE timestamp >= ? AND timestamp < ? AND deleted_at IS NULL + ORDER BY prow_job_release, id + `, start, end).Scan(&rows).Error + if err != nil { + return nil, fmt.Errorf("querying PostgreSQL Prow job runs: %w", err) + } + result := make(map[string]map[BuildID]struct{}) + for _, row := range rows { + if result[row.Release] == nil { + result[row.Release] = make(map[BuildID]struct{}) + } + result[row.Release][BuildID(row.ID)] = struct{}{} + } + return result, nil +} + +type BigQuery struct { + client *bqcachedclient.Client +} + +func NewBigQuery(client *bqcachedclient.Client) *BigQuery { + return &BigQuery{client: client} +} + +var bigQueryIdentifier = regexp.MustCompile(`^[A-Za-z0-9_-]+$`) + +func (b *BigQuery) ProwJobs(ctx context.Context, start, end time.Time) ([]BQJob, error) { + if b == nil || b.client == nil || b.client.BQ == nil { + return nil, fmt.Errorf("BigQuery client is not initialized") + } + project := b.client.BQ.Project() + dataset := b.client.Dataset + if !bigQueryIdentifier.MatchString(project) || !bigQueryIdentifier.MatchString(dataset) { + return nil, fmt.Errorf("invalid BigQuery project or dataset identifier") + } + query := b.client.Query(ctx, bqlabel.VerifyProwJobs, fmt.Sprintf(` + SELECT + prowjob_job_name, + IFNULL(CAST(prowjob_build_id AS STRING), '') AS prowjob_build_id, + prowjob_annotations, + IFNULL(CAST(pr_number AS STRING), '') AS pr_number + FROM %s + WHERE TIMESTAMP(prowjob_start) >= @start + AND TIMESTAMP(prowjob_start) < @end + AND prowjob_url IS NOT NULL + AND prowjob_state NOT IN ('pending', 'triggered') + ORDER BY TIMESTAMP(prowjob_start), prowjob_build_id + `, "`"+project+"."+dataset+".jobs`")) + query.Parameters = []bigquery.QueryParameter{ + {Name: "start", Value: start}, + {Name: "end", Value: end}, + } + iteratorRows, err := query.Read(ctx) + if err != nil { + return nil, fmt.Errorf("querying BigQuery Prow jobs: %w", err) + } + type bqRow struct { + JobName string `bigquery:"prowjob_job_name"` + BuildID string `bigquery:"prowjob_build_id"` + Annotations []string `bigquery:"prowjob_annotations"` + PRNumber string `bigquery:"pr_number"` + } + result := make([]BQJob, 0) + for { + var row bqRow + err := iteratorRows.Next(&row) + if err == iterator.Done { + break + } + if err != nil { + return nil, fmt.Errorf("reading BigQuery Prow jobs: %w", err) + } + annotations := make(map[string]string, len(row.Annotations)) + for _, annotation := range row.Annotations { + parts := strings.SplitN(annotation, "=", 2) + if len(parts) == 2 { + annotations[parts[0]] = parts[1] + } + } + result = append(result, BQJob{ + BuildID: row.BuildID, JobName: row.JobName, Annotations: annotations, HasRefs: row.PRNumber != "", + }) + } + return result, nil +} + +func CompareBuildIDs(release string, date civil.Date, bqIDs, postgresIDs map[BuildID]struct{}, malformed []string) (Summary, []Discrepancy) { + discrepancies := make([]Discrepancy, 0) + malformedSet := sets.New[string]() + for _, value := range malformed { + malformedSet.Insert(strings.TrimSpace(value)) + } + for _, value := range sets.List(malformedSet) { + kind := "malformed-build-id" + detail := "BigQuery build ID is not an unsigned integer" + if value == "" { + kind = "missing-build-id" + detail = "BigQuery build ID is blank" + } + discrepancies = append(discrepancies, Discrepancy{ + Check: CheckBQCompleteness, Release: release, Date: date, + Kind: kind, Key: value, Detail: detail, + }) + } + for _, id := range sortedBuildIDs(bqIDs) { + if _, ok := postgresIDs[id]; !ok { + discrepancies = append(discrepancies, Discrepancy{ + Check: CheckBQCompleteness, Release: release, Date: date, + Kind: "missing-in-postgres", Key: fmt.Sprint(uint64(id)), Expected: "present", Actual: "missing", + }) + } + } + for _, id := range sortedBuildIDs(postgresIDs) { + if _, ok := bqIDs[id]; !ok { + discrepancies = append(discrepancies, Discrepancy{ + Check: CheckBQCompleteness, Release: release, Date: date, + Kind: "missing-in-bigquery", Key: fmt.Sprint(uint64(id)), Expected: "present", Actual: "missing", + }) + } + } + return summary(CheckBQCompleteness, release, date, len(bqIDs), len(postgresIDs), len(discrepancies)), discrepancies +} + +func sortedBuildIDs(values map[BuildID]struct{}) []BuildID { + ids := make([]BuildID, 0, len(values)) + for id := range values { + ids = append(ids, id) + } + sort.Slice(ids, func(i, j int) bool { return ids[i] < ids[j] }) + return ids +} diff --git a/pkg/db/verify/comparison.go b/pkg/db/verify/comparison.go new file mode 100644 index 0000000000..b6ada3ad09 --- /dev/null +++ b/pkg/db/verify/comparison.go @@ -0,0 +1,95 @@ +package verify + +import ( + "fmt" + "sort" + + "cloud.google.com/go/civil" +) + +func compareCounts(check Check, release string, date civil.Date, expectedRows, actualRows []DailyRow) (Summary, []Discrepancy) { + expectedRows = scopedRows(expectedRows, release, date) + actualRows = scopedRows(actualRows, release, date) + expected := rowsByKey(expectedRows) + actual := rowsByKey(actualRows) + discrepancies := make([]Discrepancy, 0) + keys := unionKeys(expected, actual) + for _, key := range keys { + expectedCounts, hasExpected := expected[key] + actualCounts, hasActual := actual[key] + if !hasExpected { + discrepancies = append(discrepancies, Discrepancy{ + Check: check, Release: release, Date: date, Kind: "unexpected-row", Key: key.String(), + Expected: "missing", Actual: "present", + }) + continue + } + if !hasActual { + discrepancies = append(discrepancies, Discrepancy{ + Check: check, Release: release, Date: date, Kind: "missing-row", Key: key.String(), + Expected: "present", Actual: "missing", + }) + continue + } + fields := []struct { + name string + expected int64 + actual int64 + }{ + {"successes", expectedCounts.Successes, actualCounts.Successes}, + {"failures", expectedCounts.Failures, actualCounts.Failures}, + {"flakes", expectedCounts.Flakes, actualCounts.Flakes}, + {"runs", expectedCounts.Runs, actualCounts.Runs}, + } + for _, field := range fields { + if field.expected != field.actual { + discrepancies = append(discrepancies, Discrepancy{ + Check: check, Release: release, Date: date, Kind: "count-mismatch", Key: key.String(), Field: field.name, + Expected: fmt.Sprint(field.expected), Actual: fmt.Sprint(field.actual), + }) + } + } + } + return summary(check, release, date, len(expected), len(actual), len(discrepancies)), discrepancies +} + +func scopedRows(rows []DailyRow, release string, date civil.Date) []DailyRow { + result := make([]DailyRow, len(rows)) + copy(result, rows) + for i := range result { + result[i].Key.Release = release + result[i].Key.Date = date + } + return result +} + +func summary(check Check, release string, date civil.Date, expected, actual, discrepancyCount int) Summary { + return Summary{ + Check: check, Release: release, Date: date, Passed: discrepancyCount == 0, + ExpectedRows: expected, ActualRows: actual, Discrepancies: discrepancyCount, + } +} + +func rowsByKey(rows []DailyRow) map[SummaryKey]Counts { + result := make(map[SummaryKey]Counts, len(rows)) + for _, row := range rows { + result[row.Key] = row.Counts + } + return result +} + +func unionKeys(a, b map[SummaryKey]Counts) []SummaryKey { + set := make(map[SummaryKey]struct{}, len(a)+len(b)) + for key := range a { + set[key] = struct{}{} + } + for key := range b { + set[key] = struct{}{} + } + keys := make([]SummaryKey, 0, len(set)) + for key := range set { + keys = append(keys, key) + } + sort.Slice(keys, func(i, j int) bool { return keys[i].String() < keys[j].String() }) + return keys +} diff --git a/pkg/db/verify/cumulative_summaries.go b/pkg/db/verify/cumulative_summaries.go new file mode 100644 index 0000000000..4d5d53e36a --- /dev/null +++ b/pkg/db/verify/cumulative_summaries.go @@ -0,0 +1,100 @@ +package verify + +import ( + "context" + "fmt" + + "cloud.google.com/go/civil" +) + +type CumulativeSummariesVerifier struct { + PostgreSQL CumulativeSummariesReader +} + +func (v *CumulativeSummariesVerifier) Verify(ctx context.Context, scope Scope) Result { + result := Result{} + for _, release := range scope.Releases { + rows, err := v.PostgreSQL.CumulativeRows(ctx, release, scope.Date) + if err != nil { + result.Summaries = append(result.Summaries, operationalSummary(CheckCumulativeSummaries, release, scope.Date, err)) + continue + } + summary, discrepancies := CompareCumulative(release, scope.Date, rows) + result.Summaries = append(result.Summaries, summary) + result.Discrepancies = append(result.Discrepancies, discrepancies...) + } + return result +} + +func (p *PostgreSQL) CumulativeRows(ctx context.Context, release string, date civil.Date) (CumulativeRows, error) { + if err := p.validate(); err != nil { + return CumulativeRows{}, err + } + type row struct { + Source string + TestID uint64 + ProwJobID uint64 + SuiteID uint64 + Lifecycle string + Successes int64 + Failures int64 + Flakes int64 + Runs int64 + } + var rows []row + err := p.dbc.DB.WithContext(ctx).Raw(` + SELECT 'previous' AS source, test_id, prow_job_id, suite_id, lifecycle, + prefix_sum_successes AS successes, prefix_sum_failures AS failures, + prefix_sum_flakes AS flakes, prefix_sum_runs AS runs + FROM test_cumulative_summaries + WHERE date = ? AND release = ? + UNION ALL + SELECT 'daily' AS source, test_id, prow_job_id, suite_id, lifecycle, + successes, failures, flakes, runs + FROM test_daily_totals + WHERE date = ? AND release = ? + UNION ALL + SELECT 'target' AS source, test_id, prow_job_id, suite_id, lifecycle, + prefix_sum_successes AS successes, prefix_sum_failures AS failures, + prefix_sum_flakes AS flakes, prefix_sum_runs AS runs + FROM test_cumulative_summaries + WHERE date = ? AND release = ? + ORDER BY source, test_id, prow_job_id, suite_id, lifecycle + `, date.AddDays(-1), release, date, release, date, release).Scan(&rows).Error + if err != nil { + return CumulativeRows{}, fmt.Errorf("querying cumulative inputs for release %s: %w", release, err) + } + result := CumulativeRows{} + for _, row := range rows { + value := DailyRow{ + Key: SummaryKey{TestID: row.TestID, ProwJobID: row.ProwJobID, SuiteID: row.SuiteID, Lifecycle: row.Lifecycle}, + Counts: Counts{Successes: row.Successes, Failures: row.Failures, Flakes: row.Flakes, Runs: row.Runs}, + } + switch row.Source { + case "previous": + result.Previous = append(result.Previous, value) + case "daily": + result.Daily = append(result.Daily, value) + case "target": + result.Target = append(result.Target, value) + } + } + return result, nil +} + +func CompareCumulative(release string, date civil.Date, rows CumulativeRows) (Summary, []Discrepancy) { + previous := rowsByKey(scopedRows(rows.Previous, release, date)) + daily := rowsByKey(scopedRows(rows.Daily, release, date)) + expected := make(map[SummaryKey]Counts, len(previous)+len(daily)) + for key, counts := range previous { + expected[key] = counts + } + for key, counts := range daily { + expected[key] = expected[key].Add(counts) + } + expectedRows := make([]DailyRow, 0, len(expected)) + for key, counts := range expected { + expectedRows = append(expectedRows, DailyRow{Key: key, Counts: counts}) + } + return compareCounts(CheckCumulativeSummaries, release, date, expectedRows, rows.Target) +} diff --git a/pkg/db/verify/daily_totals.go b/pkg/db/verify/daily_totals.go new file mode 100644 index 0000000000..84bc3889f7 --- /dev/null +++ b/pkg/db/verify/daily_totals.go @@ -0,0 +1,100 @@ +package verify + +import ( + "context" + "fmt" + "time" + + "cloud.google.com/go/civil" +) + +type DailyTotalsVerifier struct { + PostgreSQL DailyTotalsReader +} + +func (v *DailyTotalsVerifier) Verify(ctx context.Context, scope Scope) Result { + result := Result{} + for _, release := range scope.Releases { + raw, stored, err := v.PostgreSQL.DailyRows(ctx, release, scope.Date) + if err != nil { + result.Summaries = append(result.Summaries, operationalSummary(CheckDailyTotals, release, scope.Date, err)) + continue + } + summary, discrepancies := CompareDaily(release, scope.Date, raw, stored) + result.Summaries = append(result.Summaries, summary) + result.Discrepancies = append(result.Discrepancies, discrepancies...) + } + return result +} + +func (p *PostgreSQL) DailyRows(ctx context.Context, release string, date civil.Date) ([]DailyRow, []DailyRow, error) { + if err := p.validate(); err != nil { + return nil, nil, err + } + start := date.In(time.UTC) + end := date.AddDays(1).In(time.UTC) + raw, err := p.queryDaily(ctx, ` + SELECT + pjrt.test_id, + pjrt.prow_job_id, + COALESCE(pjrt.suite_id, 0) AS suite_id, + pjrt.lifecycle, + COUNT(*) FILTER (WHERE pjrt.status = 1) AS successes, + COUNT(*) FILTER (WHERE pjrt.status = 12) AS failures, + COUNT(*) FILTER (WHERE pjrt.status = 13) AS flakes, + COUNT(*) AS runs + FROM prow_job_run_tests pjrt + JOIN prow_job_runs pjr + ON pjr.id = pjrt.prow_job_run_id + AND pjr.prow_job_release = pjrt.prow_job_run_release + AND pjr.timestamp = pjrt.prow_job_run_timestamp + WHERE pjrt.prow_job_run_timestamp >= ? + AND pjrt.prow_job_run_timestamp < ? + AND pjrt.prow_job_run_release = ? + AND (pjr.labels IS NULL OR NOT (pjr.labels @> ARRAY['InfraFailure'])) + GROUP BY pjrt.test_id, pjrt.prow_job_id, COALESCE(pjrt.suite_id, 0), pjrt.lifecycle + ORDER BY pjrt.test_id, pjrt.prow_job_id, COALESCE(pjrt.suite_id, 0), pjrt.lifecycle + `, start, end, release) + if err != nil { + return nil, nil, fmt.Errorf("querying raw daily totals for release %s: %w", release, err) + } + summaryRows, err := p.queryDaily(ctx, ` + SELECT test_id, prow_job_id, suite_id, lifecycle, successes, failures, flakes, runs + FROM test_daily_totals + WHERE date = ? AND release = ? + ORDER BY test_id, prow_job_id, suite_id, lifecycle + `, date, release) + if err != nil { + return nil, nil, fmt.Errorf("querying stored daily totals for release %s: %w", release, err) + } + return raw, summaryRows, nil +} + +func (p *PostgreSQL) queryDaily(ctx context.Context, query string, args ...any) ([]DailyRow, error) { + type row struct { + TestID uint64 + ProwJobID uint64 + SuiteID uint64 + Lifecycle string + Successes int64 + Failures int64 + Flakes int64 + Runs int64 + } + var rows []row + if err := p.dbc.DB.WithContext(ctx).Raw(query, args...).Scan(&rows).Error; err != nil { + return nil, err + } + result := make([]DailyRow, 0, len(rows)) + for _, row := range rows { + result = append(result, DailyRow{ + Key: SummaryKey{TestID: row.TestID, ProwJobID: row.ProwJobID, SuiteID: row.SuiteID, Lifecycle: row.Lifecycle}, + Counts: Counts{Successes: row.Successes, Failures: row.Failures, Flakes: row.Flakes, Runs: row.Runs}, + }) + } + return result, nil +} + +func CompareDaily(release string, date civil.Date, rawRows, summaryRows []DailyRow) (Summary, []Discrepancy) { + return compareCounts(CheckDailyTotals, release, date, rawRows, summaryRows) +} diff --git a/pkg/db/verify/runner.go b/pkg/db/verify/runner.go new file mode 100644 index 0000000000..b0303f7f4b --- /dev/null +++ b/pkg/db/verify/runner.go @@ -0,0 +1,106 @@ +package verify + +import ( + "context" + "fmt" + "strings" + + "cloud.google.com/go/civil" + "k8s.io/apimachinery/pkg/util/sets" + + v1config "github.com/openshift/sippy/pkg/apis/config/v1" + "github.com/openshift/sippy/pkg/releaseoverride" +) + +type Runner struct { + PostgreSQL PostgreSQLReader + BigQuery BigQueryReader + BigQueryInitializationError error + Config *v1config.SippyConfig + SyntheticReleaseOverrides *releaseoverride.SyntheticReleaseOverrides +} + +func (r *Runner) Run(ctx context.Context, options Options) Result { + result := Result{} + releases, err := r.PostgreSQL.Releases(ctx) + if err != nil { + for _, check := range options.Checks { + result.Summaries = append(result.Summaries, operationalSummary(check, options.Release, options.Date, err)) + } + result.Sort() + return result + } + releases = normalizeReleases(releases) + if options.Release != "" { + if !containsRelease(releases, options.Release) { + err := fmt.Errorf("release %q was not found in release definitions or historical Prow jobs", options.Release) + for _, check := range options.Checks { + result.Summaries = append(result.Summaries, operationalSummary(check, options.Release, options.Date, err)) + } + result.Sort() + return result + } + releases = []string{options.Release} + } + if len(releases) == 0 { + for _, check := range options.Checks { + result.Summaries = append(result.Summaries, operationalSummary(check, options.Release, options.Date, fmt.Errorf("no releases found"))) + } + result.Sort() + return result + } + + scope := Scope{Date: options.Date, Releases: releases} + for _, check := range options.Checks { + verifier := r.verifier(check) + if verifier == nil { + err := fmt.Errorf("unsupported verification check %q", check) + for _, release := range releases { + result.Summaries = append(result.Summaries, operationalSummary(check, release, options.Date, err)) + } + continue + } + checkResult := verifier.Verify(ctx, scope) + result.Summaries = append(result.Summaries, checkResult.Summaries...) + result.Discrepancies = append(result.Discrepancies, checkResult.Discrepancies...) + } + result.Sort() + return result +} + +func (r *Runner) verifier(check Check) Verifier { + switch check { + case CheckBQCompleteness: + return &BQCompletenessVerifier{ + PostgreSQL: r.PostgreSQL, + BigQuery: r.BigQuery, + BigQueryInitializationError: r.BigQueryInitializationError, + Config: r.Config, + SyntheticReleaseOverrides: r.SyntheticReleaseOverrides, + } + case CheckDailyTotals: + return &DailyTotalsVerifier{PostgreSQL: r.PostgreSQL} + case CheckCumulativeSummaries: + return &CumulativeSummariesVerifier{PostgreSQL: r.PostgreSQL} + default: + return nil + } +} + +func operationalSummary(check Check, release string, date civil.Date, err error) Summary { + return Summary{Check: check, Release: release, Date: date, Passed: false, Error: err.Error()} +} + +func normalizeReleases(releases []string) []string { + set := sets.New[string]() + for _, release := range releases { + if release = strings.TrimSpace(release); release != "" { + set.Insert(release) + } + } + return sets.List(set) +} + +func containsRelease(releases []string, wanted string) bool { + return sets.New[string](releases...).Has(wanted) +} diff --git a/pkg/db/verify/storage.go b/pkg/db/verify/storage.go new file mode 100644 index 0000000000..682602614d --- /dev/null +++ b/pkg/db/verify/storage.go @@ -0,0 +1,74 @@ +package verify + +import ( + "context" + "fmt" + "time" + + "cloud.google.com/go/civil" + + "github.com/openshift/sippy/pkg/db" +) + +type ReleaseReader interface { + Releases(context.Context) ([]string, error) +} + +type ProwJobRunReader interface { + ProwJobRunIDs(context.Context, time.Time, time.Time) (map[string]map[BuildID]struct{}, error) +} + +type DailyTotalsReader interface { + DailyRows(context.Context, string, civil.Date) ([]DailyRow, []DailyRow, error) +} + +type CumulativeSummariesReader interface { + CumulativeRows(context.Context, string, civil.Date) (CumulativeRows, error) +} + +type PostgreSQLReader interface { + ReleaseReader + ProwJobRunReader + DailyTotalsReader + CumulativeSummariesReader +} + +type BigQueryReader interface { + ProwJobs(context.Context, time.Time, time.Time) ([]BQJob, error) +} + +type PostgreSQL struct { + dbc *db.DB +} + +func NewPostgreSQL(dbc *db.DB) *PostgreSQL { + return &PostgreSQL{dbc: dbc} +} + +func (p *PostgreSQL) validate() error { + if p == nil || p.dbc == nil || p.dbc.DB == nil { + return fmt.Errorf("PostgreSQL client is not initialized") + } + return nil +} + +func (p *PostgreSQL) Releases(ctx context.Context) ([]string, error) { + if err := p.validate(); err != nil { + return nil, err + } + var releases []string + err := p.dbc.DB.WithContext(ctx).Raw(` + SELECT release + FROM release_definitions + WHERE deleted_at IS NULL AND release <> '' + UNION + SELECT DISTINCT release + FROM prow_jobs + WHERE deleted_at IS NULL AND release <> '' + ORDER BY release + `).Scan(&releases).Error + if err != nil { + return nil, fmt.Errorf("querying verification releases: %w", err) + } + return releases, nil +} diff --git a/pkg/db/verify/types.go b/pkg/db/verify/types.go new file mode 100644 index 0000000000..740dfe2594 --- /dev/null +++ b/pkg/db/verify/types.go @@ -0,0 +1,245 @@ +// Package verify provides read-only storage adapters and pure comparisons for +// checking the integrity of Sippy's daily data pipeline. +package verify + +import ( + "context" + "fmt" + "sort" + "strings" + + "cloud.google.com/go/civil" + log "github.com/sirupsen/logrus" + "k8s.io/apimachinery/pkg/util/sets" +) + +// Check identifies a daily data integrity check. +type Check string + +const ( + CheckBQCompleteness Check = "bq-completeness" + CheckDailyTotals Check = "daily-totals" + CheckCumulativeSummaries Check = "cumulative-summaries" +) + +// AllChecks lists supported checks in their canonical execution order. +var AllChecks = []Check{CheckBQCompleteness, CheckDailyTotals, CheckCumulativeSummaries} + +// Options selects the date, checks, and optional release for a run. +type Options struct { + Date civil.Date + Checks []Check + Release string +} + +// Scope is the normalized date and release set passed to each verifier. +type Scope struct { + Date civil.Date + Releases []string +} + +// Verifier runs one integrity check over the complete selected date and +// release scope. +type Verifier interface { + Verify(context.Context, Scope) Result +} + +func ParseChecks(values []string) ([]Check, error) { + if len(values) == 0 { + return append([]Check(nil), AllChecks...), nil + } + valid := sets.New[Check](AllChecks...) + selected := sets.New[Check]() + for _, value := range values { + check := Check(value) + if !valid.Has(check) { + allowed := make([]string, len(AllChecks)) + for i := range AllChecks { + allowed[i] = string(AllChecks[i]) + } + return nil, fmt.Errorf("invalid --check %q: must be one of %s", value, strings.Join(allowed, ", ")) + } + selected.Insert(check) + } + checks := make([]Check, 0, len(selected)) + for _, check := range AllChecks { + if selected.Has(check) { + checks = append(checks, check) + } + } + return checks, nil +} + +func ContainsCheck(checks []Check, wanted Check) bool { + return sets.New[Check](checks...).Has(wanted) +} + +// BuildID is the normalized unsigned Prow build identifier. +type BuildID uint64 + +// Summary is the bounded outcome record for one check, release, and date. +type Summary struct { + Check Check + Release string + Date civil.Date + Passed bool + ExpectedRows int + ActualRows int + Discrepancies int + Error string +} + +func (s Summary) Fields() log.Fields { + return log.Fields{ + "check": s.Check, + "release": s.Release, + "date": s.Date.String(), + "passed": s.Passed, + "expectedRows": s.ExpectedRows, + "actualRows": s.ActualRows, + "discrepancies": s.Discrepancies, + "error": s.Error, + } +} + +// Discrepancy describes one deterministic difference found by a check. +type Discrepancy struct { + Check Check + Release string + Date civil.Date + Kind string + Key string + Field string + Expected string + Actual string + Detail string +} + +func (d Discrepancy) Fields() log.Fields { + return log.Fields{ + "check": d.Check, + "release": d.Release, + "date": d.Date.String(), + "kind": d.Kind, + "key": d.Key, + "field": d.Field, + "expected": d.Expected, + "actual": d.Actual, + "detail": d.Detail, + } +} + +// Result contains all summary and discrepancy records from a run. +type Result struct { + Summaries []Summary + Discrepancies []Discrepancy +} + +func (r Result) Passed() bool { + if len(r.Summaries) == 0 { + return false + } + for _, summary := range r.Summaries { + if !summary.Passed { + return false + } + } + return true +} + +func (r *Result) Sort() { + sort.SliceStable(r.Summaries, func(i, j int) bool { + a, b := r.Summaries[i], r.Summaries[j] + if a.Check != b.Check { + return a.Check < b.Check + } + if a.Release != b.Release { + return a.Release < b.Release + } + return a.Date.Before(b.Date) + }) + sort.SliceStable(r.Discrepancies, func(i, j int) bool { + a, b := r.Discrepancies[i], r.Discrepancies[j] + av := []string{string(a.Check), a.Release, a.Date.String(), a.Kind, a.Key, a.Field, a.Expected, a.Actual, a.Detail} + bv := []string{string(b.Check), b.Release, b.Date.String(), b.Kind, b.Key, b.Field, b.Expected, b.Actual, b.Detail} + for k := range av { + if av[k] != bv[k] { + return av[k] < bv[k] + } + } + return false + }) +} + +func (r Result) Log(logger log.FieldLogger) { + for _, summary := range r.Summaries { + entry := logger.WithFields(summary.Fields()) + if summary.Passed { + entry.Info("verification summary") + } else { + entry.Error("verification summary") + } + } + for _, discrepancy := range r.Discrepancies { + logger.WithFields(discrepancy.Fields()).Error("verification discrepancy") + } +} + +// SummaryKey uniquely identifies a daily or cumulative test summary row. +type SummaryKey struct { + Release string + Date civil.Date + TestID uint64 + ProwJobID uint64 + SuiteID uint64 + Lifecycle string +} + +func (k SummaryKey) String() string { + return fmt.Sprintf("release=%q,date=%s,test_id=%d,prow_job_id=%d,suite_id=%d,lifecycle=%q", k.Release, k.Date, k.TestID, k.ProwJobID, k.SuiteID, k.Lifecycle) +} + +// Counts contains the status counters compared by summary checks. +type Counts struct { + Successes int64 + Failures int64 + Flakes int64 + Runs int64 +} + +func (c Counts) Add(other Counts) Counts { + return Counts{ + Successes: c.Successes + other.Successes, + Failures: c.Failures + other.Failures, + Flakes: c.Flakes + other.Flakes, + Runs: c.Runs + other.Runs, + } +} + +// DailyRow associates a summary key with its counters. +type DailyRow struct { + Key SummaryKey + Counts Counts +} + +// CumulativeRows contains the three inputs needed for cumulative validation. +type CumulativeRows struct { + // Previous contains cumulative rows from the day before the target date. + Previous []DailyRow + // Daily contains daily total rows from the target date. + Daily []DailyRow + // Target contains stored cumulative rows from the target date. + Target []DailyRow +} + +// BQJob contains the BigQuery Prow job fields needed for release attribution and ID comparison. +type BQJob struct { + // BuildID is the raw BigQuery build ID and may be blank or malformed. + BuildID string + // JobName is the Prow job name used for release attribution. + JobName string + // Annotations contains parsed Prow job annotations used for attribution. + Annotations map[string]string + // HasRefs reports whether the source job has pull request refs. + HasRefs bool +} diff --git a/pkg/db/verify/verify_test.go b/pkg/db/verify/verify_test.go new file mode 100644 index 0000000000..3b84b64756 --- /dev/null +++ b/pkg/db/verify/verify_test.go @@ -0,0 +1,373 @@ +package verify + +import ( + "context" + "errors" + "reflect" + "testing" + "time" + + "cloud.google.com/go/civil" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + v1config "github.com/openshift/sippy/pkg/apis/config/v1" + bqcachedclient "github.com/openshift/sippy/pkg/bigquery" +) + +var testDate = civil.Date{Year: 2026, Month: 8, Day: 25} + +func TestParseChecks(t *testing.T) { + tests := []struct { + name string + values []string + want []Check + wantErr string + }{ + {name: "default all", want: AllChecks}, + {name: "one", values: []string{"daily-totals"}, want: []Check{CheckDailyTotals}}, + {name: "repeatable canonical order and deduplication", values: []string{"cumulative-summaries", "bq-completeness", "cumulative-summaries"}, want: []Check{CheckBQCompleteness, CheckCumulativeSummaries}}, + {name: "invalid", values: []string{"unknown"}, wantErr: `invalid --check "unknown"`}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := ParseChecks(tt.values) + if tt.wantErr != "" { + require.ErrorContains(t, err, tt.wantErr) + return + } + require.NoError(t, err) + assert.Equal(t, tt.want, got) + }) + } +} + +func TestContainsCheck(t *testing.T) { + checks := []Check{CheckBQCompleteness, CheckCumulativeSummaries} + assert.True(t, ContainsCheck(checks, CheckBQCompleteness)) + assert.False(t, ContainsCheck(checks, CheckDailyTotals)) + assert.False(t, ContainsCheck(nil, CheckDailyTotals)) +} + +func TestCompareBuildIDs(t *testing.T) { + tests := []struct { + name string + bq map[BuildID]struct{} + pg map[BuildID]struct{} + malformed []string + wantKinds []string + }{ + {name: "equal", bq: idSet(1, 2), pg: idSet(1, 2)}, + {name: "both directions", bq: idSet(1), pg: idSet(2), wantKinds: []string{"missing-in-bigquery", "missing-in-postgres"}}, + {name: "malformed", malformed: []string{" bad "}, wantKinds: []string{"malformed-build-id"}}, + {name: "blank and whitespace normalize and deduplicate", malformed: []string{"", " ", "\t"}, wantKinds: []string{"missing-build-id"}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + summary, discrepancies := CompareBuildIDs("4.20", testDate, tt.bq, tt.pg, tt.malformed) + assert.Equal(t, len(tt.wantKinds) == 0, summary.Passed) + assert.Equal(t, len(tt.wantKinds), summary.Discrepancies) + kinds := make([]string, len(discrepancies)) + for i := range discrepancies { + kinds[i] = discrepancies[i].Kind + } + assert.ElementsMatch(t, tt.wantKinds, kinds) + if tt.name == "malformed" { + require.Len(t, discrepancies, 1) + assert.Equal(t, "bad", discrepancies[0].Key) + } + }) + } +} + +func TestCompareDaily(t *testing.T) { + key := SummaryKey{TestID: 1, ProwJobID: 2, SuiteID: 0, Lifecycle: "blocking"} + base := DailyRow{Key: key, Counts: Counts{Successes: 3, Failures: 2, Flakes: 1, Runs: 7}} + tests := []struct { + name string + expected []DailyRow + actual []DailyRow + wantKind string + wantField string + }{ + {name: "equal", expected: []DailyRow{base}, actual: []DailyRow{base}}, + {name: "raw only", expected: []DailyRow{base}, wantKind: "missing-row"}, + {name: "summary only", actual: []DailyRow{base}, wantKind: "unexpected-row"}, + {name: "success mismatch", expected: []DailyRow{base}, actual: []DailyRow{{Key: key, Counts: Counts{Successes: 4, Failures: 2, Flakes: 1, Runs: 7}}}, wantKind: "count-mismatch", wantField: "successes"}, + {name: "failure mismatch", expected: []DailyRow{base}, actual: []DailyRow{{Key: key, Counts: Counts{Successes: 3, Failures: 3, Flakes: 1, Runs: 7}}}, wantKind: "count-mismatch", wantField: "failures"}, + {name: "flake mismatch", expected: []DailyRow{base}, actual: []DailyRow{{Key: key, Counts: Counts{Successes: 3, Failures: 2, Flakes: 2, Runs: 7}}}, wantKind: "count-mismatch", wantField: "flakes"}, + {name: "run mismatch", expected: []DailyRow{base}, actual: []DailyRow{{Key: key, Counts: Counts{Successes: 3, Failures: 2, Flakes: 1, Runs: 8}}}, wantKind: "count-mismatch", wantField: "runs"}, + {name: "lifecycle is in key", expected: []DailyRow{base}, actual: []DailyRow{{Key: SummaryKey{TestID: 1, ProwJobID: 2, SuiteID: 0, Lifecycle: "informing"}, Counts: base.Counts}}, wantKind: "missing-row"}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + summary, discrepancies := CompareDaily("4.20", testDate, tt.expected, tt.actual) + if tt.wantKind == "" { + assert.True(t, summary.Passed) + assert.Empty(t, discrepancies) + return + } + assert.False(t, summary.Passed) + require.NotEmpty(t, discrepancies) + assert.Equal(t, tt.wantKind, discrepancies[0].Kind) + assert.Equal(t, tt.wantField, discrepancies[0].Field) + }) + } +} + +func TestCompareCumulative(t *testing.T) { + key := SummaryKey{TestID: 1, ProwJobID: 2, SuiteID: 3, Lifecycle: "blocking"} + previous := DailyRow{Key: key, Counts: Counts{Successes: 10, Failures: 2, Flakes: 1, Runs: 13}} + daily := DailyRow{Key: key, Counts: Counts{Successes: 2, Failures: 1, Runs: 3}} + tests := []struct { + name string + rows CumulativeRows + pass bool + }{ + {name: "first day", rows: CumulativeRows{Daily: []DailyRow{daily}, Target: []DailyRow{daily}}, pass: true}, + {name: "accumulation", rows: CumulativeRows{Previous: []DailyRow{previous}, Daily: []DailyRow{daily}, Target: []DailyRow{{Key: key, Counts: previous.Counts.Add(daily.Counts)}}}, pass: true}, + {name: "carry forward", rows: CumulativeRows{Previous: []DailyRow{previous}, Target: []DailyRow{previous}}, pass: true}, + {name: "broken carry forward", rows: CumulativeRows{Previous: []DailyRow{previous}}, pass: false}, + {name: "daily only missing target", rows: CumulativeRows{Daily: []DailyRow{daily}}, pass: false}, + {name: "unexpected target", rows: CumulativeRows{Target: []DailyRow{daily}}, pass: false}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + summary, discrepancies := CompareCumulative("4.20", testDate, tt.rows) + assert.Equal(t, tt.pass, summary.Passed) + assert.Equal(t, tt.pass, len(discrepancies) == 0) + }) + } +} + +func TestRunnerRunsAllSelectedChecksAfterFailures(t *testing.T) { + pg := &fakePostgreSQL{ + releases: []string{"pseudo", "4.20", "pseudo"}, + dailyErr: map[string]error{"4.20": errors.New("daily unavailable")}, + cumulative: map[string]CumulativeRows{}, + } + bq := &fakeBigQuery{jobs: []BQJob{{JobName: "job", BuildID: "bad"}}} + runner := Runner{ + PostgreSQL: pg, + BigQuery: bq, + Config: &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{ + "4.20": {Jobs: map[string]bool{"job": true}}, + }}, + } + result := runner.Run(context.Background(), Options{Date: testDate, Checks: AllChecks}) + assert.False(t, result.Passed()) + assert.Equal(t, 1, bq.calls, "BQ is queried once globally") + assert.Equal(t, 1, pg.runIDCalls, "PostgreSQL completeness is queried once globally") + assert.Equal(t, []string{"4.20", "pseudo"}, pg.dailyCalls) + assert.Equal(t, []string{"4.20", "pseudo"}, pg.cumulativeCalls) + assert.Len(t, result.Summaries, 6, "one summary per check and release") + assert.True(t, sortIsStable(result.Discrepancies)) +} + +func TestRunnerBQInitializationFailureDoesNotSkipPostgreSQLChecks(t *testing.T) { + pg := &fakePostgreSQL{releases: []string{"4.20"}, cumulative: map[string]CumulativeRows{}} + runner := Runner{PostgreSQL: pg, BigQueryInitializationError: errors.New("credentials missing")} + result := runner.Run(context.Background(), Options{Date: testDate, Checks: AllChecks}) + require.Len(t, result.Summaries, 3) + assert.Equal(t, "credentials missing", result.Summaries[0].Error) + assert.Equal(t, []string{"4.20"}, pg.dailyCalls) + assert.Equal(t, []string{"4.20"}, pg.cumulativeCalls) +} + +func TestRunnerDispatchesOnlySelectedChecks(t *testing.T) { + tests := []struct { + name string + check Check + wantBQCalls int + wantRunIDCalls int + wantDailyCalls []string + wantCumulativeCalls []string + }{ + { + name: "BigQuery completeness", + check: CheckBQCompleteness, wantBQCalls: 1, wantRunIDCalls: 1, + }, + { + name: "daily totals", + check: CheckDailyTotals, wantDailyCalls: []string{"4.20", "pseudo"}, + }, + { + name: "cumulative summaries", + check: CheckCumulativeSummaries, wantCumulativeCalls: []string{"4.20", "pseudo"}, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + pg := &fakePostgreSQL{releases: []string{"pseudo", "4.20"}} + bq := &fakeBigQuery{} + runner := Runner{PostgreSQL: pg, BigQuery: bq} + + result := runner.Run(context.Background(), Options{Date: testDate, Checks: []Check{tt.check}}) + + assert.True(t, result.Passed()) + assert.Len(t, result.Summaries, 2) + assert.Equal(t, tt.wantBQCalls, bq.calls) + assert.Equal(t, tt.wantRunIDCalls, pg.runIDCalls) + assert.Equal(t, tt.wantDailyCalls, pg.dailyCalls) + assert.Equal(t, tt.wantCumulativeCalls, pg.cumulativeCalls) + }) + } +} + +func TestRunnerNormalizesAndDeduplicatesBQBuildIDs(t *testing.T) { + pg := &fakePostgreSQL{ + releases: []string{"4.20"}, + runIDs: map[string]map[BuildID]struct{}{"4.20": idSet(1)}, + } + bq := &fakeBigQuery{jobs: []BQJob{ + {JobName: "job", BuildID: "1"}, + {JobName: "job", BuildID: "001"}, + {JobName: "job", BuildID: " bad "}, + {JobName: "job", BuildID: " bad "}, + }} + runner := Runner{ + PostgreSQL: pg, + BigQuery: bq, + Config: &v1config.SippyConfig{Releases: map[string]v1config.ReleaseConfig{ + "4.20": {Jobs: map[string]bool{"job": true}}, + }}, + } + result := runner.Run(context.Background(), Options{Date: testDate, Checks: []Check{CheckBQCompleteness}}) + require.Len(t, result.Summaries, 1) + assert.Equal(t, 1, result.Summaries[0].ExpectedRows) + assert.Equal(t, 1, result.Summaries[0].ActualRows) + assert.Equal(t, 1, result.Summaries[0].Discrepancies) + require.Len(t, result.Discrepancies, 1) + assert.Equal(t, "malformed-build-id", result.Discrepancies[0].Kind) + assert.Equal(t, "bad", result.Discrepancies[0].Key) + assert.Equal(t, time.Date(2026, 8, 25, 0, 0, 0, 0, time.UTC), bq.start) + assert.Equal(t, time.Date(2026, 8, 26, 0, 0, 0, 0, time.UTC), bq.end) +} + +func TestRunnerReleaseSelection(t *testing.T) { + pg := &fakePostgreSQL{releases: []string{"4.20", "pseudo"}} + runner := Runner{PostgreSQL: pg} + result := runner.Run(context.Background(), Options{Date: testDate, Checks: []Check{CheckDailyTotals}, Release: "pseudo"}) + assert.True(t, result.Passed()) + assert.Equal(t, []string{"pseudo"}, pg.dailyCalls) + + missing := runner.Run(context.Background(), Options{Date: testDate, Checks: []Check{CheckDailyTotals}, Release: "unknown"}) + assert.False(t, missing.Passed()) + require.Len(t, missing.Summaries, 1) + assert.ErrorContains(t, errors.New(missing.Summaries[0].Error), "was not found") +} + +func TestResultSortAndFields(t *testing.T) { + result := Result{ + Summaries: []Summary{ + {Check: CheckDailyTotals, Release: "z", Date: testDate, Passed: true}, + {Check: CheckBQCompleteness, Release: "a", Date: testDate, Passed: true}, + }, + Discrepancies: []Discrepancy{ + {Check: CheckDailyTotals, Release: "z", Date: testDate, Kind: "z"}, + {Check: CheckDailyTotals, Release: "a", Date: testDate, Kind: "a"}, + }, + } + result.Sort() + assert.Equal(t, CheckBQCompleteness, result.Summaries[0].Check) + assert.Equal(t, "a", result.Discrepancies[0].Release) + assert.Equal(t, "daily-totals", string(result.Summaries[1].Fields()["check"].(Check))) + assert.Contains(t, result.Discrepancies[0].Fields(), "expected") + assert.True(t, result.Passed()) + result.Summaries[0].Passed = false + assert.False(t, result.Passed()) +} + +func TestReadersRejectNilClients(t *testing.T) { + ctx := context.Background() + start := testDate.In(time.UTC) + end := testDate.AddDays(1).In(time.UTC) + + postgresReaders := []struct { + name string + call func(*PostgreSQL) error + }{ + {name: "releases", call: func(p *PostgreSQL) error { _, err := p.Releases(ctx); return err }}, + {name: "Prow job run IDs", call: func(p *PostgreSQL) error { _, err := p.ProwJobRunIDs(ctx, start, end); return err }}, + {name: "daily rows", call: func(p *PostgreSQL) error { _, _, err := p.DailyRows(ctx, "4.20", testDate); return err }}, + {name: "cumulative rows", call: func(p *PostgreSQL) error { _, err := p.CumulativeRows(ctx, "4.20", testDate); return err }}, + } + for _, tt := range postgresReaders { + t.Run(tt.name, func(t *testing.T) { + require.ErrorContains(t, tt.call(NewPostgreSQL(nil)), "PostgreSQL client is not initialized") + require.ErrorContains(t, tt.call(nil), "PostgreSQL client is not initialized") + }) + } + + bigQueryReaders := []*BigQuery{nil, NewBigQuery(nil), NewBigQuery(&bqcachedclient.Client{})} + for _, reader := range bigQueryReaders { + _, err := reader.ProwJobs(ctx, start, end) + require.ErrorContains(t, err, "BigQuery client is not initialized") + } +} + +func idSet(ids ...BuildID) map[BuildID]struct{} { + result := make(map[BuildID]struct{}, len(ids)) + for _, id := range ids { + result[id] = struct{}{} + } + return result +} + +func sortIsStable(discrepancies []Discrepancy) bool { + copyOf := append([]Discrepancy(nil), discrepancies...) + result := Result{Discrepancies: copyOf} + result.Sort() + return reflect.DeepEqual(discrepancies, result.Discrepancies) +} + +type fakePostgreSQL struct { + releases []string + releasesErr error + runIDs map[string]map[BuildID]struct{} + runIDCalls int + daily map[string][2][]DailyRow + dailyErr map[string]error + dailyCalls []string + cumulative map[string]CumulativeRows + cumulativeErr map[string]error + cumulativeCalls []string +} + +func (f *fakePostgreSQL) Releases(context.Context) ([]string, error) { + return f.releases, f.releasesErr +} + +func (f *fakePostgreSQL) ProwJobRunIDs(context.Context, time.Time, time.Time) (map[string]map[BuildID]struct{}, error) { + f.runIDCalls++ + return f.runIDs, nil +} + +func (f *fakePostgreSQL) DailyRows(_ context.Context, release string, _ civil.Date) ([]DailyRow, []DailyRow, error) { + f.dailyCalls = append(f.dailyCalls, release) + if err := f.dailyErr[release]; err != nil { + return nil, nil, err + } + rows := f.daily[release] + return rows[0], rows[1], nil +} + +func (f *fakePostgreSQL) CumulativeRows(_ context.Context, release string, _ civil.Date) (CumulativeRows, error) { + f.cumulativeCalls = append(f.cumulativeCalls, release) + return f.cumulative[release], f.cumulativeErr[release] +} + +type fakeBigQuery struct { + jobs []BQJob + err error + calls int + start time.Time + end time.Time +} + +func (f *fakeBigQuery) ProwJobs(_ context.Context, start, end time.Time) ([]BQJob, error) { + f.calls++ + f.start = start + f.end = end + return f.jobs, f.err +} diff --git a/test/integration/verify_test.go b/test/integration/verify_test.go new file mode 100644 index 0000000000..2fa5a42cba --- /dev/null +++ b/test/integration/verify_test.go @@ -0,0 +1,249 @@ +package integration + +import ( + "context" + "testing" + "time" + + "cloud.google.com/go/civil" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + processingv1 "github.com/openshift/sippy/pkg/apis/sippyprocessing/v1" + "github.com/openshift/sippy/pkg/db/models" + dbverify "github.com/openshift/sippy/pkg/db/verify" + intutil "github.com/openshift/sippy/test/integration/util" +) + +func TestVerifyReleaseEnumerationIncludesHistoricalPseudoReleaseWithoutTargetData(t *testing.T) { + dbc := intutil.NewTestDB(t, pgContainer) + intutil.CreateReleaseDefinition(t, dbc, "4.20", 4, 20) + intutil.CreateProwJob(t, dbc, "historical-job", "stale-pseudo", nil) + deleted := intutil.CreateProwJob(t, dbc, "deleted-job", "deleted-pseudo", nil) + require.NoError(t, dbc.DB.Delete(&deleted).Error) + intutil.CreateProwJob(t, dbc, "empty-release-job", "", nil) + + store := dbverify.NewPostgreSQL(dbc) + releases, err := store.Releases(context.Background()) + require.NoError(t, err) + assert.Equal(t, []string{"4.20", "stale-pseudo"}, releases) + + date := civil.Date{Year: 2026, Month: 8, Day: 25} + result := (&dbverify.Runner{PostgreSQL: store}).Run(context.Background(), dbverify.Options{ + Date: date, Checks: []dbverify.Check{dbverify.CheckDailyTotals, dbverify.CheckCumulativeSummaries}, Release: "stale-pseudo", + }) + require.Len(t, result.Summaries, 2) + assert.True(t, result.Passed()) + for _, summary := range result.Summaries { + assert.Equal(t, "stale-pseudo", summary.Release) + assert.Zero(t, summary.ExpectedRows) + assert.Zero(t, summary.ActualRows) + } +} + +func TestVerifyPostgreSQLProwStartHalfOpenUTCBoundaries(t *testing.T) { + dbc := intutil.NewTestDB(t, pgContainer) + job := intutil.CreateProwJob(t, dbc, "boundary-job", "4.20", nil) + date := civil.Date{Year: 2026, Month: 8, Day: 25} + start := date.In(time.UTC) + end := date.AddDays(1).In(time.UTC) + before := intutil.CreateProwJobRun(t, dbc, job.ID, "4.20", start.Add(-time.Nanosecond), true, processingv1.JobSucceeded) + atStart := intutil.CreateProwJobRun(t, dbc, job.ID, "4.20", start, true, processingv1.JobSucceeded) + atEnd := intutil.CreateProwJobRun(t, dbc, job.ID, "4.20", end, true, processingv1.JobSucceeded) + insideFromOffset := intutil.CreateProwJobRun(t, dbc, job.ID, "4.20", time.Date(2026, 8, 24, 20, 30, 0, 0, time.FixedZone("minus-four", -4*60*60)), true, processingv1.JobSucceeded) + + ids, err := dbverify.NewPostgreSQL(dbc).ProwJobRunIDs(context.Background(), start, end) + require.NoError(t, err) + assert.Contains(t, ids["4.20"], dbverify.BuildID(atStart.ID)) + assert.Contains(t, ids["4.20"], dbverify.BuildID(insideFromOffset.ID)) + assert.NotContains(t, ids["4.20"], dbverify.BuildID(before.ID)) + assert.NotContains(t, ids["4.20"], dbverify.BuildID(atEnd.ID)) +} + +func TestVerifyDailyRowsProductionSemanticsAndMismatches(t *testing.T) { + dbc := intutil.NewTestDB(t, pgContainer) + release := "pseudo-with-data" + date := civil.Date{Year: 2026, Month: 8, Day: 25} + start := date.In(time.UTC) + job := intutil.CreateProwJob(t, dbc, "daily-job", release, nil) + test := intutil.CreateTest(t, dbc, "daily test") + suite := intutil.CreateSuite(t, dbc, "suite") + + statuses := []int{ + int(processingv1.TestStatusSuccess), + int(processingv1.TestStatusFailure), + int(processingv1.TestStatusFlake), + int(processingv1.TestStatusRunning), + } + for i, status := range statuses { + timestamp := start.Add(time.Duration(i) * time.Hour) + run := intutil.CreateProwJobRun(t, dbc, job.ID, release, timestamp, status == int(processingv1.TestStatusSuccess), processingv1.JobSucceeded) + intutil.CreateProwJobRunTest(t, dbc, run.ID, job.ID, test.ID, release, timestamp, status) + } + + // A distinct suite/lifecycle proves they are both part of the key. + informingTime := start.Add(5 * time.Hour) + informingRun := intutil.CreateProwJobRun(t, dbc, job.ID, release, informingTime, true, processingv1.JobSucceeded) + intutil.CreateProwJobRunTest(t, dbc, informingRun.ID, job.ID, test.ID, release, informingTime, int(processingv1.TestStatusSuccess), intutil.WithSuiteID(suite.ID), intutil.WithLifecycle("informing")) + + // InfraFailure results and rows whose composite run key does not match are excluded. + infraTime := start.Add(6 * time.Hour) + infraRun := intutil.CreateProwJobRun(t, dbc, job.ID, release, infraTime, false, processingv1.JobInternalInfrastructureFailure, intutil.WithLabels("InfraFailure")) + intutil.CreateProwJobRunTest(t, dbc, infraRun.ID, job.ID, test.ID, release, infraTime, int(processingv1.TestStatusFailure)) + intutil.CreateProwJobRunTest(t, dbc, informingRun.ID, job.ID, test.ID, release, start.Add(7*time.Hour), int(processingv1.TestStatusFailure)) + + // The day is half-open: both joined rows are outside the target day. + beforeTime := start.Add(-time.Nanosecond) + beforeRun := intutil.CreateProwJobRun(t, dbc, job.ID, release, beforeTime, true, processingv1.JobSucceeded) + intutil.CreateProwJobRunTest(t, dbc, beforeRun.ID, job.ID, test.ID, release, beforeTime, int(processingv1.TestStatusSuccess)) + endTime := date.AddDays(1).In(time.UTC) + endRun := intutil.CreateProwJobRun(t, dbc, job.ID, release, endTime, true, processingv1.JobSucceeded) + intutil.CreateProwJobRunTest(t, dbc, endRun.ID, job.ID, test.ID, release, endTime, int(processingv1.TestStatusSuccess)) + + nullSuite := models.TestDailyTotal{ + Release: release, Date: date, TestID: test.ID, ProwJobID: job.ID, SuiteID: 0, Lifecycle: "blocking", + Successes: 1, Failures: 1, Flakes: 1, Runs: 4, + } + informing := models.TestDailyTotal{ + Release: release, Date: date, TestID: test.ID, ProwJobID: job.ID, SuiteID: suite.ID, Lifecycle: "informing", + Successes: 1, Runs: 1, + } + require.NoError(t, dbc.DB.Create(&nullSuite).Error) + require.NoError(t, dbc.DB.Create(&informing).Error) + + store := dbverify.NewPostgreSQL(dbc) + raw, stored, err := store.DailyRows(context.Background(), release, date) + require.NoError(t, err) + summary, discrepancies := dbverify.CompareDaily(release, date, raw, stored) + assert.True(t, summary.Passed) + assert.Empty(t, discrepancies) + require.Len(t, raw, 2) + assert.Equal(t, uint64(0), raw[0].Key.SuiteID, "null suite is normalized to zero") + runnerResult := (&dbverify.Runner{PostgreSQL: store}).Run(context.Background(), dbverify.Options{ + Date: date, Checks: []dbverify.Check{dbverify.CheckDailyTotals}, Release: release, + }) + require.Len(t, runnerResult.Summaries, 1) + assert.True(t, runnerResult.Passed(), "data-bearing pseudo-release is applicable") + + // All four counters are independently detected. + require.NoError(t, dbc.DB.Model(&nullSuite).Where( + "release = ? AND date = ? AND test_id = ? AND prow_job_id = ? AND suite_id = ? AND lifecycle = ?", + nullSuite.Release, nullSuite.Date, nullSuite.TestID, nullSuite.ProwJobID, nullSuite.SuiteID, nullSuite.Lifecycle, + ).Updates(map[string]any{ + "successes": 2, "failures": 2, "flakes": 2, "runs": 5, + }).Error) + _, stored, err = store.DailyRows(context.Background(), release, date) + require.NoError(t, err) + _, discrepancies = dbverify.CompareDaily(release, date, raw, stored) + fields := make([]string, 0) + for _, discrepancy := range discrepancies { + if discrepancy.Kind == "count-mismatch" { + fields = append(fields, discrepancy.Field) + } + } + assert.ElementsMatch(t, []string{"successes", "failures", "flakes", "runs"}, fields) + + // Summary-only and raw-only keys are compared in both directions. + extraSummary := models.TestDailyTotal{ + Release: release, Date: date, TestID: test.ID + 100, ProwJobID: job.ID, SuiteID: 0, Lifecycle: "blocking", Runs: 1, + } + require.NoError(t, dbc.DB.Create(&extraSummary).Error) + rawOnlyTest := intutil.CreateTest(t, dbc, "raw only") + rawOnlyTime := start.Add(8 * time.Hour) + rawOnlyRun := intutil.CreateProwJobRun(t, dbc, job.ID, release, rawOnlyTime, true, processingv1.JobSucceeded) + intutil.CreateProwJobRunTest(t, dbc, rawOnlyRun.ID, job.ID, rawOnlyTest.ID, release, rawOnlyTime, int(processingv1.TestStatusSuccess)) + raw, stored, err = store.DailyRows(context.Background(), release, date) + require.NoError(t, err) + _, discrepancies = dbverify.CompareDaily(release, date, raw, stored) + kinds := make([]string, len(discrepancies)) + for i := range discrepancies { + kinds[i] = discrepancies[i].Kind + } + assert.Contains(t, kinds, "unexpected-row") + assert.Contains(t, kinds, "missing-row") +} + +func TestVerifyCumulativeFirstDayAccumulationCarryForwardAndBreakage(t *testing.T) { + dbc := intutil.NewTestDB(t, pgContainer) + release := "4.20" + date := civil.Date{Year: 2026, Month: 8, Day: 25} + job := intutil.CreateProwJob(t, dbc, "cumulative-job", release, nil) + test := intutil.CreateTest(t, dbc, "cumulative test") + + createDaily := func(testID uint, lifecycle string, counts dbverify.Counts) { + require.NoError(t, dbc.DB.Create(&models.TestDailyTotal{ + Release: release, Date: date, TestID: testID, ProwJobID: job.ID, SuiteID: 0, Lifecycle: lifecycle, + Successes: fixtureInt32(counts.Successes), Failures: fixtureInt32(counts.Failures), + Flakes: fixtureInt32(counts.Flakes), Runs: fixtureInt32(counts.Runs), + }).Error) + } + createCumulative := func(targetDate civil.Date, testID uint, lifecycle string, counts dbverify.Counts) models.TestCumulativeSummary { + row := models.TestCumulativeSummary{ + Release: release, Date: targetDate, TestID: testID, ProwJobID: job.ID, SuiteID: 0, Lifecycle: lifecycle, + PrefixSumSuccesses: counts.Successes, PrefixSumFailures: counts.Failures, PrefixSumFlakes: counts.Flakes, PrefixSumRuns: counts.Runs, + } + require.NoError(t, dbc.DB.Create(&row).Error) + return row + } + + previousCounts := dbverify.Counts{Successes: 10, Failures: 2, Flakes: 1, Runs: 13} + dailyCounts := dbverify.Counts{Successes: 2, Failures: 1, Flakes: 1, Runs: 4} + createCumulative(date.AddDays(-1), test.ID, "accumulation", previousCounts) + createDaily(test.ID, "accumulation", dailyCounts) + createCumulative(date, test.ID, "accumulation", previousCounts.Add(dailyCounts)) + + carryTest := intutil.CreateTest(t, dbc, "carry forward") + carryCounts := dbverify.Counts{Successes: 5, Failures: 1, Runs: 6} + createCumulative(date.AddDays(-1), carryTest.ID, "blocking", carryCounts) + carryTarget := createCumulative(date, carryTest.ID, "blocking", carryCounts) + + firstTest := intutil.CreateTest(t, dbc, "first day") + firstCounts := dbverify.Counts{Successes: 1, Runs: 1} + createDaily(firstTest.ID, "blocking", firstCounts) + firstTarget := createCumulative(date, firstTest.ID, "blocking", firstCounts) + + store := dbverify.NewPostgreSQL(dbc) + rows, err := store.CumulativeRows(context.Background(), release, date) + require.NoError(t, err) + summary, discrepancies := dbverify.CompareCumulative(release, date, rows) + assert.True(t, summary.Passed) + assert.Empty(t, discrepancies) + + // A broken carry-forward is a concrete prefix counter mismatch. + require.NoError(t, dbc.DB.Model(&carryTarget).Where( + "release = ? AND date = ? AND test_id = ? AND prow_job_id = ? AND suite_id = ? AND lifecycle = ?", + carryTarget.Release, carryTarget.Date, carryTarget.TestID, carryTarget.ProwJobID, carryTarget.SuiteID, carryTarget.Lifecycle, + ).Update("prefix_sum_runs", carryCounts.Runs+1).Error) + rows, err = store.CumulativeRows(context.Background(), release, date) + require.NoError(t, err) + _, discrepancies = dbverify.CompareCumulative(release, date, rows) + require.Condition(t, func() bool { + for _, discrepancy := range discrepancies { + if discrepancy.Key == (dbverify.SummaryKey{Release: release, Date: date, TestID: uint64(carryTest.ID), ProwJobID: uint64(job.ID), Lifecycle: "blocking"}).String() && discrepancy.Field == "runs" { + return true + } + } + return false + }) + + // Removing the first-day target reports the daily-only key as missing. + require.NoError(t, dbc.DB.Where( + "release = ? AND date = ? AND test_id = ? AND prow_job_id = ? AND suite_id = ? AND lifecycle = ?", + firstTarget.Release, firstTarget.Date, firstTarget.TestID, firstTarget.ProwJobID, firstTarget.SuiteID, firstTarget.Lifecycle, + ).Delete(&models.TestCumulativeSummary{}).Error) + rows, err = store.CumulativeRows(context.Background(), release, date) + require.NoError(t, err) + _, discrepancies = dbverify.CompareCumulative(release, date, rows) + require.Condition(t, func() bool { + for _, discrepancy := range discrepancies { + if discrepancy.Kind == "missing-row" && discrepancy.Key == (dbverify.SummaryKey{Release: release, Date: date, TestID: uint64(firstTest.ID), ProwJobID: uint64(job.ID), Lifecycle: "blocking"}).String() { + return true + } + } + return false + }) +} + +func fixtureInt32(value int64) int32 { + return int32(value) //nolint:gosec // Test fixtures only use small non-negative counter constants. +}