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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
93 changes: 19 additions & 74 deletions pkg/api/componentreadiness/component_report.go
Original file line number Diff line number Diff line change
Expand Up @@ -327,104 +327,49 @@ func (c *ComponentReportGenerator) GenerateReport(ctx context.Context) (crtype.C
return crtype.ComponentReport{}, errs
}
report.GeneratedAt = componentReportTestStatus.GeneratedAt
log.Infof("GenerateReport completed in %s with %d sample results and %d base results from db", time.Since(before), sampleLen, len(componentReportTestStatus.BaseStatus))
log.WithField("duration", time.Since(before).String()).
WithField("sampleResults", sampleLen).
WithField("baseResults", len(componentReportTestStatus.BaseStatus)).
Info("GenerateReport completed")

return report, nil
}

// getTestStatus orchestrates the actual fetching of junit test run data for both basis and sample.
// goroutines are used to concurrently request the data for basis, sample, and various other edge cases.
func (c *ComponentReportGenerator) getTestStatus(ctx context.Context) (crstatus.ReportTestStatus, []error) {
before := time.Now()
fLog := log.WithField("func", "getTestStatus")

var baseStatus, sampleStatus map[string]crstatus.TestStatus
baseStatusCh := make(chan map[string]crstatus.TestStatus) // TODO: not hooked up yet, just in place for the interface for now
var baseErrs, sampleErrs []error
wg := &sync.WaitGroup{}

// channels for status as we may collect status from multiple queries run in separate goroutines
sampleStatusCh := make(chan map[string]crstatus.TestStatus)
errCh := make(chan error)
statusDoneCh := make(chan struct{}) // To signal when all processing is done
statusErrsDoneCh := make(chan struct{}) // To signal when all processing is done

// generate inputs to the channels
c.middlewares.Query(ctx, wg, baseStatusCh, sampleStatusCh, errCh)
goInterruptible(ctx, wg, func() { baseStatus, baseErrs = c.dataProvider.QueryBaseTestStatus(ctx, c.ReqOptions) })
goInterruptible(ctx, wg, func() {
fLog.Infof("running sample query with includeVariants: %+v", c.ReqOptions.VariantOption.IncludeVariants)
status, errs := c.dataProvider.QuerySampleTestStatus(ctx, c.ReqOptions, c.ReqOptions.VariantOption.IncludeVariants, c.ReqOptions.SampleRelease.Start, c.ReqOptions.SampleRelease.End)
fLog.Infof("received %d test statuses and %d errors from sample query", len(status), len(errs))
sampleStatusCh <- status
for _, err := range errs {

var baseStatus, sampleStatus map[string]crstatus.TestStatus
wg.Go(func() {
var queryErrs []error
baseStatus, sampleStatus, queryErrs = c.dataProvider.QueryTestStatus(ctx, c.ReqOptions)
for _, err := range queryErrs {
errCh <- err
}
})

// clean up channels after all queries are done
c.middlewares.Query(ctx, wg, errCh)

go func() {
wg.Wait()
close(baseStatusCh)
close(sampleStatusCh)
close(errCh)
}()

// manage output from the channels
go func() {
for status := range sampleStatusCh {
fLog.Infof("received %d test statuses over channel", len(status))
for k, v := range status {
if sampleStatus == nil {
fLog.Warnf("initializing sampleStatus map")
sampleStatus = make(map[string]crstatus.TestStatus)
}
if v2, ok := sampleStatus[k]; ok {
fLog.Warnf("sampleStatus already had key: %+v", k)
fLog.Warnf("sampleStatus new value: %+v", v)
fLog.Warnf("sampleStatus old value: %+v", v2)
}
sampleStatus[k] = v
}
}
close(statusDoneCh)
}()

go func() {
for err := range errCh {
sampleErrs = append(sampleErrs, err)
}
close(statusErrsDoneCh)
}()

<-statusDoneCh
<-statusErrsDoneCh
fLog.Infof("total test statuses: %d", len(sampleStatus))

var errs []error
if len(baseErrs) != 0 || len(sampleErrs) != 0 {
errs = append(errs, baseErrs...)
errs = append(errs, sampleErrs...)
for err := range errCh {
errs = append(errs, err)
}
log.Infof("getTestStatus completed in %s with %d sample results and %d base results",
time.Since(before), len(sampleStatus), len(baseStatus))

log.WithField("duration", time.Since(before)).
WithField("sampleResults", len(sampleStatus)).
WithField("baseResults", len(baseStatus)).
Info("getTestStatus completed")
now := time.Now()
return crstatus.ReportTestStatus{BaseStatus: baseStatus, SampleStatus: sampleStatus, GeneratedAt: &now}, errs
}

func goInterruptible(ctx context.Context, wg *sync.WaitGroup, closure func()) {
wg.Add(1)
go func() {
defer wg.Done()
select {
case <-ctx.Done():
return
default:
closure()
}
}()
}

var componentAndCapabilityGetter func(stats crstatus.TestStatus) (string, []string)

func testToComponentAndCapability(stats crstatus.TestStatus) (string, []string) {
Expand Down
27 changes: 23 additions & 4 deletions pkg/api/componentreadiness/dataprovider/bigquery/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"sort"
"strconv"
"strings"
"sync"
"time"

"cloud.google.com/go/bigquery"
Expand Down Expand Up @@ -67,15 +68,15 @@ func (p *BigQueryProvider) QueryBaseTestStatus(ctx context.Context, reqOptions r
return result.BaseStatus, nil
}

func (p *BigQueryProvider) QuerySampleTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions,
includeVariants map[string][]string,
start, end time.Time) (map[string]crstatus.TestStatus, []error) {
func (p *BigQueryProvider) querySampleTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions) (map[string]crstatus.TestStatus, []error) {
allJobVariants, errs := p.QueryJobVariants(ctx, reqOptions)
if len(errs) > 0 {
return nil, errs
}

generator := NewSampleQueryGenerator(p.client, reqOptions, allJobVariants, includeVariants, start, end)
generator := NewSampleQueryGenerator(p.client, reqOptions, allJobVariants,
reqOptions.VariantOption.IncludeVariants,
reqOptions.SampleRelease.Start, reqOptions.SampleRelease.End)
result, errs := apiPkg.GetDataFromCacheOrGenerate[crstatus.ReportTestStatus](
ctx, p.client.Cache, reqOptions.CacheOption,
apiPkg.NewCacheSpec(generator, "SampleTestStatus~", &reqOptions.SampleRelease.End),
Expand All @@ -86,6 +87,24 @@ func (p *BigQueryProvider) QuerySampleTestStatus(ctx context.Context, reqOptions
return result.SampleStatus, nil
}

func (p *BigQueryProvider) QueryTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions) (baseStatus, sampleStatus map[string]crstatus.TestStatus, errs []error) {
var baseErrs, sampleErrs []error
var wg sync.WaitGroup
wg.Go(func() {
baseStatus, baseErrs = p.QueryBaseTestStatus(ctx, reqOptions)
})
wg.Go(func() {
sampleStatus, sampleErrs = p.querySampleTestStatus(ctx, reqOptions)
})
wg.Wait()
errs = append(errs, baseErrs...)
errs = append(errs, sampleErrs...)
if len(errs) > 0 {
return nil, nil, errs
}
return baseStatus, sampleStatus, nil
}

// --- TestDetailsQuerier ---

func (p *BigQueryProvider) QueryBaseJobRunTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions) (map[string][]crstatus.TestDetailsSummary, []error) {
Expand Down
8 changes: 4 additions & 4 deletions pkg/api/componentreadiness/dataprovider/interface.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,10 @@ type TestStatusQuerier interface {
// QueryBaseTestStatus returns test status for the basis release.
QueryBaseTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions) (map[string]crstatus.TestStatus, []error)

// QuerySampleTestStatus returns test status for the sample release.
QuerySampleTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions,
includeVariants map[string][]string,
start, end time.Time) (map[string]crstatus.TestStatus, []error)
// QueryTestStatus returns both base and sample test status.
// Providers may execute this as a single optimized query or by
// delegating to QueryBaseTestStatus and a provider-internal sample query.
QueryTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions) (baseStatus, sampleStatus map[string]crstatus.TestStatus, errs []error)
}

// TestDetailsQuerier fetches per-job test breakdowns used for test details reports.
Expand Down
4 changes: 2 additions & 2 deletions pkg/api/componentreadiness/dataprovider/mixed/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,8 +64,8 @@ func (p *MixedProvider) QueryBaseTestStatus(ctx context.Context, reqOptions reqo
return p.providerFor(reqOptions).QueryBaseTestStatus(ctx, reqOptions)
}

func (p *MixedProvider) QuerySampleTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions, includeVariants map[string][]string, start, end time.Time) (map[string]crstatus.TestStatus, []error) {
return p.providerFor(reqOptions).QuerySampleTestStatus(ctx, reqOptions, includeVariants, start, end)
func (p *MixedProvider) QueryTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions) (baseStatus, sampleStatus map[string]crstatus.TestStatus, errs []error) {
return p.providerFor(reqOptions).QueryTestStatus(ctx, reqOptions)
}

func (p *MixedProvider) QueryBaseJobRunTestStatus(ctx context.Context, reqOptions reqopts.RequestOptions) (map[string][]crstatus.TestDetailsSummary, []error) {
Expand Down
Loading