[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 <build@mattermost.com>
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
c9504925e6
Коммит
c049748b88
@@ -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)
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
|
||||
Ссылка в новой задаче
Block a user