diff --git a/backend/src/app/http/websocket.clj b/backend/src/app/http/websocket.clj index a9631be1d3..fbe6dc3b1f 100644 --- a/backend/src/app/http/websocket.clj +++ b/backend/src/app/http/websocket.clj @@ -113,6 +113,7 @@ (let [psub (::profile-subscription @state) fsub (::file-subscription @state) tsub (::team-subscription @state) + osub (::organization-subscription @state) msg {:type :disconnect :profile-id profile-id :session-id session-id}] @@ -127,6 +128,11 @@ (sp/close! ch) (mbus/purge! msgbus [ch])) + ;; Close organization subscription if exists + (when-let [ch (:channel osub)] + (sp/close! ch) + (mbus/purge! msgbus [ch])) + ;; Close file subscription if exists (when-let [{:keys [topic channel]} fsub] (sp/close! channel) @@ -152,6 +158,29 @@ (sp/close! ch) (mbus/purge! msgbus [ch])))) +(defmethod handle-message :subscribe-organization + [{:keys [::mbus/msgbus]} {:keys [::ws/id ::ws/state ::ws/output-ch ::session-id]} {:keys [organization-id] :as params}] + (l/trace :fn "handle-message" :event "subscribe-organization" :organization-id organization-id :conn-id id) + (let [prev-subs (get @state ::organization-subscription) + channel (sp/chan :buf (sp/dropping-buffer 64) + :xf (remove #(= (:session-id %) session-id)))] + (sp/pipe channel output-ch false) + (mbus/sub! msgbus :topic organization-id :chan channel) + (let [subs {:organization-id organization-id :channel channel :topic organization-id}] + (swap! state assoc ::organization-subscription subs)) + ;; Close previous subscription if exists (e.g. on team switch) + (when-let [ch (:channel prev-subs)] + (sp/close! ch) + (mbus/purge! msgbus [ch])))) + +(defmethod handle-message :unsubscribe-organization + [{:keys [::mbus/msgbus]} {:keys [::ws/id ::ws/state ::session-id]} {:keys [organization-id] :as params}] + (l/trace :fn "handle-message" :event "unsubscribe-organization" :organization-id organization-id :conn-id id) + (when-let [osub (::organization-subscription @state)] + (sp/close! (:channel osub)) + (mbus/purge! msgbus [(:channel osub)]) + (swap! state dissoc ::organization-subscription))) + (defmethod handle-message :subscribe-file [{:keys [::mbus/msgbus ::db/pool]} {:keys [::ws/id ::ws/state ::ws/output-ch ::session-id ::profile-id]} {:keys [file-id] :as params}] 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/rpc_management_nitrate_test.clj b/backend/test/backend_tests/rpc_management_nitrate_test.clj index 7857bdc1a5..f396b9b863 100644 --- a/backend/test/backend_tests/rpc_management_nitrate_test.clj +++ b/backend/test/backend_tests/rpc_management_nitrate_test.clj @@ -242,7 +242,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))) @@ -454,7 +454,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)))) @@ -463,6 +463,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)] @@ -570,7 +588,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 4d9619d630..5007e5ee11 100644 --- a/frontend/src/app/main/data/dashboard.cljs +++ b/frontend/src/app/main/data/dashboard.cljs @@ -51,17 +51,28 @@ ptk/WatchEvent (watch [_ state stream] (let [stopper (rx/filter (ptk/type? ::finalize) stream) - profile-id (:profile-id state)] + profile-id (:profile-id state) + current-team (dm/get-in state [:teams team-id]) + organization-id (dm/get-in current-team [:organization :id]) + + subscriptions (cond-> [{:type :subscribe-team :team-id team-id}] + (some? organization-id) + (conj {:type :subscribe-organization :organization-id organization-id}))] (->> (rx/merge (rx/of (fetch-projects team-id) (df/fetch-fonts team-id)) + (->> (rx/from subscriptions) + (rx/map dws/send)) (->> 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..4c0954c5df 100644 --- a/frontend/src/app/main/data/workspace/notifications.cljs +++ b/frontend/src/app/main/data/workspace/notifications.cljs @@ -52,12 +52,16 @@ (watch [_ state stream] (let [stopper (rx/filter (ptk/type? ::finalize) stream) profile-id (:profile-id state) + current-team (dm/get-in state [:teams team-id]) + organization-id (dm/get-in current-team [:organization :id]) - initmsg [{:type :subscribe-file - :file-id file-id - :version (obj/get global "penpotVersion")} - {:type :subscribe-team - :team-id team-id}] + initmsg (cond-> [{:type :subscribe-file + :file-id file-id + :version (obj/get global "penpotVersion")} + {:type :subscribe-team + :team-id team-id}] + (some? organization-id) + (conj {:type :subscribe-organization :organization-id organization-id})) endmsg {:type :unsubscribe-file :file-id file-id} @@ -75,7 +79,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..07188dc5e3 100644 --- a/frontend/test/frontend_tests/data/dashboard_test.cljs +++ b/frontend/test/frontend_tests/data/dashboard_test.cljs @@ -9,6 +9,8 @@ [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] @@ -82,3 +84,49 @@ (fn [] (done'))))) done)))) + +(t/deftest dashboard-initializes-with-team-and-organization-subscriptions + (t/async done + (let [team-id (uuid/next) + org-id (uuid/next) + state {:profile-id (uuid/next) + :current-team-id team-id + :teams {team-id {:id team-id :organization {:id org-id}}}} + events (atom []) + event (dd/initialize team-id)] + (mock/with-mocks + {dws/send (mock/stub (fn [msg] (swap! events conj msg)))} + (fn [done'] + (->> (ptk/watch event state (rx/empty)) + (rx/reduce conj []) + (rx/subs! + (fn [_]) + (fn [error] (t/is false (str error)) (done')) + (fn [] + (t/is (some #(= :subscribe-team (:type %)) @events)) + (t/is (some #(and (= :subscribe-organization (:type %)) + (= org-id (:organization-id %))) @events)) + (done'))))) + done)))) + +(t/deftest dashboard-without-organization-only-subscribes-to-team + (t/async done + (let [team-id (uuid/next) + state {:profile-id (uuid/next) + :current-team-id team-id + :teams {team-id {:id team-id}}} + events (atom []) + event (dd/initialize team-id)] + (mock/with-mocks + {dws/send (mock/stub (fn [msg] (swap! events conj msg)))} + (fn [done'] + (->> (ptk/watch event state (rx/empty)) + (rx/reduce conj []) + (rx/subs! + (fn [_]) + (fn [error] (t/is false (str error)) (done')) + (fn [] + (t/is (some #(= :subscribe-team (:type %)) @events)) + (t/is (not (some #(= :subscribe-organization (:type %))) @events)) + (done'))))) + done))))