[MM-54132] Use annotated logger for log messages from jobs (#24275)

Этот коммит содержится в:
Ben Schumacher
2023-09-07 08:50:22 +02:00
коммит произвёл GitHub
родитель fcfcbd9909
Коммит 30b12f199b
120 изменённых файлов: 1064 добавлений и 1111 удалений

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

@@ -12,7 +12,7 @@ import (
const schedFreq = 10 * time.Minute
func MakeScheduler(jobServer *jobs.JobServer) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
return *cfg.MetricsSettings.Enable
}

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

@@ -10,11 +10,9 @@ import (
"github.com/mattermost/mattermost/server/v8/einterfaces"
)
const (
JobName = "ActiveUsers"
)
func MakeWorker(jobServer *jobs.JobServer, store store.Store, getMetrics func() einterfaces.MetricsInterface) *jobs.SimpleWorker {
const workerName = "ActiveUsers"
func MakeWorker(jobServer *jobs.JobServer, store store.Store, getMetrics func() einterfaces.MetricsInterface) model.Worker {
isEnabled := func(cfg *model.Config) bool {
return *cfg.MetricsSettings.Enable
}
@@ -31,6 +29,6 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store, getMetrics func()
}
return nil
}
worker := jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -9,6 +9,7 @@ import (
"time"
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/request"
)
type PeriodicScheduler struct {
@@ -18,6 +19,8 @@ type PeriodicScheduler struct {
enabledFunc func(cfg *model.Config) bool
}
var _ Scheduler = (*DailyScheduler)(nil)
func NewPeriodicScheduler(jobs *JobServer, jobType string, period time.Duration, enabledFunc func(cfg *model.Config) bool) *PeriodicScheduler {
return &PeriodicScheduler{
period: period,
@@ -36,8 +39,8 @@ func (scheduler *PeriodicScheduler) NextScheduleTime(_ *model.Config, _ time.Tim
return &nextTime
}
func (scheduler *PeriodicScheduler) ScheduleJob(_ *model.Config /* pendingJobs */, _ bool /* lastSuccessfulJob */, _ *model.Job) (*model.Job, *model.AppError) {
return scheduler.jobs.CreateJob(scheduler.jobType, nil)
func (scheduler *PeriodicScheduler) ScheduleJob(c *request.Context, _ *model.Config /* pendingJobs */, _ bool /* lastSuccessfulJob */, _ *model.Job) (*model.Job, *model.AppError) {
return scheduler.jobs.CreateJob(c, scheduler.jobType, nil)
}
type DailyScheduler struct {
@@ -47,6 +50,8 @@ type DailyScheduler struct {
enabledFunc func(cfg *model.Config) bool
}
var _ Scheduler = (*DailyScheduler)(nil)
func NewDailyScheduler(jobs *JobServer, jobType string, startTimeFunc func(cfg *model.Config) *time.Time, enabledFunc func(cfg *model.Config) bool) *DailyScheduler {
return &DailyScheduler{
startTimeFunc: startTimeFunc,
@@ -69,8 +74,8 @@ func (scheduler *DailyScheduler) NextScheduleTime(cfg *model.Config, now time.Ti
return GenerateNextStartDateTime(now, *scheduledTime)
}
func (scheduler *DailyScheduler) ScheduleJob(_ *model.Config /* pendingJobs */, _ bool /* lastSuccessfulJob */, _ *model.Job) (*model.Job, *model.AppError) {
return scheduler.jobs.CreateJob(scheduler.jobType, nil)
func (scheduler *DailyScheduler) ScheduleJob(c *request.Context, _ *model.Config /* pendingJobs */, _ bool /* lastSuccessfulJob */, _ *model.Job) (*model.Job, *model.AppError) {
return scheduler.jobs.CreateJob(c, scheduler.jobType, nil)
}
const jitterRange = 2000 // milliseconds

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

@@ -8,6 +8,7 @@ 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 {
@@ -16,6 +17,7 @@ type SimpleWorker struct {
stopped chan bool
jobs chan model.Job
jobServer *JobServer
logger mlog.LoggerIFace
execute func(job *model.Job) error
isEnabled func(cfg *model.Config) bool
}
@@ -27,6 +29,7 @@ func NewSimpleWorker(name string, jobServer *JobServer, execute func(job *model.
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", name)),
execute: execute,
isEnabled: isEnabled,
}
@@ -34,27 +37,29 @@ func NewSimpleWorker(name string, jobServer *JobServer, execute func(job *model.
}
func (worker *SimpleWorker) Run() {
mlog.Debug("Worker started", mlog.String("worker", worker.name))
worker.logger.Debug("Worker started")
defer func() {
mlog.Debug("Worker finished", mlog.String("worker", worker.name))
worker.logger.Debug("Worker finished")
worker.stopped <- true
}()
for {
select {
case <-worker.stop:
mlog.Debug("Worker received stop signal", mlog.String("worker", worker.name))
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
mlog.Debug("Worker received a new candidate job.", mlog.String("worker", worker.name))
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
}
func (worker *SimpleWorker) Stop() {
mlog.Debug("Worker stopping", mlog.String("worker", worker.name))
worker.logger.Debug("Worker stopping")
worker.stop <- true
<-worker.stopped
}
@@ -69,48 +74,47 @@ func (worker *SimpleWorker) IsEnabled(cfg *model.Config) bool {
func (worker *SimpleWorker) DoJob(job *model.Job) {
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
mlog.Warn("SimpleWorker experienced an error while trying to claim job",
mlog.String("worker", worker.name),
mlog.String("job_id", job.Id),
mlog.Err(err))
job.Logger.Warn("SimpleWorker experienced an error while trying to claim job", mlog.Err(err))
return
} else if !claimed {
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(job.Id)
job, appErr = worker.jobServer.GetJob(c, job.Id)
if appErr != nil {
mlog.Error("SimpleWorker: job execution error", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(appErr))
job.Logger.Error("SimpleWorker: job execution error", mlog.Err(appErr))
worker.setJobError(job, appErr)
}
err := worker.execute(job)
if err != nil {
mlog.Error("SimpleWorker: job execution error", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("SimpleWorker: job execution error", mlog.Err(err))
worker.setJobError(job, model.NewAppError("DoJob", "app.job.error", nil, "", http.StatusInternalServerError).Wrap(err))
return
}
mlog.Info("SimpleWorker: Job is complete", mlog.String("worker", worker.name), mlog.String("job_id", job.Id))
job.Logger.Info("SimpleWorker: Job is complete")
worker.setJobSuccess(job)
}
func (worker *SimpleWorker) setJobSuccess(job *model.Job) {
if err := worker.jobServer.SetJobProgress(job, 100); err != nil {
mlog.Error("Worker: Failed to update progress for job", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
worker.setJobError(job, err)
}
if err := worker.jobServer.SetJobSuccess(job); err != nil {
mlog.Error("SimpleWorker: Failed to set success for job", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("SimpleWorker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
}
}
func (worker *SimpleWorker) setJobError(job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
mlog.Error("SimpleWorker: Failed to set job error", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("SimpleWorker: Failed to set job error", mlog.Err(err))
}
}

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

@@ -12,7 +12,7 @@ import (
const schedFreq = 1 * time.Hour
func MakeScheduler(jobServer *jobs.JobServer) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
return true
}

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

@@ -22,7 +22,7 @@ type AppIface interface {
RemoveFile(path string) *model.AppError
}
func MakeWorker(jobServer *jobs.JobServer, store store.Store) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, store store.Store) *jobs.SimpleWorker {
isEnabled := func(cfg *model.Config) bool {
return true
}

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

@@ -12,7 +12,7 @@ import (
const schedFreq = 10 * time.Minute
func MakeScheduler(jobServer *jobs.JobServer) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
return *cfg.ServiceSettings.ExtendSessionLengthWithActivity
}

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

@@ -8,11 +8,9 @@ import (
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
const (
JobName = "ExpiryNotify"
)
func MakeWorker(jobServer *jobs.JobServer, notifySessionsExpired func() error) *jobs.SimpleWorker {
const workerName = "ExpiryNotify"
func MakeWorker(jobServer *jobs.JobServer, notifySessionsExpired func() error) model.Worker {
isEnabled := func(cfg *model.Config) bool {
return *cfg.ServiceSettings.ExtendSessionLengthWithActivity
}
@@ -21,5 +19,5 @@ func MakeWorker(jobServer *jobs.JobServer, notifySessionsExpired func() error) m
return notifySessionsExpired()
}
return jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled)
return jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
}

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

@@ -12,7 +12,7 @@ import (
const schedFreq = 24 * time.Hour
func MakeScheduler(jobServer *jobs.JobServer) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
return *cfg.ExportSettings.Directory != "" && *cfg.ExportSettings.RetentionDays > 0
}

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

@@ -15,8 +15,6 @@ import (
"github.com/mattermost/mattermost/server/v8/platform/services/configservice"
)
const jobName = "ExportDelete"
type AppIface interface {
configservice.ConfigService
ListExportDirectory(path string) ([]string, *model.AppError)
@@ -24,7 +22,9 @@ type AppIface interface {
RemoveExportFile(path string) *model.AppError
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
const workerName = "ExportDelete"
isEnabled := func(cfg *model.Config) bool {
return *cfg.ExportSettings.Directory != "" && *cfg.ExportSettings.RetentionDays > 0
}
@@ -43,7 +43,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
filename := filepath.Base(exports[i])
modTime, appErr := app.ExportFileModTime(filepath.Join(exportPath, filename))
if appErr != nil {
mlog.Debug("Worker: Failed to get file modification time",
job.Logger.Debug("Worker: Failed to get file modification time",
mlog.Err(appErr), mlog.String("export", exports[i]))
errors.Append(appErr)
continue
@@ -52,7 +52,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
if time.Now().After(modTime.Add(retentionTime)) {
// remove file data from storage.
if appErr := app.RemoveExportFile(exports[i]); appErr != nil {
mlog.Debug("Worker: Failed to remove file",
job.Logger.Debug("Worker: Failed to remove file",
mlog.Err(appErr), mlog.String("export", exports[i]))
errors.Append(appErr)
continue
@@ -61,10 +61,10 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
}
if err := errors.ErrorOrNil(); err != nil {
mlog.Warn("Worker: errors occurred", mlog.String("job-name", jobName), mlog.Err(err))
job.Logger.Warn("Worker: errors occurred", mlog.Err(err))
}
return nil
}
worker := jobs.NewSimpleWorker(jobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -15,8 +15,6 @@ import (
"github.com/mattermost/mattermost/server/v8/platform/services/configservice"
)
const jobName = "ExportProcess"
type AppIface interface {
configservice.ConfigService
WriteExportFileContext(ctx context.Context, fr io.Reader, path string) (int64, *model.AppError)
@@ -24,7 +22,9 @@ type AppIface interface {
Log() *mlog.Logger
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
const workerName = "ExportProcess"
isEnabled := func(cfg *model.Config) bool { return true }
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
@@ -59,8 +59,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
}
}()
logger := app.Log().With(mlog.String("job_id", job.Id))
appErr := app.BulkExport(request.EmptyContext(logger), wr, outPath, job, opts)
appErr := app.BulkExport(request.EmptyContext(job.Logger), wr, outPath, job, opts)
wr.Close() // Close never returns an error
if appErr != nil {
@@ -69,6 +68,6 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
return nil
}
worker := jobs.NewSimpleWorker(jobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -19,13 +19,13 @@ var ignoredFiles = map[string]bool{
"mkv": true,
}
const jobName = "ExtractContent"
type AppIface interface {
ExtractContentFromFileInfo(fileInfo *model.FileInfo) error
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) *jobs.SimpleWorker {
const workerName = "ExtractContent"
isEnabled := func(cfg *model.Config) bool {
return true
}
@@ -65,10 +65,10 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) mode
}
for _, fileInfo := range fileInfos {
if !ignoredFiles[fileInfo.Extension] {
mlog.Debug("extracting file", mlog.String("filename", fileInfo.Name), mlog.String("filepath", fileInfo.Path))
job.Logger.Debug("Extracting file", mlog.String("filename", fileInfo.Name), mlog.String("filepath", fileInfo.Path))
err = app.ExtractContentFromFileInfo(fileInfo)
if err != nil {
mlog.Warn("Failed to extract file content", mlog.Err(err), mlog.String("file_info_id", fileInfo.Id))
job.Logger.Warn("Failed to extract file content", mlog.Err(err), mlog.String("file_info_id", fileInfo.Id))
nErrs++
}
nFiles++
@@ -85,10 +85,10 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) mode
job.Data["processed"] = strconv.Itoa(nFiles)
if err := jobServer.UpdateInProgressJobData(job); err != nil {
mlog.Error("Worker: Failed to update job data", mlog.String("worker", model.JobTypeExtractContent), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("Worker: Failed to update job data", mlog.Err(err))
}
return nil
}
worker := jobs.NewSimpleWorker(jobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -12,7 +12,7 @@ import (
const schedFreq = 24 * time.Hour
func MakeScheduler(jobServer *jobs.JobServer, license *model.License) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer, license *model.License) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
return model.BuildEnterpriseReady == "true" && license == nil
}

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

@@ -12,7 +12,6 @@ import (
)
const (
JobName = "HostedPurchaseScreening"
// 3 days matches the expecation given in portal purchase flow.
waitForScreeningDuration = 3 * 24 * time.Hour
)
@@ -22,7 +21,9 @@ type ScreenTimeStore interface {
PermanentDeleteByName(name string) (*model.System, error)
}
func MakeWorker(jobServer *jobs.JobServer, license *model.License, screenTimeStore ScreenTimeStore) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, license *model.License, screenTimeStore ScreenTimeStore) *jobs.SimpleWorker {
const workerName = "HostedPurchaseScreening"
isEnabled := func(_ *model.Config) bool {
return !license.IsCloud()
}
@@ -44,6 +45,6 @@ func MakeWorker(jobServer *jobs.JobServer, license *model.License, screenTimeSto
}
return nil
}
worker := jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -12,7 +12,7 @@ import (
const schedFreq = 24 * time.Hour
func MakeScheduler(jobServer *jobs.JobServer) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
return *cfg.ImportSettings.Directory != "" && *cfg.ImportSettings.RetentionDays > 0
}

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

@@ -17,8 +17,6 @@ import (
"github.com/mattermost/mattermost/server/v8/platform/services/configservice"
)
const jobName = "ImportDelete"
type AppIface interface {
configservice.ConfigService
ListDirectory(path string) ([]string, *model.AppError)
@@ -26,7 +24,9 @@ type AppIface interface {
RemoveFile(path string) *model.AppError
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) *jobs.SimpleWorker {
const workerName = "ImportDelete"
isEnabled := func(cfg *model.Config) bool {
return *cfg.ImportSettings.Directory != "" && *cfg.ImportSettings.RetentionDays > 0
}
@@ -45,7 +45,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) model.Wo
filename := filepath.Base(imports[i])
modTime, appErr := app.FileModTime(filepath.Join(importPath, filename))
if appErr != nil {
mlog.Debug("Worker: Failed to get file modification time",
job.Logger.Debug("Worker: Failed to get file modification time",
mlog.Err(appErr), mlog.String("import", imports[i]))
multipleErrors.Append(appErr)
continue
@@ -60,7 +60,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) model.Wo
if len(filename) > minLen && filepath.Ext(filename) == model.IncompleteUploadSuffix {
uploadID := filename[:26]
if storeErr := s.UploadSession().Delete(uploadID); storeErr != nil {
mlog.Debug("Worker: Failed to delete UploadSession",
job.Logger.Debug("Worker: Failed to delete UploadSession",
mlog.Err(storeErr), mlog.String("upload_id", uploadID))
multipleErrors.Append(storeErr)
continue
@@ -71,13 +71,13 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) model.Wo
info, storeErr := s.FileInfo().GetByPath(filePath)
var nfErr *store.ErrNotFound
if storeErr != nil && !errors.As(storeErr, &nfErr) {
mlog.Debug("Worker: Failed to get FileInfo",
job.Logger.Debug("Worker: Failed to get FileInfo",
mlog.Err(storeErr), mlog.String("path", filePath))
multipleErrors.Append(storeErr)
continue
} else if storeErr == nil {
if storeErr = s.FileInfo().PermanentDelete(info.Id); storeErr != nil {
mlog.Debug("Worker: Failed to delete FileInfo",
job.Logger.Debug("Worker: Failed to delete FileInfo",
mlog.Err(storeErr), mlog.String("file_id", info.Id))
multipleErrors.Append(storeErr)
continue
@@ -87,7 +87,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) model.Wo
// remove file data from storage.
if appErr := app.RemoveFile(imports[i]); appErr != nil {
mlog.Debug("Worker: Failed to remove file",
job.Logger.Debug("Worker: Failed to remove file",
mlog.Err(appErr), mlog.String("import", imports[i]))
multipleErrors.Append(appErr)
continue
@@ -96,10 +96,10 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) model.Wo
}
if err := multipleErrors.ErrorOrNil(); err != nil {
mlog.Warn("Worker: errors occurred", mlog.String("job-name", jobName), mlog.Err(err))
job.Logger.Warn("Worker: errors occurred", mlog.Err(err))
}
return nil
}
worker := jobs.NewSimpleWorker(jobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -20,8 +20,6 @@ import (
"github.com/mattermost/mattermost/server/v8/platform/shared/filestore"
)
const jobName = "ImportProcess"
type AppIface interface {
configservice.ConfigService
RemoveFile(path string) *model.AppError
@@ -32,8 +30,10 @@ type AppIface interface {
Log() *mlog.Logger
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
appContext := request.EmptyContext(app.Log())
func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
const workerName = "ImportProcess"
appContext := request.EmptyContext(jobServer.Logger())
isEnabled := func(cfg *model.Config) bool {
return true
}
@@ -113,6 +113,6 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
}
return nil
}
worker := jobs.NewSimpleWorker(jobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -4,7 +4,6 @@
package jobs
import (
"context"
"errors"
"fmt"
"net/http"
@@ -14,6 +13,7 @@ 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/store"
)
@@ -21,8 +21,8 @@ const (
CancelWatcherPollingInterval = 5000
)
func (srv *JobServer) CreateJob(jobType string, jobData map[string]string) (*model.Job, *model.AppError) {
job, appErr := srv._createJob(jobType, jobData)
func (srv *JobServer) CreateJob(c *request.Context, jobType string, jobData map[string]string) (*model.Job, *model.AppError) {
job, appErr := srv._createJob(c, jobType, jobData)
if appErr != nil {
return nil, appErr
}
@@ -34,8 +34,8 @@ func (srv *JobServer) CreateJob(jobType string, jobData map[string]string) (*mod
return job, nil
}
func (srv *JobServer) CreateJobOnce(jobType string, jobData map[string]string) (*model.Job, *model.AppError) {
job, appErr := srv._createJob(jobType, jobData)
func (srv *JobServer) CreateJobOnce(c *request.Context, jobType string, jobData map[string]string) (*model.Job, *model.AppError) {
job, appErr := srv._createJob(c, jobType, jobData)
if appErr != nil {
return nil, appErr
}
@@ -47,7 +47,7 @@ func (srv *JobServer) CreateJobOnce(jobType string, jobData map[string]string) (
return job, nil
}
func (srv *JobServer) _createJob(jobType string, jobData map[string]string) (*model.Job, *model.AppError) {
func (srv *JobServer) _createJob(c *request.Context, jobType string, jobData map[string]string) (*model.Job, *model.AppError) {
job := model.Job{
Id: model.NewId(),
Type: jobType,
@@ -56,6 +56,8 @@ func (srv *JobServer) _createJob(jobType string, jobData map[string]string) (*mo
Data: jobData,
}
job.InitLogger(c.Logger())
if err := job.IsValid(); err != nil {
return nil, err
}
@@ -67,8 +69,8 @@ func (srv *JobServer) _createJob(jobType string, jobData map[string]string) (*mo
return &job, nil
}
func (srv *JobServer) GetJob(id string) (*model.Job, *model.AppError) {
job, err := srv.Store.Job().Get(id)
func (srv *JobServer) GetJob(c *request.Context, id string) (*model.Job, *model.AppError) {
job, err := srv.Store.Job().Get(c, id)
if err != nil {
var nfErr *store.ErrNotFound
switch {
@@ -212,9 +214,15 @@ func (srv *JobServer) HandleJobPanic(job *model.Job) {
return
}
var logger mlog.LoggerIFace = job.Logger
if job.Logger == nil {
// Fall back to JobServer logger
logger = srv.logger
}
sb := &strings.Builder{}
pprof.Lookup("goroutine").WriteTo(sb, 2)
mlog.Error("Unhandled panic in job", mlog.Any("panic", r), mlog.Any("job", job), mlog.String("stack", sb.String()))
logger.Error("Unhandled panic in job", mlog.Any("panic", r), mlog.Any("job", job), mlog.String("stack", sb.String()))
rerr, ok := r.(error)
if !ok {
@@ -223,20 +231,20 @@ func (srv *JobServer) HandleJobPanic(job *model.Job) {
appErr := srv.SetJobError(job, model.NewAppError("HandleJobPanic", "app.job.update.app_error", nil, "", http.StatusInternalServerError)).Wrap(rerr)
if appErr != nil {
mlog.Error("Failed to set the job status to 'failed'", mlog.Err(appErr), mlog.Any("job", job))
logger.Error("Failed to set the job status to 'failed'", mlog.Err(appErr), mlog.Any("job", job))
}
panic(r)
}
func (srv *JobServer) RequestCancellation(jobId string) *model.AppError {
func (srv *JobServer) RequestCancellation(c *request.Context, jobId string) *model.AppError {
updated, 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 srv.metrics != nil {
job, err := srv.GetJob(jobId)
job, err := srv.GetJob(c, jobId)
if err != nil {
return model.NewAppError("RequestCancellation", "app.job.update.app_error", nil, "", http.StatusInternalServerError).Wrap(err)
}
@@ -259,17 +267,17 @@ func (srv *JobServer) RequestCancellation(jobId string) *model.AppError {
return model.NewAppError("RequestCancellation", "jobs.request_cancellation.status.error", nil, "id="+jobId, http.StatusInternalServerError)
}
func (srv *JobServer) CancellationWatcher(ctx context.Context, jobId string, cancelChan chan struct{}) {
func (srv *JobServer) CancellationWatcher(c *request.Context, jobId string, cancelChan chan struct{}) {
for {
select {
case <-ctx.Done():
mlog.Debug("CancellationWatcher for Job Aborting as job has finished.", mlog.String("job_id", jobId))
case <-c.Context().Done():
c.Logger().Debug("CancellationWatcher for Job Aborting as job has finished.", mlog.String("job_id", jobId))
return
case <-time.After(CancelWatcherPollingInterval * time.Millisecond):
mlog.Debug("CancellationWatcher for Job started polling.", mlog.String("job_id", jobId))
jobStatus, err := srv.Store.Job().Get(jobId)
c.Logger().Debug("CancellationWatcher for Job started polling.", mlog.String("job_id", jobId))
jobStatus, err := srv.Store.Job().Get(c, jobId)
if err != nil {
mlog.Warn("Error getting job", mlog.String("job_id", jobId), mlog.Err(err))
c.Logger().Warn("Error getting job", mlog.String("job_id", jobId), mlog.Err(err))
continue
}
if jobStatus.Status == model.JobStatusCancelRequested {
@@ -298,8 +306,8 @@ func (srv *JobServer) CheckForPendingJobsByType(jobType string) (bool, *model.Ap
return count > 0, nil
}
func (srv *JobServer) GetJobsByTypeAndStatus(jobType string, status string) ([]*model.Job, *model.AppError) {
jobs, err := srv.Store.Job().GetAllByTypeAndStatus(jobType, status)
func (srv *JobServer) GetJobsByTypeAndStatus(c *request.Context, jobType string, status string) ([]*model.Job, *model.AppError) {
jobs, err := srv.Store.Job().GetAllByTypeAndStatus(c, jobType, status)
if err != nil {
return nil, model.NewAppError("GetJobsByTypeAndStatus", "app.job.get_all_jobs_by_type_and_status.app_error", nil, "", http.StatusInternalServerError).Wrap(err)
}

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

@@ -9,9 +9,12 @@ import (
"net/http"
"testing"
"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"
"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/store"
"github.com/mattermost/mattermost/server/v8/channels/store/storetest"
"github.com/mattermost/mattermost/server/v8/channels/utils/testutils"
@@ -54,7 +57,7 @@ func makeTeamEditionJobServer(t *testing.T) (*JobServer, *storetest.Store) {
mockStore.AssertExpectations(t)
})
jobServer := NewJobServer(configService, mockStore, nil)
jobServer := NewJobServer(configService, mockStore, nil, mlog.CreateConsoleTestLogger(t))
return jobServer, mockStore
}
@@ -498,11 +501,13 @@ func TestUpdateInProgressJobData(t *testing.T) {
func TestHandleJobPanic(t *testing.T) {
t.Run("no panic", func(t *testing.T) {
jobServer, _, _ := makeJobServer(t)
logger := mlog.CreateConsoleTestLogger(t)
job := &model.Job{
Type: model.JobTypeImportProcess,
Status: model.JobStatusInProgress,
}
job.InitLogger(logger)
f := func() {
defer jobServer.HandleJobPanic(job)
@@ -515,11 +520,13 @@ func TestHandleJobPanic(t *testing.T) {
t.Run("with panic string", func(t *testing.T) {
jobServer, mockStore, metrics := makeJobServer(t)
logger := mlog.CreateConsoleTestLogger(t)
job := &model.Job{
Type: model.JobTypeImportProcess,
Status: model.JobStatusInProgress,
}
job.InitLogger(logger)
f := func() {
defer jobServer.HandleJobPanic(job)
@@ -535,11 +542,13 @@ func TestHandleJobPanic(t *testing.T) {
t.Run("with panic error", func(t *testing.T) {
jobServer, mockStore, metrics := makeJobServer(t)
logger := mlog.CreateConsoleTestLogger(t)
job := &model.Job{
Type: model.JobTypeImportProcess,
Status: model.JobStatusInProgress,
}
job.InitLogger(logger)
f := func() {
defer jobServer.HandleJobPanic(job)
@@ -555,12 +564,13 @@ func TestHandleJobPanic(t *testing.T) {
}
func TestRequestCancellation(t *testing.T) {
ctx := request.EmptyContext(mlog.CreateConsoleTestLogger(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"})
err := jobServer.RequestCancellation("job_id")
err := jobServer.RequestCancellation(ctx, "job_id")
expectErrorId(t, "app.job.update.app_error", err)
})
@@ -568,9 +578,9 @@ func TestRequestCancellation(t *testing.T) {
jobServer, mockStore, _ := makeJobServer(t)
mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(true, nil)
mockStore.JobStore.On("Get", "job_id").Return(nil, &store.ErrNotFound{})
mockStore.JobStore.On("Get", mock.AnythingOfType("*request.Context"), "job_id").Return(nil, &store.ErrNotFound{})
err := jobServer.RequestCancellation("job_id")
err := jobServer.RequestCancellation(ctx, "job_id")
expectErrorId(t, "app.job.update.app_error", err)
})
@@ -583,10 +593,10 @@ func TestRequestCancellation(t *testing.T) {
}
mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(true, nil)
mockStore.JobStore.On("Get", "job_id").Return(job, nil)
mockStore.JobStore.On("Get", mock.AnythingOfType("*request.Context"), "job_id").Return(job, nil)
mockMetrics.On("DecrementJobActive", "job_type")
err := jobServer.RequestCancellation("job_id")
err := jobServer.RequestCancellation(ctx, "job_id")
require.Nil(t, err)
})
@@ -595,7 +605,7 @@ func TestRequestCancellation(t *testing.T) {
mockStore.JobStore.On("UpdateStatusOptimistically", "job_id", model.JobStatusPending, model.JobStatusCanceled).Return(true, nil)
err := jobServer.RequestCancellation("job_id")
err := jobServer.RequestCancellation(ctx, "job_id")
require.Nil(t, err)
})
@@ -605,7 +615,7 @@ func TestRequestCancellation(t *testing.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"})
err := jobServer.RequestCancellation("job_id")
err := jobServer.RequestCancellation(ctx, "job_id")
expectErrorId(t, "app.job.update.app_error", err)
})
@@ -615,7 +625,7 @@ func TestRequestCancellation(t *testing.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)
err := jobServer.RequestCancellation("job_id")
err := jobServer.RequestCancellation(ctx, "job_id")
require.Nil(t, err)
})
@@ -625,7 +635,7 @@ func TestRequestCancellation(t *testing.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)
err := jobServer.RequestCancellation("job_id")
err := jobServer.RequestCancellation(ctx, "job_id")
expectErrorId(t, "jobs.request_cancellation.status.error", err)
})
}

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

@@ -9,6 +9,7 @@ import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/public/shared/request"
)
// Default polling interval for jobs termination.
@@ -67,7 +68,7 @@ func (watcher *Watcher) Stop() {
}
func (watcher *Watcher) PollAndNotify() {
jobs, err := watcher.srv.Store.Job().GetAllByStatus(model.JobStatusPending)
jobs, err := watcher.srv.Store.Job().GetAllByStatus(request.EmptyContext(watcher.srv.logger), model.JobStatusPending)
if err != nil {
mlog.Error("Error occurred getting all pending statuses.", mlog.Err(err))
return

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

@@ -14,7 +14,7 @@ import (
const schedFreq = 2 * time.Hour
func MakeScheduler(jobServer *jobs.JobServer, license *model.License) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer, license *model.License) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
enabled := license != nil && *license.Features.Cloud
mlog.Debug("Scheduler: isEnabled: "+strconv.FormatBool(enabled), mlog.String("scheduler", model.JobTypeLastAccessibleFile))

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

@@ -8,15 +8,13 @@ import (
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
const (
JobName = "LastAccessibleFile"
)
type AppIface interface {
ComputeLastAccessibleFileTime() error
}
func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) *jobs.SimpleWorker {
const workerName = "LastAccessibleFile"
isEnabled := func(_ *model.Config) bool {
return license != nil && *license.Features.Cloud
}
@@ -25,6 +23,6 @@ func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface)
return app.ComputeLastAccessibleFileTime()
}
worker := jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -14,7 +14,7 @@ import (
const schedFreq = 30 * time.Minute
func MakeScheduler(jobServer *jobs.JobServer, license *model.License) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer, license *model.License) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
enabled := license != nil && *license.Features.Cloud
mlog.Debug("Scheduler: isEnabled: "+strconv.FormatBool(enabled), mlog.String("scheduler", model.JobTypeLastAccessiblePost))

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

@@ -8,15 +8,13 @@ import (
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
const (
JobName = "LastAccessiblePost"
)
type AppIface interface {
ComputeLastAccessiblePostTime() error
}
func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) *jobs.SimpleWorker {
const workerName = "LastAccessiblePost"
isEnabled := func(_ *model.Config) bool {
return license != nil && license.Features != nil && *license.Features.Cloud
}
@@ -25,6 +23,6 @@ func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface)
return app.ComputeLastAccessiblePostTime()
}
worker := jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -7,7 +7,10 @@ import (
"testing"
"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/store"
"github.com/stretchr/testify/require"
)
func Setup(tb testing.TB) store.Store {
@@ -16,17 +19,15 @@ func Setup(tb testing.TB) store.Store {
return store
}
func deleteAllJobsByTypeAndMigrationKey(store store.Store, jobType string, migrationKey string) {
jobs, err := store.Job().GetAllByType(model.JobTypeMigrations)
if err != nil {
panic(err)
}
func deleteAllJobsByTypeAndMigrationKey(t *testing.T, store store.Store, jobType string, migrationKey string) {
ctx := request.EmptyContext(mlog.CreateConsoleTestLogger(t))
jobs, err := store.Job().GetAllByType(ctx, model.JobTypeMigrations)
require.NoError(t, err)
for _, job := range jobs {
if key, ok := job.Data[JobDataKeyMigration]; ok && key == migrationKey {
if _, err = store.Job().Delete(job.Id); err != nil {
panic(err)
}
_, err = store.Job().Delete(job.Id)
require.NoError(t, err)
}
}
}

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

@@ -7,6 +7,7 @@ import (
"net/http"
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/request"
"github.com/mattermost/mattermost/server/v8/channels/store"
)
@@ -25,12 +26,12 @@ func MakeMigrationsList() []string {
}
}
func GetMigrationState(migration string, store store.Store) (string, *model.Job, *model.AppError) {
func GetMigrationState(c *request.Context, migration string, store store.Store) (string, *model.Job, *model.AppError) {
if _, err := store.System().GetByName(migration); err == nil {
return MigrationStateCompleted, nil, nil
}
jobs, err := store.Job().GetAllByType(model.JobTypeMigrations)
jobs, err := store.Job().GetAllByType(c, model.JobTypeMigrations)
if err != nil {
return "", nil, model.NewAppError("GetMigrationState", "app.job.get_all.app_error", nil, "", http.StatusInternalServerError).Wrap(err)
}

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

@@ -10,6 +10,8 @@ import (
"github.com/stretchr/testify/require"
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/public/shared/request"
)
func TestGetMigrationState(t *testing.T) {
@@ -17,13 +19,14 @@ func TestGetMigrationState(t *testing.T) {
t.SkipNow()
}
store := Setup(t)
ctx := request.EmptyContext(mlog.CreateConsoleTestLogger(t))
migrationKey := model.NewId()
deleteAllJobsByTypeAndMigrationKey(store, model.JobTypeMigrations, migrationKey)
deleteAllJobsByTypeAndMigrationKey(t, store, model.JobTypeMigrations, migrationKey)
// Test with no job yet.
state, job, err := GetMigrationState(migrationKey, store)
state, job, err := GetMigrationState(ctx, migrationKey, store)
assert.Nil(t, err)
assert.Nil(t, job)
assert.Equal(t, "unscheduled", state)
@@ -36,7 +39,7 @@ func TestGetMigrationState(t *testing.T) {
nErr := store.System().Save(&system)
assert.NoError(t, nErr)
state, job, err = GetMigrationState(migrationKey, store)
state, job, err = GetMigrationState(ctx, migrationKey, store)
assert.Nil(t, err)
assert.Nil(t, job)
assert.Equal(t, "completed", state)
@@ -58,7 +61,7 @@ func TestGetMigrationState(t *testing.T) {
j1, nErr = store.Job().Save(j1)
require.NoError(t, nErr)
state, job, err = GetMigrationState(migrationKey, store)
state, job, err = GetMigrationState(ctx, migrationKey, store)
assert.Nil(t, err)
assert.Equal(t, j1.Id, job.Id)
assert.Equal(t, "in_progress", state)
@@ -77,7 +80,7 @@ func TestGetMigrationState(t *testing.T) {
j2, nErr = store.Job().Save(j2)
require.NoError(t, nErr)
state, job, err = GetMigrationState(migrationKey, store)
state, job, err = GetMigrationState(ctx, migrationKey, store)
assert.Nil(t, err)
assert.Equal(t, j2.Id, job.Id)
assert.Equal(t, "in_progress", state)
@@ -96,7 +99,7 @@ func TestGetMigrationState(t *testing.T) {
j3, nErr = store.Job().Save(j3)
require.NoError(t, nErr)
state, job, err = GetMigrationState(migrationKey, store)
state, job, err = GetMigrationState(ctx, migrationKey, store)
assert.Nil(t, err)
assert.Equal(t, j3.Id, job.Id)
assert.Equal(t, "unscheduled", state)

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

@@ -8,6 +8,7 @@ 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"
)
@@ -22,7 +23,9 @@ type Scheduler struct {
allMigrationsCompleted bool
}
func MakeScheduler(jobServer *jobs.JobServer, store store.Store) model.Scheduler {
var _ jobs.Scheduler = (*Scheduler)(nil)
func MakeScheduler(jobServer *jobs.JobServer, store store.Store) *Scheduler {
return &Scheduler{jobServer, store, false}
}
@@ -41,27 +44,14 @@ func (scheduler *Scheduler) NextScheduleTime(cfg *model.Config, now time.Time, p
}
//nolint:unparam
func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError) {
mlog.Debug("Scheduling Job", mlog.String("scheduler", model.JobTypeMigrations))
func (scheduler *Scheduler) ScheduleJob(c *request.Context, cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError) {
c.Logger().Debug("Scheduling Job", mlog.String("scheduler", model.JobTypeMigrations))
// Work through the list of migrations in order. Schedule the first one that isn't done (assuming it isn't in progress already).
for _, key := range MakeMigrationsList() {
state, job, err := GetMigrationState(key, scheduler.store)
state, job, err := GetMigrationState(c, key, scheduler.store)
if err != nil {
mlog.Error("Failed to determine status of migration: ", mlog.String("scheduler", model.JobTypeMigrations), mlog.String("migration_key", key), mlog.Err(err))
return nil, nil
}
if state == MigrationStateInProgress {
// Check the migration job isn't wedged.
if job != nil && job.LastActivityAt < model.GetMillis()-MigrationJobWedgedTimeoutMilliseconds && job.CreateAt < model.GetMillis()-MigrationJobWedgedTimeoutMilliseconds {
mlog.Warn("Job appears to be wedged. Rescheduling another instance.", mlog.String("scheduler", model.JobTypeMigrations), mlog.String("wedged_job_id", job.Id), mlog.String("migration_key", key))
if err := scheduler.jobServer.SetJobError(job, nil); err != nil {
mlog.Error("Worker: Failed to set job error", mlog.String("scheduler", model.JobTypeMigrations), mlog.String("job_id", job.Id), mlog.Err(err))
}
return scheduler.createJob(key, job)
}
c.Logger().Error("Failed to determine status of migration: ", mlog.String("scheduler", model.JobTypeMigrations), mlog.String("migration_key", key), mlog.Err(err))
return nil, nil
}
@@ -70,23 +60,36 @@ func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, las
continue
}
if state == MigrationStateUnscheduled {
mlog.Debug("Scheduling a new job for migration.", mlog.String("scheduler", model.JobTypeMigrations), mlog.String("migration_key", key))
return scheduler.createJob(key, job)
if state == MigrationStateInProgress {
// Check the migration job isn't wedged.
if job != nil && job.LastActivityAt < model.GetMillis()-MigrationJobWedgedTimeoutMilliseconds && job.CreateAt < model.GetMillis()-MigrationJobWedgedTimeoutMilliseconds {
job.Logger.Warn("Job appears to be wedged. Rescheduling another instance.", mlog.String("scheduler", model.JobTypeMigrations), mlog.String("wedged_job_id", job.Id), mlog.String("migration_key", key))
if err := scheduler.jobServer.SetJobError(job, nil); err != nil {
job.Logger.Error("Worker: Failed to set job error", mlog.String("scheduler", model.JobTypeMigrations), mlog.Err(err))
}
return scheduler.createJob(c, key, job)
}
return nil, nil
}
mlog.Error("Unknown migration state. Not doing anything.", mlog.String("migration_state", state))
if state == MigrationStateUnscheduled {
job.Logger.Debug("Scheduling a new job for migration.", mlog.String("scheduler", model.JobTypeMigrations), mlog.String("migration_key", key))
return scheduler.createJob(c, key, job)
}
job.Logger.Error("Unknown migration state. Not doing anything.", mlog.String("migration_state", state))
return nil, nil
}
// If we reached here, then there aren't any migrations left to run.
scheduler.allMigrationsCompleted = true
mlog.Debug("All migrations are complete.", mlog.String("scheduler", model.JobTypeMigrations))
c.Logger().Debug("All migrations are complete.", mlog.String("scheduler", model.JobTypeMigrations))
return nil, nil
}
func (scheduler *Scheduler) createJob(migrationKey string, lastJob *model.Job) (*model.Job, *model.AppError) {
func (scheduler *Scheduler) createJob(c *request.Context, migrationKey string, lastJob *model.Job) (*model.Job, *model.AppError) {
var lastDone string
if lastJob != nil {
lastDone = lastJob.Data[JobDataKeyMigrationLastDone]
@@ -97,7 +100,7 @@ func (scheduler *Scheduler) createJob(migrationKey string, lastJob *model.Job) (
JobDataKeyMigrationLastDone: lastDone,
}
job, err := scheduler.jobServer.CreateJob(model.JobTypeMigrations, data)
job, err := scheduler.jobServer.CreateJob(c, model.JobTypeMigrations, data)
if err != nil {
return nil, err
}

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

@@ -11,6 +11,7 @@ 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"
)
@@ -25,17 +26,20 @@ type Worker struct {
stopped chan bool
jobs chan model.Job
jobServer *jobs.JobServer
logger mlog.LoggerIFace
store store.Store
closed int32
}
func MakeWorker(jobServer *jobs.JobServer, store store.Store) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, store store.Store) *Worker {
const workerName = "Migrations"
worker := Worker{
name: "Migrations",
name: workerName,
stop: make(chan struct{}),
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
store: store,
}
@@ -47,20 +51,22 @@ func (worker *Worker) Run() {
if atomic.CompareAndSwapInt32(&worker.closed, 1, 0) {
worker.stop = make(chan struct{})
}
mlog.Debug("Worker started", mlog.String("worker", worker.name))
worker.logger.Debug("Worker started")
defer func() {
mlog.Debug("Worker finished", mlog.String("worker", worker.name))
worker.logger.Debug("Worker finished")
worker.stopped <- true
}()
for {
select {
case <-worker.stop:
mlog.Debug("Worker received stop signal", mlog.String("worker", worker.name))
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
mlog.Debug("Worker received a new candidate job.", mlog.String("worker", worker.name))
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
@@ -71,7 +77,7 @@ func (worker *Worker) Stop() {
if !atomic.CompareAndSwapInt32(&worker.closed, 0, 1) {
return
}
mlog.Debug("Worker stopping", mlog.String("worker", worker.name))
worker.logger.Debug("Worker stopping")
close(worker.stop)
<-worker.stopped
}
@@ -88,47 +94,45 @@ func (worker *Worker) DoJob(job *model.Job) {
defer worker.jobServer.HandleJobPanic(job)
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
mlog.Info("Worker experienced an error while trying to claim job",
mlog.String("worker", worker.name),
mlog.String("job_id", job.Id),
mlog.String("error", err.Error()))
job.Logger.Info("Worker experienced an error while trying to claim job", mlog.Err(err))
return
} else if !claimed {
return
}
cancelContext := request.EmptyContext(worker.logger)
cancelCtx, cancelCancelWatcher := context.WithCancel(context.Background())
cancelWatcherChan := make(chan struct{}, 1)
go worker.jobServer.CancellationWatcher(cancelCtx, job.Id, cancelWatcherChan)
cancelContext.SetContext(cancelCtx)
go worker.jobServer.CancellationWatcher(cancelContext, job.Id, cancelWatcherChan)
defer cancelCancelWatcher()
for {
select {
case <-cancelWatcherChan:
mlog.Debug("Worker: Job has been canceled via CancellationWatcher", mlog.String("worker", worker.name), mlog.String("job_id", job.Id))
job.Logger.Debug("Worker: Job has been canceled via CancellationWatcher")
worker.setJobCanceled(job)
return
case <-worker.stop:
mlog.Debug("Worker: Job has been canceled via Worker Stop", mlog.String("worker", worker.name), mlog.String("job_id", job.Id))
job.Logger.Debug("Worker: Job has been canceled via Worker Stop")
worker.setJobCanceled(job)
return
case <-time.After(TimeBetweenBatches * time.Millisecond):
done, progress, err := worker.runMigration(job.Data[JobDataKeyMigration], job.Data[JobDataKeyMigrationLastDone])
if err != nil {
mlog.Error("Worker: Failed to run migration", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to run migration", mlog.Err(err))
worker.setJobError(job, err)
return
} else if done {
mlog.Info("Worker: Job is complete", mlog.String("worker", worker.name), mlog.String("job_id", job.Id))
job.Logger.Info("Worker: Job is complete")
worker.setJobSuccess(job)
return
} else {
job.Data[JobDataKeyMigrationLastDone] = progress
if err := worker.jobServer.UpdateInProgressJobData(job); err != nil {
mlog.Error("Worker: Failed to update migration status data for job", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to update migration status data for job", mlog.Err(err))
worker.setJobError(job, err)
return
}
@@ -139,20 +143,20 @@ func (worker *Worker) DoJob(job *model.Job) {
func (worker *Worker) setJobSuccess(job *model.Job) {
if err := worker.jobServer.SetJobSuccess(job); err != nil {
mlog.Error("Worker: Failed to set success for job", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
}
}
func (worker *Worker) setJobError(job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
mlog.Error("Worker: Failed to set job error", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err))
}
}
func (worker *Worker) setJobCanceled(job *model.Job) {
if err := worker.jobServer.SetJobCanceled(job); err != nil {
mlog.Error("Worker: Failed to mark job as canceled", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to mark job as canceled", mlog.Err(err))
}
}

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

@@ -14,7 +14,7 @@ import (
const installPluginSchedFreq = 24 * time.Hour
func MakeInstallPluginScheduler(jobServer *jobs.JobServer, license *model.License, jobType string) model.Scheduler {
func MakeInstallPluginScheduler(jobServer *jobs.JobServer, license *model.License, jobType string) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
enabled := jobType == model.JobTypeInstallPluginNotifyAdmin
mlog.Debug("Scheduler: isEnabled: "+strconv.FormatBool(enabled), mlog.String("scheduler", jobType))

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

@@ -14,7 +14,7 @@ import (
const schedFreq = 24 * time.Hour
func MakeScheduler(jobServer *jobs.JobServer, license *model.License, jobType string) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer, license *model.License, jobType string) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
enabled := license != nil && *license.Features.Cloud
mlog.Debug("Scheduler: isEnabled: "+strconv.FormatBool(enabled), mlog.String("scheduler", jobType))

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

@@ -18,7 +18,7 @@ type AppIface interface {
DoCheckForAdminNotifications(trial bool) *model.AppError
}
func MakeUpgradeNotifyWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) model.Worker {
func MakeUpgradeNotifyWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) *jobs.SimpleWorker {
isEnabled := func(_ *model.Config) bool {
return license != nil && license.Features != nil && *license.Features.Cloud
}
@@ -36,7 +36,7 @@ func MakeUpgradeNotifyWorker(jobServer *jobs.JobServer, license *model.License,
return worker
}
func MakeTrialNotifyWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) model.Worker {
func MakeTrialNotifyWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) *jobs.SimpleWorker {
isEnabled := func(_ *model.Config) bool {
return license != nil && license.Features != nil && *license.Features.Cloud
}
@@ -54,7 +54,7 @@ func MakeTrialNotifyWorker(jobServer *jobs.JobServer, license *model.License, ap
return worker
}
func MakeInstallPluginNotifyWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
func MakeInstallPluginNotifyWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
isEnabled := func(_ *model.Config) bool {
return true
}

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

@@ -12,7 +12,7 @@ import (
const schedFreq = 24 * time.Hour
func MakeScheduler(jobServer *jobs.JobServer) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer) *jobs.PeriodicScheduler {
isEnabled := func(cfg *model.Config) bool {
return true
}

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

@@ -19,16 +19,19 @@ type Worker struct {
stopped chan bool
jobs chan model.Job
jobServer *jobs.JobServer
logger mlog.LoggerIFace
app AppIface
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface) *Worker {
const workerName = "Plugins"
worker := Worker{
name: "Plugins",
name: workerName,
stop: make(chan bool, 1),
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
app: app,
}
@@ -36,27 +39,29 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
}
func (worker *Worker) Run() {
mlog.Debug("Worker started", mlog.String("worker", worker.name))
worker.logger.Debug("Worker started")
defer func() {
mlog.Debug("Worker finished", mlog.String("worker", worker.name))
worker.logger.Debug("Worker finished")
worker.stopped <- true
}()
for {
select {
case <-worker.stop:
mlog.Debug("Worker received stop signal", mlog.String("worker", worker.name))
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
mlog.Debug("Worker received a new candidate job.", mlog.String("worker", worker.name))
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
}
func (worker *Worker) Stop() {
mlog.Debug("Worker stopping", mlog.String("worker", worker.name))
worker.logger.Debug("Worker stopping")
worker.stop <- true
<-worker.stopped
}
@@ -71,34 +76,31 @@ func (worker *Worker) IsEnabled(cfg *model.Config) bool {
func (worker *Worker) DoJob(job *model.Job) {
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
mlog.Info("Worker experienced an error while trying to claim job",
mlog.String("worker", worker.name),
mlog.String("job_id", job.Id),
mlog.String("error", err.Error()))
job.Logger.Info("Worker experienced an error while trying to claim job", mlog.Err(err))
return
} else if !claimed {
return
}
if err := worker.app.DeleteAllExpiredPluginKeys(); err != nil {
mlog.Error("Worker: Failed to delete expired keys", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to delete expired keys", mlog.Err(err))
worker.setJobError(job, err)
return
}
mlog.Info("Worker: Job is complete", mlog.String("worker", worker.name), mlog.String("job_id", job.Id))
job.Logger.Info("Worker: Job is complete")
worker.setJobSuccess(job)
}
func (worker *Worker) setJobSuccess(job *model.Job) {
if err := worker.jobServer.SetJobSuccess(job); err != nil {
mlog.Error("Worker: Failed to set success for job", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
}
}
func (worker *Worker) setJobError(job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
mlog.Error("Worker: Failed to set job error", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err))
}
}

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

@@ -19,7 +19,7 @@ func (scheduler *Scheduler) NextScheduleTime(cfg *model.Config, _ time.Time, _ b
return &nextTime
}
func MakeScheduler(jobServer *jobs.JobServer, licenseFunc func() *model.License) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer, licenseFunc func() *model.License) *Scheduler {
enabledFunc := func(_ *model.Config) bool {
l := licenseFunc()
return l != nil && (l.SkuShortName == model.LicenseShortSkuProfessional || l.SkuShortName == model.LicenseShortSkuEnterprise)

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

@@ -8,16 +8,14 @@ import (
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
const (
JobName = "PostPersistentNotifications"
)
type AppIface interface {
SendPersistentNotifications() error
IsPersistentNotificationsEnabled() bool
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
const workerName = "PostPersistentNotifications"
isEnabled := func(_ *model.Config) bool {
return app.IsPersistentNotificationsEnabled()
}
@@ -25,6 +23,6 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
defer jobServer.HandleJobPanic(job)
return app.SendPersistentNotifications()
}
worker := jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -19,7 +19,7 @@ func (scheduler *Scheduler) NextScheduleTime(cfg *model.Config, _ time.Time, pen
return &nextTime
}
func MakeScheduler(jobServer *jobs.JobServer) model.Scheduler {
func MakeScheduler(jobServer *jobs.JobServer) *Scheduler {
isEnabled := func(cfg *model.Config) bool {
return *cfg.AnnouncementSettings.AdminNoticesEnabled || *cfg.AnnouncementSettings.UserNoticesEnabled
}

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

@@ -9,13 +9,13 @@ import (
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
const jobName = "ProductNotices"
type AppIface interface {
UpdateProductNotices() *model.AppError
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
const workerName = "ProductNotices"
isEnabled := func(cfg *model.Config) bool {
return *cfg.AnnouncementSettings.AdminNoticesEnabled || *cfg.AnnouncementSettings.UserNoticesEnabled
}
@@ -23,11 +23,11 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker {
defer jobServer.HandleJobPanic(job)
if err := app.UpdateProductNotices(); err != nil {
mlog.Error("Worker: Failed to fetch product notices", mlog.String("worker", model.JobTypeProductNotices), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("Worker: Failed to fetch product notices", mlog.Err(err))
return err
}
return nil
}
worker := jobs.NewSimpleWorker(jobName, jobServer, execute, isEnabled)
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)
return worker
}

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

@@ -31,18 +31,21 @@ type ResendInvitationEmailWorker struct {
stopped chan bool
jobs chan model.Job
jobServer *jobs.JobServer
logger mlog.LoggerIFace
app AppIface
store store.Store
telemetryService *telemetry.TelemetryService
}
func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store, telemetryService *telemetry.TelemetryService) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store, telemetryService *telemetry.TelemetryService) *ResendInvitationEmailWorker {
const workerName = "ResendInvitationEmail"
worker := ResendInvitationEmailWorker{
name: model.JobTypeResendInvitationEmail,
name: workerName,
stop: make(chan bool, 1),
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
app: app,
store: store,
telemetryService: telemetryService,
@@ -51,20 +54,20 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store, tele
}
func (rseworker *ResendInvitationEmailWorker) Run() {
mlog.Debug("Worker started", mlog.String("worker", rseworker.name))
rseworker.logger.Debug("Worker started")
defer func() {
mlog.Debug("Worker finished", mlog.String("worker", rseworker.name))
rseworker.logger.Debug("Worker finished")
rseworker.stopped <- true
}()
for {
select {
case <-rseworker.stop:
mlog.Debug("Worker received stop signal", mlog.String("worker", rseworker.name))
rseworker.logger.Debug("Worker received stop signal")
return
case job := <-rseworker.jobs:
mlog.Debug("Worker received a new candidate job.", mlog.String("worker", rseworker.name))
job.Logger.Debug("Worker received a new candidate job")
rseworker.DoJob(&job)
}
}
@@ -75,7 +78,7 @@ func (rseworker *ResendInvitationEmailWorker) IsEnabled(cfg *model.Config) bool
}
func (rseworker *ResendInvitationEmailWorker) Stop() {
mlog.Debug("Worker stopping", mlog.String("worker", rseworker.name))
rseworker.logger.Debug("Worker stopping")
rseworker.stop <- true
<-rseworker.stopped
}
@@ -96,14 +99,14 @@ func (rseworker *ResendInvitationEmailWorker) DoJob(job *model.Job) {
func (rseworker *ResendInvitationEmailWorker) setJobSuccess(job *model.Job) {
if err := rseworker.jobServer.SetJobSuccess(job); err != nil {
mlog.Error("Worker: Failed to set success for job", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to set success for job", mlog.Err(err))
rseworker.setJobError(job, err)
}
}
func (rseworker *ResendInvitationEmailWorker) setJobError(job *model.Job, appError *model.AppError) {
if err := rseworker.jobServer.SetJobError(job, appError); err != nil {
mlog.Error("Worker: Failed to set job error", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error()))
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err))
}
}
@@ -179,14 +182,14 @@ func (rseworker *ResendInvitationEmailWorker) ResendEmails(job *model.Job, inter
emailList, err := rseworker.cleanEmailData(emailListData)
if err != nil {
appErr := model.NewAppError("worker: "+rseworker.name, "job_id: "+job.Id, nil, "", http.StatusInternalServerError).Wrap(err)
mlog.Error("Worker: Failed to clean emails string data", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", appErr.Error()))
job.Logger.Error("Worker: Failed to clean emails string data", mlog.Err(appErr))
rseworker.setJobError(job, appErr)
}
channelList, err := rseworker.cleanChannelsData(channelListData)
if err != nil {
appErr := model.NewAppError("worker: "+rseworker.name, "job_id: "+job.Id, nil, "", http.StatusInternalServerError).Wrap(err)
mlog.Error("Worker: Failed to clean channel string data", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", appErr.Error()))
job.Logger.Error("Worker: Failed to clean channel string data", mlog.Err(appErr))
rseworker.setJobError(job, appErr)
}
@@ -202,7 +205,7 @@ func (rseworker *ResendInvitationEmailWorker) ResendEmails(job *model.Job, inter
_, appErr := rseworker.app.InviteNewUsersToTeamGracefully(&memberInvite, teamID, job.Data["senderID"], interval)
if appErr != nil {
mlog.Error("Worker: Failed to send emails", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", appErr.Error()))
job.Logger.Error("Worker: Failed to send emails", mlog.Err(appErr))
rseworker.setJobError(job, appErr)
}
rseworker.telemetryService.SendTelemetry("track_invite_email_resend", map[string]any{interval: interval})

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

@@ -12,20 +12,20 @@ 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"
)
const (
JobName = "S3PathMigration"
timeBetweenBatches = 1 * time.Second
)
type S3PathMigrationWorker struct {
name string
jobServer *jobs.JobServer
logger mlog.LoggerIFace
store store.Store
fileBackend *filestore.S3FileBackend
@@ -34,15 +34,17 @@ type S3PathMigrationWorker struct {
jobs chan model.Job
}
func MakeWorker(jobServer *jobs.JobServer, store store.Store, fileBackend filestore.FileBackend) model.Worker {
func MakeWorker(jobServer *jobs.JobServer, store store.Store, fileBackend filestore.FileBackend) *S3PathMigrationWorker {
// If the type cast fails, it will be nil
// which is checked later.
s3Backend, _ := fileBackend.(*filestore.S3FileBackend)
const workerName = "S3PathMigration"
worker := &S3PathMigrationWorker{
name: workerName,
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
store: store,
fileBackend: s3Backend,
name: JobName,
stop: make(chan bool, 1),
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
@@ -51,30 +53,32 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store, fileBackend filest
}
func (worker *S3PathMigrationWorker) Run() {
mlog.Debug("Worker started", mlog.String("worker", worker.name))
worker.logger.Debug("Worker started")
// We have to re-assign the stop channel again, because
// it might happen that the job was restarted due to a config change.
worker.stop = make(chan bool, 1)
defer func() {
mlog.Debug("Worker finished", mlog.String("worker", worker.name))
worker.logger.Debug("Worker finished")
worker.stopped <- true
}()
for {
select {
case <-worker.stop:
mlog.Debug("Worker received stop signal", mlog.String("worker", worker.name))
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
mlog.Debug("Worker received a new candidate job.", mlog.String("worker", worker.name))
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
}
func (worker *S3PathMigrationWorker) Stop() {
mlog.Debug("Worker stopping", mlog.String("worker", worker.name))
worker.logger.Debug("Worker stopping")
close(worker.stop)
<-worker.stopped
}
@@ -104,10 +108,7 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
defer worker.jobServer.HandleJobPanic(job)
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
mlog.Warn("S3PathMigrationWorker experienced an error while trying to claim job",
mlog.String("worker", worker.name),
mlog.String("job_id", job.Id),
mlog.Err(err))
job.Logger.Warn("S3PathMigrationWorker experienced an error while trying to claim job", mlog.Err(err))
return
} else if !claimed {
return
@@ -115,16 +116,18 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
if worker.fileBackend == nil {
err := errors.New("no S3 file backend found")
mlog.Error("S3PathMigrationWorker: ", mlog.Err(err))
job.Logger.Error("S3PathMigrationWorker: ", mlog.Err(err))
worker.setJobError(job, model.NewAppError("DoJob", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
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(job.Id)
job, appErr = worker.jobServer.GetJob(c, job.Id)
if appErr != nil {
mlog.Error("S3PathMigrationWorker: job execution error", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(appErr))
job.Logger.Error("S3PathMigrationWorker: job execution error", mlog.Err(appErr))
worker.setJobError(job, appErr)
return
}
@@ -135,14 +138,14 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
doneCount, appErr := worker.getJobMetadata(job, "done_file_count")
if appErr != nil {
mlog.Error("S3PathMigrationWorker: failed to get done file count", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(appErr))
job.Logger.Error("S3PathMigrationWorker: failed to get done file count", mlog.Err(appErr))
worker.setJobError(job, appErr)
return
}
startTime, appErr := worker.getJobMetadata(job, "start_create_at")
if appErr != nil {
mlog.Error("S3PathMigrationWorker: failed to get start create_at", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(appErr))
job.Logger.Error("S3PathMigrationWorker: failed to get start create_at", mlog.Err(appErr))
worker.setJobError(job, appErr)
return
}
@@ -158,14 +161,9 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
for {
select {
case <-worker.stop:
mlog.Info("Worker: S3 Migration has been canceled via Worker Stop. Setting the job back to pending.",
mlog.String("workername", worker.name),
mlog.String("job_id", job.Id))
job.Logger.Info("Worker: S3 Migration has been canceled via Worker Stop. Setting the job back to pending.")
if err := worker.jobServer.SetJobPending(job); err != nil {
mlog.Error("Worker: Failed to mark job as pending",
mlog.String("workername", worker.name),
mlog.String("job_id", job.Id),
mlog.Err(err))
worker.logger.Error("Worker: Failed to mark job as pending", mlog.Err(err))
}
return
case <-time.After(timeBetweenBatches):
@@ -177,11 +175,11 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
files, err = worker.store.FileInfo().GetFilesBatchForIndexing(int64(startTime), startFileID, true, pageSize)
if err != nil {
if tries > 3 {
mlog.Error("Worker: Failed to get files after multiple retries. Exiting")
job.Logger.Error("Worker: Failed to get files after multiple retries. Exiting")
worker.setJobError(job, model.NewAppError("DoJob", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
return
}
mlog.Warn("Failed to get file info for s3 migration. Retrying .. ", mlog.Err(err))
job.Logger.Warn("Failed to get file info for s3 migration. Retrying .. ", mlog.Err(err))
// Wait a bit before trying again.
time.Sleep(15 * time.Second)
@@ -191,29 +189,29 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
}
if len(files) == 0 {
mlog.Info("S3PathMigrationWorker: Job is complete", mlog.String("worker", worker.name), mlog.String("job_id", job.Id))
job.Logger.Info("S3PathMigrationWorker: Job is complete")
worker.setJobSuccess(job)
worker.markAsComplete()
worker.markAsComplete(job)
return
}
// Iterate through the rows in each page.
for _, f := range files {
mlog.Debug("Processing file ID", mlog.String("id", f.Id), mlog.String("worker", worker.name), mlog.String("job_id", job.Id))
job.Logger.Debug("Processing file ID", mlog.String("id", f.Id))
// We do not fail the job if a single image failed to encode.
if f.Path != "" {
if err := worker.fileBackend.DecodeFilePathIfNeeded(f.Path); err != nil {
mlog.Warn("Failed to encode S3 file path", mlog.String("path", f.Path), mlog.String("id", f.Id), mlog.Err(err))
job.Logger.Warn("Failed to encode S3 file path", mlog.String("path", f.Path), mlog.String("id", f.Id), mlog.Err(err))
}
}
if f.PreviewPath != "" {
if err := worker.fileBackend.DecodeFilePathIfNeeded(f.PreviewPath); err != nil {
mlog.Warn("Failed to encode S3 file path", mlog.String("path", f.PreviewPath), mlog.String("id", f.Id), mlog.Err(err))
job.Logger.Warn("Failed to encode S3 file path", mlog.String("path", f.PreviewPath), mlog.String("id", f.Id), mlog.Err(err))
}
}
if f.ThumbnailPath != "" {
if err := worker.fileBackend.DecodeFilePathIfNeeded(f.ThumbnailPath); err != nil {
mlog.Warn("Failed to encode S3 file path", mlog.String("path", f.ThumbnailPath), mlog.String("id", f.Id), mlog.Err(err))
job.Logger.Warn("Failed to encode S3 file path", mlog.String("path", f.ThumbnailPath), mlog.String("id", f.Id), mlog.Err(err))
}
}
}
@@ -234,7 +232,7 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
}
}
func (worker *S3PathMigrationWorker) markAsComplete() {
func (worker *S3PathMigrationWorker) markAsComplete(job *model.Job) {
system := model.System{
Name: model.MigrationKeyS3Path,
Value: "true",
@@ -245,24 +243,24 @@ func (worker *S3PathMigrationWorker) markAsComplete() {
// it will just fall through everything because all files would have
// converted. The actual job is idempotent, so there won't be a problem.
if err := worker.jobServer.Store.System().Save(&system); err != nil {
mlog.Error("Worker: Failed to mark s3 path migration as completed in the systems table.", mlog.String("workername", worker.name), mlog.Err(err))
job.Logger.Error("Worker: Failed to mark s3 path migration as completed in the systems table.", mlog.Err(err))
}
}
func (worker *S3PathMigrationWorker) setJobSuccess(job *model.Job) {
if err := worker.jobServer.SetJobProgress(job, 100); err != nil {
mlog.Error("Worker: Failed to update progress for job", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
worker.setJobError(job, err)
}
if err := worker.jobServer.SetJobSuccess(job); err != nil {
mlog.Error("S3PathMigrationWorker: Failed to set success for job", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("S3PathMigrationWorker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
}
}
func (worker *S3PathMigrationWorker) setJobError(job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
mlog.Error("S3PathMigrationWorker: Failed to set job error", mlog.String("worker", worker.name), mlog.String("job_id", job.Id), mlog.Err(err))
job.Logger.Error("S3PathMigrationWorker: Failed to set job error", mlog.Err(err))
}
}

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

@@ -10,8 +10,15 @@ 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 Scheduler interface {
Enabled(cfg *model.Config) bool
NextScheduleTime(cfg *model.Config, now time.Time, pendingJobs bool, lastSuccessfulJob *model.Job) *time.Time
ScheduleJob(c *request.Context, cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError)
}
type Schedulers struct {
stop chan bool
stopped chan bool
@@ -22,7 +29,7 @@ type Schedulers struct {
isLeader bool
running bool
schedulers map[string]model.Scheduler
schedulers map[string]Scheduler
nextRunTimes map[string]*time.Time
}
@@ -32,7 +39,7 @@ var (
ErrSchedulersUninitialized = errors.New("job schedulers are not initialized")
)
func (schedulers *Schedulers) AddScheduler(name string, scheduler model.Scheduler) {
func (schedulers *Schedulers) AddScheduler(name string, scheduler Scheduler) {
schedulers.schedulers[name] = scheduler
}
@@ -80,7 +87,8 @@ func (schedulers *Schedulers) Start() {
if scheduler == nil || !schedulers.isLeader || !scheduler.Enabled(cfg) {
continue
}
if _, err := schedulers.scheduleJob(cfg, name, scheduler); err != nil {
c := request.EmptyContext(schedulers.jobs.Logger())
if _, err := schedulers.scheduleJob(c, cfg, name, scheduler); err != nil {
mlog.Error("Failed to schedule job", mlog.String("scheduler", name), mlog.Err(err))
continue
}
@@ -147,7 +155,7 @@ func (schedulers *Schedulers) setNextRunTime(cfg *model.Config, name string, now
mlog.Debug("Next run time for scheduler", mlog.String("scheduler_name", name), mlog.String("next_runtime", fmt.Sprintf("%v", schedulers.nextRunTimes[name])))
}
func (schedulers *Schedulers) scheduleJob(cfg *model.Config, name string, scheduler model.Scheduler) (*model.Job, *model.AppError) {
func (schedulers *Schedulers) scheduleJob(c *request.Context, cfg *model.Config, name string, scheduler Scheduler) (*model.Job, *model.AppError) {
pendingJobs, err := schedulers.jobs.CheckForPendingJobsByType(name)
if err != nil {
return nil, err
@@ -158,7 +166,7 @@ func (schedulers *Schedulers) scheduleJob(cfg *model.Config, name string, schedu
return nil, err
}
return scheduler.ScheduleJob(cfg, pendingJobs, lastSuccessfulJob)
return scheduler.ScheduleJob(c, cfg, pendingJobs, lastSuccessfulJob)
}
func (schedulers *Schedulers) handleConfigChange(_, newConfig *model.Config) {

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

@@ -12,6 +12,7 @@ import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/plugin/plugintest/mock"
"github.com/mattermost/mattermost/server/public/shared/request"
"github.com/mattermost/mattermost/server/v8/channels/store/storetest"
"github.com/mattermost/mattermost/server/v8/channels/utils/testutils"
)
@@ -29,7 +30,7 @@ func (scheduler *MockScheduler) NextScheduleTime(cfg *model.Config, now time.Tim
return &nextTime
}
func (scheduler *MockScheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError) {
func (scheduler *MockScheduler) ScheduleJob(c *request.Context, cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError) {
return nil, nil
}

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

@@ -8,6 +8,7 @@ import (
"time"
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/store"
"github.com/mattermost/mattermost/server/v8/einterfaces"
"github.com/mattermost/mattermost/server/v8/platform/services/configservice"
@@ -17,6 +18,7 @@ type JobServer struct {
ConfigService configservice.ConfigService
Store store.Store
metrics einterfaces.MetricsInterface
logger mlog.LoggerIFace
// mut is used to protect the following fields from concurrent access.
mut sync.Mutex
@@ -24,11 +26,12 @@ type JobServer struct {
schedulers *Schedulers
}
func NewJobServer(configService configservice.ConfigService, store store.Store, metrics einterfaces.MetricsInterface) *JobServer {
func NewJobServer(configService configservice.ConfigService, store store.Store, metrics einterfaces.MetricsInterface, logger mlog.LoggerIFace) *JobServer {
srv := &JobServer{
ConfigService: configService,
Store: store,
metrics: metrics,
logger: logger,
}
srv.initWorkers()
srv.initSchedulers()
@@ -47,7 +50,7 @@ func (srv *JobServer) initSchedulers() {
clusterLeaderChanged: make(chan bool, 1),
jobs: srv,
isLeader: true,
schedulers: make(map[string]model.Scheduler),
schedulers: make(map[string]Scheduler),
nextRunTimes: make(map[string]*time.Time),
}
@@ -58,7 +61,11 @@ func (srv *JobServer) Config() *model.Config {
return srv.ConfigService.Config()
}
func (srv *JobServer) RegisterJobType(name string, worker model.Worker, scheduler model.Scheduler) {
func (srv *JobServer) Logger() mlog.LoggerIFace {
return srv.logger
}
func (srv *JobServer) RegisterJobType(name string, worker model.Worker, scheduler Scheduler) {
srv.mut.Lock()
defer srv.mut.Unlock()
if worker != nil {