Beginner with missionary here. Looking for a bit of guidance on how to use the library.
I'm trying to learn about how to use missionary by writing my own networking protocol, where I imagine callbacks from netty as values that will go onto continuous flows. So I get interop, callbacks, errors. The works.
Please keep in mind I have no idea what I'm doing. So my nomenclature is probably wrong or confused.
I'm aware that aleph has some netty support, but I'm looking for a learning project, where I want to start from a lower level, instead of throwing ring at the problem.
For documentation support I use the search function here in this slack channel and the missionary documentation. I also found some secret documentation at https://gorgeous-sorbet-a5a2bf.netlify.app/ that was missing an index.
So after much head scratching I discovered m/observe , which enables integration with other libraries by doing some variation of the following:
(def <resource
(m/observe (fn [emit!]
(library-stuff (some-interop emit!))
(fn cleanup []
(do-cleanup)))))
For me this convention takes a while getting used to, where you have a emit! that does all your transfer of values and the return value is the cleanup. But I believe that this is the right thing to use when gathering side effects from non-missionary code.
Most of the netty library deals with creating resources and structuring callbacks for them.
So I can easily make a scenario where I for example try to send a message to a port that is closed, which throws an exception. So I wish for some nice resource management from missionary on that side.
And maybe I want efficiency as well, so some resources I will have to share, instead of constantly re-creating copies. Netty is pretty fast, but when every object needs to be recreated my code slows to a crawl.
let's say I have a channel resource, created with m/observe
A channel is created with a callback. But I want it to go to 100 consumers. I'm imagining this to be a signal, but I only call (emit! channel) once when the netty callback is fired. What is the way/mechanism to share it? I cant fork 100 times because I emitted once. So what I was looking for was something that was "give the most recent value from the stream" or the like.
I could just fork on it a single time, and pass our references to the returned value but then supervision is gone. I would like to fork on it inside the 100 consumers.
If something happens and the channel throws, I want the exception to be raised and the channel to be closed, cancelling all subscribers as well.
So there I am, a bit stuck, trying to distribute a channel resource to multiple consumers, just so I can do some testing.
Thanks for any help.
For the rest of my path, just for contextualizing where I'm trying to go:
I was imagining a resource that represented the io of send, another one for io receive,
Then setting up logic of the protocol, reacting to received messages and sending responses using those signals.
And if any one failed I want to make some recovery logic.Read at least the beginning of https://www.dustingetz.com/#/page/signals%20vs%20streams%2C%20in%20terms%20of%20backpressure%20(2023). Then, based on your new knowledge decide which one you need (I suspect streams). (m/stream (m/observe ..)) returns a stream that manages the observe, i.e. 100 subscribers will get values from the 1 observe.
I've read that a few times already, it's nice, but in terms of linking those words and the library in itself it does very little. What I'm reading from that article is that a stream is blocking for both producer and consumer, and enforces a 1:1 event/read relationship. So In my example I can produce exactly one message with the channel I create from m/observe. But I might be reading it wrong.
Lemme demo some code first and I'll get right back in.
in a stream of events you usually cannot skip an event, therefore m/stream waits for all consumers to read the current event. Contrast to a signal (e.g. current mouse coordinates) where consumers can read independently of each other because they can decide when they need to know the current mouse position
This is my interpretation of creating multiple consumers. It is not yet working, but it illustrates my thinking on concrete code. How would I make the 10 update the atom with their own numbers? Lots of code here:
(require '[missionary.core :as m])
(defn demo-interop
"I always die after 500 ms"
[emit!]
(future
(Thread/sleep 10)
(println "emit")
(emit! (atom [::emit])) ;; some stateful thing
(Thread/sleep 500)
(println "die!")
(emit! (ex-info "ded" {})))
(fn []
(println "end!, cleanup?")))
(def <resource
(m/observe demo-interop))
(comment
;; if you (take 3) it deadlocks before printing "die!" why?
(def q (m/? (m/reduce conj nil (m/eduction (take 2) (m/stream <resource)))))
(count q)
;;=> 2
;; using (m/reduce .. ) to avoid deadlocks in the repl...
;; prints ::ready once. expected.
(let [test-program-spec
(m/reduce
(fn [_ v] (println "got" v) v)
nil
(m/eduction (take 5) (m/stream <resource)))
cancel (test-program-spec (fn [_]) (fn [_])) ;; never really got this
]
(Thread/sleep 510) ;; waiting longer causes the error to appear
(cancel) ;; not cancelling it causes the error to appear.
)
;; lets try to share
(let [test-program-spec
(m/reduce
(fn [_ v] (println "got" v) v)
nil
(m/eduction (take 12)
(let [origin-publisher (m/stream <resource)
;; create multiple consumers of the stream somehow
;; NON working. print statement does nothing.
subs (mapv (fn [n]
;; tried m/stream,
;; what do I do here?
;; (m/stream origin-publisher) ;; nope!
(m/ap
(let [q (m/? origin-publisher)] ;; park?
(m/?> (if (instance? clojure.lang.Atom q)
;; just to do something side-effecty
(do
(swap! q conj n)
q)
q)))))
(range 10))]
(apply m/join vector subs))))
cancel (test-program-spec (fn [_]) (fn [_])) ;; never really got this
]
(Thread/sleep 400)
(cancel))
)If your event stream has continuous time semantics, i.e. only the most recent value matters, then you should use m/signal not m/stream. m/stream has few use cases and is hard to get right. Here is a working example you can improve upon :
(defn demo-interop
"I always die after 500 ms"
[emit!]
(partial future-cancel
(future
(Thread/sleep 10)
(println "emit1")
(emit! :foo)
(Thread/sleep 20)
(println "emit2")
(emit! :bar)
(Thread/sleep 500)
(println "die!")
(emit! nil))))
(def shared-resource
(->> (m/observe demo-interop) ;; your managed resource
(m/eduction (take-while some?)) ;; terminate the event stream on nil
(m/reductions {} nil) ;; turn the event stream into a succession of states
(m/relieve) ;; discard the previous values in case the consumer is slower than the producer
(m/signal))) ;; assign an identity to the pipeline, it is a memoized box that can be shared and reused
(defn observer [i]
(m/reduce (fn [_ x] (println "Observer" i "-" x)) nil shared-resource))
(comment
;; start process with 10 observers
(def cancel
((->> (range 10)
(map observer)
(apply m/join vector))
prn prn))
(cancel)
)Nice, thanks. That is very much cleaner than my code. I'm going to experiment some more with this. It will probably take some time. So for UIs and such signals are definitely the use case. Some quick follow up questions that immediately comes to mind. I can foresee some applications where every event matters, but maybe it's nicer to just use a queue for those? Examples: • Data recording • Transactional events like handshakes What is your take on that?
You can build discrete (i.e. every event matters) dataflow pipelines in missionary, you will get backpressure and supervision by default, and you should go a long way before needing explicit queues.