Adds status sync to shared channels (#28020)

* Adds status sync to shared channels

To allow for status to be synced, this changes add a new type of
shared channel internal task. This task has its contents pre-fetched
and stored in the `existingMsg` property, and it is keyed with a user
ID besides a channel ID, so it doesn't conflict with channel-driven
synchronization tasks.

All status synchronizations are triggered from the app layer, so there
is no need of watching for new WebSocket events. Although right now
we're only syncing one user status per message, the changes account
for a list of statuses in case we want to batch them in the future.

The feature is gated by a configuration property and can be disabled
independently of the rest of Shared Channels if it's necessary. It is
backwards compatible as well, and should cause no problems with
servers running older Mattermost versions.

* Adds status sync error management and retry

* Adds DisableSharedChannelsStatusSync to the telemetry report

---------

Co-authored-by: Mattermost Build <build@mattermost.com>
Этот коммит содержится в:
Miguel de la Cruz
2024-09-10 23:39:07 +02:00
коммит произвёл GitHub
родитель c74f830b76
Коммит b898d13e55
15 изменённых файлов: 240 добавлений и 36 удалений

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

@@ -1057,6 +1057,7 @@ type AppIface interface {
SaveAcknowledgementForPost(c request.CTX, postID, userID string) (*model.PostAcknowledgement, *model.AppError) SaveAcknowledgementForPost(c request.CTX, postID, userID string) (*model.PostAcknowledgement, *model.AppError)
SaveAdminNotification(userId string, notifyData *model.NotifyAdminToUpgradeRequest) *model.AppError SaveAdminNotification(userId string, notifyData *model.NotifyAdminToUpgradeRequest) *model.AppError
SaveAdminNotifyData(data *model.NotifyAdminData) (*model.NotifyAdminData, *model.AppError) SaveAdminNotifyData(data *model.NotifyAdminData) (*model.NotifyAdminData, *model.AppError)
SaveAndBroadcastStatus(status *model.Status)
SaveBrandImage(rctx request.CTX, imageData *multipart.FileHeader) *model.AppError SaveBrandImage(rctx request.CTX, imageData *multipart.FileHeader) *model.AppError
SaveComplianceReport(rctx request.CTX, job *model.Compliance) (*model.Compliance, *model.AppError) SaveComplianceReport(rctx request.CTX, job *model.Compliance) (*model.Compliance, *model.AppError)
SaveReactionForPost(c request.CTX, reaction *model.Reaction) (*model.Reaction, *model.AppError) SaveReactionForPost(c request.CTX, reaction *model.Reaction) (*model.Reaction, *model.AppError)

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

@@ -15305,6 +15305,21 @@ func (a *OpenTracingAppLayer) SaveAdminNotifyData(data *model.NotifyAdminData) (
return resultVar0, resultVar1 return resultVar0, resultVar1
} }
func (a *OpenTracingAppLayer) SaveAndBroadcastStatus(status *model.Status) {
origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.SaveAndBroadcastStatus")
a.ctx = newCtx
a.app.Srv().Store().SetContext(newCtx)
defer func() {
a.app.Srv().Store().SetContext(origCtx)
a.ctx = origCtx
}()
defer span.Finish()
a.app.SaveAndBroadcastStatus(status)
}
func (a *OpenTracingAppLayer) SaveBrandImage(rctx request.CTX, imageData *multipart.FileHeader) *model.AppError { func (a *OpenTracingAppLayer) SaveBrandImage(rctx request.CTX, imageData *multipart.FileHeader) *model.AppError {
origCtx := a.ctx origCtx := a.ctx
span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.SaveBrandImage") span, newCtx := tracing.StartSpanWithParentByContext(a.ctx, "app.SaveBrandImage")

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

@@ -14,6 +14,7 @@ type SharedChannelServiceIFace interface {
Start() error Start() error
NotifyChannelChanged(channelId string) NotifyChannelChanged(channelId string)
NotifyUserProfileChanged(userID string) NotifyUserProfileChanged(userID string)
NotifyUserStatusChanged(status *model.Status)
SendChannelInvite(channel *model.Channel, userId string, rc *model.RemoteCluster, options ...sharedchannel.InviteOption) error SendChannelInvite(channel *model.Channel, userId string, rc *model.RemoteCluster, options ...sharedchannel.InviteOption) error
Active() bool Active() bool
InviteRemoteToChannel(channelID, remoteID, userID string, shareIfNotShared bool) error InviteRemoteToChannel(channelID, remoteID, userID string, shareIfNotShared bool) error

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

@@ -279,6 +279,9 @@ func (ps *PlatformService) SetStatusLastActivityAt(userID string, activityAt int
ps.AddStatusCacheSkipClusterSend(status) ps.AddStatusCacheSkipClusterSend(status)
ps.SetStatusAwayIfNeeded(userID, false) ps.SetStatusAwayIfNeeded(userID, false)
if ps.sharedChannelService != nil {
ps.sharedChannelService.NotifyUserStatusChanged(status)
}
} }
func (ps *PlatformService) UpdateLastActivityAtIfNeeded(session model.Session) { func (ps *PlatformService) UpdateLastActivityAtIfNeeded(session model.Session) {
@@ -346,6 +349,9 @@ func (ps *PlatformService) SetStatusOnline(userID string, manual bool) {
mlog.Error("Failed to save status", mlog.String("user_id", userID), mlog.Err(err), mlog.String("user_id", userID)) mlog.Error("Failed to save status", mlog.String("user_id", userID), mlog.Err(err), mlog.String("user_id", userID))
} }
} }
if ps.sharedChannelService != nil {
ps.sharedChannelService.NotifyUserStatusChanged(status)
}
} }
if broadcast { if broadcast {
@@ -366,6 +372,9 @@ func (ps *PlatformService) SetStatusOffline(userID string, manual bool) {
status = &model.Status{UserId: userID, Status: model.StatusOffline, Manual: manual, LastActivityAt: model.GetMillis(), ActiveChannel: ""} status = &model.Status{UserId: userID, Status: model.StatusOffline, Manual: manual, LastActivityAt: model.GetMillis(), ActiveChannel: ""}
ps.SaveAndBroadcastStatus(status) ps.SaveAndBroadcastStatus(status)
if ps.sharedChannelService != nil {
ps.sharedChannelService.NotifyUserStatusChanged(status)
}
} }
func (ps *PlatformService) SetStatusAwayIfNeeded(userID string, manual bool) { func (ps *PlatformService) SetStatusAwayIfNeeded(userID string, manual bool) {
@@ -398,19 +407,22 @@ func (ps *PlatformService) SetStatusAwayIfNeeded(userID string, manual bool) {
status.ActiveChannel = "" status.ActiveChannel = ""
ps.SaveAndBroadcastStatus(status) ps.SaveAndBroadcastStatus(status)
if ps.sharedChannelService != nil {
ps.sharedChannelService.NotifyUserStatusChanged(status)
}
} }
// SetStatusDoNotDisturbTimed takes endtime in unix epoch format in UTC // SetStatusDoNotDisturbTimed takes endtime in unix epoch format in UTC
// and sets status of given userId to dnd which will be restored back after endtime // and sets status of given userId to dnd which will be restored back after endtime
func (ps *PlatformService) SetStatusDoNotDisturbTimed(userId string, endtime int64) { func (ps *PlatformService) SetStatusDoNotDisturbTimed(userID string, endtime int64) {
if !*ps.Config().ServiceSettings.EnableUserStatuses { if !*ps.Config().ServiceSettings.EnableUserStatuses {
return return
} }
status, err := ps.GetStatus(userId) status, err := ps.GetStatus(userID)
if err != nil { if err != nil {
status = &model.Status{UserId: userId, Status: model.StatusOffline, Manual: false, LastActivityAt: 0, ActiveChannel: ""} status = &model.Status{UserId: userID, Status: model.StatusOffline, Manual: false, LastActivityAt: 0, ActiveChannel: ""}
} }
status.PrevStatus = status.Status status.PrevStatus = status.Status
@@ -420,6 +432,9 @@ func (ps *PlatformService) SetStatusDoNotDisturbTimed(userId string, endtime int
status.DNDEndTime = endtime status.DNDEndTime = endtime
ps.SaveAndBroadcastStatus(status) ps.SaveAndBroadcastStatus(status)
if ps.sharedChannelService != nil {
ps.sharedChannelService.NotifyUserStatusChanged(status)
}
} }
func (ps *PlatformService) SetStatusDoNotDisturb(userID string) { func (ps *PlatformService) SetStatusDoNotDisturb(userID string) {
@@ -437,6 +452,9 @@ func (ps *PlatformService) SetStatusDoNotDisturb(userID string) {
status.Manual = true status.Manual = true
ps.SaveAndBroadcastStatus(status) ps.SaveAndBroadcastStatus(status)
if ps.sharedChannelService != nil {
ps.sharedChannelService.NotifyUserStatusChanged(status)
}
} }
func (ps *PlatformService) SetStatusOutOfOffice(userID string) { func (ps *PlatformService) SetStatusOutOfOffice(userID string) {
@@ -454,6 +472,9 @@ func (ps *PlatformService) SetStatusOutOfOffice(userID string) {
status.Manual = true status.Manual = true
ps.SaveAndBroadcastStatus(status) ps.SaveAndBroadcastStatus(status)
if ps.sharedChannelService != nil {
ps.sharedChannelService.NotifyUserStatusChanged(status)
}
} }
func (ps *PlatformService) isUserAway(lastActivityAt int64) bool { func (ps *PlatformService) isUserAway(lastActivityAt int64) bool {

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

@@ -19,7 +19,7 @@ import (
func (a *App) getSharedChannelsService() (SharedChannelServiceIFace, error) { func (a *App) getSharedChannelsService() (SharedChannelServiceIFace, error) {
scService := a.Srv().GetSharedChannelSyncService() scService := a.Srv().GetSharedChannelSyncService()
if scService == nil || !scService.Active() { if scService == nil || !scService.Active() {
return nil, model.NewAppError("InviteRemoteToChannel", "api.command_share.service_disabled", return nil, model.NewAppError("getSharedChannelsService", "api.command_share.service_disabled",
nil, "", http.StatusBadRequest) nil, "", http.StatusBadRequest)
} }
return scService, nil return scService, nil

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

@@ -15,6 +15,7 @@ type SharedChannelServiceIFace interface {
Start() error Start() error
NotifyChannelChanged(channelId string) NotifyChannelChanged(channelId string)
NotifyUserProfileChanged(userID string) NotifyUserProfileChanged(userID string)
NotifyUserStatusChanged(status *model.Status)
SendChannelInvite(channel *model.Channel, userId string, rc *model.RemoteCluster, options ...sharedchannel.InviteOption) error SendChannelInvite(channel *model.Channel, userId string, rc *model.RemoteCluster, options ...sharedchannel.InviteOption) error
Active() bool Active() bool
InviteRemoteToChannel(channelID, remoteID, userID string, shareIfNotShared bool) error InviteRemoteToChannel(channelID, remoteID, userID string, shareIfNotShared bool) error

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

@@ -50,6 +50,10 @@ func (a *App) SetStatusOutOfOffice(userID string) {
a.Srv().Platform().SetStatusOutOfOffice(userID) a.Srv().Platform().SetStatusOutOfOffice(userID)
} }
func (a *App) SaveAndBroadcastStatus(status *model.Status) {
a.Srv().Platform().SaveAndBroadcastStatus(status)
}
func (a *App) GetStatusFromCache(userID string) *model.Status { func (a *App) GetStatusFromCache(userID string) *model.Status {
return a.Srv().Platform().GetStatusFromCache(userID) return a.Srv().Platform().GetStatusFromCache(userID)
} }
@@ -66,9 +70,14 @@ func (a *App) UpdateDNDStatusOfUsers() {
mlog.Warn("Failed to fetch dnd statues from store", mlog.String("err", err.Error())) mlog.Warn("Failed to fetch dnd statues from store", mlog.String("err", err.Error()))
return return
} }
scs, _ := a.getSharedChannelsService()
for i := range statuses { for i := range statuses {
a.Srv().Platform().AddStatusCache(statuses[i]) a.Srv().Platform().AddStatusCache(statuses[i])
a.Srv().Platform().BroadcastStatus(statuses[i]) a.Srv().Platform().BroadcastStatus(statuses[i])
if scs != nil {
scs.NotifyUserStatusChanged(statuses[i])
}
} }
} }

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

@@ -474,6 +474,11 @@ func (_m *MockAppIface) Publish(message *model.WebSocketEvent) {
_m.Called(message) _m.Called(message)
} }
// SaveAndBroadcastStatus provides a mock function with given fields: status
func (_m *MockAppIface) SaveAndBroadcastStatus(status *model.Status) {
_m.Called(status)
}
// SaveReactionForPost provides a mock function with given fields: c, reaction // SaveReactionForPost provides a mock function with given fields: c, reaction
func (_m *MockAppIface) SaveReactionForPost(c request.CTX, reaction *model.Reaction) (*model.Reaction, *model.AppError) { func (_m *MockAppIface) SaveReactionForPost(c request.CTX, reaction *model.Reaction) (*model.Reaction, *model.AppError) {
ret := _m.Called(c, reaction) ret := _m.Called(c, reaction)

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

@@ -62,6 +62,7 @@ type AppIface interface {
DeletePost(c request.CTX, postID, deleteByID string) (*model.Post, *model.AppError) DeletePost(c request.CTX, postID, deleteByID string) (*model.Post, *model.AppError)
SaveReactionForPost(c request.CTX, reaction *model.Reaction) (*model.Reaction, *model.AppError) SaveReactionForPost(c request.CTX, reaction *model.Reaction) (*model.Reaction, *model.AppError)
DeleteReactionForPost(c request.CTX, reaction *model.Reaction) *model.AppError DeleteReactionForPost(c request.CTX, reaction *model.Reaction) *model.AppError
SaveAndBroadcastStatus(status *model.Status)
PatchChannelModerationsForChannel(c request.CTX, channel *model.Channel, channelModerationsPatch []*model.ChannelModerationPatch) ([]*model.ChannelModeration, *model.AppError) PatchChannelModerationsForChannel(c request.CTX, channel *model.Channel, channelModerationsPatch []*model.ChannelModerationPatch) ([]*model.ChannelModeration, *model.AppError)
CreateUploadSession(c request.CTX, us *model.UploadSession) (*model.UploadSession, *model.AppError) CreateUploadSession(c request.CTX, us *model.UploadSession) (*model.UploadSession, *model.AppError)
FileReader(path string) (filestore.ReadCloseSeeker, *model.AppError) FileReader(path string) (filestore.ReadCloseSeeker, *model.AppError)

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

@@ -65,6 +65,7 @@ func (scs *Service) processSyncMessage(c request.CTX, syncMsg *model.SyncMsg, rc
mlog.Int("user_count", len(syncMsg.Users)), mlog.Int("user_count", len(syncMsg.Users)),
mlog.Int("post_count", len(syncMsg.Posts)), mlog.Int("post_count", len(syncMsg.Posts)),
mlog.Int("reaction_count", len(syncMsg.Reactions)), mlog.Int("reaction_count", len(syncMsg.Reactions)),
mlog.Int("status_count", len(syncMsg.Statuses)),
) )
if targetChannel, err = scs.server.GetStore().Channel().Get(syncMsg.ChannelId, true); err != nil { if targetChannel, err = scs.server.GetStore().Channel().Get(syncMsg.ChannelId, true); err != nil {
@@ -175,6 +176,10 @@ func (scs *Service) processSyncMessage(c request.CTX, syncMsg *model.SyncMsg, rc
} }
} }
for _, status := range syncMsg.Statuses {
scs.app.SaveAndBroadcastStatus(status)
}
response.SetPayload(syncResp) response.SetPayload(syncResp)
return nil return nil

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

@@ -18,23 +18,31 @@ import (
type syncTask struct { type syncTask struct {
id string id string
channelID string channelID string
userID string
remoteID string remoteID string
AddedAt time.Time AddedAt time.Time
// existingMsg is used to add information to the task on creation
// instead of waiting until the task is processed to fetch it. If
// a new task with the same ID is scheduled, its existingMsg will
// replace the previous one
existingMsg *model.SyncMsg
retryCount int retryCount int
retryMsg *model.SyncMsg retryMsg *model.SyncMsg
schedule time.Time schedule time.Time
} }
func newSyncTask(channelID string, remoteID string, retryMsg *model.SyncMsg) syncTask { func newSyncTask(channelID, userID string, remoteID string, existingMsg, retryMsg *model.SyncMsg) syncTask {
var retryID string var retryID string
if retryMsg != nil { if retryMsg != nil {
retryID = retryMsg.Id retryID = retryMsg.Id
} }
return syncTask{ return syncTask{
id: channelID + remoteID + retryID, // combination of ids to avoid duplicates id: channelID + userID + remoteID + retryID, // combination of ids to avoid duplicates
channelID: channelID, channelID: channelID,
userID: userID,
remoteID: remoteID, // empty means update all remote clusters remoteID: remoteID, // empty means update all remote clusters
existingMsg: existingMsg,
retryMsg: retryMsg, retryMsg: retryMsg,
schedule: time.Now(), schedule: time.Now(),
} }
@@ -53,13 +61,13 @@ func (scs *Service) NotifyChannelChanged(channelID string) {
return return
} }
task := newSyncTask(channelID, "", nil) task := newSyncTask(channelID, "", "", nil, nil)
task.schedule = time.Now().Add(NotifyMinimumDelay) task.schedule = time.Now().Add(NotifyMinimumDelay)
scs.addTask(task) scs.addTask(task)
} }
// NotifyUserProfileChanged is called to indicate that a user belonging to at least one // NotifyUserProfileChanged is called to indicate that a user has modified their user
// shared channel has modified their user profile (name, username, email, custom status, profile image) // profile (name, username, email, custom status, profile image)
func (scs *Service) NotifyUserProfileChanged(userID string) { func (scs *Service) NotifyUserProfileChanged(userID string) {
if rcs := scs.server.GetRemoteClusterService(); rcs == nil { if rcs := scs.server.GetRemoteClusterService(); rcs == nil {
return return
@@ -80,15 +88,63 @@ func (scs *Service) NotifyUserProfileChanged(userID string) {
notified := make(map[string]struct{}) notified := make(map[string]struct{})
for _, user := range scusers { for _, user := range scusers {
// update every channel + remote combination they belong to. // update every user + remote combination they belong to.
// Redundant updates (ie. to same remote for multiple channels) will be // Redundant updates (ie. to same remote for multiple channels) will be
// filtered out. // filtered out.
combo := user.ChannelId + user.RemoteId
combo := user.UserId + user.RemoteId
if _, ok := notified[combo]; ok { if _, ok := notified[combo]; ok {
continue continue
} }
notified[combo] = struct{}{} notified[combo] = struct{}{}
task := newSyncTask(user.ChannelId, user.RemoteId, nil) task := newSyncTask(user.ChannelId, "", user.RemoteId, nil, nil)
task.schedule = time.Now().Add(NotifyMinimumDelay)
scs.addTask(task)
}
}
// NotifyUserStatusChanged is called to indicate that a user has modified their status
func (scs *Service) NotifyUserStatusChanged(status *model.Status) {
if rcs := scs.server.GetRemoteClusterService(); rcs == nil {
return
}
if *scs.server.Config().ConnectedWorkspacesSettings.DisableSharedChannelsStatusSync {
return
}
if status.UserId == "" {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Received invalid status for sync",
mlog.String("userID", status.UserId),
)
return
}
scusers, err := scs.server.GetStore().SharedChannel().GetUsersForUser(status.UserId)
if err != nil {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Failed to fetch shared channel users",
mlog.String("userID", status.UserId),
mlog.Err(err),
)
return
}
if len(scusers) == 0 {
return
}
existingMsg := &model.SyncMsg{Statuses: []*model.Status{status}}
notified := make(map[string]struct{})
for _, user := range scusers {
// update every user + remote combination they belong to.
// Redundant updates (ie. to same remote for multiple channels) will be
// filtered out.
combo := user.UserId + user.RemoteId
if _, ok := notified[combo]; ok {
continue
}
notified[combo] = struct{}{}
task := newSyncTask(user.ChannelId, user.UserId, user.RemoteId, existingMsg, nil)
task.schedule = time.Now().Add(NotifyMinimumDelay) task.schedule = time.Now().Add(NotifyMinimumDelay)
scs.addTask(task) scs.addTask(task)
} }
@@ -115,7 +171,7 @@ func (scs *Service) ForceSyncForRemote(rc *model.RemoteCluster) {
} }
for _, scr := range scrs { for _, scr := range scrs {
task := newSyncTask(scr.ChannelId, rc.RemoteId, nil) task := newSyncTask(scr.ChannelId, "", rc.RemoteId, nil, nil)
task.schedule = time.Now().Add(NotifyMinimumDelay) task.schedule = time.Now().Add(NotifyMinimumDelay)
scs.addTask(task) scs.addTask(task)
} }
@@ -125,7 +181,12 @@ func (scs *Service) ForceSyncForRemote(rc *model.RemoteCluster) {
func (scs *Service) addTask(task syncTask) { func (scs *Service) addTask(task syncTask) {
task.AddedAt = time.Now() task.AddedAt = time.Now()
scs.mux.Lock() scs.mux.Lock()
if _, ok := scs.tasks[task.id]; !ok { if originalTask, ok := scs.tasks[task.id]; ok {
// if the task was already scheduled, we only update the
// existingMsg in case there is new information
originalTask.existingMsg = task.existingMsg
scs.tasks[task.id] = originalTask
} else {
scs.tasks[task.id] = task scs.tasks[task.id] = task
} }
scs.mux.Unlock() scs.mux.Unlock()
@@ -327,6 +388,7 @@ func (scs *Service) handlePostError(postId string, task syncTask, rc *model.Remo
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "error fetching post for sync retry", scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "error fetching post for sync retry",
mlog.String("remote", rc.DisplayName), mlog.String("remote", rc.DisplayName),
mlog.String("post_id", postId), mlog.String("post_id", postId),
mlog.Err(err),
) )
return return
} }
@@ -334,7 +396,38 @@ func (scs *Service) handlePostError(postId string, task syncTask, rc *model.Remo
syncMsg := model.NewSyncMsg(task.channelID) syncMsg := model.NewSyncMsg(task.channelID)
syncMsg.Posts = []*model.Post{post} syncMsg.Posts = []*model.Post{post}
scs.addTask(newSyncTask(task.channelID, task.remoteID, syncMsg)) scs.addTask(newSyncTask(task.channelID, task.userID, task.remoteID, nil, syncMsg))
}
func (scs *Service) handleStatusError(userId string, task syncTask, rc *model.RemoteCluster) {
if task.retryMsg != nil && len(task.retryMsg.Statuses) == 1 && task.retryMsg.Statuses[0].UserId == userId {
// this was a retry for specific status that failed previously. Try again if within MaxRetries.
if task.incRetry() {
scs.addTask(task)
} else {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "error syncing status",
mlog.String("remote", rc.DisplayName),
mlog.String("user_id", userId),
)
}
return
}
// this status failed as part of a group of statuses. Retry as an individual status.
status, err := scs.server.GetStore().Status().Get(userId)
if err != nil {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "error fetching status for sync retry",
mlog.String("remote", rc.DisplayName),
mlog.String("user_id", userId),
mlog.Err(err),
)
return
}
syncMsg := model.NewSyncMsg(task.channelID)
syncMsg.Statuses = []*model.Status{status}
scs.addTask(newSyncTask(task.channelID, task.userID, task.remoteID, nil, syncMsg))
} }
// notifyRemoteOffline creates an ephemeral post to the author for any posts created recently to remotes // notifyRemoteOffline creates an ephemeral post to the author for any posts created recently to remotes

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

@@ -35,6 +35,7 @@ type syncData struct {
profileImages map[string]*model.User profileImages map[string]*model.User
posts []*model.Post posts []*model.Post
reactions []*model.Reaction reactions []*model.Reaction
statuses []*model.Status
attachments []attachment attachments []attachment
resultRepeat bool resultRepeat bool
@@ -68,6 +69,13 @@ func (sd *syncData) isCursorChanged() bool {
sd.scr.LastPostUpdateAt != sd.resultNextCursor.LastPostUpdateAt || sd.scr.LastPostUpdateID != sd.resultNextCursor.LastPostUpdateID sd.scr.LastPostUpdateAt != sd.resultNextCursor.LastPostUpdateAt || sd.scr.LastPostUpdateID != sd.resultNextCursor.LastPostUpdateID
} }
func (sd *syncData) setDataFromMsg(msg *model.SyncMsg) {
sd.users = msg.Users
sd.posts = msg.Posts
sd.reactions = msg.Reactions
sd.statuses = msg.Statuses
}
// syncForRemote updates a remote cluster with any new posts/reactions for a specific // syncForRemote updates a remote cluster with any new posts/reactions for a specific
// channel. If many changes are found, only the oldest X changes are sent and the channel // channel. If many changes are found, only the oldest X changes are sent and the channel
// is re-added to the task map. This ensures no channels are starved for updates even if some // is re-added to the task map. This ensures no channels are starved for updates even if some
@@ -115,21 +123,30 @@ func (scs *Service) syncForRemote(task syncTask, rc *model.RemoteCluster) error
return err return err
} }
sd := newSyncData(task, rc, scr)
// if this is retrying a failed msg, just send it again. // if this is retrying a failed msg, just send it again.
if task.retryMsg != nil { if task.retryMsg != nil {
sd := newSyncData(task, rc, scr) sd.setDataFromMsg(task.retryMsg)
sd.users = task.retryMsg.Users
sd.posts = task.retryMsg.Posts
sd.reactions = task.retryMsg.Reactions
return scs.sendSyncData(sd) return scs.sendSyncData(sd)
} }
sd := newSyncData(task, rc, scr) // if this has an already existing msg, just send it right away
if task.existingMsg != nil {
sd.setDataFromMsg(task.existingMsg)
return scs.sendSyncData(sd)
}
// if we don't have a channelID at this point, we cannot fetch new
// data from the database
if task.channelID == "" {
return fmt.Errorf("task doesn't have prefetched data nor a channel ID set")
}
// schedule another sync if the repeat flag is set at some point. // schedule another sync if the repeat flag is set at some point.
defer func(rpt *bool) { defer func(rpt *bool) {
if *rpt { if *rpt {
scs.addTask(newSyncTask(task.channelID, task.remoteID, nil)) scs.addTask(newSyncTask(task.channelID, task.userID, task.remoteID, nil, nil))
} }
}(&sd.resultRepeat) }(&sd.resultRepeat)
@@ -516,6 +533,13 @@ func (scs *Service) sendSyncData(sd *syncData) error {
} }
} }
// send statuses
if len(sd.statuses) != 0 {
if err := scs.sendStatusSyncData(sd); err != nil {
merr.Append(fmt.Errorf("cannot send status sync data: %w", err))
}
}
// send user profile images // send user profile images
if len(sd.profileImages) != 0 { if len(sd.profileImages) != 0 {
scs.sendProfileImageSyncData(sd) scs.sendProfileImageSyncData(sd)
@@ -624,6 +648,25 @@ func (scs *Service) sendReactionSyncData(sd *syncData) error {
}) })
} }
// sendStatusSyncData sends the collected status updates to the remote cluster.
func (scs *Service) sendStatusSyncData(sd *syncData) error {
msg := model.NewSyncMsg(sd.task.channelID)
msg.Statuses = sd.statuses
return scs.sendSyncMsgToRemote(msg, sd.rc, func(syncResp model.SyncResponse, errResp error) {
if len(syncResp.StatusErrors) != 0 {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Response indicates error from status(es) sync",
mlog.String("remote_id", sd.rc.RemoteId),
mlog.Array("user_ids", syncResp.StatusErrors),
)
for _, userID := range syncResp.StatusErrors {
scs.handleStatusError(userID, sd.task, sd.rc)
}
}
})
}
// sendProfileImageSyncData sends the collected user profile image updates to the remote cluster. // sendProfileImageSyncData sends the collected user profile image updates to the remote cluster.
func (scs *Service) sendProfileImageSyncData(sd *syncData) { func (scs *Service) sendProfileImageSyncData(sd *syncData) {
for _, user := range sd.profileImages { for _, user := range sd.profileImages {

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

@@ -893,6 +893,7 @@ func (ts *TelemetryService) trackConfig() {
ts.SendTelemetry(TrackConfigConnectedWorkspaces, map[string]any{ ts.SendTelemetry(TrackConfigConnectedWorkspaces, map[string]any{
"enable_shared_channels": *cfg.ConnectedWorkspacesSettings.EnableSharedChannels, "enable_shared_channels": *cfg.ConnectedWorkspacesSettings.EnableSharedChannels,
"enable_remote_cluster_service": *cfg.ConnectedWorkspacesSettings.EnableRemoteClusterService && cfg.FeatureFlags.EnableRemoteClusterService, "enable_remote_cluster_service": *cfg.ConnectedWorkspacesSettings.EnableRemoteClusterService && cfg.FeatureFlags.EnableRemoteClusterService,
"disable_shared_channels_status_sync": *cfg.ConnectedWorkspacesSettings.DisableSharedChannelsStatusSync,
}) })
// Convert feature flags to map[string]any for sending // Convert feature flags to map[string]any for sending

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

@@ -3245,6 +3245,7 @@ func (w *WranglerSettings) IsValid() *AppError {
type ConnectedWorkspacesSettings struct { type ConnectedWorkspacesSettings struct {
EnableSharedChannels *bool EnableSharedChannels *bool
EnableRemoteClusterService *bool EnableRemoteClusterService *bool
DisableSharedChannelsStatusSync *bool
} }
func (c *ConnectedWorkspacesSettings) SetDefaults(isUpdate bool, e ExperimentalSettings) { func (c *ConnectedWorkspacesSettings) SetDefaults(isUpdate bool, e ExperimentalSettings) {
@@ -3263,6 +3264,10 @@ func (c *ConnectedWorkspacesSettings) SetDefaults(isUpdate bool, e ExperimentalS
c.EnableRemoteClusterService = NewPointer(false) c.EnableRemoteClusterService = NewPointer(false)
} }
} }
if c.DisableSharedChannelsStatusSync == nil {
c.DisableSharedChannelsStatusSync = NewPointer(false)
}
} }
type GlobalRelayMessageExportSettings struct { type GlobalRelayMessageExportSettings struct {

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

@@ -275,6 +275,7 @@ type SyncMsg struct {
Users map[string]*User `json:"users,omitempty"` Users map[string]*User `json:"users,omitempty"`
Posts []*Post `json:"posts,omitempty"` Posts []*Post `json:"posts,omitempty"`
Reactions []*Reaction `json:"reactions,omitempty"` Reactions []*Reaction `json:"reactions,omitempty"`
Statuses []*Status `json:"statuses,omitempty"`
} }
func NewSyncMsg(channelID string) *SyncMsg { func NewSyncMsg(channelID string) *SyncMsg {
@@ -311,6 +312,8 @@ type SyncResponse struct {
ReactionsLastUpdateAt int64 `json:"reactions_last_update_at"` ReactionsLastUpdateAt int64 `json:"reactions_last_update_at"`
ReactionErrors []string `json:"reaction_errors"` ReactionErrors []string `json:"reaction_errors"`
StatusErrors []string `json:"status_errors"` // user IDs for which the status sync failed
} }
// RegisterPluginOpts is passed by plugins to the `RegisterPluginForSharedChannels` plugin API // RegisterPluginOpts is passed by plugins to the `RegisterPluginForSharedChannels` plugin API