[MM-55726] Create batch report worker, add batch report job for exporting users to CSV (#25832)
* Split out migration logic and create generic BatchWorker * WIP * WIP * POC batch reporting * Oops * Job hookup * Working export to file * PR feedback * Merge'd * Fix error handling * Add API to start report, translations, couple fixes * Add DMs to send reports to users * Merge'd * Update types * A bit of cleanup * Some fixes * Add missing API doc * PR feedback * Fix generated * Fix bug with post creation * PR feedback * Add some tests * PR feedback * Fix lint * Some test changes * Fix tests * Add comment to explain why we forcibly stop * Rework of some tests * Batch report test * Restrict batch exports to Pro and Enterprise licenses * Fix erroneous comment --------- Co-authored-by: Mattermost Build <build@mattermost.com>
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
eac9a39677
Коммит
f7446d7443
@@ -5,7 +5,6 @@ package jobs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost/server/public/model"
|
||||
@@ -29,83 +28,57 @@ type BatchMigrationWorkerAppIFace interface {
|
||||
// server in order to retry a failed migration job. Refactoring the job infrastructure is left as
|
||||
// a future exercise.
|
||||
type BatchMigrationWorker struct {
|
||||
jobServer *JobServer
|
||||
logger mlog.LoggerIFace
|
||||
store store.Store
|
||||
app BatchMigrationWorkerAppIFace
|
||||
|
||||
stop chan struct{}
|
||||
stopped chan bool
|
||||
closed atomic.Bool
|
||||
jobs chan model.Job
|
||||
|
||||
migrationKey string
|
||||
timeBetweenBatches time.Duration
|
||||
doMigrationBatch func(data model.StringMap, store store.Store) (model.StringMap, bool, error)
|
||||
*BatchWorker
|
||||
app BatchMigrationWorkerAppIFace
|
||||
migrationKey string
|
||||
doMigrationBatch func(data model.StringMap, store store.Store) (model.StringMap, bool, error)
|
||||
}
|
||||
|
||||
// MakeBatchMigrationWorker creates a worker to process the given migration batch function.
|
||||
func MakeBatchMigrationWorker(jobServer *JobServer, store store.Store, app BatchMigrationWorkerAppIFace, migrationKey string, timeBetweenBatches time.Duration, doMigrationBatch func(data model.StringMap, store store.Store) (model.StringMap, bool, error)) model.Worker {
|
||||
func MakeBatchMigrationWorker(
|
||||
jobServer *JobServer,
|
||||
store store.Store,
|
||||
app BatchMigrationWorkerAppIFace,
|
||||
migrationKey string,
|
||||
timeBetweenBatches time.Duration,
|
||||
doMigrationBatch func(data model.StringMap, store store.Store) (model.StringMap, bool, error),
|
||||
) *BatchMigrationWorker {
|
||||
worker := &BatchMigrationWorker{
|
||||
jobServer: jobServer,
|
||||
logger: jobServer.Logger().With(mlog.String("worker_name", migrationKey)),
|
||||
store: store,
|
||||
app: app,
|
||||
stop: make(chan struct{}),
|
||||
stopped: make(chan bool, 1),
|
||||
jobs: make(chan model.Job),
|
||||
migrationKey: migrationKey,
|
||||
timeBetweenBatches: timeBetweenBatches,
|
||||
doMigrationBatch: doMigrationBatch,
|
||||
app: app,
|
||||
migrationKey: migrationKey,
|
||||
doMigrationBatch: doMigrationBatch,
|
||||
}
|
||||
worker.BatchWorker = MakeBatchWorker(jobServer, store, timeBetweenBatches, worker.doBatch)
|
||||
return worker
|
||||
}
|
||||
|
||||
// Run starts the worker dedicated to the unique migration batch job it will be given to process.
|
||||
func (worker *BatchMigrationWorker) Run() {
|
||||
worker.logger.Debug("Worker started")
|
||||
// We have to re-assign the stop channel again, because
|
||||
// it might happen that the job was restarted due to a config change.
|
||||
if worker.closed.CompareAndSwap(true, false) {
|
||||
worker.stop = make(chan struct{})
|
||||
func (worker *BatchMigrationWorker) doBatch(rctx *request.Context, job *model.Job) bool {
|
||||
// Ensure the cluster remains in sync, otherwise we restart the job to
|
||||
// ensure a complete migration. Technically, the cluster could go out of
|
||||
// sync briefly within a batch, but we accept that risk.
|
||||
if !worker.checkIsClusterInSync(rctx) {
|
||||
worker.logger.Warn("Worker: Resetting job")
|
||||
worker.resetJob(worker.logger, job)
|
||||
return true
|
||||
}
|
||||
|
||||
defer func() {
|
||||
worker.logger.Debug("Worker finished")
|
||||
worker.stopped <- true
|
||||
}()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-worker.stop:
|
||||
worker.logger.Debug("Worker received stop signal")
|
||||
return
|
||||
case job := <-worker.jobs:
|
||||
worker.DoJob(&job)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Stop interrupts the worker even if the migration has not yet completed.
|
||||
func (worker *BatchMigrationWorker) Stop() {
|
||||
// Set to close, and if already closed before, then return.
|
||||
if !worker.closed.CompareAndSwap(false, true) {
|
||||
return
|
||||
nextData, done, err := worker.doMigrationBatch(job.Data, worker.store)
|
||||
if err != nil {
|
||||
worker.logger.Error("Worker: Failed to do migration batch. Exiting", mlog.Err(err))
|
||||
worker.setJobError(worker.logger, job, model.NewAppError("doMigrationBatch", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
|
||||
return true
|
||||
} else if done {
|
||||
worker.logger.Info("Worker: Job is complete")
|
||||
worker.setJobSuccess(worker.logger, job)
|
||||
worker.markAsComplete()
|
||||
return true
|
||||
}
|
||||
|
||||
worker.logger.Debug("Worker stopping")
|
||||
close(worker.stop)
|
||||
<-worker.stopped
|
||||
}
|
||||
job.Data = nextData
|
||||
|
||||
// JobChannel is the means by which the jobs infrastructure provides the worker the job to execute.
|
||||
func (worker *BatchMigrationWorker) JobChannel() chan<- model.Job {
|
||||
return worker.jobs
|
||||
}
|
||||
|
||||
// IsEnabled is always true for batch migrations.
|
||||
func (worker *BatchMigrationWorker) IsEnabled(_ *model.Config) bool {
|
||||
return true
|
||||
// Migrations currently don't support reporting meaningful progress.
|
||||
worker.jobServer.SetJobProgress(job, 0)
|
||||
return false
|
||||
}
|
||||
|
||||
// checkIsClusterInSync returns true if all nodes in the cluster are running the same version,
|
||||
@@ -128,108 +101,6 @@ func (worker *BatchMigrationWorker) checkIsClusterInSync(rctx request.CTX) bool
|
||||
return true
|
||||
}
|
||||
|
||||
// DoJob executes the job picked up through the job channel.
|
||||
//
|
||||
// Note that this is a lot of distracting machinery here to claim the job, then double check the
|
||||
// status, and keep the status up to date in line with job infrastrcuture semantics. Unless an
|
||||
// error occurs, this worker should hold onto the job until its completed.
|
||||
func (worker *BatchMigrationWorker) DoJob(job *model.Job) {
|
||||
logger := worker.logger.With(mlog.Any("job", job))
|
||||
logger.Debug("Worker received a new candidate job.")
|
||||
defer worker.jobServer.HandleJobPanic(logger, job)
|
||||
|
||||
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
|
||||
logger.Warn("Worker experienced an error while trying to claim job", mlog.Err(err))
|
||||
return
|
||||
} else if !claimed {
|
||||
return
|
||||
}
|
||||
|
||||
c := request.EmptyContext(logger)
|
||||
var appErr *model.AppError
|
||||
|
||||
// We get the job again because ClaimJob changes the job status.
|
||||
job, appErr = worker.jobServer.GetJob(c, job.Id)
|
||||
if appErr != nil {
|
||||
worker.logger.Error("Worker: job execution error", mlog.Err(appErr))
|
||||
worker.setJobError(logger, job, appErr)
|
||||
return
|
||||
}
|
||||
|
||||
if job.Data == nil {
|
||||
job.Data = make(model.StringMap)
|
||||
}
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-worker.stop:
|
||||
logger.Info("Worker: Migration has been canceled via Worker Stop. Setting the job back to pending.")
|
||||
if err := worker.jobServer.SetJobPending(job); err != nil {
|
||||
worker.logger.Error("Worker: Failed to mark job as pending", mlog.Err(err))
|
||||
}
|
||||
return
|
||||
case <-time.After(worker.timeBetweenBatches):
|
||||
// Ensure the cluster remains in sync, otherwise we restart the job to
|
||||
// ensure a complete migration. Technically, the cluster could go out of
|
||||
// sync briefly within a batch, but we accept that risk.
|
||||
if !worker.checkIsClusterInSync(c) {
|
||||
worker.logger.Warn("Worker: Resetting job")
|
||||
worker.resetJob(logger, job)
|
||||
return
|
||||
}
|
||||
|
||||
nextData, done, err := worker.doMigrationBatch(job.Data, worker.store)
|
||||
if err != nil {
|
||||
worker.logger.Error("Worker: Failed to do migration batch. Exiting", mlog.Err(err))
|
||||
worker.setJobError(logger, job, model.NewAppError("doMigrationBatch", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
|
||||
return
|
||||
} else if done {
|
||||
logger.Info("Worker: Job is complete")
|
||||
worker.setJobSuccess(logger, job)
|
||||
worker.markAsComplete()
|
||||
return
|
||||
}
|
||||
|
||||
job.Data = nextData
|
||||
|
||||
// Migrations currently don't support reporting meaningful progress.
|
||||
worker.jobServer.SetJobProgress(job, 0)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// resetJob erases the data tracking the next batch to execute and returns the job status to
|
||||
// pending to allow the job infrastructure to requeue it.
|
||||
func (worker *BatchMigrationWorker) resetJob(logger mlog.LoggerIFace, job *model.Job) {
|
||||
job.Data = nil
|
||||
job.Progress = 0
|
||||
job.Status = model.JobStatusPending
|
||||
|
||||
if _, err := worker.store.Job().UpdateOptimistically(job, model.JobStatusInProgress); err != nil {
|
||||
worker.logger.Error("Worker: Failed to reset job data. May resume instead of restarting.", mlog.Err(err))
|
||||
}
|
||||
}
|
||||
|
||||
// setJobSuccess records the job as successful.
|
||||
func (worker *BatchMigrationWorker) setJobSuccess(logger mlog.LoggerIFace, job *model.Job) {
|
||||
if err := worker.jobServer.SetJobProgress(job, 100); err != nil {
|
||||
logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
|
||||
worker.setJobError(logger, job, err)
|
||||
}
|
||||
|
||||
if err := worker.jobServer.SetJobSuccess(job); err != nil {
|
||||
logger.Error("Worker: Failed to set success for job", mlog.Err(err))
|
||||
worker.setJobError(logger, job, err)
|
||||
}
|
||||
}
|
||||
|
||||
// setJobError puts the job into an error state, preventing the job from running again.
|
||||
func (worker *BatchMigrationWorker) setJobError(logger mlog.LoggerIFace, job *model.Job, appError *model.AppError) {
|
||||
if err := worker.jobServer.SetJobError(job, appError); err != nil {
|
||||
logger.Error("Worker: Failed to set job error", mlog.Err(err))
|
||||
}
|
||||
}
|
||||
|
||||
// markAsComplete records a discrete migration key to prevent this job from ever running again.
|
||||
func (worker *BatchMigrationWorker) markAsComplete() {
|
||||
system := model.System{
|
||||
|
||||
@@ -45,51 +45,18 @@ func (ma *MockApp) SetOutOfSync() {
|
||||
}
|
||||
|
||||
func TestBatchMigrationWorker(t *testing.T) {
|
||||
waitDone := func(t *testing.T, done chan bool, msg string) {
|
||||
t.Helper()
|
||||
|
||||
require.Eventually(t, func() bool {
|
||||
select {
|
||||
case <-done:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}, 5*time.Second, 100*time.Millisecond, msg)
|
||||
}
|
||||
|
||||
setupBatchWorker := func(t *testing.T, th *TestHelper, mockApp *MockApp, doMigrationBatch func(model.StringMap, store.Store) (model.StringMap, bool, error)) (model.Worker, *model.Job) {
|
||||
t.Helper()
|
||||
|
||||
migrationKey := model.NewId()
|
||||
timeBetweenBatches := 1 * time.Second
|
||||
|
||||
worker := jobs.MakeBatchMigrationWorker(
|
||||
th.Server.Jobs,
|
||||
th.Server.Store(),
|
||||
mockApp,
|
||||
migrationKey,
|
||||
timeBetweenBatches,
|
||||
model.NewId(),
|
||||
1*time.Second,
|
||||
doMigrationBatch,
|
||||
)
|
||||
th.Server.Jobs.RegisterJobType(migrationKey, worker, nil)
|
||||
|
||||
job, appErr := th.Server.Jobs.CreateJob(th.Context, migrationKey, nil)
|
||||
require.Nil(t, appErr)
|
||||
|
||||
done := make(chan bool)
|
||||
go func() {
|
||||
defer close(done)
|
||||
worker.Run()
|
||||
}()
|
||||
|
||||
// When ending the test, ensure we wait for the worker to finish.
|
||||
t.Cleanup(func() {
|
||||
waitDone(t, done, "worker did not stop running")
|
||||
})
|
||||
|
||||
// Give the worker time to start running
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
job := th.SetupBatchWorker(t, worker.BatchWorker)
|
||||
|
||||
return worker, job
|
||||
}
|
||||
@@ -106,18 +73,6 @@ func TestBatchMigrationWorker(t *testing.T) {
|
||||
waitDone(t, stopped, "worker did not stop")
|
||||
}
|
||||
|
||||
waitForJobStatus := func(t *testing.T, th *TestHelper, job *model.Job, status string) {
|
||||
t.Helper()
|
||||
|
||||
require.Eventuallyf(t, func() bool {
|
||||
actualJob, appErr := th.Server.Jobs.GetJob(th.Context, job.Id)
|
||||
require.Nil(t, appErr)
|
||||
require.Equal(t, job.Id, actualJob.Id)
|
||||
|
||||
return actualJob.Status == status
|
||||
}, 5*time.Second, 250*time.Millisecond, "job never transitioned to %s", status)
|
||||
}
|
||||
|
||||
assertJobReset := func(t *testing.T, th *TestHelper, job *model.Job) {
|
||||
actualJob, appErr := th.Server.Jobs.GetJob(th.Context, job.Id)
|
||||
require.Nil(t, appErr)
|
||||
@@ -144,6 +99,34 @@ func TestBatchMigrationWorker(t *testing.T) {
|
||||
return data
|
||||
}
|
||||
|
||||
t.Run("done after three batches", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
mockApp := &MockApp{}
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
worker, job = setupBatchWorker(t, th, mockApp, func(data model.StringMap, s store.Store) (model.StringMap, bool, error) {
|
||||
batchNumber := getBatchNumberFromData(t, data)
|
||||
require.LessOrEqual(t, batchNumber, 3, "only 3 batches should have run")
|
||||
|
||||
if batchNumber >= 3 {
|
||||
go worker.Stop() // Shut down the worker when the job is done
|
||||
return getDataFromBatchNumber(batchNumber), true, nil
|
||||
}
|
||||
|
||||
batchNumber++
|
||||
return getDataFromBatchNumber(batchNumber), false, nil
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForJobStatus(t, job, model.JobStatusSuccess)
|
||||
th.WaitForBatchNumber(t, job, 3)
|
||||
})
|
||||
|
||||
t.Run("clusters not in sync before first batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
@@ -165,65 +148,12 @@ func TestBatchMigrationWorker(t *testing.T) {
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
waitForJobStatus(t, th, job, model.JobStatusPending)
|
||||
th.WaitForJobStatus(t, job, model.JobStatusPending)
|
||||
assertJobReset(t, th, job)
|
||||
|
||||
stopWorker(t, worker)
|
||||
})
|
||||
|
||||
t.Run("stop after first batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
mockApp := &MockApp{}
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
worker, job = setupBatchWorker(t, th, mockApp, func(data model.StringMap, s store.Store) (model.StringMap, bool, error) {
|
||||
batchNumber := getBatchNumberFromData(t, data)
|
||||
|
||||
require.Equal(t, 1, batchNumber, "only batch 1 should have run")
|
||||
|
||||
// Shut down the worker after the first batch to prevent subsequent ones.
|
||||
go worker.Stop()
|
||||
|
||||
batchNumber++
|
||||
|
||||
return getDataFromBatchNumber(batchNumber), false, nil
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
waitForJobStatus(t, th, job, model.JobStatusPending)
|
||||
})
|
||||
|
||||
t.Run("stop after second batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
mockApp := &MockApp{}
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
worker, job = setupBatchWorker(t, th, mockApp, func(data model.StringMap, s store.Store) (model.StringMap, bool, error) {
|
||||
batchNumber := getBatchNumberFromData(t, data)
|
||||
|
||||
require.LessOrEqual(t, batchNumber, 2, "only batches 1 and 2 should have run")
|
||||
|
||||
// Shut down the worker after the first batch to prevent subsequent ones.
|
||||
go worker.Stop()
|
||||
batchNumber++
|
||||
|
||||
return getDataFromBatchNumber(batchNumber), false, nil
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
waitForJobStatus(t, th, job, model.JobStatusPending)
|
||||
})
|
||||
|
||||
t.Run("clusters not in sync after first batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
@@ -248,59 +178,9 @@ func TestBatchMigrationWorker(t *testing.T) {
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
waitForJobStatus(t, th, job, model.JobStatusPending)
|
||||
th.WaitForJobStatus(t, job, model.JobStatusPending)
|
||||
assertJobReset(t, th, job)
|
||||
|
||||
stopWorker(t, worker)
|
||||
})
|
||||
|
||||
t.Run("done after first batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
mockApp := &MockApp{}
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
worker, job = setupBatchWorker(t, th, mockApp, func(data model.StringMap, s store.Store) (model.StringMap, bool, error) {
|
||||
batchNumber := getBatchNumberFromData(t, data)
|
||||
require.Equal(t, 1, batchNumber, "only batch 1 should have run")
|
||||
|
||||
// Shut down the worker after the first batch to prevent subsequent ones.
|
||||
go worker.Stop()
|
||||
batchNumber++
|
||||
|
||||
return getDataFromBatchNumber(batchNumber), true, nil
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
waitForJobStatus(t, th, job, model.JobStatusSuccess)
|
||||
})
|
||||
|
||||
t.Run("done after three batches", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
mockApp := &MockApp{}
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
worker, job = setupBatchWorker(t, th, mockApp, func(data model.StringMap, s store.Store) (model.StringMap, bool, error) {
|
||||
batchNumber := getBatchNumberFromData(t, data)
|
||||
require.LessOrEqual(t, batchNumber, 3, "only 3 batches should have run")
|
||||
|
||||
// Shut down the worker after the first batch to prevent subsequent ones.
|
||||
go worker.Stop()
|
||||
batchNumber++
|
||||
|
||||
return getDataFromBatchNumber(batchNumber), true, nil
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
waitForJobStatus(t, th, job, model.JobStatusSuccess)
|
||||
})
|
||||
}
|
||||
|
||||
139
server/channels/jobs/batch_report_worker.go
Обычный файл
139
server/channels/jobs/batch_report_worker.go
Обычный файл
@@ -0,0 +1,139 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package jobs
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost/server/public/model"
|
||||
"github.com/mattermost/mattermost/server/public/shared/mlog"
|
||||
"github.com/mattermost/mattermost/server/public/shared/request"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/store"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
type BatchReportWorkerAppIFace interface {
|
||||
SaveReportChunk(format string, prefix string, count int, reportData []model.ReportableObject) *model.AppError
|
||||
CompileReportChunks(format string, prefix string, numberOfChunks int, headers []string) *model.AppError
|
||||
SendReportToUser(rctx request.CTX, userID string, jobId string, format string) *model.AppError
|
||||
CleanupReportChunks(format string, prefix string, numberOfChunks int) *model.AppError
|
||||
}
|
||||
|
||||
type BatchReportWorker struct {
|
||||
*BatchWorker
|
||||
app BatchReportWorkerAppIFace
|
||||
reportFormat string
|
||||
headers []string
|
||||
getData func(jobData model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error)
|
||||
}
|
||||
|
||||
func MakeBatchReportWorker(
|
||||
jobServer *JobServer,
|
||||
store store.Store,
|
||||
app BatchReportWorkerAppIFace,
|
||||
timeBetweenBatches time.Duration,
|
||||
reportFormat string,
|
||||
headers []string,
|
||||
getData func(jobData model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error),
|
||||
) *BatchReportWorker {
|
||||
worker := &BatchReportWorker{
|
||||
app: app,
|
||||
reportFormat: reportFormat,
|
||||
headers: headers,
|
||||
getData: getData,
|
||||
}
|
||||
worker.BatchWorker = MakeBatchWorker(jobServer, store, timeBetweenBatches, worker.doBatch)
|
||||
return worker
|
||||
}
|
||||
|
||||
func (worker *BatchReportWorker) doBatch(rctx *request.Context, job *model.Job) bool {
|
||||
reportData, nextData, done, err := worker.getData(job.Data)
|
||||
if err != nil {
|
||||
worker.logger.Error("Worker: Failed to get data for report batch. Exiting", mlog.Err(err))
|
||||
worker.setJobError(worker.logger, job, model.NewAppError("doBatch", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
|
||||
return true
|
||||
} else if done {
|
||||
if err = worker.complete(rctx, job); err != nil {
|
||||
worker.logger.Error("Worker: Failed to finish the batch report. Exiting", mlog.Err(err))
|
||||
worker.setJobError(worker.logger, job, model.NewAppError("doBatch", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
|
||||
} else {
|
||||
worker.logger.Info("Worker: Report job complete")
|
||||
worker.setJobSuccess(worker.logger, job)
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
err = worker.processChunk(job, reportData)
|
||||
if err != nil {
|
||||
worker.logger.Error("Worker: Failed to save report batch. Exiting", mlog.Err(err))
|
||||
worker.setJobError(worker.logger, job, model.NewAppError("doBatch", model.NoTranslation, nil, "", http.StatusInternalServerError).Wrap(err))
|
||||
return true
|
||||
}
|
||||
|
||||
job.Data = nextData
|
||||
|
||||
// We might be able to add progress for this type of job in the future
|
||||
// But for now we can just set to 0
|
||||
worker.jobServer.SetJobProgress(job, 0)
|
||||
return false
|
||||
}
|
||||
|
||||
func getFileCount(jobData model.StringMap) (int, error) {
|
||||
if jobData["file_count"] != "" {
|
||||
parsedFileCount, parseErr := strconv.Atoi(jobData["file_count"])
|
||||
if parseErr != nil {
|
||||
return 0, errors.Wrap(parseErr, "failed to parse file_count")
|
||||
}
|
||||
return parsedFileCount, nil
|
||||
}
|
||||
|
||||
// Assume it hasn't been set
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func (worker *BatchReportWorker) processChunk(job *model.Job, reportData []model.ReportableObject) error {
|
||||
fileCount, err := getFileCount(job.Data)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
appErr := worker.app.SaveReportChunk(worker.reportFormat, job.Id, fileCount, reportData)
|
||||
if appErr != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
fileCount++
|
||||
job.Data["file_count"] = strconv.Itoa(fileCount)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (worker *BatchReportWorker) complete(rctx request.CTX, job *model.Job) error {
|
||||
requestingUserId := job.Data["requesting_user_id"]
|
||||
if requestingUserId == "" {
|
||||
return errors.New("No user to send the report to")
|
||||
}
|
||||
fileCount, err := getFileCount(job.Data)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
appErr := worker.app.CompileReportChunks(worker.reportFormat, job.Id, fileCount, worker.headers)
|
||||
if appErr != nil {
|
||||
return appErr
|
||||
}
|
||||
|
||||
defer func() {
|
||||
worker.app.CleanupReportChunks(worker.reportFormat, job.Id, fileCount)
|
||||
}()
|
||||
|
||||
if appErr = worker.app.SendReportToUser(rctx, requestingUserId, job.Id, worker.reportFormat); appErr != nil {
|
||||
return appErr
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
149
server/channels/jobs/batch_report_worker_test.go
Обычный файл
149
server/channels/jobs/batch_report_worker_test.go
Обычный файл
@@ -0,0 +1,149 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package jobs_test
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strconv"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost/server/public/model"
|
||||
"github.com/mattermost/mattermost/server/public/shared/request"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/jobs"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
type ReportMockApp struct{}
|
||||
|
||||
func (rma *ReportMockApp) SaveReportChunk(format string, prefix string, count int, reportData []model.ReportableObject) *model.AppError {
|
||||
return nil
|
||||
}
|
||||
func (rma *ReportMockApp) CompileReportChunks(format string, prefix string, numberOfChunks int, headers []string) *model.AppError {
|
||||
return nil
|
||||
}
|
||||
func (rma *ReportMockApp) SendReportToUser(rctx request.CTX, userID string, jobId string, format string) *model.AppError {
|
||||
return nil
|
||||
}
|
||||
func (rma *ReportMockApp) CleanupReportChunks(format string, prefix string, numberOfChunks int) *model.AppError {
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestBatchReportWorker(t *testing.T) {
|
||||
setupBatchWorker := func(
|
||||
t *testing.T,
|
||||
th *TestHelper,
|
||||
getData func(jobData model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error),
|
||||
) (*jobs.BatchReportWorker, *model.Job) {
|
||||
t.Helper()
|
||||
|
||||
worker := jobs.MakeBatchReportWorker(
|
||||
th.Server.Jobs,
|
||||
th.Server.Store(),
|
||||
&ReportMockApp{},
|
||||
1*time.Second,
|
||||
"csv",
|
||||
[]string{},
|
||||
getData)
|
||||
job := th.SetupBatchWorker(t, worker.BatchWorker)
|
||||
return worker, job
|
||||
}
|
||||
|
||||
waitForFileCount := func(t *testing.T, th *TestHelper, job *model.Job, fileCount int) {
|
||||
t.Helper()
|
||||
|
||||
require.Eventuallyf(t, func() bool {
|
||||
actualJob, appErr := th.Server.Jobs.GetJob(th.Context, job.Id)
|
||||
require.Nil(t, appErr)
|
||||
require.Equal(t, job.Id, actualJob.Id)
|
||||
|
||||
finalFileCount, err := strconv.Atoi(actualJob.Data["file_count"])
|
||||
require.NoError(t, err)
|
||||
return finalFileCount == fileCount
|
||||
}, 5*time.Second, 250*time.Millisecond, "job did not stop at batch %d", fileCount)
|
||||
}
|
||||
|
||||
getFileCountFromData := func(t *testing.T, data model.StringMap) int {
|
||||
t.Helper()
|
||||
|
||||
if data["file_count"] == "" {
|
||||
return 0
|
||||
}
|
||||
|
||||
fileCount, err := strconv.Atoi(data["file_count"])
|
||||
require.NoError(t, err)
|
||||
|
||||
return fileCount
|
||||
}
|
||||
|
||||
createData := func(th *TestHelper, data model.StringMap) model.StringMap {
|
||||
data["requesting_user_id"] = th.SystemAdminUser.Id
|
||||
return data
|
||||
}
|
||||
|
||||
t.Run("should finish when the report is done, incrementing file count along the way", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
|
||||
iterations := 0
|
||||
|
||||
worker, job = setupBatchWorker(t, th, func(data model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error) {
|
||||
fileCount := getFileCountFromData(t, data)
|
||||
require.Equal(t, iterations, fileCount)
|
||||
require.LessOrEqual(t, fileCount, 3, "only 3 batches should have run")
|
||||
|
||||
iterations++
|
||||
|
||||
if fileCount >= 3 {
|
||||
go worker.Stop() // Shut down the worker when the job is done
|
||||
return []model.ReportableObject{}, createData(th, data), true, nil
|
||||
}
|
||||
|
||||
return []model.ReportableObject{}, createData(th, data), false, nil
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForJobStatus(t, job, model.JobStatusSuccess)
|
||||
waitForFileCount(t, th, job, 3)
|
||||
})
|
||||
|
||||
t.Run("should fail job when get data throws an error", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
worker, job = setupBatchWorker(t, th, func(data model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error) {
|
||||
go worker.Stop() // Shut down the worker right after this
|
||||
return []model.ReportableObject{}, createData(th, data), false, errors.New("failed to fetch data")
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForJobStatus(t, job, model.JobStatusError)
|
||||
})
|
||||
|
||||
t.Run("should fail if there is no user id to send the report to", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
var worker model.Worker
|
||||
var job *model.Job
|
||||
worker, job = setupBatchWorker(t, th, func(data model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error) {
|
||||
go worker.Stop() // Shut down the worker right after this
|
||||
return []model.ReportableObject{}, make(model.StringMap), true, nil
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForJobStatus(t, job, model.JobStatusError)
|
||||
})
|
||||
}
|
||||
174
server/channels/jobs/batch_worker.go
Обычный файл
174
server/channels/jobs/batch_worker.go
Обычный файл
@@ -0,0 +1,174 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package jobs
|
||||
|
||||
import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost/server/public/model"
|
||||
"github.com/mattermost/mattermost/server/public/shared/mlog"
|
||||
"github.com/mattermost/mattermost/server/public/shared/request"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/store"
|
||||
)
|
||||
|
||||
type BatchWorker struct {
|
||||
jobServer *JobServer
|
||||
logger mlog.LoggerIFace
|
||||
store store.Store
|
||||
|
||||
stop chan struct{}
|
||||
stopped chan bool
|
||||
closed atomic.Bool
|
||||
jobs chan model.Job
|
||||
|
||||
timeBetweenBatches time.Duration
|
||||
doBatch func(rctx *request.Context, job *model.Job) bool
|
||||
}
|
||||
|
||||
// MakeBatchWorker creates a worker to process the given batch function.
|
||||
func MakeBatchWorker(
|
||||
jobServer *JobServer,
|
||||
store store.Store,
|
||||
timeBetweenBatches time.Duration,
|
||||
doBatch func(rctx *request.Context, job *model.Job) bool,
|
||||
) *BatchWorker {
|
||||
return &BatchWorker{
|
||||
jobServer: jobServer,
|
||||
logger: jobServer.Logger(),
|
||||
store: store,
|
||||
stop: make(chan struct{}),
|
||||
stopped: make(chan bool, 1),
|
||||
jobs: make(chan model.Job),
|
||||
timeBetweenBatches: timeBetweenBatches,
|
||||
doBatch: doBatch,
|
||||
}
|
||||
}
|
||||
|
||||
// Run starts the worker dedicated to the unique migration batch job it will be given to process.
|
||||
func (worker *BatchWorker) Run() {
|
||||
worker.logger.Debug("Worker started")
|
||||
// We have to re-assign the stop channel again, because
|
||||
// it might happen that the job was restarted due to a config change.
|
||||
if worker.closed.CompareAndSwap(true, false) {
|
||||
worker.stop = make(chan struct{})
|
||||
}
|
||||
|
||||
defer func() {
|
||||
worker.logger.Debug("Worker finished")
|
||||
worker.stopped <- true
|
||||
}()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-worker.stop:
|
||||
worker.logger.Debug("Worker received stop signal")
|
||||
return
|
||||
case job := <-worker.jobs:
|
||||
worker.DoJob(&job)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Stop interrupts the worker even if the migration has not yet completed.
|
||||
func (worker *BatchWorker) Stop() {
|
||||
// Set to close, and if already closed before, then return.
|
||||
if !worker.closed.CompareAndSwap(false, true) {
|
||||
return
|
||||
}
|
||||
|
||||
worker.logger.Debug("Worker stopping")
|
||||
close(worker.stop)
|
||||
<-worker.stopped
|
||||
}
|
||||
|
||||
// JobChannel is the means by which the jobs infrastructure provides the worker the job to execute.
|
||||
func (worker *BatchWorker) JobChannel() chan<- model.Job {
|
||||
return worker.jobs
|
||||
}
|
||||
|
||||
// IsEnabled is always true for batches.
|
||||
func (worker *BatchWorker) IsEnabled(_ *model.Config) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
// DoJob executes the job picked up through the job channel.
|
||||
//
|
||||
// Note that this is a lot of distracting machinery here to claim the job, then double check the
|
||||
// status, and keep the status up to date in line with job infrastrcuture semantics. Unless an
|
||||
// error occurs, this worker should hold onto the job until its completed.
|
||||
func (worker *BatchWorker) DoJob(job *model.Job) {
|
||||
logger := worker.logger.With(mlog.Any("job", job))
|
||||
logger.Debug("Worker received a new candidate job.")
|
||||
defer worker.jobServer.HandleJobPanic(logger, job)
|
||||
|
||||
if claimed, err := worker.jobServer.ClaimJob(job); err != nil {
|
||||
logger.Warn("Worker experienced an error while trying to claim job", mlog.Err(err))
|
||||
return
|
||||
} else if !claimed {
|
||||
return
|
||||
}
|
||||
|
||||
c := request.EmptyContext(logger)
|
||||
var appErr *model.AppError
|
||||
|
||||
// We get the job again because ClaimJob changes the job status.
|
||||
job, appErr = worker.jobServer.GetJob(c, job.Id)
|
||||
if appErr != nil {
|
||||
worker.logger.Error("Worker: job execution error", mlog.Err(appErr))
|
||||
worker.setJobError(logger, job, appErr)
|
||||
return
|
||||
}
|
||||
|
||||
if job.Data == nil {
|
||||
job.Data = make(model.StringMap)
|
||||
}
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-worker.stop:
|
||||
logger.Info("Worker: Batch has been canceled via Worker Stop. Setting the job back to pending.")
|
||||
if err := worker.jobServer.SetJobPending(job); err != nil {
|
||||
worker.logger.Error("Worker: Failed to mark job as pending", mlog.Err(err))
|
||||
}
|
||||
return
|
||||
case <-time.After(worker.timeBetweenBatches):
|
||||
if stop := worker.doBatch(c, job); stop {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// resetJob erases the data tracking the next batch to execute and returns the job status to
|
||||
// pending to allow the job infrastructure to requeue it.
|
||||
func (worker *BatchWorker) resetJob(logger mlog.LoggerIFace, job *model.Job) {
|
||||
job.Data = nil
|
||||
job.Progress = 0
|
||||
job.Status = model.JobStatusPending
|
||||
|
||||
if _, err := worker.store.Job().UpdateOptimistically(job, model.JobStatusInProgress); err != nil {
|
||||
worker.logger.Error("Worker: Failed to reset job data. May resume instead of restarting.", mlog.Err(err))
|
||||
}
|
||||
}
|
||||
|
||||
// setJobSuccess records the job as successful.
|
||||
func (worker *BatchWorker) setJobSuccess(logger mlog.LoggerIFace, job *model.Job) {
|
||||
if err := worker.jobServer.SetJobProgress(job, 100); err != nil {
|
||||
logger.Error("Worker: Failed to update progress for job", mlog.Err(err))
|
||||
worker.setJobError(logger, job, err)
|
||||
}
|
||||
|
||||
if err := worker.jobServer.SetJobSuccess(job); err != nil {
|
||||
logger.Error("Worker: Failed to set success for job", mlog.Err(err))
|
||||
worker.setJobError(logger, job, err)
|
||||
}
|
||||
}
|
||||
|
||||
// setJobError puts the job into an error state, preventing the job from running again.
|
||||
func (worker *BatchWorker) setJobError(logger mlog.LoggerIFace, job *model.Job, appError *model.AppError) {
|
||||
if err := worker.jobServer.SetJobError(job, appError); err != nil {
|
||||
logger.Error("Worker: Failed to set job error", mlog.Err(err))
|
||||
}
|
||||
}
|
||||
147
server/channels/jobs/batch_worker_test.go
Обычный файл
147
server/channels/jobs/batch_worker_test.go
Обычный файл
@@ -0,0 +1,147 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package jobs_test
|
||||
|
||||
import (
|
||||
"strconv"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost/server/public/model"
|
||||
"github.com/mattermost/mattermost/server/public/shared/request"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/jobs"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestBatchWorker(t *testing.T) {
|
||||
createBatchWorker := func(t *testing.T, th *TestHelper, doBatch func(rctx *request.Context, job *model.Job) bool) (*jobs.BatchWorker, *model.Job) {
|
||||
t.Helper()
|
||||
|
||||
worker := jobs.MakeBatchWorker(th.Server.Jobs, th.Server.Store(), 1*time.Second, doBatch)
|
||||
job := th.SetupBatchWorker(t, worker)
|
||||
return worker, job
|
||||
}
|
||||
|
||||
getBatchNumberFromData := func(t *testing.T, data model.StringMap) int {
|
||||
t.Helper()
|
||||
|
||||
batchNumber, err := strconv.Atoi(data["batch_number"])
|
||||
require.NoError(t, err)
|
||||
|
||||
return batchNumber
|
||||
}
|
||||
|
||||
incrementBatchNumber := func(t *testing.T, th *TestHelper, job *model.Job) {
|
||||
t.Helper()
|
||||
|
||||
batchNumber, err := strconv.Atoi(job.Data["batch_number"])
|
||||
require.NoError(t, err)
|
||||
|
||||
batchNumber++
|
||||
job.Data["batch_number"] = strconv.Itoa(batchNumber)
|
||||
th.Server.Jobs.SetJobProgress(job, 0)
|
||||
}
|
||||
|
||||
t.Run("stop after first batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
var worker *jobs.BatchWorker
|
||||
worker, job := createBatchWorker(t, th, func(rctx *request.Context, job *model.Job) bool {
|
||||
batchNumber := getBatchNumberFromData(t, job.Data)
|
||||
|
||||
require.Equal(t, 1, batchNumber, "only batch 1 should have run")
|
||||
|
||||
// Shut down the worker after the first batch to prevent subsequent ones.
|
||||
if batchNumber >= 1 {
|
||||
go worker.Stop()
|
||||
} else {
|
||||
incrementBatchNumber(t, th, job)
|
||||
}
|
||||
|
||||
return false
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForJobStatus(t, job, model.JobStatusPending)
|
||||
th.WaitForBatchNumber(t, job, 1)
|
||||
})
|
||||
|
||||
t.Run("stop after second batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
var worker *jobs.BatchWorker
|
||||
worker, job := createBatchWorker(t, th, func(rctx *request.Context, job *model.Job) bool {
|
||||
batchNumber := getBatchNumberFromData(t, job.Data)
|
||||
|
||||
require.LessOrEqual(t, batchNumber, 2, "only batches 1 and 2 should have run")
|
||||
|
||||
// Shut down the worker after the second batch to prevent subsequent ones.
|
||||
if batchNumber >= 2 {
|
||||
go worker.Stop()
|
||||
} else {
|
||||
incrementBatchNumber(t, th, job)
|
||||
}
|
||||
|
||||
return false
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForJobStatus(t, job, model.JobStatusPending)
|
||||
th.WaitForBatchNumber(t, job, 2)
|
||||
})
|
||||
|
||||
t.Run("done after first batch", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
var worker *jobs.BatchWorker
|
||||
worker, job := createBatchWorker(t, th, func(rctx *request.Context, job *model.Job) bool {
|
||||
batchNumber := getBatchNumberFromData(t, job.Data)
|
||||
require.Equal(t, 1, batchNumber, "only batch 1 should have run")
|
||||
|
||||
if batchNumber >= 1 {
|
||||
go worker.Stop() // Shut down the worker when the job is done
|
||||
return true
|
||||
}
|
||||
|
||||
incrementBatchNumber(t, th, job)
|
||||
return false
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForBatchNumber(t, job, 1)
|
||||
})
|
||||
|
||||
t.Run("done after three batches", func(t *testing.T) {
|
||||
th := Setup(t).InitBasic()
|
||||
defer th.TearDown()
|
||||
|
||||
var worker *jobs.BatchWorker
|
||||
worker, job := createBatchWorker(t, th, func(rctx *request.Context, job *model.Job) bool {
|
||||
batchNumber := getBatchNumberFromData(t, job.Data)
|
||||
require.LessOrEqual(t, batchNumber, 3, "only 3 batches should have run")
|
||||
|
||||
if batchNumber >= 3 {
|
||||
go worker.Stop() // Shut down the worker when the job is done
|
||||
return true
|
||||
}
|
||||
|
||||
incrementBatchNumber(t, th, job)
|
||||
return false
|
||||
})
|
||||
|
||||
// Queue the work to be done
|
||||
worker.JobChannel() <- *job
|
||||
|
||||
th.WaitForBatchNumber(t, job, 3)
|
||||
})
|
||||
}
|
||||
107
server/channels/jobs/export_users_to_csv/export_users_to_csv.go
Обычный файл
107
server/channels/jobs/export_users_to_csv/export_users_to_csv.go
Обычный файл
@@ -0,0 +1,107 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package export_users_to_csv
|
||||
|
||||
import (
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost/server/public/model"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/jobs"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/store"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
const (
|
||||
timeBetweenBatches = 1 * time.Second
|
||||
)
|
||||
|
||||
type ExportUsersToCSVAppIFace interface {
|
||||
jobs.BatchReportWorkerAppIFace
|
||||
GetUsersForReporting(filter *model.UserReportOptions) ([]*model.UserReport, *model.AppError)
|
||||
}
|
||||
|
||||
// MakeWorker creates a batch report worker to generate CSV user reports.
|
||||
func MakeWorker(jobServer *jobs.JobServer, store store.Store, app ExportUsersToCSVAppIFace) model.Worker {
|
||||
return jobs.MakeBatchReportWorker(
|
||||
jobServer,
|
||||
store,
|
||||
app,
|
||||
timeBetweenBatches,
|
||||
"csv",
|
||||
[]string{
|
||||
"Id",
|
||||
"Username",
|
||||
"Email",
|
||||
"CreateAt",
|
||||
"Name",
|
||||
"Roles",
|
||||
"LastLogin",
|
||||
"LastStatusAt",
|
||||
"LastPostDate",
|
||||
"DaysActive",
|
||||
"TotalPosts",
|
||||
},
|
||||
getData(app),
|
||||
)
|
||||
}
|
||||
|
||||
// parseJobMetadata parses the opaque job metadata to return the information needed to decide which
|
||||
// batch to process next.
|
||||
func parseJobMetadata(data model.StringMap) (*model.UserReportOptions, error) {
|
||||
startAt, err := strconv.ParseInt(data["start_at"], 10, 64)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
endAt, err := strconv.ParseInt(data["end_at"], 10, 64)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
options := model.UserReportOptions{
|
||||
ReportingBaseOptions: model.ReportingBaseOptions{
|
||||
SortColumn: "Username",
|
||||
PageSize: 100,
|
||||
FromColumnValue: data["last_column_value"],
|
||||
FromId: data["last_user_id"],
|
||||
StartAt: startAt,
|
||||
EndAt: endAt,
|
||||
},
|
||||
}
|
||||
|
||||
return &options, nil
|
||||
}
|
||||
|
||||
// makeJobMetadata encodes the information needed to decide which batch to process next back into
|
||||
// the opaque job metadata.
|
||||
func makeJobMetadata(jobData model.StringMap, lastColumnValue string, userID string) model.StringMap {
|
||||
jobData["last_column_value"] = lastColumnValue
|
||||
jobData["last_user_id"] = userID
|
||||
return jobData
|
||||
}
|
||||
|
||||
func getData(app ExportUsersToCSVAppIFace) func(jobData model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error) {
|
||||
return func(jobData model.StringMap) ([]model.ReportableObject, model.StringMap, bool, error) {
|
||||
filter, err := parseJobMetadata(jobData)
|
||||
if err != nil {
|
||||
return nil, nil, false, errors.Wrap(err, "failed to parse job metadata")
|
||||
}
|
||||
|
||||
users, appErr := app.GetUsersForReporting(filter)
|
||||
if appErr != nil {
|
||||
return nil, nil, false, errors.Wrapf(err, "failed to get the next batch (column_value=%v, user_id=%v)", filter.FromColumnValue, filter.FromId)
|
||||
}
|
||||
|
||||
if len(users) == 0 {
|
||||
return nil, nil, true, nil
|
||||
}
|
||||
|
||||
reportableObjects := []model.ReportableObject{}
|
||||
for i := 0; i < len(users); i++ {
|
||||
reportableObjects = append(reportableObjects, users[i])
|
||||
}
|
||||
|
||||
return reportableObjects, makeJobMetadata(jobData, users[len(users)-1].Username, users[len(users)-1].Id), false, nil
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ package jobs_test
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -14,8 +15,10 @@ import (
|
||||
"github.com/mattermost/mattermost/server/public/shared/mlog"
|
||||
"github.com/mattermost/mattermost/server/public/shared/request"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/app"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/jobs"
|
||||
"github.com/mattermost/mattermost/server/v8/channels/store"
|
||||
"github.com/mattermost/mattermost/server/v8/config"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
type TestHelper struct {
|
||||
@@ -221,3 +224,73 @@ func (th *TestHelper) TearDown() {
|
||||
os.RemoveAll(th.tempWorkspace)
|
||||
}
|
||||
}
|
||||
|
||||
func (th *TestHelper) SetupBatchWorker(t *testing.T, worker *jobs.BatchWorker) *model.Job {
|
||||
t.Helper()
|
||||
|
||||
jobId := model.NewId()
|
||||
th.Server.Jobs.RegisterJobType(jobId, worker, nil)
|
||||
|
||||
jobData := make(model.StringMap)
|
||||
jobData["batch_number"] = "1"
|
||||
job, appErr := th.Server.Jobs.CreateJob(th.Context, jobId, jobData)
|
||||
|
||||
if appErr != nil {
|
||||
panic(appErr)
|
||||
}
|
||||
|
||||
done := make(chan bool)
|
||||
go func() {
|
||||
defer close(done)
|
||||
worker.Run()
|
||||
}()
|
||||
|
||||
// When ending the test, ensure we wait for the worker to finish.
|
||||
t.Cleanup(func() {
|
||||
waitDone(t, done, "worker did not stop running")
|
||||
})
|
||||
|
||||
// Give the worker time to start running
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
|
||||
return job
|
||||
}
|
||||
|
||||
func (th *TestHelper) WaitForJobStatus(t *testing.T, job *model.Job, status string) {
|
||||
t.Helper()
|
||||
|
||||
require.Eventuallyf(t, func() bool {
|
||||
actualJob, appErr := th.Server.Jobs.GetJob(th.Context, job.Id)
|
||||
require.Nil(t, appErr)
|
||||
require.Equal(t, job.Id, actualJob.Id)
|
||||
|
||||
return actualJob.Status == status
|
||||
}, 5*time.Second, 250*time.Millisecond, "job never transitioned to %s", status)
|
||||
}
|
||||
|
||||
func (th *TestHelper) WaitForBatchNumber(t *testing.T, job *model.Job, batchNumber int) {
|
||||
t.Helper()
|
||||
|
||||
require.Eventuallyf(t, func() bool {
|
||||
actualJob, appErr := th.Server.Jobs.GetJob(th.Context, job.Id)
|
||||
require.Nil(t, appErr)
|
||||
require.Equal(t, job.Id, actualJob.Id)
|
||||
|
||||
finalBatchNumber, err := strconv.Atoi(actualJob.Data["batch_number"])
|
||||
require.NoError(t, err)
|
||||
return finalBatchNumber == batchNumber
|
||||
}, 5*time.Second, 250*time.Millisecond, "job did not stop at batch %d", batchNumber)
|
||||
}
|
||||
|
||||
func waitDone(t *testing.T, done chan bool, msg string) {
|
||||
t.Helper()
|
||||
|
||||
require.Eventually(t, func() bool {
|
||||
select {
|
||||
case <-done:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}, 5*time.Second, 100*time.Millisecond, msg)
|
||||
}
|
||||
|
||||
Ссылка в новой задаче
Block a user