bluefin-internal 0.10.0.0 → 0.10.1.0
raw patch · 6 files changed
+292/−5 lines, 6 filesdep ~asyncdep ~basePVP ok
version bump matches the API change (PVP)
Dependency ranges changed: async, base
API changes (from Hackage documentation)
+ Bluefin.Internal: keepAlive :: a -> IO b -> IO b
+ Bluefin.Internal: runPureEffAsyncSafe :: (forall (es :: Effects). () => Eff es a) -> a
+ Bluefin.Internal: runPureEffAsyncSafeRestarting :: (forall (es :: Effects). () => Eff es a) -> a
+ Bluefin.Internal: runPureEffPoisonable :: (forall (es :: Effects). () => Eff es a) -> a
Files
- bluefin-internal.cabal +4/−2
- src/Bluefin/Internal.hs +89/−3
- src/Bluefin/Internal/DslBuilderEff.hs +6/−0
- test/Main.hs +32/−0
- test/Test/RunPureEff.hs +144/−0
- test/Test/SpecH.hs +17/−0
bluefin-internal.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.0 name: bluefin-internal-version: 0.10.0.0+version: 0.10.1.0 license: MIT license-file: LICENSE author: Tom Ellis@@ -77,7 +77,7 @@ hs-source-dirs: src build-depends: async >= 2.2 && < 2.3,- base >= 4.14 && < 4.23,+ base >= 4.15 && < 4.23, unliftio-core < 0.3, primitive >= 0.8 && < 0.10, transformers < 0.7,@@ -107,7 +107,9 @@ hs-source-dirs: test main-is: Main.hs other-modules: Test.GeneralBracket,+ Test.RunPureEff, Test.SpecH build-depends: base,+ async, bluefin-internal
src/Bluefin/Internal.hs view
@@ -27,8 +27,9 @@ ) import Bluefin.Internal.Vault (Vault) import Bluefin.Internal.Vault qualified as Vault+import Control.Concurrent (forkIO, forkIOWithUnmask, killThread, myThreadId, throwTo) import Control.Concurrent.Async qualified as Async-import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)+import Control.Concurrent.MVar (newEmptyMVar, putMVar, readMVar, takeMVar) import Control.Exception qualified import Control.Monad (forever) import Control.Monad.Base (MonadBase (liftBase))@@ -40,13 +41,16 @@ import Control.Monad.Trans.Reader qualified as Reader import Data.Coerce (Coercible, coerce) import Data.Foldable (for_)+import Data.Function (fix) import Data.IORef (IORef, modifyIORef', newIORef, readIORef, writeIORef) import Data.Kind (Type) import Data.Proxy (Proxy (Proxy)) import Data.Type.Coercion (Coercion (Coercion))-import GHC.Exts (Any, Proxy#, proxy#)+import GHC.Exts (Any, Proxy#, keepAlive#, proxy#) import GHC.Generics (Generic, M1, Rec1, (:*:))+import GHC.IO (IO (..)) import System.IO.Unsafe (unsafePerformIO)+import System.Mem.Weak (addFinalizer) import Unsafe.Coerce (unsafeCoerce) import Prelude hiding (drop, head, read, return) @@ -277,7 +281,89 @@ -- | Run an 'Eff' that doesn't contain any unhandled effects. runPureEff :: (forall es. Eff es a) -> a-runPureEff e = unsafePerformIO (runEff (\_ -> e))+runPureEff = runPureEffPoisonable++-- | Run an 'Eff' that doesn't contain any unhandled effects. If the+-- evaluation of the result of runPureEffPoisonable is interrupted by+-- an asynchronous exception the thunk can become poisoned and unable+-- to be resumed (subsequent evaluation throws the async exception+-- again).+--+-- See https://github.com/tomjaguarpaw/bluefin/issues/30+runPureEffPoisonable :: (forall es. Eff es a) -> a+runPureEffPoisonable e = unsafePerformIO (runEff (\_ -> e))++-- | Run an 'Eff' that doesn't contain any unhandled effects. The+-- computation runs in a dedicated worker thread, so an asynchronous+-- exception received by a thread demanding the result does not+-- interrupt the computation itself. (An exception delivered to the+-- worker thread would be rethrown and could poison the thunk, but+-- that cannot happen unless the worker thread's ID is looked up by+-- some out-of-band means -- don't do that!).+--+-- The worker thread continues to work even if the thread forcing it+-- is killed, which may be surprising. If no references to the thunk+-- for the result of runPureEffAsyncSafe remain, the worker thread is+-- killed.+--+-- A proper fix to this issue probably belongs in GHC.+runPureEffAsyncSafe :: (forall es. Eff es a) -> a+runPureEffAsyncSafe e = unsafePerformIO $ do+ result <- newEmptyMVar+ owner <- newIORef ()+ _ <- Control.Exception.mask_ $ forkIOWithUnmask $ \unmask -> do+ tid <- myThreadId+ addFinalizer owner (killThread tid)+ r <- Control.Exception.try @Control.Exception.SomeException . unmask $ do+ runEff (\_ -> e)+ putMVar result r+ r <- keepAlive owner (readMVar result)+ either Control.Exception.throwIO pure r++keepAlive :: a -> IO b -> IO b+keepAlive a (IO action) = IO $ \s -> keepAlive# a s action++-- | Like 'runPureEffAsyncSafe', but starts a fresh worker after an+-- asynchronous exception interrupts a demand for the result. This+-- means that if the thread forcing the thunk is killed the work done+-- so far is discarded, and restarted from scratch in the next+-- evaluation.+--+-- This should not be used. It exists only as an example of a (worse)+-- alternative approach. Use 'runPureEffAsyncSafe' instead.+runPureEffAsyncSafeRestarting :: (forall es. Eff es a) -> a+runPureEffAsyncSafeRestarting effBody = unsafePerformIO $+ Control.Exception.mask $ \restore -> do+ tidVar <- newEmptyMVar+ done <- newEmptyMVar++ let body = do+ tid <- forkIO $ do+ r <-+ Control.Exception.try @Control.Exception.SomeException $+ restore $+ runEff (\_ -> effBody)+ putMVar done r+ putMVar tidVar tid+ takeMVar done++ r <- fix $ \again -> do+ attempted <-+ restore $+ Control.Exception.try @Control.Exception.SomeException body+ case attempted of+ Left ex -> do+ tid <- takeMVar tidVar+ killThread tid+ _ <- takeMVar done+ myself <- myThreadId+ throwTo myself ex+ again+ Right result -> pure result++ case r of+ Left l -> Control.Exception.throwIO l+ Right r' -> pure r' unsafeCoerceEff :: Eff t r -> Eff t' r unsafeCoerceEff = coerce
src/Bluefin/Internal/DslBuilderEff.hs view
@@ -29,6 +29,12 @@ -- | ͘ Eff es r runDslBuilderEff h f = makeOp (unMkDslBuilderEff f h)+{-# INLINE [0] runDslBuilderEff #-}+-- GHC's simplifier phase numbers count down toward 0. INLINE [0] keeps this+-- wrapper intact until phase 0, the final phase, so earlier simplifications+-- can work with the call before its body is exposed; it then strongly+-- encourages inlining to remove the wrapper. See+-- https://ghc.gitlab.haskell.org/ghc/doc/users_guide/exts/pragmas.html#phase-control -- oneShot is essential for good performance. I don't fully understand -- why.
test/Main.hs view
@@ -5,11 +5,19 @@ import Bluefin.Internal import Bluefin.Internal.Vault qualified as Vault+import Control.Exception (AsyncException (ThreadKilled)) import Control.Monad (forever, when) import Data.Foldable (for_) import Data.IORef (readIORef) import Data.Maybe (isNothing) import Test.GeneralBracket (test_generalBracket)+import Test.RunPureEff+ ( assertInterruptedBracketOutcome,+ isClean,+ isPoisoned,+ isRanAtLeastTwice,+ test_runPureEffAsyncSafeReapsWorker,+ ) import Test.SpecH (SpecH, assertEqual, runSpecH) import Prelude hiding (break, read) @@ -17,6 +25,30 @@ main = runEff $ \io -> do runSpecH io $ \y -> do let assertEqual' = assertEqual y++ assertInterruptedBracketOutcome+ io+ y+ "runPureEff retains bracket's rethrown exception"+ isPoisoned+ runPureEff+ assertInterruptedBracketOutcome+ io+ y+ "runPureEffAsyncSafe survives an interrupted bracket"+ isClean+ runPureEffAsyncSafe+ assertInterruptedBracketOutcome+ io+ y+ "runPureEffAsyncSafeRestarting restarts interrupted work"+ isRanAtLeastTwice+ runPureEffAsyncSafeRestarting+ workerException <- effIO io test_runPureEffAsyncSafeReapsWorker+ assertEqual'+ "runPureEffAsyncSafe worker receives ThreadKilled after result thunk GC"+ (Just ThreadKilled)+ workerException assertEqual' "oddsUntilFirstGreaterThan5" oddsUntilFirstGreaterThan5 [1, 3, 5, 7] assertEqual' "index 1" ([0, 1, 2, 3] !? 2) (Just 2)
+ test/Test/RunPureEff.hs view
@@ -0,0 +1,144 @@+{-# LANGUAGE BangPatterns #-}+{-# LANGUAGE NoMonoLocalBinds #-}+{-# LANGUAGE NoMonomorphismRestriction #-}++module Test.RunPureEff where++import Bluefin.Internal+import Control.Concurrent (threadDelay, throwTo)+import Control.Concurrent.Async (asyncThreadId, waitCatch, withAsync)+import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar, tryPutMVar)+import Control.Exception (AsyncException (ThreadKilled), SomeException, evaluate)+import Control.Exception qualified as Exception+import Control.Monad (forever, unless)+import Data.IORef (atomicModifyIORef', newIORef, readIORef, writeIORef)+import System.Mem (performMajorGC)+import System.Timeout (timeout)+import Test.SpecH (SpecH, assertSatisfies)++-- Interrupt a thread while it forces a shared thunk which runs inside+-- a bracket, then force the thunk again and check from what point the+-- computation resumed.+test_runPureEffAsyncSafeSurvivesInterruptedBracket ::+ (forall r. (forall e. Eff e r) -> r) ->+ IO InterruptedBracketResult+test_runPureEffAsyncSafeSurvivesInterruptedBracket run = do+ started <- newEmptyMVar+ continue <- newEmptyMVar+ forceThunk <- do+ modified <- newIORef False+ ref <- newIORef $ run $ bracket (pure ()) (\() -> pure ()) $ \() -> do+ let prog = do+ previouslyModified <-+ atomicModifyIORef' modified $ \previous ->+ (True, previous)+ -- Unless the IORef was previously modified, this is our+ -- first run. Therefore, wait for the caller to unblock us,+ -- and throw us an asynchronous exception.+ unless previouslyModified $ do+ _ <- tryPutMVar started ()+ takeMVar continue+ pure previouslyModified+ unsafeProvideIO (\io -> effIO io prog)+ pure (evaluate =<< readIORef ref)+ interrupted <- withAsync forceThunk $ \worker -> do+ takeMVar started+ throwTo (asyncThreadId worker) ThreadKilled+ waitCatch worker+ case interrupted of+ Left eInterrupted+ | Just ThreadKilled <- Exception.fromException eInterrupted -> do+ resumed <- do+ -- If the first evaluation is resumed, it will need+ -- unblocking. If evaluation started from scratch it+ -- won't block on continue anyway, because+ -- previouslyModified is true.+ putMVar continue ()+ Exception.try @SomeException $ do+ forceThunk+ case resumed of+ Left eResumed+ | Just ThreadKilled <- Exception.fromException eResumed ->+ pure ThunkPoisoned+ | otherwise ->+ pure (UnexpectedException eResumed)+ Right previouslyModified ->+ pure $ case previouslyModified of+ False -> ThunkClean+ True -> RanAtLeastTwice+ | otherwise -> pure (UnexpectedException eInterrupted)+ Right _ -> pure FinishedEarly++-- Drop the result thunk after cancelling its forcing thread, then collect it+-- and check that its weak finalizer kills the computation worker.+test_runPureEffAsyncSafeReapsWorker :: IO (Maybe AsyncException)+test_runPureEffAsyncSafeReapsWorker = do+ started <- newEmptyMVar+ caught <- newEmptyMVar+ (release, forceThunk) <- do+ shared <- newIORef @(Maybe ()) $ Just $ runPureEffAsyncSafe $ do+ unsafeProvideIO $ \io -> do+ effIO io $ do+ putMVar started ()+ Exception.handle @SomeException+ (putMVar caught)+ (forever (threadDelay 1_000_000))+ pure+ ( writeIORef shared Nothing,+ maybe (fail "result thunk released") evaluate =<< readIORef shared+ )+ withAsync forceThunk $ \forcingThunk -> do+ takeMVar started+ throwTo (asyncThreadId forcingThunk) ThreadKilled+ waitCatch forcingThunk >>= \case+ Right () ->+ fail "forcing thunk completed instead of being killed"+ Left e+ | -- We expect forcingThurk to be killed by ThreadKilled,+ -- because that's what we just threw to it.+ Just ThreadKilled <- Exception.fromException e ->+ pure ()+ | -- If we were killed by anything else, that should be+ -- reported+ otherwise ->+ Exception.throwIO e++ release+ performMajorGC+ caughtException <- timeout 1000000 (takeMVar caught)+ pure (Exception.fromException =<< caughtException)++data InterruptedBracketResult+ = FinishedEarly+ | UnexpectedException !SomeException+ | ThunkPoisoned+ | ThunkClean+ | RanAtLeastTwice+ deriving stock (Show)++isClean :: InterruptedBracketResult -> Bool+isClean = \case+ ThunkClean -> True+ _ -> False++isPoisoned :: InterruptedBracketResult -> Bool+isPoisoned = \case+ ThunkPoisoned -> True+ _ -> False++isRanAtLeastTwice :: InterruptedBracketResult -> Bool+isRanAtLeastTwice = \case+ RanAtLeastTwice -> True+ _ -> False++assertInterruptedBracketOutcome ::+ (e1 <: es, e2 <: es) =>+ IOE e1 ->+ SpecH e2 ->+ String ->+ (InterruptedBracketResult -> Bool) ->+ (forall r. (forall e. Eff e r) -> r) ->+ Eff es ()+assertInterruptedBracketOutcome io y name predicate run = do+ actual <- effIO io $ test_runPureEffAsyncSafeSurvivesInterruptedBracket run+ assertSatisfies y name predicate actual
test/Test/SpecH.hs view
@@ -23,6 +23,23 @@ yield y2 ("But got: " ++ show c2) ) +assertSatisfies ::+ (e <: es, Show a) =>+ SpecH e ->+ String ->+ (a -> Bool) ->+ a ->+ Eff es ()+assertSatisfies y n predicate actual =+ yield+ y+ ( n,+ if predicate actual+ then Nothing+ else Just $ dslBuilder $ \y2 ->+ yield y2 ("Predicate was not satisfied by: " ++ show actual)+ )+ type SpecInfo r = DslBuilder (Stream String) r runTests ::