MM-57786 Fix Shared Channels plugin api (#26753)
* always ping on plugin registration; SharedChannel.IsValid allow no team for GM * wait for services to start before ping * ping plugin remotes synchronously on startup * remove the waitForInterClusterServices stuff * don't set remoteid when inviting remote to channel * Update server/public/model/remote_cluster_test.go Co-authored-by: Ibrahim Serdar Acikgoz <serdaracikgoz86@gmail.com> * address review comments --------- Co-authored-by: Mattermost Build <build@mattermost.com> Co-authored-by: Ibrahim Serdar Acikgoz <serdaracikgoz86@gmail.com>
Этот коммит содержится в:
@@ -73,7 +73,7 @@ type mockApp struct {
|
||||
pingCounts map[string]int
|
||||
}
|
||||
|
||||
func newMockApp(t *testing.T, offlinePluginIDs []string) *mockApp {
|
||||
func newMockApp(_ *testing.T, offlinePluginIDs []string) *mockApp {
|
||||
return &mockApp{
|
||||
offlinePluginIDs: offlinePluginIDs,
|
||||
pingCounts: make(map[string]int),
|
||||
|
||||
@@ -34,6 +34,23 @@ func (rcs *Service) PingNow(rc *model.RemoteCluster) {
|
||||
}
|
||||
}
|
||||
|
||||
// pingAllNow emits a ping to all remotes immediately without waiting for next ping loop.
|
||||
func (rcs *Service) pingAllNow(filter model.RemoteClusterQueryFilter) {
|
||||
// get all remotes, including any previously offline.
|
||||
remotes, err := rcs.server.GetStore().RemoteCluster().GetAll(filter)
|
||||
if err != nil {
|
||||
rcs.server.Log().Log(mlog.LvlRemoteClusterServiceError, "Ping all remote clusters failed (could not get list of remotes)", mlog.Err(err))
|
||||
return
|
||||
}
|
||||
|
||||
for _, rc := range remotes {
|
||||
// filter out unconfirmed invites so we don't ping them without permission
|
||||
if rc.IsConfirmed() {
|
||||
rcs.PingNow(rc)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// pingLoop periodically sends a ping to all remote clusters.
|
||||
func (rcs *Service) pingLoop(done <-chan struct{}) {
|
||||
pingChan := make(chan *model.RemoteCluster, MaxConcurrentSends*2)
|
||||
@@ -53,24 +70,8 @@ func (rcs *Service) pingGenerator(pingChan chan *model.RemoteCluster, done <-cha
|
||||
pingFreq := rcs.GetPingFreq()
|
||||
start := time.Now()
|
||||
|
||||
// get all remotes, including any previously offline.
|
||||
remotes, err := rcs.server.GetStore().RemoteCluster().GetAll(model.RemoteClusterQueryFilter{})
|
||||
if err != nil {
|
||||
rcs.server.Log().Log(mlog.LvlRemoteClusterServiceError, "Ping remote cluster failed (could not get list of remotes)", mlog.Err(err))
|
||||
select {
|
||||
case <-time.After(pingFreq):
|
||||
continue
|
||||
case <-done:
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
for _, rc := range remotes {
|
||||
// filter out unconfirmed invites so we don't ping them without permission
|
||||
if rc.IsConfirmed() {
|
||||
pingChan <- rc
|
||||
}
|
||||
}
|
||||
// ping all remotes, including any previously offline.
|
||||
rcs.pingAllNow(model.RemoteClusterQueryFilter{})
|
||||
|
||||
// try to maintain frequency
|
||||
elapsed := time.Since(start)
|
||||
|
||||
@@ -261,6 +261,10 @@ func (rcs *Service) resume() {
|
||||
rcs.done = make(chan struct{})
|
||||
|
||||
if !disablePing {
|
||||
// first ping all the plugin remotes immediately, synchronously.
|
||||
rcs.pingAllNow(model.RemoteClusterQueryFilter{OnlyPlugins: true})
|
||||
|
||||
// start the async ping loop
|
||||
rcs.pingLoop(rcs.done)
|
||||
}
|
||||
|
||||
|
||||
@@ -72,6 +72,37 @@ func (scs *Service) SendChannelInvite(channel *model.Channel, userId string, rc
|
||||
|
||||
msg := model.NewRemoteClusterMsg(TopicChannelInvite, json)
|
||||
|
||||
// onInvite is called after invite is sent, whether to a remote cluster or plugin.
|
||||
onInvite := func(_ model.RemoteClusterMsg, rc *model.RemoteCluster, resp *remotecluster.Response, err error) {
|
||||
if err != nil || !resp.IsSuccess() {
|
||||
scs.sendEphemeralPost(channel.Id, userId, fmt.Sprintf("Error sending channel invite for %s: %s", rc.DisplayName, combineErrors(err, resp.Err)))
|
||||
return
|
||||
}
|
||||
|
||||
scr := &model.SharedChannelRemote{
|
||||
ChannelId: sc.ChannelId,
|
||||
CreatorId: userId,
|
||||
RemoteId: rc.RemoteId,
|
||||
IsInviteAccepted: true,
|
||||
IsInviteConfirmed: true,
|
||||
LastPostCreateAt: model.GetMillis(),
|
||||
LastPostUpdateAt: model.GetMillis(),
|
||||
}
|
||||
if _, err = scs.server.GetStore().SharedChannel().SaveRemote(scr); err != nil {
|
||||
scs.sendEphemeralPost(channel.Id, userId, fmt.Sprintf("Error confirming channel invite for %s: %v", rc.DisplayName, err))
|
||||
return
|
||||
}
|
||||
scs.NotifyChannelChanged(sc.ChannelId)
|
||||
scs.sendEphemeralPost(channel.Id, userId, fmt.Sprintf("`%s` has been added to channel.", rc.DisplayName))
|
||||
}
|
||||
|
||||
if rc.IsPlugin() {
|
||||
// for now plugins are considered fully invited automatically
|
||||
// TODO: MM-57537 create plugin hook that passes invitation to plugins if BitflagOptionAutoInvited is not set
|
||||
onInvite(msg, rc, &remotecluster.Response{Status: remotecluster.ResponseStatusOK}, nil)
|
||||
return nil
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), remotecluster.SendTimeout)
|
||||
defer cancel()
|
||||
|
||||
|
||||
@@ -121,7 +121,7 @@ func (scs *Service) InviteRemoteToChannel(channelID, remoteID, userID string, sh
|
||||
ChannelId: channelID,
|
||||
CreatorId: userID,
|
||||
Home: true,
|
||||
RemoteId: remoteID,
|
||||
RemoteId: "", // channel originates locally
|
||||
}
|
||||
if _, err = scs.ShareChannel(sc); err != nil {
|
||||
return model.NewAppError("InviteRemoteToChannel", "api.command_share.share_channel.error",
|
||||
|
||||
@@ -219,7 +219,7 @@ func (scs *Service) upsertSyncUser(c request.CTX, user *model.User, channel *mod
|
||||
return userSaved, nil
|
||||
}
|
||||
|
||||
func (scs *Service) insertSyncUser(rctx request.CTX, user *model.User, channel *model.Channel, rc *model.RemoteCluster) (*model.User, error) {
|
||||
func (scs *Service) insertSyncUser(rctx request.CTX, user *model.User, _ *model.Channel, rc *model.RemoteCluster) (*model.User, error) {
|
||||
var err error
|
||||
var userSaved *model.User
|
||||
var suffix string
|
||||
@@ -270,7 +270,7 @@ func (scs *Service) insertSyncUser(rctx request.CTX, user *model.User, channel *
|
||||
return nil, fmt.Errorf("error inserting sync user %s: %w", user.Id, err)
|
||||
}
|
||||
|
||||
func (scs *Service) updateSyncUser(rctx request.CTX, patch *model.UserPatch, user *model.User, channel *model.Channel, rc *model.RemoteCluster) (*model.User, error) {
|
||||
func (scs *Service) updateSyncUser(rctx request.CTX, patch *model.UserPatch, user *model.User, _ *model.Channel, rc *model.RemoteCluster) (*model.User, error) {
|
||||
var err error
|
||||
var update *model.UserUpdate
|
||||
var suffix string
|
||||
|
||||
@@ -107,6 +107,10 @@ func (scs *Service) syncForRemote(task syncTask, rc *model.RemoteCluster) error
|
||||
if scr, err = scs.server.GetStore().SharedChannel().SaveRemote(scr); err != nil {
|
||||
return fmt.Errorf("cannot auto-create shared channel remote (channel_id=%s, remote_id=%s): %w", task.channelID, rc.RemoteId, err)
|
||||
}
|
||||
scs.server.Log().Log(mlog.LvlSharedChannelServiceDebug, "Auto-invited remote to channel (BitflagOptionAutoInvited)",
|
||||
mlog.String("remote", rc.DisplayName),
|
||||
mlog.String("channel_id", task.channelID),
|
||||
)
|
||||
} else if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
Ссылка в новой задаче
Block a user