hspec-meta 2.11.17 → 2.11.18
raw patch · 22 files changed
+1489/−1303 lines, 22 filesdep ~QuickCheckPVP ok
version bump matches the API change (PVP)
Dependency ranges changed: QuickCheck
API changes (from Hackage documentation)
Files
- hspec-core/src/DList.hs +29/−0
- hspec-core/src/Test/Hspec/Core/Compat.hs +8/−3
- hspec-core/src/Test/Hspec/Core/Config/Definition.hs +2/−1
- hspec-core/src/Test/Hspec/Core/Example.hs +6/−3
- hspec-core/src/Test/Hspec/Core/Extension.hs +1/−0
- hspec-core/src/Test/Hspec/Core/Extension/Config.hs +2/−0
- hspec-core/src/Test/Hspec/Core/Extension/Item.hs +2/−0
- hspec-core/src/Test/Hspec/Core/Extension/Option.hs +2/−0
- hspec-core/src/Test/Hspec/Core/Extension/Spec.hs +2/−0
- hspec-core/src/Test/Hspec/Core/Extension/Tree.hs +2/−0
- hspec-core/src/Test/Hspec/Core/Format.hs +2/−0
- hspec-core/src/Test/Hspec/Core/Formatters/Diff.hs +29/−3
- hspec-core/src/Test/Hspec/Core/Formatters/Internal.hs +26/−24
- hspec-core/src/Test/Hspec/Core/Runner.hs +7/−4
- hspec-core/src/Test/Hspec/Core/Runner/Eval.hs +131/−138
- hspec-core/src/Test/Hspec/Core/Runner/JobQueue.hs +137/−77
- hspec-core/src/Test/Hspec/Core/Util.hs +0/−7
- hspec-meta.cabal +8/−5
- vendor/async-2.2.5/Control/Concurrent/Async.hs +3/−870
- vendor/async-2.2.5/Control/Concurrent/Async/Internal.hs +914/−0
- vendor/stm-2.5.0.1/Control/Concurrent/STM/TMVar.hs +0/−168
- vendor/stm-2.5.3.1/Control/Concurrent/STM/TMVar.hs +176/−0
+ hspec-core/src/DList.hs view
@@ -0,0 +1,29 @@+{-# LANGUAGE CPP #-}+module DList where++import Prelude ()+import Test.Hspec.Core.Compat++newtype DList a = DList ([a] -> [a])++instance Monoid (DList a) where+ mempty = DList id+#if MIN_VERSION_base(4,11,0)+instance Semigroup (DList a) where+#endif+ DList f+#if MIN_VERSION_base(4,11,0)+ <>+#else+ `mappend`+#endif+ DList g = DList (f . g)++singleton :: a -> DList a+singleton x = DList (x :)++fromList :: [a] -> DList a+fromList xs = DList (xs ++)++toList :: DList a -> [a]+toList (DList f) = f []
hspec-core/src/Test/Hspec/Core/Compat.hs view
@@ -89,12 +89,17 @@ import Control.Concurrent -#ifndef __MHS__ import GHC.IO.Exception++#if MIN_VERSION_base(4,13,0)+import GHC.IORef as Imports (atomicSwapIORef) #else-import System.IO.Error+atomicSwapIORef :: IORef a -> a -> IO a+atomicSwapIORef ref new = atomicModifyIORef ref $ \ old -> (new, old) #endif- ( ioe_type, IOErrorType(..) )++atomicReadIORef :: IORef b -> IO b+atomicReadIORef ref = atomicModifyIORef' ref $ \ a -> (a, a) isUnsupportedOperation :: IOError -> Bool isUnsupportedOperation e = ioe_type e == UnsupportedOperation
hspec-core/src/Test/Hspec/Core/Config/Definition.hs view
@@ -34,7 +34,7 @@ import Test.Hspec.Core.Compat import System.Directory (getTemporaryDirectory, removeFile)-import System.IO (openTempFile, hClose)+import System.IO (openTempFile, hClose, hFlush, stdout) import System.Process (system) import Test.Hspec.Core.Annotations (Annotations)@@ -184,6 +184,7 @@ tmp <- getTemporaryDirectory withTempFile tmp "hspec-expected" expected $ \ expectedFile -> do withTempFile tmp "hspec-actual" actual $ \ actualFile -> do+ hFlush stdout void . system $ unwords [command, expectedFile, actualFile] withTempFile :: FilePath -> FilePath -> String -> (FilePath -> IO a) -> IO a
hspec-core/src/Test/Hspec/Core/Example.hs view
@@ -23,7 +23,7 @@ , safeEvaluateExample -- END RE-EXPORTED from Test.Hspec.Core.Spec , safeEvaluateResultStatus-, exceptionToResultStatus+, exceptionToResult , toLocation , hunitFailureToResult ) where@@ -179,10 +179,13 @@ safeEvaluate :: IO Result -> IO Result safeEvaluate action = do- r <- safeTry $ forceResult <$> action+ r <- try $ fmap forceResult action >>= evaluate case r of- Left e -> Result "" <$> exceptionToResultStatus e+ Left e -> exceptionToResult e Right result -> return result++exceptionToResult :: SomeException -> IO Result+exceptionToResult err = Result "" <$> exceptionToResultStatus err safeEvaluateResultStatus :: IO ResultStatus -> IO ResultStatus safeEvaluateResultStatus action = do
hspec-core/src/Test/Hspec/Core/Extension.hs view
@@ -1,5 +1,6 @@ {-# OPTIONS_GHC -fno-warn-deprecations #-} -- | Stability: unstable+-- @since 2.11.10 module Test.Hspec.Core.Extension {-# WARNING "This API is experimental." #-} ( -- * Lifecycle of a test run {- |
hspec-core/src/Test/Hspec/Core/Extension/Config.hs view
@@ -1,4 +1,6 @@ {-# OPTIONS_GHC -fno-warn-deprecations #-}+-- | Stability: unstable+-- @since 2.11.10 module Test.Hspec.Core.Extension.Config {-# WARNING "This API is experimental." #-} ( -- * Types Config(..)
hspec-core/src/Test/Hspec/Core/Extension/Item.hs view
@@ -1,3 +1,5 @@+-- | Stability: unstable+-- @since 2.11.10 module Test.Hspec.Core.Extension.Item {-# WARNING "This API is experimental." #-} ( -- * Types Item(..)
hspec-core/src/Test/Hspec/Core/Extension/Option.hs view
@@ -1,3 +1,5 @@+-- | Stability: unstable+-- @since 2.11.10 module Test.Hspec.Core.Extension.Option {-# WARNING "This API is experimental." #-} ( Option , flag
hspec-core/src/Test/Hspec/Core/Extension/Spec.hs view
@@ -1,3 +1,5 @@+-- | Stability: unstable+-- @since 2.11.10 module Test.Hspec.Core.Extension.Spec {-# WARNING "This API is experimental." #-} ( mapItems ) where
hspec-core/src/Test/Hspec/Core/Extension/Tree.hs view
@@ -1,3 +1,5 @@+-- | Stability: unstable+-- @since 2.11.10 module Test.Hspec.Core.Extension.Tree {-# WARNING "This API is experimental." #-} ( SpecTree , mapItems
hspec-core/src/Test/Hspec/Core/Format.hs view
@@ -103,6 +103,8 @@ data Signal = Ok | NotOk SomeException monadic :: MonadIO m => (m () -> IO ()) -> (Event -> m ()) -> IO Format+-- `monadic` is not used internally anymore; we may want to deprecate it+-- eventually monadic run format = do mvar <- newEmptyMVar done <- newEmptyMVar
hspec-core/src/Test/Hspec/Core/Formatters/Diff.hs view
@@ -123,16 +123,42 @@ | otherwise = (LinesBoth xs :) diff :: String -> String -> [Diff]-diff expected actual = diffs+diff expected actual = case commonPrefixAndSuffix expectedChunks actualChunks of+ -- Avoid the general diff algorithm when the only change is an insertion or+ -- deletion between common boundaries. This case can contain many chunks.+ (prefix, [], insertion, suffix) -> chunk Both prefix ++ chunk Second insertion ++ chunk Both suffix+ (prefix, deletion, [], suffix) -> chunk Both prefix ++ chunk First deletion ++ chunk Both suffix+ (prefix, xs, ys, suffix) -> chunk Both prefix ++ diffs xs ys ++ chunk Both suffix where- diffs :: [Diff]- diffs = map (toDiff . fmap concat) $ Diff.getGroupedDiff expectedChunks actualChunks+ diffs :: [String] -> [String] -> [Diff]+ diffs xs ys = map (toDiff . fmap concat) $ Diff.getGroupedDiff xs ys expectedChunks :: [String] expectedChunks = partition expected actualChunks :: [String] actualChunks = partition actual++ chunk :: (String -> Diff) -> [String] -> [Diff]+ chunk constructor = \ case+ [] -> []+ xs -> [constructor (concat xs)]++commonPrefixAndSuffix :: Eq a => [a] -> [a] -> ([a], [a], [a], [a])+commonPrefixAndSuffix expected actual = (prefix, expectedMiddle, actualMiddle, suffix)+ where+ (prefix, expectedSuffix, actualSuffix) = commonPrefix expected actual+ (suffix, expectedMiddle, actualMiddle) = commonSuffix expectedSuffix actualSuffix++ commonPrefix :: Eq a => [a] -> [a] -> ([a], [a], [a])+ commonPrefix = go []+ where+ go acc (x : xs) (y : ys) | x == y = go (x : acc) xs ys+ go acc xs ys = (reverse acc, xs, ys)++ commonSuffix :: Eq a => [a] -> [a] -> ([a], [a], [a])+ commonSuffix xs ys = case commonPrefix (reverse xs) (reverse ys) of+ (s, e, a) -> (reverse s, reverse e, reverse a) toDiff :: Diff.Diff String -> Diff toDiff d = case d of
hspec-core/src/Test/Hspec/Core/Formatters/Internal.hs view
@@ -51,6 +51,7 @@ #ifdef TEST , runFormatM+, newFormatterState , splitLines #endif ) where@@ -100,21 +101,24 @@ } formatterToFormat :: Formatter -> FormatConfig -> IO Format-formatterToFormat Formatter{..} config = monadic (runFormatM config) $ \ case- Started -> formatterStarted- GroupStarted path -> formatterGroupStarted path- GroupDone path -> formatterGroupDone path- Progress path progress -> formatterProgress path progress- ItemStarted path -> formatterItemStarted path- ItemDone path item -> do- case itemResult item of- Success {} -> increaseSuccessCount- Pending {} -> increasePendingCount- Failure loc err -> addFailure $ FailureRecord (loc <|> itemLocation item) path err- formatterItemDone path item- Done _ -> formatterDone- where- addFailure r = modify $ \ s -> s { stateFailMessages = r : stateFailMessages s }+formatterToFormat Formatter{..} config = do+ ref <- newFormatterState config >>= newIORef+ return $ runFormatM ref . \ case+ Started -> formatterStarted+ GroupStarted path -> formatterGroupStarted path+ GroupDone path -> formatterGroupDone path+ Progress path progress -> formatterProgress path progress+ ItemStarted path -> formatterItemStarted path+ ItemDone path item -> do+ case itemResult item of+ Success {} -> increaseSuccessCount+ Pending {} -> increasePendingCount+ Failure loc err -> addFailure $ FailureRecord (loc <|> itemLocation item) path err+ formatterItemDone path item+ Done _ -> formatterDone+ where+ addFailure :: FailureRecord -> FormatM ()+ addFailure r = modify $ \ s -> s { stateFailMessages = r : stateFailMessages s } -- | Get the number of failed examples encountered so far. getFailCount :: FormatM Int@@ -216,14 +220,16 @@ -- NOTE: We use an IORef here, so that the state persists when UserInterrupt is -- thrown.-newtype FormatM a = FormatM (ReaderT (IORef FormatterState) IO a)+newtype FormatM a = FormatM { unFormatM :: ReaderT (IORef FormatterState) IO a } deriving (Functor, Applicative, Monad, MonadIO) -runFormatM :: FormatConfig -> FormatM a -> IO a-runFormatM config (FormatM action) = withLineBuffering $ do+runFormatM :: IORef FormatterState -> FormatM a -> IO a+runFormatM ref action = runReaderT (unFormatM action) ref++newFormatterState :: FormatConfig -> IO FormatterState+newFormatterState config = do time <- getMonotonicTime cpuTime <- if formatConfigPrintCpuTime config then Just <$> CPUTime.getCPUTime else pure Nothing- let progress = formatConfigReportProgress config && not (formatConfigHtmlOutput config) state = FormatterState {@@ -235,11 +241,7 @@ , stateConfig = config { formatConfigReportProgress = progress } , stateColor = Nothing }- newIORef state >>= runReaderT action--withLineBuffering :: IO a -> IO a-withLineBuffering action = bracket (IO.hGetBuffering stdout) (IO.hSetBuffering stdout) $ \ _ -> do- IO.hSetBuffering stdout IO.LineBuffering >> action+ return state -- | Increase the counter for successful examples increaseSuccessCount :: FormatM ()
hspec-core/src/Test/Hspec/Core/Runner.hs view
@@ -114,7 +114,6 @@ import qualified Test.QuickCheck as QC import Test.Hspec.Core.Util (Path)-import Test.Hspec.Core.Clock import Test.Hspec.Core.Spec hiding (pruneTree, pruneForest) import Test.Hspec.Core.Tree (formatDefaultDescription) import Test.Hspec.Core.Config@@ -170,7 +169,7 @@ removeCleanup _ = pass markSuccess :: EvalItem -> EvalItem- markSuccess item = item {evalItemAction = \ _ -> return (0, Result "" Success)}+ markSuccess item = item {evalItemAction = \ _ -> return (Result "" Success)} -- | Run a given spec and write a report to `stdout`. -- Exit with `exitFailure` if at least one spec item fails.@@ -373,7 +372,7 @@ concurrentJobs <- maybe getDefaultConcurrentJobs return $ configConcurrentJobs config - results <- fmap toSpecResult . withHiddenCursor (progressReporting colorMode) stdout $ do+ results <- fmap toSpecResult . withLineBuffering . withHiddenCursor (progressReporting colorMode) stdout $ do let formatConfig = FormatConfig { formatConfigUseColor = shouldUseColor colorMode@@ -458,7 +457,7 @@ evalItemDescription = requirement , evalItemLocation = loc , evalItemConcurrency = if isParallelizable == Just True then Concurrent else Sequential- , evalItemAction = \ progress -> measure $ e params withUnit progress+ , evalItemAction = \ progress -> e params withUnit progress } withUnit :: ActionWith () -> IO ()@@ -476,6 +475,10 @@ doNotLeakCommandLineArgumentsToExamples :: IO a -> IO a doNotLeakCommandLineArgumentsToExamples = withArgs []++withLineBuffering :: IO a -> IO a+withLineBuffering action = bracket (hGetBuffering stdout) (hSetBuffering stdout) $ \ _ -> do+ hSetBuffering stdout LineBuffering >> action withHiddenCursor :: ProgressReporting -> Handle -> IO a -> IO a withHiddenCursor progress h = case progress of
hspec-core/src/Test/Hspec/Core/Runner/Eval.hs view
@@ -1,8 +1,10 @@ {-# LANGUAGE CPP #-} {-# LANGUAGE DeriveTraversable #-}+{-# LANGUAGE NamedFieldPuns #-} {-# LANGUAGE RecordWildCards #-} {-# LANGUAGE ConstraintKinds #-} {-# LANGUAGE LambdaCase #-}+{-# LANGUAGE ScopedTypeVariables #-} module Test.Hspec.Core.Runner.Eval ( EvalConfig(..) , ColorMode(..)@@ -21,7 +23,6 @@ import Control.Monad.IO.Class (liftIO) import Control.Monad.Trans.Reader-import Control.Monad.Trans.Class import Test.Hspec.Core.Util import Test.Hspec.Core.Spec (Progress, FailureReason(..), Result(..), ResultStatus(..), ProgressCallback)@@ -30,17 +31,21 @@ import qualified Test.Hspec.Core.Format as Format import Test.Hspec.Core.Clock import Test.Hspec.Core.Example.Location-import Test.Hspec.Core.Example (safeEvaluateResultStatus, exceptionToResultStatus)+import Test.Hspec.Core.Example (safeEvaluateResultStatus, exceptionToResult) import qualified Data.List.NonEmpty as NonEmpty import Data.List.NonEmpty (NonEmpty(..)) -import Test.Hspec.Core.Runner.JobQueue+import DList (DList(..))+import qualified DList +import Test.Hspec.Core.Runner.JobQueue (JobQueue, Open, Concurrency(..), AbortEarly(..), Report(..))+import qualified Test.Hspec.Core.Runner.JobQueue as JobQueue+ data Tree c a =- Node String (NonEmpty (Tree c a))- | NodeWithCleanup (Maybe (String, Location)) c (NonEmpty (Tree c a))- | Leaf a+ Node !String !(NonEmpty (Tree c a))+ | NodeWithCleanup !(Maybe (String, Location)) !c !(NonEmpty (Tree c a))+ | Leaf !a deriving (Eq, Show, Functor, Foldable, Traversable) data EvalConfig = EvalConfig {@@ -54,7 +59,6 @@ data Env = Env { envConfig :: EvalConfig-, envAbort :: IORef Bool , envResults :: IORef [(Path, Format.Item)] } @@ -65,29 +69,11 @@ type EvalM = ReaderT Env IO -abort :: EvalM ()-abort = do- ref <- asks envAbort- liftIO $ writeIORef ref True--shouldAbort :: EvalM Bool-shouldAbort = do- ref <- asks envAbort- liftIO $ readIORef ref- addResult :: Path -> Format.Item -> EvalM () addResult path item = do ref <- asks envResults liftIO $ modifyIORef ref ((path, item) :) -reportItem :: Path -> Maybe Location -> EvalM (Seconds, Result) -> EvalM ()-reportItem path loc action = do- reportItemStarted path- action >>= reportResult path loc--reportItemStarted :: Path -> EvalM ()-reportItemStarted = formatEvent . Format.ItemStarted- reportItemDone :: Path -> Format.Item -> EvalM () reportItemDone path item = do addResult path item@@ -99,8 +85,8 @@ Pending{} -> False Failure{} -> True -reportResult :: Path -> Maybe Location -> (Seconds, Result) -> EvalM ()-reportResult path loc (duration, result) = do+formatResult :: Path -> Maybe Location -> (Seconds, Result) -> EvalM ()+formatResult path loc (duration, result) = do mode <- asks (evalConfigColorMode . envConfig) case result of Result info status -> reportItemDone path $ Format.Item loc duration info $ case status of@@ -118,17 +104,11 @@ Error _ _ -> err #endif -groupStarted :: Path -> EvalM ()-groupStarted = formatEvent . Format.GroupStarted--groupDone :: Path -> EvalM ()-groupDone = formatEvent . Format.GroupDone- data EvalItem = EvalItem { evalItemDescription :: String , evalItemLocation :: Maybe Location , evalItemConcurrency :: Concurrency-, evalItemAction :: ProgressCallback -> IO (Seconds, Result)+, evalItemAction :: ProgressCallback -> IO Result } type EvalTree = Tree (IO ()) EvalItem@@ -136,36 +116,79 @@ -- | Evaluate all examples of a given spec and produce a report. runFormatter :: EvalConfig -> [EvalTree] -> IO [(Path, Format.Item)] runFormatter config specs = do- withJobQueue (evalConfigConcurrentJobs config) $ \ queue -> do- withTimer 0.05 $ \ timer -> do- env <- mkEnv- runningSpecs_ <- enqueueItems queue specs+ queue <- JobQueue.new+ enqueuedSpecs <- enqueueItems queue specs+ jobs <- JobQueue.finalize queue+ withTimer 0.05 $ \ timer -> do+ env <- mkEnv+ let+ applyReportResult :: Item -> WithReportResult () Item+ applyReportResult = WithReportResult formatResult - let- applyReportProgress :: RunningItem_ IO -> RunningItem IO- applyReportProgress item = fmap (. reportProgress timer) item+ runningSpecs :: [Tree () (WithReportResult_ Item)]+ runningSpecs = applyCleanup abortEarly $ map (fmap applyReportResult) enqueuedSpecs - runningSpecs :: [RunningTree () EvalM]- runningSpecs = applyCleanup abortEarly $ map (fmap applyReportProgress) runningSpecs_+ runEvalM :: EvalM a -> IO a+ runEvalM = flip runReaderT env - abortEarly :: Result -> Bool- abortEarly result = evalConfigFailFast config && isFailure result+ eval :: forall item progress a. ([String] -> item -> DList (Report progress a)) -> [(Tree () item)] -> [Report progress a]+ eval evalItem = DList.toList . foldMap foldSpec+ where+ foldSpec :: Tree () item -> DList (Report progress a)+ foldSpec = foldTree FoldTree {+ onGroupStarted = DList.singleton . Report . groupStarted+ , onGroupDone = DList.singleton . Report . groupDone+ , onCleanup+ , onLeaf = evalItem+ } - getResults :: IO [(Path, Format.Item)]- getResults = reverse <$> readIORef (envResults env)+ groupStarted :: Path -> IO ()+ groupStarted = format . Format.GroupStarted - formatItems :: IO ()- formatItems = runReaderT (eval runningSpecs) env+ groupDone :: Path -> IO ()+ groupDone = format . Format.GroupDone - formatDone :: IO ()- formatDone = getResults >>= format . Format.Done+ onCleanup :: Maybe (String, Location) -> [String] -> () -> DList (Report progress a)+ onCleanup _loc _groups () = mempty - format Format.Started- formatItems `finally` formatDone- getResults+ items :: [Report Progress Result]+ items = eval evalItem runningSpecs+ where+ evalItem :: [String] -> WithReportResult_ Item -> DList (Report Progress Result)+ evalItem groups (WithReportResult reportResult (Item requirement loc action)) = DList $+ (Report (reportItemStarted path) :)+ . (ReportResult action progress result :)+ where+ path :: Path+ path = (groups, requirement)++ progress :: ProgressCallback+ progress = reportProgress timer path++ result :: (Seconds, Either SomeException Result) -> IO AbortEarly+ result = traverse (either exceptionToResult return) >=> runEvalM . reportResult path loc++ reportItemStarted :: Path -> IO ()+ reportItemStarted = format . Format.ItemStarted++ abortEarly :: Result -> Bool+ abortEarly result = evalConfigFailFast config && isFailure result++ getResults :: IO [(Path, Format.Item)]+ getResults = reverse <$> readIORef (envResults env)++ formatItems :: IO ()+ formatItems = JobQueue.run (evalConfigConcurrentJobs config) jobs items++ formatDone :: IO ()+ formatDone = getResults >>= format . Format.Done++ format Format.Started+ formatItems `finally` formatDone+ getResults where mkEnv :: IO Env- mkEnv = Env config <$> newIORef False <*> newIORef []+ mkEnv = Env config <$> newIORef [] format :: Format format = evalConfigFormat config@@ -176,68 +199,74 @@ when r $ do format (Format.Progress path progress) -data Item a = Item {+type ReportResult abort = Path -> Maybe Location -> (Seconds, Result) -> EvalM abort++data WithReportResult abort a = WithReportResult {+ _reportResult :: ReportResult abort+, _item :: a+}++type WithReportResult_ = WithReportResult AbortEarly++data Item = Item { itemDescription :: String , itemLocation :: Maybe Location-, itemAction :: a-} deriving Functor+, itemAction :: JobQueue.Result Progress Result+} -type RunningItem m = Item (Path -> m (Seconds, Result))-type RunningTree c m = Tree c (RunningItem m)+applyFailFast :: (Result -> Bool) -> Tree c (WithReportResult () a) -> Tree c (WithReportResult_ a)+applyFailFast abortEarly = fmap applyToWithReportResult+ where+ applyToWithReportResult :: WithReportResult () a -> WithReportResult_ a+ applyToWithReportResult (WithReportResult report a) = WithReportResult (applyToReportResult report) a -type RunningItem_ m = Item (Job m Progress (Seconds, Result))-type RunningTree_ m = Tree (IO ()) (RunningItem_ m)+ applyToReportResult :: ReportResult () -> ReportResult AbortEarly+ applyToReportResult report path loc result@(_, r) = do+ report path loc result+ return $ if abortEarly r then AbortEarly else NoAbortEarly -applyFailFast :: (Result -> Bool) -> RunningTree () IO -> RunningTree () EvalM-applyFailFast = fmap . fmap . fmap . applyToItem- where- applyToItem abortEarly action = do- result@(_, r) <- lift action- when (abortEarly r) abort- return result+type Children c a = NonEmpty (Tree c (WithReportResult_ a)) -applyCleanup :: (Result -> Bool) -> [RunningTree (IO ()) IO] -> [RunningTree () EvalM]-applyCleanup abortEarly = map (applyFailFast abortEarly . go)+applyCleanup :: (Result -> Bool) -> [Tree (IO ()) (WithReportResult () a)] -> [Tree () (WithReportResult_ a)]+applyCleanup abortEarly = map (go . applyFailFast abortEarly) where- go :: RunningTree (IO ()) IO -> RunningTree () IO+ go :: Tree (IO ()) (WithReportResult_ a) -> Tree () (WithReportResult_ a) go t = case t of Node label xs -> Node label (go <$> xs)- NodeWithCleanup loc cleanup xs -> NodeWithCleanup loc () (applyCleanupAction abortEarly loc cleanup $ go <$> xs)+ NodeWithCleanup loc cleanup xs -> NodeWithCleanup loc () (go <$> apply loc cleanup xs) Leaf a -> Leaf a -applyCleanupAction :: (Result -> Bool) -> Maybe (String, Location) -> IO () -> NonEmpty (RunningTree () IO) -> NonEmpty (RunningTree () IO)-applyCleanupAction abortEarly loc cleanup = forLastLeaf (addCleanupOn (not . abortEarly)) . forEachLeaf (addCleanupOn abortEarly)- where- addCleanupOn p = addCleanupToItem p loc cleanup+ apply :: Maybe (String, Location) -> IO () -> Children c a -> Children c a+ apply loc cleanup = forEachLeaf (addCleanupOn abortEarly) . forLastLeaf (addCleanupOn (not . abortEarly))+ where+ addCleanupOn :: (Result -> Bool) -> WithReportResult abort a -> WithReportResult abort a+ addCleanupOn p (WithReportResult report item) = WithReportResult (addCleanup p loc cleanup report) item -forEachLeaf :: (a -> b) -> NonEmpty (Tree () a) -> NonEmpty (Tree () b)+forEachLeaf :: (a -> b) -> NonEmpty (Tree c a) -> NonEmpty (Tree c b) forEachLeaf f = fmap (fmap f) -forLastLeaf :: (a -> a) -> NonEmpty (Tree () a) -> NonEmpty (Tree () a)+forLastLeaf :: (a -> a) -> NonEmpty (Tree c a) -> NonEmpty (Tree c a) forLastLeaf p = go where go = NonEmpty.reverse . mapHead goNode . NonEmpty.reverse goNode node = case node of Node description xs -> Node description (go xs)- NodeWithCleanup loc_ () xs -> NodeWithCleanup loc_ () (go xs)+ NodeWithCleanup loc_ c xs -> NodeWithCleanup loc_ c (go xs) Leaf item -> Leaf (p item) mapHead :: (a -> a) -> NonEmpty a -> NonEmpty a mapHead f xs = case xs of y :| ys -> f y :| ys -addCleanupToItem :: (Result -> Bool) -> Maybe (String, Location) -> IO () -> RunningItem IO -> RunningItem IO-addCleanupToItem shouldRunCleanup loc cleanup item = item {- itemAction = \ path -> do- result@(t1, r1) <- itemAction item path- if shouldRunCleanup r1 then do- (t2, r2) <- measure $ safeEvaluateResultStatus (cleanup >> return Success)- let t = t1 + t2- return (t, mergeResults loc r1 r2)- else do- return result-}+addCleanup :: (Result -> Bool) -> Maybe (String, Location) -> IO () -> ReportResult abort -> ReportResult abort+addCleanup shouldRunCleanup loc cleanup reportProgress path l result@(t1, r1) = do+ if shouldRunCleanup r1 then do+ (t2, r2) <- liftIO $ measure $ safeEvaluateResultStatus (cleanup >> return Success)+ let t = t1 + t2+ reportProgress path l $ (t, mergeResults loc r1 r2)+ else do+ reportProgress path l result mergeResults :: Maybe (String, Location) -> Result -> ResultStatus -> Result mergeResults mCallSite (Result info r1) r2 = Result info $ case (r1, r2) of@@ -254,72 +283,36 @@ Just (name, _) -> Just $ "in " ++ name ++ "-hook:" Nothing -> Nothing -enqueueItems :: MonadIO m => JobQueue -> [EvalTree] -> IO [RunningTree_ m]+enqueueItems :: JobQueue Open Progress Result -> [Tree c EvalItem] -> IO [Tree c Item] enqueueItems queue = mapM (traverse $ enqueueItem queue) -enqueueItem :: MonadIO m => JobQueue -> EvalItem -> IO (RunningItem_ m)+enqueueItem :: JobQueue Open Progress Result -> EvalItem -> IO Item enqueueItem queue EvalItem{..} = do- job <- enqueueJob queue evalItemConcurrency evalItemAction+ job <- JobQueue.enqueue queue evalItemConcurrency evalItemAction return Item { itemDescription = evalItemDescription , itemLocation = evalItemLocation- , itemAction = job >=> liftIO . either exceptionToResult return+ , itemAction = job }- where- exceptionToResult :: SomeException -> IO (Seconds, Result)- exceptionToResult err = (,) 0 . Result "" <$> exceptionToResultStatus err -eval :: [RunningTree () EvalM] -> EvalM ()-eval specs = do- sequenceActions (concatMap foldSpec specs)- where- foldSpec :: RunningTree () EvalM -> [EvalM ()]- foldSpec = foldTree FoldTree {- onGroupStarted = groupStarted- , onGroupDone = groupDone- , onCleanup = runCleanup- , onLeafe = evalItem- }-- runCleanup :: Maybe (String, Location) -> [String] -> () -> EvalM ()- runCleanup _loc _groups = return-- evalItem :: [String] -> RunningItem EvalM -> EvalM ()- evalItem groups (Item requirement loc action) = do- reportItem path loc $ action path- where- path :: Path- path = (groups, requirement)- data FoldTree c a r = FoldTree { onGroupStarted :: Path -> r , onGroupDone :: Path -> r , onCleanup :: Maybe (String, Location) -> [String] -> c -> r-, onLeafe :: [String] -> a -> r+, onLeaf :: [String] -> a -> r } -foldTree :: FoldTree c a r -> Tree c a -> [r]+foldTree :: Monoid r => FoldTree c a r -> Tree c a -> r foldTree FoldTree{..} = go [] where- go rGroups (Node group xs) = start : children ++ [done]+ go rGroups (Node group xs) = start <> children <> done where path = (reverse rGroups, group) start = onGroupStarted path- children = concatMap (go (group : rGroups)) xs+ children = foldMap (go (group : rGroups)) xs done = onGroupDone path- go rGroups (NodeWithCleanup loc action xs) = children ++ [cleanup]+ go rGroups (NodeWithCleanup loc action xs) = children <> cleanup where- children = concatMap (go rGroups) xs+ children = foldMap (go rGroups) xs cleanup = onCleanup loc (reverse rGroups) action- go rGroups (Leaf a) = [onLeafe (reverse rGroups) a]--sequenceActions :: [EvalM ()] -> EvalM ()-sequenceActions = go- where- go :: [EvalM ()] -> EvalM ()- go [] = pass- go (action : actions) = do- action- shouldAbort >>= \ case- False -> go actions- True -> pass+ go rGroups (Leaf a) = onLeaf (reverse rGroups) a
hspec-core/src/Test/Hspec/Core/Runner/JobQueue.hs view
@@ -1,108 +1,168 @@-{-# LANGUAGE CPP #-}+{-# LANGUAGE LambdaCase #-}+{-# LANGUAGE ViewPatterns #-} {-# LANGUAGE ScopedTypeVariables #-}-{-# LANGUAGE ConstraintKinds #-}- module Test.Hspec.Core.Runner.JobQueue (- MonadIO-, Job-, Concurrency(..)+ Concurrency(..) , JobQueue-, withJobQueue-, enqueueJob+, Open+, Final+, new+, finalize+, run+, enqueue+, Report(..)+, AbortEarly(..)+, Result ) where import Prelude ()-import Test.Hspec.Core.Compat hiding (Monad)-import qualified Test.Hspec.Core.Compat as M+import Test.Hspec.Core.Compat hiding (all) +import Foreign.Marshal.Utils (withMany) import Control.Concurrent-import Control.Concurrent.Async (Async, async, waitCatch, cancelMany)--import Control.Monad.IO.Class (liftIO)-import qualified Control.Monad.IO.Class as M---- for compatibility with GHC < 7.10.1-type Monad m = (Functor m, Applicative m, M.Monad m)-type MonadIO m = (Monad m, M.MonadIO m)+import Control.Concurrent.Async (Async, withAsync, link, wait, cancelMany) -type Job m progress a = (progress -> m ()) -> m a+import Test.Hspec.Core.Util (safeTry)+import Test.Hspec.Core.Clock (Seconds, measure) data Concurrency = Sequential | Concurrent -data JobQueue = JobQueue {- _semaphore :: Semaphore-, _cancelQueue :: CancelQueue-}+data Open+data Final -data Semaphore = Semaphore {- _wait :: IO ()-, _signal :: IO ()-}+type QueuedJob progress a = IO (Continuation progress a) -type CancelQueue = IORef [Async ()]+newtype JobQueue final progress a = JobQueue (IORef [QueuedJob progress a]) -withJobQueue :: Int -> (JobQueue -> IO a) -> IO a-withJobQueue concurrency = bracket new cancelAll- where- new :: IO JobQueue- new = JobQueue <$> newSemaphore concurrency <*> newIORef []+new :: IO (JobQueue Open progress a)+new = JobQueue <$> newIORef [] - cancelAll :: JobQueue -> IO ()- cancelAll (JobQueue _ cancelQueue) = readIORef cancelQueue >>= cancelMany+enqueueJob :: JobQueue Open progress a -> QueuedJob progress a -> IO ()+enqueueJob (JobQueue ref) job = modifyIORef' ref (job :) -newSemaphore :: Int -> IO Semaphore-newSemaphore capacity = do- sem <- newQSem capacity- return $ Semaphore (waitQSem sem) (signalQSem sem)+finalize :: JobQueue Open progress a -> IO (JobQueue Final progress a)+finalize (JobQueue ref) = do+ jobs <- reverse <$> atomicSwapIORef ref []+ JobQueue <$> newIORef jobs -enqueueJob :: MonadIO m => JobQueue -> Concurrency -> Job IO progress a -> IO (Job m progress (Either SomeException a))-enqueueJob (JobQueue sem cancelQueue) concurrency = case concurrency of- Sequential -> runSequentially cancelQueue- Concurrent -> runConcurrently sem cancelQueue+dequeue :: JobQueue Final progress a -> IO (Maybe (QueuedJob progress a))+dequeue (JobQueue ref) = atomicModifyIORef' ref $ \ case+ job : jobs -> (jobs, Just job)+ [] -> ([], Nothing) -runSequentially :: forall m progress a. MonadIO m => CancelQueue -> Job IO progress a -> IO (Job m progress (Either SomeException a))-runSequentially cancelQueue action = do- barrier <- newEmptyMVar+data Done = Abort | Done++run :: forall progress a. Int -> JobQueue Final progress a -> [Report progress a] -> IO ()+run concurrency jobs items = do+ parent <- newEmptyMVar let- wait :: IO ()- wait = takeMVar barrier+ worker :: Continuation progress a -> IO ()+ worker = workerThread parent jobs - signal :: m ()- signal = liftIO $ putMVar barrier ()+ workers :: [IO ()]+ workers = worker (ReportingThread items) : replicate (concurrency - 1) (worker WorkerThread) - job <- runConcurrently (Semaphore wait pass) cancelQueue action- return $ \ notifyPartial -> signal >> job notifyPartial+ withAsyncs workers $ \ threads -> do+ mapM_ link threads+ takeMVar parent >>= \ case+ Abort -> cancelMany threads+ Done -> mapM_ wait threads -data Partial progress a = Partial progress | Done+type Job progress a = (progress -> IO ()) -> IO a -runConcurrently :: forall m progress a. MonadIO m => Semaphore -> CancelQueue -> Job IO progress a -> IO (Job m progress (Either SomeException a))-runConcurrently (Semaphore wait signal) cancelQueue action = do- result :: MVar (Partial progress a) <- newEmptyMVar- let- worker :: IO a- worker = bracket_ wait signal $ do- interruptible (action partialResult) `finally` done- where- partialResult :: progress -> IO ()- partialResult = replaceMVar result . Partial+data Result progress a =+ SequentialResult (Job progress a)+ | ConcurrentResult (IORef (ConcurrentJobState progress a)) - done :: IO ()- done = replaceMVar result Done+enqueue :: JobQueue Open progress a -> Concurrency -> Job progress a -> IO (Result progress a)+enqueue queue concurrency job = case concurrency of+ Sequential -> return $ SequentialResult job+ Concurrent -> enqueueConcurrent queue job - pushOnCancelQueue :: Async a -> IO ()- pushOnCancelQueue = (modifyIORef cancelQueue . (:) . void)+data ConcurrentJobState progress a =+ Running+ | Reporting [Report progress a] (progress -> IO ())+ | Completed (Seconds, Either SomeException a) - job <- bracket (async worker) pushOnCancelQueue return+enqueueConcurrent :: forall progress a. JobQueue Open progress a -> Job progress a -> IO (Result progress a)+enqueueConcurrent queue job = do+ ref <- newIORef Running let- waitForResult :: (progress -> m ()) -> m (Either SomeException a)- waitForResult notifyPartial = do- r <- liftIO (takeMVar result)- case r of- Partial progress -> notifyPartial progress >> waitForResult notifyPartial- Done -> liftIO $ waitCatch job+ reportProgress :: progress -> IO ()+ reportProgress progress = atomicReadIORef ref >>= \ case+ Reporting _ report -> report progress+ _ -> pass - return waitForResult+ reportResult :: (Seconds, Either SomeException a) -> IO (Continuation progress a)+ reportResult result = atomicSwapIORef ref (Completed result) >>= \ case+ Running -> return WorkerThread+ Reporting items _ -> return $ ReportingThread items+ Completed _ -> thisCanNeverHappen "job completed twice" -replaceMVar :: MVar a -> a -> IO ()-replaceMVar mvar p = tryTakeMVar mvar >> putMVar mvar p+ enqueueJob queue $ do+ eval (job reportProgress) >>= reportResult++ return $ ConcurrentResult ref++data Report progress a =+ Report (IO ())+ | ReportResult {+ _resultAction :: Result progress a+ , _reportProgress :: progress -> IO ()+ , _reportResult :: (Seconds, Either SomeException a) -> IO AbortEarly+ }++data AbortEarly = NoAbortEarly | AbortEarly++data Continuation progress a = WorkerThread | ReportingThread [Report progress a]++workerThread :: forall progress a. MVar Done -> JobQueue Final progress a -> Continuation progress a -> IO ()+workerThread parent jobs = loop+ where+ loop :: Continuation progress a -> IO ()+ loop = \ case+ WorkerThread -> dequeue jobs >>= \ case+ Nothing -> pass+ Just job -> job >>= loop+ ReportingThread [] -> done+ ReportingThread all@(item : items) -> case item of+ Report action -> do+ action >> keepReporting+ ReportResult result reportProgress (checkAbort -> reportResult) -> case result of+ SequentialResult action -> do+ r <- eval (action reportProgress)+ reportResult r keepReporting+ ConcurrentResult ref -> passReportingResponsibility ref (Reporting all reportProgress) >>= \ case+ Running -> continueAsWorkerThread+ Completed r -> reportResult r keepReporting+ Reporting _ _ -> thisCanNeverHappen "more than one reporting thread"+ where+ passReportingResponsibility :: IORef s -> s -> IO s+ passReportingResponsibility = atomicSwapIORef++ continueAsWorkerThread :: IO ()+ continueAsWorkerThread = loop WorkerThread++ keepReporting :: IO ()+ keepReporting = loop (ReportingThread items)++ checkAbort :: (r -> IO AbortEarly) -> r -> IO () -> IO ()+ checkAbort report r continue = report r >>= \ case+ NoAbortEarly -> continue+ AbortEarly -> abort++ abort :: IO ()+ abort = putMVar parent Abort++ done :: IO ()+ done = putMVar parent Done++eval :: IO a -> IO (Seconds, Either SomeException a)+eval = measure . safeTry++withAsyncs :: forall a r. [IO a] -> ([Async a] -> IO r) -> IO r+withAsyncs = withMany withAsync++thisCanNeverHappen :: HasCallStack => String -> a+thisCanNeverHappen = error
hspec-core/src/Test/Hspec/Core/Util.hs view
@@ -1,4 +1,3 @@-{-# LANGUAGE CPP #-} {-# LANGUAGE ViewPatterns #-} -- | Stability: unstable module Test.Hspec.Core.Util (@@ -24,13 +23,7 @@ import Test.Hspec.Core.Compat hiding (join) import Data.Char (isSpace)--#ifndef __MHS__ import GHC.IO.Exception-#else-import System.IO.Error-#endif- import Control.Concurrent.Async -- |
hspec-meta.cabal view
@@ -5,7 +5,7 @@ -- see: https://github.com/sol/hpack name: hspec-meta-version: 2.11.17+version: 2.11.18 synopsis: A version of Hspec which is used to test Hspec itself description: A stable version of Hspec which is used to test the in-development version of Hspec.@@ -34,10 +34,12 @@ other-modules: Data.Algorithm.Diff Control.Concurrent.Async+ Control.Concurrent.Async.Internal Test.Hspec Test.Hspec.Formatters Test.Hspec.QuickCheck Test.Hspec.Runner+ DList GetOpt.Declarative GetOpt.Declarative.Environment GetOpt.Declarative.Interpret@@ -95,7 +97,7 @@ ghc-options: -Wall -fno-warn-incomplete-uni-patterns build-depends: HUnit ==1.6.*- , QuickCheck >=2.13.1 && <2.19+ , QuickCheck >=2.13.1 && <2.20 , ansi-terminal >=0.6.2 , array , base >=4.9.0.0 && <5@@ -123,7 +125,7 @@ other-modules: Control.Concurrent.STM.TMVar hs-source-dirs:- vendor/stm-2.5.0.1/+ vendor/stm-2.5.3.1/ executable hspec-meta-discover main-is: hspec-discover.hs@@ -138,7 +140,7 @@ ghc-options: -Wall -fno-warn-incomplete-uni-patterns build-depends: HUnit ==1.6.*- , QuickCheck >=2.13.1 && <2.19+ , QuickCheck >=2.13.1 && <2.20 , ansi-terminal >=0.6.2 , array , base >=4.9.0.0 && <5@@ -158,6 +160,7 @@ if impl(ghc) other-modules: Control.Concurrent.Async+ Control.Concurrent.Async.Internal hs-source-dirs: vendor/async-2.2.5/ cpp-options: -DENABLE_SPEC_HOOK_ARGS@@ -168,4 +171,4 @@ other-modules: Control.Concurrent.STM.TMVar hs-source-dirs:- vendor/stm-2.5.0.1/+ vendor/stm-2.5.3.1/
vendor/async-2.2.5/Control/Concurrent/Async.hs view
@@ -1,870 +1,3 @@-{-# LANGUAGE CPP, MagicHash, UnboxedTuples, RankNTypes,- ExistentialQuantification #-}-#if __GLASGOW_HASKELL__ >= 701-{-# LANGUAGE Trustworthy #-}-#endif-#if __GLASGOW_HASKELL__ < 710-{-# LANGUAGE DeriveDataTypeable #-}-#endif-{-# OPTIONS -Wall -fno-warn-implicit-prelude -fno-warn-unused-imports #-}---------------------------------------------------------------------------------- |--- Module : Control.Concurrent.Async.Internal--- Copyright : (c) Simon Marlow 2012--- License : BSD3 (see the file LICENSE)------ Maintainer : Simon Marlow <marlowsd@gmail.com>--- Stability : provisional--- Portability : non-portable (requires concurrency)------ This module is an internal module. The public API is provided in--- "Control.Concurrent.Async". Breaking changes to this module will not be--- reflected in a major bump, and using this module may break your code--- unless you are extremely careful.-----------------------------------------------------------------------------------module Control.Concurrent.Async where--import Control.Concurrent.STM.TMVar-import Control.Exception-import Control.Concurrent-import qualified Data.Foldable as F-#if !MIN_VERSION_base(4,6,0)-import Prelude hiding (catch)-#endif-import Control.Monad-import Control.Applicative-#if !MIN_VERSION_base(4,8,0)-import Data.Monoid (Monoid(mempty,mappend))-import Data.Traversable-#endif-#if __GLASGOW_HASKELL__ < 710-import Data.Typeable-#endif-#if MIN_VERSION_base(4,8,0)-import Data.Bifunctor-#endif-#if MIN_VERSION_base(4,9,0)-import Data.Semigroup (Semigroup((<>)))-#endif--import Data.IORef--import GHC.Exts-import GHC.IO hiding (finally, onException)-import GHC.Conc---- -------------------------------------------------------------------------------- STM Async API----- | An asynchronous action spawned by 'async' or 'withAsync'.--- Asynchronous actions are executed in a separate thread, and--- operations are provided for waiting for asynchronous actions to--- complete and obtaining their results (see e.g. 'wait').----data Async a = Async- { asyncThreadId :: {-# UNPACK #-} !ThreadId- -- ^ Returns the 'ThreadId' of the thread running- -- the given 'Async'.- , _asyncWait :: STM (Either SomeException a)- }--instance Eq (Async a) where- Async a _ == Async b _ = a == b--instance Ord (Async a) where- Async a _ `compare` Async b _ = a `compare` b--instance Functor Async where- fmap f (Async a w) = Async a (fmap (fmap f) w)---- | Compare two Asyncs that may have different types by their 'ThreadId'.-compareAsyncs :: Async a -> Async b -> Ordering-compareAsyncs (Async t1 _) (Async t2 _) = compare t1 t2---- | Spawn an asynchronous action in a separate thread.------ Like for 'forkIO', the action may be left running unintentionally--- (see module-level documentation for details).------ __Use 'withAsync' style functions wherever you can instead!__-async :: IO a -> IO (Async a)-async = inline asyncUsing rawForkIO---- | Like 'async' but using 'forkOS' internally.-asyncBound :: IO a -> IO (Async a)-asyncBound = asyncUsing forkOS---- | Like 'async' but using 'forkOn' internally.-asyncOn :: Int -> IO a -> IO (Async a)-asyncOn = asyncUsing . rawForkOn---- | Like 'async' but using 'forkIOWithUnmask' internally. The child--- thread is passed a function that can be used to unmask asynchronous--- exceptions.-asyncWithUnmask :: ((forall b . IO b -> IO b) -> IO a) -> IO (Async a)-asyncWithUnmask actionWith = asyncUsing rawForkIO (actionWith unsafeUnmask)---- | Like 'asyncOn' but using 'forkOnWithUnmask' internally. The--- child thread is passed a function that can be used to unmask--- asynchronous exceptions.-asyncOnWithUnmask :: Int -> ((forall b . IO b -> IO b) -> IO a) -> IO (Async a)-asyncOnWithUnmask cpu actionWith =- asyncUsing (rawForkOn cpu) (actionWith unsafeUnmask)--asyncUsing :: (IO () -> IO ThreadId)- -> IO a -> IO (Async a)-asyncUsing doFork = \action -> do- var <- newEmptyTMVarIO- -- t <- forkFinally action (\r -> atomically $ putTMVar var r)- -- slightly faster:- t <- mask $ \restore ->- doFork $ try (restore action) >>= atomically . putTMVar var- return (Async t (readTMVar var))---- | Spawn an asynchronous action in a separate thread, and pass its--- @Async@ handle to the supplied function. When the function returns--- or throws an exception, 'uninterruptibleCancel' is called on the @Async@.------ > withAsync action inner = mask $ \restore -> do--- > a <- async (restore action)--- > restore (inner a) `finally` uninterruptibleCancel a------ This is a useful variant of 'async' that ensures an @Async@ is--- never left running unintentionally.------ Note: a reference to the child thread is kept alive until the call--- to `withAsync` returns, so nesting many `withAsync` calls requires--- linear memory.----withAsync :: IO a -> (Async a -> IO b) -> IO b-withAsync = inline withAsyncUsing rawForkIO---- | Like 'withAsync' but uses 'forkOS' internally.-withAsyncBound :: IO a -> (Async a -> IO b) -> IO b-withAsyncBound = withAsyncUsing forkOS---- | Like 'withAsync' but uses 'forkOn' internally.-withAsyncOn :: Int -> IO a -> (Async a -> IO b) -> IO b-withAsyncOn = withAsyncUsing . rawForkOn---- | Like 'withAsync' but uses 'forkIOWithUnmask' internally. The--- child thread is passed a function that can be used to unmask--- asynchronous exceptions.-withAsyncWithUnmask- :: ((forall c. IO c -> IO c) -> IO a) -> (Async a -> IO b) -> IO b-withAsyncWithUnmask actionWith =- withAsyncUsing rawForkIO (actionWith unsafeUnmask)---- | Like 'withAsyncOn' but uses 'forkOnWithUnmask' internally. The--- child thread is passed a function that can be used to unmask--- asynchronous exceptions-withAsyncOnWithUnmask- :: Int -> ((forall c. IO c -> IO c) -> IO a) -> (Async a -> IO b) -> IO b-withAsyncOnWithUnmask cpu actionWith =- withAsyncUsing (rawForkOn cpu) (actionWith unsafeUnmask)--withAsyncUsing :: (IO () -> IO ThreadId)- -> IO a -> (Async a -> IO b) -> IO b--- The bracket version works, but is slow. We can do better by--- hand-coding it:-withAsyncUsing doFork = \action inner -> do- var <- newEmptyTMVarIO- mask $ \restore -> do- t <- doFork $ try (restore action) >>= atomically . putTMVar var- let a = Async t (readTMVar var)- r <- restore (inner a) `catchAll` \e -> do- uninterruptibleCancel a- throwIO e- uninterruptibleCancel a- return r---- | Wait for an asynchronous action to complete, and return its--- value. If the asynchronous action threw an exception, then the--- exception is re-thrown by 'wait'.------ > wait = atomically . waitSTM----{-# INLINE wait #-}-wait :: Async a -> IO a-wait = tryAgain . atomically . waitSTM- where- -- See: https://github.com/simonmar/async/issues/14- tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f---- | Wait for an asynchronous action to complete, and return either--- @Left e@ if the action raised an exception @e@, or @Right a@ if it--- returned a value @a@.------ > waitCatch = atomically . waitCatchSTM----{-# INLINE waitCatch #-}-waitCatch :: Async a -> IO (Either SomeException a)-waitCatch = tryAgain . atomically . waitCatchSTM- where- -- See: https://github.com/simonmar/async/issues/14- tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f---- | Check whether an 'Async' has completed yet. If it has not--- completed yet, then the result is @Nothing@, otherwise the result--- is @Just e@ where @e@ is @Left x@ if the @Async@ raised an--- exception @x@, or @Right a@ if it returned a value @a@.------ > poll = atomically . pollSTM----{-# INLINE poll #-}-poll :: Async a -> IO (Maybe (Either SomeException a))-poll = atomically . pollSTM---- | A version of 'wait' that can be used inside an STM transaction.----waitSTM :: Async a -> STM a-waitSTM a = do- r <- waitCatchSTM a- either throwSTM return r---- | A version of 'waitCatch' that can be used inside an STM transaction.----{-# INLINE waitCatchSTM #-}-waitCatchSTM :: Async a -> STM (Either SomeException a)-waitCatchSTM (Async _ w) = w---- | A version of 'poll' that can be used inside an STM transaction.----{-# INLINE pollSTM #-}-pollSTM :: Async a -> STM (Maybe (Either SomeException a))-pollSTM (Async _ w) = (Just <$> w) `orElse` return Nothing---- | Cancel an asynchronous action by throwing the @AsyncCancelled@--- exception to it, and waiting for the `Async` thread to quit.--- Has no effect if the 'Async' has already completed.------ > cancel a = throwTo (asyncThreadId a) AsyncCancelled <* waitCatch a------ Note that 'cancel' will not terminate until the thread the 'Async'--- refers to has terminated. This means that 'cancel' will block for--- as long said thread blocks when receiving an asynchronous exception.------ For example, it could block if:------ * It's executing a foreign call, and thus cannot receive the asynchronous--- exception;--- * It's executing some cleanup handler after having received the exception,--- and the handler is blocking.-{-# INLINE cancel #-}-cancel :: Async a -> IO ()-cancel a@(Async t _) = throwTo t AsyncCancelled <* waitCatch a---- | Cancel multiple asynchronous actions by throwing the @AsyncCancelled@--- exception to each of them in turn, then waiting for all the `Async` threads--- to complete.-cancelMany :: [Async a] -> IO ()-cancelMany as = do- mapM_ (\(Async t _) -> throwTo t AsyncCancelled) as- mapM_ waitCatch as---- | The exception thrown by `cancel` to terminate a thread.-data AsyncCancelled = AsyncCancelled- deriving (Show, Eq-#if __GLASGOW_HASKELL__ < 710- ,Typeable-#endif- )--instance Exception AsyncCancelled where-#if __GLASGOW_HASKELL__ >= 708- fromException = asyncExceptionFromException- toException = asyncExceptionToException-#endif---- | Cancel an asynchronous action------ This is a variant of `cancel`, but it is not interruptible.-{-# INLINE uninterruptibleCancel #-}-uninterruptibleCancel :: Async a -> IO ()-uninterruptibleCancel = uninterruptibleMask_ . cancel---- | Cancel an asynchronous action by throwing the supplied exception--- to it.------ > cancelWith a x = throwTo (asyncThreadId a) x------ The notes about the synchronous nature of 'cancel' also apply to--- 'cancelWith'.-cancelWith :: Exception e => Async a -> e -> IO ()-cancelWith a@(Async t _) e = throwTo t e <* waitCatch a---- | Wait for any of the supplied asynchronous operations to complete.--- The value returned is a pair of the 'Async' that completed, and the--- result that would be returned by 'wait' on that 'Async'.--- The input list must be non-empty.------ If multiple 'Async's complete or have completed, then the value--- returned corresponds to the first completed 'Async' in the list.----{-# INLINE waitAnyCatch #-}-waitAnyCatch :: [Async a] -> IO (Async a, Either SomeException a)-waitAnyCatch = atomically . waitAnyCatchSTM---- | A version of 'waitAnyCatch' that can be used inside an STM transaction.------ @since 2.1.0-waitAnyCatchSTM :: [Async a] -> STM (Async a, Either SomeException a)-waitAnyCatchSTM [] =- throwSTM $ ErrorCall- "waitAnyCatchSTM: invalid argument: input list must be non-empty"-waitAnyCatchSTM asyncs =- foldr orElse retry $- map (\a -> do r <- waitCatchSTM a; return (a, r)) asyncs---- | Like 'waitAnyCatch', but also cancels the other asynchronous--- operations as soon as one has completed.----waitAnyCatchCancel :: [Async a] -> IO (Async a, Either SomeException a)-waitAnyCatchCancel asyncs =- waitAnyCatch asyncs `finally` cancelMany asyncs---- | Wait for any of the supplied @Async@s to complete. If the first--- to complete throws an exception, then that exception is re-thrown--- by 'waitAny'.--- The input list must be non-empty.------ If multiple 'Async's complete or have completed, then the value--- returned corresponds to the first completed 'Async' in the list.----{-# INLINE waitAny #-}-waitAny :: [Async a] -> IO (Async a, a)-waitAny = atomically . waitAnySTM---- | A version of 'waitAny' that can be used inside an STM transaction.------ @since 2.1.0-waitAnySTM :: [Async a] -> STM (Async a, a)-waitAnySTM [] =- throwSTM $ ErrorCall- "waitAnySTM: invalid argument: input list must be non-empty"-waitAnySTM asyncs =- foldr orElse retry $- map (\a -> do r <- waitSTM a; return (a, r)) asyncs---- | Like 'waitAny', but also cancels the other asynchronous--- operations as soon as one has completed.----waitAnyCancel :: [Async a] -> IO (Async a, a)-waitAnyCancel asyncs =- waitAny asyncs `finally` cancelMany asyncs---- | Wait for the first of two @Async@s to finish.-{-# INLINE waitEitherCatch #-}-waitEitherCatch :: Async a -> Async b- -> IO (Either (Either SomeException a)- (Either SomeException b))-waitEitherCatch left right =- tryAgain $ atomically (waitEitherCatchSTM left right)- where- -- See: https://github.com/simonmar/async/issues/14- tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f---- | A version of 'waitEitherCatch' that can be used inside an STM transaction.------ @since 2.1.0-waitEitherCatchSTM :: Async a -> Async b- -> STM (Either (Either SomeException a)- (Either SomeException b))-waitEitherCatchSTM left right =- (Left <$> waitCatchSTM left)- `orElse`- (Right <$> waitCatchSTM right)---- | Like 'waitEitherCatch', but also 'cancel's both @Async@s before--- returning.----waitEitherCatchCancel :: Async a -> Async b- -> IO (Either (Either SomeException a)- (Either SomeException b))-waitEitherCatchCancel left right =- waitEitherCatch left right `finally` cancelMany [() <$ left, () <$ right]---- | Wait for the first of two @Async@s to finish. If the @Async@--- that finished first raised an exception, then the exception is--- re-thrown by 'waitEither'.----{-# INLINE waitEither #-}-waitEither :: Async a -> Async b -> IO (Either a b)-waitEither left right = atomically (waitEitherSTM left right)---- | A version of 'waitEither' that can be used inside an STM transaction.------ @since 2.1.0-waitEitherSTM :: Async a -> Async b -> STM (Either a b)-waitEitherSTM left right =- (Left <$> waitSTM left)- `orElse`- (Right <$> waitSTM right)---- | Like 'waitEither', but the result is ignored.----{-# INLINE waitEither_ #-}-waitEither_ :: Async a -> Async b -> IO ()-waitEither_ left right = atomically (waitEitherSTM_ left right)---- | A version of 'waitEither_' that can be used inside an STM transaction.------ @since 2.1.0-waitEitherSTM_:: Async a -> Async b -> STM ()-waitEitherSTM_ left right =- (void $ waitSTM left)- `orElse`- (void $ waitSTM right)---- | Like 'waitEither', but also 'cancel's both @Async@s before--- returning.----waitEitherCancel :: Async a -> Async b -> IO (Either a b)-waitEitherCancel left right =- waitEither left right `finally` cancelMany [() <$ left, () <$ right]---- | Waits for both @Async@s to finish, but if either of them throws--- an exception before they have both finished, then the exception is--- re-thrown by 'waitBoth'.----{-# INLINE waitBoth #-}-waitBoth :: Async a -> Async b -> IO (a,b)-waitBoth left right = tryAgain $ atomically (waitBothSTM left right)- where- -- See: https://github.com/simonmar/async/issues/14- tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f---- | A version of 'waitBoth' that can be used inside an STM transaction.------ @since 2.1.0-waitBothSTM :: Async a -> Async b -> STM (a,b)-waitBothSTM left right = do- a <- waitSTM left- `orElse`- (waitSTM right >> retry)- b <- waitSTM right- return (a,b)----- -------------------------------------------------------------------------------- Linking threads--data ExceptionInLinkedThread =- forall a . ExceptionInLinkedThread (Async a) SomeException-#if __GLASGOW_HASKELL__ < 710- deriving Typeable-#endif--instance Show ExceptionInLinkedThread where- showsPrec p (ExceptionInLinkedThread (Async t _) e) =- showParen (p >= 11) $- showString "ExceptionInLinkedThread " .- showsPrec 11 t .- showString " " .- showsPrec 11 e--instance Exception ExceptionInLinkedThread where-#if __GLASGOW_HASKELL__ >= 708- fromException = asyncExceptionFromException- toException = asyncExceptionToException-#endif---- | Link the given @Async@ to the current thread, such that if the--- @Async@ raises an exception, that exception will be re-thrown in--- the current thread, wrapped in 'ExceptionInLinkedThread'.------ 'link' ignores 'AsyncCancelled' exceptions thrown in the other thread,--- so that it's safe to 'cancel' a thread you're linked to. If you want--- different behaviour, use 'linkOnly'.----link :: Async a -> IO ()-link = linkOnly (not . isCancel)---- | Link the given @Async@ to the current thread, such that if the--- @Async@ raises an exception, that exception will be re-thrown in--- the current thread, wrapped in 'ExceptionInLinkedThread'.------ The supplied predicate determines which exceptions in the target--- thread should be propagated to the source thread.----linkOnly- :: (SomeException -> Bool) -- ^ return 'True' if the exception- -- should be propagated, 'False'- -- otherwise.- -> Async a- -> IO ()-linkOnly shouldThrow a = do- me <- myThreadId- void $ forkRepeat $ do- r <- waitCatch a- case r of- Left e | shouldThrow e -> throwTo me (ExceptionInLinkedThread a e)- _otherwise -> return ()---- | Link two @Async@s together, such that if either raises an--- exception, the same exception is re-thrown in the other @Async@,--- wrapped in 'ExceptionInLinkedThread'.------ 'link2' ignores 'AsyncCancelled' exceptions, so that it's possible--- to 'cancel' either thread without cancelling the other. If you--- want different behaviour, use 'link2Only'.----link2 :: Async a -> Async b -> IO ()-link2 = link2Only (not . isCancel)---- | Link two @Async@s together, such that if either raises an--- exception, the same exception is re-thrown in the other @Async@,--- wrapped in 'ExceptionInLinkedThread'.------ The supplied predicate determines which exceptions in the target--- thread should be propagated to the source thread.----link2Only :: (SomeException -> Bool) -> Async a -> Async b -> IO ()-link2Only shouldThrow left@(Async tl _) right@(Async tr _) =- void $ forkRepeat $ do- r <- waitEitherCatch left right- case r of- Left (Left e) | shouldThrow e ->- throwTo tr (ExceptionInLinkedThread left e)- Right (Left e) | shouldThrow e ->- throwTo tl (ExceptionInLinkedThread right e)- _ -> return ()--isCancel :: SomeException -> Bool-isCancel e- | Just AsyncCancelled <- fromException e = True- | otherwise = False----- --------------------------------------------------------------------------------- | Run two @IO@ actions concurrently, and return the first to--- finish. The loser of the race is 'cancel'led.------ > race left right =--- > withAsync left $ \a ->--- > withAsync right $ \b ->--- > waitEither a b----race :: IO a -> IO b -> IO (Either a b)---- | Like 'race', but the result is ignored.----race_ :: IO a -> IO b -> IO ()----- | Run two @IO@ actions concurrently, and return both results. If--- either action throws an exception at any time, then the other--- action is 'cancel'led, and the exception is re-thrown by--- 'concurrently'.------ > concurrently left right =--- > withAsync left $ \a ->--- > withAsync right $ \b ->--- > waitBoth a b-concurrently :: IO a -> IO b -> IO (a,b)----- | Run two @IO@ actions concurrently. If both of them end with @Right@,--- return both results. If one of then ends with @Left@, interrupt the other--- action and return the @Left@.----concurrentlyE :: IO (Either e a) -> IO (Either e b) -> IO (Either e (a, b))---- | 'concurrently', but ignore the result values------ @since 2.1.1-concurrently_ :: IO a -> IO b -> IO ()--#define USE_ASYNC_VERSIONS 0--#if USE_ASYNC_VERSIONS--race left right =- withAsync left $ \a ->- withAsync right $ \b ->- waitEither a b--race_ left right = void $ race left right--concurrently left right =- withAsync left $ \a ->- withAsync right $ \b ->- waitBoth a b--concurrently_ left right = void $ concurrently left right--#else---- MVar versions of race/concurrently--- More ugly than the Async versions, but quite a bit faster.---- race :: IO a -> IO b -> IO (Either a b)-race left right = concurrently' left right collect- where- collect m = do- e <- m- case e of- Left ex -> throwIO ex- Right r -> return r---- race_ :: IO a -> IO b -> IO ()-race_ left right = void $ race left right---- concurrently :: IO a -> IO b -> IO (a,b)-concurrently left right = concurrently' left right (collect [])- where- collect [Left a, Right b] _ = return (a,b)- collect [Right b, Left a] _ = return (a,b)- collect xs m = do- e <- m- case e of- Left ex -> throwIO ex- Right r -> collect (r:xs) m---- concurrentlyE :: IO (Either e a) -> IO (Either e b) -> IO (Either e (a, b))-concurrentlyE left right = concurrently' left right (collect [])- where- collect [Left (Right a), Right (Right b)] _ = return $ Right (a,b)- collect [Right (Right b), Left (Right a)] _ = return $ Right (a,b)- collect (Left (Left ea):_) _ = return $ Left ea- collect (Right (Left eb):_) _ = return $ Left eb- collect xs m = do- e <- m- case e of- Left ex -> throwIO ex- Right r -> collect (r:xs) m--concurrently' :: IO a -> IO b- -> (IO (Either SomeException (Either a b)) -> IO r)- -> IO r-concurrently' left right collect = do- done <- newEmptyMVar- mask $ \restore -> do- -- Note: uninterruptibleMask here is because we must not allow- -- the putMVar in the exception handler to be interrupted,- -- otherwise the parent thread will deadlock when it waits for- -- the thread to terminate.- lid <- forkIO $ uninterruptibleMask_ $- restore (left >>= putMVar done . Right . Left)- `catchAll` (putMVar done . Left)- rid <- forkIO $ uninterruptibleMask_ $- restore (right >>= putMVar done . Right . Right)- `catchAll` (putMVar done . Left)-- count <- newIORef (2 :: Int)- let takeDone = do- r <- takeMVar done -- interruptible- -- Decrement the counter so we know how many takes are left.- -- Since only the parent thread is calling this, we can- -- use non-atomic modifications.- -- NB. do this *after* takeMVar, because takeMVar might be- -- interrupted.- modifyIORef count (subtract 1)- return r-- let tryAgain f = f `catch` \BlockedIndefinitelyOnMVar -> f-- stop = do- -- kill right before left, to match the semantics of- -- the version using withAsync. (#27)- uninterruptibleMask_ $ do- count' <- readIORef count- -- we only need to use killThread if there are still- -- children alive. Note: forkIO here is because the- -- child thread could be in an uninterruptible- -- putMVar.- when (count' > 0) $- void $ forkIO $ do- throwTo rid AsyncCancelled- throwTo lid AsyncCancelled- -- ensure the children are really dead- replicateM_ count' (tryAgain $ takeMVar done)-- r <- collect (tryAgain $ takeDone) `onException` stop- stop- return r--concurrently_ left right = concurrently' left right (collect 0)- where- collect 2 _ = return ()- collect i m = do- e <- m- case e of- Left ex -> throwIO ex- Right _ -> collect (i + 1 :: Int) m---#endif---- | Maps an 'IO'-performing function over any 'Traversable' data--- type, performing all the @IO@ actions concurrently, and returning--- the original data structure with the arguments replaced by the--- results.------ If any of the actions throw an exception, then all other actions are--- cancelled and the exception is re-thrown.------ For example, @mapConcurrently@ works with lists:------ > pages <- mapConcurrently getURL ["url1", "url2", "url3"]------ Take into account that @async@ will try to immediately spawn a thread--- for each element of the @Traversable@, so running this on large--- inputs without care may lead to resource exhaustion (of memory,--- file descriptors, or other limited resources).-mapConcurrently :: Traversable t => (a -> IO b) -> t a -> IO (t b)-mapConcurrently f = runConcurrently . traverse (Concurrently . f)---- | `forConcurrently` is `mapConcurrently` with its arguments flipped------ > pages <- forConcurrently ["url1", "url2", "url3"] $ \url -> getURL url------ @since 2.1.0-forConcurrently :: Traversable t => t a -> (a -> IO b) -> IO (t b)-forConcurrently = flip mapConcurrently---- | `mapConcurrently_` is `mapConcurrently` with the return value discarded;--- a concurrent equivalent of 'mapM_'.-mapConcurrently_ :: F.Foldable f => (a -> IO b) -> f a -> IO ()-mapConcurrently_ f = runConcurrently . F.foldMap (Concurrently . void . f)---- | `forConcurrently_` is `forConcurrently` with the return value discarded;--- a concurrent equivalent of 'forM_'.-forConcurrently_ :: F.Foldable f => f a -> (a -> IO b) -> IO ()-forConcurrently_ = flip mapConcurrently_---- | Perform the action in the given number of threads.------ @since 2.1.1-replicateConcurrently :: Int -> IO a -> IO [a]-replicateConcurrently cnt = runConcurrently . sequenceA . replicate cnt . Concurrently---- | Same as 'replicateConcurrently', but ignore the results.------ @since 2.1.1-replicateConcurrently_ :: Int -> IO a -> IO ()-replicateConcurrently_ cnt = runConcurrently . F.fold . replicate cnt . Concurrently . void---- --------------------------------------------------------------------------------- | A value of type @Concurrently a@ is an @IO@ operation that can be--- composed with other @Concurrently@ values, using the @Applicative@--- and @Alternative@ instances.------ Calling @runConcurrently@ on a value of type @Concurrently a@ will--- execute the @IO@ operations it contains concurrently, before--- delivering the result of type @a@.------ For example------ > (page1, page2, page3)--- > <- runConcurrently $ (,,)--- > <$> Concurrently (getURL "url1")--- > <*> Concurrently (getURL "url2")--- > <*> Concurrently (getURL "url3")----newtype Concurrently a = Concurrently { runConcurrently :: IO a }--instance Functor Concurrently where- fmap f (Concurrently a) = Concurrently $ f <$> a--instance Applicative Concurrently where- pure = Concurrently . return- Concurrently fs <*> Concurrently as =- Concurrently $ (\(f, a) -> f a) <$> concurrently fs as---- | 'Control.Alternative.empty' waits forever. 'Control.Alternative.<|>' returns the first to finish and 'cancel's the other.-instance Alternative Concurrently where- empty = Concurrently $ forever (threadDelay maxBound)- Concurrently as <|> Concurrently bs =- Concurrently $ either id id <$> race as bs--#if MIN_VERSION_base(4,9,0)--- | Only defined by @async@ for @base >= 4.9@------ @since 2.1.0-instance Semigroup a => Semigroup (Concurrently a) where- (<>) = liftA2 (<>)---- | @since 2.1.0-instance (Semigroup a, Monoid a) => Monoid (Concurrently a) where- mempty = pure mempty- mappend = (<>)-#else--- | @since 2.1.0-instance Monoid a => Monoid (Concurrently a) where- mempty = pure mempty- mappend = liftA2 mappend-#endif---- | A value of type @ConcurrentlyE e a@ is an @IO@ operation that can be--- composed with other @ConcurrentlyE@ values, using the @Applicative@ instance.------ Calling @runConcurrentlyE@ on a value of type @ConcurrentlyE e a@ will--- execute the @IO@ operations it contains concurrently, before delivering--- either the result of type @a@, or an error of type @e@ if one of the actions--- returns @Left@.------ | @since 2.2.5-newtype ConcurrentlyE e a = ConcurrentlyE { runConcurrentlyE :: IO (Either e a) }--instance Functor (ConcurrentlyE e) where- fmap f (ConcurrentlyE ea) = ConcurrentlyE $ fmap (fmap f) ea--#if MIN_VERSION_base(4,8,0)-instance Bifunctor ConcurrentlyE where- bimap f g (ConcurrentlyE ea) = ConcurrentlyE $ fmap (bimap f g) ea-#endif--instance Applicative (ConcurrentlyE e) where- pure = ConcurrentlyE . return . return- ConcurrentlyE fs <*> ConcurrentlyE eas =- ConcurrentlyE $ fmap (\(f, a) -> f a) <$> concurrentlyE fs eas--#if MIN_VERSION_base(4,9,0)--- | Either the combination of the successful results, or the first failure.-instance Semigroup a => Semigroup (ConcurrentlyE e a) where- (<>) = liftA2 (<>)--instance (Semigroup a, Monoid a) => Monoid (ConcurrentlyE e a) where- mempty = pure mempty- mappend = (<>)-#endif---- -------------------------------------------------------------------------------- | Fork a thread that runs the supplied action, and if it raises an--- exception, re-runs the action. The thread terminates only when the--- action runs to completion without raising an exception.-forkRepeat :: IO a -> IO ThreadId-forkRepeat action =- mask $ \restore ->- let go = do r <- tryAll (restore action)- case r of- Left _ -> go- _ -> return ()- in forkIO go--catchAll :: IO a -> (SomeException -> IO a) -> IO a-catchAll = catch--tryAll :: IO a -> IO (Either SomeException a)-tryAll = try---- A version of forkIO that does not include the outer exception--- handler: saves a bit of time when we will be installing our own--- exception handler.-{-# INLINE rawForkIO #-}-rawForkIO :: IO () -> IO ThreadId-rawForkIO (IO action) = IO $ \ s ->- case (fork# action s) of (# s1, tid #) -> (# s1, ThreadId tid #)--{-# INLINE rawForkOn #-}-rawForkOn :: Int -> IO () -> IO ThreadId-rawForkOn (I# cpu) (IO action) = IO $ \ s ->- case (forkOn# cpu action s) of (# s1, tid #) -> (# s1, ThreadId tid #)+module Control.Concurrent.Async (module Control.Concurrent.Async.Internal) where+import Prelude ()+import Control.Concurrent.Async.Internal
+ vendor/async-2.2.5/Control/Concurrent/Async/Internal.hs view
@@ -0,0 +1,914 @@+{-# LANGUAGE CPP, MagicHash, UnboxedTuples, RankNTypes,+ ExistentialQuantification #-}+#if __GLASGOW_HASKELL__ >= 701+{-# LANGUAGE Trustworthy #-}+#endif+#if __GLASGOW_HASKELL__ < 710+{-# LANGUAGE DeriveDataTypeable #-}+#endif+{-# OPTIONS -Wall -fno-warn-implicit-prelude -fno-warn-unused-imports #-}++-----------------------------------------------------------------------------+-- |+-- Module : Control.Concurrent.Async.Internal+-- Copyright : (c) Simon Marlow 2012+-- License : BSD3 (see the file LICENSE)+--+-- Maintainer : Simon Marlow <marlowsd@gmail.com>+-- Stability : provisional+-- Portability : non-portable (requires concurrency)+--+-- This module is an internal module. The public API is provided in+-- "Control.Concurrent.Async". Breaking changes to this module will not be+-- reflected in a major bump, and using this module may break your code+-- unless you are extremely careful.+--+-----------------------------------------------------------------------------++module Control.Concurrent.Async.Internal where++import Control.Concurrent.STM.TMVar+import Control.Exception+import Control.Concurrent+import qualified Data.Foldable as F+#if !MIN_VERSION_base(4,6,0)+import Prelude hiding (catch)+#endif+import Control.Monad+import Control.Applicative+#if !MIN_VERSION_base(4,8,0)+import Data.Monoid (Monoid(mempty,mappend))+import Data.Traversable+#endif+#if __GLASGOW_HASKELL__ < 710+import Data.Typeable+#endif+#if MIN_VERSION_base(4,8,0)+import Data.Bifunctor+#endif+#if MIN_VERSION_base(4,9,0)+import Data.Semigroup (Semigroup((<>)))+#endif++import Data.IORef++import GHC.Exts+import GHC.IO hiding (finally, onException)+import GHC.Conc++-- -----------------------------------------------------------------------------+-- STM Async API+++-- | An asynchronous action spawned by 'async' or 'withAsync'.+-- Asynchronous actions are executed in a separate thread, and+-- operations are provided for waiting for asynchronous actions to+-- complete and obtaining their results (see e.g. 'wait').+--+data Async a = Async+ { asyncThreadId :: {-# UNPACK #-} !ThreadId+ -- ^ Returns the 'ThreadId' of the thread running+ -- the given 'Async'.+ , _asyncWait :: STM (Either SomeException a)+ }++instance Eq (Async a) where+ Async a _ == Async b _ = a == b++instance Ord (Async a) where+ Async a _ `compare` Async b _ = a `compare` b++instance Functor Async where+ fmap f (Async a w) = Async a (fmap (fmap f) w)++-- | Compare two Asyncs that may have different types by their 'ThreadId'.+compareAsyncs :: Async a -> Async b -> Ordering+compareAsyncs (Async t1 _) (Async t2 _) = compare t1 t2++-- | Spawn an asynchronous action in a separate thread.+--+-- Like for 'forkIO', the action may be left running unintentionally+-- (see module-level documentation for details).+--+-- __Use 'withAsync' style functions wherever you can instead!__+async ::+ IO a -> IO (Async a)+async = inline asyncUsing rawForkIO++-- | Like 'async' but using 'forkOS' internally.+asyncBound ::+ IO a -> IO (Async a)+asyncBound = asyncUsing forkOS++-- | Like 'async' but using 'forkOn' internally.+asyncOn ::+ Int -> IO a -> IO (Async a)+asyncOn = asyncUsing . rawForkOn++-- | Like 'async' but using 'forkIOWithUnmask' internally. The child+-- thread is passed a function that can be used to unmask asynchronous+-- exceptions.+asyncWithUnmask ::+ ((forall b . IO b -> IO b) -> IO a) -> IO (Async a)+asyncWithUnmask actionWith = asyncUsing rawForkIO (actionWith unsafeUnmask)++-- | Like 'asyncOn' but using 'forkOnWithUnmask' internally. The+-- child thread is passed a function that can be used to unmask+-- asynchronous exceptions.+asyncOnWithUnmask ::+ Int -> ((forall b . IO b -> IO b) -> IO a) -> IO (Async a)+asyncOnWithUnmask cpu actionWith =+ asyncUsing (rawForkOn cpu) (actionWith unsafeUnmask)++asyncUsing ::+ (IO () -> IO ThreadId) -> IO a -> IO (Async a)+asyncUsing doFork action = do+ var <- newEmptyTMVarIO+ let action_plus = debugLabelMe >> action+ -- t <- forkFinally action (\r -> atomically $ putTMVar var r)+ -- slightly faster:+ t <- mask $ \restore ->+ doFork $ try (restore action_plus) >>= atomically . putTMVar var+ return (Async t (readTMVar var))+++-- | Spawn an asynchronous action in a separate thread, and pass its+-- @Async@ handle to the supplied function. When the function returns+-- or throws an exception, 'uninterruptibleCancel' is called on the @Async@.+--+-- > withAsync action inner = mask $ \restore -> do+-- > a <- async (restore action)+-- > restore (inner a) `finally` uninterruptibleCancel a+--+-- This is a useful variant of 'async' that ensures an @Async@ is+-- never left running unintentionally.+--+-- Note: a reference to the child thread is kept alive until the call+-- to `withAsync` returns, so nesting many `withAsync` calls requires+-- linear memory.+--+withAsync ::+ IO a -> (Async a -> IO b) -> IO b+withAsync = inline withAsyncUsing rawForkIO++-- | Like 'withAsync' but uses 'forkOS' internally.+withAsyncBound ::+ IO a -> (Async a -> IO b) -> IO b+withAsyncBound = withAsyncUsing forkOS++-- | Like 'withAsync' but uses 'forkOn' internally.+withAsyncOn ::+ Int -> IO a -> (Async a -> IO b) -> IO b+withAsyncOn = withAsyncUsing . rawForkOn++-- | Like 'withAsync' but uses 'forkIOWithUnmask' internally. The+-- child thread is passed a function that can be used to unmask+-- asynchronous exceptions.+withAsyncWithUnmask ::+ ((forall c. IO c -> IO c) -> IO a) -> (Async a -> IO b) -> IO b+withAsyncWithUnmask actionWith =+ withAsyncUsing rawForkIO (actionWith unsafeUnmask)++-- | Like 'withAsyncOn' but uses 'forkOnWithUnmask' internally. The+-- child thread is passed a function that can be used to unmask+-- asynchronous exceptions+withAsyncOnWithUnmask ::+ Int -> ((forall c. IO c -> IO c) -> IO a) -> (Async a -> IO b) -> IO b+withAsyncOnWithUnmask cpu actionWith =+ withAsyncUsing (rawForkOn cpu) (actionWith unsafeUnmask)++withAsyncUsing ::+ (IO () -> IO ThreadId) -> IO a -> (Async a -> IO b) -> IO b+-- The bracket version works, but is slow. We can do better by+-- hand-coding it:+withAsyncUsing doFork action inner = do+ var <- newEmptyTMVarIO+ mask $ \restore -> do+ let action_plus = debugLabelMe >> action+ t <- doFork $ try (restore action_plus) >>= atomically . putTMVar var+ let a = Async t (readTMVar var)+ r <- restore (inner a) `catchAll` \e -> do+ uninterruptibleCancel a+ throwIO e+ uninterruptibleCancel a+ return r++-- | Wait for an asynchronous action to complete, and return its+-- value. If the asynchronous action threw an exception, then the+-- exception is re-thrown by 'wait'.+--+-- > wait = atomically . waitSTM+--+{-# INLINE wait #-}+wait :: Async a -> IO a+wait = tryAgain . atomically . waitSTM+ where+ -- See: https://github.com/simonmar/async/issues/14+ tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f++-- | Wait for an asynchronous action to complete, and return either+-- @Left e@ if the action raised an exception @e@, or @Right a@ if it+-- returned a value @a@.+--+-- > waitCatch = atomically . waitCatchSTM+--+{-# INLINE waitCatch #-}+waitCatch :: Async a -> IO (Either SomeException a)+waitCatch = tryAgain . atomically . waitCatchSTM+ where+ -- See: https://github.com/simonmar/async/issues/14+ tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f++-- | Check whether an 'Async' has completed yet. If it has not+-- completed yet, then the result is @Nothing@, otherwise the result+-- is @Just e@ where @e@ is @Left x@ if the @Async@ raised an+-- exception @x@, or @Right a@ if it returned a value @a@.+--+-- > poll = atomically . pollSTM+--+{-# INLINE poll #-}+poll :: Async a -> IO (Maybe (Either SomeException a))+poll = atomically . pollSTM++-- | A version of 'wait' that can be used inside an STM transaction.+--+waitSTM :: Async a -> STM a+waitSTM a = do+ r <- waitCatchSTM a+ either throwSTM return r++-- | A version of 'waitCatch' that can be used inside an STM transaction.+--+{-# INLINE waitCatchSTM #-}+waitCatchSTM :: Async a -> STM (Either SomeException a)+waitCatchSTM (Async _ w) = w++-- | A version of 'poll' that can be used inside an STM transaction.+--+{-# INLINE pollSTM #-}+pollSTM :: Async a -> STM (Maybe (Either SomeException a))+pollSTM (Async _ w) = (Just <$> w) `orElse` return Nothing++-- | Cancel an asynchronous action by throwing the @AsyncCancelled@+-- exception to it, and waiting for the `Async` thread to quit.+-- Has no effect if the 'Async' has already completed.+--+-- > cancel a = throwTo (asyncThreadId a) AsyncCancelled <* waitCatch a+--+-- Note that 'cancel' will not terminate until the thread the 'Async'+-- refers to has terminated. This means that 'cancel' will block for+-- as long said thread blocks when receiving an asynchronous exception.+--+-- For example, it could block if:+--+-- * It's executing a foreign call, and thus cannot receive the asynchronous+-- exception;+-- * It's executing some cleanup handler after having received the exception,+-- and the handler is blocking.+{-# INLINE cancel #-}+cancel :: Async a -> IO ()+cancel a@(Async t _) = throwTo t AsyncCancelled <* waitCatch a++-- | Cancel multiple asynchronous actions by throwing the @AsyncCancelled@+-- exception to each of them in turn, then waiting for all the `Async` threads+-- to complete.+--+-- @since 2.2.5+cancelMany :: [Async a] -> IO ()+cancelMany as = do+ mapM_ (\(Async t _) -> throwTo t AsyncCancelled) as+ mapM_ waitCatch as++-- | The exception thrown by `cancel` to terminate a thread.+data AsyncCancelled = AsyncCancelled+ deriving (Show, Eq+#if __GLASGOW_HASKELL__ < 710+ ,Typeable+#endif+ )++instance Exception AsyncCancelled where+#if __GLASGOW_HASKELL__ >= 708+ fromException = asyncExceptionFromException+ toException = asyncExceptionToException+#endif++-- | Cancel an asynchronous action+--+-- This is a variant of `cancel`, but it is not interruptible.+{-# INLINE uninterruptibleCancel #-}+uninterruptibleCancel :: Async a -> IO ()+uninterruptibleCancel = uninterruptibleMask_ . cancel++-- | Cancel an asynchronous action by throwing the supplied exception+-- to it.+--+-- > cancelWith a x = throwTo (asyncThreadId a) x+--+-- The notes about the synchronous nature of 'cancel' also apply to+-- 'cancelWith'.+cancelWith :: Exception e => Async a -> e -> IO ()+cancelWith a@(Async t _) e = throwTo t e <* waitCatch a++-- | Wait for any of the supplied asynchronous operations to complete.+-- The value returned is a pair of the 'Async' that completed, and the+-- result that would be returned by 'wait' on that 'Async'.+-- The input list must be non-empty.+--+-- If multiple 'Async's complete or have completed, then the value+-- returned corresponds to the first completed 'Async' in the list.+--+{-# INLINE waitAnyCatch #-}+waitAnyCatch :: [Async a] -> IO (Async a, Either SomeException a)+waitAnyCatch = atomically . waitAnyCatchSTM++-- | A version of 'waitAnyCatch' that can be used inside an STM transaction.+--+-- @since 2.1.0+waitAnyCatchSTM :: [Async a] -> STM (Async a, Either SomeException a)+waitAnyCatchSTM [] =+ throwSTM $ ErrorCall+ "waitAnyCatchSTM: invalid argument: input list must be non-empty"+waitAnyCatchSTM asyncs =+ foldr orElse retry $+ map (\a -> do r <- waitCatchSTM a; return (a, r)) asyncs++-- | Like 'waitAnyCatch', but also cancels the other asynchronous+-- operations as soon as one has completed.+--+waitAnyCatchCancel :: [Async a] -> IO (Async a, Either SomeException a)+waitAnyCatchCancel asyncs =+ waitAnyCatch asyncs `finally` cancelMany asyncs++-- | Wait for any of the supplied @Async@s to complete. If the first+-- to complete throws an exception, then that exception is re-thrown+-- by 'waitAny'.+-- The input list must be non-empty.+--+-- If multiple 'Async's complete or have completed, then the value+-- returned corresponds to the first completed 'Async' in the list.+--+{-# INLINE waitAny #-}+waitAny :: [Async a] -> IO (Async a, a)+waitAny = atomically . waitAnySTM++-- | A version of 'waitAny' that can be used inside an STM transaction.+--+-- @since 2.1.0+waitAnySTM :: [Async a] -> STM (Async a, a)+waitAnySTM [] =+ throwSTM $ ErrorCall+ "waitAnySTM: invalid argument: input list must be non-empty"+waitAnySTM asyncs =+ foldr orElse retry $+ map (\a -> do r <- waitSTM a; return (a, r)) asyncs++-- | Like 'waitAny', but also cancels the other asynchronous+-- operations as soon as one has completed.+--+waitAnyCancel :: [Async a] -> IO (Async a, a)+waitAnyCancel asyncs =+ waitAny asyncs `finally` cancelMany asyncs++-- | Wait for the first of two @Async@s to finish.+{-# INLINE waitEitherCatch #-}+waitEitherCatch :: Async a -> Async b+ -> IO (Either (Either SomeException a)+ (Either SomeException b))+waitEitherCatch left right =+ tryAgain $ atomically (waitEitherCatchSTM left right)+ where+ -- See: https://github.com/simonmar/async/issues/14+ tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f++-- | A version of 'waitEitherCatch' that can be used inside an STM transaction.+--+-- @since 2.1.0+waitEitherCatchSTM :: Async a -> Async b+ -> STM (Either (Either SomeException a)+ (Either SomeException b))+waitEitherCatchSTM left right =+ (Left <$> waitCatchSTM left)+ `orElse`+ (Right <$> waitCatchSTM right)++-- | Like 'waitEitherCatch', but also 'cancel's both @Async@s before+-- returning.+--+waitEitherCatchCancel :: Async a -> Async b+ -> IO (Either (Either SomeException a)+ (Either SomeException b))+waitEitherCatchCancel left right =+ waitEitherCatch left right `finally` cancelMany [() <$ left, () <$ right]++-- | Wait for the first of two @Async@s to finish. If the @Async@+-- that finished first raised an exception, then the exception is+-- re-thrown by 'waitEither'.+--+{-# INLINE waitEither #-}+waitEither :: Async a -> Async b -> IO (Either a b)+waitEither left right = atomically (waitEitherSTM left right)++-- | A version of 'waitEither' that can be used inside an STM transaction.+--+-- @since 2.1.0+waitEitherSTM :: Async a -> Async b -> STM (Either a b)+waitEitherSTM left right =+ (Left <$> waitSTM left)+ `orElse`+ (Right <$> waitSTM right)++-- | Like 'waitEither', but the result is ignored.+--+{-# INLINE waitEither_ #-}+waitEither_ :: Async a -> Async b -> IO ()+waitEither_ left right = atomically (waitEitherSTM_ left right)++-- | A version of 'waitEither_' that can be used inside an STM transaction.+--+-- @since 2.1.0+waitEitherSTM_:: Async a -> Async b -> STM ()+waitEitherSTM_ left right =+ (void $ waitSTM left)+ `orElse`+ (void $ waitSTM right)++-- | Like 'waitEither', but also 'cancel's both @Async@s before+-- returning.+--+waitEitherCancel :: Async a -> Async b -> IO (Either a b)+waitEitherCancel left right =+ waitEither left right `finally` cancelMany [() <$ left, () <$ right]++-- | Waits for both @Async@s to finish, but if either of them throws+-- an exception before they have both finished, then the exception is+-- re-thrown by 'waitBoth'.+--+{-# INLINE waitBoth #-}+waitBoth :: Async a -> Async b -> IO (a,b)+waitBoth left right = tryAgain $ atomically (waitBothSTM left right)+ where+ -- See: https://github.com/simonmar/async/issues/14+ tryAgain f = f `catch` \BlockedIndefinitelyOnSTM -> f++-- | A version of 'waitBoth' that can be used inside an STM transaction.+--+-- @since 2.1.0+waitBothSTM :: Async a -> Async b -> STM (a,b)+waitBothSTM left right = do+ a <- waitSTM left+ `orElse`+ (waitSTM right >> retry)+ b <- waitSTM right+ return (a,b)+++-- -----------------------------------------------------------------------------+-- Linking threads++data ExceptionInLinkedThread =+ forall a . ExceptionInLinkedThread (Async a) SomeException+#if __GLASGOW_HASKELL__ < 710+ deriving Typeable+#endif++instance Show ExceptionInLinkedThread where+ showsPrec p (ExceptionInLinkedThread (Async t _) e) =+ showParen (p >= 11) $+ showString "ExceptionInLinkedThread " .+ showsPrec 11 t .+ showString " " .+ showsPrec 11 e++instance Exception ExceptionInLinkedThread where+#if __GLASGOW_HASKELL__ >= 708+ fromException = asyncExceptionFromException+ toException = asyncExceptionToException+#endif++-- | Link the given @Async@ to the current thread, such that if the+-- @Async@ raises an exception, that exception will be re-thrown in+-- the current thread, wrapped in 'ExceptionInLinkedThread'.+--+-- 'link' ignores 'AsyncCancelled' exceptions thrown in the other thread,+-- so that it's safe to 'cancel' a thread you're linked to. If you want+-- different behaviour, use 'linkOnly'.+--+link :: Async a -> IO ()+link = linkOnly (not . isCancel)++-- | Link the given @Async@ to the current thread, such that if the+-- @Async@ raises an exception, that exception will be re-thrown in+-- the current thread, wrapped in 'ExceptionInLinkedThread'.+--+-- The supplied predicate determines which exceptions in the target+-- thread should be propagated to the source thread.+--+linkOnly+ :: (SomeException -> Bool) -- ^ return 'True' if the exception+ -- should be propagated, 'False'+ -- otherwise.+ -> Async a+ -> IO ()+linkOnly shouldThrow a = do+ me <- myThreadId+ void $ forkRepeat $ do+ r <- waitCatch a+ case r of+ Left e | shouldThrow e -> throwTo me (ExceptionInLinkedThread a e)+ _otherwise -> return ()++-- | Link two @Async@s together, such that if either raises an+-- exception, the same exception is re-thrown in the other @Async@,+-- wrapped in 'ExceptionInLinkedThread'.+--+-- 'link2' ignores 'AsyncCancelled' exceptions, so that it's possible+-- to 'cancel' either thread without cancelling the other. If you+-- want different behaviour, use 'link2Only'.+--+link2 :: Async a -> Async b -> IO ()+link2 = link2Only (not . isCancel)++-- | Link two @Async@s together, such that if either raises an+-- exception, the same exception is re-thrown in the other @Async@,+-- wrapped in 'ExceptionInLinkedThread'.+--+-- The supplied predicate determines which exceptions in the target+-- thread should be propagated to the source thread.+--+link2Only :: (SomeException -> Bool) -> Async a -> Async b -> IO ()+link2Only shouldThrow left@(Async tl _) right@(Async tr _) =+ void $ forkRepeat $ do+ r <- waitEitherCatch left right+ case r of+ Left (Left e) | shouldThrow e ->+ throwTo tr (ExceptionInLinkedThread left e)+ Right (Left e) | shouldThrow e ->+ throwTo tl (ExceptionInLinkedThread right e)+ _ -> return ()++isCancel :: SomeException -> Bool+isCancel e+ | Just AsyncCancelled <- fromException e = True+ | otherwise = False+++-- -----------------------------------------------------------------------------++-- | Run two @IO@ actions concurrently, and return the first to+-- finish. The loser of the race is 'cancel'led.+--+-- > race left right =+-- > withAsync left $ \a ->+-- > withAsync right $ \b ->+-- > waitEither a b+--+race ::+ IO a -> IO b -> IO (Either a b)++-- | Like 'race', but the result is ignored.+--+race_ ::+ IO a -> IO b -> IO ()+++-- | Run two @IO@ actions concurrently, and return both results. If+-- either action throws an exception at any time, then the other+-- action is 'cancel'led, and the exception is re-thrown by+-- 'concurrently'.+--+-- > concurrently left right =+-- > withAsync left $ \a ->+-- > withAsync right $ \b ->+-- > waitBoth a b+--+-- To run more than two actions concurrently, see 'mapConcurrently'.+concurrently ::+ IO a -> IO b -> IO (a,b)+++-- | Run two @IO@ actions concurrently. If both of them end with @Right@,+-- return both results. If one of then ends with @Left@, interrupt the other+-- action and return the @Left@.+--+concurrentlyE ::+ IO (Either e a) -> IO (Either e b) -> IO (Either e (a, b))++-- | 'concurrently', but ignore the result values+--+-- @since 2.1.1+concurrently_ ::+ IO a -> IO b -> IO ()++#define USE_ASYNC_VERSIONS 0++#if USE_ASYNC_VERSIONS++race left right =+ withAsync left $ \a ->+ withAsync right $ \b ->+ waitEither a b++race_ left right = void $ race left right++concurrently left right =+ withAsync left $ \a ->+ withAsync right $ \b ->+ waitBoth a b++concurrently_ left right = void $ concurrently left right++#else++-- MVar versions of race/concurrently+-- More ugly than the Async versions, but quite a bit faster.++-- race :: IO a -> IO b -> IO (Either a b)+race left right = concurrently' left right collect+ where+ collect m = do+ e <- m+ case e of+ Left ex -> throwIO ex+ Right r -> return r++-- race_ :: IO a -> IO b -> IO ()+race_ left right = void $ race left right++-- concurrently :: IO a -> IO b -> IO (a,b)+concurrently left right = concurrently' left right (collect [])+ where+ collect [Left a, Right b] _ = return (a,b)+ collect [Right b, Left a] _ = return (a,b)+ collect xs m = do+ e <- m+ case e of+ Left ex -> throwIO ex+ Right r -> collect (r:xs) m++-- concurrentlyE :: IO (Either e a) -> IO (Either e b) -> IO (Either e (a, b))+concurrentlyE left right = concurrently' left right (collect [])+ where+ collect [Left (Right a), Right (Right b)] _ = return $ Right (a,b)+ collect [Right (Right b), Left (Right a)] _ = return $ Right (a,b)+ collect (Left (Left ea):_) _ = return $ Left ea+ collect (Right (Left eb):_) _ = return $ Left eb+ collect xs m = do+ e <- m+ case e of+ Left ex -> throwIO ex+ Right r -> collect (r:xs) m++concurrently' ::+ IO a -> IO b+ -> (IO (Either SomeException (Either a b)) -> IO r)+ -> IO r+concurrently' left right collect = do+ done <- newEmptyMVar+ mask $ \restore -> do+ -- Note: uninterruptibleMask here is because we must not allow+ -- the putMVar in the exception handler to be interrupted,+ -- otherwise the parent thread will deadlock when it waits for+ -- the thread to terminate.+ lid <- forkIO $ uninterruptibleMask_ $+ restore (left >>= putMVar done . Right . Left)+ `catchAll` (putMVar done . Left)+ rid <- forkIO $ uninterruptibleMask_ $+ restore (right >>= putMVar done . Right . Right)+ `catchAll` (putMVar done . Left)++ count <- newIORef (2 :: Int)+ let takeDone = do+ r <- takeMVar done -- interruptible+ -- Decrement the counter so we know how many takes are left.+ -- Since only the parent thread is calling this, we can+ -- use non-atomic modifications.+ -- NB. do this *after* takeMVar, because takeMVar might be+ -- interrupted.+ modifyIORef count (subtract 1)+ return r++ let tryAgain f = f `catch` \BlockedIndefinitelyOnMVar -> f++ stop = do+ -- kill right before left, to match the semantics of+ -- the version using withAsync. (#27)+ uninterruptibleMask_ $ do+ count' <- readIORef count+ -- we only need to use killThread if there are still+ -- children alive. Note: forkIO here is because the+ -- child thread could be in an uninterruptible+ -- putMVar.+ when (count' > 0) $+ void $ forkIO $ do+ throwTo rid AsyncCancelled+ throwTo lid AsyncCancelled+ -- ensure the children are really dead+ replicateM_ count' (tryAgain $ takeMVar done)++ r <- collect (tryAgain takeDone) `onException` stop+ stop+ return r++concurrently_ left right = concurrently' left right (collect 0)+ where+ collect 2 _ = return ()+ collect i m = do+ e <- m+ case e of+ Left ex -> throwIO ex+ Right _ -> collect (i + 1 :: Int) m+++#endif++-- | Maps an 'IO'-performing function over any 'Traversable' data+-- type, performing all the @IO@ actions concurrently, and returning+-- the original data structure with the arguments replaced by the+-- results.+--+-- If any of the actions throw an exception, then all other actions are+-- cancelled and the exception is re-thrown.+--+-- For example, @mapConcurrently@ works with lists:+--+-- > pages <- mapConcurrently getURL ["url1", "url2", "url3"]+--+-- If you just have a list of actions, run them concurrently with+--+-- > results <- mapConcurrently id [act1, act2, act3]+--+-- NOTE: @mapConcurrently@ will immediately spawn a thread for each+-- element of the @Traversable@, so running this on large inputs can+-- lead to resource exhaustion (of memory, file descriptors, or other+-- limited resources). To avoid unbounded resource usage, see+-- "Control.Concurrent.Stream".+mapConcurrently ::+ Traversable t => (a -> IO b) -> t a -> IO (t b)+mapConcurrently f = runConcurrently . traverse (Concurrently . f)++-- | `forConcurrently` is `mapConcurrently` with its arguments flipped+--+-- > pages <- forConcurrently ["url1", "url2", "url3"] $ \url -> getURL url+--+-- @since 2.1.0+forConcurrently ::+ Traversable t => t a -> (a -> IO b) -> IO (t b)+forConcurrently = flip mapConcurrently++-- | `mapConcurrently_` is `mapConcurrently` with the return value discarded;+-- a concurrent equivalent of 'mapM_'.+mapConcurrently_ ::+ F.Foldable f => (a -> IO b) -> f a -> IO ()+mapConcurrently_ f = runConcurrently . F.foldMap (Concurrently . void . f)++-- | `forConcurrently_` is `forConcurrently` with the return value discarded;+-- a concurrent equivalent of 'forM_'.+forConcurrently_ ::+ F.Foldable f => f a -> (a -> IO b) -> IO ()+forConcurrently_ = flip mapConcurrently_++-- | Perform the action in the given number of threads.+--+-- @since 2.1.1+replicateConcurrently ::+ Int -> IO a -> IO [a]+replicateConcurrently cnt = runConcurrently . replicateM cnt . Concurrently++-- | Same as 'replicateConcurrently', but ignore the results.+--+-- @since 2.1.1+replicateConcurrently_ ::+ Int -> IO a -> IO ()+replicateConcurrently_ cnt = runConcurrently . F.fold . replicate cnt . Concurrently . void++-- -----------------------------------------------------------------------------++-- | A value of type @Concurrently a@ is an @IO@ operation that can be+-- composed with other @Concurrently@ values, using the @Applicative@+-- and @Alternative@ instances.+--+-- Calling @runConcurrently@ on a value of type @Concurrently a@ will+-- execute the @IO@ operations it contains concurrently, before+-- delivering the result of type @a@.+--+-- For example+--+-- > (page1, page2, page3)+-- > <- runConcurrently $ (,,)+-- > <$> Concurrently (getURL "url1")+-- > <*> Concurrently (getURL "url2")+-- > <*> Concurrently (getURL "url3")+--+newtype Concurrently a = Concurrently { runConcurrently :: IO a }++instance Functor Concurrently where+ fmap f (Concurrently a) = Concurrently $ f <$> a++instance Applicative Concurrently where+ pure = Concurrently . return+ Concurrently fs <*> Concurrently as =+ Concurrently $ (\(f, a) -> f a) <$> concurrently fs as++-- | 'Control.Alternative.empty' waits forever. 'Control.Alternative.<|>' returns the first to finish and 'cancel's the other.+instance Alternative Concurrently where+ empty = Concurrently $ forever (threadDelay maxBound)+ Concurrently as <|> Concurrently bs =+ Concurrently $ either id id <$> race as bs++#if MIN_VERSION_base(4,9,0)+-- | Only defined by @async@ for @base >= 4.9@+--+-- @since 2.1.0+instance Semigroup a => Semigroup (Concurrently a) where+ (<>) = liftA2 (<>)++-- | @since 2.1.0+instance (Semigroup a, Monoid a) => Monoid (Concurrently a) where+ mempty = pure mempty+ mappend = (<>)+#else+-- | @since 2.1.0+instance Monoid a => Monoid (Concurrently a) where+ mempty = pure mempty+ mappend = liftA2 mappend+#endif++-- | A value of type @ConcurrentlyE e a@ is an @IO@ operation that can be+-- composed with other @ConcurrentlyE@ values, using the @Applicative@ instance.+--+-- Calling @runConcurrentlyE@ on a value of type @ConcurrentlyE e a@ will+-- execute the @IO@ operations it contains concurrently, before delivering+-- either the result of type @a@, or an error of type @e@ if one of the actions+-- returns @Left@.+--+-- | @since 2.2.5+newtype ConcurrentlyE e a = ConcurrentlyE { runConcurrentlyE :: IO (Either e a) }++instance Functor (ConcurrentlyE e) where+ fmap f (ConcurrentlyE ea) = ConcurrentlyE $ fmap (fmap f) ea++#if MIN_VERSION_base(4,8,0)+instance Bifunctor ConcurrentlyE where+ bimap f g (ConcurrentlyE ea) = ConcurrentlyE $ fmap (bimap f g) ea+#endif++instance Applicative (ConcurrentlyE e) where+ pure = ConcurrentlyE . return . return+ ConcurrentlyE fs <*> ConcurrentlyE eas =+ ConcurrentlyE $ fmap (\(f, a) -> f a) <$> concurrentlyE fs eas++#if MIN_VERSION_base(4,9,0)+-- | Either the combination of the successful results, or the first failure.+instance Semigroup a => Semigroup (ConcurrentlyE e a) where+ (<>) = liftA2 (<>)++instance (Semigroup a, Monoid a) => Monoid (ConcurrentlyE e a) where+ mempty = pure mempty+ mappend = (<>)+#endif++-- ----------------------------------------------------------------------------++-- | Fork a thread that runs the supplied action, and if it raises an+-- exception, re-runs the action. The thread terminates only when the+-- action runs to completion without raising an exception.+forkRepeat ::+ IO a -> IO ThreadId+forkRepeat action =+ mask $ \restore ->+ let go = do r <- tryAll (restore action)+ case r of+ Left _ -> go+ _ -> return ()+ in forkIO (debugLabelMe >> go)++catchAll :: IO a -> (SomeException -> IO a) -> IO a+catchAll = catch++tryAll :: IO a -> IO (Either SomeException a)+tryAll = try++-- A version of forkIO that does not include the outer exception+-- handler: saves a bit of time when we will be installing our own+-- exception handler.+{-# INLINE rawForkIO #-}+rawForkIO ::+ IO () -> IO ThreadId+rawForkIO action = IO $ \ s ->+ case fork# action_plus s of (# s1, tid #) -> (# s1, ThreadId tid #)+ where+ (IO action_plus) = debugLabelMe >> action++{-# INLINE rawForkOn #-}+rawForkOn ::+ Int -> IO () -> IO ThreadId+rawForkOn (I# cpu) action = IO $ \ s ->+ case forkOn# cpu action_plus s of (# s1, tid #) -> (# s1, ThreadId tid #)+ where+ (IO action_plus) = debugLabelMe >> action++debugLabelMe ::+ IO ()+debugLabelMe =+ pure ()
− vendor/stm-2.5.0.1/Control/Concurrent/STM/TMVar.hs
@@ -1,168 +0,0 @@-{-# LANGUAGE CPP, DeriveDataTypeable, MagicHash, UnboxedTuples #-}--#if __GLASGOW_HASKELL__ >= 701-{-# LANGUAGE Trustworthy #-}-#endif--{-# OPTIONS -fno-warn-implicit-prelude #-}---------------------------------------------------------------------------------- |--- Module : Control.Concurrent.STM.TMVar--- Copyright : (c) The University of Glasgow 2004--- License : BSD-style (see the file libraries/base/LICENSE)------ Maintainer : libraries@haskell.org--- Stability : experimental--- Portability : non-portable (requires STM)------ TMVar: Transactional MVars, for use in the STM monad--- (GHC only)-----------------------------------------------------------------------------------module Control.Concurrent.STM.TMVar (-#ifdef __GLASGOW_HASKELL__- -- * TMVars- TMVar,- newTMVar,- newEmptyTMVar,- newTMVarIO,- newEmptyTMVarIO,- takeTMVar,- putTMVar,- readTMVar,- tryReadTMVar,- swapTMVar,- tryTakeTMVar,- tryPutTMVar,- isEmptyTMVar,- mkWeakTMVar-#endif- ) where--#ifdef __GLASGOW_HASKELL__-import GHC.Base-import GHC.Conc-import GHC.Weak--import Data.Typeable (Typeable)--newtype TMVar a = TMVar (TVar (Maybe a)) deriving (Eq, Typeable)-{- ^-A 'TMVar' is a synchronising variable, used-for communication between concurrent threads. It can be thought of-as a box, which may be empty or full.--}---- |Create a 'TMVar' which contains the supplied value.-newTMVar :: a -> STM (TMVar a)-newTMVar a = do- t <- newTVar (Just a)- return (TMVar t)---- |@IO@ version of 'newTMVar'. This is useful for creating top-level--- 'TMVar's using 'System.IO.Unsafe.unsafePerformIO', because using--- 'atomically' inside 'System.IO.Unsafe.unsafePerformIO' isn't--- possible.-newTMVarIO :: a -> IO (TMVar a)-newTMVarIO a = do- t <- newTVarIO (Just a)- return (TMVar t)---- |Create a 'TMVar' which is initially empty.-newEmptyTMVar :: STM (TMVar a)-newEmptyTMVar = do- t <- newTVar Nothing- return (TMVar t)---- |@IO@ version of 'newEmptyTMVar'. This is useful for creating top-level--- 'TMVar's using 'System.IO.Unsafe.unsafePerformIO', because using--- 'atomically' inside 'System.IO.Unsafe.unsafePerformIO' isn't--- possible.-newEmptyTMVarIO :: IO (TMVar a)-newEmptyTMVarIO = do- t <- newTVarIO Nothing- return (TMVar t)---- |Return the contents of the 'TMVar'. If the 'TMVar' is currently--- empty, the transaction will 'retry'. After a 'takeTMVar',--- the 'TMVar' is left empty.-takeTMVar :: TMVar a -> STM a-takeTMVar (TMVar t) = do- m <- readTVar t- case m of- Nothing -> retry- Just a -> do writeTVar t Nothing; return a---- | A version of 'takeTMVar' that does not 'retry'. The 'tryTakeTMVar'--- function returns 'Nothing' if the 'TMVar' was empty, or @'Just' a@ if--- the 'TMVar' was full with contents @a@. After 'tryTakeTMVar', the--- 'TMVar' is left empty.-tryTakeTMVar :: TMVar a -> STM (Maybe a)-tryTakeTMVar (TMVar t) = do- m <- readTVar t- case m of- Nothing -> return Nothing- Just a -> do writeTVar t Nothing; return (Just a)---- |Put a value into a 'TMVar'. If the 'TMVar' is currently full,--- 'putTMVar' will 'retry'.-putTMVar :: TMVar a -> a -> STM ()-putTMVar (TMVar t) a = do- m <- readTVar t- case m of- Nothing -> do writeTVar t (Just a); return ()- Just _ -> retry---- | A version of 'putTMVar' that does not 'retry'. The 'tryPutTMVar'--- function attempts to put the value @a@ into the 'TMVar', returning--- 'True' if it was successful, or 'False' otherwise.-tryPutTMVar :: TMVar a -> a -> STM Bool-tryPutTMVar (TMVar t) a = do- m <- readTVar t- case m of- Nothing -> do writeTVar t (Just a); return True- Just _ -> return False---- | This is a combination of 'takeTMVar' and 'putTMVar'; ie. it--- takes the value from the 'TMVar', puts it back, and also returns--- it.-readTMVar :: TMVar a -> STM a-readTMVar (TMVar t) = do- m <- readTVar t- case m of- Nothing -> retry- Just a -> return a---- | A version of 'readTMVar' which does not retry. Instead it--- returns @Nothing@ if no value is available.------ @since 2.3-tryReadTMVar :: TMVar a -> STM (Maybe a)-tryReadTMVar (TMVar t) = readTVar t---- |Swap the contents of a 'TMVar' for a new value.-swapTMVar :: TMVar a -> a -> STM a-swapTMVar (TMVar t) new = do- m <- readTVar t- case m of- Nothing -> retry- Just old -> do writeTVar t (Just new); return old---- |Check whether a given 'TMVar' is empty.-isEmptyTMVar :: TMVar a -> STM Bool-isEmptyTMVar (TMVar t) = do- m <- readTVar t- case m of- Nothing -> return True- Just _ -> return False---- | Make a 'Weak' pointer to a 'TMVar', using the second argument as--- a finalizer to run when the 'TMVar' is garbage-collected.------ @since 2.4.4-mkWeakTMVar :: TMVar a -> IO () -> IO (Weak (TMVar a))-mkWeakTMVar tmv@(TMVar (TVar t#)) (IO finalizer) = IO $ \s ->- case mkWeak# t# tmv finalizer s of (# s1, w #) -> (# s1, Weak w #)-#endif
+ vendor/stm-2.5.3.1/Control/Concurrent/STM/TMVar.hs view
@@ -0,0 +1,176 @@+{-# LANGUAGE CPP, DeriveDataTypeable, MagicHash, UnboxedTuples #-}++#if __GLASGOW_HASKELL__ >= 701+{-# LANGUAGE Trustworthy #-}+#endif++{-# OPTIONS -fno-warn-implicit-prelude #-}++-----------------------------------------------------------------------------+-- |+-- Module : Control.Concurrent.STM.TMVar+-- Copyright : (c) The University of Glasgow 2004+-- License : BSD-style (see the file libraries/base/LICENSE)+--+-- Maintainer : libraries@haskell.org+-- Stability : experimental+-- Portability : non-portable (requires STM)+--+-- TMVar: Transactional MVars, for use in the STM monad+-- (GHC only)+--+-----------------------------------------------------------------------------++module Control.Concurrent.STM.TMVar (+#ifdef __GLASGOW_HASKELL__+ -- * TMVars+ TMVar,+ newTMVar,+ newEmptyTMVar,+ newTMVarIO,+ newEmptyTMVarIO,+ takeTMVar,+ putTMVar,+ readTMVar,+ writeTMVar,+ tryReadTMVar,+ swapTMVar,+ tryTakeTMVar,+ tryPutTMVar,+ isEmptyTMVar,+ mkWeakTMVar+#endif+ ) where++#ifdef __GLASGOW_HASKELL__+import GHC.Base+import GHC.Conc+import GHC.Weak++import Data.Typeable (Typeable)++newtype TMVar a = TMVar (TVar (Maybe a)) deriving (Eq, Typeable)+{- ^+A 'TMVar' is a synchronising variable, used+for communication between concurrent threads. It can be thought of+as a box, which may be empty or full.+-}++-- |Create a 'TMVar' which contains the supplied value.+newTMVar :: a -> STM (TMVar a)+newTMVar a = do+ t <- newTVar (Just a)+ return (TMVar t)++-- |@IO@ version of 'newTMVar'. This is useful for creating top-level+-- 'TMVar's using 'System.IO.Unsafe.unsafePerformIO', because using+-- 'atomically' inside 'System.IO.Unsafe.unsafePerformIO' isn't+-- possible.+newTMVarIO :: a -> IO (TMVar a)+newTMVarIO a = do+ t <- newTVarIO (Just a)+ return (TMVar t)++-- |Create a 'TMVar' which is initially empty.+newEmptyTMVar :: STM (TMVar a)+newEmptyTMVar = do+ t <- newTVar Nothing+ return (TMVar t)++-- |@IO@ version of 'newEmptyTMVar'. This is useful for creating top-level+-- 'TMVar's using 'System.IO.Unsafe.unsafePerformIO', because using+-- 'atomically' inside 'System.IO.Unsafe.unsafePerformIO' isn't+-- possible.+newEmptyTMVarIO :: IO (TMVar a)+newEmptyTMVarIO = do+ t <- newTVarIO Nothing+ return (TMVar t)++-- |Return the contents of the 'TMVar'. If the 'TMVar' is currently+-- empty, the transaction will 'retry'. After a 'takeTMVar',+-- the 'TMVar' is left empty.+takeTMVar :: TMVar a -> STM a+takeTMVar (TMVar t) = do+ m <- readTVar t+ case m of+ Nothing -> retry+ Just a -> do writeTVar t Nothing; return a++-- | A version of 'takeTMVar' that does not 'retry'. The 'tryTakeTMVar'+-- function returns 'Nothing' if the 'TMVar' was empty, or @'Just' a@ if+-- the 'TMVar' was full with contents @a@. After 'tryTakeTMVar', the+-- 'TMVar' is left empty.+tryTakeTMVar :: TMVar a -> STM (Maybe a)+tryTakeTMVar (TMVar t) = do+ m <- readTVar t+ case m of+ Nothing -> return Nothing+ Just a -> do writeTVar t Nothing; return (Just a)++-- |Put a value into a 'TMVar'. If the 'TMVar' is currently full,+-- 'putTMVar' will 'retry'.+putTMVar :: TMVar a -> a -> STM ()+putTMVar (TMVar t) a = do+ m <- readTVar t+ case m of+ Nothing -> do writeTVar t (Just a); return ()+ Just _ -> retry++-- | A version of 'putTMVar' that does not 'retry'. The 'tryPutTMVar'+-- function attempts to put the value @a@ into the 'TMVar', returning+-- 'True' if it was successful, or 'False' otherwise.+tryPutTMVar :: TMVar a -> a -> STM Bool+tryPutTMVar (TMVar t) a = do+ m <- readTVar t+ case m of+ Nothing -> do writeTVar t (Just a); return True+ Just _ -> return False++-- | This is a combination of 'takeTMVar' and 'putTMVar'; ie. it+-- takes the value from the 'TMVar', puts it back, and also returns+-- it.+readTMVar :: TMVar a -> STM a+readTMVar (TMVar t) = do+ m <- readTVar t+ case m of+ Nothing -> retry+ Just a -> return a++-- | A version of 'readTMVar' which does not retry. Instead it+-- returns @Nothing@ if no value is available.+--+-- @since 2.3+tryReadTMVar :: TMVar a -> STM (Maybe a)+tryReadTMVar (TMVar t) = readTVar t++-- |Swap the contents of a 'TMVar' for a new value.+swapTMVar :: TMVar a -> a -> STM a+swapTMVar (TMVar t) new = do+ m <- readTVar t+ case m of+ Nothing -> retry+ Just old -> do writeTVar t (Just new); return old++-- | Non-blocking write of a new value to a 'TMVar'+-- Puts if empty. Replaces if populated.+--+-- @since 2.5.1+writeTMVar :: TMVar a -> a -> STM ()+writeTMVar (TMVar t) new = writeTVar t (Just new)++-- |Check whether a given 'TMVar' is empty.+isEmptyTMVar :: TMVar a -> STM Bool+isEmptyTMVar (TMVar t) = do+ m <- readTVar t+ case m of+ Nothing -> return True+ Just _ -> return False++-- | Make a 'Weak' pointer to a 'TMVar', using the second argument as+-- a finalizer to run when the 'TMVar' is garbage-collected.+--+-- @since 2.4.4+mkWeakTMVar :: TMVar a -> IO () -> IO (Weak (TMVar a))+mkWeakTMVar tmv@(TMVar (TVar t#)) (IO finalizer) = IO $ \s ->+ case mkWeak# t# tmv finalizer s of (# s1, w #) -> (# s1, Weak w #)+#endif