* Refactor job retrieval to support multiple statuses & multiple types - Updated job retrieval functions to handle multiple job statuses. - Renamed `GetJobsByTypeAndStatus` to `GetJobsByTypesAndStatuses` for consistency across the codebase. - Adjusted related function signatures and implementations in the job store and retry layer to accommodate the new method. - Updated tests to reflect changes in job retrieval logic and ensure proper functionality. * Add compliance export create command and tests - Introduced `ComplianceExportCreateCmd` to facilitate the creation of compliance export jobs with options for date, start, and end timestamps. - Added unit tests for the new command, covering various scenarios including valid and invalid inputs. - Updated documentation to include usage examples and options for the new command. - Enhanced existing tests to ensure proper functionality of compliance export job handling. * update docs * update tests for new logic * Refactor message export job tests to use DefaultPreviousJobPageSize - Updated all test cases in worker_test.go to replace hardcoded page size of 100 with DefaultPreviousJobPageSize for consistency. - Adjusted the worker.go file to define DefaultPreviousJobPageSize and use it in job retrieval logic. - Ensured that the changes maintain the functionality of job data initialization and retrieval tests. * PR comments * PR comments, simplifications, clarifications, formatting * prefer hypen over underscore in command names * merge conflict * update mmctl docs
751 строка
26 KiB
Go
751 строка
26 KiB
Go
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
|
// See LICENSE.enterprise for license information.
|
|
|
|
package message_export
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"path"
|
|
"strconv"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/stretchr/testify/assert"
|
|
tmock "github.com/stretchr/testify/mock"
|
|
"github.com/stretchr/testify/require"
|
|
|
|
"github.com/mattermost/mattermost/server/public/model"
|
|
"github.com/mattermost/mattermost/server/public/plugin/plugintest/mock"
|
|
"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/storetest"
|
|
st "github.com/mattermost/mattermost/server/v8/channels/store/storetest"
|
|
"github.com/mattermost/mattermost/server/v8/channels/utils/testutils"
|
|
"github.com/mattermost/mattermost/server/v8/einterfaces/mocks"
|
|
"github.com/mattermost/mattermost/server/v8/enterprise/message_export/shared"
|
|
)
|
|
|
|
func TestInitJobDataNoJobData(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
job := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusPending,
|
|
Type: model.JobTypeMessageExport,
|
|
}
|
|
|
|
// mock job store doesn't return a previously successful job, forcing fallback to config
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return(nil, errors.New("test"))
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
ConfigService: &testutils.StaticConfigService{
|
|
Cfg: &model.Config{
|
|
// mock config
|
|
MessageExportSettings: model.MessageExportSettings{
|
|
EnableExport: model.NewPointer(true),
|
|
ExportFormat: model.NewPointer(model.ComplianceExportTypeActiance),
|
|
DailyRunTime: model.NewPointer("01:00"),
|
|
ExportFromTimestamp: model.NewPointer(int64(0)),
|
|
BatchSize: model.NewPointer(10000),
|
|
ChannelBatchSize: model.NewPointer(100),
|
|
ChannelHistoryBatchSize: model.NewPointer(100),
|
|
},
|
|
},
|
|
},
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
now := time.Now()
|
|
worker.initJobData(request.EmptyContext(logger), logger, job, now)
|
|
|
|
assert.Equal(t, model.ComplianceExportTypeActiance, job.Data[shared.JobDataExportType])
|
|
assert.Equal(t, strconv.Itoa(*worker.jobServer.Config().MessageExportSettings.BatchSize), job.Data[shared.JobDataBatchSize])
|
|
assert.Equal(t, strconv.FormatInt(*worker.jobServer.Config().MessageExportSettings.ExportFromTimestamp, 10), job.Data[shared.JobDataBatchStartTime])
|
|
expectedDir := path.Join(model.ComplianceExportPath, fmt.Sprintf("%s-%d-%d", now.Format(model.ComplianceExportDirectoryFormat), 0, now.UnixMilli()))
|
|
assert.Equal(t, expectedDir, job.Data[shared.JobDataExportDir])
|
|
}
|
|
|
|
func TestInitJobDataPreviousJobNoJobData(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
previousJob := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
StartAt: model.GetMillis() - 1000,
|
|
LastActivityAt: model.GetMillis() - 1000,
|
|
}
|
|
|
|
job := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusPending,
|
|
Type: model.JobTypeMessageExport,
|
|
}
|
|
|
|
// mock job store returns a previously successful job, but it doesn't have job data either, so we still fall back to config
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return([]*model.Job{previousJob}, nil)
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
ConfigService: &testutils.StaticConfigService{
|
|
Cfg: &model.Config{
|
|
// mock config
|
|
MessageExportSettings: model.MessageExportSettings{
|
|
EnableExport: model.NewPointer(true),
|
|
ExportFormat: model.NewPointer(model.ComplianceExportTypeActiance),
|
|
DailyRunTime: model.NewPointer("01:00"),
|
|
ExportFromTimestamp: model.NewPointer(int64(0)),
|
|
BatchSize: model.NewPointer(10000),
|
|
ChannelBatchSize: model.NewPointer(100),
|
|
ChannelHistoryBatchSize: model.NewPointer(100),
|
|
},
|
|
},
|
|
},
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
now := time.Now()
|
|
worker.initJobData(request.EmptyContext(logger), logger, job, now)
|
|
|
|
assert.Equal(t, model.ComplianceExportTypeActiance, job.Data[shared.JobDataExportType])
|
|
assert.Equal(t, strconv.Itoa(*worker.jobServer.Config().MessageExportSettings.BatchSize), job.Data[shared.JobDataBatchSize])
|
|
assert.Equal(t, strconv.FormatInt(*worker.jobServer.Config().MessageExportSettings.ExportFromTimestamp, 10), job.Data[shared.JobDataBatchStartTime])
|
|
expectedDir := path.Join(model.ComplianceExportPath, fmt.Sprintf("%s-%d-%d", now.Format(model.ComplianceExportDirectoryFormat), 0, now.UnixMilli()))
|
|
assert.Equal(t, expectedDir, job.Data[shared.JobDataExportDir])
|
|
}
|
|
|
|
func TestInitJobDataPreviousJobWithJobData(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
previousJob := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
StartAt: model.GetMillis() - 1000,
|
|
LastActivityAt: model.GetMillis() - 1000,
|
|
Data: map[string]string{shared.JobDataBatchStartTime: "123"},
|
|
}
|
|
|
|
job := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusPending,
|
|
Type: model.JobTypeMessageExport,
|
|
Data: map[string]string{shared.JobDataExportDir: "this-is-the-export-dir"},
|
|
}
|
|
|
|
// mock job store returns a previously successful job that has the config that we're looking for, so we use it
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return([]*model.Job{previousJob}, nil)
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
ConfigService: &testutils.StaticConfigService{
|
|
Cfg: &model.Config{
|
|
// mock config
|
|
MessageExportSettings: model.MessageExportSettings{
|
|
EnableExport: model.NewPointer(true),
|
|
ExportFormat: model.NewPointer(model.ComplianceExportTypeActiance),
|
|
DailyRunTime: model.NewPointer("01:00"),
|
|
ExportFromTimestamp: model.NewPointer(int64(0)),
|
|
BatchSize: model.NewPointer(10000),
|
|
ChannelBatchSize: model.NewPointer(100),
|
|
ChannelHistoryBatchSize: model.NewPointer(100),
|
|
},
|
|
},
|
|
},
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
now := time.Now()
|
|
worker.initJobData(request.EmptyContext(logger), logger, job, now)
|
|
|
|
assert.Equal(t, model.ComplianceExportTypeActiance, job.Data[shared.JobDataExportType])
|
|
assert.Equal(t, strconv.Itoa(*worker.jobServer.Config().MessageExportSettings.BatchSize), job.Data[shared.JobDataBatchSize])
|
|
assert.Equal(t, previousJob.Data[shared.JobDataBatchStartTime], job.Data[shared.JobDataBatchStartTime])
|
|
expectedDir := "this-is-the-export-dir"
|
|
assert.Equal(t, expectedDir, job.Data[shared.JobDataExportDir])
|
|
}
|
|
|
|
func TestInitJobDataPreviousJobWithJobDataPre105(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
previousJob := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
StartAt: model.GetMillis() - 1000,
|
|
LastActivityAt: model.GetMillis() - 1000,
|
|
Data: map[string]string{"batch_start_timestamp": "123"},
|
|
}
|
|
|
|
job := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusPending,
|
|
Type: model.JobTypeMessageExport,
|
|
Data: map[string]string{shared.JobDataExportDir: "this-is-the-export-dir"},
|
|
}
|
|
|
|
// mock job store returns a previously successful job that has the config that we're looking for, so we use it
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return([]*model.Job{previousJob}, nil)
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
ConfigService: &testutils.StaticConfigService{
|
|
Cfg: &model.Config{
|
|
// mock config
|
|
MessageExportSettings: model.MessageExportSettings{
|
|
EnableExport: model.NewPointer(true),
|
|
ExportFormat: model.NewPointer(model.ComplianceExportTypeActiance),
|
|
DailyRunTime: model.NewPointer("01:00"),
|
|
ExportFromTimestamp: model.NewPointer(int64(0)),
|
|
BatchSize: model.NewPointer(10000),
|
|
ChannelBatchSize: model.NewPointer(100),
|
|
ChannelHistoryBatchSize: model.NewPointer(100),
|
|
},
|
|
},
|
|
},
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
now := time.Now()
|
|
worker.initJobData(request.EmptyContext(logger), logger, job, now)
|
|
|
|
assert.Equal(t, model.ComplianceExportTypeActiance, job.Data[shared.JobDataExportType])
|
|
assert.Equal(t, strconv.Itoa(*worker.jobServer.Config().MessageExportSettings.BatchSize), job.Data[shared.JobDataBatchSize])
|
|
|
|
// Assert the new job picks up the <10.5 job start time:
|
|
assert.Equal(t, previousJob.Data[shared.JobDataBatchStartTime], job.Data[shared.JobDataBatchStartTime])
|
|
|
|
expectedDir := "this-is-the-export-dir"
|
|
assert.Equal(t, expectedDir, job.Data[shared.JobDataExportDir])
|
|
}
|
|
|
|
func TestDoJobNoPostsToExport(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
mockMetrics := &mocks.MetricsInterface{}
|
|
defer mockMetrics.AssertExpectations(t)
|
|
|
|
job := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusPending,
|
|
Type: model.JobTypeMessageExport,
|
|
}
|
|
retJob := *job
|
|
retJob.Status = model.JobStatusInProgress
|
|
|
|
// claim job succeeds
|
|
mockStore.JobStore.
|
|
On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).
|
|
Return(&retJob, nil)
|
|
mockMetrics.On("IncrementJobActive", model.JobTypeMessageExport)
|
|
|
|
// no previous job, data will be loaded from config
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return(nil, errors.New("test"))
|
|
|
|
// no channels with activity
|
|
mockStore.ChannelMemberHistoryStore.On("GetChannelsWithActivityDuring", mock.Anything, mock.Anything).
|
|
Return(make([]string, 0), nil)
|
|
|
|
// no posts found to export
|
|
mockStore.ComplianceStore.On("MessageExport", mock.Anything, mock.AnythingOfType("model.MessageExportCursor"), 10001).Return(
|
|
make([]*model.MessageExport, 0), model.MessageExportCursor{}, nil,
|
|
)
|
|
|
|
mockStore.PostStore.On("AnalyticsPostCount", mock.Anything).Return(
|
|
int64(shared.EstimatedPostCount), nil,
|
|
)
|
|
|
|
// job completed successfully
|
|
mockStore.JobStore.On("UpdateOptimistically", mock.AnythingOfType("*model.Job"), model.JobStatusInProgress).Return(true, nil)
|
|
mockStore.JobStore.On("UpdateStatus", job.Id, model.JobStatusSuccess).Return(job, nil)
|
|
mockMetrics.On("DecrementJobActive", model.JobTypeMessageExport)
|
|
|
|
tempDir, err := os.MkdirTemp("", "")
|
|
require.NoError(t, err)
|
|
t.Cleanup(func() {
|
|
err = os.RemoveAll(tempDir)
|
|
assert.NoError(t, err)
|
|
})
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: jobs.NewJobServer(
|
|
&testutils.StaticConfigService{
|
|
Cfg: &model.Config{
|
|
// mock config
|
|
FileSettings: model.FileSettings{
|
|
DriverName: model.NewPointer(model.ImageDriverLocal),
|
|
Directory: model.NewPointer(tempDir),
|
|
},
|
|
MessageExportSettings: model.MessageExportSettings{
|
|
EnableExport: model.NewPointer(true),
|
|
ExportFormat: model.NewPointer(model.ComplianceExportTypeActiance),
|
|
DailyRunTime: model.NewPointer("01:00"),
|
|
ExportFromTimestamp: model.NewPointer(int64(0)),
|
|
BatchSize: model.NewPointer(10000),
|
|
ChannelBatchSize: model.NewPointer(100),
|
|
ChannelHistoryBatchSize: model.NewPointer(100),
|
|
},
|
|
},
|
|
},
|
|
mockStore,
|
|
mockMetrics,
|
|
logger,
|
|
),
|
|
logger: logger,
|
|
}
|
|
|
|
// actually execute the code under test
|
|
worker.DoJob(job)
|
|
}
|
|
|
|
func TestDoJobWithDedicatedExportBackend(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
mockMetrics := &mocks.MetricsInterface{}
|
|
defer mockMetrics.AssertExpectations(t)
|
|
|
|
job := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusPending,
|
|
Type: model.JobTypeMessageExport,
|
|
}
|
|
retJob := *job
|
|
retJob.Status = model.JobStatusInProgress
|
|
|
|
// claim job succeeds
|
|
mockStore.JobStore.
|
|
On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).
|
|
Return(&retJob, nil)
|
|
mockMetrics.On("IncrementJobActive", model.JobTypeMessageExport)
|
|
|
|
// no previous job, data will be loaded from config
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return(nil, errors.New("test"))
|
|
|
|
channelId := st.NewTestID()
|
|
channelName := st.NewTestID()
|
|
channelDisplayName := st.NewTestID()
|
|
channelType := model.ChannelTypeOpen
|
|
messages := []*model.MessageExport{
|
|
{
|
|
TeamId: model.NewPointer(st.NewTestID()),
|
|
ChannelId: model.NewPointer(channelId),
|
|
ChannelName: model.NewPointer(channelName),
|
|
UserId: model.NewPointer(st.NewTestID()),
|
|
UserEmail: model.NewPointer(st.NewTestID()),
|
|
Username: model.NewPointer(st.NewTestID()),
|
|
PostId: model.NewPointer(st.NewTestID()),
|
|
PostCreateAt: model.NewPointer[int64](123),
|
|
PostUpdateAt: model.NewPointer[int64](123),
|
|
PostDeleteAt: model.NewPointer[int64](123),
|
|
PostMessage: model.NewPointer(st.NewTestID()),
|
|
},
|
|
}
|
|
|
|
// need to export at least one post to make an export directory and file
|
|
|
|
mockStore.ChannelMemberHistoryStore.On("GetChannelsWithActivityDuring", mock.Anything, mock.Anything).
|
|
Return([]string{*messages[0].ChannelId}, nil)
|
|
mockStore.ChannelStore.On("GetMany", []string{channelId}, true).
|
|
Return(model.ChannelList{{
|
|
Id: channelId,
|
|
DisplayName: channelDisplayName,
|
|
Name: channelName,
|
|
Type: channelType,
|
|
}}, nil)
|
|
|
|
mockStore.ComplianceStore.On("MessageExport", mock.Anything, mock.AnythingOfType("model.MessageExportCursor"), 10001).Return(
|
|
messages, model.MessageExportCursor{}, nil,
|
|
).Once()
|
|
mockStore.ComplianceStore.On("MessageExport", mock.Anything, mock.AnythingOfType("model.MessageExportCursor"), 10001).Return(
|
|
make([]*model.MessageExport, 0), model.MessageExportCursor{}, nil,
|
|
).Once()
|
|
mockStore.ChannelMemberHistoryStore.On("GetUsersInChannelDuring", mock.Anything, mock.Anything, mock.Anything).Return(nil, nil)
|
|
|
|
mockStore.PostStore.On("AnalyticsPostCount", mock.Anything).Return(
|
|
int64(1), nil,
|
|
)
|
|
|
|
// job completed successfully
|
|
mockStore.JobStore.On("UpdateOptimistically", mock.AnythingOfType("*model.Job"), model.JobStatusInProgress).Return(true, nil)
|
|
mockStore.JobStore.On("UpdateStatus", job.Id, model.JobStatusSuccess).Return(job, nil)
|
|
mockMetrics.On("DecrementJobActive", model.JobTypeMessageExport)
|
|
|
|
// create primary filestore directory
|
|
tempPrimaryDir, err := os.MkdirTemp("", "")
|
|
require.NoError(t, err)
|
|
defer os.RemoveAll(tempPrimaryDir)
|
|
|
|
// create dedicated filestore directory
|
|
tempDedicatedDir, err := os.MkdirTemp("", "")
|
|
require.NoError(t, err)
|
|
defer os.RemoveAll(tempDedicatedDir)
|
|
|
|
// setup worker with primary and dedicated filestores.
|
|
worker := &MessageExportWorker{
|
|
jobServer: jobs.NewJobServer(
|
|
&testutils.StaticConfigService{
|
|
Cfg: &model.Config{
|
|
// mock config
|
|
FileSettings: model.FileSettings{
|
|
DriverName: model.NewPointer(model.ImageDriverLocal),
|
|
Directory: model.NewPointer(tempPrimaryDir),
|
|
DedicatedExportStore: model.NewPointer(true),
|
|
ExportDriverName: model.NewPointer(model.ImageDriverLocal),
|
|
ExportDirectory: model.NewPointer(tempDedicatedDir),
|
|
},
|
|
MessageExportSettings: model.MessageExportSettings{
|
|
EnableExport: model.NewPointer(true),
|
|
ExportFormat: model.NewPointer(model.ComplianceExportTypeActiance),
|
|
DailyRunTime: model.NewPointer("01:00"),
|
|
ExportFromTimestamp: model.NewPointer(int64(0)),
|
|
BatchSize: model.NewPointer(10000),
|
|
ChannelBatchSize: model.NewPointer(100),
|
|
ChannelHistoryBatchSize: model.NewPointer(100),
|
|
},
|
|
},
|
|
},
|
|
mockStore,
|
|
mockMetrics,
|
|
logger,
|
|
),
|
|
logger: logger,
|
|
}
|
|
|
|
// actually execute the code under test
|
|
worker.DoJob(job)
|
|
|
|
// ensure no primary filestore files exist
|
|
files, err := os.ReadDir(tempPrimaryDir)
|
|
require.NoError(t, err)
|
|
assert.Zero(t, len(files))
|
|
|
|
// ensure some dedicated filestore files exist
|
|
files, err = os.ReadDir(tempDedicatedDir)
|
|
require.NoError(t, err)
|
|
assert.NotZero(t, len(files))
|
|
}
|
|
|
|
func TestDoJobCancel(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
|
|
mockStore := &storetest.Store{}
|
|
t.Cleanup(func() { mockStore.AssertExpectations(t) })
|
|
mockMetrics := &mocks.MetricsInterface{}
|
|
t.Cleanup(func() { mockMetrics.AssertExpectations(t) })
|
|
|
|
job := &model.Job{
|
|
Id: st.NewTestID(),
|
|
CreateAt: model.GetMillis(),
|
|
Status: model.JobStatusPending,
|
|
Type: model.JobTypeMessageExport,
|
|
}
|
|
|
|
tempDir, err := os.MkdirTemp("", "")
|
|
require.NoError(t, err)
|
|
t.Cleanup(func() { os.RemoveAll(tempDir) })
|
|
|
|
impl := MessageExportJobInterfaceImpl{
|
|
Server: &app.Server{
|
|
Jobs: jobs.NewJobServer(
|
|
&testutils.StaticConfigService{
|
|
Cfg: &model.Config{
|
|
// mock config
|
|
FileSettings: model.FileSettings{
|
|
DriverName: model.NewPointer(model.ImageDriverLocal),
|
|
Directory: model.NewPointer(tempDir),
|
|
},
|
|
MessageExportSettings: model.MessageExportSettings{
|
|
EnableExport: model.NewPointer(true),
|
|
ExportFormat: model.NewPointer(model.ComplianceExportTypeActiance),
|
|
DailyRunTime: model.NewPointer("01:00"),
|
|
ExportFromTimestamp: model.NewPointer(int64(0)),
|
|
BatchSize: model.NewPointer(10000),
|
|
ChannelBatchSize: model.NewPointer(100),
|
|
ChannelHistoryBatchSize: model.NewPointer(100),
|
|
},
|
|
},
|
|
},
|
|
mockStore,
|
|
mockMetrics,
|
|
logger,
|
|
),
|
|
},
|
|
}
|
|
worker, ok := impl.MakeWorker().(*MessageExportWorker)
|
|
require.True(t, ok)
|
|
|
|
retJob := *job
|
|
retJob.Status = model.JobStatusInProgress
|
|
|
|
// Claim job succeeds
|
|
mockStore.JobStore.
|
|
On("UpdateStatusOptimistically", job.Id, model.JobStatusPending, model.JobStatusInProgress).
|
|
Return(&retJob, nil)
|
|
mockMetrics.On("IncrementJobActive", model.JobTypeMessageExport)
|
|
|
|
// No previous job, data will be loaded from config
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return(nil, errors.New("test"))
|
|
|
|
// Job updates the system console UI, once for getting channels, once for getting activity
|
|
mockStore.JobStore.On("UpdateOptimistically", mock.AnythingOfType("*model.Job"), model.JobStatusInProgress).Return(true, nil).Times(2)
|
|
|
|
// a few calls pass
|
|
mockStore.ChannelMemberHistoryStore.On("GetChannelsWithActivityDuring", mock.Anything, mock.Anything).
|
|
Return([]string{"channel-id"}, nil)
|
|
mockStore.ChannelStore.On("GetMany", []string{"channel-id"}, true).
|
|
Return(model.ChannelList{{
|
|
Id: "channel-id",
|
|
DisplayName: "channel-display-name",
|
|
Name: "channel-name",
|
|
Type: model.ChannelTypeDirect,
|
|
}}, nil)
|
|
mockStore.ChannelMemberHistoryStore.On("GetUsersInChannelDuring", mock.Anything, mock.Anything, []string{"channel-id"}).Return([]*model.ChannelMemberHistoryResult{}, nil)
|
|
|
|
cancelled := make(chan struct{})
|
|
// Cancel the worker and return an error
|
|
mockStore.ComplianceStore.On("MessageExport", mock.Anything, mock.AnythingOfType("model.MessageExportCursor"), 10001).Run(func(args tmock.Arguments) {
|
|
worker.cancel()
|
|
|
|
rctx, ok := args.Get(0).(request.CTX)
|
|
require.True(t, ok)
|
|
assert.Error(t, rctx.Context().Err())
|
|
assert.ErrorIs(t, rctx.Context().Err(), context.Canceled)
|
|
|
|
cancelled <- struct{}{}
|
|
}).Return(
|
|
nil, model.MessageExportCursor{}, context.Canceled,
|
|
)
|
|
|
|
mockStore.PostStore.On("AnalyticsPostCount", mock.Anything).Return(
|
|
int64(shared.EstimatedPostCount), nil,
|
|
)
|
|
|
|
// Job marked as pending
|
|
mockStore.JobStore.On("UpdateStatus", job.Id, model.JobStatusPending).Return(job, nil)
|
|
mockMetrics.On("DecrementJobActive", model.JobTypeMessageExport)
|
|
|
|
go worker.Run()
|
|
|
|
worker.JobChannel() <- *job
|
|
|
|
// Wait for the cancelation
|
|
<-cancelled
|
|
|
|
// Cleanup
|
|
worker.Stop()
|
|
}
|
|
|
|
func TestGetPreviousJobNoJobs(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
// Mock the job store to return empty jobs list
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return([]*model.Job{}, nil).Once()
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
rctx := request.EmptyContext(logger)
|
|
job, err := worker.getPreviousNonCliJob(rctx)
|
|
|
|
require.NoError(t, err)
|
|
assert.Nil(t, job, "Expected nil job when no jobs are returned")
|
|
}
|
|
|
|
func TestGetPreviousJobOneRegularJob(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
regularJob := &model.Job{
|
|
Id: st.NewTestID(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
Data: map[string]string{},
|
|
}
|
|
|
|
// Mock the job store to return one regular job
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return([]*model.Job{regularJob}, nil).Once()
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
rctx := request.EmptyContext(logger)
|
|
job, err := worker.getPreviousNonCliJob(rctx)
|
|
|
|
require.NoError(t, err)
|
|
assert.Equal(t, regularJob.Id, job.Id, "Expected to get the regular job")
|
|
}
|
|
|
|
func TestGetPreviousJobOneMmctlJob(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
mmctlJob := &model.Job{
|
|
Id: st.NewTestID(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
Data: map[string]string{shared.JobDataInitiatedBy: "mmctl"},
|
|
}
|
|
|
|
// Mock the job store to return only mmctl jobs (4 jobs, not a full page)
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return([]*model.Job{mmctlJob, mmctlJob, mmctlJob, mmctlJob}, nil).Once()
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
rctx := request.EmptyContext(logger)
|
|
job, err := worker.getPreviousNonCliJob(rctx)
|
|
|
|
require.NoError(t, err)
|
|
assert.Nil(t, job, "Expected nil job when only mmctl jobs are found")
|
|
}
|
|
|
|
func TestGetPreviousJobManyJobs(t *testing.T) {
|
|
logger := mlog.CreateConsoleTestLogger(t)
|
|
mockStore := &storetest.Store{}
|
|
defer mockStore.AssertExpectations(t)
|
|
|
|
// Create DefaultPageSize mmctl jobs for first page
|
|
firstPageJobs := make([]*model.Job, DefaultPreviousJobPageSize)
|
|
for i := range DefaultPreviousJobPageSize {
|
|
firstPageJobs[i] = &model.Job{
|
|
Id: st.NewTestID(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
Data: map[string]string{shared.JobDataInitiatedBy: "mmctl"},
|
|
}
|
|
}
|
|
|
|
// Create DefaultPageSize mmctl jobs for second page
|
|
secondPageJobs := make([]*model.Job, DefaultPreviousJobPageSize)
|
|
for i := range DefaultPreviousJobPageSize {
|
|
secondPageJobs[i] = &model.Job{
|
|
Id: st.NewTestID(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
Data: map[string]string{shared.JobDataInitiatedBy: "mmctl"},
|
|
}
|
|
}
|
|
|
|
// Create 1 regular job for the third page (last job)
|
|
regularJob := &model.Job{
|
|
Id: st.NewTestID(),
|
|
Status: model.JobStatusSuccess,
|
|
Type: model.JobTypeMessageExport,
|
|
Data: map[string]string{},
|
|
}
|
|
thirdPageJobs := []*model.Job{regularJob}
|
|
|
|
// Mock the job store to return the jobs in pages
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
0, DefaultPreviousJobPageSize).Return(firstPageJobs, nil).Once()
|
|
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
1*DefaultPreviousJobPageSize, DefaultPreviousJobPageSize).Return(secondPageJobs, nil).Once()
|
|
|
|
mockStore.JobStore.On("GetAllByTypesAndStatusesPage", mock.Anything,
|
|
[]string{model.JobTypeMessageExport},
|
|
[]string{model.JobStatusWarning, model.JobStatusSuccess},
|
|
2*DefaultPreviousJobPageSize, DefaultPreviousJobPageSize).Return(thirdPageJobs, nil).Once()
|
|
|
|
worker := &MessageExportWorker{
|
|
jobServer: &jobs.JobServer{
|
|
Store: mockStore,
|
|
},
|
|
logger: logger,
|
|
}
|
|
|
|
rctx := request.EmptyContext(logger)
|
|
job, err := worker.getPreviousNonCliJob(rctx)
|
|
|
|
require.NoError(t, err)
|
|
assert.NotNil(t, job)
|
|
assert.Equal(t, regularJob.Id, job.Id, "Expected to find the regular job at the end")
|
|
}
|