pooled-io 0.0.0.1 → 0.0.1
raw patch · 4 files changed
+76/−8 lines, 4 filesdep +containersdep ~deepseqPVP ok
version bump matches the API change (PVP)
Dependencies added: containers
Dependency ranges changed: deepseq
API changes (from Hackage documentation)
+ Control.Concurrent.PooledIO.Independent: runException :: Maybe Int -> [IO ()] -> IO ()
Files
- pooled-io.cabal +9/−3
- src/Control/Concurrent/PooledIO/InOrder.hs +1/−1
- src/Control/Concurrent/PooledIO/Independent.hs +27/−1
- src/Control/Concurrent/PooledIO/Monad.hs +39/−3
pooled-io.cabal view
@@ -1,5 +1,5 @@ Name: pooled-io-Version: 0.0.0.1+Version: 0.0.1 License: BSD3 License-File: LICENSE Author: Henning Thielemann <haskell@henning-thielemann.de>@@ -29,6 +29,11 @@ Related packages: . * @lazyio@: interleave IO actions in a single thread+ .+ * @async@: start threads and wait for their results, forward exceptions,+ but do not throttle concurrency with respect to number of available cores+ .+ * @parallel-tasks@ Tested-With: GHC==7.4.1 Cabal-Version: >=1.8 Build-Type: Simple@@ -38,7 +43,7 @@ default: False Source-Repository this- Tag: 0.0.0.1+ Tag: 0.0.1 Type: darcs Location: http://code.haskell.org/~thielema/pooled-io/ @@ -49,8 +54,9 @@ Library Build-Depends: transformers >=0.2.2 && <0.5,- deepseq >=1.3 && <1.4,+ deepseq >=1.3 && <1.5, unsafe >=0.0 && <0.1,+ containers >=0.4 && <0.6, utility-ht >=0.0.9 && <0.1, base >=4 && <5
src/Control/Concurrent/PooledIO/InOrder.hs view
@@ -7,7 +7,7 @@ import qualified System.Unsafe as Unsafe -import Control.Monad.IO.Class (MonadIO, liftIO)+import Control.Monad.IO.Class (liftIO) import Control.Applicative (Applicative, pure, (<*>))
src/Control/Concurrent/PooledIO/Independent.hs view
@@ -3,9 +3,12 @@ run, runLimited, runUnlimited,+ runException, ) where -import Control.Concurrent.PooledIO.Monad (forkFinally, withNumCapabilities)+import Control.Concurrent.PooledIO.Monad+ (withNumCapabilities, chooseNumCapabilities,+ forkFinally, forkTry, takeMVarTry, runTry) import Control.Concurrent.MVar (MVar, newEmptyMVar, takeMVar) import Control.Exception (evaluate) @@ -16,6 +19,7 @@ Execute all actions parallelly but run at most @numCapabilities@ threads at once. Stop when all actions are finished.+If a thread throws an exception this terminates only the throwing thread. -} run :: [IO ()] -> IO () run = withNumCapabilities runLimited@@ -42,3 +46,25 @@ mvar <- newEmptyMVar forkFinally mvar act return mvar+++{- |+If a thread ends with an exception,+then terminate all threads and forward that exception.+@runException Nothing@ chooses to use all capabilities,+whereas @runException (Just n)@ chooses @n@ capabilities.+-}+runException :: Maybe Int -> [IO ()] -> IO ()+runException maybeNumCaps acts = do+ numCaps <- chooseNumCapabilities maybeNumCaps+ runOneBreaksAll numCaps acts++runOneBreaksAll :: Int -> [IO ()] -> IO ()+runOneBreaksAll numCaps acts = do+ let (start, queue) = splitAt numCaps acts+ n <- evaluate $ length start+ mvar <- newEmptyMVar+ runTry $ do+ mapM_ (forkTry mvar) start+ mapM_ (\act -> takeMVarTry mvar >> forkTry mvar act) queue+ replicateM_ n $ takeMVarTry mvar
src/Control/Concurrent/PooledIO/Monad.hs view
@@ -2,19 +2,22 @@ module Control.Concurrent.PooledIO.Monad where import Control.Concurrent.MVar (MVar, newEmptyMVar, takeMVar, putMVar)-import Control.Concurrent (forkIO, getNumCapabilities)+import Control.Concurrent (ThreadId, forkIO, getNumCapabilities, myThreadId, killThread) import Control.DeepSeq (NFData, ($!!))-import Control.Exception (finally)+import Control.Exception (SomeException, finally, try, throw) import qualified Control.Monad.Trans.State as MS import qualified Control.Monad.Trans.Reader as MR import qualified Control.Monad.Trans.Class as MT import Control.Monad.IO.Class (MonadIO, liftIO) -import Control.Monad (replicateM_)+import Control.Monad (replicateM_, liftM2) import Control.Functor.HT (void) +import qualified Data.Foldable as Fold+import qualified Data.Set as Set; import Data.Set (Set) + type T = MR.ReaderT (MVar ()) (MS.StateT Int IO) @@ -33,6 +36,39 @@ forkFinally :: MVar () -> IO () -> IO () forkFinally mvar act = void $ forkIO $ finally act $ putMVar mvar ()++forkTry ::+ (NFData a) =>+ MVar (ThreadId, Either SomeException a) -> IO a ->+ MS.StateT (Set ThreadId) IO ()+forkTry mvar act = do+ thread <-+ liftIO $ forkIO $+ applyStrictRight (putMVar mvar) =<< liftM2 (,) myThreadId (try act)+ MS.modify (Set.insert thread)++applyStrictRight :: (NFData a) => ((id, Either e a) -> b) -> (id, Either e a) -> b+applyStrictRight f (thread, ee) =+ case ee of+ Left e -> f (thread, Left e)+ Right a -> f . (,) thread . Right $!! a++takeMVarTry ::+ MVar (ThreadId, Either SomeException a) ->+ MS.StateT (Set ThreadId) IO a+takeMVarTry mvar = do+ (thread, ee) <- liftIO $ takeMVar mvar+ MS.modify (Set.delete thread)+ case ee of+ Left e -> do liftIO . Fold.mapM_ killThread =<< MS.get; liftIO $ throw e+ Right a -> return a++runTry :: MS.StateT (Set ThreadId) IO a -> IO a+runTry act = MS.evalStateT act Set.empty+++chooseNumCapabilities :: Maybe Int -> IO Int+chooseNumCapabilities = maybe getNumCapabilities return withNumCapabilities :: (Int -> a -> IO b) -> a -> IO b withNumCapabilities run acts = do