MM-62751: [Shared Channels] Allow remote users to be discoverable in the create DM/GM modal (#30918)

Этот коммит содержится в:
catalintomai
2025-06-13 16:51:12 +02:00
коммит произвёл GitHub
родитель 476b46d1d7
Коммит c46ed6c681
24 изменённых файлов: 2133 добавлений и 38 удалений

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

@@ -289,3 +289,19 @@ func (rcs *Service) pause() {
rcs.server.Log().Debug("Remote Cluster Service inactive")
}
// SetActive forces the service to be active or inactive
func (rcs *Service) SetActive(active bool) {
rcs.mux.Lock()
defer rcs.mux.Unlock()
if rcs.active == active {
return
}
if active {
rcs.resume()
} else {
rcs.pause()
}
}

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

@@ -24,6 +24,7 @@ const (
TopicSync = "sharedchannel_sync"
TopicChannelInvite = "sharedchannel_invite"
TopicUploadCreate = "sharedchannel_upload"
TopicGlobalUserSync = "sharedchannel_global_user_sync"
MaxRetries = 3
MaxUsersPerSync = 25
NotifyRemoteOfflineThreshold = time.Second * 10
@@ -98,6 +99,7 @@ type Service struct {
syncTopicListenerId string
inviteTopicListenerId string
uploadTopicListenerId string
globalSyncTopicListenerId string
siteURL *url.URL
}
@@ -130,6 +132,7 @@ func (scs *Service) Start() error {
scs.syncTopicListenerId = rcs.AddTopicListener(TopicSync, scs.onReceiveSyncMessage)
scs.inviteTopicListenerId = rcs.AddTopicListener(TopicChannelInvite, scs.onReceiveChannelInvite)
scs.uploadTopicListenerId = rcs.AddTopicListener(TopicUploadCreate, scs.onReceiveUploadCreate)
scs.globalSyncTopicListenerId = rcs.AddTopicListener(TopicGlobalUserSync, scs.onReceiveSyncMessage)
scs.connectionStateListenerId = rcs.AddConnectionStateListener(scs.onConnectionStateChange)
scs.mux.Unlock()
@@ -248,6 +251,9 @@ func (scs *Service) onConnectionStateChange(rc *model.RemoteCluster, online bool
// when a previously offline remote comes back online force a sync.
scs.SendPendingInvitesForRemote(rc)
scs.ForceSyncForRemote(rc)
// Schedule global user sync if feature is enabled
scs.scheduleGlobalUserSync(rc)
}
scs.server.Log().Log(mlog.LvlSharedChannelServiceDebug, "Remote cluster connection status changed",
@@ -278,3 +284,47 @@ func (scs *Service) notifyClientsForSharedChannelUpdate(channel *model.Channel)
messageWs.Add("channel_id", channel.Id)
scs.app.Publish(messageWs)
}
// isGlobalUserSyncEnabled checks if the global user sync feature is enabled
func (scs *Service) isGlobalUserSyncEnabled() bool {
cfg := scs.server.Config()
return cfg.FeatureFlags.EnableSyncAllUsersForRemoteCluster ||
(cfg.ConnectedWorkspacesSettings.SyncUsersOnConnectionOpen != nil && *cfg.ConnectedWorkspacesSettings.SyncUsersOnConnectionOpen)
}
// scheduleGlobalUserSync schedules a task to sync all users with a remote cluster
func (scs *Service) scheduleGlobalUserSync(rc *model.RemoteCluster) {
if !scs.isGlobalUserSyncEnabled() {
return
}
// Schedule the sync task
go func() {
// Create a special sync task with empty channelID
// This empty channelID is a deliberate marker for a global user sync task
task := newSyncTask("", "", rc.RemoteId, nil, nil)
task.schedule = time.Now().Add(NotifyMinimumDelay)
scs.addTask(task)
scs.server.Log().Log(mlog.LvlSharedChannelServiceDebug, "Scheduled global user sync task for remote",
mlog.String("remote", rc.DisplayName),
mlog.String("remoteId", rc.RemoteId),
)
}()
}
// OnReceiveSyncMessageForTesting exposes onReceiveSyncMessage for testing
func (scs *Service) OnReceiveSyncMessageForTesting(msg model.RemoteClusterMsg, rc *model.RemoteCluster, response *remotecluster.Response) error {
return scs.onReceiveSyncMessage(msg, rc, response)
}
// GetUserSyncBatchSizeForTesting returns the configured batch size for user syncing (exported for testing)
func (scs *Service) GetUserSyncBatchSizeForTesting() int {
return scs.getGlobalUserSyncBatchSize()
}
// HandleSyncAllUsersForTesting exposes syncAllUsers for testing
func (scs *Service) HandleSyncAllUsersForTesting(rc *model.RemoteCluster) error {
return scs.syncAllUsers(rc)
}

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

@@ -24,8 +24,8 @@ var (
)
func (scs *Service) onReceiveSyncMessage(msg model.RemoteClusterMsg, rc *model.RemoteCluster, response *remotecluster.Response) error {
if msg.Topic != TopicSync {
return fmt.Errorf("wrong topic, expected `%s`, got `%s`", TopicSync, msg.Topic)
if msg.Topic != TopicSync && msg.Topic != TopicGlobalUserSync {
return fmt.Errorf("wrong topic, expected `%s` or `%s`, got `%s`", TopicSync, TopicGlobalUserSync, msg.Topic)
}
if len(msg.Payload) == 0 {
@@ -47,6 +47,36 @@ func (scs *Service) onReceiveSyncMessage(msg model.RemoteClusterMsg, rc *model.R
return scs.processSyncMessage(request.EmptyContext(scs.server.Log()), &sm, rc, response)
}
func (scs *Service) processGlobalUserSync(c request.CTX, syncMsg *model.SyncMsg, rc *model.RemoteCluster, response *remotecluster.Response) error {
syncResp := model.SyncResponse{
UserErrors: make([]string, 0),
UsersSyncd: make([]string, 0),
}
scs.server.Log().Log(mlog.LvlSharedChannelServiceDebug, "Processing global user sync",
mlog.String("remote", rc.Name),
mlog.Int("user_count", len(syncMsg.Users)),
)
// Process all users in the sync message
for _, user := range syncMsg.Users {
if userSaved, err := scs.upsertSyncUser(c, user, nil, rc); err != nil {
syncResp.UserErrors = append(syncResp.UserErrors, user.Id)
} else {
syncResp.UsersSyncd = append(syncResp.UsersSyncd, userSaved.Id)
if syncResp.UsersLastUpdateAt < user.UpdateAt {
syncResp.UsersLastUpdateAt = user.UpdateAt
}
scs.server.Log().Log(mlog.LvlSharedChannelServiceDebug, "Global user upserted via sync",
mlog.String("remote", rc.Name),
mlog.String("user_id", user.Id),
)
}
}
return response.SetPayload(syncResp)
}
func (scs *Service) processSyncMessage(c request.CTX, syncMsg *model.SyncMsg, rc *model.RemoteCluster, response *remotecluster.Response) error {
var targetChannel *model.Channel
var team *model.Team
@@ -68,6 +98,25 @@ func (scs *Service) processSyncMessage(c request.CTX, syncMsg *model.SyncMsg, rc
mlog.Int("status_count", len(syncMsg.Statuses)),
)
// Check if this is a global user sync message (no channel ID and only users)
if syncMsg.ChannelId == "" {
if len(syncMsg.Posts) != 0 ||
len(syncMsg.Reactions) != 0 ||
len(syncMsg.Statuses) != 0 {
return fmt.Errorf("global user sync message should not contain posts, reactions or statuses")
}
if len(syncMsg.Users) == 0 {
return nil
}
// Check if feature flag is enabled
if !scs.isGlobalUserSyncEnabled() {
return nil
}
return scs.processGlobalUserSync(c, syncMsg, rc, response)
}
// For regular sync messages, we need a specific channel
if targetChannel, err = scs.server.GetStore().Channel().Get(syncMsg.ChannelId, true); err != nil {
// if the channel doesn't exist then none of these sync items are going to work.
return fmt.Errorf("channel not found processing sync message: %w", err)
@@ -239,7 +288,7 @@ func (scs *Service) upsertSyncUser(c request.CTX, user *model.User, channel *mod
// Instead of undoing what succeeded on any failure we simply do all steps each
// time. AddUserToChannel & AddUserToTeamByTeamId do not error if user was already
// added and exit quickly. Not needed for DMs where teamId is empty.
if channel.TeamId != "" {
if channel != nil && channel.TeamId != "" {
// add user to team
if err := scs.app.AddUserToTeamByTeamId(request.EmptyContext(scs.server.Log()), channel.TeamId, userSaved); err != nil {
return nil, fmt.Errorf("error adding sync user to Team: %w", err)

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

@@ -571,14 +571,14 @@ func (scs *Service) shouldUserSync(user *model.User, channelID string, rc *model
if _, err = scs.server.GetStore().SharedChannel().SaveUser(scu); err != nil {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Error adding user to shared channel users",
mlog.String("user_id", user.Id),
mlog.String("channel_id", user.Id),
mlog.String("channel_id", channelID),
mlog.String("remote_id", rc.RemoteId),
mlog.Err(err),
)
} else {
scs.server.Log().Log(mlog.LvlSharedChannelServiceDebug, "Added user to shared channel users",
mlog.String("user_id", user.Id),
mlog.String("channel_id", user.Id),
mlog.String("channel_id", channelID),
mlog.String("remote_id", rc.RemoteId),
)
}

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

@@ -39,8 +39,9 @@ type syncData struct {
statuses []*model.Status
attachments []attachment
resultRepeat bool
resultNextCursor model.GetPostsSinceForSyncCursor
resultRepeat bool
resultNextCursor model.GetPostsSinceForSyncCursor
GlobalUserSyncLastTimestamp int64
}
func newSyncData(task syncTask, rc *model.RemoteCluster, scr *model.SharedChannelRemote) *syncData {
@@ -83,6 +84,12 @@ func (sd *syncData) setDataFromMsg(msg *model.SyncMsg) {
// channels are very active.
// Returning an error forces a retry on the task.
func (scs *Service) syncForRemote(task syncTask, rc *model.RemoteCluster) error {
// Empty channelID indicates a global user sync task
// Normal syncTasks always include a valid channelID
if task.channelID == "" {
return scs.syncAllUsers(rc)
}
rcs := scs.server.GetRemoteClusterService()
if rcs == nil {
return fmt.Errorf("cannot update remote cluster %s for channel id %s; Remote Cluster Service not enabled", rc.Name, task.channelID)
@@ -241,6 +248,7 @@ func (scs *Service) fetchUsersForSync(sd *syncData) error {
return err
}
// Don't sync users back to the remote cluster they originated from
for _, u := range users {
if u.GetRemoteID() != sd.rc.RemoteId {
sd.users[u.Id] = u
@@ -495,7 +503,7 @@ func (scs *Service) filterPostsForSync(sd *syncData) {
continue
}
// don't sync a post back to the remote it came from.
// don't sync a post back to the remote cluster it came from.
if p.GetRemoteID() == sd.rc.RemoteId {
continue
}
@@ -578,6 +586,11 @@ func (scs *Service) sendUserSyncData(sd *syncData) error {
msg.Users = sd.users
err := scs.sendSyncMsgToRemote(msg, sd.rc, func(syncResp model.SyncResponse, errResp error) {
// Only update cursor on successful sync
if errResp == nil && sd.GlobalUserSyncLastTimestamp > 0 {
scs.updateGlobalSyncCursor(sd.rc, sd.GlobalUserSyncLastTimestamp)
}
for _, userID := range syncResp.UsersSyncd {
if err := scs.server.GetStore().SharedChannel().UpdateUserLastSyncAt(userID, sd.task.channelID, sd.rc.RemoteId); err != nil {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Cannot update shared channel user LastSyncAt",
@@ -691,6 +704,204 @@ func (scs *Service) sendProfileImageSyncData(sd *syncData) {
}
}
// shouldUserSyncGlobal determines if a user needs to be synchronized globally.
// Compares user's update timestamp with the remote cluster's LastGlobalUserSyncAt.
func (scs *Service) shouldUserSyncGlobal(user *model.User, rc *model.RemoteCluster) (bool, error) {
// Don't sync users back to the remote cluster they originated from
if user.IsRemote() && user.GetRemoteID() == rc.RemoteId {
return false, nil
}
// Calculate latest update time for this user (profile or picture)
latestUserUpdateTime := user.UpdateAt
if user.LastPictureUpdate > latestUserUpdateTime {
latestUserUpdateTime = user.LastPictureUpdate
}
// For initial sync (LastGlobalUserSyncAt=0), sync all users
// For incremental sync, only sync users updated after the last sync
if rc.LastGlobalUserSyncAt == 0 {
return true, nil
}
return latestUserUpdateTime > rc.LastGlobalUserSyncAt, nil
}
// syncAllUsers synchronizes all local users to a remote cluster.
// This is called when a connection with a remote cluster is established or when handling a global user sync task.
// Uses cursor-based approach with LastGlobalUserSyncAt to resume after interruptions.
func (scs *Service) syncAllUsers(rc *model.RemoteCluster) error {
// Check if feature is enabled
if !scs.server.Config().FeatureFlags.EnableSyncAllUsersForRemoteCluster {
return nil
}
if !rc.IsOnline() {
return fmt.Errorf("remote cluster %s is not online", rc.RemoteId)
}
// Start metrics tracking
metrics := scs.server.GetMetrics()
start := time.Now()
defer func() {
if metrics != nil {
metrics.IncrementSharedChannelsSyncCounter(rc.RemoteId)
metrics.ObserveSharedChannelsSyncCollectionDuration(rc.RemoteId, time.Since(start).Seconds())
}
}()
batchSize := scs.getGlobalUserSyncBatchSize()
// Create sync data with collected users
sd := &syncData{
task: syncTask{remoteID: rc.RemoteId},
rc: rc,
scr: &model.SharedChannelRemote{RemoteId: rc.RemoteId},
users: make(map[string]*model.User),
}
// Collect users to sync
users, latestTimestamp, _, hasMore, err := scs.collectUsersForGlobalSync(rc, batchSize)
if err != nil {
return err
}
// Exit early if no users to sync
if len(users) == 0 {
scs.server.Log().Log(mlog.LvlSharedChannelServiceDebug, "No users to sync for remote cluster",
mlog.String("remote_id", rc.RemoteId))
return nil
}
// Add users to sync data
sd.users = users
sd.GlobalUserSyncLastTimestamp = latestTimestamp
// Send the collected users to remote
if err := scs.sendUserSyncData(sd); err != nil {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Error sending user batch during sync",
mlog.String("remote_id", rc.RemoteId),
mlog.Err(err),
)
return fmt.Errorf("error sending user batch during sync: %w", err)
}
// Schedule next batch if needed
if hasMore {
scs.scheduleNextUserSyncBatch(rc, latestTimestamp, batchSize, len(users))
}
return nil
}
// getGlobalUserSyncBatchSize returns the configured batch size for user syncing
func (scs *Service) getGlobalUserSyncBatchSize() int {
batchSize := MaxUsersPerSync
if scs.server.Config().ConnectedWorkspacesSettings.GlobalUserSyncBatchSize != nil {
configValue := *scs.server.Config().ConnectedWorkspacesSettings.GlobalUserSyncBatchSize
if configValue > 0 && configValue <= 200 {
batchSize = configValue
}
}
return batchSize
}
// collectUsersForGlobalSync fetches users that need to be synced to the remote
func (scs *Service) collectUsersForGlobalSync(rc *model.RemoteCluster, batchSize int) (map[string]*model.User, int64, int, bool, error) {
options := &model.UserGetOptions{
Page: 0,
PerPage: 100, // Database fetch batch size
Active: true,
Sort: "update_at_asc", // Order by UpdateAt ASC to ensure cursor consistency
}
// Only use UpdatedAfter for incremental syncs, not the initial sync
// This ensures users with UpdateAt=0 are included in the first sync
if rc.LastGlobalUserSyncAt > 0 {
options.UpdatedAfter = rc.LastGlobalUserSyncAt
}
users := make(map[string]*model.User)
latestTimestamp := rc.LastGlobalUserSyncAt
totalCount := 0
// Page through database results
for {
batch, err := scs.server.GetStore().User().GetAllProfiles(options)
if err != nil {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Error fetching users for global sync",
mlog.String("remote_id", rc.RemoteId),
mlog.Err(err),
)
return nil, 0, 0, false, err
}
if len(batch) == 0 {
break // No more users to process
}
totalCount += len(batch)
// Process each user in this database page
for _, user := range batch {
// Stop if we've reached batch limit
if len(users) >= batchSize {
return users, latestTimestamp, totalCount, true, nil
}
// Skip users from remotes
if user.IsRemote() {
continue
}
// Check if user needs syncing
needsSync, _ := scs.shouldUserSyncGlobal(user, rc)
if !needsSync {
continue
}
// Add user and update cursor timestamp
users[user.Id] = user
userUpdateTime := max(user.UpdateAt, user.LastPictureUpdate)
if userUpdateTime > latestTimestamp {
latestTimestamp = userUpdateTime
}
}
// Check if we've reached the end of results
if len(batch) < options.PerPage {
break
}
// Move to next page
options.Page++
}
return users, latestTimestamp, totalCount, false, nil
}
// updateGlobalSyncCursor updates the LastGlobalUserSyncAt value for the remote cluster
func (scs *Service) updateGlobalSyncCursor(rc *model.RemoteCluster, newTimestamp int64) {
if err := scs.server.GetStore().RemoteCluster().UpdateLastGlobalUserSyncAt(rc.RemoteId, newTimestamp); err == nil {
rc.LastGlobalUserSyncAt = newTimestamp
} else {
scs.server.Log().Log(mlog.LvlSharedChannelServiceError, "Failed to update global user sync cursor",
mlog.String("remote_id", rc.RemoteId),
mlog.Err(err),
)
}
}
// scheduleNextUserSyncBatch creates a new task for the next batch of user sync
func (scs *Service) scheduleNextUserSyncBatch(rc *model.RemoteCluster, timestamp int64, batchSize, processedCount int) {
// Use timestamp as userID to make each batch task unique
// This prevents task ID collisions between different batches for the same remote
timestampStr := fmt.Sprintf("%d", timestamp)
task := newSyncTask("", timestampStr, rc.RemoteId, nil, nil)
task.schedule = time.Now().Add(NotifyMinimumDelay)
scs.addTask(task)
}
// sendSyncMsgToRemote synchronously sends the sync message to the remote cluster (or plugin).
func (scs *Service) sendSyncMsgToRemote(msg *model.SyncMsg, rc *model.RemoteCluster, f sendSyncMsgResultFunc) error {
rcs := scs.server.GetRemoteClusterService()
@@ -706,7 +917,15 @@ func (scs *Service) sendSyncMsgToRemote(msg *model.SyncMsg, rc *model.RemoteClus
if err != nil {
return err
}
rcMsg := model.NewRemoteClusterMsg(TopicSync, b)
// Use appropriate topic based on message type
topic := TopicSync
if msg.ChannelId == "" && len(msg.Users) > 0 &&
len(msg.Posts) == 0 && len(msg.Reactions) == 0 &&
len(msg.Statuses) == 0 {
topic = TopicGlobalUserSync
}
rcMsg := model.NewRemoteClusterMsg(topic, b)
ctx, cancel := context.WithTimeout(context.Background(), remotecluster.SendTimeout)
defer cancel()