[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>
Этот коммит содержится в:
Christopher Poile
2025-03-14 10:24:26 -04:00
коммит произвёл GitHub
родитель c9504925e6
Коммит c049748b88
25 изменённых файлов: 292 добавлений и 183 удалений

Просмотреть файл

@@ -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"]