Pass a logger instead of embedding on job (#24650)

* pass a logger instead of embedding on job

* leverage mlog.Millis

* use worker logger with HandleJobPanic

* rely on existing LogClone instead

* guard Job.LogClone against nil Job

* s/workername/worker_name

* Revert "rely on existing LogClone instead"

This reverts commit 17303cbac90d4b01815abca1309b78b97de368fb.

* Revert "guard Job.LogClone against nil Job"

This reverts commit f1ae22dee58d76f084582857830ffe8d4c546d7e.
Этот коммит содержится в:
Jesse Hallam
2023-10-09 11:04:55 -03:00
коммит произвёл GitHub
родитель aa172294bd
Коммит 47bfa2b66b
31 изменённых файлов: 223 добавлений и 264 удалений

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

@@ -22,7 +22,6 @@ func TestGetJob(t *testing.T) {
Id: model.NewId(),
Status: model.NewId(),
}
status.InitLogger(th.TestLogger)
_, err := th.App.Srv().Store().Job().Save(status)
require.NoError(t, err)
@@ -240,8 +239,6 @@ func TestGetJobByType(t *testing.T) {
}
for _, status := range statuses {
status.InitLogger(th.TestLogger)
_, err := th.App.Srv().Store().Job().Save(status)
require.NoError(t, err)
defer th.App.Srv().Store().Job().Delete(status.Id)
@@ -286,8 +283,6 @@ func TestGetJobsByTypes(t *testing.T) {
}
for _, status := range statuses {
status.InitLogger(th.TestLogger)
_, err := th.App.Srv().Store().Job().Save(status)
require.NoError(t, err)
defer th.App.Srv().Store().Job().Delete(status.Id)

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

@@ -5,6 +5,7 @@ package active_users
import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/jobs"
"github.com/mattermost/mattermost/server/v8/channels/store"
"github.com/mattermost/mattermost/server/v8/einterfaces"
@@ -16,8 +17,8 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store, getMetrics func()
isEnabled := func(cfg *model.Config) bool {
return *cfg.MetricsSettings.Enable
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
count, err := store.User().Count(model.UserCountOptions{IncludeDeleted: false})
if err != nil {

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

@@ -18,18 +18,18 @@ type SimpleWorker struct {
jobs chan model.Job
jobServer *JobServer
logger mlog.LoggerIFace
execute func(job *model.Job) error
execute func(logger mlog.LoggerIFace, job *model.Job) error
isEnabled func(cfg *model.Config) bool
}
func NewSimpleWorker(name string, jobServer *JobServer, execute func(job *model.Job) error, isEnabled func(cfg *model.Config) bool) *SimpleWorker {
func NewSimpleWorker(name string, jobServer *JobServer, execute func(logger mlog.LoggerIFace, job *model.Job) error, isEnabled func(cfg *model.Config) bool) *SimpleWorker {
worker := SimpleWorker{
name: name,
stop: make(chan bool, 1),
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", name)),
logger: jobServer.Logger().With(mlog.String("worker_name", name)),
execute: execute,
isEnabled: isEnabled,
}
@@ -50,9 +50,6 @@ func (worker *SimpleWorker) Run() {
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
@@ -73,8 +70,11 @@ func (worker *SimpleWorker) IsEnabled(cfg *model.Config) bool {
}
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 {
job.Logger.Warn("SimpleWorker experienced an error while trying to claim job", mlog.Err(err))
logger.Warn("SimpleWorker experienced an error while trying to claim job", mlog.Err(err))
return
} else if !claimed {
return
@@ -85,37 +85,37 @@ func (worker *SimpleWorker) DoJob(job *model.Job) {
// We get the job again because ClaimJob changes the job status.
newJob, appErr := worker.jobServer.GetJob(c, job.Id)
if appErr != nil {
job.Logger.Error("SimpleWorker: job execution error", mlog.Err(appErr))
worker.setJobError(job, appErr)
logger.Error("SimpleWorker: job execution error", mlog.Err(appErr))
worker.setJobError(logger, job, appErr)
return
}
job = newJob
err := worker.execute(job)
err := worker.execute(logger, job)
if err != nil {
job.Logger.Error("SimpleWorker: job execution error", mlog.Err(err))
worker.setJobError(job, model.NewAppError("DoJob", "app.job.error", nil, "", http.StatusInternalServerError).Wrap(err))
logger.Error("SimpleWorker: job execution error", mlog.Err(err))
worker.setJobError(logger, job, model.NewAppError("DoJob", "app.job.error", nil, "", http.StatusInternalServerError).Wrap(err))
return
}
job.Logger.Info("SimpleWorker: Job is complete")
worker.setJobSuccess(job)
logger.Info("SimpleWorker: Job is complete")
worker.setJobSuccess(logger, job)
}
func (worker *SimpleWorker) setJobSuccess(job *model.Job) {
func (worker *SimpleWorker) setJobSuccess(logger mlog.LoggerIFace, job *model.Job) {
if err := worker.jobServer.SetJobProgress(job, 100); err != nil {
job.Logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
worker.setJobError(logger, job, err)
}
if err := worker.jobServer.SetJobSuccess(job); err != nil {
job.Logger.Error("SimpleWorker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("SimpleWorker: Failed to set success for job", mlog.Err(err))
worker.setJobError(logger, job, err)
}
}
func (worker *SimpleWorker) setJobError(job *model.Job, appError *model.AppError) {
func (worker *SimpleWorker) setJobError(logger mlog.LoggerIFace, job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("SimpleWorker: Failed to set job error", mlog.Err(err))
logger.Error("SimpleWorker: Failed to set job error", mlog.Err(err))
}
}

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

@@ -17,12 +17,11 @@ func TestSimpleWorkerPanic(t *testing.T) {
jobServer, mockStore, mockMetrics := makeJobServer(t)
job := &model.Job{
Id: "job_id",
Type: "job_type",
Logger: jobServer.logger.(*mlog.Logger),
Id: "job_id",
Type: "job_type",
}
exec := func(_ *model.Job) error {
exec := func(_ mlog.LoggerIFace, _ *model.Job) error {
return nil
}

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

@@ -7,6 +7,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/jobs"
"github.com/mattermost/mattermost/server/v8/channels/store"
"github.com/mattermost/mattermost/server/v8/platform/services/configservice"
@@ -26,8 +27,8 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store) *jobs.SimpleWorker
isEnabled := func(cfg *model.Config) bool {
return true
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
return store.DesktopTokens().DeleteOlderThan(time.Now().Add(-maxAge).Unix())
}

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

@@ -5,6 +5,7 @@ package expirynotify
import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
@@ -14,8 +15,8 @@ func MakeWorker(jobServer *jobs.JobServer, notifySessionsExpired func() error) *
isEnabled := func(cfg *model.Config) bool {
return *cfg.ServiceSettings.ExtendSessionLengthWithActivity
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
return notifySessionsExpired()
}

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

@@ -28,8 +28,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
isEnabled := func(cfg *model.Config) bool {
return *cfg.ExportSettings.Directory != "" && *cfg.ExportSettings.RetentionDays > 0
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
exportPath := *app.Config().ExportSettings.Directory
retentionTime := time.Duration(*app.Config().ExportSettings.RetentionDays) * 24 * time.Hour
@@ -43,7 +43,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
filename := filepath.Base(exports[i])
modTime, appErr := app.ExportFileModTime(filepath.Join(exportPath, filename))
if appErr != nil {
job.Logger.Debug("Worker: Failed to get file modification time",
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) *jobs.SimpleWorker {
if time.Now().After(modTime.Add(retentionTime)) {
// remove file data from storage.
if appErr := app.RemoveExportFile(exports[i]); appErr != nil {
job.Logger.Debug("Worker: Failed to remove file",
logger.Debug("Worker: Failed to remove file",
mlog.Err(appErr), mlog.String("export", exports[i]))
errors.Append(appErr)
continue
@@ -61,7 +61,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
}
if err := errors.ErrorOrNil(); err != nil {
job.Logger.Warn("Worker: errors occurred", mlog.Err(err))
logger.Warn("Worker: errors occurred", mlog.Err(err))
}
return nil
}

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

@@ -26,8 +26,8 @@ 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)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
opts := model.BulkExportOpts{
CreateArchive: true,
@@ -59,7 +59,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
}
}()
appErr := app.BulkExport(request.EmptyContext(job.Logger), wr, outPath, job, opts)
appErr := app.BulkExport(request.EmptyContext(logger), wr, outPath, job, opts)
wr.Close() // Close never returns an error
if appErr != nil {

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

@@ -29,8 +29,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) *job
isEnabled := func(cfg *model.Config) bool {
return true
}
execute := func(job *model.Job) error {
jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
jobServer.HandleJobPanic(logger, job)
var err error
var fromTS int64
@@ -65,10 +65,10 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) *job
}
for _, fileInfo := range fileInfos {
if !ignoredFiles[fileInfo.Extension] {
job.Logger.Debug("Extracting file", mlog.String("filename", fileInfo.Name), mlog.String("filepath", fileInfo.Path))
logger.Debug("Extracting file", mlog.String("filename", fileInfo.Name), mlog.String("filepath", fileInfo.Path))
err = app.ExtractContentFromFileInfo(fileInfo)
if err != nil {
job.Logger.Warn("Failed to extract file content", mlog.Err(err), mlog.String("file_info_id", fileInfo.Id))
logger.Warn("Failed to extract file content", mlog.Err(err), mlog.String("file_info_id", fileInfo.Id))
nErrs++
}
nFiles++
@@ -85,7 +85,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) *job
job.Data["processed"] = strconv.Itoa(nFiles)
if err := jobServer.UpdateInProgressJobData(job); err != nil {
job.Logger.Error("Worker: Failed to update job data", mlog.Err(err))
logger.Error("Worker: Failed to update job data", mlog.Err(err))
}
return 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/jobs"
)
@@ -27,8 +28,8 @@ func MakeWorker(jobServer *jobs.JobServer, license *model.License, screenTimeSto
isEnabled := func(_ *model.Config) bool {
return !license.IsCloud()
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
now := time.Now()
screenTimeValue, err := screenTimeStore.GetByName(model.SystemHostedPurchaseNeedsScreening)

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

@@ -30,8 +30,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) *jobs.Si
isEnabled := func(cfg *model.Config) bool {
return *cfg.ImportSettings.Directory != "" && *cfg.ImportSettings.RetentionDays > 0
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
importPath := *app.Config().ImportSettings.Directory
retentionTime := time.Duration(*app.Config().ImportSettings.RetentionDays) * 24 * time.Hour
@@ -45,7 +45,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) *jobs.Si
filename := filepath.Base(imports[i])
modTime, appErr := app.FileModTime(filepath.Join(importPath, filename))
if appErr != nil {
job.Logger.Debug("Worker: Failed to get file modification time",
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) *jobs.Si
if len(filename) > minLen && filepath.Ext(filename) == model.IncompleteUploadSuffix {
uploadID := filename[:26]
if storeErr := s.UploadSession().Delete(uploadID); storeErr != nil {
job.Logger.Debug("Worker: Failed to delete UploadSession",
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) *jobs.Si
info, storeErr := s.FileInfo().GetByPath(filePath)
var nfErr *store.ErrNotFound
if storeErr != nil && !errors.As(storeErr, &nfErr) {
job.Logger.Debug("Worker: Failed to get FileInfo",
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 {
job.Logger.Debug("Worker: Failed to delete FileInfo",
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) *jobs.Si
// remove file data from storage.
if appErr := app.RemoveFile(imports[i]); appErr != nil {
job.Logger.Debug("Worker: Failed to remove file",
logger.Debug("Worker: Failed to remove file",
mlog.Err(appErr), mlog.String("import", imports[i]))
multipleErrors.Append(appErr)
continue
@@ -96,7 +96,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) *jobs.Si
}
if err := multipleErrors.ErrorOrNil(); err != nil {
job.Logger.Warn("Worker: errors occurred", mlog.Err(err))
logger.Warn("Worker: errors occurred", mlog.Err(err))
}
return nil
}

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

@@ -37,8 +37,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
isEnabled := func(cfg *model.Config) bool {
return true
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
importFileName, ok := job.Data["import_file"]
if !ok {

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

@@ -21,6 +21,19 @@ const (
CancelWatcherPollingInterval = 5000
)
// JobLoggerFields returns the logger annotations reflecting the given job metadata.
func JobLoggerFields(job *model.Job) []mlog.Field {
if job == nil {
return nil
}
return []mlog.Field{
mlog.String("job_id", job.Id),
mlog.String("job_type", job.Type),
mlog.Millis("job_create_at", job.CreateAt),
}
}
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 {
@@ -56,8 +69,6 @@ func (srv *JobServer) _createJob(c *request.Context, jobType string, jobData map
Data: jobData,
}
job.InitLogger(c.Logger())
if err := job.IsValid(); err != nil {
return nil, err
}
@@ -208,18 +219,12 @@ func (srv *JobServer) UpdateInProgressJobData(job *model.Job) *model.AppError {
// HandleJobPanic is used to handle panics during the execution of a job. It logs the panic and sets the status for the job.
// After handling, the method repanics! This method is supposed to be `defer`'d at the start of the job.
func (srv *JobServer) HandleJobPanic(job *model.Job) {
func (srv *JobServer) HandleJobPanic(logger mlog.LoggerIFace, job *model.Job) {
r := recover()
if r == nil {
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)
logger.Error("Unhandled panic in job", mlog.Any("panic", r), mlog.Any("job", job), mlog.String("stack", sb.String()))

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

@@ -501,17 +501,16 @@ 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)
jobServer, _, _ := makeJobServer(t)
job := &model.Job{
Type: model.JobTypeImportProcess,
Status: model.JobStatusInProgress,
}
job.InitLogger(logger)
f := func() {
defer jobServer.HandleJobPanic(job)
defer jobServer.HandleJobPanic(logger, job)
fmt.Println("OK")
}
@@ -520,17 +519,16 @@ func TestHandleJobPanic(t *testing.T) {
})
t.Run("with panic string", func(t *testing.T) {
jobServer, mockStore, metrics := makeJobServer(t)
logger := mlog.CreateConsoleTestLogger(t)
jobServer, mockStore, metrics := makeJobServer(t)
job := &model.Job{
Type: model.JobTypeImportProcess,
Status: model.JobStatusInProgress,
}
job.InitLogger(logger)
f := func() {
defer jobServer.HandleJobPanic(job)
defer jobServer.HandleJobPanic(logger, job)
panic("not OK")
}
@@ -542,17 +540,16 @@ func TestHandleJobPanic(t *testing.T) {
})
t.Run("with panic error", func(t *testing.T) {
jobServer, mockStore, metrics := makeJobServer(t)
logger := mlog.CreateConsoleTestLogger(t)
jobServer, mockStore, metrics := makeJobServer(t)
job := &model.Job{
Type: model.JobTypeImportProcess,
Status: model.JobStatusInProgress,
}
job.InitLogger(logger)
f := func() {
defer jobServer.HandleJobPanic(job)
defer jobServer.HandleJobPanic(logger, job)
panic(fmt.Errorf("not OK"))
}

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

@@ -5,6 +5,7 @@ package last_accessible_file
import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
@@ -18,8 +19,8 @@ func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface)
isEnabled := func(_ *model.Config) bool {
return license != nil && *license.Features.Cloud
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
return app.ComputeLastAccessibleFileTime()
}

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

@@ -5,6 +5,7 @@ package last_accessible_post
import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
@@ -18,8 +19,8 @@ func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface)
isEnabled := func(_ *model.Config) bool {
return license != nil && license.Features != nil && *license.Features.Cloud
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
return app.ComputeLastAccessiblePostTime()
}

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

@@ -55,6 +55,8 @@ func (scheduler *Scheduler) ScheduleJob(c *request.Context, cfg *model.Config, p
return nil, nil
}
logger := c.Logger().With(jobs.JobLoggerFields(job)...)
if state == MigrationStateCompleted {
// This migration is done. Continue to check the next.
continue
@@ -63,9 +65,9 @@ func (scheduler *Scheduler) ScheduleJob(c *request.Context, cfg *model.Config, p
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))
logger.Warn("Job appears to be wedged. Rescheduling another instance.", mlog.String("scheduler", model.JobTypeMigrations), 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))
logger.Error("Worker: Failed to set job error", mlog.String("scheduler", model.JobTypeMigrations), mlog.Err(err))
}
return scheduler.createJob(c, key, job)
}
@@ -74,16 +76,11 @@ func (scheduler *Scheduler) ScheduleJob(c *request.Context, cfg *model.Config, p
}
if state == MigrationStateUnscheduled {
// GetMigrationState can return a nil job
logger := scheduler.jobServer.Logger()
if job != nil {
logger = job.Logger
}
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))
logger.Error("Unknown migration state. Not doing anything.", mlog.String("migration_state", state))
return nil, nil
}

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

@@ -39,7 +39,7 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store) *Worker {
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
logger: jobServer.Logger().With(mlog.String("worker_name", workerName)),
store: store,
}
@@ -64,9 +64,6 @@ func (worker *Worker) Run() {
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
@@ -91,10 +88,13 @@ func (worker *Worker) IsEnabled(_ *model.Config) bool {
}
func (worker *Worker) DoJob(job *model.Job) {
defer worker.jobServer.HandleJobPanic(job)
logger := worker.logger.With(jobs.JobLoggerFields(job)...)
logger.Debug("Worker: Received a new candidate job.")
defer worker.jobServer.HandleJobPanic(logger, job)
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
job.Logger.Info("Worker experienced an error while trying to claim job", mlog.Err(err))
logger.Info("Worker experienced an error while trying to claim job", mlog.Err(err))
return
} else if !claimed {
return
@@ -110,30 +110,30 @@ func (worker *Worker) DoJob(job *model.Job) {
for {
select {
case <-cancelWatcherChan:
job.Logger.Debug("Worker: Job has been canceled via CancellationWatcher")
worker.setJobCanceled(job)
logger.Debug("Worker: Job has been canceled via CancellationWatcher")
worker.setJobCanceled(logger, job)
return
case <-worker.stop:
job.Logger.Debug("Worker: Job has been canceled via Worker Stop")
worker.setJobCanceled(job)
logger.Debug("Worker: Job has been canceled via Worker Stop")
worker.setJobCanceled(logger, job)
return
case <-time.After(TimeBetweenBatches * time.Millisecond):
done, progress, err := worker.runMigration(job.Data[JobDataKeyMigration], job.Data[JobDataKeyMigrationLastDone])
if err != nil {
job.Logger.Error("Worker: Failed to run migration", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("Worker: Failed to run migration", mlog.Err(err))
worker.setJobError(logger, job, err)
return
} else if done {
job.Logger.Info("Worker: Job is complete")
worker.setJobSuccess(job)
logger.Info("Worker: Job is complete")
worker.setJobSuccess(logger, job)
return
} else {
job.Data[JobDataKeyMigrationLastDone] = progress
if err := worker.jobServer.UpdateInProgressJobData(job); err != nil {
job.Logger.Error("Worker: Failed to update migration status data for job", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("Worker: Failed to update migration status data for job", mlog.Err(err))
worker.setJobError(logger, job, err)
return
}
}
@@ -141,22 +141,22 @@ func (worker *Worker) DoJob(job *model.Job) {
}
}
func (worker *Worker) setJobSuccess(job *model.Job) {
func (worker *Worker) setJobSuccess(logger mlog.LoggerIFace, job *model.Job) {
if err := worker.jobServer.SetJobSuccess(job); err != nil {
job.Logger.Error("Worker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("Worker: Failed to set success for job", mlog.Err(err))
worker.setJobError(logger, job, err)
}
}
func (worker *Worker) setJobError(job *model.Job, appError *model.AppError) {
func (worker *Worker) setJobError(logger mlog.LoggerIFace, job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err))
logger.Error("Worker: Failed to set job error", mlog.Err(err))
}
}
func (worker *Worker) setJobCanceled(job *model.Job) {
func (worker *Worker) setJobCanceled(logger mlog.LoggerIFace, job *model.Job) {
if err := worker.jobServer.SetJobCanceled(job); err != nil {
job.Logger.Error("Worker: Failed to mark job as canceled", mlog.Err(err))
logger.Error("Worker: Failed to mark job as canceled", mlog.Err(err))
}
}

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

@@ -5,6 +5,7 @@ package notify_admin
import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
@@ -22,8 +23,8 @@ func MakeUpgradeNotifyWorker(jobServer *jobs.JobServer, license *model.License,
isEnabled := func(_ *model.Config) bool {
return license != nil && license.Features != nil && *license.Features.Cloud
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
appErr := app.DoCheckForAdminNotifications(false)
if appErr != nil {
@@ -40,8 +41,8 @@ func MakeTrialNotifyWorker(jobServer *jobs.JobServer, license *model.License, ap
isEnabled := func(_ *model.Config) bool {
return license != nil && license.Features != nil && *license.Features.Cloud
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
appErr := app.DoCheckForAdminNotifications(true)
if appErr != nil {
@@ -58,8 +59,8 @@ func MakeInstallPluginNotifyWorker(jobServer *jobs.JobServer, app AppIface) *job
isEnabled := func(_ *model.Config) bool {
return true
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
appErr := app.DoCheckForAdminNotifications(false)
if appErr != nil {

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

@@ -31,7 +31,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *Worker {
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
logger: jobServer.Logger().With(mlog.String("worker_name", workerName)),
app: app,
}
@@ -52,9 +52,6 @@ func (worker *Worker) Run() {
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
@@ -75,32 +72,35 @@ func (worker *Worker) IsEnabled(cfg *model.Config) bool {
}
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 {
job.Logger.Info("Worker experienced an error while trying to claim job", mlog.Err(err))
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 {
job.Logger.Error("Worker: Failed to delete expired keys", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("Worker: Failed to delete expired keys", mlog.Err(err))
worker.setJobError(logger, job, err)
return
}
job.Logger.Info("Worker: Job is complete")
worker.setJobSuccess(job)
logger.Info("Worker: Job is complete")
worker.setJobSuccess(logger, job)
}
func (worker *Worker) setJobSuccess(job *model.Job) {
func (worker *Worker) setJobSuccess(logger mlog.LoggerIFace, job *model.Job) {
if err := worker.jobServer.SetJobSuccess(job); err != nil {
job.Logger.Error("Worker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("Worker: Failed to set success for job", mlog.Err(err))
worker.setJobError(logger, job, err)
}
}
func (worker *Worker) setJobError(job *model.Job, appError *model.AppError) {
func (worker *Worker) setJobError(logger mlog.LoggerIFace, job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err))
logger.Error("Worker: Failed to set job error", mlog.Err(err))
}
}

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

@@ -5,6 +5,7 @@ package post_persistent_notifications
import (
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/jobs"
)
@@ -19,8 +20,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
isEnabled := func(_ *model.Config) bool {
return app.IsPersistentNotificationsEnabled()
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
return app.SendPersistentNotifications()
}
worker := jobs.NewSimpleWorker(workerName, jobServer, execute, isEnabled)

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

@@ -19,11 +19,11 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) *jobs.SimpleWorker {
isEnabled := func(cfg *model.Config) bool {
return *cfg.AnnouncementSettings.AdminNoticesEnabled || *cfg.AnnouncementSettings.UserNoticesEnabled
}
execute := func(job *model.Job) error {
defer jobServer.HandleJobPanic(job)
execute := func(logger mlog.LoggerIFace, job *model.Job) error {
defer jobServer.HandleJobPanic(logger, job)
if err := app.UpdateProductNotices(); err != nil {
job.Logger.Error("Worker: Failed to fetch product notices", mlog.Err(err))
logger.Error("Worker: Failed to fetch product notices", mlog.Err(err))
return err
}
return nil

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

@@ -45,7 +45,7 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store, tele
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
logger: jobServer.Logger().With(mlog.String("worker_name", workerName)),
app: app,
store: store,
telemetryService: telemetryService,
@@ -67,7 +67,6 @@ func (rseworker *ResendInvitationEmailWorker) Run() {
rseworker.logger.Debug("Worker received stop signal")
return
case job := <-rseworker.jobs:
job.Logger.Debug("Worker received a new candidate job")
rseworker.DoJob(&job)
}
}
@@ -88,25 +87,27 @@ func (rseworker *ResendInvitationEmailWorker) JobChannel() chan<- model.Job {
}
func (rseworker *ResendInvitationEmailWorker) DoJob(job *model.Job) {
defer rseworker.jobServer.HandleJobPanic(job)
logger := rseworker.logger.With(jobs.JobLoggerFields(job)...)
logger.Debug("Worker: Received a new candidate job.")
defer rseworker.jobServer.HandleJobPanic(logger, job)
elapsedTimeSinceSchedule, DurationInMillis := rseworker.GetDurations(job)
if elapsedTimeSinceSchedule > DurationInMillis {
rseworker.ResendEmails(job, "48")
rseworker.TearDown(job)
rseworker.ResendEmails(logger, job, "48")
rseworker.TearDown(logger, job)
}
}
func (rseworker *ResendInvitationEmailWorker) setJobSuccess(job *model.Job) {
func (rseworker *ResendInvitationEmailWorker) setJobSuccess(logger mlog.LoggerIFace, job *model.Job) {
if err := rseworker.jobServer.SetJobSuccess(job); err != nil {
job.Logger.Error("Worker: Failed to set success for job", mlog.Err(err))
rseworker.setJobError(job, err)
logger.Error("Worker: Failed to set success for job", mlog.Err(err))
rseworker.setJobError(logger, job, err)
}
}
func (rseworker *ResendInvitationEmailWorker) setJobError(job *model.Job, appError *model.AppError) {
func (rseworker *ResendInvitationEmailWorker) setJobError(logger mlog.LoggerIFace, job *model.Job, appError *model.AppError) {
if err := rseworker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err))
logger.Error("Worker: Failed to set job error", mlog.Err(err))
}
}
@@ -169,12 +170,12 @@ func (rseworker *ResendInvitationEmailWorker) GetDurations(job *model.Job) (int6
}
func (rseworker *ResendInvitationEmailWorker) TearDown(job *model.Job) {
func (rseworker *ResendInvitationEmailWorker) TearDown(logger mlog.LoggerIFace, job *model.Job) {
rseworker.store.System().PermanentDeleteByName(job.Id)
rseworker.setJobSuccess(job)
rseworker.setJobSuccess(logger, job)
}
func (rseworker *ResendInvitationEmailWorker) ResendEmails(job *model.Job, interval string) {
func (rseworker *ResendInvitationEmailWorker) ResendEmails(logger mlog.LoggerIFace, job *model.Job, interval string) {
teamID := job.Data["teamID"]
emailListData := job.Data["emailList"]
channelListData := job.Data["channelList"]
@@ -182,15 +183,15 @@ 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)
job.Logger.Error("Worker: Failed to clean emails string data", mlog.Err(appErr))
rseworker.setJobError(job, appErr)
logger.Error("Worker: Failed to clean emails string data", mlog.Err(appErr))
rseworker.setJobError(logger, 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)
job.Logger.Error("Worker: Failed to clean channel string data", mlog.Err(appErr))
rseworker.setJobError(job, appErr)
logger.Error("Worker: Failed to clean channel string data", mlog.Err(appErr))
rseworker.setJobError(logger, job, appErr)
}
emailList = rseworker.removeAlreadyJoined(teamID, emailList)
@@ -205,8 +206,8 @@ func (rseworker *ResendInvitationEmailWorker) ResendEmails(job *model.Job, inter
_, appErr := rseworker.app.InviteNewUsersToTeamGracefully(&memberInvite, teamID, job.Data["senderID"], interval)
if appErr != nil {
job.Logger.Error("Worker: Failed to send emails", mlog.Err(appErr))
rseworker.setJobError(job, appErr)
logger.Error("Worker: Failed to send emails", mlog.Err(appErr))
rseworker.setJobError(logger, job, appErr)
}
rseworker.telemetryService.SendTelemetry("track_invite_email_resend", map[string]any{interval: interval})
}

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

@@ -42,7 +42,7 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store, fileBackend filest
worker := &S3PathMigrationWorker{
name: workerName,
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
logger: jobServer.Logger().With(mlog.String("worker_name", workerName)),
store: store,
fileBackend: s3Backend,
stop: make(chan bool, 1),
@@ -69,9 +69,6 @@ func (worker *S3PathMigrationWorker) Run() {
worker.logger.Debug("Worker received stop signal")
return
case job := <-worker.jobs:
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker received a new candidate job")
worker.DoJob(&job)
}
}
@@ -105,10 +102,12 @@ func (worker *S3PathMigrationWorker) getJobMetadata(job *model.Job, key string)
}
func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
defer worker.jobServer.HandleJobPanic(job)
logger := worker.logger.With(jobs.JobLoggerFields(job)...)
logger.Debug("Worker: Received a new candidate job.")
defer worker.jobServer.HandleJobPanic(logger, job)
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
job.Logger.Warn("S3PathMigrationWorker experienced an error while trying to claim job", mlog.Err(err))
logger.Warn("S3PathMigrationWorker experienced an error while trying to claim job", mlog.Err(err))
return
} else if !claimed {
return
@@ -116,8 +115,8 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
if worker.fileBackend == nil {
err := errors.New("no S3 file backend found")
job.Logger.Error("S3PathMigrationWorker: ", mlog.Err(err))
worker.setJobError(job, model.NewAppError("DoJob", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
logger.Error("S3PathMigrationWorker: ", mlog.Err(err))
worker.setJobError(logger, job, model.NewAppError("DoJob", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
return
}
@@ -127,8 +126,8 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
// We get the job again because ClaimJob changes the job status.
job, appErr = worker.jobServer.GetJob(c, job.Id)
if appErr != nil {
job.Logger.Error("S3PathMigrationWorker: job execution error", mlog.Err(appErr))
worker.setJobError(job, appErr)
logger.Error("S3PathMigrationWorker: job execution error", mlog.Err(appErr))
worker.setJobError(logger, job, appErr)
return
}
@@ -138,15 +137,15 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
doneCount, appErr := worker.getJobMetadata(job, "done_file_count")
if appErr != nil {
job.Logger.Error("S3PathMigrationWorker: failed to get done file count", mlog.Err(appErr))
worker.setJobError(job, appErr)
logger.Error("S3PathMigrationWorker: failed to get done file count", mlog.Err(appErr))
worker.setJobError(logger, job, appErr)
return
}
startTime, appErr := worker.getJobMetadata(job, "start_create_at")
if appErr != nil {
job.Logger.Error("S3PathMigrationWorker: failed to get start create_at", mlog.Err(appErr))
worker.setJobError(job, appErr)
logger.Error("S3PathMigrationWorker: failed to get start create_at", mlog.Err(appErr))
worker.setJobError(logger, job, appErr)
return
}
if startTime == 0 {
@@ -161,7 +160,7 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
for {
select {
case <-worker.stop:
job.Logger.Info("Worker: S3 Migration has been canceled via Worker Stop. Setting the job back to pending.")
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 {
worker.logger.Error("Worker: Failed to mark job as pending", mlog.Err(err))
}
@@ -175,11 +174,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 {
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))
logger.Error("Worker: Failed to get files after multiple retries. Exiting")
worker.setJobError(logger, job, model.NewAppError("DoJob", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
return
}
job.Logger.Warn("Failed to get file info for s3 migration. Retrying .. ", mlog.Err(err))
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)
@@ -189,29 +188,29 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
}
if len(files) == 0 {
job.Logger.Info("S3PathMigrationWorker: Job is complete")
worker.setJobSuccess(job)
worker.markAsComplete(job)
logger.Info("S3PathMigrationWorker: Job is complete")
worker.setJobSuccess(logger, job)
worker.markAsComplete(logger, job)
return
}
// Iterate through the rows in each page.
for _, f := range files {
job.Logger.Debug("Processing file ID", mlog.String("id", f.Id))
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 {
job.Logger.Warn("Failed to encode S3 file path", mlog.String("path", f.Path), mlog.String("id", f.Id), mlog.Err(err))
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 {
job.Logger.Warn("Failed to encode S3 file path", mlog.String("path", f.PreviewPath), mlog.String("id", f.Id), mlog.Err(err))
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 {
job.Logger.Warn("Failed to encode S3 file path", mlog.String("path", f.ThumbnailPath), mlog.String("id", f.Id), mlog.Err(err))
logger.Warn("Failed to encode S3 file path", mlog.String("path", f.ThumbnailPath), mlog.String("id", f.Id), mlog.Err(err))
}
}
}
@@ -232,7 +231,7 @@ func (worker *S3PathMigrationWorker) DoJob(job *model.Job) {
}
}
func (worker *S3PathMigrationWorker) markAsComplete(job *model.Job) {
func (worker *S3PathMigrationWorker) markAsComplete(logger mlog.LoggerIFace, job *model.Job) {
system := model.System{
Name: model.MigrationKeyS3Path,
Value: "true",
@@ -243,24 +242,24 @@ func (worker *S3PathMigrationWorker) markAsComplete(job *model.Job) {
// 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 {
job.Logger.Error("Worker: Failed to mark s3 path migration as completed in the systems table.", mlog.Err(err))
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) {
func (worker *S3PathMigrationWorker) setJobSuccess(logger mlog.LoggerIFace, job *model.Job) {
if err := worker.jobServer.SetJobProgress(job, 100); err != nil {
job.Logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
worker.setJobError(logger, job, err)
}
if err := worker.jobServer.SetJobSuccess(job); err != nil {
job.Logger.Error("S3PathMigrationWorker: Failed to set success for job", mlog.Err(err))
worker.setJobError(job, err)
logger.Error("S3PathMigrationWorker: Failed to set success for job", mlog.Err(err))
worker.setJobError(logger, job, err)
}
}
func (worker *S3PathMigrationWorker) setJobError(job *model.Job, appError *model.AppError) {
func (worker *S3PathMigrationWorker) setJobError(logger mlog.LoggerIFace, job *model.Job, appError *model.AppError) {
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("S3PathMigrationWorker: Failed to set job error", mlog.Err(err))
logger.Error("S3PathMigrationWorker: Failed to set job error", mlog.Err(err))
}
}

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

@@ -219,8 +219,6 @@ func (jss SqlJobStore) Get(c *request.Context, id string) (*model.Job, error) {
return nil, errors.Wrapf(err, "failed to get Job with id=%s", id)
}
status.InitLogger(c.Logger())
return &status, nil
}
@@ -240,9 +238,6 @@ func (jss SqlJobStore) GetAllByTypesPage(c *request.Context, jobTypes []string,
if err = jss.GetReplicaX().Select(&jobs, query, args...); err != nil {
return nil, errors.Wrapf(err, "failed to find Jobs with types")
}
for _, j := range jobs {
j.InitLogger(c.Logger())
}
return jobs, nil
}
@@ -261,9 +256,6 @@ func (jss SqlJobStore) GetAllByType(c *request.Context, jobType string) ([]*mode
if err = jss.GetReplicaX().Select(&statuses, query, args...); err != nil {
return nil, errors.Wrapf(err, "failed to find Jobs with type=%s", jobType)
}
for _, j := range statuses {
j.InitLogger(c.Logger())
}
return statuses, nil
}
@@ -282,9 +274,6 @@ func (jss SqlJobStore) GetAllByTypeAndStatus(c *request.Context, jobType string,
if err = jss.GetReplicaX().Select(&jobs, query, args...); err != nil {
return nil, errors.Wrapf(err, "failed to find Jobs with type=%s", jobType)
}
for _, j := range jobs {
j.InitLogger(c.Logger())
}
return jobs, nil
}
@@ -305,9 +294,6 @@ func (jss SqlJobStore) GetAllByTypePage(c *request.Context, jobType string, offs
if err = jss.GetReplicaX().Select(&statuses, query, args...); err != nil {
return nil, errors.Wrapf(err, "failed to find Jobs with type=%s", jobType)
}
for _, j := range statuses {
j.InitLogger(c.Logger())
}
return statuses, nil
}
@@ -326,9 +312,6 @@ func (jss SqlJobStore) GetAllByStatus(c *request.Context, status string) ([]*mod
if err = jss.GetReplicaX().Select(&statuses, query, args...); err != nil {
return nil, errors.Wrapf(err, "failed to find Jobs with status=%s", status)
}
for _, j := range statuses {
j.InitLogger(c.Logger())
}
return statuses, nil
}

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

@@ -280,7 +280,6 @@ func (s *MmctlE2ETestSuite) TestExportJobShowCmdF() {
Type: model.JobTypeExportProcess,
})
s.Require().Nil(appErr)
job.Logger = nil
time.Sleep(time.Millisecond)
@@ -370,7 +369,6 @@ func (s *MmctlE2ETestSuite) TestExportJobListCmdF() {
Type: model.JobTypeExportProcess,
})
s.Require().Nil(appErr)
job2.Logger = nil
time.Sleep(time.Millisecond)
@@ -378,7 +376,6 @@ func (s *MmctlE2ETestSuite) TestExportJobListCmdF() {
Type: model.JobTypeExportProcess,
})
s.Require().Nil(appErr)
job3.Logger = nil
err := exportJobListCmdF(c, cmd, nil)
s.Require().Nil(err)

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

@@ -76,7 +76,6 @@ func (s *MmctlE2ETestSuite) TestExtractJobShowCmdF() {
Data: map[string]string{},
})
s.Require().Nil(appErr)
job.Logger = nil
s.Run("no permissions", func() {
printer.Clean()
@@ -170,7 +169,6 @@ func (s *MmctlE2ETestSuite) TestExtractJobListCmdF() {
Data: map[string]string{},
})
s.Require().Nil(appErr)
job2.Logger = nil
time.Sleep(time.Millisecond)
@@ -179,7 +177,6 @@ func (s *MmctlE2ETestSuite) TestExtractJobListCmdF() {
Data: map[string]string{},
})
s.Require().Nil(appErr)
job3.Logger = nil
time.Sleep(time.Millisecond)

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

@@ -256,7 +256,6 @@ func (s *MmctlE2ETestSuite) TestImportJobShowCmdF() {
Data: map[string]string{"import_file": "import1.zip"},
})
s.Require().Nil(appErr)
job.Logger = nil
s.Run("no permissions", func() {
printer.Clean()
@@ -350,7 +349,6 @@ func (s *MmctlE2ETestSuite) TestImportJobListCmdF() {
Data: map[string]string{"import_file": "import2.zip"},
})
s.Require().Nil(appErr)
job2.Logger = nil
time.Sleep(time.Millisecond)
@@ -359,7 +357,6 @@ func (s *MmctlE2ETestSuite) TestImportJobListCmdF() {
Data: map[string]string{"import_file": "import3.zip"},
})
s.Require().Nil(appErr)
job3.Logger = nil
err := importJobListCmdF(c, cmd, nil)
s.Require().Nil(err)

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

@@ -48,7 +48,7 @@ func MakeWorker(jobServer *jobs.JobServer, engine *bleveengine.BleveEngine) *Ble
stopped: make(chan bool, 1),
jobs: make(chan model.Job),
jobServer: jobServer,
logger: jobServer.Logger().With(mlog.String("workername", workerName)),
logger: jobServer.Logger().With(mlog.String("worker_name", workerName)),
engine: engine,
}
}
@@ -114,9 +114,6 @@ func (worker *BleveIndexerWorker) Run() {
worker.logger.Debug("Worker: Received stop signal")
return
case job := <-worker.jobs:
job.Logger = job.Logger.With(mlog.String("workername", worker.name))
job.Logger.Debug("Worker: Received a new candidate job.")
worker.DoJob(&job)
}
}
@@ -133,21 +130,24 @@ func (worker *BleveIndexerWorker) Stop() {
}
func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
logger := worker.logger.With(jobs.JobLoggerFields(job)...)
logger.Debug("Worker: Received a new candidate job.")
claimed, err := worker.jobServer.ClaimJob(job)
if err != nil {
job.Logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(err))
logger.Warn("Worker: Error occurred while trying to claim job", mlog.Err(err))
return
}
if !claimed {
return
}
job.Logger.Info("Worker: Indexing job claimed by worker")
logger.Info("Worker: Indexing job claimed by worker")
if !worker.engine.IsActive() {
appError := model.NewAppError("BleveIndexerWorker", "bleveengine.indexer.do_job.engine_inactive", nil, "", http.StatusInternalServerError)
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("Worker: Failed to run job as ")
logger.Error("Worker: Failed to run job as ")
}
return
}
@@ -166,10 +166,10 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
if startString, ok := job.Data["start_time"]; ok {
startInt, err := strconv.ParseInt(startString, 10, 64)
if err != nil {
job.Logger.Error("Worker: Failed to parse start_time for job", mlog.String("start_time", startString), mlog.Err(err))
logger.Error("Worker: Failed to parse start_time for job", mlog.String("start_time", startString), mlog.Err(err))
appError := model.NewAppError("BleveIndexerWorker", "bleveengine.indexer.do_job.parse_start_time.error", nil, "", http.StatusInternalServerError).Wrap(err)
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err), mlog.NamedErr("set_error", appError))
logger.Error("Worker: Failed to set job error", mlog.Err(err), mlog.NamedErr("set_error", appError))
}
return
}
@@ -179,10 +179,10 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
// A user or a channel may be created before any post.
oldestEntityCreationTime, err := worker.jobServer.Store.Post().GetOldestEntityCreationTime()
if err != nil {
job.Logger.Error("Worker: Failed to fetch oldest entity for job.", mlog.String("start_time", startString), mlog.Err(err))
logger.Error("Worker: Failed to fetch oldest entity for job.", mlog.String("start_time", startString), mlog.Err(err))
appError := model.NewAppError("BleveIndexerWorker", "bleveengine.indexer.do_job.get_oldest_entity.error", nil, "", http.StatusInternalServerError).Wrap(err)
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err), mlog.NamedErr("set_error", appError))
logger.Error("Worker: Failed to set job error", mlog.Err(err), mlog.NamedErr("set_error", appError))
}
return
}
@@ -193,10 +193,10 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
if endString, ok := job.Data["end_time"]; ok {
endInt, err := strconv.ParseInt(endString, 10, 64)
if err != nil {
job.Logger.Error("Worker: Failed to parse end_time for job", mlog.String("end_time", endString), mlog.Err(err))
logger.Error("Worker: Failed to parse end_time for job", mlog.String("end_time", endString), mlog.Err(err))
appError := model.NewAppError("BleveIndexerWorker", "bleveengine.indexer.do_job.parse_end_time.error", nil, "", http.StatusInternalServerError).Wrap(err)
if err := worker.jobServer.SetJobError(job, appError); err != nil {
job.Logger.Error("Worker: Failed to set job errorv", mlog.Err(err), mlog.NamedErr("set_error", appError))
logger.Error("Worker: Failed to set job errorv", mlog.Err(err), mlog.NamedErr("set_error", appError))
}
return
}
@@ -219,7 +219,7 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
// Counting all posts may fail or timeout when the posts table is large. If this happens, log a warning, but carry
// on with the indexing job anyway. The only issue is that the progress % reporting will be inaccurate.
if count, err := worker.jobServer.Store.Post().AnalyticsPostCount(&model.PostCountOptions{}); err != nil {
job.Logger.Warn("Worker: Failed to fetch total post count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
logger.Warn("Worker: Failed to fetch total post count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
progress.TotalPostsCount = estimatedPostCount
} else {
progress.TotalPostsCount = count
@@ -227,7 +227,7 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
// Same possible fail as above can happen when counting channels
if count, err := worker.jobServer.Store.Channel().AnalyticsTypeCount("", ""); err != nil {
job.Logger.Warn("Worker: Failed to fetch total channel count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
logger.Warn("Worker: Failed to fetch total channel count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
progress.TotalChannelsCount = estimatedChannelCount
} else {
progress.TotalChannelsCount = count
@@ -238,7 +238,7 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
IncludeBotAccounts: true, // This actually doesn't join with the bots table
// since ExcludeRegularUsers is set to false
}); err != nil {
job.Logger.Warn("Worker: Failed to fetch total user count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
logger.Warn("Worker: Failed to fetch total user count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
progress.TotalUsersCount = estimatedUserCount
} else {
progress.TotalUsersCount = count
@@ -247,7 +247,7 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
// Counting all files may fail or timeout when the file_info table is large. If this happens, log a warning, but carry
// on with the indexing job anyway. The only issue is that the progress % reporting will be inaccurate.
if count, err := worker.jobServer.Store.FileInfo().CountAll(); err != nil {
job.Logger.Warn("Worker: Failed to fetch total file info count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
logger.Warn("Worker: Failed to fetch total file info count for job. An estimated value will be used for progress reporting.", mlog.Err(err))
progress.TotalFilesCount = estimatedFilesCount
} else {
progress.TotalFilesCount = count
@@ -263,25 +263,25 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
for {
select {
case <-cancelWatcherChan:
job.Logger.Info("Worker: Indexing job has been canceled via CancellationWatcher")
logger.Info("Worker: Indexing job has been canceled via CancellationWatcher")
if err := worker.jobServer.SetJobCanceled(job); err != nil {
job.Logger.Error("Worker: Failed to mark job as cancelled", mlog.Err(err))
logger.Error("Worker: Failed to mark job as cancelled", mlog.Err(err))
}
return
case <-worker.stop:
job.Logger.Info("Worker: Indexing has been canceled via Worker Stop")
logger.Info("Worker: Indexing has been canceled via Worker Stop")
if err := worker.jobServer.SetJobCanceled(job); err != nil {
job.Logger.Error("Worker: Failed to mark job as canceled", mlog.Err(err))
logger.Error("Worker: Failed to mark job as canceled", mlog.Err(err))
}
return
case <-time.After(timeBetweenBatches):
var err *model.AppError
if progress, err = worker.IndexBatch(job.Logger, progress); err != nil {
job.Logger.Error("Worker: Failed to index batch for job", mlog.Err(err))
if progress, err = worker.IndexBatch(logger, progress); err != nil {
logger.Error("Worker: Failed to index batch for job", mlog.Err(err))
if err2 := worker.jobServer.SetJobError(job, err); err2 != nil {
job.Logger.Error("Worker: Failed to set job error", mlog.Err(err2), mlog.NamedErr("set_error", err))
logger.Error("Worker: Failed to set job error", mlog.Err(err2), mlog.NamedErr("set_error", err))
}
return
}
@@ -300,21 +300,21 @@ func (worker *BleveIndexerWorker) DoJob(job *model.Job) {
job.Data["end_time"] = strconv.FormatInt(progress.EndAtTime, 10)
if err := worker.jobServer.SetJobProgress(job, progress.CurrentProgress()); err != nil {
job.Logger.Error("Worker: Failed to set progress for job", mlog.Err(err))
logger.Error("Worker: Failed to set progress for job", mlog.Err(err))
if err2 := worker.jobServer.SetJobError(job, err); err2 != nil {
job.Logger.Error("Worker: Failed to set error for job", mlog.Err(err2), mlog.NamedErr("set_error", err))
logger.Error("Worker: Failed to set error for job", mlog.Err(err2), mlog.NamedErr("set_error", err))
}
return
}
if progress.IsDone() {
if err := worker.jobServer.SetJobSuccess(job); err != nil {
job.Logger.Error("Worker: Failed to set success for job", mlog.Err(err))
logger.Error("Worker: Failed to set success for job", mlog.Err(err))
if err2 := worker.jobServer.SetJobError(job, err); err2 != nil {
job.Logger.Error("Worker: Failed to set error for job", mlog.Err(err2), mlog.NamedErr("set_error", err))
logger.Error("Worker: Failed to set error for job", mlog.Err(err2), mlog.NamedErr("set_error", err))
}
}
job.Logger.Info("Worker: Indexing job finished successfully")
logger.Info("Worker: Indexing job finished successfully")
return
}
}

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

@@ -23,14 +23,12 @@ func TestBleveIndexer(t *testing.T) {
defer mockStore.AssertExpectations(t)
t.Run("Call GetOldestEntityCreationTime for the first indexing call", func(t *testing.T) {
logger := mlog.CreateConsoleTestLogger(t)
job := &model.Job{
Id: model.NewId(),
CreateAt: model.GetMillis(),
Status: model.JobStatusPending,
Type: model.JobTypeBlevePostIndexing,
}
job.InitLogger(logger)
mockStore.JobStore.On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).Return(true, nil)
mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(true, nil)
@@ -64,6 +62,7 @@ func TestBleveIndexer(t *testing.T) {
worker := &BleveIndexerWorker{
jobServer: jobServer,
engine: bleveEngine,
logger: mlog.CreateConsoleTestLogger(t),
}
worker.DoJob(job)

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

@@ -5,9 +5,6 @@ package model
import (
"net/http"
"time"
"github.com/mattermost/mattermost/server/public/shared/mlog"
)
const (
@@ -82,8 +79,6 @@ type Job struct {
Status string `json:"status"`
Progress int64 `json:"progress"`
Data StringMap `json:"data"`
Logger *mlog.Logger `json:"-"`
}
func (j *Job) Auditable() map[string]interface{} {
@@ -124,16 +119,6 @@ func (j *Job) IsValid() *AppError {
return nil
}
// InitLogger attaches an annotated logger to a Job.
// It should always be called after creating a new Job to ensure `Job.Logger` it set.
func (j *Job) InitLogger(logger mlog.LoggerIFace) {
j.Logger = logger.With(
mlog.String("job_id", j.Id),
mlog.String("job_type", j.Type),
mlog.String("create_at", time.UnixMilli(j.CreateAt).String()),
)
}
func (j *Job) LogClone() any {
return j.Auditable()
}