From c11ad8995fcaaf76a022fb43be85c0bc604d5204 Mon Sep 17 00:00:00 2001 From: Kyriakos Z <3829551+koox00@users.noreply.github.com> Date: Fri, 2 Sep 2022 13:17:22 +0300 Subject: [PATCH] Adds OmitConnection parameter to broadcast (#20723) * Adds OmitConnection parameter to broadcast Currently we have no means to omit sending a websocket event to a specific connection id. This is needed mainly so that the initiator won't receive an event for the action it just initiated. Will be used for the global drafts feature, so that we won't update drafts through ws when a user is typing. This commit adds OmitConnection to the Broadcast struct and to the NewWebSocketEvent function signature. shouldSendEvent should return false for that specific connection. * Return early only if connection id matches the omitted Co-authored-by: Mattermod --- api4/plugin.go | 2 +- api4/user.go | 2 +- api4/websocket_test.go | 4 +- app/admin_advisor.go | 2 +- app/channel.go | 44 ++++++++++---------- app/channel_category.go | 8 ++-- app/cluster.go | 2 +- app/emoji.go | 2 +- app/group.go | 14 +++---- app/integration_action.go | 2 +- app/notification.go | 4 +- app/plugin.go | 6 +-- app/plugin_api.go | 2 +- app/plugin_statuses.go | 2 +- app/post.go | 14 +++---- app/preference.go | 8 ++-- app/reaction.go | 2 +- app/role.go | 2 +- app/server.go | 4 +- app/shared_channel_notifier_test.go | 6 +-- app/slashcommands/command_expand_collapse.go | 2 +- app/slashcommands/command_share.go | 2 +- app/status.go | 2 +- app/team.go | 12 +++--- app/teams/teams.go | 2 +- app/user.go | 30 ++++++------- app/web_conn.go | 7 +++- app/web_conn_test.go | 9 ++-- app/web_hub_test.go | 8 ++-- app/webhub_fuzz.go | 2 +- model/websocket_message.go | 23 +++++----- model/websocket_message_test.go | 10 ++--- 32 files changed, 125 insertions(+), 116 deletions(-) diff --git a/api4/plugin.go b/api4/plugin.go index 775b515a57..eb5fd11652 100644 --- a/api4/plugin.go +++ b/api4/plugin.go @@ -427,7 +427,7 @@ func setFirstAdminVisitMarketplaceStatus(c *Context, w http.ResponseWriter, r *h return } - message := model.NewWebSocketEvent(model.WebsocketFirstAdminVisitMarketplaceStatusReceived, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketFirstAdminVisitMarketplaceStatusReceived, "", "", "", nil, "") message.Add("firstAdminVisitMarketplaceStatus", firstAdminVisitMarketplaceObj.Value) c.App.Publish(message) diff --git a/api4/user.go b/api4/user.go index 55156821a8..ef35adf353 100644 --- a/api4/user.go +++ b/api4/user.go @@ -1504,7 +1504,7 @@ func updateUserActive(c *Context, w http.ResponseWriter, r *http.Request) { }) } - message := model.NewWebSocketEvent(model.WebsocketEventUserActivationStatusChange, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventUserActivationStatusChange, "", "", "", nil, "") c.App.Publish(message) ReturnStatusOK(w) diff --git a/api4/websocket_test.go b/api4/websocket_test.go index 7655a44f48..051b7ae28f 100644 --- a/api4/websocket_test.go +++ b/api4/websocket_test.go @@ -41,7 +41,7 @@ func TestWebSocketEvent(t *testing.T) { omitUser := make(map[string]bool, 1) omitUser["somerandomid"] = true - evt1 := model.NewWebSocketEvent(model.WebsocketEventTyping, "", th.BasicChannel.Id, "", omitUser) + evt1 := model.NewWebSocketEvent(model.WebsocketEventTyping, "", th.BasicChannel.Id, "", omitUser, "") evt1.Add("user_id", "somerandomid") th.App.Publish(evt1) @@ -69,7 +69,7 @@ func TestWebSocketEvent(t *testing.T) { require.True(t, eventHit, "did not receive typing event") - evt2 := model.NewWebSocketEvent(model.WebsocketEventTyping, "", "somerandomid", "", nil) + evt2 := model.NewWebSocketEvent(model.WebsocketEventTyping, "", "somerandomid", "", nil, "") th.App.Publish(evt2) time.Sleep(300 * time.Millisecond) diff --git a/app/admin_advisor.go b/app/admin_advisor.go index ca588b21f7..2fdd0c16be 100644 --- a/app/admin_advisor.go +++ b/app/admin_advisor.go @@ -213,7 +213,7 @@ func (a *App) setWarnMetricsStatusAndNotify(warnMetricId string) *model.AppError } // Inform client that this metric warning has been acked - message := model.NewWebSocketEvent(model.WebsocketWarnMetricStatusRemoved, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketWarnMetricStatusRemoved, "", "", "", nil, "") message.Add("warnMetricId", warnMetricId) a.Publish(message) diff --git a/app/channel.go b/app/channel.go index df8edb67a9..e7c0826c9a 100644 --- a/app/channel.go +++ b/app/channel.go @@ -126,7 +126,7 @@ func (a *App) JoinDefaultChannels(c request.CTX, teamID string, user *model.User a.invalidateCacheForChannelMembers(channel.Id) - message := model.NewWebSocketEvent(model.WebsocketEventUserAdded, "", channel.Id, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventUserAdded, "", channel.Id, "", nil, "") message.Add("user_id", user.Id) message.Add("team_id", channel.TeamId) a.Publish(message) @@ -209,7 +209,7 @@ func (a *App) CreateChannelWithUser(c request.CTX, channel *model.Channel, userI a.postJoinChannelMessage(c, user, channel) - message := model.NewWebSocketEvent(model.WebsocketEventChannelCreated, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelCreated, "", "", userID, nil, "") message.Add("channel_id", channel.Id) message.Add("team_id", channel.TeamId) a.Publish(message) @@ -409,7 +409,7 @@ func (a *App) handleCreationEvent(c request.CTX, userID, otherUserID string, cha }) } - message := model.NewWebSocketEvent(model.WebsocketEventDirectAdded, "", channel.Id, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventDirectAdded, "", channel.Id, "", nil, "") message.Add("creator_id", userID) message.Add("teammate_id", otherUserID) a.Publish(message) @@ -527,7 +527,7 @@ func (a *App) CreateGroupChannel(c request.CTX, userIDs []string, creatorId stri a.InvalidateCacheForUser(userID) } - message := model.NewWebSocketEvent(model.WebsocketEventGroupAdded, "", channel.Id, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventGroupAdded, "", channel.Id, "", nil, "") message.Add("teammate_ids", model.ArrayToJSON(userIDs)) a.Publish(message) @@ -653,7 +653,7 @@ func (a *App) UpdateChannel(c request.CTX, channel *model.Channel) (*model.Chann a.invalidateCacheForChannel(channel) - messageWs := model.NewWebSocketEvent(model.WebsocketEventChannelUpdated, "", channel.Id, "", nil) + messageWs := model.NewWebSocketEvent(model.WebsocketEventChannelUpdated, "", channel.Id, "", nil, "") channelJSON, jsonErr := json.Marshal(channel) if jsonErr != nil { return nil, model.NewAppError("UpdateChannel", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -724,7 +724,7 @@ func (a *App) UpdateChannelPrivacy(c request.CTX, oldChannel *model.Channel, use a.invalidateCacheForChannel(channel) - messageWs := model.NewWebSocketEvent(model.WebsocketEventChannelConverted, channel.TeamId, "", "", nil) + messageWs := model.NewWebSocketEvent(model.WebsocketEventChannelConverted, channel.TeamId, "", "", nil, "") messageWs.Add("channel_id", channel.Id) a.Publish(messageWs) @@ -779,7 +779,7 @@ func (a *App) RestoreChannel(c request.CTX, channel *model.Channel, userID strin channel.DeleteAt = 0 a.invalidateCacheForChannel(channel) - message := model.NewWebSocketEvent(model.WebsocketEventChannelRestored, channel.TeamId, "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelRestored, channel.TeamId, "", "", nil, "") message.Add("channel_id", channel.Id) a.Publish(message) @@ -1018,7 +1018,7 @@ func (a *App) PatchChannelModerationsForChannel(c request.CTX, channel *model.Ch return nil, appErr } - message := model.NewWebSocketEvent(model.WebsocketEventChannelSchemeUpdated, "", channel.Id, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelSchemeUpdated, "", channel.Id, "", nil, "") a.Publish(message) c.Logger().Info("Permission scheme created.", mlog.String("channel_id", channel.Id), mlog.String("channel_name", channel.Name)) } else { @@ -1076,7 +1076,7 @@ func (a *App) PatchChannelModerationsForChannel(c request.CTX, channel *model.Ch return nil, err } - message := model.NewWebSocketEvent(model.WebsocketEventChannelSchemeUpdated, "", channel.Id, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelSchemeUpdated, "", channel.Id, "", nil, "") a.Publish(message) memberRole = higherScopedMemberRole @@ -1279,7 +1279,7 @@ func (a *App) UpdateChannelMemberNotifyProps(c request.CTX, data map[string]stri a.invalidateCacheForChannelMembersNotifyProps(member.ChannelId) // Notify the clients that the member notify props changed - evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", member.UserId, nil) + evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", member.UserId, nil, "") memberJSON, jsonErr := json.Marshal(member) if jsonErr != nil { return nil, model.NewAppError("UpdateChannelMemberNotifyProps", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -1308,7 +1308,7 @@ func (a *App) updateChannelMember(c request.CTX, member *model.ChannelMember) (* a.InvalidateCacheForUser(member.UserId) // Notify the clients that the member notify props changed - evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", member.UserId, nil) + evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", member.UserId, nil, "") memberJSON, jsonErr := json.Marshal(member) if jsonErr != nil { return nil, model.NewAppError("updateChannelMember", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -1434,7 +1434,7 @@ func (a *App) DeleteChannel(c request.CTX, channel *model.Channel, userID string } a.invalidateCacheForChannel(channel) - message := model.NewWebSocketEvent(model.WebsocketEventChannelDeleted, channel.TeamId, "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelDeleted, channel.TeamId, "", "", nil, "") message.Add("channel_id", channel.Id) message.Add("delete_at", deleteAt) a.Publish(message) @@ -1524,7 +1524,7 @@ func (a *App) AddUserToChannel(c request.CTX, user *model.User, channel *model.C return nil, err } - message := model.NewWebSocketEvent(model.WebsocketEventUserAdded, "", channel.Id, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventUserAdded, "", channel.Id, "", nil, "") message.Add("user_id", user.Id) message.Add("team_id", channel.TeamId) a.Publish(message) @@ -2479,13 +2479,13 @@ func (a *App) removeUserFromChannel(c request.CTX, userIDToRemove string, remove }) } - message := model.NewWebSocketEvent(model.WebsocketEventUserRemoved, "", channel.Id, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventUserRemoved, "", channel.Id, "", nil, "") message.Add("user_id", userIDToRemove) message.Add("remover_id", removerUserId) a.Publish(message) // because the removed user no longer belongs to the channel we need to send a separate websocket event - userMsg := model.NewWebSocketEvent(model.WebsocketEventUserRemoved, "", "", userIDToRemove, nil) + userMsg := model.NewWebSocketEvent(model.WebsocketEventUserRemoved, "", "", userIDToRemove, nil, "") userMsg.Add("channel_id", channel.Id) userMsg.Add("remover_id", removerUserId) a.Publish(userMsg) @@ -2698,7 +2698,7 @@ func (a *App) markChannelAsUnreadFromPostCRTUnsupported(c request.CTX, postID st if jsonErr != nil { return nil, model.NewAppError("MarkChannelAsUnreadFromPost", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) } - message := model.NewWebSocketEvent(model.WebsocketEventThreadUpdated, channel.TeamId, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventThreadUpdated, channel.TeamId, "", userID, nil, "") message.Add("thread", string(payload)) a.Publish(message) } @@ -2714,7 +2714,7 @@ func (a *App) markChannelAsUnreadFromPostCRTUnsupported(c request.CTX, postID st } func (a *App) sendWebSocketPostUnreadEvent(c request.CTX, channelUnread *model.ChannelUnreadAt, postID string, withMsgCountRoot bool) { - message := model.NewWebSocketEvent(model.WebsocketEventPostUnread, channelUnread.TeamId, channelUnread.ChannelId, channelUnread.UserId, nil) + message := model.NewWebSocketEvent(model.WebsocketEventPostUnread, channelUnread.TeamId, channelUnread.ChannelId, channelUnread.UserId, nil, "") message.Add("msg_count", channelUnread.MsgCount) if withMsgCountRoot { message.Add("msg_count_root", channelUnread.MsgCountRoot) @@ -2929,7 +2929,7 @@ func (a *App) MarkChannelsAsViewed(c request.CTX, channelIDs []string, userID st if *a.Config().ServiceSettings.EnableChannelViewedMessages { for _, channelID := range channelIDs { - message := model.NewWebSocketEvent(model.WebsocketEventChannelViewed, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelViewed, "", "", userID, nil, "") message.Add("channel_id", channelID) a.Publish(message) } @@ -2941,7 +2941,7 @@ func (a *App) MarkChannelsAsViewed(c request.CTX, channelIDs []string, userID st if updateThreads && a.IsCRTEnabledForUser(c, userID) { timestamp := model.GetMillis() for _, channelID := range channelIDs { - message := model.NewWebSocketEvent(model.WebsocketEventThreadReadChanged, "", channelID, userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventThreadReadChanged, "", channelID, userID, nil, "") message.Add("timestamp", timestamp) a.Publish(message) } @@ -2996,7 +2996,7 @@ func (a *App) PermanentDeleteChannel(c request.CTX, channel *model.Channel) *mod } a.invalidateCacheForChannel(channel) - message := model.NewWebSocketEvent(model.WebsocketEventChannelDeleted, channel.TeamId, "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelDeleted, channel.TeamId, "", "", nil, "") message.Add("channel_id", channel.Id) message.Add("delete_at", deleteAt) a.Publish(message) @@ -3259,7 +3259,7 @@ func (a *App) setChannelsMuted(c request.CTX, channelIDs []string, userID string for _, member := range updated { a.invalidateCacheForChannelMembersNotifyProps(member.ChannelId) - evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", member.UserId, nil) + evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", member.UserId, nil, "") memberJSON, jsonErr := json.Marshal(member) if jsonErr != nil { @@ -3367,7 +3367,7 @@ func (a *App) forEachChannelMember(c request.CTX, channelID string, f func(model func (a *App) ClearChannelMembersCache(c request.CTX, channelID string) error { clearSessionCache := func(channelMember model.ChannelMember) error { a.ClearSessionCacheForUser(channelMember.UserId) - message := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", channelMember.UserId, nil) + message := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", channelMember.UserId, nil, "") memberJSON, jsonErr := json.Marshal(channelMember) if jsonErr != nil { return jsonErr diff --git a/app/channel_category.go b/app/channel_category.go index 1420935b79..65f80c8545 100644 --- a/app/channel_category.go +++ b/app/channel_category.go @@ -115,7 +115,7 @@ func (a *App) CreateSidebarCategory(c request.CTX, userID, teamID string, newCat return nil, model.NewAppError("CreateSidebarCategory", "app.channel.sidebar_categories.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } } - message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryCreated, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryCreated, teamID, "", userID, nil, "") message.Add("category_id", category.Id) a.Publish(message) return category, nil @@ -135,7 +135,7 @@ func (a *App) UpdateSidebarCategoryOrder(c request.CTX, userID, teamID string, c return model.NewAppError("UpdateSidebarCategoryOrder", "app.channel.sidebar_categories.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } } - message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryOrderUpdated, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryOrderUpdated, teamID, "", userID, nil, "") message.Add("order", categoryOrder) a.Publish(message) return nil @@ -147,7 +147,7 @@ func (a *App) UpdateSidebarCategories(c request.CTX, userID, teamID string, cate return nil, model.NewAppError("UpdateSidebarCategories", "app.channel.sidebar_categories.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } - message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryUpdated, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryUpdated, teamID, "", userID, nil, "") updatedCategoriesJSON, jsonErr := json.Marshal(updatedCategories) if jsonErr != nil { @@ -280,7 +280,7 @@ func (a *App) DeleteSidebarCategory(c request.CTX, userID, teamID, categoryId st } } - message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryDeleted, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryDeleted, teamID, "", userID, nil, "") message.Add("category_id", categoryId) a.Publish(message) diff --git a/app/cluster.go b/app/cluster.go index 4e3b90dbe0..ae9e98fa8e 100644 --- a/app/cluster.go +++ b/app/cluster.go @@ -48,7 +48,7 @@ func (s *clusterWrapper) PublishPluginClusterEvent(productID string, ev model.Pl } func (s *clusterWrapper) PublishWebSocketEvent(productID string, event string, payload map[string]any, broadcast *model.WebsocketBroadcast) { - ev := model.NewWebSocketEvent(fmt.Sprintf("custom_%v_%v", productID, event), "", "", "", nil) + ev := model.NewWebSocketEvent(fmt.Sprintf("custom_%v_%v", productID, event), "", "", "", nil, "") ev = ev.SetBroadcast(broadcast).SetData(payload) s.srv.Publish(ev) } diff --git a/app/emoji.go b/app/emoji.go index 930abb8fa5..25f1d6bece 100644 --- a/app/emoji.go +++ b/app/emoji.go @@ -78,7 +78,7 @@ func (a *App) CreateEmoji(sessionUserId string, emoji *model.Emoji, multiPartIma return nil, model.NewAppError("CreateEmoji", "app.emoji.create.internal_error", nil, "", http.StatusInternalServerError).Wrap(err) } - message := model.NewWebSocketEvent(model.WebsocketEventEmojiAdded, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventEmojiAdded, "", "", "", nil, "") emojiJSON, jsonErr := json.Marshal(emoji) if jsonErr != nil { return nil, model.NewAppError("CreateEmoji", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) diff --git a/app/group.go b/app/group.go index 1633d609d0..4a1829cef4 100644 --- a/app/group.go +++ b/app/group.go @@ -145,7 +145,7 @@ func (a *App) CreateGroupWithUserIds(group *model.GroupWithUserIds) (*model.Grou } } - messageWs := model.NewWebSocketEvent(model.WebsocketEventReceivedGroup, "", "", "", nil) + messageWs := model.NewWebSocketEvent(model.WebsocketEventReceivedGroup, "", "", "", nil, "") count, err := a.Srv().Store.Group().GetMemberCount(newGroup.Id) if err != nil { return nil, model.NewAppError("CreateGroupWithUserIds", "app.group.id.app_error", nil, "", http.StatusBadRequest).Wrap(err) @@ -190,7 +190,7 @@ func (a *App) UpdateGroup(group *model.Group) (*model.Group, *model.AppError) { } updatedGroup.MemberCount = model.NewInt(int(count)) - messageWs := model.NewWebSocketEvent(model.WebsocketEventReceivedGroup, "", "", "", nil) + messageWs := model.NewWebSocketEvent(model.WebsocketEventReceivedGroup, "", "", "", nil, "") groupJSON, err := json.Marshal(updatedGroup) if err != nil { @@ -381,9 +381,9 @@ func (a *App) UpsertGroupSyncable(groupSyncable *model.GroupSyncable) (*model.Gr var messageWs *model.WebSocketEvent if gs.Type == model.GroupSyncableTypeTeam { - messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupAssociatedToTeam, gs.SyncableId, "", "", nil) + messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupAssociatedToTeam, gs.SyncableId, "", "", nil, "") } else { - messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupAssociatedToChannel, "", gs.SyncableId, "", nil) + messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupAssociatedToChannel, "", gs.SyncableId, "", nil, "") } messageWs.Add("group_id", gs.GroupId) a.Publish(messageWs) @@ -482,9 +482,9 @@ func (a *App) DeleteGroupSyncable(groupID string, syncableID string, syncableTyp var messageWs *model.WebSocketEvent if gs.Type == model.GroupSyncableTypeTeam { - messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupNotAssociatedToTeam, gs.SyncableId, "", "", nil) + messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupNotAssociatedToTeam, gs.SyncableId, "", "", nil, "") } else { - messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupNotAssociatedToChannel, "", gs.SyncableId, "", nil) + messageWs = model.NewWebSocketEvent(model.WebsocketEventReceivedGroupNotAssociatedToChannel, "", gs.SyncableId, "", nil, "") } messageWs.Add("group_id", gs.GroupId) @@ -779,7 +779,7 @@ func (a *App) DeleteGroupMembers(groupID string, userIDs []string) ([]*model.Gro } func (a *App) publishGroupMemberEvent(eventName string, groupMember *model.GroupMember) *model.AppError { - messageWs := model.NewWebSocketEvent(eventName, "", "", groupMember.UserId, nil) + messageWs := model.NewWebSocketEvent(eventName, "", "", groupMember.UserId, nil, "") groupMemberJSON, jsonErr := json.Marshal(groupMember) if jsonErr != nil { return model.NewAppError("publishGroupMemberEvent", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) diff --git a/app/integration_action.go b/app/integration_action.go index 63a4278782..e85ad88f5e 100644 --- a/app/integration_action.go +++ b/app/integration_action.go @@ -597,7 +597,7 @@ func (a *App) OpenInteractiveDialog(request model.OpenDialogRequest) *model.AppE a.ch.srv.Log().Warn("Error encoding request", mlog.Err(err)) } - message := model.NewWebSocketEvent(model.WebsocketEventOpenDialog, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventOpenDialog, "", "", userID, nil, "") message.Add("dialog", string(jsonRequest)) a.Publish(message) diff --git a/app/notification.go b/app/notification.go index b1ea390836..b0870c612f 100644 --- a/app/notification.go +++ b/app/notification.go @@ -525,7 +525,7 @@ func (a *App) SendNotifications(c request.CTX, post *model.Post, team *model.Tea } } - message := model.NewWebSocketEvent(model.WebsocketEventPosted, "", post.ChannelId, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventPosted, "", post.ChannelId, "", nil, "") // Note that PreparePostForClient should've already been called by this point postJSON, jsonErr := post.ToJSON() @@ -584,7 +584,7 @@ func (a *App) SendNotifications(c request.CTX, post *model.Post, team *model.Tea continue } if a.IsCRTEnabledForUser(c, uid) { - message := model.NewWebSocketEvent(model.WebsocketEventThreadUpdated, team.Id, "", uid, nil) + message := model.NewWebSocketEvent(model.WebsocketEventThreadUpdated, team.Id, "", uid, nil, "") threadMembership := participantMemberships[uid] if threadMembership == nil { tm, err := a.Srv().Store.Thread().GetMembershipForUser(uid, post.RootId) diff --git a/app/plugin.go b/app/plugin.go index 261c395e7f..00cabfdbc2 100644 --- a/app/plugin.go +++ b/app/plugin.go @@ -146,7 +146,7 @@ func (ch *Channels) syncPluginsActiveState() { deactivated := pluginsEnvironment.Deactivate(plugin.Manifest.Id) if deactivated && plugin.Manifest.HasClient() { - message := model.NewWebSocketEvent(model.WebsocketEventPluginDisabled, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventPluginDisabled, "", "", "", nil, "") message.Add("manifest", plugin.Manifest.ClientManifest()) ch.srv.Publish(message) } @@ -503,7 +503,7 @@ func (ch *Channels) notifyIntegrationsUsageChanged() *model.AppError { return appErr } - message := model.NewWebSocketEvent(model.WebsocketEventIntegrationsUsageChanged, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventIntegrationsUsageChanged, "", "", "", nil, "") message.Add("usage", usage) message.GetBroadcast().ContainsSensitiveData = true ch.Publish(message) @@ -866,7 +866,7 @@ func (ch *Channels) notifyPluginEnabled(manifest *model.Manifest) error { } // Notify all cluster peer clients. - message := model.NewWebSocketEvent(model.WebsocketEventPluginEnabled, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventPluginEnabled, "", "", "", nil, "") message.Add("manifest", manifest.ClientManifest()) ch.srv.Publish(message) diff --git a/app/plugin_api.go b/app/plugin_api.go index ceb7e6ea2b..f922d37ee9 100644 --- a/app/plugin_api.go +++ b/app/plugin_api.go @@ -943,7 +943,7 @@ func (api *PluginAPI) KVList(page, perPage int) ([]string, *model.AppError) { } func (api *PluginAPI) PublishWebSocketEvent(event string, payload map[string]any, broadcast *model.WebsocketBroadcast) { - ev := model.NewWebSocketEvent(fmt.Sprintf("custom_%v_%v", api.id, event), "", "", "", nil) + ev := model.NewWebSocketEvent(fmt.Sprintf("custom_%v_%v", api.id, event), "", "", "", nil, "") ev = ev.SetBroadcast(broadcast).SetData(payload) api.app.Publish(ev) } diff --git a/app/plugin_statuses.go b/app/plugin_statuses.go index ad71925967..c183c6402d 100644 --- a/app/plugin_statuses.go +++ b/app/plugin_statuses.go @@ -99,7 +99,7 @@ func (ch *Channels) notifyPluginStatusesChanged() error { } // Notify any system admins. - message := model.NewWebSocketEvent(model.WebsocketEventPluginStatusesChanged, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventPluginStatusesChanged, "", "", "", nil, "") message.Add("plugin_statuses", pluginStatuses) message.GetBroadcast().ContainsSensitiveData = true ch.srv.Publish(message) diff --git a/app/post.go b/app/post.go index fc6d20bfd0..0a6c300e5c 100644 --- a/app/post.go +++ b/app/post.go @@ -510,7 +510,7 @@ func (a *App) SendEphemeralPost(c request.CTX, userID string, post *model.Post) } post.GenerateActionIds() - message := model.NewWebSocketEvent(model.WebsocketEventEphemeralMessage, "", post.ChannelId, userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventEphemeralMessage, "", post.ChannelId, userID, nil, "") post = a.PreparePostForClientWithEmbedsAndImages(c, post, true, false) post = model.AddPostActionCookies(post, a.PostActionCookieSecret()) @@ -533,7 +533,7 @@ func (a *App) UpdateEphemeralPost(c request.CTX, userID string, post *model.Post } post.GenerateActionIds() - message := model.NewWebSocketEvent(model.WebsocketEventPostEdited, "", post.ChannelId, userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventPostEdited, "", post.ChannelId, userID, nil, "") post = a.PreparePostForClientWithEmbedsAndImages(c, post, true, false) post = model.AddPostActionCookies(post, a.PostActionCookieSecret()) postJSON, jsonErr := post.ToJSON() @@ -555,7 +555,7 @@ func (a *App) DeleteEphemeralPost(userID, postID string) { UpdateAt: model.GetMillis(), } - message := model.NewWebSocketEvent(model.WebsocketEventPostDeleted, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventPostDeleted, "", "", userID, nil, "") postJSON, jsonErr := post.ToJSON() if jsonErr != nil { mlog.Warn("Failed to encode post to JSON", mlog.Err(jsonErr)) @@ -690,7 +690,7 @@ func (a *App) UpdatePost(c *request.Context, post *model.Post, safeUpdate bool) return nil, model.NewAppError("UpdatePost", "app.post.update.app_error", nil, "", http.StatusInternalServerError).Wrap(nErr) } - message := model.NewWebSocketEvent(model.WebsocketEventPostEdited, "", rpost.ChannelId, "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventPostEdited, "", rpost.ChannelId, "", nil, "") postJSON, jsonErr := rpost.ToJSON() if jsonErr != nil { return nil, model.NewAppError("UpdatePost", "app.post.marshal.app_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -1260,12 +1260,12 @@ func (a *App) DeletePost(c request.CTX, postID, deleteByID string) (*model.Post, return nil, model.NewAppError("DeletePost", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(err) } - userMessage := model.NewWebSocketEvent(model.WebsocketEventPostDeleted, "", post.ChannelId, "", nil) + userMessage := model.NewWebSocketEvent(model.WebsocketEventPostDeleted, "", post.ChannelId, "", nil, "") userMessage.Add("post", string(postJSON)) userMessage.GetBroadcast().ContainsSanitizedData = true a.Publish(userMessage) - adminMessage := model.NewWebSocketEvent(model.WebsocketEventPostDeleted, "", post.ChannelId, "", nil) + adminMessage := model.NewWebSocketEvent(model.WebsocketEventPostDeleted, "", post.ChannelId, "", nil, "") adminMessage.Add("post", string(postJSON)) adminMessage.Add("delete_by", deleteByID) adminMessage.GetBroadcast().ContainsSensitiveData = true @@ -1989,7 +1989,7 @@ func (a *App) SetPostReminder(postID, userID string, targetTime int64) *model.Ap }, } - message := model.NewWebSocketEvent(model.WebsocketEventEphemeralMessage, "", ephemeralPost.ChannelId, userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventEphemeralMessage, "", ephemeralPost.ChannelId, userID, nil, "") ephemeralPost = a.PreparePostForClientWithEmbedsAndImages(request.EmptyContext(a.Log()), ephemeralPost, true, false) ephemeralPost = model.AddPostActionCookies(ephemeralPost, a.PostActionCookieSecret()) diff --git a/app/preference.go b/app/preference.go index c6bfee0b58..4cf8903740 100644 --- a/app/preference.go +++ b/app/preference.go @@ -82,11 +82,11 @@ func (a *App) UpdatePreferences(userID string, preferences model.Preferences) *m return model.NewAppError("UpdatePreferences", "api.preference.update_preferences.update_sidebar.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } - message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryUpdated, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryUpdated, "", "", userID, nil, "") // TODO this needs to be updated to include information on which categories changed a.Publish(message) - message = model.NewWebSocketEvent(model.WebsocketEventPreferencesChanged, "", "", userID, nil) + message = model.NewWebSocketEvent(model.WebsocketEventPreferencesChanged, "", "", userID, nil, "") prefsJSON, jsonErr := json.Marshal(preferences) if jsonErr != nil { return model.NewAppError("UpdatePreferences", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -116,11 +116,11 @@ func (a *App) DeletePreferences(userID string, preferences model.Preferences) *m return model.NewAppError("DeletePreferences", "api.preference.delete_preferences.update_sidebar.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } - message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryUpdated, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventSidebarCategoryUpdated, "", "", userID, nil, "") // TODO this needs to be updated to include information on which categories changed a.Publish(message) - message = model.NewWebSocketEvent(model.WebsocketEventPreferencesDeleted, "", "", userID, nil) + message = model.NewWebSocketEvent(model.WebsocketEventPreferencesDeleted, "", "", userID, nil, "") prefsJSON, jsonErr := json.Marshal(preferences) if jsonErr != nil { return model.NewAppError("DeletePreferences", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) diff --git a/app/reaction.go b/app/reaction.go index 3e3c253ffa..b48c33536b 100644 --- a/app/reaction.go +++ b/app/reaction.go @@ -161,7 +161,7 @@ func (a *App) DeleteReactionForPost(c *request.Context, reaction *model.Reaction func (a *App) sendReactionEvent(event string, reaction *model.Reaction, post *model.Post) { // send out that a reaction has been added/removed - message := model.NewWebSocketEvent(event, "", post.ChannelId, "", nil) + message := model.NewWebSocketEvent(event, "", post.ChannelId, "", nil, "") reactionJSON, err := json.Marshal(reaction) if err != nil { a.Log().Warn("Failed to encode reaction to JSON", mlog.Err(err)) diff --git a/app/role.go b/app/role.go index 535a1b605f..156ef32a49 100644 --- a/app/role.go +++ b/app/role.go @@ -259,7 +259,7 @@ func (a *App) CheckRolesExist(roleNames []string) *model.AppError { } func (a *App) sendUpdatedRoleEvent(role *model.Role) *model.AppError { - message := model.NewWebSocketEvent(model.WebsocketEventRoleUpdated, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventRoleUpdated, "", "", "", nil, "") roleJSON, jsonErr := json.Marshal(role) if jsonErr != nil { return model.NewAppError("sendUpdatedRoleEvent", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) diff --git a/app/server.go b/app/server.go index 17dcfdcba0..135a3eb5df 100644 --- a/app/server.go +++ b/app/server.go @@ -507,7 +507,7 @@ func NewServer(options ...Option) (*Server, error) { ch := s.Channels() ch.regenerateClientConfig() - message := model.NewWebSocketEvent(model.WebsocketEventConfigChanged, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventConfigChanged, "", "", "", nil, "") appInstance := New(ServerConnector(ch)) message.Add("config", appInstance.ClientConfigWithComputed()) @@ -523,7 +523,7 @@ func NewServer(options ...Option) (*Server, error) { s.licenseListenerId = s.AddLicenseListener(func(oldLicense, newLicense *model.License) { s.Channels().regenerateClientConfig() - message := model.NewWebSocketEvent(model.WebsocketEventLicenseChanged, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventLicenseChanged, "", "", "", nil, "") message.Add("license", s.GetSanitizedClientLicense()) s.Go(func() { s.Publish(message) diff --git a/app/shared_channel_notifier_test.go b/app/shared_channel_notifier_test.go index 510fe9a19c..6800cc8957 100644 --- a/app/shared_channel_notifier_test.go +++ b/app/shared_channel_notifier_test.go @@ -33,7 +33,7 @@ func TestServerSyncSharedChannelHandler(t *testing.T) { th.App.ch.srv.SetSharedChannelSyncService(mockService) channel := th.CreateChannel(th.Context, th.BasicTeam, WithShared(true)) - websocketEvent := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, model.NewId(), channel.Id, "", nil) + websocketEvent := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, model.NewId(), channel.Id, "", nil, "") th.App.ch.srv.SharedChannelSyncHandler(websocketEvent) assert.Empty(t, mockService.channelNotifications) @@ -47,7 +47,7 @@ func TestServerSyncSharedChannelHandler(t *testing.T) { mockService.active = true th.App.ch.srv.SetSharedChannelSyncService(mockService) - websocketEvent := model.NewWebSocketEvent(model.WebsocketEventPosted, model.NewId(), model.NewId(), "", nil) + websocketEvent := model.NewWebSocketEvent(model.WebsocketEventPosted, model.NewId(), model.NewId(), "", nil, "") th.App.ch.srv.SharedChannelSyncHandler(websocketEvent) assert.Empty(t, mockService.channelNotifications) @@ -62,7 +62,7 @@ func TestServerSyncSharedChannelHandler(t *testing.T) { th.App.ch.srv.SetSharedChannelSyncService(mockService) channel := th.CreateChannel(th.Context, th.BasicTeam, WithShared(true)) - websocketEvent := model.NewWebSocketEvent(model.WebsocketEventPosted, model.NewId(), channel.Id, "", nil) + websocketEvent := model.NewWebSocketEvent(model.WebsocketEventPosted, model.NewId(), channel.Id, "", nil, "") th.App.ch.srv.SharedChannelSyncHandler(websocketEvent) assert.Len(t, mockService.channelNotifications, 1) diff --git a/app/slashcommands/command_expand_collapse.go b/app/slashcommands/command_expand_collapse.go index 7c8ab6583f..7dfada3d36 100644 --- a/app/slashcommands/command_expand_collapse.go +++ b/app/slashcommands/command_expand_collapse.go @@ -75,7 +75,7 @@ func setCollapsePreference(a *app.App, args *model.CommandArgs, isCollapse bool) return &model.CommandResponse{Text: args.T("api.command_expand_collapse.fail.app_error") + err.Error(), ResponseType: model.CommandResponseTypeEphemeral} } - socketMessage := model.NewWebSocketEvent(model.WebsocketEventPreferenceChanged, "", "", args.UserId, nil) + socketMessage := model.NewWebSocketEvent(model.WebsocketEventPreferenceChanged, "", "", args.UserId, nil, "") prefJSON, err := json.Marshal(pref) if err != nil { diff --git a/app/slashcommands/command_share.go b/app/slashcommands/command_share.go index 997411fa31..9a9d62be58 100644 --- a/app/slashcommands/command_share.go +++ b/app/slashcommands/command_share.go @@ -323,7 +323,7 @@ func (sp *ShareProvider) doStatus(a *app.App, args *model.CommandArgs, _ map[str } func notifyClientsForChannelUpdate(a *app.App, sharedChannel *model.SharedChannel) { - messageWs := model.NewWebSocketEvent(model.WebsocketEventChannelConverted, sharedChannel.TeamId, "", "", nil) + messageWs := model.NewWebSocketEvent(model.WebsocketEventChannelConverted, sharedChannel.TeamId, "", "", nil, "") messageWs.Add("channel_id", sharedChannel.ChannelId) a.Publish(messageWs) } diff --git a/app/status.go b/app/status.go index 6d420e2993..21c47ca760 100644 --- a/app/status.go +++ b/app/status.go @@ -234,7 +234,7 @@ func (a *App) BroadcastStatus(status *model.Status) { // this is considered a non-critical service and will be disabled when server busy. return } - event := model.NewWebSocketEvent(model.WebsocketEventStatusChange, "", "", status.UserId, nil) + event := model.NewWebSocketEvent(model.WebsocketEventStatusChange, "", "", status.UserId, nil, "") event.Add("status", status.Status) event.Add("user_id", status.UserId) a.Publish(event) diff --git a/app/team.go b/app/team.go index c02264c241..25f7ae9664 100644 --- a/app/team.go +++ b/app/team.go @@ -417,7 +417,7 @@ func (a *App) sendTeamEvent(team *model.Team, event string) *model.AppError { // in case of update_team event - we send the message only to members of that team teamID = team.Id } - message := model.NewWebSocketEvent(event, teamID, "", "", nil) + message := model.NewWebSocketEvent(event, teamID, "", "", nil, "") teamJSON, jsonErr := json.Marshal(team) if jsonErr != nil { return model.NewAppError("sendTeamEvent", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -568,7 +568,7 @@ func (a *App) UpdateTeamMemberSchemeRoles(teamID string, userID string, isScheme } func (a *App) sendUpdatedMemberRoleEvent(userID string, member *model.TeamMember) *model.AppError { - message := model.NewWebSocketEvent(model.WebsocketEventMemberroleUpdated, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventMemberroleUpdated, "", "", userID, nil, "") tmJSON, jsonErr := json.Marshal(member) if jsonErr != nil { return model.NewAppError("sendUpdatedMemberRoleEvent", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -859,7 +859,7 @@ func (a *App) JoinUserToTeam(c request.CTX, team *model.Team, user *model.User, }) } - message := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, "", "", user.Id, nil) + message := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, "", "", user.Id, nil, "") message.Add("team_id", team.Id) message.Add("user_id", user.Id) a.Publish(message) @@ -1084,7 +1084,7 @@ func (a *App) AddTeamMember(c request.CTX, teamID, userID string) (*model.TeamMe return nil, err } - message := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, "", "", userID, nil, "") message.Add("team_id", teamID) message.Add("user_id", userID) a.Publish(message) @@ -1113,7 +1113,7 @@ func (a *App) AddTeamMembers(c *request.Context, teamID string, userIDs []string Member: teamMember, }) - message := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, "", "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventAddedToTeam, "", "", userID, nil, "") message.Add("team_id", teamID) message.Add("user_id", userID) a.Publish(message) @@ -2115,7 +2115,7 @@ func (a *App) ClearTeamMembersCache(teamID string) error { for _, teamMember := range teamMembers { a.ClearSessionCacheForUser(teamMember.UserId) - message := model.NewWebSocketEvent(model.WebsocketEventMemberroleUpdated, "", "", teamMember.UserId, nil) + message := model.NewWebSocketEvent(model.WebsocketEventMemberroleUpdated, "", "", teamMember.UserId, nil, "") tmJSON, jsonErr := json.Marshal(teamMember) if jsonErr != nil { return jsonErr diff --git a/app/teams/teams.go b/app/teams/teams.go index db9cd7c90a..af0ef2bccf 100644 --- a/app/teams/teams.go +++ b/app/teams/teams.go @@ -192,7 +192,7 @@ func (ts *TeamService) JoinUserToTeam(team *model.Team, user *model.User) (*mode // RemoveTeamMember removes the team member from the team. This method sends // the websocket message before actually removing so the user being removed gets it. func (ts *TeamService) RemoveTeamMember(teamMember *model.TeamMember) error { - message := model.NewWebSocketEvent(model.WebsocketEventLeaveTeam, teamMember.TeamId, "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventLeaveTeam, teamMember.TeamId, "", "", nil, "") message.Add("user_id", teamMember.UserId) message.Add("team_id", teamMember.TeamId) ts.wh.Publish(message) diff --git a/app/user.go b/app/user.go index 85accfa404..a9fed62539 100644 --- a/app/user.go +++ b/app/user.go @@ -291,7 +291,7 @@ func (a *App) createUserOrGuest(c request.CTX, user *model.User, guest bool) (*m go a.UpdateViewedProductNoticesForNewUser(ruser.Id) // This message goes to everyone, so the teamID, channelID and userID are irrelevant - message := model.NewWebSocketEvent(model.WebsocketEventNewUser, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventNewUser, "", "", "", nil, "") message.Add("user_id", ruser.Id) a.Publish(message) @@ -785,7 +785,7 @@ func (a *App) SetDefaultProfileImage(c request.CTX, user *model.User) *model.App options := a.Config().GetSanitizeOptions() updatedUser.SanitizeProfile(options) - message := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", nil, "") message.Add("user", updatedUser) a.Publish(message) @@ -981,7 +981,7 @@ func (a *App) DeactivateGuests(c *request.Context) *model.AppError { a.Srv().Store.Channel().ClearCaches() a.Srv().Store.User().ClearCaches() - message := model.NewWebSocketEvent(model.WebsocketEventGuestsDeactivated, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventGuestsDeactivated, "", "", "", nil, "") a.Publish(message) return nil @@ -1078,19 +1078,19 @@ func (a *App) sendUpdatedUserEvent(user model.User) { unsanitizedCopyOfUser := user.DeepCopy() a.SanitizeProfile(adminCopyOfUser, true) - adminMessage := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", omitUsers) + adminMessage := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", omitUsers, "") adminMessage.Add("user", adminCopyOfUser) adminMessage.GetBroadcast().ContainsSensitiveData = true a.Publish(adminMessage) a.SanitizeProfile(&user, false) - message := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", omitUsers) + message := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", omitUsers, "") message.Add("user", &user) message.GetBroadcast().ContainsSanitizedData = true a.Publish(message) // send unsanitized user to event creator - sourceUserMessage := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", unsanitizedCopyOfUser.Id, nil) + sourceUserMessage := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", unsanitizedCopyOfUser.Id, nil, "") sourceUserMessage.Add("user", unsanitizedCopyOfUser) a.Publish(sourceUserMessage) } @@ -1526,7 +1526,7 @@ func (a *App) UpdateUserRolesWithUser(c request.CTX, user *model.User, newRoles a.ClearSessionCacheForUser(user.Id) if sendWebSocketEvent { - message := model.NewWebSocketEvent(model.WebsocketEventUserRoleUpdated, "", "", user.Id, nil) + message := model.NewWebSocketEvent(model.WebsocketEventUserRoleUpdated, "", "", user.Id, nil, "") message.Add("user_id", user.Id) message.Add("roles", newRoles) a.Publish(message) @@ -2225,7 +2225,7 @@ func (a *App) PromoteGuestToUser(c *request.Context, user *model.User, requestor for _, member := range channelMembers { a.invalidateCacheForChannelMembers(member.ChannelId) - evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", user.Id, nil) + evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", user.Id, nil, "") memberJSON, jsonErr := json.Marshal(member) if jsonErr != nil { return model.NewAppError("PromoteGuestToUser", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -2270,7 +2270,7 @@ func (a *App) DemoteUserToGuest(c request.CTX, user *model.User) *model.AppError for _, member := range channelMembers { a.invalidateCacheForChannelMembers(member.ChannelId) - evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", user.Id, nil) + evt := model.NewWebSocketEvent(model.WebsocketEventChannelMemberUpdated, "", "", user.Id, nil, "") memberJSON, jsonErr := json.Marshal(member) if jsonErr != nil { return model.NewAppError("DemoteUserToGuest", "api.marshal_error", nil, "", http.StatusInternalServerError).Wrap(jsonErr) @@ -2288,7 +2288,7 @@ func (a *App) PublishUserTyping(userID, channelID, parentId string) *model.AppEr omitUsers := make(map[string]bool, 1) omitUsers[userID] = true - event := model.NewWebSocketEvent(model.WebsocketEventTyping, "", channelID, "", omitUsers) + event := model.NewWebSocketEvent(model.WebsocketEventTyping, "", channelID, "", omitUsers, "") event.Add("parent_id", parentId) event.Add("user_id", userID) a.Publish(event) @@ -2309,7 +2309,7 @@ func (a *App) invalidateUserCacheAndPublish(userID string) { options := a.Config().GetSanitizeOptions() user.SanitizeProfile(options) - message := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", nil) + message := model.NewWebSocketEvent(model.WebsocketEventUserUpdated, "", "", "", nil, "") message.Add("user", user) a.Publish(message) } @@ -2467,7 +2467,7 @@ func (a *App) UpdateThreadsReadForUser(userID, teamID string) *model.AppError { if nErr != nil { return model.NewAppError("UpdateThreadsReadForUser", "app.user.update_threads_read_for_user.app_error", nil, "", http.StatusInternalServerError).Wrap(nErr) } - message := model.NewWebSocketEvent(model.WebsocketEventThreadReadChanged, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventThreadReadChanged, teamID, "", userID, nil, "") a.Publish(message) return nil } @@ -2492,7 +2492,7 @@ func (a *App) UpdateThreadFollowForUser(userID, teamID, threadID string, state b if thread != nil { replyCount = thread.ReplyCount } - message := model.NewWebSocketEvent(model.WebsocketEventThreadFollowChanged, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventThreadFollowChanged, teamID, "", userID, nil, "") message.Add("thread_id", threadID) message.Add("state", state) message.Add("reply_count", replyCount) @@ -2531,7 +2531,7 @@ func (a *App) UpdateThreadFollowForUserFromChannelAdd(c request.CTX, userID, tea return model.NewAppError("UpdateThreadFollowForUserFromChannelAdd", "app.user.update_thread_follow_for_user.app_error", nil, "", http.StatusInternalServerError).Wrap(err) } - message := model.NewWebSocketEvent(model.WebsocketEventThreadUpdated, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventThreadUpdated, teamID, "", userID, nil, "") userThread, err := a.Srv().Store.Thread().GetThreadForUser(teamID, tm, true) if err != nil { var errNotFound *store.ErrNotFound @@ -2623,7 +2623,7 @@ func (a *App) UpdateThreadReadForUser(c request.CTX, currentSessionId, userID, t a.clearPushNotification(currentSessionId, userID, post.ChannelId, threadID) } - message := model.NewWebSocketEvent(model.WebsocketEventThreadReadChanged, teamID, "", userID, nil) + message := model.NewWebSocketEvent(model.WebsocketEventThreadReadChanged, teamID, "", userID, nil, "") message.Add("thread_id", threadID) message.Add("timestamp", timestamp) message.Add("unread_mentions", membership.UnreadMentions) diff --git a/app/web_conn.go b/app/web_conn.go index 9970d60dab..198e0edbd0 100644 --- a/app/web_conn.go +++ b/app/web_conn.go @@ -662,7 +662,7 @@ func (wc *WebConn) IsAuthenticated() bool { } func (wc *WebConn) createHelloMessage() *model.WebSocketEvent { - msg := model.NewWebSocketEvent(model.WebsocketEventHello, "", "", wc.UserId, nil) + msg := model.NewWebSocketEvent(model.WebsocketEventHello, "", "", wc.UserId, nil, "") msg.Add("server_version", fmt.Sprintf("%v.%v.%v.%v", model.CurrentVersion, model.BuildNumber, wc.App.ClientConfigHash(), @@ -749,6 +749,11 @@ func (wc *WebConn) shouldSendEvent(msg *model.WebSocketEvent) bool { return wc.GetConnectionID() == msg.GetBroadcast().ConnectionId } + // if the connection is omitted don't send the message + if wc.GetConnectionID() == msg.GetBroadcast().OmitConnectionId { + return false + } + // If the event is destined to a specific user if msg.GetBroadcast().UserId != "" { return wc.UserId == msg.GetBroadcast().UserId diff --git a/app/web_conn_test.go b/app/web_conn_test.go index 67c9306b7c..a28673e6e7 100644 --- a/app/web_conn_test.go +++ b/app/web_conn_test.go @@ -122,10 +122,11 @@ func TestWebConnShouldSendEvent(t *testing.T) { {"should only send to admin", &model.WebsocketBroadcast{ContainsSensitiveData: true}, false, false, true, false}, {"should only send to non-admins", &model.WebsocketBroadcast{ContainsSanitizedData: true}, true, true, false, true}, {"should send to nobody", &model.WebsocketBroadcast{ContainsSensitiveData: true, ContainsSanitizedData: true}, false, false, false, false}, + {"should omit basic user 2 by connection id", &model.WebsocketBroadcast{OmitConnectionId: user2ConnID}, true, false, true, true}, // needs more cases to get full coverage } - event := model.NewWebSocketEvent("some_event", "", "", "", nil) + event := model.NewWebSocketEvent("some_event", "", "", "", nil, "") for _, c := range cases { t.Run(c.Description, func(t *testing.T) { event = event.SetBroadcast(c.Broadcast) @@ -178,11 +179,11 @@ func TestWebConnShouldSendEvent(t *testing.T) { assert.True(t, adminUserWc.shouldSendEvent(event), "expected admin") }) - event2 := model.NewWebSocketEvent(model.WebsocketEventUpdateTeam, th.BasicTeam.Id, "", "", nil) + event2 := model.NewWebSocketEvent(model.WebsocketEventUpdateTeam, th.BasicTeam.Id, "", "", nil, "") assert.True(t, basicUserWc.shouldSendEvent(event2)) assert.True(t, basicUser2Wc.shouldSendEvent(event2)) - event3 := model.NewWebSocketEvent(model.WebsocketEventUpdateTeam, "wrongId", "", "", nil) + event3 := model.NewWebSocketEvent(model.WebsocketEventUpdateTeam, "wrongId", "", "", nil, "") assert.False(t, basicUserWc.shouldSendEvent(event3)) } @@ -368,7 +369,7 @@ func TestWebConnDrainDeadQueue(t *testing.T) { defer wc.WebSocket.Close() for i := 0; i < limit; i++ { - msg := model.NewWebSocketEvent("", "", "", "", map[string]bool{}) + msg := model.NewWebSocketEvent("", "", "", "", map[string]bool{}, "") msg = msg.SetSequence(int64(i)) wc.addToDeadQueue(msg) } diff --git a/app/web_hub_test.go b/app/web_hub_test.go index ee14d43fd2..f4f347cb7f 100644 --- a/app/web_hub_test.go +++ b/app/web_hub_test.go @@ -102,7 +102,7 @@ func TestHubStopRaceCondition(t *testing.T) { hub.UpdateActivity("userId", "sessionToken", 0) for i := 0; i <= broadcastQueueSize; i++ { - hub.Broadcast(model.NewWebSocketEvent("", "", "", "", nil)) + hub.Broadcast(model.NewWebSocketEvent("", "", "", "", nil, "")) } hub.InvalidateUser("userId") @@ -195,7 +195,7 @@ func TestHubSessionRevokeRace(t *testing.T) { go func() { for i := 0; i <= broadcastQueueSize; i++ { - hub.Broadcast(model.NewWebSocketEvent("", "teamID", "", "", nil)) + hub.Broadcast(model.NewWebSocketEvent("", "teamID", "", "", nil, "")) } close(done) }() @@ -404,10 +404,10 @@ func TestReliableWebSocketSend(t *testing.T) { th := SetupWithClusterMock(t, testCluster) defer th.TearDown() - ev := model.NewWebSocketEvent("test_unreliable_event", "", "", "", nil) + ev := model.NewWebSocketEvent("test_unreliable_event", "", "", "", nil, "") ev = ev.SetBroadcast(&model.WebsocketBroadcast{}) th.App.Publish(ev) - ev2 := model.NewWebSocketEvent("test_reliable_event", "", "", "", nil) + ev2 := model.NewWebSocketEvent("test_reliable_event", "", "", "", nil, "") ev2 = ev2.SetBroadcast(&model.WebsocketBroadcast{ ReliableClusterSend: true, }) diff --git a/app/webhub_fuzz.go b/app/webhub_fuzz.go index a305b3cfc0..35de30fb35 100644 --- a/app/webhub_fuzz.go +++ b/app/webhub_fuzz.go @@ -229,7 +229,7 @@ func Fuzz(data []byte) int { msg := model.NewWebSocketEvent(input.event, input.selectTeamID, input.selectChannelID, - input.createUserID, nil) + input.createUserID, nil, "") for k, v := range input.attachment { msg.Add(k, v) } diff --git a/model/websocket_message.go b/model/websocket_message.go index fcb0416ed7..41c802dbfc 100644 --- a/model/websocket_message.go +++ b/model/websocket_message.go @@ -86,11 +86,12 @@ type WebSocketMessage interface { } type WebsocketBroadcast struct { - OmitUsers map[string]bool `json:"omit_users"` // broadcast is omitted for users listed here - UserId string `json:"user_id"` // broadcast only occurs for this user - ChannelId string `json:"channel_id"` // broadcast only occurs for users in this channel - TeamId string `json:"team_id"` // broadcast only occurs for users in this team - ConnectionId string `json:"connection_id"` // broadcast only occurs for this connection + OmitUsers map[string]bool `json:"omit_users"` // broadcast is omitted for users listed here + UserId string `json:"user_id"` // broadcast only occurs for this user + ChannelId string `json:"channel_id"` // broadcast only occurs for users in this channel + TeamId string `json:"team_id"` // broadcast only occurs for users in this team + ConnectionId string `json:"connection_id"` // broadcast only occurs for this connection + OmitConnectionId string `json:"omit_connection_id"` // broadcast is omitted for this connection ContainsSanitizedData bool `json:"-"` ContainsSensitiveData bool `json:"-"` // ReliableClusterSend indicates whether or not the message should @@ -113,6 +114,7 @@ func (wb *WebsocketBroadcast) copy() *WebsocketBroadcast { c.UserId = wb.UserId c.ChannelId = wb.ChannelId c.TeamId = wb.TeamId + c.OmitConnectionId = wb.OmitConnectionId c.ContainsSanitizedData = wb.ContainsSanitizedData c.ContainsSensitiveData = wb.ContainsSensitiveData @@ -185,15 +187,16 @@ func (ev *WebSocketEvent) Add(key string, value any) { ev.data[key] = value } -func NewWebSocketEvent(event, teamId, channelId, userId string, omitUsers map[string]bool) *WebSocketEvent { +func NewWebSocketEvent(event, teamId, channelId, userId string, omitUsers map[string]bool, omitConnectionId string) *WebSocketEvent { return &WebSocketEvent{ event: event, data: make(map[string]any), broadcast: &WebsocketBroadcast{ - TeamId: teamId, - ChannelId: channelId, - UserId: userId, - OmitUsers: omitUsers}, + TeamId: teamId, + ChannelId: channelId, + UserId: userId, + OmitUsers: omitUsers, + OmitConnectionId: omitConnectionId}, } } diff --git a/model/websocket_message_test.go b/model/websocket_message_test.go index 95cad01d95..14c49fa3ff 100644 --- a/model/websocket_message_test.go +++ b/model/websocket_message_test.go @@ -13,7 +13,7 @@ import ( func TestWebSocketEvent(t *testing.T) { userId := NewId() - m := NewWebSocketEvent("some_event", NewId(), NewId(), userId, nil) + m := NewWebSocketEvent("some_event", NewId(), NewId(), userId, nil, "") m.Add("RootId", NewId()) user := &User{ Id: userId, @@ -32,7 +32,7 @@ func TestWebSocketEvent(t *testing.T) { } func TestWebSocketEventImmutable(t *testing.T) { - m := NewWebSocketEvent("some_event", NewId(), NewId(), NewId(), nil) + m := NewWebSocketEvent("some_event", NewId(), NewId(), NewId(), nil, "") new := m.SetEvent("new_event") if new == m { @@ -111,7 +111,7 @@ func TestWebSocketResponse(t *testing.T) { } func TestWebSocketEvent_PrecomputeJSON(t *testing.T) { - event := NewWebSocketEvent(WebsocketEventPosted, "foo", "bar", "baz", nil) + event := NewWebSocketEvent(WebsocketEventPosted, "foo", "bar", "baz", nil, "") event = event.SetSequence(7) before, err := event.ToJSON() @@ -126,7 +126,7 @@ func TestWebSocketEvent_PrecomputeJSON(t *testing.T) { var stringSink []byte func BenchmarkWebSocketEvent_ToJSON(b *testing.B) { - event := NewWebSocketEvent(WebsocketEventPosted, "foo", "bar", "baz", nil) + event := NewWebSocketEvent(WebsocketEventPosted, "foo", "bar", "baz", nil, "") for i := 0; i < 100; i++ { event.GetData()[NewId()] = NewId() } @@ -217,7 +217,7 @@ func TestWebSocketEventDeepCopy(t *testing.T) { ContainsSensitiveData: true, } - ev := NewWebSocketEvent("test", "team", "channel", "user", omitUsers) + ev := NewWebSocketEvent("test", "team", "channel", "user", omitUsers, "") ev.Add("post", &Post{}) ev.SetBroadcast(broadcast)