* refactor: Move job cleanup to separate controller and fix history write - Created JobCleanupController in apps/provisioning/pkg/controller - Separated cleanup logic from ConcurrentJobDriver - Fixed bug where expired jobs were not written to history - Added comprehensive tests with 93.8% coverage - Removed cleanup interval parameter from ConcurrentJobDriver - Cleanup now properly follows Complete + WriteJob pattern Fixes expired jobs being lost instead of archived * refactor: Update lease renewal interval to use jobExpiry variable - Changed the lease renewal interval in the GetPostStartHooks method to utilize the jobExpiry variable for improved clarity and maintainability. * Format code * Fix Unix milliseconds * fix: correct Unix timestamp assertions and remove duplicate test expectations - Changed Unix() to UnixMilli() for correct millisecond timestamp validation - Removed duplicate store.AssertExpectations(t) calls throughout tests
451 lines
14 KiB
Go
451 lines
14 KiB
Go
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
|
|
}
|
|
}
|