From 489eaa4605d19929dfe46600e1be1465eb872b46 Mon Sep 17 00:00:00 2001 From: Allan Guwatudde Date: Wed, 31 Mar 2021 20:20:53 +0300 Subject: [PATCH] [MM-32639] - Resend user invite emails (#17113) Co-authored-by: Mattermod --- api4/team.go | 18 +++ app/app.go | 4 + app/enterprise.go | 7 + einterfaces/jobs/resend_invitation_email.go | 11 ++ .../ResendInvitationEmailJobInterface.go | 47 ++++++ imports/placeholder.go | 3 + jobs/jobs_watcher.go | 7 + .../resend_invitation_email.go | 18 +++ jobs/resend_invitation_email/scheduler.go | 42 +++++ jobs/resend_invitation_email/worker.go | 148 ++++++++++++++++++ jobs/schedulers.go | 4 + jobs/server.go | 1 + jobs/workers.go | 13 ++ model/job.go | 2 + 14 files changed, 325 insertions(+) create mode 100644 einterfaces/jobs/resend_invitation_email.go create mode 100644 einterfaces/mocks/ResendInvitationEmailJobInterface.go create mode 100644 jobs/resend_invitation_email/resend_invitation_email.go create mode 100644 jobs/resend_invitation_email/scheduler.go create mode 100644 jobs/resend_invitation_email/worker.go diff --git a/api4/team.go b/api4/team.go index 6f7a1e2e66..78a381ae49 100644 --- a/api4/team.go +++ b/api4/team.go @@ -1220,6 +1220,24 @@ func inviteUsersToTeam(c *Context, w http.ResponseWriter, r *http.Request) { emailList, invitesOverLimit, _ = c.App.GetErrorListForEmailsOverLimit(emailList, cloudUserLimit) } } + + // we get the emailList after it has finished checks like the emails over the list + + scheduledAt := model.GetMillis() + jobData := map[string]string{ + "emailList": model.ArrayToJson(emailList), + "teamID": c.Params.TeamId, + "senderID": c.App.Session().UserId, + "scheduledAt": strconv.FormatInt(scheduledAt, 10), + } + + // we then manually schedule the job + _, e := c.App.Srv().Jobs.CreateJob(model.JOB_TYPE_RESEND_INVITATION_EMAIL, jobData) + if e != nil { + c.Err = model.NewAppError("Api4.inviteUsersToTeam", e.Id, nil, e.Error(), e.StatusCode) + return + } + var invitesWithError []*model.EmailInviteWithError var err *model.AppError if emailList != nil { diff --git a/app/app.go b/app/app.go index 6d65ef9245..d868815967 100644 --- a/app/app.go +++ b/app/app.go @@ -140,6 +140,10 @@ func (a *App) initJobs() { a.srv.Jobs.Cloud = jobsCloudInterface(a.srv) } + if jobsResendInvitationEmailInterface != nil { + a.srv.Jobs.ResendInvitationEmails = jobsResendInvitationEmailInterface(a) + } + a.srv.Jobs.InitWorkers() a.srv.Jobs.InitSchedulers() } diff --git a/app/enterprise.go b/app/enterprise.go index d4a5c84601..ba4fea08af 100644 --- a/app/enterprise.go +++ b/app/enterprise.go @@ -96,6 +96,13 @@ func RegisterJobsActiveUsersInterface(f func(*App) tjobs.ActiveUsersJobInterface jobsActiveUsersInterface = f } +var jobsResendInvitationEmailInterface func(*App) ejobs.ResendInvitationEmailJobInterface + +// RegisterJobsResendInvitationEmailInterface is used to register or initialize the jobsResendInvitationEmailInterface +func RegisterJobsResendInvitationEmailInterface(f func(*App) ejobs.ResendInvitationEmailJobInterface) { + jobsResendInvitationEmailInterface = f +} + var jobsCloudInterface func(*Server) ejobs.CloudJobInterface func RegisterJobsCloudInterface(f func(*Server) ejobs.CloudJobInterface) { diff --git a/einterfaces/jobs/resend_invitation_email.go b/einterfaces/jobs/resend_invitation_email.go new file mode 100644 index 0000000000..3b4aefbb8c --- /dev/null +++ b/einterfaces/jobs/resend_invitation_email.go @@ -0,0 +1,11 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. +package jobs + +import "github.com/mattermost/mattermost-server/v5/model" + +// ResendInvitationEmailJobInterface defines the interface for the job to resend invitation emails +type ResendInvitationEmailJobInterface interface { + MakeWorker() model.Worker + MakeScheduler() model.Scheduler +} diff --git a/einterfaces/mocks/ResendInvitationEmailJobInterface.go b/einterfaces/mocks/ResendInvitationEmailJobInterface.go new file mode 100644 index 0000000000..f8a8ea3b0b --- /dev/null +++ b/einterfaces/mocks/ResendInvitationEmailJobInterface.go @@ -0,0 +1,47 @@ +// Code generated by mockery v1.0.0. DO NOT EDIT. + +// Regenerate this file using `make einterfaces-mocks`. + +package mocks + +import ( + model "github.com/mattermost/mattermost-server/v5/model" + mock "github.com/stretchr/testify/mock" +) + +// ResendInvitationEmailJobInterface is an autogenerated mock type for the ResendInvitationEmailJobInterface type +type ResendInvitationEmailJobInterface struct { + mock.Mock +} + +// MakeScheduler provides a mock function with given fields: +func (_m *ResendInvitationEmailJobInterface) MakeScheduler() model.Scheduler { + ret := _m.Called() + + var r0 model.Scheduler + if rf, ok := ret.Get(0).(func() model.Scheduler); ok { + r0 = rf() + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(model.Scheduler) + } + } + + return r0 +} + +// MakeWorker provides a mock function with given fields: +func (_m *ResendInvitationEmailJobInterface) MakeWorker() model.Worker { + ret := _m.Called() + + var r0 model.Worker + if rf, ok := ret.Get(0).(func() model.Worker); ok { + r0 = rf() + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(model.Worker) + } + } + + return r0 +} diff --git a/imports/placeholder.go b/imports/placeholder.go index 5b482e803a..a3d82e252c 100644 --- a/imports/placeholder.go +++ b/imports/placeholder.go @@ -33,4 +33,7 @@ import ( // This is a placeholder so this package can be imported in Team Edition when it will be otherwise empty. _ "github.com/mattermost/mattermost-server/v5/jobs/export_delete" + + // This is a placeholder so this package can be imported in Team Edition when it will be otherwise empty. + _ "github.com/mattermost/mattermost-server/v5/jobs/resend_invitation_email" ) diff --git a/jobs/jobs_watcher.go b/jobs/jobs_watcher.go index e5237c8800..b3fcfaf981 100644 --- a/jobs/jobs_watcher.go +++ b/jobs/jobs_watcher.go @@ -184,6 +184,13 @@ func (watcher *Watcher) PollAndNotify() { default: } } + } else if job.Type == model.JOB_TYPE_RESEND_INVITATION_EMAIL { + if watcher.workers.ResendInvitationEmail != nil { + select { + case watcher.workers.ResendInvitationEmail.JobChannel() <- *job: + default: + } + } } } } diff --git a/jobs/resend_invitation_email/resend_invitation_email.go b/jobs/resend_invitation_email/resend_invitation_email.go new file mode 100644 index 0000000000..4257cb5bc0 --- /dev/null +++ b/jobs/resend_invitation_email/resend_invitation_email.go @@ -0,0 +1,18 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. +package resend_invitation_email + +import ( + "github.com/mattermost/mattermost-server/v5/app" + ejobs "github.com/mattermost/mattermost-server/v5/einterfaces/jobs" +) + +type ResendInvitationEmailJobInterfaceImpl struct { + App *app.App +} + +func init() { + app.RegisterJobsResendInvitationEmailInterface(func(a *app.App) ejobs.ResendInvitationEmailJobInterface { + return &ResendInvitationEmailJobInterfaceImpl{a} + }) +} diff --git a/jobs/resend_invitation_email/scheduler.go b/jobs/resend_invitation_email/scheduler.go new file mode 100644 index 0000000000..579b0a7543 --- /dev/null +++ b/jobs/resend_invitation_email/scheduler.go @@ -0,0 +1,42 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. +package resend_invitation_email + +import ( + "time" + + "github.com/mattermost/mattermost-server/v5/app" + "github.com/mattermost/mattermost-server/v5/model" +) + +const ResendInvitationEmailJob = "ResendInvitationEmailJob" + +type ResendInvitationEmailScheduler struct { + App *app.App +} + +func (rse *ResendInvitationEmailJobInterfaceImpl) MakeScheduler() model.Scheduler { + return &ResendInvitationEmailScheduler{rse.App} +} + +func (s *ResendInvitationEmailScheduler) Name() string { + return ResendInvitationEmailJob + "Scheduler" +} + +func (s *ResendInvitationEmailScheduler) JobType() string { + return model.JOB_TYPE_RESEND_INVITATION_EMAIL +} + +func (s *ResendInvitationEmailScheduler) Enabled(cfg *model.Config) bool { + return *cfg.ServiceSettings.EnableEmailInvitations +} + +func (s *ResendInvitationEmailScheduler) NextScheduleTime(cfg *model.Config, now time.Time, pendingJobs bool, lastSuccessfulJob *model.Job) *time.Time { + t := time.Now().Add(5 * time.Second) + return &t +} + +func (s *ResendInvitationEmailScheduler) ScheduleJob(cfg *model.Config, pendingJobs bool, lastSuccessfulJob *model.Job) (*model.Job, *model.AppError) { + // noop because we manually schedule the job in api4.inviteUsersToTeam handler + return nil, nil +} diff --git a/jobs/resend_invitation_email/worker.go b/jobs/resend_invitation_email/worker.go new file mode 100644 index 0000000000..fd4e266d26 --- /dev/null +++ b/jobs/resend_invitation_email/worker.go @@ -0,0 +1,148 @@ +// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved. +// See LICENSE.txt for license information. +package resend_invitation_email + +import ( + "encoding/json" + "net/http" + "os" + "strconv" + + "github.com/mattermost/mattermost-server/v5/app" + "github.com/mattermost/mattermost-server/v5/model" + "github.com/mattermost/mattermost-server/v5/shared/mlog" +) + +const TwentyFourHoursInMillis int64 = 86400000 + +type ResendInvitationEmailWorker struct { + name string + stop chan bool + stopped chan bool + jobs chan model.Job + App *app.App +} + +func (rse *ResendInvitationEmailJobInterfaceImpl) MakeWorker() model.Worker { + worker := ResendInvitationEmailWorker{ + name: ResendInvitationEmailJob, + stop: make(chan bool, 1), + stopped: make(chan bool, 1), + jobs: make(chan model.Job), + App: rse.App, + } + return &worker +} + +func (rseworker *ResendInvitationEmailWorker) Run() { + mlog.Debug("Worker started", mlog.String("worker", rseworker.name)) + + defer func() { + mlog.Debug("Worker finished", mlog.String("worker", rseworker.name)) + rseworker.stopped <- true + }() + + for { + select { + case <-rseworker.stop: + mlog.Debug("Worker received stop signal", mlog.String("worker", rseworker.name)) + return + case job := <-rseworker.jobs: + mlog.Debug("Worker received a new candidate job.", mlog.String("worker", rseworker.name)) + rseworker.DoJob(&job) + } + } +} + +func (rseworker *ResendInvitationEmailWorker) Stop() { + mlog.Debug("Worker stopping", mlog.String("worker", rseworker.name)) + rseworker.stop <- true + <-rseworker.stopped +} + +func (rseworker *ResendInvitationEmailWorker) JobChannel() chan<- model.Job { + return rseworker.jobs +} + +func (rseworker *ResendInvitationEmailWorker) cleanEmailData(emailStringData string) ([]string, error) { + // emailStringData looks like this ["user1@gmail.com","user2@gmail.com"] + emails := []string{} + err := json.Unmarshal([]byte(emailStringData), &emails) + if err != nil { + return nil, err + } + + return emails, nil +} + +func (rseworker *ResendInvitationEmailWorker) removeAlreadyJoined(teamID string, emailList []string) []string { + var notJoinedYet []string + for _, email := range emailList { + // check if the user with this email is on the system already + user, appErr := rseworker.App.GetUserByEmail(email) + if appErr != nil { + notJoinedYet = append(notJoinedYet, email) + continue + } + // now we check if they are part of the team already + userID := []string{user.Id} + members, appErr := rseworker.App.GetTeamMembersByIds(teamID, userID, nil) + if len(members) == 0 || appErr != nil { + notJoinedYet = append(notJoinedYet, email) + } + } + + return notJoinedYet +} + +func (rseworker *ResendInvitationEmailWorker) DoJob(job *model.Job) { + scheduledAt, _ := strconv.ParseInt(job.Data["scheduledAt"], 10, 64) + now := model.GetMillis() + + elapsedTimeSinceSchedule := now - scheduledAt + + var DurationInMillis int64 + + duration := os.Getenv("MM_RESEND_INVITATION_EMAIL_JOB_DURATION") + + DurationInMillis, parseError := strconv.ParseInt(duration, 10, 64) + if parseError != nil { + // default to 24 hours + DurationInMillis = TwentyFourHoursInMillis + } + + if elapsedTimeSinceSchedule > DurationInMillis { + teamID := job.Data["teamID"] + emailListData := job.Data["emailList"] + + emailList, err := rseworker.cleanEmailData(emailListData) + if err != nil { + appErr := model.NewAppError("worker: "+rseworker.name, "job_id: "+job.Id, nil, err.Error(), http.StatusInternalServerError) + mlog.Error("Worker: Failed to clean emails string data", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", appErr.Error())) + rseworker.setJobError(job, appErr) + } + + emailList = rseworker.removeAlreadyJoined(teamID, emailList) + + _, appErr := rseworker.App.InviteNewUsersToTeamGracefully(emailList, teamID, job.Data["senderID"]) + if appErr != nil { + mlog.Error("Worker: Failed to send emails", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", appErr.Error())) + rseworker.setJobError(job, appErr) + } + rseworker.setJobSuccess(job) + } + +} + +func (rseworker *ResendInvitationEmailWorker) setJobSuccess(job *model.Job) { + if err := rseworker.App.Srv().Jobs.SetJobSuccess(job); err != nil { + mlog.Error("Worker: Failed to set success for job", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error())) + rseworker.setJobError(job, err) + } +} + +func (rseworker *ResendInvitationEmailWorker) setJobError(job *model.Job, appError *model.AppError) { + if err := rseworker.App.Srv().Jobs.SetJobError(job, appError); err != nil { + mlog.Error("Worker: Failed to set job error", mlog.String("worker", rseworker.name), mlog.String("job_id", job.Id), mlog.String("error", err.Error())) + } +} diff --git a/jobs/schedulers.go b/jobs/schedulers.go index 4c829ac9a2..7a581d0f18 100644 --- a/jobs/schedulers.go +++ b/jobs/schedulers.go @@ -89,6 +89,10 @@ func (srv *JobServer) InitSchedulers() error { schedulers.schedulers = append(schedulers.schedulers, cloudInterface.MakeScheduler()) } + if resendInvitationEmailInterface := srv.ResendInvitationEmails; resendInvitationEmailInterface != nil { + schedulers.schedulers = append(schedulers.schedulers, resendInvitationEmailInterface.MakeScheduler()) + } + if importDeleteInterface := srv.ImportDelete; importDeleteInterface != nil { schedulers.schedulers = append(schedulers.schedulers, importDeleteInterface.MakeScheduler()) } diff --git a/jobs/server.go b/jobs/server.go index 7e3a352421..5ccc3dc46d 100644 --- a/jobs/server.go +++ b/jobs/server.go @@ -35,6 +35,7 @@ type JobServer struct { ExportProcess tjobs.ExportProcessInterface ExportDelete tjobs.ExportDeleteInterface Cloud ejobs.CloudJobInterface + ResendInvitationEmails ejobs.ResendInvitationEmailJobInterface // mut is used to protect the following fields from concurrent access. mut sync.Mutex diff --git a/jobs/workers.go b/jobs/workers.go index a7d1bfb995..5bd5d1567f 100644 --- a/jobs/workers.go +++ b/jobs/workers.go @@ -31,6 +31,7 @@ type Workers struct { ExportProcess model.Worker ExportDelete model.Worker Cloud model.Worker + ResendInvitationEmail model.Worker listenerId string running bool @@ -119,6 +120,10 @@ func (srv *JobServer) InitWorkers() error { workers.Cloud = cloudInterface.MakeWorker() } + if resendInvitationEmailInterface := srv.ResendInvitationEmails; resendInvitationEmailInterface != nil { + workers.ResendInvitationEmail = resendInvitationEmailInterface.MakeWorker() + } + srv.workers = workers return nil @@ -192,6 +197,10 @@ func (workers *Workers) Start() { go workers.Cloud.Run() } + if workers.ResendInvitationEmail != nil { + go workers.ResendInvitationEmail.Run() + } + go workers.Watcher.Start() workers.listenerId = workers.ConfigService.AddConfigListener(workers.handleConfigChange) @@ -321,6 +330,10 @@ func (workers *Workers) Stop() { workers.Cloud.Stop() } + if workers.ResendInvitationEmail != nil { + workers.ResendInvitationEmail.Stop() + } + workers.running = false mlog.Info("Stopped workers") diff --git a/model/job.go b/model/job.go index 95fe972b5e..9b3274cf70 100644 --- a/model/job.go +++ b/model/job.go @@ -27,6 +27,7 @@ const ( JOB_TYPE_EXPORT_PROCESS = "export_process" JOB_TYPE_EXPORT_DELETE = "export_delete" JOB_TYPE_CLOUD = "cloud" + JOB_TYPE_RESEND_INVITATION_EMAIL = "resend_invitation_email" JOB_STATUS_PENDING = "pending" JOB_STATUS_IN_PROGRESS = "in_progress" @@ -75,6 +76,7 @@ func (j *Job) IsValid() *AppError { case JOB_TYPE_EXPORT_PROCESS: case JOB_TYPE_EXPORT_DELETE: case JOB_TYPE_CLOUD: + case JOB_TYPE_RESEND_INVITATION_EMAIL: default: return NewAppError("Job.IsValid", "model.job.is_valid.type.app_error", nil, "id="+j.Id, http.StatusBadRequest) }