From c049748b88636196d76bc5a7ca8b0bfc65cad368 Mon Sep 17 00:00:00 2001 From: Christopher Poile Date: Fri, 14 Mar 2025 10:24:26 -0400 Subject: [PATCH] [MM-63314] Fix ClaimJob in HA environments (#30383) * ClaimJob now returns newly claimed job * internal code affected by change * test changes required * two branches: for mysql, use transaction; for postgres, use returning * two branches: for mysql, use transaction; for postgres, use returning * use same millis value for LastActivityAt and StartAt * blank commit --------- Co-authored-by: Mattermost Build --- server/channels/app/export_test.go | 14 +- server/channels/jobs/base_workers.go | 20 +-- server/channels/jobs/base_workers_test.go | 4 +- server/channels/jobs/batch_worker.go | 19 +-- server/channels/jobs/jobs.go | 20 +-- server/channels/jobs/jobs_test.go | 140 +++++++++++++----- server/channels/jobs/migrations/worker.go | 8 +- server/channels/jobs/plugins/worker.go | 8 +- .../s3_path_migration/s3_path_migration.go | 20 +-- .../channels/store/retrylayer/retrylayer.go | 2 +- server/channels/store/sqlstore/job_store.go | 92 +++++++++--- server/channels/store/sqlstore/utils.go | 3 +- server/channels/store/store.go | 2 +- server/channels/store/storetest/job_store.go | 32 ++-- .../store/storetest/mocks/JobStore.go | 12 +- .../channels/store/timerlayer/timerlayer.go | 2 +- .../elasticsearch/common/indexing_job.go | 6 +- .../elasticsearch/aggregation_job.go | 7 +- .../elasticsearch/aggregation_job_test.go | 3 +- .../opensearch/aggregation_job.go | 7 +- .../opensearch/aggregation_job_test.go | 3 +- server/enterprise/message_export/worker.go | 9 +- .../enterprise/message_export/worker_test.go | 23 ++- .../bleveengine/indexer/indexing_job.go | 10 +- .../bleveengine/indexer/indexing_job_test.go | 9 +- 25 files changed, 292 insertions(+), 183 deletions(-) diff --git a/server/channels/app/export_test.go b/server/channels/app/export_test.go index 9bd58a0120..e84e690544 100644 --- a/server/channels/app/export_test.go +++ b/server/channels/app/export_test.go @@ -941,27 +941,25 @@ func TestExportFileWarnings(t *testing.T) { job, appErr := th.App.Srv().Jobs.CreateJob(th.Context, model.JobTypeExportProcess, nil) require.Nil(t, appErr) - ok, appErr := th.App.Srv().Jobs.ClaimJob(job) - require.Nil(t, appErr) - require.True(t, ok) - job, appErr = th.App.Srv().Jobs.GetJob(th.Context, job.Id) + newJob, appErr := th.App.Srv().Jobs.ClaimJob(job) require.Nil(t, appErr) + require.NotNil(t, newJob) opts := model.BulkExportOpts{ IncludeAttachments: true, CreateArchive: true, } - appErr = th.App.BulkExport(th.Context, exportFile, dir, job, opts) + appErr = th.App.BulkExport(th.Context, exportFile, dir, newJob, opts) // should not get an error for the missing file require.Nil(t, appErr) // should get a warning instead: testlib.AssertLog(t, buffer, mlog.LvlWarn.Name, "Unable to export file attachment") - // should get info in the job data: - job, appErr = th.App.Srv().Jobs.GetJob(th.Context, job.Id) + // should get info in the newJob data: + newJob, appErr = th.App.Srv().Jobs.GetJob(th.Context, newJob.Id) require.Nil(t, appErr) - warnings, ok := job.Data["num_warnings"] + warnings, ok := newJob.Data["num_warnings"] require.True(t, ok) require.Equal(t, "1", warnings) diff --git a/server/channels/jobs/base_workers.go b/server/channels/jobs/base_workers.go index 72ccadc5eb..1a71326fc0 100644 --- a/server/channels/jobs/base_workers.go +++ b/server/channels/jobs/base_workers.go @@ -8,7 +8,6 @@ import ( "github.com/mattermost/mattermost/server/public/model" "github.com/mattermost/mattermost/server/public/shared/mlog" - "github.com/mattermost/mattermost/server/public/shared/request" ) type SimpleWorker struct { @@ -73,23 +72,14 @@ func (worker *SimpleWorker) DoJob(job *model.Job) { logger := worker.logger.With(JobLoggerFields(job)...) logger.Debug("Worker: Received a new candidate job.") - if claimed, err := worker.jobServer.ClaimJob(job); err != nil { - logger.Warn("SimpleWorker experienced an error while trying to claim job", mlog.Err(err)) - return - } else if !claimed { - return - } - - c := request.EmptyContext(worker.logger) - - // We get the job again because ClaimJob changes the job status. - newJob, appErr := worker.jobServer.GetJob(c, job.Id) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) if appErr != nil { - logger.Error("SimpleWorker: job execution error", mlog.Err(appErr)) - worker.setJobError(logger, job, appErr) + logger.Warn("SimpleWorker experienced an error while trying to claim job", mlog.Err(appErr)) + return + } else if job == nil { return } - job = newJob err := worker.execute(logger, job) if err != nil { diff --git a/server/channels/jobs/base_workers_test.go b/server/channels/jobs/base_workers_test.go index e7fa4c2ebc..c5d4e49a48 100644 --- a/server/channels/jobs/base_workers_test.go +++ b/server/channels/jobs/base_workers_test.go @@ -29,9 +29,9 @@ func TestSimpleWorkerPanic(t *testing.T) { return true } - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress).Return(&model.Job{Id: "job_id", Type: "job_type"}, nil) mockStore.JobStore.On("UpdateOptimistically", mock.AnythingOfType("*model.Job"), model.JobStatusInProgress).Return(true, nil) - mockStore.JobStore.On("Get", mock.AnythingOfType("*request.Context"), "job_id").Return(nil, errors.New("test")) + mockStore.JobStore.On("UpdateStatus", "job_id", "success").Return(nil, errors.New("test")) mockMetrics.On("IncrementJobActive", "job_type") mockMetrics.On("DecrementJobActive", "job_type") sWorker := NewSimpleWorker("test", jobServer, exec, isEnabled) diff --git a/server/channels/jobs/batch_worker.go b/server/channels/jobs/batch_worker.go index 0214eb1c55..ead8cbb149 100644 --- a/server/channels/jobs/batch_worker.go +++ b/server/channels/jobs/batch_worker.go @@ -119,21 +119,12 @@ func (worker *BatchWorker) DoJob(job *model.Job) { logger.Debug("Worker received a new candidate job.") defer worker.jobServer.HandleJobPanic(logger, job) - if claimed, err := worker.jobServer.ClaimJob(job); err != nil { - logger.Warn("Worker experienced an error while trying to claim job", mlog.Err(err)) - return - } else if !claimed { - return - } - - c := request.EmptyContext(logger) var appErr *model.AppError - - // We get the job again because ClaimJob changes the job status. - job, appErr = worker.jobServer.GetJob(c, job.Id) + job, appErr = worker.jobServer.ClaimJob(job) if appErr != nil { - worker.logger.Error("Worker: job execution error", mlog.Err(appErr)) - worker.setJobError(logger, job, appErr) + logger.Warn("Worker experienced an error while trying to claim job", mlog.Err(appErr)) + return + } else if job == nil { return } @@ -141,6 +132,8 @@ func (worker *BatchWorker) DoJob(job *model.Job) { job.Data = make(model.StringMap) } + c := request.EmptyContext(logger) + for { select { case <-worker.stopCh: diff --git a/server/channels/jobs/jobs.go b/server/channels/jobs/jobs.go index d457ecc1a9..fc7a468eb2 100644 --- a/server/channels/jobs/jobs.go +++ b/server/channels/jobs/jobs.go @@ -95,17 +95,17 @@ func (srv *JobServer) GetJob(c request.CTX, id string) (*model.Job, *model.AppEr return job, nil } -func (srv *JobServer) ClaimJob(job *model.Job) (bool, *model.AppError) { - updated, err := srv.Store.Job().UpdateStatusOptimistically(job.Id, model.JobStatusPending, model.JobStatusInProgress) +func (srv *JobServer) ClaimJob(job *model.Job) (*model.Job, *model.AppError) { + newJob, err := srv.Store.Job().UpdateStatusOptimistically(job.Id, model.JobStatusPending, model.JobStatusInProgress) if err != nil { - return false, model.NewAppError("ClaimJob", "app.job.update.app_error", nil, "", http.StatusInternalServerError).Wrap(err) + return nil, model.NewAppError("ClaimJob", "app.job.update.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } - if updated && srv.metrics != nil { - srv.metrics.IncrementJobActive(job.Type) + if newJob != nil && srv.metrics != nil { + srv.metrics.IncrementJobActive(newJob.Type) } - return updated, nil + return newJob, nil } func (srv *JobServer) SetJobProgress(job *model.Job, progress int64) *model.AppError { @@ -244,11 +244,11 @@ func (srv *JobServer) HandleJobPanic(logger mlog.LoggerIFace, job *model.Job) { } func (srv *JobServer) RequestCancellation(c request.CTX, jobId string) *model.AppError { - updated, err := srv.Store.Job().UpdateStatusOptimistically(jobId, model.JobStatusPending, model.JobStatusCanceled) + newJob, err := srv.Store.Job().UpdateStatusOptimistically(jobId, model.JobStatusPending, model.JobStatusCanceled) if err != nil { return model.NewAppError("RequestCancellation", "app.job.update.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } - if updated { + if newJob != nil { if srv.metrics != nil { job, err := srv.GetJob(c, jobId) if err != nil { @@ -261,12 +261,12 @@ func (srv *JobServer) RequestCancellation(c request.CTX, jobId string) *model.Ap return nil } - updated, err = srv.Store.Job().UpdateStatusOptimistically(jobId, model.JobStatusInProgress, model.JobStatusCancelRequested) + newJob, err = srv.Store.Job().UpdateStatusOptimistically(jobId, model.JobStatusInProgress, model.JobStatusCancelRequested) if err != nil { return model.NewAppError("RequestCancellation", "app.job.update.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } - if updated { + if newJob != nil { return nil } diff --git a/server/channels/jobs/jobs_test.go b/server/channels/jobs/jobs_test.go index ba40b7dccc..0860343900 100644 --- a/server/channels/jobs/jobs_test.go +++ b/server/channels/jobs/jobs_test.go @@ -71,12 +71,16 @@ func TestClaimJob(t *testing.T) { Id: "job_id", Type: "job_type", } + retJob := *job + retJob.Status = model.JobStatusInProgress - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress).Return(false, &model.AppError{Message: "message"}) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress). + Return(&retJob, &model.AppError{Message: "message"}) - updated, err := jobServer.ClaimJob(job) - expectErrorId(t, "app.job.update.app_error", err) - require.False(t, updated) + newJob, appErr := jobServer.ClaimJob(job) + expectErrorId(t, "app.job.update.app_error", appErr) + require.Nil(t, newJob) }) t.Run("no existing job to update", func(t *testing.T) { @@ -87,11 +91,13 @@ func TestClaimJob(t *testing.T) { Type: "job_type", } - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress).Return(false, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress). + Return(nil, nil) - updated, err := jobServer.ClaimJob(job) - require.Nil(t, err) - require.False(t, updated) + newJob, appErr := jobServer.ClaimJob(job) + require.Nil(t, appErr) + require.Nil(t, newJob) }) t.Run("pending job updated", func(t *testing.T) { @@ -101,13 +107,17 @@ func TestClaimJob(t *testing.T) { Id: "job_id", Type: "job_type", } + retJob := *job + retJob.Status = model.JobStatusInProgress - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress). + Return(&retJob, nil) mockMetrics.On("IncrementJobActive", "job_type") - updated, err := jobServer.ClaimJob(job) + newJob, err := jobServer.ClaimJob(job) require.Nil(t, err) - require.True(t, updated) + require.NotNil(t, newJob) }) t.Run("pending job updated, nil metrics service", func(t *testing.T) { @@ -117,12 +127,16 @@ func TestClaimJob(t *testing.T) { Id: "job_id", Type: "job_type", } + retJob := *job + retJob.Status = model.JobStatusInProgress - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusInProgress). + Return(&retJob, nil) - updated, err := jobServer.ClaimJob(job) - require.Nil(t, err) - require.True(t, updated) + newJob, appErr := jobServer.ClaimJob(job) + require.Nil(t, appErr) + require.NotNil(t, newJob) }) } @@ -136,10 +150,9 @@ func TestSetJobProgress(t *testing.T) { Type: "job_type", } - job.Status = model.JobStatusInProgress - job.Progress = progress - - mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(false, &model.AppError{Message: "message"}) + mockStore.JobStore. + On("UpdateOptimistically", job, model.JobStatusInProgress). + Return(false, &model.AppError{Message: "message"}) err := jobServer.SetJobProgress(job, progress) expectErrorId(t, "app.job.update.app_error", err) @@ -157,7 +170,9 @@ func TestSetJobProgress(t *testing.T) { job.Status = model.JobStatusInProgress job.Progress = progress - mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore. + On("UpdateOptimistically", job, model.JobStatusInProgress). + Return(true, nil) err := jobServer.SetJobProgress(job, progress) require.Nil(t, err) @@ -173,7 +188,9 @@ func TestSetJobWarning(t *testing.T) { Type: "job_type", } - mockStore.JobStore.On("UpdateStatus", "job_id", model.JobStatusWarning).Return(job, &model.AppError{Message: "message"}) + mockStore.JobStore. + On("UpdateStatus", "job_id", model.JobStatusWarning). + Return(nil, &model.AppError{Message: "message"}) err := jobServer.SetJobWarning(job) expectErrorId(t, "app.job.update.app_error", err) @@ -186,8 +203,12 @@ func TestSetJobWarning(t *testing.T) { Id: "job_id", Type: "job_type", } + retJob := *job + retJob.Status = model.JobStatusWarning - mockStore.JobStore.On("UpdateStatus", "job_id", model.JobStatusWarning).Return(job, nil) + mockStore.JobStore. + On("UpdateStatus", "job_id", model.JobStatusWarning). + Return(&retJob, nil) err := jobServer.SetJobWarning(job) require.Nil(t, err) @@ -249,7 +270,9 @@ func TestSetJobError(t *testing.T) { Type: "job_type", } - mockStore.JobStore.On("UpdateStatus", "job_id", model.JobStatusError).Return(job, &model.AppError{Message: "message"}) + mockStore.JobStore. + On("UpdateStatus", "job_id", model.JobStatusError). + Return(nil, &model.AppError{Message: "message"}) err := jobServer.SetJobError(job, nil) expectErrorId(t, "app.job.update.app_error", err) @@ -263,7 +286,9 @@ func TestSetJobError(t *testing.T) { Type: "job_type", } - mockStore.JobStore.On("UpdateStatus", "job_id", model.JobStatusError).Return(job, nil) + mockStore.JobStore. + On("UpdateStatus", "job_id", model.JobStatusError). + Return(job, nil) mockMetrics.On("DecrementJobActive", "job_type") err := jobServer.SetJobError(job, nil) @@ -298,7 +323,9 @@ func TestSetJobError(t *testing.T) { Data: map[string]string{"error": jobError.Message}, } - mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(false, &model.AppError{Message: "message"}) + mockStore.JobStore. + On("UpdateOptimistically", job, model.JobStatusInProgress). + Return(false, &model.AppError{Message: "message"}) err := jobServer.SetJobError(job, jobError) expectErrorId(t, "app.job.update.app_error", err) @@ -566,7 +593,9 @@ func TestRequestCancellation(t *testing.T) { t.Run("error cancelling", func(t *testing.T) { jobServer, mockStore, _ := makeJobServer(t) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(false, &model.AppError{Message: "message"}) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled). + Return(nil, &model.AppError{Message: "message"}) err := jobServer.RequestCancellation(ctx, "job_id") expectErrorId(t, "app.job.update.app_error", err) @@ -575,8 +604,17 @@ func TestRequestCancellation(t *testing.T) { t.Run("cancelled, job not found", func(t *testing.T) { jobServer, mockStore, _ := makeJobServer(t) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(true, nil) - mockStore.JobStore.On("Get", mock.AnythingOfType("*request.Context"), "job_id").Return(nil, &store.ErrNotFound{}) + job := &model.Job{ + Id: "job_id", + Type: "job_type", + } + + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled). + Return(job, nil) + mockStore.JobStore. + On("Get", mock.AnythingOfType("*request.Context"), "job_id"). + Return(nil, &store.ErrNotFound{}) err := jobServer.RequestCancellation(ctx, "job_id") expectErrorId(t, "app.job.update.app_error", err) @@ -590,8 +628,12 @@ func TestRequestCancellation(t *testing.T) { Type: "job_type", } - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(true, nil) - mockStore.JobStore.On("Get", mock.AnythingOfType("*request.Context"), "job_id").Return(job, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled). + Return(job, nil) + mockStore.JobStore. + On("Get", mock.AnythingOfType("*request.Context"), "job_id"). + Return(job, nil) mockMetrics.On("DecrementJobActive", "job_type") err := jobServer.RequestCancellation(ctx, "job_id") @@ -601,7 +643,14 @@ func TestRequestCancellation(t *testing.T) { t.Run("cancelled, success, nil metrics service", func(t *testing.T) { jobServer, mockStore := makeTeamEditionJobServer(t) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(true, nil) + job := &model.Job{ + Id: "job_id", + Type: "job_type", + } + + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled). + Return(job, nil) err := jobServer.RequestCancellation(ctx, "job_id") require.Nil(t, err) @@ -610,8 +659,12 @@ func TestRequestCancellation(t *testing.T) { t.Run("unable to cancel, requesting cancellation instead, error setting status", func(t *testing.T) { jobServer, mockStore, _ := makeJobServer(t) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(false, nil) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusInProgress, model.JobStatusCancelRequested).Return(false, &model.AppError{Message: "message"}) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled). + Return(nil, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusInProgress, model.JobStatusCancelRequested). + Return(nil, &model.AppError{Message: "message"}) err := jobServer.RequestCancellation(ctx, "job_id") expectErrorId(t, "app.job.update.app_error", err) @@ -620,8 +673,17 @@ func TestRequestCancellation(t *testing.T) { t.Run("unable to cancel, requesting cancellation instead, success", func(t *testing.T) { jobServer, mockStore, _ := makeJobServer(t) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(false, nil) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusInProgress, model.JobStatusCancelRequested).Return(true, nil) + job := &model.Job{ + Id: "job_id", + Type: "job_type", + } + + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled). + Return(nil, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusInProgress, model.JobStatusCancelRequested). + Return(job, nil) err := jobServer.RequestCancellation(ctx, "job_id") require.Nil(t, err) @@ -630,8 +692,12 @@ func TestRequestCancellation(t *testing.T) { t.Run("unable to cancel, requesting cancellation instead, unexpected state", func(t *testing.T) { jobServer, mockStore, _ := makeJobServer(t) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(false, nil) - mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusInProgress, model.JobStatusCancelRequested).Return(false, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled). + Return(nil, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", "job_id", model.JobStatusInProgress, model.JobStatusCancelRequested). + Return(nil, nil) err := jobServer.RequestCancellation(ctx, "job_id") expectErrorId(t, "jobs.request_cancellation.status.error", err) diff --git a/server/channels/jobs/migrations/worker.go b/server/channels/jobs/migrations/worker.go index 5cf60ad087..5d8eb8681e 100644 --- a/server/channels/jobs/migrations/worker.go +++ b/server/channels/jobs/migrations/worker.go @@ -93,10 +93,12 @@ func (worker *Worker) DoJob(job *model.Job) { defer worker.jobServer.HandleJobPanic(logger, job) - if claimed, err := worker.jobServer.ClaimJob(job); err != nil { - logger.Info("Worker experienced an error while trying to claim job", mlog.Err(err)) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) + if appErr != nil { + logger.Warn("Worker experienced an error while trying to claim job", mlog.Err(appErr)) return - } else if !claimed { + } else if job == nil { return } diff --git a/server/channels/jobs/plugins/worker.go b/server/channels/jobs/plugins/worker.go index 875c096100..d31bcf68fe 100644 --- a/server/channels/jobs/plugins/worker.go +++ b/server/channels/jobs/plugins/worker.go @@ -75,10 +75,12 @@ func (worker *Worker) DoJob(job *model.Job) { logger := worker.logger.With(jobs.JobLoggerFields(job)...) logger.Debug("Worker: Received a new candidate job.") - if claimed, err := worker.jobServer.ClaimJob(job); err != nil { - logger.Info("Worker experienced an error while trying to claim job", mlog.Err(err)) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) + if appErr != nil { + logger.Warn("Worker experienced an error while trying to claim job", mlog.Err(appErr)) return - } else if !claimed { + } else if job == nil { return } diff --git a/server/channels/jobs/s3_path_migration/s3_path_migration.go b/server/channels/jobs/s3_path_migration/s3_path_migration.go index 28203261a4..ae439c78a3 100644 --- a/server/channels/jobs/s3_path_migration/s3_path_migration.go +++ b/server/channels/jobs/s3_path_migration/s3_path_migration.go @@ -12,7 +12,6 @@ import ( "github.com/mattermost/mattermost/server/public/model" "github.com/mattermost/mattermost/server/public/shared/mlog" - "github.com/mattermost/mattermost/server/public/shared/request" "github.com/mattermost/mattermost/server/v8/channels/jobs" "github.com/mattermost/mattermost/server/v8/channels/store" "github.com/mattermost/mattermost/server/v8/platform/shared/filestore" @@ -106,10 +105,12 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) { logger.Debug("Worker: Received a new candidate job.") defer worker.jobServer.HandleJobPanic(logger, job) - if claimed, err := worker.jobServer.ClaimJob(job); err != nil { - logger.Warn("S3PathMigrationWorker experienced an error while trying to claim job", mlog.Err(err)) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) + if appErr != nil { + logger.Warn("S3PathMigrationWorker experienced an error while trying to claim job", mlog.Err(appErr)) return - } else if !claimed { + } else if job == nil { return } @@ -120,17 +121,6 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) { return } - c := request.EmptyContext(worker.logger) - - var appErr *model.AppError - // We get the job again because ClaimJob changes the job status. - job, appErr = worker.jobServer.GetJob(c, job.Id) - if appErr != nil { - logger.Error("S3PathMigrationWorker: job execution error", mlog.Err(appErr)) - worker.setJobError(logger, job, appErr) - return - } - // Check if there is metadata for that job. // If there isn't, it will be empty by default, which is the right value. startFileID := job.Data["start_file_id"] diff --git a/server/channels/store/retrylayer/retrylayer.go b/server/channels/store/retrylayer/retrylayer.go index 8f49be7067..a59bbea5e8 100644 --- a/server/channels/store/retrylayer/retrylayer.go +++ b/server/channels/store/retrylayer/retrylayer.go @@ -6384,7 +6384,7 @@ func (s *RetryLayerJobStore) UpdateStatus(id string, status string) (*model.Job, } -func (s *RetryLayerJobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (bool, error) { +func (s *RetryLayerJobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (*model.Job, error) { tries := 0 for { diff --git a/server/channels/store/sqlstore/job_store.go b/server/channels/store/sqlstore/job_store.go index fecb56865c..c91876e8b8 100644 --- a/server/channels/store/sqlstore/job_store.go +++ b/server/channels/store/sqlstore/job_store.go @@ -172,34 +172,86 @@ func (jss SqlJobStore) UpdateStatus(id string, status string) (*model.Job, error return job, nil } -func (jss SqlJobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (bool, error) { +func (jss SqlJobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (*model.Job, error) { + lastActivityAndStartTime := model.GetMillis() + + if jss.DriverName() == model.DatabaseDriverMysql { + tx, err := jss.GetMaster().Beginx() + if err != nil { + return nil, errors.Wrap(err, "begin_transaction") + } + defer finalizeTransactionX(tx, &err) + + builder := jss.getQueryBuilder(). + Update("Jobs"). + Set("LastActivityAt", lastActivityAndStartTime). + Set("Status", newStatus). + Where(sq.Eq{"Id": id, "Status": currentStatus}) + + if newStatus == model.JobStatusInProgress { + builder = builder.Set("StartAt", lastActivityAndStartTime) + } + + sqlResult, err := tx.ExecBuilder(builder) + if err != nil { + return nil, errors.Wrapf(err, "failed to update Job with id=%s", id) + } + rows, err := sqlResult.RowsAffected() + if err != nil { + return nil, errors.Wrap(err, "unable to get rows affected") + } + if rows != 1 { + return nil, nil + } + + getBuilder := jss.getQueryBuilder(). + Select("*"). + From("Jobs"). + Where(sq.Eq{"Id": id, "Status": newStatus}) + + var job model.Job + if err = tx.GetBuilder(&job, getBuilder); err != nil { + if errors.Is(err, sql.ErrNoRows) { + return nil, store.NewErrNotFound("Job", id) + } + return nil, errors.Wrapf(err, "failed to get Job with id=%s", id) + } + + err = tx.Commit() + if err != nil { + return nil, errors.Wrap(err, "commit_transaction") + } + + return &job, nil + } + + // For PostgreSQL, use RETURNING to get the updated job in a single query builder := jss.getQueryBuilder(). Update("Jobs"). - Set("LastActivityAt", model.GetMillis()). + Set("LastActivityAt", lastActivityAndStartTime). Set("Status", newStatus). - Where(sq.Eq{"Id": id, "Status": currentStatus}) + Where(sq.Eq{"Id": id, "Status": currentStatus}). + Suffix("RETURNING *") if newStatus == model.JobStatusInProgress { - builder = builder.Set("StartAt", model.GetMillis()) - } - query, args, err := builder.ToSql() - if err != nil { - return false, errors.Wrap(err, "job_tosql") + builder = builder.Set("StartAt", lastActivityAndStartTime) } - sqlResult, err := jss.GetMaster().Exec(query, args...) - if err != nil { - return false, errors.Wrapf(err, "failed to update Job with id=%s", id) - } - rows, err := sqlResult.RowsAffected() - if err != nil { - return false, errors.Wrap(err, "unable to get rows affected") - } - if rows != 1 { - return false, nil + var job []*model.Job + if err := jss.GetMaster().SelectBuilder(&job, builder); err != nil { + if errors.Is(err, sql.ErrNoRows) { + return nil, store.NewErrNotFound("Job", id) + } + return nil, errors.Wrapf(err, "failed to update Job with id=%s", id) } - return true, nil + // we are updating by id, so we should only ever update 1 job + if len(job) != 1 { + // no row was updated, but no error above, so to remain consistent we return nil, nil + return nil, nil + } + + return job[0], nil } func (jss SqlJobStore) Get(c request.CTX, id string) (*model.Job, error) { @@ -213,7 +265,7 @@ func (jss SqlJobStore) Get(c request.CTX, id string) (*model.Job, error) { var status model.Job if err = jss.GetReplica().Get(&status, query, args...); err != nil { - if err == sql.ErrNoRows { + if errors.Is(err, sql.ErrNoRows) { return nil, store.NewErrNotFound("Job", id) } return nil, errors.Wrapf(err, "failed to get Job with id=%s", id) diff --git a/server/channels/store/sqlstore/utils.go b/server/channels/store/sqlstore/utils.go index 4aed34836d..9056c6f42e 100644 --- a/server/channels/store/sqlstore/utils.go +++ b/server/channels/store/sqlstore/utils.go @@ -5,6 +5,7 @@ package sqlstore import ( "database/sql" + "errors" "fmt" "io" "net/url" @@ -54,7 +55,7 @@ func MapStringsToQueryParams(list []string, paramPrefix string) (string, map[str // finalizeTransactionX ensures a transaction is closed after use, rolling back if not already committed. func finalizeTransactionX(transaction *sqlxTxWrapper, perr *error) { // Rollback returns sql.ErrTxDone if the transaction was already closed. - if err := transaction.Rollback(); err != nil && err != sql.ErrTxDone { + if err := transaction.Rollback(); err != nil && !errors.Is(err, sql.ErrTxDone) { *perr = merror.Append(*perr, err) } } diff --git a/server/channels/store/store.go b/server/channels/store/store.go index d3f6fbd382..1b198488b0 100644 --- a/server/channels/store/store.go +++ b/server/channels/store/store.go @@ -783,7 +783,7 @@ type JobStore interface { SaveOnce(job *model.Job) (*model.Job, error) UpdateOptimistically(job *model.Job, currentStatus string) (bool, error) UpdateStatus(id string, status string) (*model.Job, error) - UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (bool, error) + UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (*model.Job, error) Get(c request.CTX, id string) (*model.Job, error) GetAllByType(c request.CTX, jobType string) ([]*model.Job, error) GetAllByTypeAndStatus(c request.CTX, jobType string, status string) ([]*model.Job, error) diff --git a/server/channels/store/storetest/job_store.go b/server/channels/store/storetest/job_store.go index 54fc07477a..4d42d938e6 100644 --- a/server/channels/store/storetest/job_store.go +++ b/server/channels/store/storetest/job_store.go @@ -600,9 +600,9 @@ func testJobUpdateStatusUpdateStatusOptimistically(t *testing.T, rctx request.CT time.Sleep(2 * time.Millisecond) - updated, err := ss.Job().UpdateStatusOptimistically(job.Id, model.JobStatusInProgress, model.JobStatusSuccess) + updatedJob, err := ss.Job().UpdateStatusOptimistically(job.Id, model.JobStatusInProgress, model.JobStatusSuccess) require.NoError(t, err) - require.False(t, updated) + require.Nil(t, updatedJob) received, err = ss.Job().Get(rctx, job.Id) require.NoError(t, err) @@ -612,30 +612,26 @@ func testJobUpdateStatusUpdateStatusOptimistically(t *testing.T, rctx request.CT time.Sleep(2 * time.Millisecond) - updated, err = ss.Job().UpdateStatusOptimistically(job.Id, model.JobStatusPending, model.JobStatusInProgress) + updatedJob, err = ss.Job().UpdateStatusOptimistically(job.Id, model.JobStatusPending, model.JobStatusInProgress) require.NoError(t, err) - require.True(t, updated, "should have succeeded") + require.NotNil(t, updatedJob, "should have succeeded") var startAtSet int64 - received, err = ss.Job().Get(rctx, job.Id) - require.NoError(t, err) - require.Equal(t, model.JobStatusInProgress, received.Status) - require.NotEqual(t, 0, received.StartAt) - require.Greater(t, received.LastActivityAt, lastUpdateAt) - lastUpdateAt = received.LastActivityAt - startAtSet = received.StartAt + require.Equal(t, model.JobStatusInProgress, updatedJob.Status) + require.NotEqual(t, 0, updatedJob.StartAt) + require.Greater(t, updatedJob.LastActivityAt, lastUpdateAt) + lastUpdateAt = updatedJob.LastActivityAt + startAtSet = updatedJob.StartAt time.Sleep(2 * time.Millisecond) - updated, err = ss.Job().UpdateStatusOptimistically(job.Id, model.JobStatusInProgress, model.JobStatusSuccess) + updatedJob, err = ss.Job().UpdateStatusOptimistically(job.Id, model.JobStatusInProgress, model.JobStatusSuccess) require.NoError(t, err) - require.True(t, updated, "should have succeeded") + require.NotNil(t, updatedJob, "should have succeeded") - received, err = ss.Job().Get(rctx, job.Id) - require.NoError(t, err) - require.Equal(t, model.JobStatusSuccess, received.Status) - require.Equal(t, startAtSet, received.StartAt) - require.Greater(t, received.LastActivityAt, lastUpdateAt) + require.Equal(t, model.JobStatusSuccess, updatedJob.Status) + require.Equal(t, startAtSet, updatedJob.StartAt) + require.Greater(t, updatedJob.LastActivityAt, lastUpdateAt) } func testJobDelete(t *testing.T, rctx request.CTX, ss store.Store) { diff --git a/server/channels/store/storetest/mocks/JobStore.go b/server/channels/store/storetest/mocks/JobStore.go index 3c7411e326..bed3b6c8f5 100644 --- a/server/channels/store/storetest/mocks/JobStore.go +++ b/server/channels/store/storetest/mocks/JobStore.go @@ -478,22 +478,24 @@ func (_m *JobStore) UpdateStatus(id string, status string) (*model.Job, error) { } // UpdateStatusOptimistically provides a mock function with given fields: id, currentStatus, newStatus -func (_m *JobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (bool, error) { +func (_m *JobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (*model.Job, error) { ret := _m.Called(id, currentStatus, newStatus) if len(ret) == 0 { panic("no return value specified for UpdateStatusOptimistically") } - var r0 bool + var r0 *model.Job var r1 error - if rf, ok := ret.Get(0).(func(string, string, string) (bool, error)); ok { + if rf, ok := ret.Get(0).(func(string, string, string) (*model.Job, error)); ok { return rf(id, currentStatus, newStatus) } - if rf, ok := ret.Get(0).(func(string, string, string) bool); ok { + if rf, ok := ret.Get(0).(func(string, string, string) *model.Job); ok { r0 = rf(id, currentStatus, newStatus) } else { - r0 = ret.Get(0).(bool) + if ret.Get(0) != nil { + r0 = ret.Get(0).(*model.Job) + } } if rf, ok := ret.Get(1).(func(string, string, string) error); ok { diff --git a/server/channels/store/timerlayer/timerlayer.go b/server/channels/store/timerlayer/timerlayer.go index 3641b4f6c4..89524e0344 100644 --- a/server/channels/store/timerlayer/timerlayer.go +++ b/server/channels/store/timerlayer/timerlayer.go @@ -5107,7 +5107,7 @@ func (s *TimerLayerJobStore) UpdateStatus(id string, status string) (*model.Job, return result, err } -func (s *TimerLayerJobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (bool, error) { +func (s *TimerLayerJobStore) UpdateStatusOptimistically(id string, currentStatus string, newStatus string) (*model.Job, error) { start := time.Now() result, err := s.JobStore.UpdateStatusOptimistically(id, currentStatus, newStatus) diff --git a/server/enterprise/elasticsearch/common/indexing_job.go b/server/enterprise/elasticsearch/common/indexing_job.go index 851c1c8779..5482ccd81e 100644 --- a/server/enterprise/elasticsearch/common/indexing_job.go +++ b/server/enterprise/elasticsearch/common/indexing_job.go @@ -219,12 +219,12 @@ func (worker *IndexerWorker) DoJob(job *model.Job) { logger.Debug("Worker: Received a new candidate job.") defer worker.jobServer.HandleJobPanic(logger, job) - claimed, appErr := worker.jobServer.ClaimJob(job) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) if appErr != nil { logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(appErr)) return - } - if !claimed { + } else if job == nil { return } diff --git a/server/enterprise/elasticsearch/elasticsearch/aggregation_job.go b/server/enterprise/elasticsearch/elasticsearch/aggregation_job.go index 198a6bdcc6..31c73f5198 100644 --- a/server/enterprise/elasticsearch/elasticsearch/aggregation_job.go +++ b/server/enterprise/elasticsearch/elasticsearch/aggregation_job.go @@ -143,13 +143,12 @@ func (worker *ElasticsearchAggregatorWorker) DoJob(job *model.Job) { logger.Debug("Worker: Received a new candidate job.") defer worker.jobServer.HandleJobPanic(logger, job) - claimed, appErr := worker.jobServer.ClaimJob(job) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) if appErr != nil { logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(appErr)) return - } - - if !claimed { + } else if job == nil { return } diff --git a/server/enterprise/elasticsearch/elasticsearch/aggregation_job_test.go b/server/enterprise/elasticsearch/elasticsearch/aggregation_job_test.go index e4bcfb6892..3722ff6fb4 100644 --- a/server/enterprise/elasticsearch/elasticsearch/aggregation_job_test.go +++ b/server/enterprise/elasticsearch/elasticsearch/aggregation_job_test.go @@ -44,7 +44,8 @@ func TestElasticsearchAggregation(t *testing.T) { mockJobStore.On("UpdateStatusOptimistically", mock.AnythingOfType("string"), model.JobStatusPending, - model.JobStatusInProgress).Return(true, nil) + model.JobStatusInProgress). + Return(&model.Job{}, nil) mockJobStore.On("GetAllByType", mock.AnythingOfType("string")).Return([]*model.Job{{ Id: "abcxyz123", Type: "EnterpriseElasticsearchIndexer", diff --git a/server/enterprise/elasticsearch/opensearch/aggregation_job.go b/server/enterprise/elasticsearch/opensearch/aggregation_job.go index f83a4e416c..9d35db25e3 100644 --- a/server/enterprise/elasticsearch/opensearch/aggregation_job.go +++ b/server/enterprise/elasticsearch/opensearch/aggregation_job.go @@ -142,13 +142,12 @@ func (worker *OpensearchAggregatorWorker) DoJob(job *model.Job) { logger.Debug("Worker: Received a new candidate job.") defer worker.jobServer.HandleJobPanic(logger, job) - claimed, appErr := worker.jobServer.ClaimJob(job) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) if appErr != nil { logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(appErr)) return - } - - if !claimed { + } else if job == nil { return } diff --git a/server/enterprise/elasticsearch/opensearch/aggregation_job_test.go b/server/enterprise/elasticsearch/opensearch/aggregation_job_test.go index e2706dd14c..ff6054f23b 100644 --- a/server/enterprise/elasticsearch/opensearch/aggregation_job_test.go +++ b/server/enterprise/elasticsearch/opensearch/aggregation_job_test.go @@ -60,7 +60,8 @@ func TestElasticsearchAggregation(t *testing.T) { mockJobStore.On("UpdateStatusOptimistically", mock.AnythingOfType("string"), model.JobStatusPending, - model.JobStatusInProgress).Return(true, nil) + model.JobStatusInProgress). + Return(&model.Job{}, nil) mockJobStore.On("GetAllByType", mock.AnythingOfType("string")).Return([]*model.Job{{ Id: "abcxyz123", Type: "EnterpriseElasticsearchIndexer", diff --git a/server/enterprise/message_export/worker.go b/server/enterprise/message_export/worker.go index 851af7e088..72170e174b 100644 --- a/server/enterprise/message_export/worker.go +++ b/server/enterprise/message_export/worker.go @@ -142,13 +142,12 @@ func (w *MessageExportWorker) DoJob(job *model.Job) { logger.Debug("Worker: Received a new candidate job.") defer w.jobServer.HandleJobPanic(logger, job) - claimed, appErr := w.jobServer.ClaimJob(job) + var appErr *model.AppError + job, appErr = w.jobServer.ClaimJob(job) if appErr != nil { - logger.Info("Worker: Error occurred while trying to claim job", mlog.Err(appErr)) + logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(appErr)) return - } - - if !claimed { + } else if job == nil { return } diff --git a/server/enterprise/message_export/worker_test.go b/server/enterprise/message_export/worker_test.go index 348c83fe03..390f45aa14 100644 --- a/server/enterprise/message_export/worker_test.go +++ b/server/enterprise/message_export/worker_test.go @@ -263,9 +263,13 @@ func TestDoJobNoPostsToExport(t *testing.T) { Status: model.JobStatusPending, Type: model.JobTypeMessageExport, } + retJob := *job + retJob.Status = model.JobStatusInProgress // claim job succeeds - mockStore.JobStore.On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress). + Return(&retJob, nil) mockMetrics.On("IncrementJobActive", model.JobTypeMessageExport) // no previous job, data will be loaded from config @@ -285,7 +289,7 @@ func TestDoJobNoPostsToExport(t *testing.T) { ) // job completed successfully - mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore.On("UpdateOptimistically", mock.AnythingOfType("*model.Job"), model.JobStatusInProgress).Return(true, nil) mockStore.JobStore.On("UpdateStatus", job.Id, model.JobStatusSuccess).Return(job, nil) mockMetrics.On("DecrementJobActive", model.JobTypeMessageExport) @@ -342,9 +346,13 @@ func TestDoJobWithDedicatedExportBackend(t *testing.T) { Status: model.JobStatusPending, Type: model.JobTypeMessageExport, } + retJob := *job + retJob.Status = model.JobStatusInProgress // claim job succeeds - mockStore.JobStore.On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress). + Return(&retJob, nil) mockMetrics.On("IncrementJobActive", model.JobTypeMessageExport) // no previous job, data will be loaded from config @@ -395,7 +403,7 @@ func TestDoJobWithDedicatedExportBackend(t *testing.T) { ) // job completed successfully - mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore.On("UpdateOptimistically", mock.AnythingOfType("*model.Job"), model.JobStatusInProgress).Return(true, nil) mockStore.JobStore.On("UpdateStatus", job.Id, model.JobStatusSuccess).Return(job, nil) mockMetrics.On("DecrementJobActive", model.JobTypeMessageExport) @@ -503,8 +511,13 @@ func TestDoJobCancel(t *testing.T) { worker, ok := impl.MakeWorker().(*MessageExportWorker) require.True(t, ok) + retJob := *job + retJob.Status = model.JobStatusInProgress + // Claim job succeeds - mockStore.JobStore.On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress). + Return(&retJob, nil) mockMetrics.On("IncrementJobActive", model.JobTypeMessageExport) // No previous job, data will be loaded from config diff --git a/server/platform/services/searchengine/bleveengine/indexer/indexing_job.go b/server/platform/services/searchengine/bleveengine/indexer/indexing_job.go index 890e394100..519c179aea 100644 --- a/server/platform/services/searchengine/bleveengine/indexer/indexing_job.go +++ b/server/platform/services/searchengine/bleveengine/indexer/indexing_job.go @@ -150,12 +150,12 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) { logger := worker.logger.With(jobs.JobLoggerFields(job)...) logger.Debug("Worker: Received a new candidate job.") - claimed, err := worker.jobServer.ClaimJob(job) - if err != nil { - logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(err)) + var appErr *model.AppError + job, appErr = worker.jobServer.ClaimJob(job) + if appErr != nil { + logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(appErr)) return - } - if !claimed { + } else if job == nil { return } diff --git a/server/platform/services/searchengine/bleveengine/indexer/indexing_job_test.go b/server/platform/services/searchengine/bleveengine/indexer/indexing_job_test.go index 796a30cb96..02ee77fc8c 100644 --- a/server/platform/services/searchengine/bleveengine/indexer/indexing_job_test.go +++ b/server/platform/services/searchengine/bleveengine/indexer/indexing_job_test.go @@ -8,6 +8,7 @@ import ( "os" "testing" + "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" "github.com/mattermost/mattermost/server/public/model" @@ -29,9 +30,13 @@ func TestBleveIndexer(t *testing.T) { Status: model.JobStatusPending, Type: model.JobTypeBlevePostIndexing, } + retJob := *job + retJob.Status = model.JobStatusInProgress - mockStore.JobStore.On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).Return(true, nil) - mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(true, nil) + mockStore.JobStore. + On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress). + Return(&retJob, nil) + mockStore.JobStore.On("UpdateOptimistically", mock.AnythingOfType("*model.Job"), model.JobStatusInProgress).Return(true, nil) mockStore.PostStore.On("GetOldestEntityCreationTime").Return(int64(1), errors.New("")) // intentionally return error to return from function tempDir, err := os.MkdirTemp("", "setupConfigFile")