Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 18 additions & 5 deletions src/ebb/impl/ambiguous.clj
Original file line number Diff line number Diff line change
Expand Up @@ -312,16 +312,29 @@
(defn- cancel [^Process ps]
(when (.-live ps)
(set! (.-live ps) false)
;; THE FLOWS FIRST, THEN THE PARKS. Cancelling a park can resume its
;; branch synchronously on this fiber: the branch fails at its `?` and
;; pumps, and if the flow its fork draws from is still live the pump
;; pulls the next value and forks again, whose park is cancelled on
;; arrival, which pumps again -- and control never comes back to the
;; loop below that would have cancelled the flow. Over an infinite seed
;; that was a process spinning forever inside this function at full CPU
;; (samizdat's supervisor stream, stopped: 88,641 iterations in the two
;; seconds after its cancel, the reduce never settling); over a finite one
;; it drained the whole seed before stopping. With the flows cancelled
;; first the next pull terminates the choice and the pump runs dry.
;; Ambiguous.java cancels its choice ring before its token for the same
;; reason. Pinned by ebb.ap-cancel-seed-test.
(doseq [^Choice ch (.-choices ps)]
(when (.-live ch)
(set! (.-live ch) false)
(when-some [it (.-iterator ch)] (it))))
;; NOT named `parks`: a local sharing a name with a mutable field reads the
;; FIELD, so this snapshot came back as the [] just assigned and every
;; park-cancel was silently skipped. See doc/conformance.md.
(let [v-parks (.-parks ps)]
(set! (.-parks ps) [])
(doseq [c v-parks] (c)))
(doseq [^Choice ch (.-choices ps)]
(when (.-live ch)
(set! (.-live ch) false)
(when-some [it (.-iterator ch)] (it)))))
(doseq [c v-parks] (c))))
nil)

;; ------------------------------------------------------------------- the body
Expand Down
61 changes: 61 additions & 0 deletions test/ebb/ap_cancel_seed_test.clj
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
;; A cancelled ap must stop drawing from the flow it forked on. Found in
;; samizdat (karamazov-3cll.9): a supervisor stream is a reduce over
;; (ap (let [_ (?> (seed (repeat nil)))] (? (sleep poll-ms)) (drain)))
;; and stopping it cancelled the reduce, which cancelled the ap -- whose
;; parks were then cancelled on arrival while the fork kept pulling the next
;; value from an infinite seed. Measured: 88,641 iterations in the two
;; seconds after the cancel, and the reduce never settled. Missionary's
;; Ambiguous cancels the forked flow along with the parks; ebb's did not.
(ns ebb.ap-cancel-seed-test
(:require [clojure.test :as t]
[ebb.core :as m]))

(defn- start [task]
(let [done (promise)
cancel (task (fn [v] (deliver done [:ok v])) (fn [e] (deliver done [:err e])))]
{:done done :cancel cancel}))

(defn- observed
"A flow that counts the cancels reaching it."
[flow counter]
(fn [n t]
(let [ps (flow n t)]
(reify clojure.lang.IFn
(invoke [_] (swap! counter inc) (ps))
clojure.lang.IDeref
(deref [_] @ps)))))

(t/deftest cancelling-an-ap-cancels-the-flow-it-forked-on
(let [seed-cancels (atom 0)
runs (atom 0)
flow (m/ap (let [x (m/?> (observed (m/seed (repeat 1)) seed-cancels))]
(swap! runs inc)
(m/? (m/sleep 5))
x))
{:keys [done cancel]} (start (m/reduce (fn [_ _] nil) nil flow))]
(Thread/sleep 60)
(t/is (pos? @runs) "the flow was running")
(let [before @runs]
(cancel)
(Thread/sleep 100)
(t/is (= 1 @seed-cancels) "the cancel reached the seed the fork draws from")
(let [after @runs]
(Thread/sleep 300)
(t/is (= after @runs) "no further forks once cancelled")
(t/is (< (- after before) 3) "and at most the branch already in flight ran"))
(let [[tag e] (deref done 1000 [:unsettled nil])]
(t/is (= :err tag) "the reduce settled")
(t/is (m/cancelled? e) "with Cancelled")))))

(t/deftest cancelling-an-ap-parked-in-a-fork-terminates-it
;; The same shape without a park in the body: a finite seed the consumer
;; cancels part-way through must also end.
(let [flow (m/ap (let [x (m/?> (m/seed (range 100000)))]
(m/? (m/sleep 1))
x))
{:keys [done cancel]} (start (m/reduce (fn [acc x] (conj acc x)) [] flow))]
(Thread/sleep 30)
(cancel)
(let [[tag e] (deref done 1000 [:unsettled nil])]
(t/is (= :err tag))
(t/is (m/cancelled? e)))))