packages feed

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