missionary 2024-11-02

I have a weird bug. I have this fn that is working:

(defn add-formula-cell [dag cell-id formula-fn input-cell-id-vec]
  (assert dag "dag needs to be non nil")
  (assert (vector? input-cell-id-vec) "input-cell-id-vec needs to be a vector")
  (let [input-cells (map #(get-cell-or-throw dag %) input-cell-id-vec)
        ;_ (println "all input cells are good!")
        formula-fn-wrapped (fn [& args]
                             (if (some-input-no-value? args)
                               (create-no-val cell-id)
                               (try
                                 (let [start (. System (nanoTime))
                                       ;`result (calculate dag formula-fn args)
                                       result (m/? (m/via m/cpu (calculate dag formula-fn args)))
                                       stime (str "\r\ncell " cell-id
                                                  " calculated in "
                                                  (/ (double (- (. System (nanoTime)) start)) 1000000.0)
                                                  " msecs")]
                                   (when (:logger dag)
                                     (trace/write-text (:logger dag) stime))
                                   result)
                                 (catch Exception ex
                                   (when (:logger dag)
                                     (trace/write-ex (:logger dag) cell-id ex))
                                   (throw ex)))))
        formula-cell (apply m/latest formula-fn-wrapped input-cells)
        ;formula-cell-wrapped (m/stream formula-cell)
        formula-cell-wrapped (m/signal formula-cell)]
    (add-cell dag cell-id formula-cell-wrapped)))
 
And now instead of running the function on m/latest that is syncronous, I want to rewrite it that it can be asyncronosu (so return a missionary task)
(defn add-formula-cell [dag cell-id formula-fn input-cell-id-vec sp?]
  (assert dag "dag needs to be non nil")
  (assert (vector? input-cell-id-vec) "input-cell-id-vec needs to be a vector")
  (let [input-cells (map #(get-cell-or-throw dag %) input-cell-id-vec)
        ;_ (println "all input cells are good!")
        input-f (apply m/latest vector input-cells)
        formula-result-f (m/ap 
                          (m/amb (create-no-val cell-id))
                          (let [args (seq (m/?> input-f))]
                             (println "args: " args "sp?: " sp?)
                             (if (some-input-no-value? args)
                               (create-no-val cell-id)
                               (try
                                 (let [start (. System (nanoTime))
                                       ;`result (calculate dag formula-fn args)
                                       result (m/? (m/via m/cpu
                                                    (if sp? 
                                                     (m/? (calculate dag formula-fn args))
                                                     (calculate dag formula-fn args))))
                                       stime (str "\r\ncell " cell-id
                                                  " calculated in "
                                                  (/ (double (- (. System (nanoTime)) start)) 1000000.0)
                                                  " msecs")]
                                   (when (:logger dag)
                                     (trace/write-text (:logger dag) stime))
                                   (println "flow result: " result)
                                   result)
                                 (catch Exception ex
                                   (when (:logger dag)
                                     (trace/write-ex (:logger dag) cell-id ex))
                                   (throw ex))))))
        formula-result-f-wrapped (m/signal formula-result-f)
        ]
    (add-cell dag cell-id formula-result-f-wrapped))) 
And now when I wait for this result, It prints "flow result 5"; but the return value of the flow is: [{}, true, false, false, 5] and it is of java.lang.Object type. I believe this is coming from m/signal. The ap process definition is ok. I dont get it. My working example runs a calc-fn in m/latest. And my change is that the m/latest just consolidates the args to a vector, and then it runs a m/ap that does (m/?> on the calc-fn) I didnt think it was such a huge change.

m/signal expects the flow to be initialized, i.e. the process is immediately able to produce a value. If it's not the case, try adding a m/reductions stage

(def d (m/signal (m/reductions {} nil
                   (m/ap (let [v (m/?> a)]
                           (println "args: " v)
                           (m/? (slow-mult v)))))))

Thanks @leonoel are ap events when run with (m/reductions {} nil) then the same thing as if they are generated with cp? I am trying to understand if there is a difference.

The issues arise once I have m/? that has at least m/sleep in it somewehre.

When everything is synchronous, it works.

(defn run! [f]
  (m/? 
    (m/reduce (fn [r v]
                 v
                 ) nil 
    (m/eduction (take 1) f))))

(def a (m/signal (m/latest vector (m/seed [1]) (m/seed [2]))))

(run! a)
;; => [1 2]

(defn mult [[a b]]
  (m/sp (* a b)))

(def c (m/signal (m/ap (let [v (m/?> a)] 
                 (println "args: " v)
                 (m/? (mult v)) 
                 ))))

(run! c)
;; => 2


(defn slow-mult [[a b]]
  (m/sp (m/? (m/sleep 500))
        (* a b)))

(def d (m/signal (m/ap (let [v (m/?> a)]
                         (println "args: " v)
                         (m/? (slow-mult v))))))

(run! d)
;; => [{}, true, true, true, #error {
;;     :cause "Sleep cancelled."
;;     :via
;;     [{:type missionary.Cancelled
;;       :message "Sleep cancelled."}]
;;     :trace
;;     []}]

I think this is the root issue.

I am trying to build a calculation tree.

For cases where certain cells do not have values, but continuous process requires them, I have created a empty value record:

(defrecord no-val [cell-id])

(defn create-no-val [cell-id]
  (no-val. cell-id))

(defn is-no-val? [v]
  (instance? no-val v))

And when I query cell values, I filter out empty values:

(defn take-first-val [f]
  ; flows dont implement deref
  (m/eduction
   (remove is-no-val?)
   (take 1)
   f))

(defn current-v
  "gets the first valid value from the flow"
  [f]
  (m/reduce (fn [r v]
              (println "current v: " v " r: " r)
              v) nil
            (take-first-val f)))

Asynchronous operators (i.e. ? and ?>) are not allowed in cp. Therefore, cp flows are always initialized, so they can be passed to latest or signal without a reductions stage. cp also leverages the non-asynchronous property to defer computation on consumer demand, whereas ap computes ASAP. My current understanding is that ap and cp are actually the same thing, cp is just the special case of ap with no asynchronous operators and the continuous time semantics can be inferred from the static structure of the code, more details here https://github.com/leonoel/missionary/issues/109

Thank you so much for this explanation @leonoel. Your documentation is pretty good, but on this difference it was quite minimalistic. Thanks!!!