diff --git a/pkg/registry/apis/provisioning/jobs/driver.go b/pkg/registry/apis/provisioning/jobs/driver.go index 407b95bdce3..a3adaeb5550 100644 --- a/pkg/registry/apis/provisioning/jobs/driver.go +++ b/pkg/registry/apis/provisioning/jobs/driver.go @@ -3,6 +3,7 @@ package jobs import ( "context" "errors" + "fmt" "time" "k8s.io/apiserver/pkg/endpoints/request" @@ -48,10 +49,12 @@ var _ Store = (*persistentStore)(nil) // There may be multiple jobDrivers running in parallel. // The jobDriver deals with cleaning up upon death and ensuring that jobs remain claimable. type jobDriver struct { - // Timeout for processing a job. This should be the same or less than a claim expiry. - timeout time.Duration + // Timeout for processing a job. This must be less than a claim expiry. + jobTimeout time.Duration + // CleanupInterval is the time between cleanup runs. cleanupInterval time.Duration + // JobInterval is the time between job ticks. This should be relatively low. jobInterval time.Duration @@ -69,21 +72,25 @@ type jobDriver struct { } func NewJobDriver( - timeout, cleanupInterval, jobInterval time.Duration, + jobTimeout, cleanupInterval, jobInterval time.Duration, store Store, repoGetter RepoGetter, historicJobs History, workers ...Worker, -) *jobDriver { +) (*jobDriver, error) { + if cleanupInterval < jobTimeout { + return nil, fmt.Errorf("the cleanup interval must be larger than the jobTimeout (cleanup:%s < job:%s)", + cleanupInterval.String(), jobTimeout.String()) + } return &jobDriver{ - timeout: timeout, + jobTimeout: jobTimeout, cleanupInterval: cleanupInterval, jobInterval: jobInterval, store: store, repoGetter: repoGetter, historicJobs: historicJobs, workers: workers, - } + }, nil } // Run drives jobs to completion. This is a blocking function. @@ -104,8 +111,13 @@ func (d *jobDriver) Run(ctx context.Context) { panic("unreachable?: failed to grant provisioning identity: " + err.Error()) } + // Remove old jobs + if err = d.store.Cleanup(ctx); err != nil { + logger.Error("failed to clean up old jobs at start", "error", err) + } + // Drive without waiting on startup. - d.startDriving(ctx) + d.processJobsUntilDoneOrError(ctx) for { select { @@ -113,28 +125,30 @@ func (d *jobDriver) Run(ctx context.Context) { if err := d.store.Cleanup(ctx); err != nil { logger.Error("failed to cleanup jobs", "error", err) } + + // These events do not queue if the worker is already running case <-jobTicker.C: - d.startDriving(ctx) + d.processJobsUntilDoneOrError(ctx) case <-d.store.InsertNotifications(): - d.startDriving(ctx) + d.processJobsUntilDoneOrError(ctx) } } } -func (d *jobDriver) startDriving(ctx context.Context) { - timeoutCtx, cancel := context.WithTimeout(ctx, d.timeout) - defer cancel() - for timeoutCtx.Err() == nil { - if err := d.drive(timeoutCtx); err != nil { - if !errors.Is(err, context.Canceled) && !errors.Is(err, ErrNoJobs) { +// This will keep processing jobs until there are none left (or we hit an error) +func (d *jobDriver) processJobsUntilDoneOrError(ctx context.Context) { + for { + err := d.claimAndProcessOneJob(ctx) + if err != nil { + if !errors.Is(err, ErrNoJobs) { logging.FromContext(ctx).Error("failed to drive jobs", "error", err) } - break + return } } } -func (d *jobDriver) drive(ctx context.Context) error { +func (d *jobDriver) claimAndProcessOneJob(ctx context.Context) error { logger := logging.FromContext(ctx) // Claim a job to work on. @@ -158,13 +172,21 @@ func (d *jobDriver) drive(ctx context.Context) error { return apifmt.Errorf("failed to grant provisioning identity: %w", err) } + jobctx, cancel := context.WithTimeout(ctx, d.jobTimeout) + defer cancel() // Ensure resources are released when the function returns + // Process the job. start := time.Now() job.Status.Started = start.UnixMilli() - err = d.processJob(ctx, job) // NOTE: We pass in a pointer here such that the job status can be kept in Complete without re-fetching. + err = d.processJob(jobctx, job) // NOTE: We pass in a pointer here such that the job status can be kept in Complete without re-fetching. end := time.Now() logger.Debug("job processed", "duration", end.Sub(start), "error", err) + // Capture job timeout + if jobctx.Err() != nil && err == nil { + err = jobctx.Err() + } + // Mark the job as failed and remove from queue if err != nil { job.Status.State = provisioning.JobStateError diff --git a/pkg/registry/apis/provisioning/jobs/persistentstore.go b/pkg/registry/apis/provisioning/jobs/persistentstore.go index 04fac903be0..abfdf2cecc7 100644 --- a/pkg/registry/apis/provisioning/jobs/persistentstore.go +++ b/pkg/registry/apis/provisioning/jobs/persistentstore.go @@ -83,6 +83,7 @@ type persistentStore struct { // clock is a function that returns the current time. clock func() time.Time + // expiry is the time after which a job is considered abandoned. // If a job is abandoned, it will have its claim cleaned up periodically. expiry time.Duration @@ -93,7 +94,7 @@ type persistentStore struct { notifications chan struct{} } -func NewStore( +func NewJobStore( jobStore jobStorage, expiry time.Duration, ) (*persistentStore, error) { @@ -204,6 +205,7 @@ func (s *persistentStore) Claim(ctx context.Context) (job *provisioning.Job, rol } // Rollback the claim. + refetchedJob = refetchedJob.DeepCopy() delete(refetchedJob.Labels, LabelJobClaim) refetchedJob.Status.State = provisioning.JobStatePending diff --git a/pkg/registry/apis/provisioning/register.go b/pkg/registry/apis/provisioning/register.go index 2f6ed45cba3..a65f2008813 100644 --- a/pkg/registry/apis/provisioning/register.go +++ b/pkg/registry/apis/provisioning/register.go @@ -330,7 +330,7 @@ func (b *APIBuilder) UpdateAPIGroupInfo(apiGroupInfo *genericapiserver.APIGroupI return fmt.Errorf("failed to create job storage: %w", err) } - b.jobs, err = jobs.NewStore(realJobStore, time.Second*30) + b.jobs, err = jobs.NewJobStore(realJobStore, 30*time.Second) // FIXME: this timeout if err != nil { return fmt.Errorf("failed to create job store: %w", err) } @@ -583,8 +583,15 @@ func (b *APIBuilder) GetPostStartHooks() (map[string]genericapiserver.PostStartH commenter := pullrequest.NewCommenter() pullRequestWorker := pullrequest.NewPullRequestWorker(evaluator, commenter) - driver := jobs.NewJobDriver(time.Second*28, time.Second*30, time.Second*30, b.jobs, b, b.jobHistory, + driver, err := jobs.NewJobDriver( + time.Minute*20, // Max time for each job + time.Minute*22, // Cleanup any checked out jobs. FIXME: this is slow if things crash/fail! + time.Second*30, // Periodically look for new jobs + b.jobs, b, b.jobHistory, exportWorker, syncWorker, migrationWorker, pullRequestWorker) + if err != nil { + return err + } go driver.Run(postStartHookCtx.Context) repoController, err := controller.NewRepositoryController(