From a037e0aae1aa33023262a617b9bc84b79b5ea88c Mon Sep 17 00:00:00 2001 From: Matthew Staebler Date: Wed, 15 Jul 2026 22:58:15 -0400 Subject: [PATCH] Fix sippy_prow_jobs_loaded metric to report total queried count The bulk writer refactor changed the fetch loop to only iterate over new job runs, which caused jobsImportedCount to drop from the full BigQuery result set to just the filtered new entries. Set the count once from len(prowJobs) to restore the original metric behavior. Co-Authored-By: Claude Opus 4.6 --- pkg/dataloader/prowloader/prow.go | 12 +++--------- 1 file changed, 3 insertions(+), 9 deletions(-) diff --git a/pkg/dataloader/prowloader/prow.go b/pkg/dataloader/prowloader/prow.go index d89b46a301..b9d6d634ef 100644 --- a/pkg/dataloader/prowloader/prow.go +++ b/pkg/dataloader/prowloader/prow.go @@ -12,7 +12,6 @@ import ( "regexp" "strconv" "sync" - "sync/atomic" "time" "cloud.google.com/go/bigquery" @@ -72,8 +71,6 @@ type ProwLoader struct { releaseRegexps map[string][]*regexp.Regexp config *v1config.SippyConfig ghCommenter *commenter.GitHubCommenter - jobsImportedCount atomic.Int32 - jobsProcessedCount atomic.Int32 gcsClient *storage.Client promPusher *push.Pusher loadSince *time.Time @@ -259,6 +256,8 @@ func (pl *ProwLoader) Load() { pl.errors = append(pl.errors, errors.Wrap(err, "error pre-fetching labels from BigQuery")) } + prowLoaderQueriedMetricGauge.Set(float64(len(prowJobs))) + // Match jobs to releases and bulk-upsert ProwJob definitions before // the concurrent processing loop. The prowJobCache is read-only after // this point. @@ -274,7 +273,6 @@ func (pl *ProwLoader) Load() { queue := make(chan *prow.ProwJob) results := make(chan *jobRunResult, len(entries)) fetchErrsCh := make(chan error, len(entries)) - total := len(entries) go func() { defer close(queue) @@ -306,8 +304,6 @@ func (pl *ProwLoader) Load() { if result != nil { results <- result } - pl.jobsImportedCount.Add(1) - log.Infof("%d of %d job runs processed", pl.jobsImportedCount.Load(), total) } }(fetchCtx) } @@ -339,9 +335,7 @@ func (pl *ProwLoader) Load() { log.Infof("finished importing new job runs in %+v", time.Since(start)) if pl.promPusher != nil { - prowLoaderQueriedMetricGauge.Set(float64(pl.jobsImportedCount.Load())) pl.promPusher.Collector(prowLoaderQueriedMetricGauge) - prowLoaderProcessedMetricGauge.Set(float64(pl.jobsProcessedCount.Load())) pl.promPusher.Collector(prowLoaderProcessedMetricGauge) } } @@ -1037,7 +1031,6 @@ func (pl *ProwLoader) fetchJobRunResult(ctx context.Context, pj *prow.ProwJob) ( return nil, err } - pl.jobsProcessedCount.Add(1) return result, nil } @@ -1240,6 +1233,7 @@ func (pl *ProwLoader) accumulateAndWriteJobRuns(ctx context.Context, results <-c if total > 0 { log.WithField("runs", total).Info("all job run batches committed") } + prowLoaderProcessedMetricGauge.Set(float64(total)) return nil }