[MM-53291] Data retention improvements (#24253)

* adding new migration for RetentionIdsForDeletion, changing logic for deleting orphaned reactions. Updating delete user and channel endpoints to remove respective reactions
Этот коммит содержится в:
Ben Cooke
2023-09-06 08:25:27 -04:00
коммит произвёл GitHub
родитель d13429aa92
Коммит 791ee40568
27 изменённых файлов: 820 добавлений и 275 удалений

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

@@ -175,6 +175,7 @@ func (s SqlChannelMemberHistoryStore) PermanentDeleteBatchForRetentionPolicies(n
NowMillis: now,
GlobalPolicyEndTime: globalPolicyEndTime,
Limit: limit,
StoreDeletedIds: false,
}, s.SqlStore, cursor)
}

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

@@ -11,6 +11,7 @@ import (
"regexp"
"strings"
"sync"
"time"
sq "github.com/mattermost/squirrel"
"github.com/pkg/errors"
@@ -929,24 +930,31 @@ func (s *SqlPostStore) Delete(postID string, time int64, deleteByID string) (err
return nil
}
func (s *SqlPostStore) permanentDelete(postId string) (err error) {
var post model.Post
func (s *SqlPostStore) permanentDelete(postIds []string) (err error) {
transaction, err := s.GetMasterX().Beginx()
if err != nil {
return errors.Wrap(err, "begin_transaction")
}
defer finalizeTransactionX(transaction, &err)
err = transaction.Get(&post, "SELECT * FROM Posts WHERE Id = ?", postId)
if err != nil && err != sql.ErrNoRows {
return errors.Wrapf(err, "failed to get Post with id=%s", postId)
}
if err = s.permanentDeleteThreads(transaction, post.Id); err != nil {
return errors.Wrapf(err, "failed to cleanup threads for Post with id=%s", postId)
if err = s.permanentDeleteThreads(transaction, postIds); err != nil {
return err
}
if _, err = transaction.NamedExec("DELETE FROM Posts WHERE Id = :id OR RootId = :rootid", map[string]any{"id": postId, "rootid": postId}); err != nil {
return errors.Wrapf(err, "failed to delete Post with id=%s", postId)
if err = s.permanentDeleteReactions(transaction, postIds); err != nil {
return err
}
query := s.getQueryBuilder().
Delete("Posts").
Where(
sq.Or{
sq.Eq{"Id": postIds},
sq.Eq{"RootId": postIds},
},
)
if _, err = transaction.ExecBuilder(query); err != nil {
return errors.Wrap(err, "failed to delete Posts")
}
if err = transaction.Commit(); err != nil {
@@ -980,10 +988,17 @@ func (s *SqlPostStore) permanentDeleteAllCommentByUser(userId string) (err error
return errors.Wrapf(err, "failed to delete Posts with userId=%s", userId)
}
postIds := []string{}
for _, ids := range results {
if err = s.updateThreadAfterReplyDeletion(transaction, ids.RootId, userId); err != nil {
return err
}
postIds = append(postIds, ids.Id)
}
// Delete all the reactions on the comments
if err = s.permanentDeleteReactions(transaction, postIds); err != nil {
return err
}
if err = transaction.Commit(); err != nil {
@@ -1005,22 +1020,20 @@ func (s *SqlPostStore) PermanentDeleteByUser(userId string) error {
// Now attempt to delete all the root posts for a user. This will also
// delete all the comments for each post
found := true
count := 0
for found {
for {
var ids []string
err := s.GetMasterX().Select(&ids, "SELECT Id FROM Posts WHERE UserId = ? LIMIT 1000", userId)
if err != nil {
return errors.Wrapf(err, "failed to find Posts with userId=%s", userId)
}
found = false
for _, id := range ids {
found = true
if err = s.permanentDelete(id); err != nil {
return err
}
if len(ids) == 0 {
break
}
if err = s.permanentDelete(ids); err != nil {
return err
}
// This is a fail safe, give up if more than 10k messages
@@ -1035,6 +1048,7 @@ func (s *SqlPostStore) PermanentDeleteByUser(userId string) error {
// Permanent deletes all channel root posts and comments,
// deletes all threads and thread memberships
// deletes all reactions
// no thread comment cleanup needed, since we are deleting threads and thread memberships
func (s *SqlPostStore) PermanentDeleteByChannel(channelId string) (err error) {
transaction, err := s.GetMasterX().Beginx()
@@ -1043,20 +1057,39 @@ func (s *SqlPostStore) PermanentDeleteByChannel(channelId string) (err error) {
}
defer finalizeTransactionX(transaction, &err)
results := []postIds{}
err = transaction.Select(&results, "SELECT Id, RootId, UserId FROM Posts WHERE ChannelId = ?", channelId)
if err != nil {
return errors.Wrapf(err, "failed to fetch Posts with channelId=%s", channelId)
}
id := ""
for {
ids := []string{}
err = transaction.Select(&ids, "SELECT Id FROM Posts WHERE ChannelId = ? AND Id > ? ORDER BY Id ASC LIMIT 500", channelId, id)
if err != nil {
return errors.Wrapf(err, "failed to fetch Posts with channelId=%s", channelId)
}
for _, ids := range results {
if err = s.permanentDeleteThreads(transaction, ids.Id); err != nil {
if len(ids) == 0 {
break
}
id = ids[len(ids)-1]
if err = s.permanentDeleteThreads(transaction, ids); err != nil {
return err
}
}
time.Sleep(10 * time.Millisecond)
if _, err = transaction.Exec("DELETE FROM Posts WHERE ChannelId = ?", channelId); err != nil {
return errors.Wrapf(err, "failed to delete Posts with channelId=%s", channelId)
if err = s.permanentDeleteReactions(transaction, ids); err != nil {
return err
}
time.Sleep(10 * time.Millisecond)
query := s.getQueryBuilder().
Delete("Posts").
Where(
sq.Eq{"Id": ids},
)
if _, err = transaction.ExecBuilder(query); err != nil {
return errors.Wrap(err, "failed to delete Posts")
}
time.Sleep(10 * time.Millisecond)
}
if err = transaction.Commit(); err != nil {
@@ -2419,48 +2452,10 @@ func (s *SqlPostStore) PermanentDeleteBatchForRetentionPolicies(now, globalPolic
NowMillis: now,
GlobalPolicyEndTime: globalPolicyEndTime,
Limit: limit,
StoreDeletedIds: true,
}, s.SqlStore, cursor)
}
// DeleteOrphanedRows removes entries from Posts when a corresponding channel no longer exists.
func (s *SqlPostStore) DeleteOrphanedRows(limit int) (deleted int64, err error) {
var query string
// We need the extra level of nesting to deal with MySQL's locking
if s.DriverName() == model.DatabaseDriverMysql {
// MySQL fails to do a proper antijoin if the selecting column
// and the joining column are different. In that case, doing a subquery
// leads to a faster plan because MySQL materializes the sub-query
// and does a covering index scan on Posts table. More details on the PR with
// this commit.
query = `
DELETE FROM Posts WHERE Id IN (
SELECT * FROM (
SELECT Posts.Id FROM Posts
WHERE Posts.ChannelId NOT IN (SELECT Id FROM Channels USE INDEX (PRIMARY))
LIMIT ?
) AS A
)`
} else {
query = `
DELETE FROM Posts WHERE Id IN (
SELECT * FROM (
SELECT Posts.Id FROM Posts
LEFT JOIN Channels ON Posts.ChannelId = Channels.Id
WHERE Channels.Id IS NULL
LIMIT ?
) AS A
)`
}
result, err := s.GetMasterX().Exec(query, limit)
if err != nil {
return
}
deleted, err = result.RowsAffected()
return
}
func (s *SqlPostStore) PermanentDeleteBatch(endTime int64, limit int64) (int64, error) {
var query string
if s.DriverName() == "postgres" {
@@ -2779,16 +2774,40 @@ func (s *SqlPostStore) GetOldestEntityCreationTime() (int64, error) {
}
// Deletes a thread and a thread membership if the postId is a root post
func (s *SqlPostStore) permanentDeleteThreads(transaction *sqlxTxWrapper, postId string) error {
if _, err := transaction.Exec("DELETE FROM Threads WHERE PostId = ?", postId); err != nil {
func (s *SqlPostStore) permanentDeleteThreads(transaction *sqlxTxWrapper, postIds []string) error {
query := s.getQueryBuilder().
Delete("Threads").
Where(
sq.Eq{"PostId": postIds},
)
if _, err := transaction.ExecBuilder(query); err != nil {
return errors.Wrap(err, "failed to delete Threads")
}
if _, err := transaction.Exec("DELETE FROM ThreadMemberships WHERE PostId = ?", postId); err != nil {
query = s.getQueryBuilder().
Delete("ThreadMemberships").
Where(
sq.Eq{"PostId": postIds},
)
if _, err := transaction.ExecBuilder(query); err != nil {
return errors.Wrap(err, "failed to delete ThreadMemberships")
}
return nil
}
func (s *SqlPostStore) permanentDeleteReactions(transaction *sqlxTxWrapper, postIds []string) error {
query := s.getQueryBuilder().
Delete("Reactions").
Where(
sq.Eq{"PostId": postIds},
)
if _, err := transaction.ExecBuilder(query); err != nil {
return errors.Wrap(err, "failed to delete Reactions")
}
return nil
}
// deleteThread marks a thread as deleted at the given time.
func (s *SqlPostStore) deleteThread(transaction *sqlxTxWrapper, postId string, deleteAtTime int64) error {
queryString, args, err := s.getQueryBuilder().

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

@@ -4,6 +4,8 @@
package sqlstore
import (
"time"
sq "github.com/mattermost/squirrel"
"github.com/mattermost/mattermost/server/public/model"
@@ -205,24 +207,91 @@ func (s *SqlReactionStore) DeleteAllWithEmojiName(emojiName string) error {
return nil
}
// DeleteOrphanedRows removes entries from Reactions when a corresponding post no longer exists.
func (s *SqlReactionStore) DeleteOrphanedRows(limit int) (deleted int64, err error) {
// We need the extra level of nesting to deal with MySQL's locking
const query = `
DELETE FROM Reactions WHERE PostId IN (
SELECT * FROM (
SELECT PostId FROM Reactions
LEFT JOIN Posts ON Reactions.PostId = Posts.Id
WHERE Posts.Id IS NULL
LIMIT ?
) AS A
)`
result, err := s.GetMasterX().Exec(query, limit)
func (s *SqlReactionStore) permanentDeleteReactions(userId string, postIds *[]string) error {
txn, err := s.GetMasterX().Beginx()
if err != nil {
return
return err
}
deleted, err = result.RowsAffected()
return
defer finalizeTransactionX(txn, &err)
err = txn.Select(postIds, "SELECT PostId FROM Reactions WHERE UserId = ?", userId)
if err != nil {
return errors.Wrapf(err, "failed to get Reactions with userId=%s", userId)
}
query := s.getQueryBuilder().
Delete("Reactions").
Where(sq.And{
sq.Eq{"PostId": postIds},
sq.Eq{"UserId": userId},
})
_, err = txn.ExecBuilder(query)
if err != nil {
return errors.Wrapf(err, "failed to delete reactions with userId=%s", userId)
}
if err = txn.Commit(); err != nil {
return err
}
return nil
}
func (s SqlReactionStore) PermanentDeleteByUser(userId string) error {
now := model.GetMillis()
postIds := []string{}
err := s.permanentDeleteReactions(userId, &postIds)
if err != nil {
return err
}
transaction, err := s.GetMasterX().Beginx()
if err != nil {
return err
}
defer finalizeTransactionX(transaction, &err)
for _, postId := range postIds {
_, err = transaction.Exec(UpdatePostHasReactionsOnDeleteQuery, now, postId, postId)
if err != nil {
mlog.Warn("Unable to update Post.HasReactions while removing reactions",
mlog.String("post_id", postId),
mlog.Err(err))
}
time.Sleep(10 * time.Millisecond)
}
if err = transaction.Commit(); err != nil {
return err
}
return nil
}
func (s *SqlReactionStore) DeleteOrphanedRowsByIds(r *model.RetentionIdsForDeletion) error {
txn, err := s.GetMasterX().Beginx()
if err != nil {
return err
}
defer finalizeTransactionX(txn, &err)
query := s.getQueryBuilder().
Delete("Reactions").
Where(
sq.Eq{"PostId": r.Ids},
)
_, err = txn.ExecBuilder(query)
if err != nil {
return errors.Wrapf(err, "failed to delete orphaned reactions with RetentionIdsForDeletion Id=%s", r.Id)
}
err = deleteFromRetentionIdsTx(txn, r.Id)
if err != nil {
return err
}
if err = txn.Commit(); err != nil {
return err
}
return nil
}
func (s *SqlReactionStore) PermanentDeleteBatch(endTime int64, limit int64) (int64, error) {

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

@@ -5,6 +5,7 @@ package sqlstore
import (
"database/sql"
"encoding/json"
"fmt"
"strconv"
"strings"
@@ -817,6 +818,92 @@ func (s *SqlRetentionPolicyStore) GetChannelPoliciesCountForUser(userID string)
return count, nil
}
func scanRetentionIdsForDeletion(rows *sql.Rows, isPostgres bool) ([]*model.RetentionIdsForDeletion, error) {
idsForDeletion := []*model.RetentionIdsForDeletion{}
for rows.Next() {
var row model.RetentionIdsForDeletion
if isPostgres {
if err := rows.Scan(
&row.Id, &row.TableName, pq.Array(&row.Ids),
); err != nil {
return nil, errors.Wrap(err, "unable to scan columns")
}
} else {
var ids []byte
if err := rows.Scan(
&row.Id, &row.TableName, &ids,
); err != nil {
return nil, errors.Wrap(err, "unable to scan columns")
}
if err := json.Unmarshal(ids, &row.Ids); err != nil {
return nil, errors.Wrap(err, "failed to unmarshal ids")
}
}
idsForDeletion = append(idsForDeletion, &row)
}
if err := rows.Err(); err != nil {
return nil, errors.Wrap(err, "error while iterating over rows")
}
return idsForDeletion, nil
}
func (s *SqlRetentionPolicyStore) GetIdsForDeletionByTableName(tableName string, limit int) ([]*model.RetentionIdsForDeletion, error) {
query := s.getQueryBuilder().
Select("*").
From("RetentionIdsForDeletion").
Where(
sq.Eq{"TableName": tableName},
).
Limit(uint64(limit))
queryString, args, err := query.ToSql()
if err != nil {
return nil, errors.Wrap(err, "get_ids_for_deletion_tosql")
}
rows, err := s.GetReplicaX().DB.Query(queryString, args...)
if err != nil {
return nil, errors.Wrap(err, "failed to get ids for deletion")
}
defer rows.Close()
isPostgres := s.DriverName() == model.DatabaseDriverPostgres
idsForDeletion, err := scanRetentionIdsForDeletion(rows, isPostgres)
if err != nil {
return nil, errors.Wrap(err, "failed to scan ids for deletion")
}
return idsForDeletion, nil
}
func insertRetentionIdsForDeletion(txn *sqlxTxWrapper, row *model.RetentionIdsForDeletion, s *SqlStore) error {
row.PreSave()
insertBuilder := s.getQueryBuilder().
Insert("RetentionIdsForDeletion").
Columns("Id", "TableName", "Ids")
if s.DriverName() == model.DatabaseDriverPostgres {
insertBuilder = insertBuilder.
Values(row.Id, row.TableName, pq.Array(row.Ids))
} else {
jsonIds, err := json.Marshal(row.Ids)
if err != nil {
return err
}
insertBuilder = insertBuilder.
Values(row.Id, row.TableName, jsonIds)
}
insertQuery, insertArgs, err := insertBuilder.ToSql()
if err != nil {
return err
}
if _, err = txn.Exec(insertQuery, insertArgs...); err != nil {
return err
}
return nil
}
// RetentionPolicyBatchDeletionInfo gives information on how to delete records
// under a retention policy; see `genericPermanentDeleteBatchForRetentionPolicies`.
//
@@ -845,6 +932,7 @@ type RetentionPolicyBatchDeletionInfo struct {
NowMillis int64
GlobalPolicyEndTime int64
Limit int64
StoreDeletedIds bool
}
// genericPermanentDeleteBatchForRetentionPolicies is a helper function for tables
@@ -958,31 +1046,115 @@ func genericRetentionPoliciesDeletion(
if err != nil {
return 0, errors.Wrap(err, r.Table+"_tosql")
}
if s.DriverName() == model.DatabaseDriverPostgres {
primaryKeysStr := "(" + strings.Join(r.PrimaryKeys, ",") + ")"
query = `
DELETE FROM ` + r.Table + ` WHERE ` + primaryKeysStr + ` IN (
` + query + `
)`
} else {
// MySQL does not support the LIMIT clause in a subquery with IN
clauses := make([]string, len(r.PrimaryKeys))
for i, key := range r.PrimaryKeys {
clauses[i] = r.Table + "." + key + " = A." + key
if r.StoreDeletedIds {
txn, err := s.GetMasterX().Beginx()
if err != nil {
return 0, err
}
defer finalizeTransactionX(txn, &err)
if s.DriverName() == model.DatabaseDriverPostgres {
primaryKeysStr := "(" + strings.Join(r.PrimaryKeys, ",") + ")"
query = fmt.Sprintf("DELETE FROM %s WHERE %s IN (%s) RETURNING %s.%s", r.Table, primaryKeysStr, query, r.Table, r.PrimaryKeys[0])
var rows *sql.Rows
rows, err = txn.Query(query, args...)
if err != nil {
return 0, errors.Wrap(err, "failed to delete "+r.Table)
}
defer rows.Close()
ids := []string{}
for rows.Next() {
var id string
if err = rows.Scan(&id); err != nil {
return 0, errors.Wrap(err, "unable to scan from rows")
}
ids = append(ids, id)
}
if err = rows.Err(); err != nil {
return 0, errors.Wrap(err, "failed while iterating over rows")
}
rowsAffected = int64(len(ids))
if len(ids) > 0 {
retentionIdsRow := model.RetentionIdsForDeletion{
TableName: r.Table,
Ids: ids,
}
err = insertRetentionIdsForDeletion(txn, &retentionIdsRow, s)
if err != nil {
return 0, err
}
}
} else {
retentionIdsRow := model.RetentionIdsForDeletion{
TableName: r.Table,
Ids: []string{},
}
// 1. Select rows that will be deleted
if err = txn.Select(&retentionIdsRow.Ids, query, args...); err != nil {
return 0, err
}
if len(retentionIdsRow.Ids) > 0 {
// 2. Insert selected ids into RetentionIdsForDeletion table
err = insertRetentionIdsForDeletion(txn, &retentionIdsRow, s)
if err != nil {
return 0, err
}
query = getDeleteQueriesForMySQL(r, query)
// 3. Delete from Parent table
var result sql.Result
result, err = txn.Exec(query, args...)
if err != nil {
return 0, errors.Wrap(err, "failed to delete "+r.Table)
}
rowsAffected, err = result.RowsAffected()
if err != nil {
return 0, errors.Wrap(err, "failed to get rows affected for "+r.Table)
}
}
}
if err = txn.Commit(); err != nil {
return 0, err
}
} else {
if s.DriverName() == model.DatabaseDriverPostgres {
primaryKeysStr := "(" + strings.Join(r.PrimaryKeys, ",") + ")"
query = fmt.Sprintf("DELETE FROM %s WHERE %s IN (%s)", r.Table, primaryKeysStr, query)
} else {
query = getDeleteQueriesForMySQL(r, query)
}
result, err := s.GetMasterX().Exec(query, args...)
if err != nil {
return 0, errors.Wrap(err, "failed to delete "+r.Table)
}
rowsAffected, err = result.RowsAffected()
if err != nil {
return 0, errors.Wrap(err, "failed to get rows affected for "+r.Table)
}
joinClause := strings.Join(clauses, " AND ")
query = `
DELETE ` + r.Table + ` FROM ` + r.Table + ` INNER JOIN (
` + query + `
) AS A ON ` + joinClause
}
result, err := s.GetMasterX().Exec(query, args...)
if err != nil {
return 0, errors.Wrap(err, "failed to delete "+r.Table)
}
rowsAffected, err = result.RowsAffected()
if err != nil {
return 0, errors.Wrap(err, "failed to get rows affected for "+r.Table)
}
return
}
func getDeleteQueriesForMySQL(r RetentionPolicyBatchDeletionInfo, query string) string {
// MySQL does not support the LIMIT clause in a subquery with IN
clauses := make([]string, len(r.PrimaryKeys))
for i, key := range r.PrimaryKeys {
clauses[i] = r.Table + "." + key + " = A." + key
}
joinClause := strings.Join(clauses, " AND ")
return fmt.Sprintf("DELETE %s FROM %s INNER JOIN (%s) AS A ON %s", r.Table, r.Table, query, joinClause)
}
func deleteFromRetentionIdsTx(txn *sqlxTxWrapper, id string) (err error) {
if _, err := txn.Exec("DELETE FROM RetentionIdsForDeletion WHERE Id = ?", id); err != nil {
return errors.Wrap(err, "Failed to delete from RetentionIdsForDeletion")
}
return nil
}

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

@@ -919,6 +919,7 @@ func (s *SqlThreadStore) PermanentDeleteBatchForRetentionPolicies(now, globalPol
NowMillis: now,
GlobalPolicyEndTime: globalPolicyEndTime,
Limit: limit,
StoreDeletedIds: false,
}, s.SqlStore, cursor)
}
@@ -939,66 +940,32 @@ func (s *SqlThreadStore) PermanentDeleteBatchThreadMembershipsForRetentionPolici
NowMillis: now,
GlobalPolicyEndTime: globalPolicyEndTime,
Limit: limit,
StoreDeletedIds: false,
}, s.SqlStore, cursor)
}
// DeleteOrphanedRows removes orphaned rows from Threads and ThreadMemberships
func (s *SqlThreadStore) DeleteOrphanedRows(limit int) (deleted int64, err error) {
var threadsQuery string
// We need the extra level of nesting to deal with MySQL's locking
if s.DriverName() == model.DatabaseDriverMysql {
// MySQL fails to do a proper antijoin if the selecting column
// and the joining column are different. In that case, doing a subquery
// leads to a faster plan because MySQL materializes the sub-query
// and does a covering index scan on Threads table. More details on the PR with
// this commit.
threadsQuery = `
DELETE FROM Threads WHERE PostId IN (
SELECT * FROM (
SELECT Threads.PostId FROM Threads
WHERE Threads.ChannelId NOT IN (SELECT Id FROM Channels USE INDEX(PRIMARY))
LIMIT ?
) AS A
)`
} else {
threadsQuery = `
DELETE FROM Threads WHERE PostId IN (
SELECT * FROM (
SELECT Threads.PostId FROM Threads
LEFT JOIN Channels ON Threads.ChannelId = Channels.Id
WHERE Channels.Id IS NULL
LIMIT ?
) AS A
)`
}
// We only delete a thread membership if the entire thread no longer exists,
// not if the root post has been deleted
const threadMembershipsQuery = `
DELETE FROM ThreadMemberships WHERE PostId IN (
SELECT * FROM (
SELECT ThreadMemberships.PostId FROM ThreadMemberships
LEFT JOIN Threads ON ThreadMemberships.PostId = Threads.PostId
WHERE Threads.PostId IS NULL
LIMIT ?
) AS A
)`
result, err := s.GetMasterX().Exec(threadsQuery, limit)
DELETE FROM ThreadMemberships WHERE PostId IN (
SELECT * FROM (
SELECT ThreadMemberships.PostId FROM ThreadMemberships
LEFT JOIN Threads ON ThreadMemberships.PostId = Threads.PostId
WHERE Threads.PostId IS NULL
LIMIT ?
) AS A
)`
result, err := s.GetMasterX().Exec(threadMembershipsQuery, limit)
if err != nil {
return
}
rpcDeleted, err := result.RowsAffected()
deleted, err = result.RowsAffected()
if err != nil {
return
}
result, err = s.GetMasterX().Exec(threadMembershipsQuery, limit)
if err != nil {
return
}
rptDeleted, err := result.RowsAffected()
if err != nil {
return
}
deleted = rpcDeleted + rptDeleted
return
}