packages feed

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 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