From b2c3fd872a707ac30a1f3707f73645df82df6ce7 Mon Sep 17 00:00:00 2001 From: Andrey Antukh Date: Thu, 1 Oct 2026 13:16:12 +0200 Subject: [PATCH] :recycle: Scope organization notifications to specific WebSocket topics (#11460) * :recycle: Scope organization notifications to specific WebSocket topics Publish team/org notifications to team-id and organization-id topics instead of broadcasting to all connections via uuid/zero. Dashboard and workspace now subscribe to their current team and organization on initialization, receiving only relevant events. Backend: added subscribe-organization/unsubscribe-organization WebSocket handlers and updated :close to clean up organization subscriptions. Modified notify-team-change, notify-organization-deletion, and notify-organization-change-sso to publish to specific topics. Frontend: dashboard and workspace now subscribe to team-id and organization-id (when applicable) on initialization, with nil-guard in topic filters. Closes #11455 AI-assisted-by: longcat-2.0 * :wrench: Remove unused session-id binding in unsubscribe-organization Clj-kondo lint fix. AI-assisted-by: longcat-2.0 * :bug: Fix permission check on team ws connection --------- Co-authored-by: alonso.torres --- backend/src/app/http/websocket.clj | 57 +++++++++++++------ backend/src/app/main.clj | 11 ++-- backend/src/app/rpc/notifications.clj | 12 ++-- .../backend_tests/http_websocket_test.clj | 45 +++++++++++++++ .../rpc_management_nitrate_test.clj | 26 +++++++-- frontend/src/app/main/data/dashboard.cljs | 18 ++++-- .../main/data/workspace/notifications.cljs | 5 +- .../frontend_tests/data/dashboard_test.cljs | 48 ++++++++++++++++ 8 files changed, 185 insertions(+), 37 deletions(-) diff --git a/backend/src/app/http/websocket.clj b/backend/src/app/http/websocket.clj index 18bf25fd8a..50bd64f115 100644 --- a/backend/src/app/http/websocket.clj +++ b/backend/src/app/http/websocket.clj @@ -8,6 +8,7 @@ "A penpot notification service for file cooperative edition." (:require [app.binfile.common :as bfc] + [app.common.data.macros :as dm] [app.common.exceptions :as ex] [app.common.logging :as l] [app.common.pprint :as pp] @@ -18,6 +19,7 @@ [app.http.session :as session] [app.metrics :as mtx] [app.msgbus :as mbus] + [app.nitrate :as nitrate] [app.rpc.commands.files :as files] [app.rpc.commands.teams :as teams] [app.util.websocket :as ws] @@ -42,13 +44,13 @@ (defn repl-get-connections-for-file [file-id] (->> (vals @state) - (filter #(= file-id (-> % deref ::file-subscription :file-id))) + (filter #(= file-id (-> % ::ws/state deref ::file-subscription :file-id))) (map ::ws/id))) (defn repl-get-connections-for-team [team-id] (->> (vals @state) - (filter #(= team-id (-> % deref ::team-subscription :team-id))) + (filter #(= team-id (-> % ::ws/state deref ::team-subscription :team-id))) (map ::ws/id))) (defn repl-close-connection @@ -60,15 +62,17 @@ (defn repl-get-connection-info [id] (when-let [wsp (get @state id)] - {:id id - :created-at (::created-at wsp) - :profile-id (::profile-id wsp) - :session-id (::session-id wsp) - :user-agent (::ws/user-agent wsp) - :ip-addr (::ws/remote-addr wsp) - :last-activity-at (::ws/last-activity-at wsp) - :subscribed-file (-> wsp ::file-subscription :file-id) - :subscribed-team (-> wsp ::team-subscription :team-id)})) + (let [subs (some-> wsp ::ws/state deref)] + {:id id + :created-at (::created-at wsp) + :profile-id (::profile-id wsp) + :session-id (::session-id wsp) + :user-agent (::ws/user-agent wsp) + :ip-addr (::ws/remote-addr wsp) + :last-activity-at (::ws/last-activity-at wsp) + :subscribed-file (-> subs ::file-subscription :file-id) + :subscribed-team (-> subs ::team-subscription :team-id) + :subscribed-org (-> subs ::team-subscription :organization-id)}))) (defn repl-print-connection-info [id] @@ -133,18 +137,39 @@ (mbus/purge! msgbus [channel]) (mbus/pub! msgbus :topic topic :message msg)))) +(defn- get-team-organization-id + "Returns the id of the organization that owns `team-id`, or nil when + the team has no organization or nitrate cannot be reached." + [cfg team-id] + (try + (-> (nitrate/call cfg :get-team-organization {:team-id team-id}) + (dm/get-in [:organization :id])) + (catch Throwable cause + (l/warn :hint "unable to resolve team organization" + :team-id team-id + :cause cause) + nil))) + (defmethod handle-message :subscribe-team [cfg {:keys [::ws/id ::ws/state ::ws/output-ch ::session-id ::profile-id]} {:keys [team-id] :as params}] (l/trace :fn "handle-message" :event "subscribe-team" :team-id team-id :conn-id id) (teams/check-read-permissions! cfg profile-id team-id) - (let [prev-subs (get @state ::team-subscription) - channel (sp/chan :buf (sp/dropping-buffer 64) - :xf (remove #(= (:session-id %) session-id)))] + (let [prev-subs (get @state ::team-subscription) + organization-id (get-team-organization-id cfg team-id) + ;; Resolved server-side so a client only hears its readable team's org + topics (cond-> [team-id] + (some? organization-id) + (conj organization-id)) + channel (sp/chan :buf (sp/dropping-buffer 64) + :xf (remove #(= (:session-id %) session-id)))] (sp/pipe channel output-ch false) - (mbus/sub! (::mbus/msgbus cfg) :topic team-id :chan channel) + (mbus/sub! (::mbus/msgbus cfg) :topics topics :chan channel) - (let [subs {:team-id team-id :channel channel :topic team-id}] + (let [subs {:team-id team-id + :organization-id organization-id + :channel channel + :topic team-id}] (swap! state assoc ::team-subscription subs)) ;; Close previous subscription if exists diff --git a/backend/src/app/main.clj b/backend/src/app/main.clj index de72b9a0ef..d9ce4bf7c2 100644 --- a/backend/src/app/main.clj +++ b/backend/src/app/main.clj @@ -379,11 +379,12 @@ ::setup/props (ig/ref ::setup/props)} ::http.ws/routes - {::db/pool (ig/ref ::db/pool) - ::mtx/metrics (ig/ref ::mtx/metrics) - ::mbus/msgbus (ig/ref ::mbus/msgbus) - ::setup/props (ig/ref ::setup/props) - ::session/manager (ig/ref ::session/manager)} + {::db/pool (ig/ref ::db/pool) + ::mtx/metrics (ig/ref ::mtx/metrics) + ::mbus/msgbus (ig/ref ::mbus/msgbus) + ::setup/props (ig/ref ::setup/props) + ::session/manager (ig/ref ::session/manager) + :app.nitrate/client (ig/ref :app.nitrate/client)} :app.http.assets/routes {::http.assets/path (cf/get :assets-path) diff --git a/backend/src/app/rpc/notifications.clj b/backend/src/app/rpc/notifications.clj index ad68556493..00d6b8548e 100644 --- a/backend/src/app/rpc/notifications.clj +++ b/backend/src/app/rpc/notifications.clj @@ -6,16 +6,14 @@ (ns app.rpc.notifications (:require - [app.common.uuid :as uuid] [app.msgbus :as mbus])) (defn notify-team-change [cfg team notification] - (let [msgbus (::mbus/msgbus cfg)] + (let [msgbus (::mbus/msgbus cfg) + team-id (:id team)] (mbus/pub! msgbus - ;;TODO There is a bug on dashboard with teams notifications. - ;;For now we send it to uuid/zero instead of team-id - :topic uuid/zero + :topic team-id :message {:type :team-organization-change :team team :notification notification}))) @@ -37,7 +35,7 @@ [cfg organization-id organization-name teams deleted-teams] (let [msgbus (::mbus/msgbus cfg)] (mbus/pub! msgbus - :topic uuid/zero + :topic organization-id :message {:type :organization-deleted :organization-id organization-id :organization-name organization-name @@ -48,6 +46,6 @@ [cfg organization-id] (let [msgbus (::mbus/msgbus cfg)] (mbus/pub! msgbus - :topic uuid/zero + :topic organization-id :message {:type :organization-change-sso :organization-id organization-id}))) diff --git a/backend/test/backend_tests/http_websocket_test.clj b/backend/test/backend_tests/http_websocket_test.clj index 22683e93c2..6ffc10e744 100644 --- a/backend/test/backend_tests/http_websocket_test.clj +++ b/backend/test/backend_tests/http_websocket_test.clj @@ -10,6 +10,7 @@ [app.db :as db] [app.http.websocket :as ws] [app.msgbus :as mbus] + [app.nitrate :as nitrate] [app.rpc :as-alias rpc] [app.rpc.commands.files :as files] [app.rpc.commands.teams :as teams] @@ -68,6 +69,50 @@ (t/testing "permission check passes for authorized user" (t/is (nil? (teams/check-read-permissions! cfg (:id profile1) (:id team))))))) +(defn- subscribed-topics + "Runs :subscribe-team for `profile-id` on `team-id` with `nitrate-call` + standing in for nitrate; returns the topics and the stored subscription." + [profile-id team-id nitrate-call] + (let [state (atom {}) + output-ch (sp/chan :buf (sp/dropping-buffer 64)) + wsp (make-wsp profile-id state output-ch) + calls (atom [])] + (with-redefs [nitrate/call nitrate-call + mbus/sub! (fn [_ & {:keys [topics]}] + (swap! calls conj topics))] + ((get-method ws/handle-message :subscribe-team) + th/*system* wsp {:team-id team-id})) + (some-> @state ::ws/team-subscription :channel sp/close!) + {:topics @calls + :subscription (::ws/team-subscription @state)})) + +(t/deftest subscribe-team-subscribes-to-team-organization + (let [profile (th/create-profile* 1 {:is-active true}) + team-id (:id (th/create-team* 1 {:profile-id (:id profile)})) + org-id (uuid/next)] + + (t/testing "adds the organization topic when the team has one" + (let [{:keys [topics subscription]} + (subscribed-topics (:id profile) team-id + (fn [_ method params] + (when (= :get-team-organization method) + {:id (:team-id params) + :organization {:id org-id}})))] + (t/is (= [[team-id org-id]] topics)) + (t/is (= org-id (:organization-id subscription))))) + + (t/testing "subscribes only to the team when it has no organization" + (let [{:keys [topics subscription]} + (subscribed-topics (:id profile) team-id (constantly nil))] + (t/is (= [[team-id]] topics)) + (t/is (nil? (:organization-id subscription))))) + + (t/testing "subscribes to the team when nitrate fails" + (let [{:keys [topics]} + (subscribed-topics (:id profile) team-id + (fn [& _] (throw (ex-info "nitrate down" {}))))] + (t/is (= [[team-id]] topics)))))) + (t/deftest pointer-update-validates-file-id (let [profile (th/create-profile* 1 {:is-active true}) file (th/create-file* 1 {:profile-id (:id profile) diff --git a/backend/test/backend_tests/rpc_management_nitrate_test.clj b/backend/test/backend_tests/rpc_management_nitrate_test.clj index 3aa81b7a2b..a754e45b1c 100644 --- a/backend/test/backend_tests/rpc_management_nitrate_test.clj +++ b/backend/test/backend_tests/rpc_management_nitrate_test.clj @@ -243,7 +243,7 @@ :organization organization}))] (t/is (th/success? out)) (t/is (= 1 (count @calls))) - (t/is (= uuid/zero (-> @calls first :topic))) + (t/is (= team-id (-> @calls first :topic))) (let [msg (-> @calls first :message)] (t/is (= :team-organization-change (:type msg))) (t/is (= nil (:notification msg))) @@ -267,7 +267,7 @@ :organization {:name organization-name}}))] (t/is (th/success? out)) (t/is (= 1 (count @calls))) - (t/is (= uuid/zero (-> @calls first :topic))) + (t/is (= team-id (-> @calls first :topic))) (let [msg (-> @calls first :message)] (t/is (= :team-organization-change (:type msg))) (t/is (= "dashboard.team-no-longer-belong-organization" (:notification msg))) @@ -475,7 +475,7 @@ ;; --- Verify: exactly one organization-deleted event is published on the message bus --- (t/is (:called? @mbus-mock)) (let [msg (apply hash-map (rest (:call-args @mbus-mock)))] - (t/is (= uuid/zero (:topic msg))) + (t/is (= organization-id (:topic msg))) (t/is (= :organization-deleted (:type (:message msg)))) (t/is (= organization-id (:organization-id (:message msg)))) (t/is (= organization-name (:organization-name (:message msg)))) @@ -484,6 +484,24 @@ (t/is (= #{(:id empty-team)} (set (:deleted-teams (:message msg)))))))))) +(t/deftest notify-organization-change-sso-publishes-event + (let [organization-id (uuid/random) + calls (atom []) + out (with-redefs [mbus/pub! (fn [_cfg & {:keys [topic message]}] + (swap! calls conj {:topic topic + :message message}))] + (th/management-command! {::th/type :notify-organization-sso-change + ::rpc/profile-id (uuid/random) + :organization-id organization-id + :updated-props true + :announce-activation false}))] + (t/is (th/success? out)) + (t/is (= 1 (count @calls))) + (t/is (= organization-id (-> @calls first :topic))) + (let [msg (-> @calls first :message)] + (t/is (= :organization-change-sso (:type msg))) + (t/is (= organization-id (:organization-id msg)))))) + (t/deftest notify-user-organizations-deletion-renames-or-deletes-teams-and-publishes-per-organization-events ;; --- Deferred owned-organizations: nil during setup, filled before RPC --- (let [owned-organizations-ref (atom nil)] @@ -591,7 +609,7 @@ ;; --- Verify: one organization-deleted event per organization, all on correct topic --- (t/is (= 2 (count msgs))) - (t/is (every? #(= uuid/zero (:topic %)) + (t/is (every? #(contains? #{organization-1-id organization-2-id} (:topic %)) (->> (:call-args-list @mbus-mock) (map #(apply hash-map (rest %)))))) (t/is (= #{:organization-deleted} (set (map :type msgs)))) diff --git a/frontend/src/app/main/data/dashboard.cljs b/frontend/src/app/main/data/dashboard.cljs index 37d9f9fc44..b0e47eafea 100644 --- a/frontend/src/app/main/data/dashboard.cljs +++ b/frontend/src/app/main/data/dashboard.cljs @@ -50,18 +50,28 @@ (ptk/reify ::initialize ptk/WatchEvent (watch [_ state stream] - (let [stopper (rx/filter (ptk/type? ::finalize) stream) - profile-id (:profile-id state)] + (let [stopper (rx/filter (ptk/type? ::finalize) stream) + profile-id (:profile-id state) + organization-id (dm/get-in state [:teams team-id :organization :id]) + initmsg {:type :subscribe-team :team-id team-id}] (->> (rx/merge (rx/of (fetch-projects team-id) - (df/fetch-fonts team-id)) + (df/fetch-fonts team-id) + (dws/send initmsg)) + ;; On reconnect, send again the subscription message + (->> stream + (rx/filter (ptk/type? ::dws/opened)) + (rx/map #(dws/send initmsg))) (->> stream (rx/filter (ptk/type? ::dws/message)) (rx/map deref) (rx/filter (fn [{:keys [topic] :as msg}] (or (= topic uuid/zero) - (= topic profile-id)))) + (= topic profile-id) + (= topic team-id) + (when (some? organization-id) + (= topic organization-id))))) (rx/map process-message))) (rx/take-until stopper)))))) diff --git a/frontend/src/app/main/data/workspace/notifications.cljs b/frontend/src/app/main/data/workspace/notifications.cljs index 55f9735814..b3eead45ef 100644 --- a/frontend/src/app/main/data/workspace/notifications.cljs +++ b/frontend/src/app/main/data/workspace/notifications.cljs @@ -52,6 +52,7 @@ (watch [_ state stream] (let [stopper (rx/filter (ptk/type? ::finalize) stream) profile-id (:profile-id state) + organization-id (dm/get-in state [:teams team-id :organization :id]) initmsg [{:type :subscribe-file :file-id file-id @@ -75,7 +76,9 @@ (or (= topic uuid/zero) (= topic profile-id) (= topic team-id) - (= topic file-id)))) + (= topic file-id) + (when (some? organization-id) + (= topic organization-id))))) (rx/map process-message)) ;; On reconnect, send again the subscription messages diff --git a/frontend/test/frontend_tests/data/dashboard_test.cljs b/frontend/test/frontend_tests/data/dashboard_test.cljs index 001ea2f924..12fefdb476 100644 --- a/frontend/test/frontend_tests/data/dashboard_test.cljs +++ b/frontend/test/frontend_tests/data/dashboard_test.cljs @@ -9,10 +9,13 @@ [app.common.uuid :as uuid] [app.config :as cf] [app.main.data.common :as dcm] + [app.main.data.dashboard :as dd] + [app.main.data.websocket :as dws] [app.main.repo :as rp] [app.main.router :as rt] [beicon.v2.core :as rx] [cljs.test :as t :include-macros true] + [frontend-tests.helpers.async :as async] [frontend-tests.helpers.mock :as mock] [potok.v2.core :as ptk])) @@ -82,3 +85,48 @@ (fn [] (done'))))) done)))) + +(defn- sent-messages + "Runs the watch of `event` against `stream` with `dws/send` stubbed + and resolves to the messages it sends over the websocket." + [event state stream] + (let [sent (atom [])] + (-> (mock/with-mocks* + {dws/send (mock/stub (fn [msg] (swap! sent conj msg) msg))} + (await (async/observe (ptk/watch event state stream)))) + (.then (fn [_] @sent))))) + +(t/deftest ^:async dashboard-initialize-subscribes-to-team + (let [team-id (uuid/next) + state {:profile-id (uuid/next) + :teams {team-id {:id team-id :organization {:id (uuid/next)}}}} + sent (await (sent-messages (dd/initialize team-id) state (rx/empty)))] + (t/is (= [{:type :subscribe-team :team-id team-id}] sent)))) + +(t/deftest ^:async dashboard-initialize-resubscribes-on-reconnect + (let [team-id (uuid/next) + state {:profile-id (uuid/next) + :teams {team-id {:id team-id}}} + stream (rx/of (ptk/data-event ::dws/opened {}) + (ptk/data-event ::dws/opened {})) + sent (await (sent-messages (dd/initialize team-id) state stream))] + (t/is (= 3 (count sent))) + (t/is (every? #(= {:type :subscribe-team :team-id team-id} %) sent)))) + +(t/deftest ^:async dashboard-initialize-accepts-team-organization-messages + (let [team-id (uuid/next) + org-id (uuid/next) + state {:profile-id (uuid/next) + :teams {team-id {:id team-id :organization {:id org-id}}}} + message (fn [topic] + (ptk/data-event ::dws/message + {:type :organization-change-sso + :topic topic + :organization-id org-id})) + stream (rx/of (message org-id) (message (uuid/next))) + processed (atom [])] + (await (mock/with-mocks* + {dws/send (mock/stub identity)} + (await (async/observe (ptk/watch (dd/initialize team-id) state stream) + :on-next #(swap! processed conj (ptk/type %)))))) + (t/is (= 1 (count (filter #{::dcm/handle-organization-change-sso} @processed))))))