[MM-32639] - Resend user invite emails (#17113)

Co-authored-by: Mattermod <mattermod@users.noreply.github.com>
Этот коммит содержится в:
Allan Guwatudde
2021-03-31 20:20:53 +03:00
коммит произвёл GitHub
родитель ab5925c4de
Коммит 489eaa4605
14 изменённых файлов: 325 добавлений и 0 удалений

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

@@ -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 {

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

@@ -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()
}

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

@@ -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) {

11
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
}

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

@@ -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
}

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

@@ -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"
)

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

@@ -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:
}
}
}
}
}

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

@@ -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}
})
}

42
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
}

148
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()))
}
}

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

@@ -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())
}

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

@@ -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

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

@@ -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")

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

@@ -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)
}