[MM-33853] remove CSV row limit in compliance monitoring (#17185)

Этот коммит содержится в:
Max Erenberg
2021-04-23 09:19:13 -04:00
коммит произвёл GitHub
родитель ff657bfdef
Коммит 9cd50a4e22
9 изменённых файлов: 229 добавлений и 72 удалений

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

@@ -37,6 +37,19 @@ type Compliance struct {
type Compliances []Compliance
// ComplianceExportCursor is used for paginated iteration of posts
// for compliance export.
// We need to keep track of the last post ID in addition to the last post
// CreateAt to break ties when two posts have the same CreateAt.
type ComplianceExportCursor struct {
LastChannelsQueryPostCreateAt int64
LastChannelsQueryPostID string
ChannelsQueryCompleted bool
LastDirectMessagesQueryPostCreateAt int64
LastDirectMessagesQueryPostID string
DirectMessagesQueryCompleted bool
}
func (c *Compliance) ToJson() string {
b, _ := json.Marshal(c)
return string(b)

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

@@ -2301,6 +2301,7 @@ type ComplianceSettings struct {
Enable *bool `access:"compliance_compliance_monitoring"`
Directory *string `access:"compliance_compliance_monitoring"` // telemetry: none
EnableDaily *bool `access:"compliance_compliance_monitoring"`
BatchSize *int `access:"compliance_compliance_monitoring"` // telemetry: none
}
func (s *ComplianceSettings) SetDefaults() {
@@ -2315,6 +2316,10 @@ func (s *ComplianceSettings) SetDefaults() {
if s.EnableDaily == nil {
s.EnableDaily = NewBool(false)
}
if s.BatchSize == nil {
s.BatchSize = NewInt(30000)
}
}
type LocalizationSettings struct {

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

@@ -2747,7 +2747,7 @@ func (s *OpenTracingLayerCommandWebhookStore) TryUse(id string, limit int) error
return err
}
func (s *OpenTracingLayerComplianceStore) ComplianceExport(compliance *model.Compliance) ([]*model.CompliancePost, error) {
func (s *OpenTracingLayerComplianceStore) ComplianceExport(compliance *model.Compliance, cursor model.ComplianceExportCursor, limit int) ([]*model.CompliancePost, model.ComplianceExportCursor, error) {
origCtx := s.Root.Store.Context()
span, newCtx := tracing.StartSpanWithParentByContext(s.Root.Store.Context(), "ComplianceStore.ComplianceExport")
s.Root.Store.SetContext(newCtx)
@@ -2756,13 +2756,13 @@ func (s *OpenTracingLayerComplianceStore) ComplianceExport(compliance *model.Com
}()
defer span.Finish()
result, err := s.ComplianceStore.ComplianceExport(compliance)
result, resultVar1, err := s.ComplianceStore.ComplianceExport(compliance, cursor, limit)
if err != nil {
span.LogFields(spanlog.Error(err))
ext.Error.Set(span, true)
}
return result, err
return result, resultVar1, err
}
func (s *OpenTracingLayerComplianceStore) Get(id string) (*model.Compliance, error) {

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

@@ -2934,21 +2934,21 @@ func (s *RetryLayerCommandWebhookStore) TryUse(id string, limit int) error {
}
func (s *RetryLayerComplianceStore) ComplianceExport(compliance *model.Compliance) ([]*model.CompliancePost, error) {
func (s *RetryLayerComplianceStore) ComplianceExport(compliance *model.Compliance, cursor model.ComplianceExportCursor, limit int) ([]*model.CompliancePost, model.ComplianceExportCursor, error) {
tries := 0
for {
result, err := s.ComplianceStore.ComplianceExport(compliance)
result, resultVar1, err := s.ComplianceStore.ComplianceExport(compliance, cursor, limit)
if err == nil {
return result, nil
return result, resultVar1, nil
}
if !isRepeatableError(err) {
return result, err
return result, resultVar1, err
}
tries++
if tries >= 3 {
err = errors.Wrap(err, "giving up after 3 consecutive repeatable transaction failures")
return result, err
return result, resultVar1, err
}
}

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

@@ -81,50 +81,49 @@ func (s SqlComplianceStore) Get(id string) (*model.Compliance, error) {
return obj.(*model.Compliance), nil
}
func (s SqlComplianceStore) ComplianceExport(job *model.Compliance) ([]*model.CompliancePost, error) {
props := map[string]interface{}{"StartTime": job.StartAt, "EndTime": job.EndAt}
func (s SqlComplianceStore) ComplianceExport(job *model.Compliance, cursor model.ComplianceExportCursor, limit int) ([]*model.CompliancePost, model.ComplianceExportCursor, error) {
props := map[string]interface{}{"EndTime": job.EndAt, "Limit": limit}
keywordQuery := ""
keywords := strings.Fields(strings.TrimSpace(strings.ToLower(strings.Replace(job.Keywords, ",", " ", -1))))
if len(keywords) > 0 {
clauses := make([]string, len(keywords))
keywordQuery = "AND ("
for index, keyword := range keywords {
for i, keyword := range keywords {
keyword = sanitizeSearchTerm(keyword, "\\")
if index >= 1 {
keywordQuery += " OR LOWER(Posts.Message) LIKE :Keyword" + strconv.Itoa(index)
} else {
keywordQuery += "LOWER(Posts.Message) LIKE :Keyword" + strconv.Itoa(index)
}
props["Keyword"+strconv.Itoa(index)] = "%" + keyword + "%"
clauses[i] = "LOWER(Posts.Message) LIKE :Keyword" + strconv.Itoa(i)
props["Keyword"+strconv.Itoa(i)] = "%" + keyword + "%"
}
keywordQuery += ")"
keywordQuery = "AND (" + strings.Join(clauses, " OR ") + ")"
}
emailQuery := ""
emails := strings.Fields(strings.TrimSpace(strings.ToLower(strings.Replace(job.Emails, ",", " ", -1))))
if len(emails) > 0 {
clauses := make([]string, len(emails))
emailQuery = "AND ("
for index, email := range emails {
if index >= 1 {
emailQuery += " OR Users.Email = :Email" + strconv.Itoa(index)
} else {
emailQuery += "Users.Email = :Email" + strconv.Itoa(index)
}
props["Email"+strconv.Itoa(index)] = email
for i, email := range emails {
clauses[i] = "Users.Email = :Email" + strconv.Itoa(i)
props["Email"+strconv.Itoa(i)] = email
}
emailQuery += ")"
emailQuery = "AND (" + strings.Join(clauses, " OR ") + ")"
}
query :=
`(SELECT
// The idea is to first iterate over the channel posts, and then when we run out of those,
// start iterating over the direct message posts.
var channelPosts []*model.CompliancePost
channelsQuery := ""
if !cursor.ChannelsQueryCompleted {
if cursor.LastChannelsQueryPostCreateAt == 0 {
cursor.LastChannelsQueryPostCreateAt = job.StartAt
}
props["LastPostCreateAt"] = cursor.LastChannelsQueryPostCreateAt
props["LastPostId"] = cursor.LastChannelsQueryPostID
channelsQuery = `
SELECT
Teams.Name AS TeamName,
Teams.DisplayName AS TeamDisplayName,
Channels.Name AS ChannelName,
@@ -151,17 +150,43 @@ func (s SqlComplianceStore) ComplianceExport(job *model.Compliance) ([]*model.Co
Channels,
Users,
Posts
LEFT JOIN Bots ON Bots.UserId = Posts.UserId
LEFT JOIN
Bots ON Bots.UserId = Posts.UserId
WHERE
Teams.Id = Channels.TeamId
AND Posts.ChannelId = Channels.Id
AND Posts.UserId = Users.Id
AND Posts.CreateAt > :StartTime
AND Posts.CreateAt <= :EndTime
AND (
Posts.CreateAt > :LastPostCreateAt
OR (Posts.CreateAt = :LastPostCreateAt AND Posts.Id > :LastPostId)
)
AND Posts.CreateAt < :EndTime
` + emailQuery + `
` + keywordQuery + `)
UNION ALL
(SELECT
` + keywordQuery + `
ORDER BY Posts.CreateAt, Posts.Id
LIMIT :Limit`
if _, err := s.GetReplica().Select(&channelPosts, channelsQuery, props); err != nil {
return nil, cursor, errors.Wrap(err, "unable to export compliance")
}
if len(channelPosts) < limit {
cursor.ChannelsQueryCompleted = true
} else {
cursor.LastChannelsQueryPostCreateAt = channelPosts[len(channelPosts)-1].PostCreateAt
cursor.LastChannelsQueryPostID = channelPosts[len(channelPosts)-1].PostId
}
}
var directMessagePosts []*model.CompliancePost
directMessagesQuery := ""
if !cursor.DirectMessagesQueryCompleted && len(channelPosts) < limit {
if cursor.LastDirectMessagesQueryPostCreateAt == 0 {
cursor.LastDirectMessagesQueryPostCreateAt = job.StartAt
}
props["LastPostCreateAt"] = cursor.LastDirectMessagesQueryPostCreateAt
props["LastPostId"] = cursor.LastDirectMessagesQueryPostID
props["Limit"] = limit - len(channelPosts)
directMessagesQuery = `
SELECT
'direct-messages' AS TeamName,
'Direct Messages' AS TeamDisplayName,
Channels.Name AS ChannelName,
@@ -187,24 +212,34 @@ func (s SqlComplianceStore) ComplianceExport(job *model.Compliance) ([]*model.Co
Channels,
Users,
Posts
LEFT JOIN Bots ON Bots.UserId = Posts.UserId
LEFT JOIN
Bots ON Bots.UserId = Posts.UserId
WHERE
Channels.TeamId = ''
AND Posts.ChannelId = Channels.Id
AND Posts.UserId = Users.Id
AND Posts.CreateAt > :StartTime
AND Posts.CreateAt <= :EndTime
AND (
Posts.CreateAt > :LastPostCreateAt
OR (Posts.CreateAt = :LastPostCreateAt AND Posts.Id > :LastPostId)
)
AND Posts.CreateAt < :EndTime
` + emailQuery + `
` + keywordQuery + `)
ORDER BY PostCreateAt
LIMIT 30000`
` + keywordQuery + `
ORDER BY Posts.CreateAt, Posts.Id
LIMIT :Limit`
var cposts []*model.CompliancePost
if _, err := s.GetReplica().Select(&cposts, query, props); err != nil {
return nil, errors.Wrap(err, "unable to export compliance")
if _, err := s.GetReplica().Select(&directMessagePosts, directMessagesQuery, props); err != nil {
return nil, cursor, errors.Wrap(err, "unable to export compliance")
}
if len(directMessagePosts) < limit {
cursor.DirectMessagesQueryCompleted = true
} else {
cursor.LastDirectMessagesQueryPostCreateAt = directMessagePosts[len(directMessagePosts)-1].PostCreateAt
cursor.LastDirectMessagesQueryPostID = directMessagePosts[len(directMessagePosts)-1].PostId
}
}
return cposts, nil
return append(channelPosts, directMessagePosts...), cursor, nil
}
func (s SqlComplianceStore) MessageExport(after int64, limit int) ([]*model.MessageExport, error) {

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

@@ -468,7 +468,7 @@ type ComplianceStore interface {
Update(compliance *model.Compliance) (*model.Compliance, error)
Get(id string) (*model.Compliance, error)
GetAll(offset, limit int) (model.Compliances, error)
ComplianceExport(compliance *model.Compliance) ([]*model.CompliancePost, error)
ComplianceExport(compliance *model.Compliance, cursor model.ComplianceExportCursor, limit int) ([]*model.CompliancePost, model.ComplianceExportCursor, error)
MessageExport(after int64, limit int) ([]*model.MessageExport, error)
}

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

@@ -98,6 +98,9 @@ func testComplianceStore(t *testing.T, ss store.Store) {
func testComplianceExport(t *testing.T, ss store.Store) {
time.Sleep(100 * time.Millisecond)
const (
limit = 30000
)
t1 := &model.Team{}
t1.DisplayName = "DisplayName"
@@ -166,47 +169,67 @@ func testComplianceExport(t *testing.T, ss store.Store) {
time.Sleep(100 * time.Millisecond)
cr1 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o2a.CreateAt + 1}
cposts, nErr := ss.Compliance().ComplianceExport(cr1)
cposts, _, nErr := ss.Compliance().ComplianceExport(cr1, model.ComplianceExportCursor{}, limit)
require.NoError(t, nErr)
assert.Len(t, cposts, 4)
assert.Equal(t, cposts[0].PostId, o1.Id)
assert.Equal(t, cposts[3].PostId, o2a.Id)
cr2 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o2a.CreateAt + 1, Emails: u2.Email}
cposts, nErr = ss.Compliance().ComplianceExport(cr2)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr2, model.ComplianceExportCursor{}, limit)
require.NoError(t, nErr)
assert.Len(t, cposts, 1)
assert.Equal(t, cposts[0].PostId, o2a.Id)
cr3 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o2a.CreateAt + 1, Emails: u2.Email + ", " + u1.Email}
cposts, nErr = ss.Compliance().ComplianceExport(cr3)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr3, model.ComplianceExportCursor{}, limit)
require.NoError(t, nErr)
assert.Len(t, cposts, 4)
assert.Equal(t, cposts[0].PostId, o1.Id)
assert.Equal(t, cposts[3].PostId, o2a.Id)
cr4 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o2a.CreateAt + 1, Keywords: o2a.Message}
cposts, nErr = ss.Compliance().ComplianceExport(cr4)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr4, model.ComplianceExportCursor{}, limit)
require.NoError(t, nErr)
assert.Len(t, cposts, 1)
assert.Equal(t, cposts[0].PostId, o2a.Id)
cr5 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o2a.CreateAt + 1, Keywords: o2a.Message + " " + o1.Message}
cposts, nErr = ss.Compliance().ComplianceExport(cr5)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr5, model.ComplianceExportCursor{}, limit)
require.NoError(t, nErr)
assert.Len(t, cposts, 2)
assert.Equal(t, cposts[0].PostId, o1.Id)
cr6 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o2a.CreateAt + 1, Emails: u2.Email + ", " + u1.Email, Keywords: o2a.Message + " " + o1.Message}
cposts, nErr = ss.Compliance().ComplianceExport(cr6)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr6, model.ComplianceExportCursor{}, limit)
require.NoError(t, nErr)
assert.Len(t, cposts, 2)
assert.Equal(t, cposts[0].PostId, o1.Id)
assert.Equal(t, cposts[1].PostId, o2a.Id)
t.Run("multiple batches", func(t *testing.T) {
cr7 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o2a.CreateAt + 1}
cursor := model.ComplianceExportCursor{}
cposts, cursor, nErr = ss.Compliance().ComplianceExport(cr7, cursor, 2)
require.NoError(t, nErr)
assert.Len(t, cposts, 2)
assert.Equal(t, cposts[0].PostId, o1.Id)
assert.Equal(t, cposts[1].PostId, o1a.Id)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr7, cursor, 3)
require.NoError(t, nErr)
assert.Len(t, cposts, 2)
assert.Equal(t, cposts[0].PostId, o2.Id)
assert.Equal(t, cposts[1].PostId, o2a.Id)
})
}
func testComplianceExportDirectMessages(t *testing.T, ss store.Store) {
defer cleanupStoreState(t, ss)
time.Sleep(100 * time.Millisecond)
const (
limit = 30000
)
t1 := &model.Team{}
t1.DisplayName = "DisplayName"
@@ -285,11 +308,85 @@ func testComplianceExportDirectMessages(t *testing.T, ss store.Store) {
time.Sleep(100 * time.Millisecond)
cr1 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o3.CreateAt + 1, Emails: u1.Email}
cposts, nErr := ss.Compliance().ComplianceExport(cr1)
cposts, _, nErr := ss.Compliance().ComplianceExport(cr1, model.ComplianceExportCursor{}, limit)
require.NoError(t, nErr)
assert.Len(t, cposts, 4)
assert.Equal(t, cposts[0].PostId, o1.Id)
assert.Equal(t, cposts[len(cposts)-1].PostId, o3.Id)
t.Run("mix of channel and direct messages", func(t *testing.T) {
// This will "cross the boundary" between the two queries
cursor := model.ComplianceExportCursor{}
cr2 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o3.CreateAt + 1, Emails: u1.Email}
cposts, cursor, nErr = ss.Compliance().ComplianceExport(cr2, cursor, 2)
require.NoError(t, nErr)
assert.Len(t, cposts, 2)
assert.Equal(t, cposts[0].PostId, o1.Id)
assert.Equal(t, cposts[len(cposts)-1].PostId, o1a.Id)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr2, cursor, 2)
require.NoError(t, nErr)
assert.Len(t, cposts, 2)
assert.Equal(t, cposts[0].PostId, o2.Id)
assert.Equal(t, cposts[len(cposts)-1].PostId, o3.Id)
// This will exhaust the first query before moving to the next one
cursor = model.ComplianceExportCursor{}
cr3 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: o1.CreateAt - 1, EndAt: o3.CreateAt + 1, Emails: u1.Email}
cposts, cursor, nErr = ss.Compliance().ComplianceExport(cr3, cursor, 3)
require.NoError(t, nErr)
assert.Len(t, cposts, 3)
assert.Equal(t, cposts[0].PostId, o1.Id)
assert.Equal(t, cposts[len(cposts)-1].PostId, o2.Id)
cposts, _, nErr = ss.Compliance().ComplianceExport(cr3, cursor, 2)
require.NoError(t, nErr)
assert.Len(t, cposts, 1)
assert.Equal(t, cposts[0].PostId, o3.Id)
})
t.Run("timestamp collision", func(t *testing.T) {
time.Sleep(100 * time.Millisecond)
nowMillis := model.GetMillis()
createPost := func(createAt int64) {
post := &model.Post{}
post.ChannelId = c1.Id
post.UserId = u1.Id
post.CreateAt = createAt
post.Message = "zz" + model.NewId() + "b"
post, nErr = ss.Post().Save(post)
require.NoError(t, nErr)
}
for i := 0; i < 3; i++ {
createPost(nowMillis)
}
for i := 0; i < 2; i++ {
createPost(nowMillis + 1)
}
cursor := model.ComplianceExportCursor{}
cr4 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: nowMillis, EndAt: nowMillis + 2}
cposts, cursor, nErr = ss.Compliance().ComplianceExport(cr4, cursor, 2)
require.NoError(t, nErr)
assert.Len(t, cposts, 2)
cr5 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: nowMillis, EndAt: nowMillis + 2}
cposts, _, nErr = ss.Compliance().ComplianceExport(cr5, cursor, 3)
require.NoError(t, nErr)
assert.Len(t, cposts, 3)
// range should be [inclusive, exclusive)
cursor = model.ComplianceExportCursor{}
cr6 := &model.Compliance{Desc: "test" + model.NewId(), StartAt: nowMillis, EndAt: nowMillis + 1}
cposts, _, nErr = ss.Compliance().ComplianceExport(cr6, cursor, 5)
require.NoError(t, nErr)
assert.Len(t, cposts, 3)
})
}
func testMessageExportPublicChannel(t *testing.T, ss store.Store) {

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

@@ -14,27 +14,34 @@ type ComplianceStore struct {
mock.Mock
}
// ComplianceExport provides a mock function with given fields: compliance
func (_m *ComplianceStore) ComplianceExport(compliance *model.Compliance) ([]*model.CompliancePost, error) {
ret := _m.Called(compliance)
// ComplianceExport provides a mock function with given fields: compliance, cursor, limit
func (_m *ComplianceStore) ComplianceExport(compliance *model.Compliance, cursor model.ComplianceExportCursor, limit int) ([]*model.CompliancePost, model.ComplianceExportCursor, error) {
ret := _m.Called(compliance, cursor, limit)
var r0 []*model.CompliancePost
if rf, ok := ret.Get(0).(func(*model.Compliance) []*model.CompliancePost); ok {
r0 = rf(compliance)
if rf, ok := ret.Get(0).(func(*model.Compliance, model.ComplianceExportCursor, int) []*model.CompliancePost); ok {
r0 = rf(compliance, cursor, limit)
} else {
if ret.Get(0) != nil {
r0 = ret.Get(0).([]*model.CompliancePost)
}
}
var r1 error
if rf, ok := ret.Get(1).(func(*model.Compliance) error); ok {
r1 = rf(compliance)
var r1 model.ComplianceExportCursor
if rf, ok := ret.Get(1).(func(*model.Compliance, model.ComplianceExportCursor, int) model.ComplianceExportCursor); ok {
r1 = rf(compliance, cursor, limit)
} else {
r1 = ret.Error(1)
r1 = ret.Get(1).(model.ComplianceExportCursor)
}
return r0, r1
var r2 error
if rf, ok := ret.Get(2).(func(*model.Compliance, model.ComplianceExportCursor, int) error); ok {
r2 = rf(compliance, cursor, limit)
} else {
r2 = ret.Error(2)
}
return r0, r1, r2
}
// Get provides a mock function with given fields: id

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

@@ -2525,10 +2525,10 @@ func (s *TimerLayerCommandWebhookStore) TryUse(id string, limit int) error {
return err
}
func (s *TimerLayerComplianceStore) ComplianceExport(compliance *model.Compliance) ([]*model.CompliancePost, error) {
func (s *TimerLayerComplianceStore) ComplianceExport(compliance *model.Compliance, cursor model.ComplianceExportCursor, limit int) ([]*model.CompliancePost, model.ComplianceExportCursor, error) {
start := timemodule.Now()
result, err := s.ComplianceStore.ComplianceExport(compliance)
result, resultVar1, err := s.ComplianceStore.ComplianceExport(compliance, cursor, limit)
elapsed := float64(timemodule.Since(start)) / float64(timemodule.Second)
if s.Root.Metrics != nil {
@@ -2538,7 +2538,7 @@ func (s *TimerLayerComplianceStore) ComplianceExport(compliance *model.Complianc
}
s.Root.Metrics.ObserveStoreMethodDuration("ComplianceStore.ComplianceExport", success, elapsed)
}
return result, err
return result, resultVar1, err
}
func (s *TimerLayerComplianceStore) Get(id string) (*model.Compliance, error) {