diff --git a/jobs/active_users/worker.go b/jobs/active_users/worker.go index 8f562dd6a5..7dc9abcfdf 100644 --- a/jobs/active_users/worker.go +++ b/jobs/active_users/worker.go @@ -19,6 +19,8 @@ func MakeWorker(jobServer *jobs.JobServer, store store.Store, getMetrics func() return *cfg.MetricsSettings.Enable } execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + count, err := store.User().Count(model.UserCountOptions{IncludeDeleted: false}) if err != nil { return err diff --git a/jobs/expirynotify/worker.go b/jobs/expirynotify/worker.go index 98b9a0b47b..af9a09a8e7 100644 --- a/jobs/expirynotify/worker.go +++ b/jobs/expirynotify/worker.go @@ -17,6 +17,8 @@ func MakeWorker(jobServer *jobs.JobServer, notifySessionsExpired func() error) m return *cfg.ServiceSettings.ExtendSessionLengthWithActivity } execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + return notifySessionsExpired() } return jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled) diff --git a/jobs/export_delete/worker.go b/jobs/export_delete/worker.go index 8a8c509426..39cf155978 100644 --- a/jobs/export_delete/worker.go +++ b/jobs/export_delete/worker.go @@ -28,6 +28,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker { return *cfg.ExportSettings.Directory != "" && *cfg.ExportSettings.RetentionDays > 0 } execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + exportPath := *app.Config().ExportSettings.Directory retentionTime := time.Duration(*app.Config().ExportSettings.RetentionDays) * 24 * time.Hour exports, appErr := app.ListDirectory(exportPath) diff --git a/jobs/export_process/worker.go b/jobs/export_process/worker.go index de5ee69489..a9deaa8634 100644 --- a/jobs/export_process/worker.go +++ b/jobs/export_process/worker.go @@ -26,6 +26,8 @@ type AppIface interface { func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker { isEnabled := func(cfg *model.Config) bool { return true } execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + opts := model.BulkExportOpts{ CreateArchive: true, } diff --git a/jobs/extract_content/worker.go b/jobs/extract_content/worker.go index c2eb9ad0af..714f7ff5c9 100644 --- a/jobs/extract_content/worker.go +++ b/jobs/extract_content/worker.go @@ -30,6 +30,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, store store.Store) mode return true } execute := func(job *model.Job) error { + jobServer.HandleJobPanic(job) + var err error var fromTS int64 = 0 var toTS int64 = model.GetMillis() diff --git a/jobs/import_delete/worker.go b/jobs/import_delete/worker.go index 6dc46b4ae8..429248e8ff 100644 --- a/jobs/import_delete/worker.go +++ b/jobs/import_delete/worker.go @@ -30,6 +30,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface, s store.Store) model.Wo return *cfg.ImportSettings.Directory != "" && *cfg.ImportSettings.RetentionDays > 0 } execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + importPath := *app.Config().ImportSettings.Directory retentionTime := time.Duration(*app.Config().ImportSettings.RetentionDays) * 24 * time.Hour imports, appErr := app.ListDirectory(importPath) diff --git a/jobs/import_process/worker.go b/jobs/import_process/worker.go index 5747d04b93..c5414ff4a0 100644 --- a/jobs/import_process/worker.go +++ b/jobs/import_process/worker.go @@ -38,6 +38,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker { return true } execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + importFileName, ok := job.Data["import_file"] if !ok { return model.NewAppError("ImportProcessWorker", "import_process.worker.do_job.missing_file", nil, "", http.StatusBadRequest) diff --git a/jobs/jobs.go b/jobs/jobs.go index e23ca766c6..3ec15cc0fe 100644 --- a/jobs/jobs.go +++ b/jobs/jobs.go @@ -6,7 +6,10 @@ package jobs import ( "context" "errors" + "fmt" "net/http" + "runtime/pprof" + "strings" "time" "github.com/mattermost/mattermost-server/v6/model" @@ -179,6 +182,31 @@ func (srv *JobServer) UpdateInProgressJobData(job *model.Job) *model.AppError { return nil } +// 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) { + r := recover() + if r == nil { + return + } + + 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())) + + rerr, ok := r.(error) + if !ok { + rerr = fmt.Errorf("job panic: %v", r) + } + + 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)) + } + + panic(r) +} + func (srv *JobServer) RequestCancellation(jobId string) *model.AppError { updated, err := srv.Store.Job().UpdateStatusOptimistically(jobId, model.JobStatusPending, model.JobStatusCanceled) if err != nil { diff --git a/jobs/jobs_test.go b/jobs/jobs_test.go index 04c315d85e..79a484a735 100644 --- a/jobs/jobs_test.go +++ b/jobs/jobs_test.go @@ -5,6 +5,7 @@ package jobs import ( "errors" + "fmt" "net/http" "testing" @@ -494,6 +495,65 @@ func TestUpdateInProgressJobData(t *testing.T) { }) } +func TestHandleJobPanic(t *testing.T) { + t.Run("no panic", func(t *testing.T) { + jobServer, _, _ := makeJobServer(t) + + job := &model.Job{ + Type: model.JobTypeImportProcess, + Status: model.JobStatusInProgress, + } + + f := func() { + defer jobServer.HandleJobPanic(job) + fmt.Println("OK") + } + + require.NotPanics(t, f) + require.Equal(t, model.JobStatusInProgress, job.Status) + }) + + t.Run("with panic string", func(t *testing.T) { + jobServer, mockStore, metrics := makeJobServer(t) + + job := &model.Job{ + Type: model.JobTypeImportProcess, + Status: model.JobStatusInProgress, + } + + f := func() { + defer jobServer.HandleJobPanic(job) + panic("not OK") + } + + mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(true, nil) + metrics.On("DecrementJobActive", model.JobTypeImportProcess) + + require.Panics(t, f) + require.Equal(t, model.JobStatusError, job.Status) + }) + + t.Run("with panic error", func(t *testing.T) { + jobServer, mockStore, metrics := makeJobServer(t) + + job := &model.Job{ + Type: model.JobTypeImportProcess, + Status: model.JobStatusInProgress, + } + + f := func() { + defer jobServer.HandleJobPanic(job) + panic(fmt.Errorf("not OK")) + } + + mockStore.JobStore.On("UpdateOptimistically", job, model.JobStatusInProgress).Return(true, nil) + metrics.On("DecrementJobActive", model.JobTypeImportProcess) + + require.Panics(t, f) + require.Equal(t, model.JobStatusError, job.Status) + }) +} + func TestRequestCancellation(t *testing.T) { t.Run("error cancelling", func(t *testing.T) { jobServer, mockStore, _ := makeJobServer(t) diff --git a/jobs/last_accessible_file/worker.go b/jobs/last_accessible_file/worker.go index 177c1016b0..5cb18e4f7c 100644 --- a/jobs/last_accessible_file/worker.go +++ b/jobs/last_accessible_file/worker.go @@ -20,7 +20,9 @@ func MakeWorker(jobServer *jobs.JobServer, license *model.License, app AppIface) isEnabled := func(_ *model.Config) bool { return license != nil && *license.Features.Cloud } - execute := func(_ *model.Job) error { + execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + return app.ComputeLastAccessibleFileTime() } worker := jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled) diff --git a/jobs/last_accessible_post/worker.go b/jobs/last_accessible_post/worker.go index 2fc09c38b3..266c2b0da8 100644 --- a/jobs/last_accessible_post/worker.go +++ b/jobs/last_accessible_post/worker.go @@ -20,7 +20,9 @@ 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(_ *model.Job) error { + execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + return app.ComputeLastAccessiblePostTime() } worker := jobs.NewSimpleWorker(JobName, jobServer, execute, isEnabled) diff --git a/jobs/migrations/worker.go b/jobs/migrations/worker.go index 482709399a..03fe9e2396 100644 --- a/jobs/migrations/worker.go +++ b/jobs/migrations/worker.go @@ -85,6 +85,8 @@ func (worker *Worker) IsEnabled(_ *model.Config) bool { } 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), diff --git a/jobs/notify_admin/worker.go b/jobs/notify_admin/worker.go index 283ce43adf..363e7628aa 100644 --- a/jobs/notify_admin/worker.go +++ b/jobs/notify_admin/worker.go @@ -21,7 +21,9 @@ func MakeUpgradeNotifyWorker(jobServer *jobs.JobServer, license *model.License, isEnabled := func(_ *model.Config) bool { return license != nil && license.Features != nil && *license.Features.Cloud } - execute := func(_ *model.Job) error { + execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + appErr := app.DoCheckForAdminNotifications(false) if appErr != nil { return appErr @@ -37,7 +39,9 @@ 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(_ *model.Job) error { + execute := func(job *model.Job) error { + defer jobServer.HandleJobPanic(job) + appErr := app.DoCheckForAdminNotifications(true) if appErr != nil { return appErr diff --git a/jobs/product_notices/worker.go b/jobs/product_notices/worker.go index 20935dfc80..8fe8d3a949 100644 --- a/jobs/product_notices/worker.go +++ b/jobs/product_notices/worker.go @@ -20,6 +20,8 @@ func MakeWorker(jobServer *jobs.JobServer, app AppIface) model.Worker { return *cfg.AnnouncementSettings.AdminNoticesEnabled || *cfg.AnnouncementSettings.UserNoticesEnabled } execute := func(job *model.Job) error { + 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)) return err diff --git a/jobs/resend_invitation_email/worker.go b/jobs/resend_invitation_email/worker.go index 67fd37fa65..3c9c764fb9 100644 --- a/jobs/resend_invitation_email/worker.go +++ b/jobs/resend_invitation_email/worker.go @@ -85,6 +85,8 @@ func (rseworker *ResendInvitationEmailWorker) JobChannel() chan<- model.Job { } func (rseworker *ResendInvitationEmailWorker) DoJob(job *model.Job) { + defer rseworker.jobServer.HandleJobPanic(job) + elapsedTimeSinceSchedule, DurationInMillis := rseworker.GetDurations(job) if elapsedTimeSinceSchedule > DurationInMillis { rseworker.ResendEmails(job, "48")