From d25d926462fd59339ed0b6664377760093cef372 Mon Sep 17 00:00:00 2001 From: Andres Martinez Gotor Date: Fri, 25 Jul 2025 13:42:16 +0200 Subject: [PATCH] Advisor: Fix paginated requests for checks (#108583) --- .../pkg/app/checkscheduler/checkscheduler.go | 39 ++++++++++-- .../app/checkscheduler/checkscheduler_test.go | 63 +++++++++++++++++++ 2 files changed, 97 insertions(+), 5 deletions(-) diff --git a/apps/advisor/pkg/app/checkscheduler/checkscheduler.go b/apps/advisor/pkg/app/checkscheduler/checkscheduler.go index 9e382c5315f..92fe70fe192 100644 --- a/apps/advisor/pkg/app/checkscheduler/checkscheduler.go +++ b/apps/advisor/pkg/app/checkscheduler/checkscheduler.go @@ -102,6 +102,12 @@ func (r *Runner) Run(ctx context.Context) error { } else { lastCreated = time.Now() } + } else { + // Run an initial cleanup to remove old checks + err = r.cleanupChecks(ctxWithoutCancel, logger) + if err != nil { + logger.Error("Error cleaning up old check reports", "error", err) + } } } @@ -132,17 +138,35 @@ func (r *Runner) Run(ctx context.Context) error { } } +func (r *Runner) listChecks(ctx context.Context, logger logging.Logger) ([]resource.Object, error) { + list, err := r.client.List(ctx, r.namespace, resource.ListOptions{}) + if err != nil { + return nil, err + } + + checks := list.GetItems() + for list.GetContinue() != "" { + logger.Debug("List has continue token, listing next page", "continue", list.GetContinue()) + list, err = r.client.List(ctx, r.namespace, resource.ListOptions{Continue: list.GetContinue()}) + if err != nil { + return nil, err + } + checks = append(checks, list.GetItems()...) + } + return checks, nil +} + // checkLastCreated returns the creation time of the last check created // regardless of its ID. This assumes that the checks are created in batches // so a batch will have a similar creation time. // In case it finds an unprocessed check from a previous run, it will set it to error. func (r *Runner) checkLastCreated(ctx context.Context, log logging.Logger) (time.Time, error) { - list, err := r.client.List(ctx, r.namespace, resource.ListOptions{}) + checkList, err := r.listChecks(ctx, log) if err != nil { return time.Time{}, err } lastCreated := time.Time{} - for _, item := range list.GetItems() { + for _, item := range checkList { itemCreated := item.GetCreationTimestamp().Time if itemCreated.After(lastCreated) { lastCreated = itemCreated @@ -209,14 +233,16 @@ func (r *Runner) createChecks(ctx context.Context, logger logging.Logger) error // cleanupChecks deletes the olders checks if the number of checks exceeds the limit. func (r *Runner) cleanupChecks(ctx context.Context, logger logging.Logger) error { - list, err := r.client.List(ctx, r.namespace, resource.ListOptions{Limit: -1}) + checkList, err := r.listChecks(ctx, logger) if err != nil { return err } + logger.Debug("Cleaning up checks", "numChecks", len(checkList)) + // organize checks by type checksByType := map[string][]resource.Object{} - for _, check := range list.GetItems() { + for _, check := range checkList { labels := check.GetLabels() checkType, ok := labels[checks.TypeLabel] if !ok { @@ -226,8 +252,10 @@ func (r *Runner) cleanupChecks(ctx context.Context, logger logging.Logger) error checksByType[checkType] = append(checksByType[checkType], check) } - for _, checks := range checksByType { + for checkType, checks := range checksByType { + logger.Debug("Checking checks", "checkType", checkType, "numChecks", len(checks)) if len(checks) > r.maxHistory { + logger.Debug("Deleting old checks", "checkType", checkType, "maxHistory", r.maxHistory, "numChecks", len(checks)) // Sort checks by creation time sort.Slice(checks, func(i, j int) bool { ti := checks[i].GetCreationTimestamp().Time @@ -242,6 +270,7 @@ func (r *Runner) cleanupChecks(ctx context.Context, logger logging.Logger) error if err != nil { return fmt.Errorf("error deleting check: %w", err) } + logger.Debug("Deleted check", "check", check.GetStaticMetadata().Identifier()) } } } diff --git a/apps/advisor/pkg/app/checkscheduler/checkscheduler_test.go b/apps/advisor/pkg/app/checkscheduler/checkscheduler_test.go index 69b4e809b4c..de522be7adc 100644 --- a/apps/advisor/pkg/app/checkscheduler/checkscheduler_test.go +++ b/apps/advisor/pkg/app/checkscheduler/checkscheduler_test.go @@ -103,6 +103,69 @@ func TestRunner_checkLastCreated_UnprocessedCheck(t *testing.T) { assert.Equal(t, expectedAnnotations, patchOperation.Value) } +func TestRunner_checkLastCreated_PaginatedResponse(t *testing.T) { + // Create checks with different creation times + past := time.Now().Add(-1 * time.Hour) + now := time.Now() + + mockClient := &MockClient{ + listFunc: func(ctx context.Context, namespace string, options resource.ListOptions) (resource.ListObject, error) { + if options.Continue == "" { + // First page - return oldest and middle checks with continue token + return &advisorv0alpha1.CheckList{ + ListMeta: metav1.ListMeta{ + Continue: "continue-token-123", + }, + Items: []advisorv0alpha1.Check{ + { + ObjectMeta: metav1.ObjectMeta{ + Name: "check-1", + CreationTimestamp: metav1.NewTime(past), + Annotations: map[string]string{ + checks.StatusAnnotation: "completed", + }, + }, + }, + { + ObjectMeta: metav1.ObjectMeta{ + Name: "check-2", + CreationTimestamp: metav1.NewTime(past), + Annotations: map[string]string{ + checks.StatusAnnotation: "completed", + }, + }, + }, + }, + }, nil + } + // Second page - verify continue token is passed and return newest check + assert.Equal(t, "continue-token-123", options.Continue) + return &advisorv0alpha1.CheckList{ + Items: []advisorv0alpha1.Check{ + { + ObjectMeta: metav1.ObjectMeta{ + Name: "check-3", + CreationTimestamp: metav1.NewTime(now), + Annotations: map[string]string{ + checks.StatusAnnotation: "completed", + }, + }, + }, + }, + }, nil + }, + } + + runner := &Runner{ + client: mockClient, + log: &logging.NoOpLogger{}, + } + + lastCreated, err := runner.checkLastCreated(context.Background(), &logging.NoOpLogger{}) + assert.NoError(t, err) + assert.Equal(t, now.Truncate(time.Second), lastCreated.Truncate(time.Second)) +} + func TestRunner_createChecks_ErrorOnCreate(t *testing.T) { mockCheckService := &MockCheckService{checks: []checks.Check{&mockCheck{id: "check-1"}}}