From 209b1862cc2d935e76b54f7f3d084f544b72418c Mon Sep 17 00:00:00 2001 From: Yogthos Date: Sun, 6 Sep 2026 21:28:29 -0400 Subject: [PATCH] fix(ap): cancel the flows a fork draws from before its parks Cancelling an ap cancelled its pending parks first and its forked flows second. Cancelling a park can resume the branch synchronously on the owner fiber: the branch fails at its `?` and pumps, the pump pulls the next value from a flow that is still live and forks again, the new park is cancelled on arrival, and control never reaches the loop that would have cancelled the flow. Over an infinite seed the process spun forever inside `cancel` at full CPU and the consumer never settled; over a finite seed it drained the whole seed before stopping. Found through samizdat's supervisor stream, a reduce over (ap (let [_ (?> (seed (repeat nil)))] (? (sleep ms)) (drain))): stopping it measured 88,641 iterations in the two seconds after the cancel. Cancel the flows first, then the parks, as Ambiguous.java cancels its choice ring before its token. With the flow cancelled, the next pull terminates the choice and the pump runs dry. Pinned by ebb.ap-cancel-seed-test: the cancel reaches the seed, no further fork runs, and the reduce settles with Cancelled. --- src/ebb/impl/ambiguous.clj | 23 +++++++++--- test/ebb/ap_cancel_seed_test.clj | 61 ++++++++++++++++++++++++++++++++ 2 files changed, 79 insertions(+), 5 deletions(-) create mode 100644 test/ebb/ap_cancel_seed_test.clj diff --git a/src/ebb/impl/ambiguous.clj b/src/ebb/impl/ambiguous.clj index 0d2e987..1582f8f 100644 --- a/src/ebb/impl/ambiguous.clj +++ b/src/ebb/impl/ambiguous.clj @@ -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 diff --git a/test/ebb/ap_cancel_seed_test.clj b/test/ebb/ap_cancel_seed_test.clj new file mode 100644 index 0000000..65e3fab --- /dev/null +++ b/test/ebb/ap_cancel_seed_test.clj @@ -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)))))