MM-57326: [Shared Channels] Message priority, acknowledgement and persistent notifications need to be synced (#30736)

Этот коммит содержится в:
catalintomai
2025-06-16 02:30:21 +02:00
коммит произвёл GitHub
родитель fa1c77d9b0
Коммит 85391de22a
27 изменённых файлов: 2667 добавлений и 101 удалений

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

@@ -32,12 +32,13 @@ type syncData struct {
rc *model.RemoteCluster
scr *model.SharedChannelRemote
users map[string]*model.User
profileImages map[string]*model.User
posts []*model.Post
reactions []*model.Reaction
statuses []*model.Status
attachments []attachment
users map[string]*model.User
profileImages map[string]*model.User
posts []*model.Post
reactions []*model.Reaction
acknowledgements []*model.PostAcknowledgement
statuses []*model.Status
attachments []attachment
resultRepeat bool
resultNextCursor model.GetPostsSinceForSyncCursor
@@ -59,7 +60,7 @@ func newSyncData(task syncTask, rc *model.RemoteCluster, scr *model.SharedChanne
}
func (sd *syncData) isEmpty() bool {
return len(sd.users) == 0 && len(sd.profileImages) == 0 && len(sd.posts) == 0 && len(sd.reactions) == 0 && len(sd.attachments) == 0
return len(sd.users) == 0 && len(sd.profileImages) == 0 && len(sd.posts) == 0 && len(sd.reactions) == 0 && len(sd.acknowledgements) == 0 && len(sd.attachments) == 0
}
func (sd *syncData) isCursorChanged() bool {
@@ -75,6 +76,7 @@ func (sd *syncData) setDataFromMsg(msg *model.SyncMsg) {
sd.users = msg.Users
sd.posts = msg.Posts
sd.reactions = msg.Reactions
sd.acknowledgements = msg.Acknowledgements
sd.statuses = msg.Statuses
}
@@ -185,6 +187,11 @@ func (scs *Service) syncForRemote(task syncTask, rc *model.RemoteCluster) error
return fmt.Errorf("cannot fetch reactions for sync %v: %w", sd, err)
}
// fetch acknowledgements for posts
if err := scs.fetchAcknowledgementsForSync(sd); err != nil {
return fmt.Errorf("cannot fetch acknowledgements for sync %v: %w", sd, err)
}
// fetch users associated with posts & reactions
if err := scs.fetchPostUsersForSync(sd); err != nil {
return fmt.Errorf("cannot fetch post users for sync %v: %w", sd, err)
@@ -218,6 +225,7 @@ func (scs *Service) syncForRemote(task syncTask, rc *model.RemoteCluster) error
mlog.Int("images", len(sd.profileImages)),
mlog.Int("posts", len(sd.posts)),
mlog.Int("reactions", len(sd.reactions)),
mlog.Int("acknowledgements", len(sd.acknowledgements)),
mlog.Int("attachments", len(sd.attachments)),
)
@@ -319,6 +327,13 @@ func (scs *Service) fetchPostsForSync(sd *syncData) error {
sd.posts = appendPosts(sd.posts, posts, scs.server.GetStore().Post(), cursor.LastPostUpdateAt, scs.server.Log())
}
// Populate metadata for all posts before syncing
for i, post := range sd.posts {
if post != nil {
sd.posts[i] = scs.app.PreparePostForClient(request.EmptyContext(scs.server.Log()), post, false, false, true)
}
}
sd.resultNextCursor = nextCursor
sd.resultRepeat = count >= maxPostsPerSync
@@ -377,6 +392,28 @@ func (scs *Service) fetchReactionsForSync(sd *syncData) error {
return merr.ErrorOrNil()
}
// fetchAcknowledgementsForSync populates the sync data with any new acknowledgements since the last sync.
func (scs *Service) fetchAcknowledgementsForSync(sd *syncData) error {
start := time.Now()
defer func() {
if metrics := scs.server.GetMetrics(); metrics != nil {
metrics.ObserveSharedChannelsSyncCollectionStepDuration(sd.rc.RemoteId, "Acknowledgements", time.Since(start).Seconds())
}
}()
merr := merror.New()
for _, post := range sd.posts {
// any acknowledgements originating from the remote cluster are filtered out
acknowledgements, err := scs.server.GetStore().PostAcknowledgement().GetForPostSince(post.Id, sd.scr.LastPostUpdateAt, sd.rc.RemoteId, true)
if err != nil {
merr.Append(fmt.Errorf("could not get acknowledgements for post %s: %w", post.Id, err))
continue
}
sd.acknowledgements = append(sd.acknowledgements, acknowledgements...)
}
return merr.ErrorOrNil()
}
// fetchPostUsersForSync populates the sync data with all users associated with posts.
func (scs *Service) fetchPostUsersForSync(sd *syncData) error {
start := time.Now()
@@ -402,6 +439,10 @@ func (scs *Service) fetchPostUsersForSync(sd *syncData) error {
userIDs[reaction.UserId] = p2mm{}
}
for _, acknowledgement := range sd.acknowledgements {
userIDs[acknowledgement.UserId] = p2mm{}
}
for _, post := range sd.posts {
// add author
userIDs[post.UserId] = p2mm{}
@@ -495,7 +536,9 @@ func (scs *Service) filterPostsForSync(sd *syncData) {
// - new posts (EditAt == 0)
// - edited posts (EditAt >= LastPostUpdateAt)
// - deleted posts (DeleteAt > 0)
if p.EditAt > 0 && p.EditAt < sd.scr.LastPostUpdateAt && p.DeleteAt == 0 {
// - posts with metadata changes (acknowledgements/priority)
hasMetadataChanges := p.Metadata != nil && (p.Metadata.Acknowledgements != nil || p.Metadata.Priority != nil)
if p.EditAt > 0 && p.EditAt < sd.scr.LastPostUpdateAt && p.DeleteAt == 0 && !hasMetadataChanges {
continue
}
@@ -514,6 +557,7 @@ func (scs *Service) filterPostsForSync(sd *syncData) {
filtered = append(filtered, p)
}
sd.posts = filtered
}
@@ -552,6 +596,13 @@ func (scs *Service) sendSyncData(sd *syncData) error {
scs.updateCursorForRemote(sd.scr.Id, sd.rc, sd.resultNextCursor)
}
// send acknowledgements
if len(sd.acknowledgements) != 0 {
if err := scs.sendAcknowledgementSyncData(sd); err != nil {
merr.Append(fmt.Errorf("cannot send acknowledgement sync data: %w", err))
}
}
// send reactions
if len(sd.reactions) != 0 {
if err := scs.sendReactionSyncData(sd); err != nil {
@@ -679,6 +730,29 @@ func (scs *Service) sendReactionSyncData(sd *syncData) error {
})
}
// sendAcknowledgementSyncData sends the collected acknowledgement updates to the remote cluster.
func (scs *Service) sendAcknowledgementSyncData(sd *syncData) error {
start := time.Now()
defer func() {
if metrics := scs.server.GetMetrics(); metrics != nil {
metrics.ObserveSharedChannelsSyncSendStepDuration(sd.rc.RemoteId, "Acknowledgements", time.Since(start).Seconds())
}
}()
msg := model.NewSyncMsg(sd.task.channelID)
msg.Acknowledgements = sd.acknowledgements
return scs.sendSyncMsgToRemote(msg, sd.rc, func(syncResp model.SyncResponse, errResp error) {
if len(syncResp.AcknowledgementErrors) != 0 {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Response indicates error for acknowledgement(s) sync",
mlog.String("channel_id", sd.task.channelID),
mlog.String("remote_id", sd.rc.RemoteId),
mlog.Array("acknowledgement_posts", syncResp.AcknowledgementErrors),
)
}
})
}
// sendStatusSyncData sends the collected status updates to the remote cluster.
func (scs *Service) sendStatusSyncData(sd *syncData) error {
msg := model.NewSyncMsg(sd.task.channelID)