mirror of
https://github.com/penpot/penpot.git
synced 2025-08-06 09:28:26 +02:00
✨ Adds support to rx streams on workers framework
This commit is contained in:
parent
b648fb7446
commit
1a70071405
4 changed files with 76 additions and 51 deletions
|
@ -14,14 +14,25 @@
|
|||
(declare handle-response)
|
||||
(defrecord Worker [instance stream])
|
||||
|
||||
(defn- send-message! [worker {sender-id :sender-id :as message}]
|
||||
(let [data (t/encode message)
|
||||
instance (:instance worker)]
|
||||
(.postMessage instance data)
|
||||
(->> (:stream worker)
|
||||
(rx/filter #(= (:reply-to %) sender-id))
|
||||
(rx/take 1)
|
||||
(rx/map handle-response))))
|
||||
(defn- send-message!
|
||||
([worker message]
|
||||
(send-message! worker message nil))
|
||||
|
||||
([worker {sender-id :sender-id :as message} {:keys [many?] :or {many? false}}]
|
||||
(let [take-messages
|
||||
(fn [ob]
|
||||
(if many?
|
||||
(rx/take-while #(not (:completed %)) ob)
|
||||
(rx/take 1 ob)))
|
||||
|
||||
data (t/encode message)
|
||||
instance (:instance worker)]
|
||||
|
||||
(.postMessage instance data)
|
||||
(->> (:stream worker)
|
||||
(rx/filter #(= (:reply-to %) sender-id))
|
||||
(take-messages)
|
||||
(rx/map handle-response)))))
|
||||
|
||||
(defn ask!
|
||||
[worker message]
|
||||
|
@ -30,6 +41,14 @@
|
|||
{:sender-id (uuid/next)
|
||||
:payload message}))
|
||||
|
||||
(defn ask-many!
|
||||
[worker message]
|
||||
(send-message!
|
||||
worker
|
||||
{:sender-id (uuid/next)
|
||||
:payload message}
|
||||
{:many? true}))
|
||||
|
||||
(defn ask-buffered!
|
||||
[worker message]
|
||||
(send-message!
|
||||
|
|
Loading…
Add table
Add a link
Reference in a new issue