package jobs import ( "context" "errors" "strings" "sync" "time" "go.opentelemetry.io/otel/attribute" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apiserver/pkg/endpoints/request" "github.com/grafana/grafana-app-sdk/logging" "github.com/grafana/grafana/apps/provisioning/pkg/apifmt" provisioning "github.com/grafana/grafana/apps/provisioning/pkg/apis/provisioning/v0alpha1" "github.com/grafana/grafana/pkg/apimachinery/identity" "github.com/grafana/grafana/pkg/infra/tracing" ) // Store is an abstraction for the storage API. // This exists to allow for unit testing. // //go:generate mockery --name Store --structname MockStore --inpackage --filename store_mock.go --with-expecter type Store interface { // Claim takes a job from storage, marks it as ours, and returns it. // // Any job which has not been claimed by another worker is fair game. // // If err is not nil, the job and rollback values are always nil. // The err may be ErrNoJobs if there are no jobs to claim. Claim(ctx context.Context) (job *provisioning.Job, rollback func(), err error) // Complete marks a job as completed and removes it from the active job store. // Callers are responsible for writing the job to history after calling this. Complete(ctx context.Context, job *provisioning.Job) error // Update saves the job back to the store. Update(ctx context.Context, job *provisioning.Job) (*provisioning.Job, error) // RenewLease renews the lease for a claimed job, extending its expiry time. // Returns an error if the lease cannot be renewed (e.g., job was completed or lease expired). RenewLease(ctx context.Context, job *provisioning.Job) error // Get retrieves a job by name for conflict resolution. Get(ctx context.Context, namespace, name string) (*provisioning.Job, error) // ListExpiredJobs lists jobs with expired leases (claim timestamp older than the given time). // Returns jobs in batches up to the specified limit. ListExpiredJobs(ctx context.Context, expiredBefore time.Time, limit int) ([]*provisioning.Job, error) } // jobDriver drives jobs to completion and manages the job queue. // There may be multiple jobDrivers running in parallel. // The jobDriver processes jobs but does not handle cleanup - that's handled by ConcurrentJobDriver. type jobDriver struct { // Timeout for processing a job. This must be less than a claim expiry. jobTimeout time.Duration // JobInterval is the time between job ticks. This should be relatively low. jobInterval time.Duration // LeaseRenewalInterval is how often to renew job leases. leaseRenewalInterval time.Duration // Store is the job storage backend. store Store // RepoGetter lets us access repositories to pass to the worker. repoGetter RepoGetter // save info about finished jobs historicJobs HistoryWriter // Workers process the job. // Only the first worker who supports the job will process it; the rest are ignored. workers []Worker // notifications channel for job create events notifications chan struct{} // Mutex to protect concurrent access to job processing mu sync.Mutex // currentJob is the job currently being processed currentJob *provisioning.Job } func NewJobDriver( jobTimeout, jobInterval, leaseRenewalInterval time.Duration, store Store, repoGetter RepoGetter, historicJobs HistoryWriter, notifications chan struct{}, workers ...Worker, ) (*jobDriver, error) { return &jobDriver{ jobTimeout: jobTimeout, jobInterval: jobInterval, leaseRenewalInterval: leaseRenewalInterval, store: store, repoGetter: repoGetter, historicJobs: historicJobs, workers: workers, notifications: notifications, }, nil } // Run drives jobs to completion. This is a blocking function. // It will run until the context is canceled or an error occurs. // This is a thread-safe function; it may be called from multiple goroutines. // // Note: This function intentionally does NOT create a tracing span because it runs indefinitely // until shutdown. Individual job processing operations already have their own spans. func (d *jobDriver) Run(ctx context.Context) error { jobTicker := time.NewTicker(d.jobInterval) defer jobTicker.Stop() logger := logging.FromContext(ctx).With("logger", "job-driver") ctx = logging.Context(ctx, logger) ctx, _, err := identity.WithProvisioningIdentity(ctx, "*") // "*" grants us access to all namespaces. if err != nil { return apifmt.Errorf("failed to grant provisioning identity: %w", err) } // Drive without waiting on startup. d.processJobsUntilDoneOrError(ctx) for { select { case <-ctx.Done(): logger.Info("job driver stopped") return nil // Context cancellation is expected during shutdown case <-jobTicker.C: d.processJobsUntilDoneOrError(ctx) case <-d.notifications: d.processJobsUntilDoneOrError(ctx) } } } // This will keep processing jobs until there are none left (or we hit an error) func (d *jobDriver) processJobsUntilDoneOrError(ctx context.Context) { for { // Check if context is cancelled before attempting to claim jobs if ctx.Err() != nil { return } err := d.claimAndProcessOneJob(ctx) if err != nil { if !errors.Is(err, ErrNoJobs) { logging.FromContext(ctx).Error("failed to drive jobs", "error", err) } return } } } func (d *jobDriver) claimAndProcessOneJob(ctx context.Context) error { ctx, span := tracing.Start(ctx, "provisioning.jobs.claim_and_process_one_job") defer span.End() logger := logging.FromContext(ctx) // Claim a job to work on. claimedJob, rollback, err := d.store.Claim(ctx) if err != nil { if !errors.Is(err, ErrNoJobs) { span.RecordError(err) } return apifmt.Errorf("failed to claim job: %w", err) } // Ensure that the job is cleaned up if we fail to complete it. // The rollback function does not care about cancellations. defer rollback() namespace := claimedJob.GetNamespace() logger = logger.With("job", claimedJob.GetName(), "namespace", namespace) ctx = logging.Context(ctx, logger) d.currentJob = claimedJob span.SetAttributes( attribute.String("job.name", claimedJob.GetName()), attribute.String("job.namespace", namespace), attribute.String("job.repository", claimedJob.Spec.Repository), attribute.String("job.action", string(claimedJob.Spec.Action)), ) // Now that we have a job, we need to augment our namespace to grant ourselves permission to work on it. // Incidentally, this also limits our permissions to only the namespace of the job. ctx = request.WithNamespace(ctx, namespace) ctx, _, err = identity.WithProvisioningIdentity(ctx, namespace) if err != nil { 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 // Set up lease renewal goroutine leaseRenewalCtx, cancelLeaseRenewal := context.WithCancel(jobctx) leaseExpired := make(chan struct{}) go d.leaseRenewalLoop(leaseRenewalCtx, logger, leaseExpired) defer cancelLeaseRenewal() recorder := newJobProgressRecorder(d.onProgress()) recorder.SetMessage(ctx, "start job") // Process the job with lease loss detection err = d.processJobWithLeaseCheck(jobctx, recorder, leaseExpired) end := time.Now() logger.Debug("job processed", "duration", end.Sub(recorder.Started()), "error", err) // Check if parent context was cancelled (graceful shutdown) if ctx.Err() != nil { logger.Debug("context cancel - job will retry") // Don't complete the job - let it be retried by another worker d.mu.Lock() d.currentJob = nil d.mu.Unlock() return nil } // Capture job timeout (but not parent context cancellation) if jobctx.Err() != nil && err == nil && ctx.Err() == nil { err = jobctx.Err() } // Record job processing error on span if err != nil { span.RecordError(err) } // Complete the job d.mu.Lock() d.currentJob.Status = recorder.Complete(ctx, err) defer func() { d.currentJob = nil d.mu.Unlock() }() // Save the finished job err = d.historicJobs.WriteJob(ctx, d.currentJob.DeepCopy()) if err != nil { // We're not going to return this as it is not critical. Not ideal, but not critical. logger.Warn("failed to write historic job", "error", err) } // Mark the job as completed. if err := d.store.Complete(ctx, d.currentJob); err != nil { span.RecordError(err) return apifmt.Errorf("failed to complete job '%s' in '%s': %w", d.currentJob.GetName(), d.currentJob.GetNamespace(), err) } logger.Info("job complete") return nil } // leaseRenewalLoop continuously renews the lease for a job until the context is cancelled. // If lease renewal fails persistently, it signals via the leaseExpired channel. // // Note: This function intentionally does NOT create a tracing span because it runs indefinitely // for the lifetime of a job. Individual RenewLease calls already have their own spans. func (d *jobDriver) leaseRenewalLoop(ctx context.Context, logger logging.Logger, leaseExpired chan struct{}) { ticker := time.NewTicker(d.leaseRenewalInterval) defer ticker.Stop() logger.Debug("start lease renewal loop", "renewal_interval", d.leaseRenewalInterval) consecutiveFailures := 0 maxFailures := 3 // Allow a few failures before giving up for { select { case <-ctx.Done(): logger.Debug("lease renewal loop stopped") return case <-ticker.C: d.mu.Lock() if d.currentJob == nil { d.mu.Unlock() return } err := d.store.RenewLease(ctx, d.currentJob) d.mu.Unlock() if err != nil { consecutiveFailures++ if apierrors.IsNotFound(err) || strings.Contains(err.Error(), "job no longer exists") { logger.Error("job no longer exists - lease expired", "error", err) close(leaseExpired) return } logger.Warn("failed to renew lease", "error", err, "consecutive_failures", consecutiveFailures) if consecutiveFailures >= maxFailures { logger.Error("too many consecutive lease renewal failures - job will be aborted", "consecutive_failures", consecutiveFailures, "max_failures", maxFailures) close(leaseExpired) return } } else { if consecutiveFailures > 0 { logger.Debug("lease renewal recovered", "previous_failures", consecutiveFailures) } consecutiveFailures = 0 } } } } // processJobWithLeaseCheck processes a job but aborts if the lease expires or context is cancelled. func (d *jobDriver) processJobWithLeaseCheck(ctx context.Context, recorder JobProgressRecorder, leaseExpired <-chan struct{}) error { // Run the job processing in a goroutine so we can monitor lease expiry resultChan := make(chan error, 1) go func() { resultChan <- d.processJob(ctx, recorder) }() select { case err := <-resultChan: return err case <-leaseExpired: return apifmt.Errorf("job aborted due to lease expiry") case <-ctx.Done(): // Return context error directly - caller will determine if this is due to graceful shutdown // or job timeout based on which context was cancelled return ctx.Err() } } func (d *jobDriver) processJob(ctx context.Context, recorder JobProgressRecorder) error { ctx, span := tracing.Start(ctx, "provisioning.jobs.process_job") defer span.End() logger := logging.FromContext(ctx) d.mu.Lock() if d.currentJob == nil { d.mu.Unlock() return nil } // Here it's safe to copy as only job spec is used for processing job := d.currentJob.DeepCopy() repoName := d.currentJob.Spec.Repository namespace := d.currentJob.Namespace d.mu.Unlock() span.SetAttributes( attribute.String("job.repository", repoName), attribute.String("job.action", string(job.Spec.Action)), ) for _, worker := range d.workers { if !worker.IsSupported(ctx, *job) { continue } repo, err := d.repoGetter.GetRepository(ctx, namespace, repoName) if err != nil { span.RecordError(err) return apifmt.Errorf("failed to get repository '%s': %w", repoName, err) } r := repo.Config() if r.DeletionTimestamp != nil && !r.DeletionTimestamp.IsZero() { logger.Info("repository marked for deletion - skip job", "name", r.Name, "namespace", r.Namespace, "deletionTimestamp", r.DeletionTimestamp, ) return nil } err = worker.Process(ctx, repo, *job, recorder) if err != nil { span.RecordError(err) } return err } err := apifmt.Errorf("no workers were registered to handle the job") span.RecordError(err) return err } func (d *jobDriver) onProgress() ProgressFn { return func(ctx context.Context, status provisioning.JobStatus) error { ctx, span := tracing.Start(ctx, "provisioning.jobs.update_progress") defer span.End() logging.FromContext(ctx).Debug("job progress", "status", status) const maxRetries = 3 for attempt := 0; attempt < maxRetries; attempt++ { d.mu.Lock() if d.currentJob == nil { d.mu.Unlock() return nil } // Use the current job for the first attempt; on retry attempts, fetch fresh data from the store to resolve conflicts if attempt > 0 { // Fetch the latest version to resolve conflicts latest, err := d.store.Get(ctx, d.currentJob.GetNamespace(), d.currentJob.GetName()) if err != nil { d.mu.Unlock() if apierrors.IsNotFound(err) { // Job was completed/deleted, nothing to update return nil } return apifmt.Errorf("failed to fetch job for progress update: %w", err) } *d.currentJob = *latest } job := d.currentJob // Update status on the current job job.Status = status updated, err := d.store.Update(ctx, job) if err != nil { if apierrors.IsConflict(err) && attempt < maxRetries-1 { // Conflict detected, retry with fresh data logging.FromContext(ctx).Debug("progress update conflict, retrying", "attempt", attempt+1) continue } d.mu.Unlock() return apifmt.Errorf("failed to update job progress: %w", err) } // Update succeeded, update our local copy *d.currentJob = *updated d.mu.Unlock() span.SetAttributes( attribute.String("job.state", string(status.State)), attribute.Int("attempt", attempt+1), ) return nil } err := apifmt.Errorf("failed to update job progress after %d attempts", maxRetries) span.RecordError(err) return err } }