diff --git a/app/app.go b/app/app.go index 45da3ce875..8dcd561633 100644 --- a/app/app.go +++ b/app/app.go @@ -129,9 +129,6 @@ func (a *App) initJobs() { if jobsLdapSyncInterface != nil { a.srv.Jobs.LdapSync = jobsLdapSyncInterface(a) } - if jobsMigrationsInterface != nil { - a.srv.Jobs.Migrations = jobsMigrationsInterface(a) - } if jobsPluginsInterface != nil { a.srv.Jobs.Plugins = jobsPluginsInterface(a) } diff --git a/app/enterprise.go b/app/enterprise.go index db903d5e32..6cca286df9 100644 --- a/app/enterprise.go +++ b/app/enterprise.go @@ -72,9 +72,9 @@ func RegisterJobsLdapSyncInterface(f func(*App) ejobs.LdapSyncInterface) { jobsLdapSyncInterface = f } -var jobsMigrationsInterface func(*App) tjobs.MigrationsJobInterface +var jobsMigrationsInterface func(*Server) tjobs.MigrationsJobInterface -func RegisterJobsMigrationsJobInterface(f func(*App) tjobs.MigrationsJobInterface) { +func RegisterJobsMigrationsJobInterface(f func(*Server) tjobs.MigrationsJobInterface) { jobsMigrationsInterface = f } diff --git a/app/server.go b/app/server.go index 94667c36b1..1407b991ac 100644 --- a/app/server.go +++ b/app/server.go @@ -1290,4 +1290,7 @@ func (s *Server) initJobs() { if jobsBleveIndexerInterface != nil { s.Jobs.BleveIndexer = jobsBleveIndexerInterface(s) } + if jobsMigrationsInterface != nil { + s.Jobs.Migrations = jobsMigrationsInterface(s) + } } diff --git a/migrations/advanced_permissions_phase_2.go b/migrations/advanced_permissions_phase_2.go index cdbccb6ace..99e4cd5301 100644 --- a/migrations/advanced_permissions_phase_2.go +++ b/migrations/advanced_permissions_phase_2.go @@ -71,7 +71,7 @@ func (worker *Worker) runAdvancedPermissionsPhase2Migration(lastDone string) (bo if progress.CurrentTable == "TeamMembers" { // Run a TeamMembers migration batch. - if result, err := worker.app.Srv().Store.Team().MigrateTeamMembers(progress.LastTeamId, progress.LastUserId); err != nil { + if result, err := worker.srv.Store.Team().MigrateTeamMembers(progress.LastTeamId, progress.LastUserId); err != nil { return false, progress.ToJson(), err } else { if result == nil { @@ -86,7 +86,7 @@ func (worker *Worker) runAdvancedPermissionsPhase2Migration(lastDone string) (bo } } else if progress.CurrentTable == "ChannelMembers" { // Run a ChannelMembers migration batch. - if data, err := worker.app.Srv().Store.Channel().MigrateChannelMembers(progress.LastChannelId, progress.LastUserId); err != nil { + if data, err := worker.srv.Store.Channel().MigrateChannelMembers(progress.LastChannelId, progress.LastUserId); err != nil { return false, progress.ToJson(), err } else { if data == nil { diff --git a/migrations/migrations.go b/migrations/migrations.go index d738edd5fc..7dfc9117bf 100644 --- a/migrations/migrations.go +++ b/migrations/migrations.go @@ -20,12 +20,12 @@ const ( ) type MigrationsJobInterfaceImpl struct { - App *app.App + srv *app.Server } func init() { - app.RegisterJobsMigrationsJobInterface(func(a *app.App) tjobs.MigrationsJobInterface { - return &MigrationsJobInterfaceImpl{a} + app.RegisterJobsMigrationsJobInterface(func(s *app.Server) tjobs.MigrationsJobInterface { + return &MigrationsJobInterfaceImpl{s} }) } diff --git a/migrations/scheduler.go b/migrations/scheduler.go index 807f19e2ff..acbb0bac57 100644 --- a/migrations/scheduler.go +++ b/migrations/scheduler.go @@ -17,12 +17,12 @@ const ( ) type Scheduler struct { - App *app.App + srv *app.Server allMigrationsCompleted bool } func (m *MigrationsJobInterfaceImpl) MakeScheduler() model.Scheduler { - return &Scheduler{m.App, false} + return &Scheduler{m.srv, false} } func (scheduler *Scheduler) Name() string { @@ -51,7 +51,7 @@ func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, las // 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.App.Srv().Store) + state, job, err := GetMigrationState(key, scheduler.srv.Store) if err != nil { mlog.Error("Failed to determine status of migration: ", mlog.String("scheduler", scheduler.Name()), mlog.String("migration_key", key), mlog.String("error", err.Error())) return nil, nil @@ -61,10 +61,10 @@ func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, las // Check the migration job isn't wedged. if job != nil && job.LastActivityAt < model.GetMillis()-MIGRATION_JOB_WEDGED_TIMEOUT_MILLISECONDS && job.CreateAt < model.GetMillis()-MIGRATION_JOB_WEDGED_TIMEOUT_MILLISECONDS { mlog.Warn("Job appears to be wedged. Rescheduling another instance.", mlog.String("scheduler", scheduler.Name()), mlog.String("wedged_job_id", job.Id), mlog.String("migration_key", key)) - if err := scheduler.App.Srv().Jobs.SetJobError(job, nil); err != nil { + if err := scheduler.srv.Jobs.SetJobError(job, nil); err != nil { mlog.Error("Worker: Failed to set job error", mlog.String("scheduler", scheduler.Name()), mlog.String("job_id", job.Id), mlog.String("error", err.Error())) } - return scheduler.createJob(key, job, scheduler.App.Srv().Store) + return scheduler.createJob(key, job, scheduler.srv.Store) } return nil, nil @@ -77,7 +77,7 @@ func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, las if state == MIGRATION_STATE_UNSCHEDULED { mlog.Debug("Scheduling a new job for migration.", mlog.String("scheduler", scheduler.Name()), mlog.String("migration_key", key)) - return scheduler.createJob(key, job, scheduler.App.Srv().Store) + return scheduler.createJob(key, job, scheduler.srv.Store) } mlog.Error("Unknown migration state. Not doing anything.", mlog.String("migration_state", state)) @@ -102,7 +102,7 @@ func (scheduler *Scheduler) createJob(migrationKey string, lastJob *model.Job, s JOB_DATA_KEY_MIGRATION_LAST_DONE: lastDone, } - if job, err := scheduler.App.Srv().Jobs.CreateJob(model.JOB_TYPE_MIGRATIONS, data); err != nil { + if job, err := scheduler.srv.Jobs.CreateJob(model.JOB_TYPE_MIGRATIONS, data); err != nil { return nil, err } else { return job, nil diff --git a/migrations/worker.go b/migrations/worker.go index 1b20f13c98..f1067aaea7 100644 --- a/migrations/worker.go +++ b/migrations/worker.go @@ -24,7 +24,7 @@ type Worker struct { stopped chan bool jobs chan model.Job jobServer *jobs.JobServer - app *app.App + srv *app.Server } func (m *MigrationsJobInterfaceImpl) MakeWorker() model.Worker { @@ -33,8 +33,8 @@ func (m *MigrationsJobInterfaceImpl) MakeWorker() model.Worker { stop: make(chan bool, 1), stopped: make(chan bool, 1), jobs: make(chan model.Job), - jobServer: m.App.Srv().Jobs, - app: m.App, + jobServer: m.srv.Jobs, + srv: m.srv, } return &worker @@ -83,7 +83,7 @@ func (worker *Worker) DoJob(job *model.Job) { cancelCtx, cancelCancelWatcher := context.WithCancel(context.Background()) cancelWatcherChan := make(chan interface{}, 1) - go worker.app.Srv().Jobs.CancellationWatcher(cancelCtx, job.Id, cancelWatcherChan) + go worker.srv.Jobs.CancellationWatcher(cancelCtx, job.Id, cancelWatcherChan) defer cancelCancelWatcher() @@ -111,7 +111,7 @@ func (worker *Worker) DoJob(job *model.Job) { return } else { job.Data[JOB_DATA_KEY_MIGRATION_LAST_DONE] = progress - if err := worker.app.Srv().Jobs.UpdateInProgressJobData(job); err != nil { + if err := worker.srv.Jobs.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())) worker.setJobError(job, err) return @@ -122,20 +122,20 @@ func (worker *Worker) DoJob(job *model.Job) { } func (worker *Worker) setJobSuccess(job *model.Job) { - if err := worker.app.Srv().Jobs.SetJobSuccess(job); err != nil { + if err := worker.srv.Jobs.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())) worker.setJobError(job, err) } } func (worker *Worker) setJobError(job *model.Job, appError *model.AppError) { - if err := worker.app.Srv().Jobs.SetJobError(job, appError); err != nil { + if err := worker.srv.Jobs.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())) } } func (worker *Worker) setJobCanceled(job *model.Job) { - if err := worker.app.Srv().Jobs.SetJobCanceled(job); err != nil { + if err := worker.srv.Jobs.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())) } } @@ -157,7 +157,7 @@ func (worker *Worker) runMigration(key string, lastDone string) (bool, string, * } if done { - if saveErr := worker.app.Srv().Store.System().Save(&model.System{Name: key, Value: "true"}); saveErr != nil { + if saveErr := worker.srv.Store.System().Save(&model.System{Name: key, Value: "true"}); saveErr != nil { return false, "", saveErr } }