Change the App dependency with Server in migrations package (#14804)
Co-authored-by: Mattermod <mattermod@users.noreply.github.com>
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
a9ba052207
Коммит
e2e352e223
@@ -129,9 +129,6 @@ func (a *App) initJobs() {
|
|||||||
if jobsLdapSyncInterface != nil {
|
if jobsLdapSyncInterface != nil {
|
||||||
a.srv.Jobs.LdapSync = jobsLdapSyncInterface(a)
|
a.srv.Jobs.LdapSync = jobsLdapSyncInterface(a)
|
||||||
}
|
}
|
||||||
if jobsMigrationsInterface != nil {
|
|
||||||
a.srv.Jobs.Migrations = jobsMigrationsInterface(a)
|
|
||||||
}
|
|
||||||
if jobsPluginsInterface != nil {
|
if jobsPluginsInterface != nil {
|
||||||
a.srv.Jobs.Plugins = jobsPluginsInterface(a)
|
a.srv.Jobs.Plugins = jobsPluginsInterface(a)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -72,9 +72,9 @@ func RegisterJobsLdapSyncInterface(f func(*App) ejobs.LdapSyncInterface) {
|
|||||||
jobsLdapSyncInterface = f
|
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
|
jobsMigrationsInterface = f
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1290,4 +1290,7 @@ func (s *Server) initJobs() {
|
|||||||
if jobsBleveIndexerInterface != nil {
|
if jobsBleveIndexerInterface != nil {
|
||||||
s.Jobs.BleveIndexer = jobsBleveIndexerInterface(s)
|
s.Jobs.BleveIndexer = jobsBleveIndexerInterface(s)
|
||||||
}
|
}
|
||||||
|
if jobsMigrationsInterface != nil {
|
||||||
|
s.Jobs.Migrations = jobsMigrationsInterface(s)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -71,7 +71,7 @@ func (worker *Worker) runAdvancedPermissionsPhase2Migration(lastDone string) (bo
|
|||||||
|
|
||||||
if progress.CurrentTable == "TeamMembers" {
|
if progress.CurrentTable == "TeamMembers" {
|
||||||
// Run a TeamMembers migration batch.
|
// 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
|
return false, progress.ToJson(), err
|
||||||
} else {
|
} else {
|
||||||
if result == nil {
|
if result == nil {
|
||||||
@@ -86,7 +86,7 @@ func (worker *Worker) runAdvancedPermissionsPhase2Migration(lastDone string) (bo
|
|||||||
}
|
}
|
||||||
} else if progress.CurrentTable == "ChannelMembers" {
|
} else if progress.CurrentTable == "ChannelMembers" {
|
||||||
// Run a ChannelMembers migration batch.
|
// 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
|
return false, progress.ToJson(), err
|
||||||
} else {
|
} else {
|
||||||
if data == nil {
|
if data == nil {
|
||||||
|
|||||||
@@ -20,12 +20,12 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type MigrationsJobInterfaceImpl struct {
|
type MigrationsJobInterfaceImpl struct {
|
||||||
App *app.App
|
srv *app.Server
|
||||||
}
|
}
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
app.RegisterJobsMigrationsJobInterface(func(a *app.App) tjobs.MigrationsJobInterface {
|
app.RegisterJobsMigrationsJobInterface(func(s *app.Server) tjobs.MigrationsJobInterface {
|
||||||
return &MigrationsJobInterfaceImpl{a}
|
return &MigrationsJobInterfaceImpl{s}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -17,12 +17,12 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type Scheduler struct {
|
type Scheduler struct {
|
||||||
App *app.App
|
srv *app.Server
|
||||||
allMigrationsCompleted bool
|
allMigrationsCompleted bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MigrationsJobInterfaceImpl) MakeScheduler() model.Scheduler {
|
func (m *MigrationsJobInterfaceImpl) MakeScheduler() model.Scheduler {
|
||||||
return &Scheduler{m.App, false}
|
return &Scheduler{m.srv, false}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (scheduler *Scheduler) Name() string {
|
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).
|
// 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() {
|
for _, key := range MakeMigrationsList() {
|
||||||
state, job, err := GetMigrationState(key, scheduler.App.Srv().Store)
|
state, job, err := GetMigrationState(key, scheduler.srv.Store)
|
||||||
if err != nil {
|
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()))
|
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
|
return nil, nil
|
||||||
@@ -61,10 +61,10 @@ func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, las
|
|||||||
// Check the migration job isn't wedged.
|
// 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 {
|
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))
|
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()))
|
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
|
return nil, nil
|
||||||
@@ -77,7 +77,7 @@ func (scheduler *Scheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, las
|
|||||||
|
|
||||||
if state == MIGRATION_STATE_UNSCHEDULED {
|
if state == MIGRATION_STATE_UNSCHEDULED {
|
||||||
mlog.Debug("Scheduling a new job for migration.", mlog.String("scheduler", scheduler.Name()), mlog.String("migration_key", key))
|
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))
|
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,
|
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
|
return nil, err
|
||||||
} else {
|
} else {
|
||||||
return job, nil
|
return job, nil
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ type Worker struct {
|
|||||||
stopped chan bool
|
stopped chan bool
|
||||||
jobs chan model.Job
|
jobs chan model.Job
|
||||||
jobServer *jobs.JobServer
|
jobServer *jobs.JobServer
|
||||||
app *app.App
|
srv *app.Server
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *MigrationsJobInterfaceImpl) MakeWorker() model.Worker {
|
func (m *MigrationsJobInterfaceImpl) MakeWorker() model.Worker {
|
||||||
@@ -33,8 +33,8 @@ func (m *MigrationsJobInterfaceImpl) MakeWorker() model.Worker {
|
|||||||
stop: make(chan bool, 1),
|
stop: make(chan bool, 1),
|
||||||
stopped: make(chan bool, 1),
|
stopped: make(chan bool, 1),
|
||||||
jobs: make(chan model.Job),
|
jobs: make(chan model.Job),
|
||||||
jobServer: m.App.Srv().Jobs,
|
jobServer: m.srv.Jobs,
|
||||||
app: m.App,
|
srv: m.srv,
|
||||||
}
|
}
|
||||||
|
|
||||||
return &worker
|
return &worker
|
||||||
@@ -83,7 +83,7 @@ func (worker *Worker) DoJob(job *model.Job) {
|
|||||||
|
|
||||||
cancelCtx, cancelCancelWatcher := context.WithCancel(context.Background())
|
cancelCtx, cancelCancelWatcher := context.WithCancel(context.Background())
|
||||||
cancelWatcherChan := make(chan interface{}, 1)
|
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()
|
defer cancelCancelWatcher()
|
||||||
|
|
||||||
@@ -111,7 +111,7 @@ func (worker *Worker) DoJob(job *model.Job) {
|
|||||||
return
|
return
|
||||||
} else {
|
} else {
|
||||||
job.Data[JOB_DATA_KEY_MIGRATION_LAST_DONE] = progress
|
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()))
|
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)
|
worker.setJobError(job, err)
|
||||||
return
|
return
|
||||||
@@ -122,20 +122,20 @@ func (worker *Worker) DoJob(job *model.Job) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (worker *Worker) setJobSuccess(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()))
|
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)
|
worker.setJobError(job, err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (worker *Worker) setJobError(job *model.Job, appError *model.AppError) {
|
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()))
|
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) {
|
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()))
|
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 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
|
return false, "", saveErr
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user