broadcast-chan 0.2.1.2 → 0.3.0
raw patch · 20 files changed
+1530/−43 lines, 20 filesdep +clockdep +conduitdep +containersdep ~asyncdep ~unliftio-corePVP ok
version bump matches the API change (PVP)
Dependencies added: clock, conduit, containers, foldl, monad-loops, optparse-applicative, paramtree, pipes, pipes-safe, random, resourcet, tagged, tasty, tasty-golden, tasty-hunit, temporary, text
Dependency ranges changed: async, unliftio-core
API changes (from Hackage documentation)
- BroadcastChan: Drop :: Action
- BroadcastChan: Handle :: (a -> SomeException -> m Action) -> Handler m a
- BroadcastChan: In :: Direction
- BroadcastChan: Out :: Direction
- BroadcastChan: Retry :: Action
- BroadcastChan: Simple :: Action -> Handler m a
- BroadcastChan: Terminate :: Action
- BroadcastChan: closeBChan :: MonadIO m => BroadcastChan In a -> m Bool
- BroadcastChan: data Action
- BroadcastChan: data BroadcastChan (dir :: Direction) a
- BroadcastChan: data Direction
- BroadcastChan: data Handler m a
- BroadcastChan: foldBChan :: (MonadIO m, MonadIO n) => (x -> a -> x) -> x -> (x -> b) -> BroadcastChan d a -> n (m b)
- BroadcastChan: foldBChanM :: (MonadIO m, MonadIO n) => (x -> a -> m x) -> m x -> (x -> m b) -> BroadcastChan d a -> n (m b)
- BroadcastChan: getBChanContents :: BroadcastChan dir a -> IO [a]
- BroadcastChan: isClosedBChan :: MonadIO m => BroadcastChan dir a -> m Bool
- BroadcastChan: newBChanListener :: MonadIO m => BroadcastChan dir a -> m (BroadcastChan Out a)
- BroadcastChan: newBroadcastChan :: MonadIO m => m (BroadcastChan In a)
- BroadcastChan: parFoldMap :: (Foldable f, MonadUnliftIO m) => Handler m a -> Int -> (a -> m b) -> (r -> b -> r) -> r -> f a -> m r
- BroadcastChan: parFoldMapM :: forall a b f m r. (Foldable f, MonadUnliftIO m) => Handler m a -> Int -> (a -> m b) -> (r -> b -> m r) -> r -> f a -> m r
- BroadcastChan: parMapM_ :: (Foldable f, MonadUnliftIO m) => Handler m a -> Int -> (a -> m ()) -> f a -> m ()
- BroadcastChan: readBChan :: MonadIO m => BroadcastChan Out a -> m (Maybe a)
- BroadcastChan: type In = 'In
- BroadcastChan: type Out = 'Out
- BroadcastChan: writeBChan :: MonadIO m => BroadcastChan In a -> a -> m Bool
- BroadcastChan.Extra: Bracket :: IO [Weak ThreadId] -> ([Weak ThreadId] -> IO ()) -> m r -> BracketOnError m r
- BroadcastChan.Extra: Drop :: Action
- BroadcastChan.Extra: Handle :: (a -> SomeException -> m Action) -> Handler m a
- BroadcastChan.Extra: Retry :: Action
- BroadcastChan.Extra: Simple :: Action -> Handler m a
- BroadcastChan.Extra: Terminate :: Action
- BroadcastChan.Extra: ThreadBracket :: IO () -> IO () -> IO () -> ThreadBracket
- BroadcastChan.Extra: [action] :: BracketOnError m r -> m r
- BroadcastChan.Extra: [allocate] :: BracketOnError m r -> IO [Weak ThreadId]
- BroadcastChan.Extra: [cleanupForkError] :: ThreadBracket -> IO ()
- BroadcastChan.Extra: [cleanupFork] :: ThreadBracket -> IO ()
- BroadcastChan.Extra: [cleanup] :: BracketOnError m r -> [Weak ThreadId] -> IO ()
- BroadcastChan.Extra: [setupFork] :: ThreadBracket -> IO ()
- BroadcastChan.Extra: data Action
- BroadcastChan.Extra: data BracketOnError m r
- BroadcastChan.Extra: data Handler m a
- BroadcastChan.Extra: data ThreadBracket
- BroadcastChan.Extra: instance GHC.Classes.Eq BroadcastChan.Extra.Action
- BroadcastChan.Extra: instance GHC.Exception.Type.Exception BroadcastChan.Extra.Shutdown
- BroadcastChan.Extra: instance GHC.Show.Show BroadcastChan.Extra.Action
- BroadcastChan.Extra: instance GHC.Show.Show BroadcastChan.Extra.Shutdown
- BroadcastChan.Extra: mapHandler :: (m Action -> n Action) -> Handler m a -> Handler n a
- BroadcastChan.Extra: runParallel :: forall a b m n r. (MonadIO m, MonadIO n) => Either (b -> n r) (r -> b -> n r) -> Handler IO a -> Int -> (a -> IO b) -> ((a -> m ()) -> (a -> m (Maybe b)) -> n r) -> n (BracketOnError n r)
- BroadcastChan.Extra: runParallelWith :: forall a b m n r. (MonadIO m, MonadIO n) => ThreadBracket -> Either (b -> n r) (r -> b -> n r) -> Handler IO a -> Int -> (a -> IO b) -> ((a -> m ()) -> (a -> m (Maybe b)) -> n r) -> n (BracketOnError n r)
- BroadcastChan.Extra: runParallelWith_ :: (MonadIO m, MonadIO n) => ThreadBracket -> Handler IO a -> Int -> (a -> IO ()) -> ((a -> m ()) -> n r) -> n (BracketOnError n r)
- BroadcastChan.Extra: runParallel_ :: (MonadIO m, MonadIO n) => Handler IO a -> Int -> (a -> IO ()) -> ((a -> m ()) -> n r) -> n (BracketOnError n r)
- BroadcastChan.Prelude: forM_ :: MonadIO m => BroadcastChan Out a -> (a -> m b) -> m ()
- BroadcastChan.Prelude: mapM_ :: MonadIO m => (a -> m b) -> BroadcastChan Out a -> m ()
- BroadcastChan.Throw: Drop :: Action
- BroadcastChan.Throw: Handle :: (a -> SomeException -> m Action) -> Handler m a
- BroadcastChan.Throw: In :: Direction
- BroadcastChan.Throw: Out :: Direction
- BroadcastChan.Throw: ReadFailed :: BChanError
- BroadcastChan.Throw: Retry :: Action
- BroadcastChan.Throw: Simple :: Action -> Handler m a
- BroadcastChan.Throw: Terminate :: Action
- BroadcastChan.Throw: WriteFailed :: BChanError
- BroadcastChan.Throw: closeBChan :: MonadIO m => BroadcastChan In a -> m Bool
- BroadcastChan.Throw: data Action
- BroadcastChan.Throw: data BChanError
- BroadcastChan.Throw: data BroadcastChan (dir :: Direction) a
- BroadcastChan.Throw: data Direction
- BroadcastChan.Throw: data Handler m a
- BroadcastChan.Throw: foldBChan :: (MonadIO m, MonadIO n) => (x -> a -> x) -> x -> (x -> b) -> BroadcastChan d a -> n (m b)
- BroadcastChan.Throw: foldBChanM :: (MonadIO m, MonadIO n) => (x -> a -> m x) -> m x -> (x -> m b) -> BroadcastChan d a -> n (m b)
- BroadcastChan.Throw: getBChanContents :: BroadcastChan dir a -> IO [a]
- BroadcastChan.Throw: instance GHC.Classes.Eq BroadcastChan.Throw.BChanError
- BroadcastChan.Throw: instance GHC.Exception.Type.Exception BroadcastChan.Throw.BChanError
- BroadcastChan.Throw: instance GHC.Read.Read BroadcastChan.Throw.BChanError
- BroadcastChan.Throw: instance GHC.Show.Show BroadcastChan.Throw.BChanError
- BroadcastChan.Throw: isClosedBChan :: MonadIO m => BroadcastChan dir a -> m Bool
- BroadcastChan.Throw: newBChanListener :: MonadIO m => BroadcastChan dir a -> m (BroadcastChan Out a)
- BroadcastChan.Throw: newBroadcastChan :: MonadIO m => m (BroadcastChan In a)
- BroadcastChan.Throw: parFoldMap :: (Foldable f, MonadUnliftIO m) => Handler m a -> Int -> (a -> m b) -> (r -> b -> r) -> r -> f a -> m r
- BroadcastChan.Throw: parFoldMapM :: forall a b f m r. (Foldable f, MonadUnliftIO m) => Handler m a -> Int -> (a -> m b) -> (r -> b -> m r) -> r -> f a -> m r
- BroadcastChan.Throw: parMapM_ :: (Foldable f, MonadUnliftIO m) => Handler m a -> Int -> (a -> m ()) -> f a -> m ()
- BroadcastChan.Throw: readBChan :: BroadcastChan Out a -> IO a
- BroadcastChan.Throw: type In = 'In
- BroadcastChan.Throw: type Out = 'Out
- BroadcastChan.Throw: writeBChan :: BroadcastChan In a -> a -> IO ()
+ BroadcastChan.Test: (@=?) :: (Eq a, Show a, HasCallStack) => a -> a -> Assertion
+ BroadcastChan.Test: (@?) :: IO Bool -> String -> Assertion
+ BroadcastChan.Test: (@?=) :: (Eq a, Show a, HasCallStack) => a -> a -> Assertion
+ BroadcastChan.Test: HUnitFailure :: Maybe SrcLoc -> String -> HUnitFailure
+ BroadcastChan.Test: assert :: Assertable t => t -> Assertion
+ BroadcastChan.Test: assertBool :: HasCallStack => String -> Bool -> Assertion
+ BroadcastChan.Test: assertEqual :: (Eq a, Show a, HasCallStack) => String -> a -> a -> Assertion
+ BroadcastChan.Test: assertFailure :: HasCallStack => String -> IO a
+ BroadcastChan.Test: assertString :: HasCallStack => String -> Assertion
+ BroadcastChan.Test: assertionPredicate :: AssertionPredicable t => t -> IO Bool
+ BroadcastChan.Test: class () => Assertable t
+ BroadcastChan.Test: class () => AssertionPredicable t
+ BroadcastChan.Test: class Monad m => MonadIO (m :: Type -> Type)
+ BroadcastChan.Test: data () => HUnitFailure
+ BroadcastChan.Test: expect :: (Eq e, Exception e) => e -> IO a -> Assertion
+ BroadcastChan.Test: genStreamTests :: (Eq r, Show r) => String -> ([Int] -> (Int -> IO Int) -> IO r) -> (Handler IO Int -> [Int] -> (Int -> IO Int) -> Int -> IO r) -> TestTree
+ BroadcastChan.Test: infix 0 @?
+ BroadcastChan.Test: infix 1 @?=
+ BroadcastChan.Test: instance GHC.Classes.Eq BroadcastChan.Test.SlowTests
+ BroadcastChan.Test: instance GHC.Classes.Eq BroadcastChan.Test.TestException
+ BroadcastChan.Test: instance GHC.Classes.Ord BroadcastChan.Test.SlowTests
+ BroadcastChan.Test: instance GHC.Exception.Type.Exception BroadcastChan.Test.TestException
+ BroadcastChan.Test: instance GHC.Show.Show BroadcastChan.Test.TestException
+ BroadcastChan.Test: instance Test.Tasty.Options.IsOption BroadcastChan.Test.SlowTests
+ BroadcastChan.Test: liftIO :: MonadIO m => IO a -> m a
+ BroadcastChan.Test: mapHandler :: (m Action -> n Action) -> Handler m a -> Handler n a
+ BroadcastChan.Test: runTests :: String -> [TestTree] -> IO ()
+ BroadcastChan.Test: testCase :: TestName -> Assertion -> TestTree
+ BroadcastChan.Test: testCaseInfo :: TestName -> IO String -> TestTree
+ BroadcastChan.Test: testCaseSteps :: TestName -> ((String -> IO ()) -> Assertion) -> TestTree
+ BroadcastChan.Test: type Assertion = IO ()
+ BroadcastChan.Test: type AssertionPredicate = IO Bool
+ BroadcastChan.Test: type HasCallStack = ?callStack :: CallStack
+ BroadcastChan.Test: withLoggedOutput :: FilePath -> (Handle -> IO r) -> IO (r, Text)
Files
- BroadcastChan.hs +8/−1
- BroadcastChan/Extra.hs +18/−6
- BroadcastChan/Internal.hs +50/−0
- BroadcastChan/Prelude.hs +6/−2
- BroadcastChan/Throw.hs +27/−6
- CHANGELOG.md +18/−1
- README.md +1/−3
- benchmarks/Utils.hs +1/−0
- broadcast-chan-test/BroadcastChan/Test.hs +392/−0
- broadcast-chan.cabal +164/−24
- conduit/BroadcastChan/Conduit.hs +42/−0
- conduit/BroadcastChan/Conduit/Internal.hs +117/−0
- conduit/BroadcastChan/Conduit/Throw.hs +43/−0
- pipes/BroadcastChan/Pipes.hs +43/−0
- pipes/BroadcastChan/Pipes/Internal.hs +71/−0
- pipes/BroadcastChan/Pipes/Throw.hs +44/−0
- tests/Basic.hs +366/−0
- tests/ConduitTest.hs +37/−0
- tests/IOTest.hs +47/−0
- tests/PipeTest.hs +35/−0
BroadcastChan.hs view
@@ -5,7 +5,7 @@ ------------------------------------------------------------------------------- -- | -- Module : BroadcastChan--- Copyright : (C) 2014-2021 Merijn Verstraaten+-- Copyright : (C) 2014-2022 Merijn Verstraaten -- License : BSD-style (see the file LICENSE) -- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl> -- Stability : experimental@@ -64,6 +64,7 @@ , newBChanListener -- * Basic Operations , readBChan+ , tryReadBChan , writeBChan , closeBChan , isClosedBChan@@ -104,6 +105,8 @@ -- -- This function does __NOT__ guarantee that elements are processed in a -- deterministic order!+--+-- @since 0.2.0 parMapM_ :: (F.Foldable f, MonadUnliftIO m) => Handler m a@@ -131,6 +134,8 @@ -- -- This function does __NOT__ guarantee that elements are processed in a -- deterministic order!+--+-- @since 0.2.0 parFoldMap :: (F.Foldable f, MonadUnliftIO m) => Handler m a@@ -153,6 +158,8 @@ -- -- This function does __NOT__ guarantee that elements are processed in a -- deterministic order!+--+-- @since 0.2.0 parFoldMapM :: forall a b f m r . (F.Foldable f, MonadUnliftIO m)
BroadcastChan/Extra.hs view
@@ -6,7 +6,7 @@ ------------------------------------------------------------------------------- -- | -- Module : BroadcastChan.Extra--- Copyright : (C) 2014-2021 Merijn Verstraaten+-- Copyright : (C) 2014-2022 Merijn Verstraaten -- License : BSD-style (see the file LICENSE) -- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl> -- Stability : experimental@@ -70,6 +70,8 @@ instance Exception Shutdown -- | Action to take when an exception occurs while processing an element.+--+-- @since 0.2.0 data Action = Drop -- ^ Drop the current element and continue processing.@@ -81,6 +83,8 @@ deriving (Eq, Show) -- | Exception handler for parallel processing.+--+-- @since 0.2.0 data Handler m a = Simple Action -- ^ Always take the specified 'Action'.@@ -90,6 +94,8 @@ -- | Allocation, cleanup, and work actions for parallel processing. These -- should be passed to an appropriate @bracketOnError@ function.+--+-- @since 0.2.0 data BracketOnError m r = Bracket { allocate :: IO [Weak ThreadId]@@ -109,14 +115,14 @@ -- care not to leak resources in these cases. In other words, IFF 'setupFork' -- succeeds then this library will ensure the corresponding cleanup runs. ----- @since 0.2.1+-- @since 0.3.0 data ThreadBracket = ThreadBracket { setupFork :: IO () -- ^ Setup action to run before spawning a new thread. , cleanupFork :: IO () -- ^ Normal cleanup action upon thread termination.- , cleanupForkError :: IO ()+ , cleanupForkError :: SomeException -> IO () -- ^ Exceptional cleanup action in case thread terminates due to an -- exception. }@@ -125,10 +131,12 @@ noopBracket = ThreadBracket { setupFork = return () , cleanupFork = return ()- , cleanupForkError = return ()+ , cleanupForkError = const (return ()) } -- | Convenience function for changing the monad the exception handler runs in.+--+-- @since 0.2.0 mapHandler :: (m Action -> n Action) -> Handler m a -> Handler n a mapHandler _ (Simple act) = Simple act mapHandler mmorph (Handle f) = Handle $ \a exc -> mmorph (f a exc)@@ -183,9 +191,9 @@ signalQSemN shutdownSem 1 case exit of Left exc- | Just Shutdown <- fromException exc -> cleanupForkError+ | Just Shutdown <- fromException exc -> cleanupForkError exc | otherwise -> do- cleanupForkError+ cleanupForkError exc reportErr <- tryTakeMVar excMVar case reportErr of Nothing -> return ()@@ -247,6 +255,8 @@ -- setting up 'Control.Concurrent.forkIO' threads and exception handlers. The -- 'cleanup' action ensures all threads are terminate in case of an exception. -- Finally, 'action' performs the actual parallel processing of elements.+--+-- @since 0.2.0 runParallel :: forall a b m n r . (MonadIO m, MonadIO n)@@ -344,6 +354,8 @@ -- setting up 'Control.Concurrent.forkIO' threads and exception handlers. Th -- 'cleanup' action ensures all threads are terminate in case of an exception. -- Finally, 'action' performs the actual parallel processing of elements.+--+-- @since 0.2.0 runParallel_ :: (MonadIO m, MonadIO n) => Handler IO a
BroadcastChan/Internal.hs view
@@ -11,18 +11,26 @@ -- | Used with DataKinds as phantom type indicating whether a 'BroadcastChan' -- value is a read or write end.+--+-- @since 0.2.0 data Direction = In -- ^ Indicates a write 'BroadcastChan' | Out -- ^ Indicates a read 'BroadcastChan' -- | Alias for the v'In' type from the 'Direction' kind, allows users to write -- the @'BroadcastChan' v'In' a@ type without enabling @DataKinds@.+--+-- @since 0.2.0 type In = 'In -- | Alias for the v'Out' type from the 'Direction' kind, allows users to write -- the @'BroadcastChan' v'Out' a@ type without enabling @DataKinds@.+--+-- @since 0.2.0 type Out = 'Out -- | The abstract type representing the read or write end of a 'BroadcastChan'.+--+-- @since 0.2.0 newtype BroadcastChan (dir :: Direction) a = BChan (MVar (Stream a)) deriving (Eq) @@ -31,6 +39,8 @@ data ChItem a = ChItem a {-# UNPACK #-} !(Stream a) | Closed -- | Creates a new 'BroadcastChan' write end.+--+-- @since 0.2.0 newBroadcastChan :: MonadIO m => m (BroadcastChan In a) newBroadcastChan = liftIO $ do hole <- newEmptyMVar@@ -40,6 +50,8 @@ -- | Close a 'BroadcastChan', disallowing further writes. Returns 'True' if the -- 'BroadcastChan' was closed. Returns 'False' if the 'BroadcastChan' was -- __already__ closed.+--+-- @since 0.2.0 closeBChan :: MonadIO m => BroadcastChan In a -> m Bool closeBChan (BChan writeVar) = liftIO . mask_ $ do old_hole <- takeMVar writeVar@@ -62,6 +74,8 @@ -- -- @True@ indicates the channel is both closed and empty, meaning reads -- will always fail.+--+-- @since 0.2.0 isClosedBChan :: MonadIO m => BroadcastChan dir a -> m Bool isClosedBChan (BChan mvar) = liftIO $ do old_hole <- readMVar mvar@@ -78,6 +92,8 @@ -- message was written, 'False' is the channel is closed. -- See @BroadcastChan.Throw.@'BroadcastChan.Throw.writeBChan' for an -- exception throwing variant.+--+-- @since 0.2.0 writeBChan :: MonadIO m => BroadcastChan In a -> a -> m Bool writeBChan (BChan writeVar) val = liftIO $ do new_hole <- newEmptyMVar@@ -95,6 +111,8 @@ -- 'Nothing' if the 'BroadcastChan' is closed and empty. -- See @BroadcastChan.Throw.@'BroadcastChan.Throw.readBChan' for an exception -- throwing variant.+--+-- @since 0.2.0 readBChan :: MonadIO m => BroadcastChan Out a -> m (Maybe a) readBChan (BChan readVar) = liftIO $ do modifyMVarMasked readVar $ \read_end -> do -- Note [modifyMVarMasked]@@ -106,6 +124,30 @@ Closed -> return (read_end, Nothing) {-# INLINE readBChan #-} +-- | A non-blocking version of 'readBChan'. Returns immediately with 'Nothing'+-- if the 'BroadcastChan' is empty.+--+-- The inner 'Maybe' value correspond to the return value of 'readBChan'. That+-- is, 'Nothing' if the 'BroadcastChan' is closed and empty and 'Just'+-- otherwise.+--+-- See @BroadcastChan.Throw.@'BroadcastChan.Throw.tryReadBChan' for an+-- exception throwing variant.+--+-- @since 0.3.0+tryReadBChan :: MonadIO m => BroadcastChan Out a -> m (Maybe (Maybe a))+tryReadBChan (BChan readVar) = liftIO $ do+ modifyMVarMasked readVar $ \read_end -> do -- Note [modifyMVarMasked]+ -- Use tryReadMVar here, not tryTakeMVar,+ -- else newBChanListener doesn't work+ result <- tryReadMVar read_end+ case result of+ Nothing -> return (read_end, Nothing)+ Just Closed -> return (read_end, Just Nothing)+ Just (ChItem val new_read_end) ->+ return (new_read_end, Just (Just val))+{-# INLINE tryReadBChan #-}+ -- Note [modifyMVarMasked] -- This prevents a theoretical deadlock if an asynchronous exception -- happens during the readMVar while the MVar is empty. In that case@@ -124,6 +166,8 @@ -- ['BroadcastChan' v'Out':]: -- -- Will receive all currently unread messages and all future messages.+--+-- @since 0.2.0 newBChanListener :: MonadIO m => BroadcastChan dir a -> m (BroadcastChan Out a) newBChanListener (BChan mvar) = liftIO $ do hole <- readMVar mvar@@ -148,6 +192,8 @@ -- the list resulting from this function is __not__ affected by reads on -- the input channel. Every message that is unread or written after the -- 'IO' action completes __will__ end up in the result list.+--+-- @since 0.2.0 getBChanContents :: BroadcastChan dir a -> IO [a] getBChanContents = newBChanListener >=> go where@@ -182,6 +228,8 @@ -- After the outer action completes the fold is unaffected by other -- (concurrent) reads performed on the original channel. So it's safe to -- reuse the channel.+--+-- @since 0.2.0 foldBChan :: (MonadIO m, MonadIO n) => (x -> a -> x)@@ -206,6 +254,8 @@ -- @"Control.Foldl".'Control.Foldl.impurely' 'foldBChanM' :: ('MonadIO' m, 'MonadIO' n) => 'Control.Foldl.FoldM' m a b -> 'BroadcastChan' d a -> n (m b)@ -- -- Has the same behaviour and guarantees as 'foldBChan'.+--+-- @since 0.2.0 foldBChanM :: (MonadIO m, MonadIO n) => (x -> a -> m x)
BroadcastChan/Prelude.hs view
@@ -2,7 +2,7 @@ ------------------------------------------------------------------------------- -- | -- Module : BroadcastChan.Prelude--- Copyright : (C) 2014-2021 Merijn Verstraaten+-- Copyright : (C) 2014-2022 Merijn Verstraaten -- License : BSD-style (see the file LICENSE) -- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl> -- Stability : experimental@@ -22,11 +22,15 @@ import BroadcastChan -- | 'mapM_' with it's arguments flipped.+--+-- @since 0.2.1 forM_ :: MonadIO m => BroadcastChan Out a -> (a -> m b) -> m () forM_ = flip mapM_ -- | Map a monadic function over the elements of a 'BroadcastChan', ignoring--- the results.+-- the results1+--+-- @since 0.2.1 mapM_ :: MonadIO m => (a -> m b) -> BroadcastChan Out a -> m () mapM_ f ch = do result <- liftIO $ readBChan ch
BroadcastChan/Throw.hs view
@@ -3,7 +3,7 @@ ------------------------------------------------------------------------------- -- | -- Module : BroadcastChan.Throw--- Copyright : (C) 2014-2021 Merijn Verstraaten+-- Copyright : (C) 2014-2022 Merijn Verstraaten -- License : BSD-style (see the file LICENSE) -- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl> -- Stability : experimental@@ -18,6 +18,7 @@ module BroadcastChan.Throw ( BChanError(..) , readBChan+ , tryReadBChan , writeBChan -- * Re-exports from "BroadcastChan" -- ** Datatypes@@ -45,13 +46,16 @@ ) where import Control.Monad (when)+import Control.Monad.IO.Unlift (MonadIO(..)) import Control.Exception (Exception, throwIO) import Data.Typeable (Typeable) -import BroadcastChan hiding (writeBChan, readBChan)+import BroadcastChan hiding (writeBChan, readBChan, tryReadBChan) import qualified BroadcastChan as Internal -- | Exception type for 'BroadcastChan' operations.+--+-- @since 0.2.0 data BChanError = WriteFailed -- ^ Attempted to write to closed 'BroadcastChan' | ReadFailed -- ^ Attempted to read from an empty closed 'BroadcastChan'@@ -61,18 +65,35 @@ -- | Like 'Internal.readBChan', but throws a 'ReadFailed' exception when -- reading from a closed and empty 'BroadcastChan'.-readBChan :: BroadcastChan Out a -> IO a+--+-- @since 0.2.0+readBChan :: MonadIO m => BroadcastChan Out a -> m a readBChan ch = do result <- Internal.readBChan ch case result of- Nothing -> throwIO ReadFailed+ Nothing -> liftIO $ throwIO ReadFailed Just x -> return x {-# INLINE readBChan #-}+--+-- | Like 'Internal.tryReadBChan', but throws a 'ReadFailed' exception when+-- reading from a closed and empty 'BroadcastChan'.+--+-- @since 0.3.0+tryReadBChan :: MonadIO m => BroadcastChan Out a -> m (Maybe a)+tryReadBChan ch = do+ result <- Internal.tryReadBChan ch+ case result of+ Nothing -> return Nothing+ Just Nothing -> liftIO $ throwIO ReadFailed+ Just v -> return v+{-# INLINE tryReadBChan #-} -- | Like 'Internal.writeBChan', but throws a 'WriteFailed' exception when -- writing to closed 'BroadcastChan'.-writeBChan :: BroadcastChan In a -> a -> IO ()+--+-- @since 0.2.0+writeBChan :: MonadIO m => BroadcastChan In a -> a -> m () writeBChan ch val = do success <- Internal.writeBChan ch val- when (not success) $ throwIO WriteFailed+ when (not success) . liftIO $ throwIO WriteFailed {-# INLINE writeBChan #-}
CHANGELOG.md view
@@ -1,3 +1,20 @@+0.3.0 [2025.03.16]+------------------+* Fixed Haddock links.+* Generalised `readBChan`/`writeBChan` in `BroadcastChan.Throw` to use+ `MonadIO`.+* Add non-blocking `tryReadBChan` in `BroadcastChan` and `BroadcastChan.Throw`,+ fixes [#12](https://github.com/merijn/broadcast-chan/issues/12).+* Absorb `broadcast-chan-tests`, `broadcast-chan-pipes`, and+ `broadcast-chan-conduit` into `broadcast-chan` as public sub-libraries.+* `ThreadBracket` in `BroadcastChan.Extra` has it's signature changed to+ support resourcet-1.3. Now takes `SomeException` as argument for+ `cleanupForkError`, existing code can simply ignore this argument.++0.2.1.2 [2022.08.24]+--------------------+* Revision updating bounds for GHC 9.4.+ 0.2.1.2 [2021.12.01] -------------------- * Update bounds for GHC 9.0 and 9.2.@@ -6,7 +23,7 @@ 0.2.1.1 [2020.03.05] ---------------------* Updated imports to support `unliftio-core` 0.2.x+* Updated imports to support `unliftio-core` 0.2.x. 0.2.1 [2019.11.17] ------------------
README.md view
@@ -2,9 +2,7 @@ ================================================================ [](https://en.wikipedia.org/wiki/BSD_License) [](https://hackage.haskell.org/package/broadcast-chan)-[](https://matrix.hackage.haskell.org/#/package/broadcast-chan)-[](https://www.stackage.org/package/broadcast-chan)-[](https://travis-ci.com/merijn/broadcast-chan)+[](https://github.com/merijn/broadcast-chan/actions/workflows/haskell-ci.yml/) A closable, fair, single-wakeup channel that avoids the 0 reader space leak that `Control.Concurrent.Chan` from base suffers from.
benchmarks/Utils.hs view
@@ -1,3 +1,4 @@+import Prelude hiding (Foldable(..)) import Control.Concurrent import Control.Monad (forM_) import Data.List (foldl')
+ broadcast-chan-test/BroadcastChan/Test.hs view
@@ -0,0 +1,392 @@+{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE Trustworthy #-}+{-# LANGUAGE ScopedTypeVariables #-}+-------------------------------------------------------------------------------+-- |+-- Module : BroadcastChan.Test+-- Copyright : (C) 2014-2022 Merijn Verstraaten+-- License : BSD-style (see the file LICENSE)+-- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl>+-- Stability : experimental+-- Portability : haha+--+-- Module containing testing helpers shared across all broadcast-chan packages.+-------------------------------------------------------------------------------+module BroadcastChan.Test+ ( (@?)+ , expect+ , genStreamTests+ , runTests+ , withLoggedOutput+ , MonadIO(..)+ , mapHandler+ -- * Re-exports of @tasty@ and @tasty-hunit@+ , module Test.Tasty+ , module Test.Tasty.HUnit+ ) where++import Prelude hiding (seq)+import Control.Concurrent (forkIO, setNumCapabilities, threadDelay)+import Control.Concurrent.Async (wait, withAsync)+import Control.Concurrent.MVar+import Control.Concurrent.QSemN+import Control.Concurrent.STM+import Control.Monad (void, when)+import Control.Monad.IO.Class (MonadIO(liftIO))+import Control.Exception (Exception, throwIO, try)+import Data.Bifunctor (second)+import Data.IntSet (IntSet)+import qualified Data.IntSet as IS+import Data.List (sort)+import Data.Map (Map)+import qualified Data.Map as M+import Data.Proxy (Proxy(Proxy))+import Data.Text (Text)+import qualified Data.Text as T+import qualified Data.Text.IO as T+import Data.Tagged (Tagged, untag)+import Data.Typeable (Typeable)+import Options.Applicative (flag', long, help)+import System.Clock+ (Clock(Monotonic), TimeSpec, diffTimeSpec, getTime, toNanoSecs)+import System.Environment (setEnv)+import System.IO (Handle, SeekMode(AbsoluteSeek), hPrint, hSeek)+import System.IO.Temp (withSystemTempFile)+import Test.Tasty+import Test.Tasty.Ingredients.Basic (consoleTestReporter, listingTests)+import Test.Tasty.Golden.Advanced (goldenTest)+import Test.Tasty.HUnit hiding ((@?))+import qualified Test.Tasty.HUnit as HUnit+import Test.Tasty.Options++import BroadcastChan.Extra (Action(..), Handler(..), mapHandler)+import ParamTree++data TestException = TestException deriving (Eq, Show, Typeable)+instance Exception TestException++infix 0 @?+-- | Monomorphised version of 'Test.Tasty.HUnit.@?' to avoid ambiguous type+-- errors when combined with predicates that are @MonadIO m => m Bool@.+--+-- @since 0.2.0+(@?) :: IO Bool -> String -> Assertion+(@?) = (HUnit.@?)++-- | Test which fails if the expected exception is not thrown by the 'IO'+-- action.+--+-- @since 0.2.0+expect :: (Eq e, Exception e) => e -> IO a -> Assertion+expect err act = do+ result <- try act+ case result of+ Left e | e == err -> return ()+ | otherwise -> assertFailure $+ "Expected: " ++ show err ++ "\nGot: " ++ show e+ Right _ -> assertFailure $ "Expected exception, got success."++-- | Pauses a number of microseconds before returning its input.+doNothing :: Int -> a -> IO a+doNothing threadPause x = do+ threadDelay threadPause+ return x++-- | Print a value, then return it.+doPrint :: Show a => Handle -> a -> IO a+doPrint hnd x = do+ hPrint hnd x+ return x++doDrop :: Show a => (a -> Bool) -> Handle -> a -> IO a+doDrop predicate hnd val+ | predicate val = throwIO TestException+ | otherwise = doPrint hnd val++doRace :: MVar () -> QSemN -> a -> IO a+doRace mvar sem _ = do+ result <- tryReadMVar mvar+ case result of+ Nothing -> signalQSemN sem 1 >> readMVar mvar+ Just () -> return ()+ throwIO TestException++fromTimeSpec :: Fractional n => TimeSpec -> n+fromTimeSpec = fromIntegral . toNanoSecs++speedupTest+ :: forall r . (Eq r, Show r)+ => IO (TVar (Map ([Int], Int) (MVar (r, Double))))+ -> ([Int] -> (Int -> IO Int) -> IO r)+ -> ([Int] -> (Int -> IO Int) -> Int -> IO r)+ -> Int+ -> [Int]+ -> Int+ -> String+ -> TestTree+speedupTest getCache seqSink parSink n inputs pause name = testCase name $ do+ withAsync cachedSequential $ \seqAsync ->+ withAsync (timed $ parSink inputs testFun n) $ \parAsync -> do+ (seqResult, seqTime) <- wait seqAsync+ (parResult, parTime) <- wait parAsync+ seqResult @=? parResult+ let lowerBound :: Double+ lowerBound = seqTime / (fromIntegral n + 1)++ upperBound :: Double+ upperBound = seqTime / (fromIntegral n - 1)++ errorMsg :: String+ errorMsg = mconcat+ [ "Parallel time should be 1/"+ , show n+ , "th of sequential time!\n"+ , "Actual time was 1/"+ , show (round $ seqTime / parTime :: Int)+ , "th (", show (seqTime/parTime), ")"+ ]+ assertBool errorMsg $ lowerBound < parTime && parTime < upperBound+ where+ testFun :: Int -> IO Int+ testFun = doNothing pause++ timed :: IO a -> IO (a, Double)+ timed = fmap (\(x, t) -> (x, fromTimeSpec t)) . withTime++ cachedSequential :: IO (r, Double)+ cachedSequential = do+ mvar <- newEmptyMVar+ cacheTVar <- getCache+ result <- atomically $ do+ cacheMap <- readTVar cacheTVar+ let (oldVal, newMap) = M.insertLookupWithKey+ (\_ _ v -> v)+ (inputs, pause)+ mvar+ cacheMap+ writeTVar cacheTVar newMap+ return oldVal++ case result of+ Just var -> readMVar var+ Nothing -> do+ timed (seqSink inputs testFun) >>= putMVar mvar+ readMVar mvar++-- | Run an IO action while logging the output to a @Handle@. Returns the+-- result and the logged output.+--+-- @since 0.2.0+withLoggedOutput :: FilePath -> (Handle -> IO r) -> IO (r, Text)+withLoggedOutput filename act = withSystemTempFile filename $ \_ hnd ->+ (,) <$> act hnd <*> rewindAndRead hnd+ where+ rewindAndRead :: Handle -> IO Text+ rewindAndRead hnd = do+ hSeek hnd AbsoluteSeek 0+ T.hGetContents hnd++nonDeterministicGolden+ :: forall r+ . (Eq r, Show r)+ => String+ -> (Handle -> IO r)+ -> (Handle -> IO r)+ -> TestTree+nonDeterministicGolden label controlAction testAction =+ goldenTest label (normalise control) (normalise test) diff update+ where+ normalise :: MonadIO m => IO (a, Text) -> m (a, Text)+ normalise = liftIO . fmap (second (T.strip . T.unlines . sort . T.lines))++ control :: IO (r, Text)+ control = withLoggedOutput "control.out" controlAction++ test :: IO (r, Text)+ test = withLoggedOutput "test.out" testAction++ diff :: (r, Text) -> (r, Text) -> IO (Maybe String)+ diff (controlResult, controlOutput) (testResult, testOutput) =+ return $ resultDiff `mappend` outputDiff+ where+ resultDiff :: Maybe String+ resultDiff+ | controlResult == testResult = Nothing+ | otherwise = Just "Results differ!\n"++ outputDiff :: Maybe String+ outputDiff+ | controlOutput == testOutput = Nothing+ | otherwise = Just . mconcat $+ [ "Outputs differ!\n"+ , "Expected:\n\"", T.unpack controlOutput, "\"\n\n"+ , "Got:\n\"", T.unpack testOutput, "\"\n"+ ]++ update :: (r, Text) -> IO ()+ update _ = return ()++outputTest+ :: forall r . (Eq r, Show r)+ => ([Int] -> (Int -> IO Int) -> IO r)+ -> ([Int] -> (Int -> IO Int) -> Int -> IO r)+ -> Int+ -> [Int]+ -> String+ -> TestTree+outputTest seqSink parSink threads inputs label =+ nonDeterministicGolden label seqTest parTest+ where+ seqTest :: Handle -> IO r+ seqTest = seqSink inputs . doPrint++ parTest :: Handle -> IO r+ parTest hndl = parSink inputs (doPrint hndl) threads++dropTest+ :: (Eq r, Show r)+ => ([Int] -> (Int -> IO Int) -> IO r)+ -> (Handler IO Int -> [Int] -> (Int -> IO Int) -> Int -> IO r)+ -> TestTree+dropTest seqImpl parImpl = nonDeterministicGolden "drop"+ (seqImpl filteredInputs . doPrint)+ (\hnd -> parImpl (Simple Drop) inputs (doDrop even hnd) 2)+ where+ inputs = [1..100]+ filteredInputs = filter (not . even) inputs++terminationTest+ :: (Handler IO Int -> [Int] -> (Int -> IO Int) -> Int -> IO r) -> TestTree+terminationTest parImpl = testCase "termination" $+ expect TestException . void $+ withSystemTempFile "terminate.out" $ \_ hndl ->+ parImpl (Simple Terminate) [1..100] (doDrop even hndl) 4++raceTest+ :: (Handler IO Int -> [Int] -> (Int -> IO Int) -> Int -> IO r) -> TestTree+raceTest parImpl = testCase "race" $+ expect TestException . void $ do+ sem <- newQSemN 0+ mvar <- newEmptyMVar+ forkIO $ do+ waitQSemN sem parCount+ putMVar mvar ()+ parImpl (Simple Terminate) [1..100] (doRace mvar sem) parCount+ where+ parCount :: Int+ parCount = 4++retryTest+ :: (Eq r, Show r)+ => ([Int] -> (Int -> IO Int) -> IO r)+ -> (Handler IO Int -> [Int] -> (Int -> IO Int) -> Int -> IO r)+ -> TestTree+retryTest seqImpl parImpl = withRetryCheck $ \getRetryCheck ->+ nonDeterministicGolden+ "retry"+ (seqImpl seqInputs . doPrint)+ (\h -> parImpl (Simple Retry) parInputs (dropAfterPrint getRetryCheck h) 4)+ where+ withRetryCheck = withResource alloc clean+ where+ alloc = updateRetry <$> newMVar IS.empty+ clean _ = return ()++ parInputs = [1..100]+ seqInputs = parInputs ++ filter even parInputs++ updateRetry :: MVar IntSet -> Int -> IO Bool+ updateRetry mvar val = modifyMVar mvar updateSet+ where+ updateSet :: IntSet -> IO (IntSet, Bool)+ updateSet set+ | IS.member val set = return (set, False)+ | otherwise = return (IS.insert val set, True)++ dropAfterPrint :: IO (Int -> IO Bool) -> Handle -> Int -> IO Int+ dropAfterPrint checkPresence hnd val = do+ hPrint hnd val+ when (even val) $ do+ isNotPresent <- checkPresence >>= ($ val)+ when isNotPresent $ throwIO TestException+ return val++newtype SlowTests = SlowTests Bool+ deriving (Eq, Ord, Typeable)++instance IsOption SlowTests where+ defaultValue = SlowTests False+ parseValue = fmap SlowTests . safeRead+ optionName = return "slow-tests"+ optionHelp = return "Run slow tests."+ optionCLParser = flag' (SlowTests True) $ mconcat+ [ long (untag (optionName :: Tagged SlowTests String))+ , help (untag (optionHelp :: Tagged SlowTests String))+ ]++-- | Takes a name, a sequential sink, and a parallel sink and generates tasty+-- tests from these.+--+-- The parallel and sequential sink should perform the same tasks so their+-- results can be compared to check correctness.+--+-- The sinks should take a list of input data, a function processing the data,+-- and return a result that can be compared for equality.+--+-- Furthermore the parallel sink should take a number indicating how many+-- concurrent consumers should be used.+--+-- @since 0.2.0+genStreamTests+ :: (Eq r, Show r)+ => String -- ^ Name to group tests under+ -> ([Int] -> (Int -> IO Int) -> IO r) -- ^ Sequential sink+ -> (Handler IO Int -> [Int] -> (Int -> IO Int) -> Int -> IO r)+ -- ^ Parallel sink+ -> TestTree+genStreamTests name seq par = askOption $ \(SlowTests slow) ->+ withResource (newTVarIO M.empty) (const $ return ()) $ \getCache ->+ let+ testTree = growTree (Just ".") testGroup+ threads = simpleParam "threads" [1,2,5]+ bigInputs | slow = derivedParam (enumFromTo 0) "inputs" [600]+ | otherwise = derivedParam (enumFromTo 0) "inputs" [300]+ smallInputs = derivedParam (enumFromTo 0) "inputs" [0,1,2]+ pause = simpleParam "pause" [10^(4 :: Int)]++ in testGroup name+ [ testTree "output" (outputTest seq (par term)) $+ threads . paramSets [ smallInputs, bigInputs ]+ , testTree "speedup" (speedupTest getCache seq (par term)) $+ threads . bigInputs . pause+ , testGroup "exceptions"+ [ dropTest seq par+ , terminationTest par+ , raceTest par+ , retryTest seq par+ ]+ ]+ where+ term = Simple Terminate++-- | Run a list of 'TestTree''s and group them under the specified name.+--+-- @since 0.2.0+runTests :: String -> [TestTree] -> IO ()+runTests name tests = do+ setNumCapabilities 5+ setEnv "TASTY_NUM_THREADS" "100"+ defaultMainWithIngredients ingredients $ testGroup name tests+ where+ ingredients =+ [ includingOptions [Option (Proxy :: Proxy SlowTests)]+ , listingTests+ , consoleTestReporter+ ]++withTime :: IO a -> IO (a, TimeSpec)+withTime act = do+ start <- getTime Monotonic+ r <- act+ end <- getTime Monotonic+ return $ (r, diffTimeSpec start end)
broadcast-chan.cabal view
@@ -1,13 +1,13 @@-Cabal-Version: 2.2+Cabal-Version: 3.4 Name: broadcast-chan-Version: 0.2.1.2+Version: 0.3.0 Homepage: https://github.com/merijn/broadcast-chan Bug-Reports: https://github.com/merijn/broadcast-chan/issues Author: Merijn Verstraaten Maintainer: Merijn Verstraaten <merijn@inconsistent.nl>-Copyright: Copyright © 2014-2021 Merijn Verstraaten+Copyright: Copyright © 2014-2025 Merijn Verstraaten License: BSD-3-Clause License-File: LICENSE@@ -15,10 +15,11 @@ Category: System Build-Type: Simple Tested-With: GHC == 8.0.2, GHC == 8.2.2, GHC == 8.4.4, GHC == 8.6.5,- GHC == 8.8.4, GHC == 8.10.7, GHC == 9.0.1, GHC == 9.2.1+ GHC == 8.8.4, GHC == 8.10.7, GHC == 9.0.2, GHC == 9.2.7,+ GHC == 9.4.8, GHC == 9.6.6, GHC == 9.8.4, GHC == 9.10.1,+ GHC == 9.12.1 -Extra-Source-Files: README.md- , CHANGELOG.md+Extra-Doc-Files: README.md, CHANGELOG.md Synopsis: Closable, fair, single-wakeup channel type that avoids 0 reader space leaks.@@ -27,16 +28,16 @@ __WARNING:__ While the code in this library should be fairly stable and production, the API is something I'm still working on. API changes will follow the PVP, but __expect__ breaking API changes in future versions!- .+ A closable, fair, single-wakeup channel that avoids the 0 reader space leak that @"Control.Concurrent.Chan"@ from base suffers from.- .+ The @Chan@ type from @"Control.Concurrent.Chan"@ consists of both a read and write end combined into a single value. This means there is always at least 1 read end for a @Chan@, which keeps any values written to it alive. This is a problem for applications/libraries that want to have a channel that can have zero listeners.- .+ Suppose we have an library that produces events and we want to let users register to receive events. If we use a channel and write all events to it, we would like to drop and garbage collect any events that take place when@@ -44,30 +45,30 @@ makes this impossible. We end up with a @Chan@ that forever accumulates more and more events that will never get removed, resulting in a memory leak.- .+ @"BroadcastChan"@ splits channels into separate read and write ends. Any message written to a a channel with no existing read end is immediately dropped so it can be garbage collected. Once a read end is created, all messages written to the channel will be accessible to that read end.- .+ Once all read ends for a channel have disappeared and been garbage collected, the channel will return to dropping messages as soon as they are written.- .+ __Why should I use "BroadcastChan" over "Control.Concurrent.Chan"?__- .+ * @"BroadcastChan"@ is closable,- .+ * @"BroadcastChan"@ has no 0 reader space leak,- .+ * @"BroadcastChan"@ has comparable or better performance.- .+ __Why should I use "BroadcastChan" over various (closable) STM channels?__- .+ * @"BroadcastChan"@ is single-wakeup,- .+ * @"BroadcastChan"@ is fair,- .+ * @"BroadcastChan"@ performs better under contention. Flag sync@@ -100,10 +101,149 @@ Trustworthy TupleSections - Build-Depends: base >= 4.8 && < 4.17+ Build-Depends: base >= 4.8 && < 4.22 , transformers >= 0.2 && < 0.7 , unliftio-core >= 0.1.1 && < 0.3 +Library conduit+ Visibility: public+ Default-Language: Haskell2010+ GHC-Options: -Wall -O2 -Wno-unused-do-bind+ Exposed-Modules: BroadcastChan.Conduit+ BroadcastChan.Conduit.Throw+ Other-Modules: BroadcastChan.Conduit.Internal++ Hs-Source-Dirs: conduit++ Other-Extensions: CPP+ NamedFieldPuns+ Safe+ ScopedTypeVariables+ Trustworthy++ Build-Depends: base >= 4.8 && < 4.22+ , broadcast-chan+ , conduit >= 1.2 && < 1.4+ , resourcet >= 1.1 && < 1.4+ , transformers >= 0.2 && < 0.7+ , unliftio-core >= 0.1 && < 0.3++Library pipes+ Visibility: public+ Default-Language: Haskell2010+ GHC-Options: -Wall -O2 -Wno-unused-do-bind+ Exposed-Modules: BroadcastChan.Pipes+ BroadcastChan.Pipes.Throw+ Other-Modules: BroadcastChan.Pipes.Internal++ Hs-Source-Dirs: pipes++ Other-Extensions: NamedFieldPuns+ Safe+ ScopedTypeVariables++ Build-Depends: base >= 4.8 && < 4.22+ , broadcast-chan+ , pipes >= 4.1.6 && < 4.4+ , pipes-safe >= 2.3.1 && < 2.4++Library test+ Default-Language: Haskell2010+ GHC-Options: -Wall -O2 -Wno-unused-do-bind+ Exposed-Modules: BroadcastChan.Test+ Hs-Source-Dirs: broadcast-chan-test++ Other-Extensions: DeriveDataTypeable+ ScopedTypeVariables+ Trustworthy++ Build-Depends: base >= 4.8 && < 4.22+ , async >= 2.1 && < 2.3+ , broadcast-chan+ , clock >= 0.7 && < 0.9+ , containers >= 0.4 && < 0.9+ , optparse-applicative >= 0.12 && < 0.19+ , paramtree >= 0.1.1 && < 0.2+ , stm >= 2.4 && < 2.6+ , tagged == 0.8.*+ , tasty >= 0.11 && < 1.6+ , tasty-golden >= 2.0 && < 2.4+ , tasty-hunit >= 0.9 && < 0.11+ , temporary >= 1.2 && < 1.4+ , text >= 1.0 && < 1.3 || ^>= 2.0 || ^>= 2.1++ if impl(ghc <= 7.10)+ Build-Depends: transformers >= 0.2 && < 0.7++Common basic+ Default-Language: Haskell2010+ GHC-Options: -Wall -Wno-unused-do-bind+ Hs-Source-Dirs: tests+ Build-Depends: base >= 4.8 && < 4.22+ , broadcast-chan+ , broadcast-chan:test+ , foldl >= 1.4 && < 1.5+ , monad-loops >= 0.4.3 && < 0.5+ , random >= 1.1 && < 1.4++Test-Suite basic+ Import: basic+ Type: exitcode-stdio-1.0+ Main-Is: Basic.hs+ GHC-Options: -threaded -with-rtsopts=-qg++Test-Suite basic-unthreaded+ Import: basic+ Type: exitcode-stdio-1.0+ Main-Is: Basic.hs++Common parallel-io+ Default-Language: Haskell2010+ GHC-Options: -Wall -O2 -Wno-unused-do-bind+ Hs-Source-Dirs: tests+ Build-Depends: base >= 4.8 && < 4.22+ , broadcast-chan+ , broadcast-chan:test+ , containers >= 0.4 && < 0.9++Test-Suite parallel-io+ Import: parallel-io+ Type: exitcode-stdio-1.0+ Main-Is: IOTest.hs+ GHC-Options: -threaded -with-rtsopts=-qg++Test-Suite parallel-io-unthreaded+ Import: parallel-io+ Type: exitcode-stdio-1.0+ Main-Is: IOTest.hs++Test-Suite conduit-test+ Default-Language: Haskell2010+ Type: exitcode-stdio-1.0+ Main-Is: ConduitTest.hs+ GHC-Options: -Wall -Wno-unused-do-bind -threaded -with-rtsopts=-qg+ Hs-Source-Dirs: tests++ Build-Depends: base >= 4.8 && < 4.22+ , broadcast-chan:conduit+ , broadcast-chan:test+ , containers >= 0.4 && < 0.9+ , conduit >= 1.2 && < 1.4++Test-Suite pipes-test+ Default-Language: Haskell2010+ Type: exitcode-stdio-1.0+ Main-Is: PipeTest.hs+ GHC-Options: -Wall -Wno-unused-do-bind -threaded -with-rtsopts=-qg+ Hs-Source-Dirs: tests+ Build-Depends: base >= 4.8 && < 4.22+ , broadcast-chan:pipes+ , broadcast-chan:test+ , containers >= 0.4 && < 0.9+ , foldl >= 1.0.4 && < 1.5+ , pipes >= 4.1.6 && < 4.4+ , pipes-safe >= 2.3.1 && < 2.4+ Common concurrent-benchmarks Default-Language: Haskell2010 @@ -115,10 +255,10 @@ Other-Extensions: BangPatterns - Build-Depends: base+ Build-Depends: base >= 4.8 && < 4.22 , async >= 2.0 && < 2.3- , criterion >= 1.2 && < 1.6- , deepseq >= 1.1 && < 1.5+ , criterion >= 1.2 && < 1.7+ , deepseq >= 1.1 && < 1.6 , stm >= 2.4 && < 2.6 Benchmark sync@@ -153,7 +293,7 @@ GHC-Options: -threaded -with-rtsopts=-qg Hs-Source-Dirs: benchmarks - Build-Depends: base+ Build-Depends: base >= 4.8 && < 4.22 , broadcast-chan Source-Repository head
+ conduit/BroadcastChan/Conduit.hs view
@@ -0,0 +1,42 @@+{-# LANGUAGE Safe #-}+-------------------------------------------------------------------------------+-- |+-- Module : BroadcastChan.Conduit+-- Copyright : (C) 2014-2022 Merijn Verstraaten+-- License : BSD-style (see the file LICENSE)+-- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl>+-- Stability : experimental+-- Portability : haha+--+-- This module is identical to "BroadcastChan", but replaces the parallel+-- processing operations with functions for creating conduits and sinks that+-- process in parallel.+-------------------------------------------------------------------------------+module BroadcastChan.Conduit+ ( Action(..)+ , Handler(..)+ , parMapM+ , parMapM_+ -- * Re-exports from "BroadcastChan"+ -- ** Datatypes+ , BroadcastChan+ , Direction(..)+ , In+ , Out+ -- ** Construction+ , newBroadcastChan+ , newBChanListener+ -- ** Basic Operations+ , readBChan+ , writeBChan+ , closeBChan+ , isClosedBChan+ , getBChanContents+ -- ** Foldl combinators+ -- | Combinators for use with Tekmo's @foldl@ package.+ , foldBChan+ , foldBChanM+ ) where++import BroadcastChan hiding (parMapM_)+import BroadcastChan.Conduit.Internal
+ conduit/BroadcastChan/Conduit/Internal.hs view
@@ -0,0 +1,117 @@+{-# LANGUAGE CPP #-}+{-# LANGUAGE NamedFieldPuns #-}+{-# LANGUAGE Trustworthy #-}+{-# LANGUAGE ScopedTypeVariables #-}+module BroadcastChan.Conduit.Internal (parMapM, parMapM_) where++import Control.Exception (SomeException)+import Control.Monad ((>=>))+import Control.Monad.Trans.Resource (MonadResource)+import qualified Control.Monad.Trans.Resource as Resource+import qualified Control.Monad.Trans.Resource.Internal as ResourceI+import Control.Monad.Trans.Class (lift)+import Control.Monad.IO.Unlift (MonadUnliftIO, UnliftIO(..), askUnliftIO)+import Data.Acquire (ReleaseType(..), allocateAcquire, mkAcquireType)+import Data.Conduit (ConduitM, (.|), awaitForever, yield)+import qualified Data.Conduit.List as C+import Data.Foldable (traverse_)+import Data.Void (Void)++import BroadcastChan.Extra (BracketOnError(..), Handler, ThreadBracket(..))+import qualified BroadcastChan.Extra as Extra++releaseWithException :: SomeException -> ReleaseType+#if !MIN_VERSION_resourcet(1,3,0)+releaseWithException _ = ReleaseException+#else+releaseWithException exc = ReleaseExceptionWith exc+#endif++bracketOnError :: MonadResource m => IO a -> (a -> IO ()) -> m r -> m r+bracketOnError alloc clean work =+ allocateAcquire (mkAcquireType alloc cleanup) >>= const work+ where+#if !MIN_VERSION_resourcet(1,3,0)+ cleanup x ReleaseException = clean x+#else+ cleanup x (ReleaseExceptionWith _) = clean x+#endif+ cleanup _ _ = return ()++-- | Create a conduit that processes inputs in parallel.+--+-- This function does __NOT__ guarantee that input elements are processed or+-- output in a deterministic order!+--+-- @since 0.2.0+parMapM+ :: (MonadResource m, MonadUnliftIO m)+ => Handler m a+ -- ^ Exception handler+ -> Int+ -- ^ Number of parallel threads to use+ -> (a -> m b)+ -- ^ Function to run in parallel+ -> ConduitM a b m ()+parMapM hnd threads workFun = do+ UnliftIO runInIO <- lift askUnliftIO++ resourceState <- Resource.liftResourceT Resource.getInternalState++ let threadBracket = ThreadBracket+ { setupFork = ResourceI.stateAlloc resourceState+ , cleanupFork = ResourceI.stateCleanup ReleaseNormal resourceState+ , cleanupForkError = \exc ->+ ResourceI.stateCleanup (releaseWithException exc) resourceState+ }++ Bracket{allocate,cleanup,action} <- Extra.runParallelWith+ threadBracket+ (Left yield)+ (Extra.mapHandler runInIO hnd)+ threads+ (runInIO . workFun)+ body++ bracketOnError allocate cleanup action+ where+ body :: Monad m => (a -> m ()) -> (a -> m (Maybe b)) -> ConduitM a b m ()+ body buffer process = do+ C.isolate threads .| C.mapM_ buffer+ awaitForever $ lift . process >=> traverse_ yield++-- | Create a conduit sink that consumes inputs in parallel.+--+-- This function does __NOT__ guarantee that input elements are processed or+-- output in a deterministic order!+--+-- @since 0.2.0+parMapM_+ :: (MonadResource m, MonadUnliftIO m)+ => Handler m a+ -- ^ Exception handler+ -> Int+ -- ^ Number of parallel threads to use+ -> (a -> m ())+ -- ^ Function to run in parallel+ -> ConduitM a Void m ()+parMapM_ hnd threads workFun = do+ UnliftIO runInIO <- lift askUnliftIO++ resourceState <- Resource.liftResourceT Resource.getInternalState++ let threadBracket = ThreadBracket+ { setupFork = ResourceI.stateAlloc resourceState+ , cleanupFork = ResourceI.stateCleanup ReleaseNormal resourceState+ , cleanupForkError = \exc ->+ ResourceI.stateCleanup (releaseWithException exc) resourceState+ }++ Bracket{allocate,cleanup,action} <- Extra.runParallelWith_+ threadBracket+ (Extra.mapHandler runInIO hnd)+ threads+ (runInIO . workFun)+ C.mapM_++ bracketOnError allocate cleanup action
+ conduit/BroadcastChan/Conduit/Throw.hs view
@@ -0,0 +1,43 @@+{-# LANGUAGE Safe #-}+-------------------------------------------------------------------------------+-- |+-- Module : BroadcastChan.Conduit.Throw+-- Copyright : (C) 2014-2022 Merijn Verstraaten+-- License : BSD-style (see the file LICENSE)+-- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl>+-- Stability : experimental+-- Portability : haha+--+-- This module is identical to "BroadcastChan.Throw", but replaces the parallel+-- processing operations with functions for creating conduits and sinks that+-- process in parallel.+-------------------------------------------------------------------------------+module BroadcastChan.Conduit.Throw+ ( Action(..)+ , Handler(..)+ , parMapM+ , parMapM_+ -- * Re-exports from "BroadcastChan.Throw"+ -- ** Datatypes+ , BroadcastChan+ , Direction(..)+ , In+ , Out+ , BChanError(..)+ -- ** Construction+ , newBroadcastChan+ , newBChanListener+ -- ** Basic Operations+ , readBChan+ , writeBChan+ , closeBChan+ , isClosedBChan+ , getBChanContents+ -- ** Foldl combinators+ -- | Combinators for use with Tekmo's @foldl@ package.+ , foldBChan+ , foldBChanM+ ) where++import BroadcastChan.Throw hiding (parMapM_)+import BroadcastChan.Conduit.Internal
+ pipes/BroadcastChan/Pipes.hs view
@@ -0,0 +1,43 @@+{-# LANGUAGE Safe #-}+-------------------------------------------------------------------------------+-- |+-- Module : BroadcastChan.Pipes+-- Copyright : (C) 2014-2022 Merijn Verstraaten+-- License : BSD-style (see the file LICENSE)+-- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl>+-- Stability : experimental+-- Portability : haha+--+-- This module is identical to "BroadcastChan", but replaces the parallel+-- processing operations with functions for creating producers and effects that+-- process in parallel.+-------------------------------------------------------------------------------+module BroadcastChan.Pipes+ ( Action(..)+ , Handler(..)+ , parMapM+ , parMapM_+ -- * Re-exports from "BroadcastChan"+ -- ** Datatypes+ , BroadcastChan+ , Direction(..)+ , In+ , Out+ -- ** Construction+ , newBroadcastChan+ , newBChanListener+ -- ** Basic Operations+ , readBChan+ , tryReadBChan+ , writeBChan+ , closeBChan+ , isClosedBChan+ , getBChanContents+ -- ** Foldl combinators+ -- | Combinators for use with Tekmo's @foldl@ package.+ , foldBChan+ , foldBChanM+ ) where++import BroadcastChan hiding (parMapM_)+import BroadcastChan.Pipes.Internal
+ pipes/BroadcastChan/Pipes/Internal.hs view
@@ -0,0 +1,71 @@+{-# LANGUAGE NamedFieldPuns #-}+{-# LANGUAGE Safe #-}+{-# LANGUAGE ScopedTypeVariables #-}+module BroadcastChan.Pipes.Internal (parMapM, parMapM_) where++import Control.Monad ((>=>), replicateM)+import Data.Foldable (traverse_)+import Pipes+import qualified Pipes.Prelude as P+import Pipes.Safe (MonadSafe)+import qualified Pipes.Safe as Safe++import BroadcastChan.Extra+ (BracketOnError(..), Handler, runParallel, runParallel_)++bracketOnError :: MonadSafe m => IO a -> (a -> IO b) -> m c -> m c+bracketOnError alloc clean =+ Safe.bracketOnError (liftIO alloc) (liftIO . clean) . const++-- | Create a producer that processes its inputs in parallel.+--+-- This function does __NOT__ guarantee that input elements are processed or+-- output in a deterministic order!+--+-- @since 0.2.0+parMapM+ :: forall a b m+ . MonadSafe m+ => Handler IO a+ -- ^ Exception handler+ -> Int+ -- ^ Number of parallel threads to use+ -> (a -> IO b)+ -- ^ Function to run in parallel+ -> Producer a m ()+ -- ^ Input producer+ -> Producer b m ()+parMapM hndl i f prod = do+ Bracket{allocate,cleanup,action} <- runParallel (Left yield) hndl i f body+ bracketOnError allocate cleanup action+ where+ body :: (a -> m ()) -> (a -> m (Maybe b)) -> Producer b m ()+ body buffer process = prod >-> work+ where+ work :: Pipe a b m ()+ work = do+ replicateM i (await >>= lift . buffer)+ for cat $ lift . process >=> traverse_ yield++-- | Create an Effect that processes its inputs in parallel.+--+-- This function does __NOT__ guarantee that input elements are processed or+-- output in a deterministic order!+--+-- @since 0.2.0+parMapM_+ :: MonadSafe m+ => Handler IO a+ -- ^ Exception handler+ -> Int+ -- ^ Number of parallel threads to use+ -> (a -> IO ())+ -- ^ Function to run in parallel+ -> Producer a m r+ -- ^ Input producer+ -> Effect m r+parMapM_ hndl i f prod = do+ Bracket{allocate,cleanup,action} <- runParallel_ hndl i f workProd+ bracketOnError allocate cleanup action+ where+ workProd buffer = prod >-> P.mapM_ buffer
+ pipes/BroadcastChan/Pipes/Throw.hs view
@@ -0,0 +1,44 @@+{-# LANGUAGE Safe #-}+-------------------------------------------------------------------------------+-- |+-- Module : BroadcastChan.Pipes.Throw+-- Copyright : (C) 2014-2022 Merijn Verstraaten+-- License : BSD-style (see the file LICENSE)+-- Maintainer : Merijn Verstraaten <merijn@inconsistent.nl>+-- Stability : experimental+-- Portability : haha+--+-- This module is identical to "BroadcastChan.Throw", but replaces the parallel+-- processing operations with functions for creating producers and effects that+-- process in parallel.+-------------------------------------------------------------------------------+module BroadcastChan.Pipes.Throw+ ( Action(..)+ , Handler(..)+ , parMapM+ , parMapM_+ -- * Re-exports from "BroadcastChan.Throw"+ -- ** Datatypes+ , BroadcastChan+ , Direction(..)+ , In+ , Out+ , BChanError(..)+ -- ** Construction+ , newBroadcastChan+ , newBChanListener+ -- ** Basic Operations+ , readBChan+ , tryReadBChan+ , writeBChan+ , closeBChan+ , isClosedBChan+ , getBChanContents+ -- ** Foldl combinators+ -- | Combinators for use with Tekmo's @foldl@ package.+ , foldBChan+ , foldBChanM+ ) where++import BroadcastChan.Throw hiding (parMapM_)+import BroadcastChan.Pipes.Internal
+ tests/Basic.hs view
@@ -0,0 +1,366 @@+import Control.Foldl (Fold, FoldM)+import qualified Control.Foldl as Foldl+import Control.Monad (forM_)+import Control.Monad.Loops (unfoldM)+import Data.Maybe (isNothing)+import System.IO (Handle, hPrint)+import System.Random (getStdGen, randomIO, randoms)+import System.Timeout (timeout)++import BroadcastChan+import BroadcastChan.Test+import BroadcastChan.Throw (BChanError(..))+import qualified BroadcastChan.Throw as Throw++shouldn'tBlock :: IO a -> IO a+shouldn'tBlock act = do+ result <- timeout 2000000 act+ case result of+ Nothing -> assertFailure "Shouldn't block!"+ Just a -> return a++shouldBlock :: IO a -> IO Bool+shouldBlock act = do+ result <- timeout 2000000 act+ case result of+ Nothing -> return True+ Just _ -> return False++checkedWrite :: BroadcastChan In a -> a -> IO ()+checkedWrite chan val = writeBChan chan val @? "Write shouldn't fail"++randomList :: Int -> IO [Int]+randomList n = take n . randoms <$> getStdGen++data ChanContent b = ChanEmpty (IO b -> Assertion) | ChanNonEmpty (Int -> b)+data ChanState = ChanOpen | ChanClosed++readBoilerPlate+ :: (Eq b, Show b)+ => String+ -> ChanState+ -> ChanContent b+ -> (BroadcastChan Out Int -> IO b)+ -> TestTree+readBoilerPlate name state content getResult = testCase name $ do+ inChan <- newBroadcastChan+ outChan <- newBChanListener inChan+ val <- randomIO :: IO Int++ let handleRead = case content of+ ChanNonEmpty conv ->+ (>>= assertEqual "Read should match write" (conv val))+ ChanEmpty conv -> conv++ case content of+ ChanEmpty{} -> return ()+ ChanNonEmpty{} -> Throw.writeBChan inChan val++ case state of+ ChanOpen -> return ()+ ChanClosed -> () <$ closeBChan inChan++ handleRead $ getResult outChan++readTests :: TestTree+readTests = testGroup "read tests"+ [ readBoilerPlate "read non-empty" ChanOpen nonEmpty $+ shouldn'tBlock . readBChan+ , readBoilerPlate "read non-empty (throw)" ChanOpen nonEmptyThrow $+ shouldn'tBlock . Throw.readBChan+ , readBoilerPlate "read non-empty closed" ChanClosed nonEmpty $+ shouldn'tBlock . readBChan+ , readBoilerPlate "read non-empty closed (throw)" ChanClosed nonEmptyThrow $+ shouldn'tBlock . Throw.readBChan+ , readBoilerPlate "read empty" ChanOpen emptyBlock $+ shouldBlock . readBChan+ , readBoilerPlate "read empty (throw)" ChanOpen emptyBlock $+ shouldBlock . Throw.readBChan+ , readBoilerPlate "read empty closed" ChanClosed emptyNonBlock $+ fmap isNothing . shouldn'tBlock . readBChan+ , readBoilerPlate "read empty closed (throw)" ChanClosed emptyThrow $+ shouldn'tBlock . Throw.readBChan+ ]+ where+ nonEmpty = ChanNonEmpty Just+ nonEmptyThrow = ChanNonEmpty id+ emptyBlock = ChanEmpty $ (@? "Read should block")+ emptyNonBlock = ChanEmpty $ (@? "Read shouldn't block")+ emptyThrow = ChanEmpty $ expect ReadFailed++tryReadTests :: TestTree+tryReadTests = testGroup "try read tests"+ [ readBoilerPlate "try read non-empty" ChanOpen nonEmpty $+ shouldn'tBlock . tryReadBChan+ , readBoilerPlate "try read non-empty (throw)" ChanOpen nonEmptyThrow $+ shouldn'tBlock . Throw.tryReadBChan+ , readBoilerPlate "try read non-empty closed" ChanClosed nonEmpty $+ shouldn'tBlock . tryReadBChan+ , readBoilerPlate "try read non-empty closed (throw)" ChanClosed nonEmptyThrow $+ shouldn'tBlock . Throw.tryReadBChan+ , readBoilerPlate "try read empty" ChanOpen empty $+ shouldn'tBlock . tryReadBChan+ , readBoilerPlate "try read empty (throw)" ChanOpen empty $+ shouldn'tBlock . Throw.tryReadBChan+ , readBoilerPlate "try read empty closed" ChanClosed emptyClosed $+ shouldn'tBlock . tryReadBChan+ , readBoilerPlate "try read empty closed (throw)" ChanClosed emptyThrow $+ shouldn'tBlock . Throw.tryReadBChan+ ]+ where+ nonEmpty = ChanNonEmpty $ Just . Just+ nonEmptyThrow = ChanNonEmpty Just+ empty = ChanEmpty $ (@? "Read shouldn't block") . fmap isNothing+ emptyClosed = ChanEmpty $+ (>>= assertEqual "Expect successful nothing" (Just Nothing))+ emptyThrow = ChanEmpty $ expect ReadFailed++writeTests :: TestTree+writeTests = testGroup "write tests"+ [ writeClosed+ , writeClosedThrow+ , writeBeforeListener "write before listener" $ checkedWrite+ , writeBeforeListener "write before listener (throw)" $ Throw.writeBChan+ , writeBroadCast "write broadcast" $ checkedWrite+ , writeBroadCast "write broadcast (throw)" $ Throw.writeBChan+ ]+ where+ writeClosed :: TestTree+ writeClosed = testCase "write closed" $ do+ chan <- newBroadcastChan+ closeBChan chan+ not <$> writeBChan chan () @? "Write should fail"++ writeClosedThrow :: TestTree+ writeClosedThrow = testCase "write closed (throw)" $ do+ chan <- newBroadcastChan+ closeBChan chan+ expect WriteFailed $ Throw.writeBChan chan ()++ writeBeforeListener+ :: String -> (BroadcastChan In Int -> Int -> IO ()) -> TestTree+ writeBeforeListener name write = testCase name $ do+ inChan <- newBroadcastChan+ forM_ [1..10] $ write inChan+ closeBChan inChan+ outChan <- newBChanListener inChan+ isNothing <$> readBChan outChan @? "Read should fail"++ writeBroadCast+ :: String -> (BroadcastChan In Int -> Int -> IO ()) -> TestTree+ writeBroadCast name write = testCase name $ do+ inChan <- newBroadcastChan+ outChan1 <- newBChanListener inChan+ outChan2 <- newBChanListener inChan+ inputs <- randomList 10++ forM_ inputs $ write inChan+ closeBChan inChan+ result1 <- unfoldM $ readBChan outChan1+ result2 <- unfoldM $ readBChan outChan2+ assertEqual "Result should equal input" inputs result1+ assertEqual "Results should be equal" result1 result2++closedTests :: TestTree+closedTests = testGroup "closed tests"+ [ noBlockUnclosedIn+ , noBlockClosedIn+ , noBlockUnclosedOut+ , noBlockClosedOut+ , noBlockClosedEmptyOut+ ]+ where+ noBlockUnclosedIn :: TestTree+ noBlockUnclosedIn = testCase "no block unclosed in" $ do+ chan <- newBroadcastChan+ not <$> shouldn'tBlock (isClosedBChan chan) @? "Shouldn't be closed"++ noBlockClosedIn :: TestTree+ noBlockClosedIn = testCase "no block closed in" $ do+ chan <- newBroadcastChan+ closeBChan chan+ shouldn'tBlock (isClosedBChan chan) @? "Should be closed"++ noBlockUnclosedOut :: TestTree+ noBlockUnclosedOut = testCase "no block unclosed out" $ do+ inChan <- newBroadcastChan+ outChan <- newBChanListener inChan+ not <$> shouldn'tBlock (isClosedBChan outChan) @? "Shouldn't be closed"++ noBlockClosedOut :: TestTree+ noBlockClosedOut = testCase "no block closed out" $ do+ inChan <- newBroadcastChan+ outChan <- newBChanListener inChan+ Throw.writeBChan inChan ()+ closeBChan inChan+ not <$> shouldn'tBlock (isClosedBChan outChan) @? "Shouldn't be closed"++ noBlockClosedEmptyOut :: TestTree+ noBlockClosedEmptyOut = testCase "no block closed empty out" $ do+ inChan <- newBroadcastChan+ outChan <- newBChanListener inChan+ closeBChan inChan+ shouldn'tBlock (isClosedBChan outChan) @? "Should be closed"++chanContentsTests :: TestTree+chanContentsTests = testGroup "getBChanContents"+ [ noBlockOnEmptyIn+ , noBlockOnEmptyOut+ , noBlockOnFilledIn+ , noBlockOnFilledOut+ , checkConcurrencyOut+ ]+ where+ noBlockOnEmptyIn :: TestTree+ noBlockOnEmptyIn = testCase "no block on empty in" $ do+ chan <- newBroadcastChan+ results <- shouldn'tBlock $ getBChanContents chan+ closeBChan chan+ shouldn'tBlock $ assertBool "Should be empty list!" (null results)++ noBlockOnEmptyOut :: TestTree+ noBlockOnEmptyOut = testCase "no block on empty out" $ do+ inChan <- newBroadcastChan+ outChan <- newBChanListener inChan+ results <- shouldn'tBlock $ getBChanContents outChan+ closeBChan inChan+ shouldn'tBlock $ assertBool "Should be empty list!" (null results)++ noBlockOnFilledIn :: TestTree+ noBlockOnFilledIn = testCase "no block on filled in" $ do+ inChan <- newBroadcastChan+ throwawayInputs <- randomList 10+ forM_ throwawayInputs $ Throw.writeBChan inChan+ results <- shouldn'tBlock $ getBChanContents inChan+ inputs <- randomList 10+ forM_ inputs $ Throw.writeBChan inChan+ closeBChan inChan+ assertEqual "Should be only inputs after action" inputs results++ noBlockOnFilledOut :: TestTree+ noBlockOnFilledOut = testCase "no block on filled out" $ do+ inChan <- newBroadcastChan+ outChan <- newBChanListener inChan+ inputsBefore <- randomList 10+ forM_ inputsBefore $ Throw.writeBChan inChan+ results <- shouldn'tBlock $ getBChanContents outChan+ inputsAfter <- randomList 10+ forM_ inputsAfter $ Throw.writeBChan inChan+ closeBChan inChan+ let allInputs = inputsBefore ++ inputsAfter+ assertEqual "Should be both inputs" allInputs results++ checkConcurrencyOut :: TestTree+ checkConcurrencyOut = testCase "interleaved with readBChan" $ do+ inChan <- newBroadcastChan+ outChan <- newBChanListener inChan+ inputs <- randomList 10+ forM_ inputs $ Throw.writeBChan inChan+ closeBChan inChan+ contents <- getBChanContents outChan+ forM_ contents $ \val -> do+ result <- readBChan outChan+ case result of+ Nothing -> assertFailure "Element missing!"+ Just v -> assertEqual "Should be equal" val v++foldlTests :: TestTree+foldlTests = testGroup "foldl tests"+ [ foldBChanIn+ , foldBChanOut+ , foldBChanMIn+ , foldBChanMOut+ ]+ where+ pureFold :: Fold a b -> BroadcastChan d a -> IO (IO b)+ pureFold = Foldl.purely foldBChan++ printList :: Show a => Handle -> FoldM IO a [a]+ printList hnd = Foldl.premapM doPrint $ Foldl.generalize Foldl.list+ where+ doPrint val = val <$ hPrint hnd val++ impureFold :: FoldM IO a b -> BroadcastChan d a -> IO (IO b)+ impureFold = Foldl.impurely foldBChanM++ foldBChanIn :: TestTree+ foldBChanIn = testCase "foldBChan in" $ do+ inChan <- newBroadcastChan+ inputsBefore <- randomList 10+ forM_ inputsBefore $ Throw.writeBChan inChan+ foldList <- shouldn'tBlock $ pureFold Foldl.list inChan++ inputsAfter <- randomList 10+ forM_ inputsAfter $ Throw.writeBChan inChan+ closeBChan inChan+ (inputsAfter==) <$> shouldn'tBlock foldList @? "Lists should be equal"++ foldBChanOut :: TestTree+ foldBChanOut = testCase "foldBChan out" $ do+ inChan <- newBroadcastChan+ inputsBefore <- randomList 10+ forM_ inputsBefore $ Throw.writeBChan inChan+ outChan <- newBChanListener inChan+ inputsBetween <- randomList 10+ forM_ inputsBetween $ Throw.writeBChan inChan++ foldList <- shouldn'tBlock $ pureFold Foldl.list outChan++ inputsAfter <- randomList 10+ forM_ inputsAfter $ Throw.writeBChan inChan+ closeBChan inChan++ let expected = inputsBetween ++ inputsAfter+ (expected==) <$> shouldn'tBlock foldList @? "Lists should be equal"++ foldBChanMIn :: TestTree+ foldBChanMIn = testCase "foldBChanM in" $ do+ inChan <- newBroadcastChan+ inputsBefore <- randomList 10+ forM_ inputsBefore $ Throw.writeBChan inChan+ inputsAfter <- randomList 10++ control <- withLoggedOutput "foldBChanControl.out" $ \hnd -> do+ Foldl.foldM (printList hnd) inputsAfter++ validation <- withLoggedOutput "foldBChanM.out" $ \hnd -> do+ foldPrintList <- shouldn'tBlock $ impureFold (printList hnd) inChan+ forM_ inputsAfter $ Throw.writeBChan inChan+ closeBChan inChan+ foldPrintList++ assertEqual "Results and output should be equal" control validation++ foldBChanMOut :: TestTree+ foldBChanMOut = testCase "foldBChanM out" $ do+ inChan <- newBroadcastChan+ inputsBefore <- randomList 10+ forM_ inputsBefore $ Throw.writeBChan inChan+ outChan <- newBChanListener inChan+ inputsBetween <- randomList 10+ forM_ inputsBetween $ Throw.writeBChan inChan+ inputsAfter <- randomList 10+ let expected = inputsBetween ++ inputsAfter++ control <- withLoggedOutput "foldBChanControl.out" $ \hnd -> do+ Foldl.foldM (printList hnd) expected++ validation <- withLoggedOutput "foldBChanM.out" $ \hnd -> do+ foldPrintList <- shouldn'tBlock $+ impureFold (printList hnd) outChan++ forM_ inputsAfter $ Throw.writeBChan inChan+ closeBChan inChan+ foldPrintList++ assertEqual "Results and output should be equal" control validation++main :: IO ()+main = runTests "basic"+ [ readTests+ , tryReadTests+ , writeTests+ , closedTests+ , chanContentsTests+ , foldlTests+ ]
+ tests/ConduitTest.hs view
@@ -0,0 +1,37 @@+import Control.Monad (void)+import Data.Set (Set)+import qualified Data.Set as S+import Data.Conduit+import qualified Data.Conduit.List as C++import BroadcastChan.Conduit+import BroadcastChan.Test++sequentialSink :: [a] -> (a -> IO b) -> IO ()+sequentialSink inputs f =+ runConduitRes $ C.sourceList inputs .| C.mapM_ (liftIO . void . f)++parallelSink :: Handler IO a -> [a] -> (a -> IO b) -> Int -> IO ()+parallelSink hnd inputs f n = runConduitRes $+ C.sourceList inputs .| parMapM_ handler n (liftIO . void . f)+ where+ handler = mapHandler liftIO hnd++sequentialFold :: Ord b => [a] -> (a -> IO b) -> IO (Set b)+sequentialFold inputs f = runConduitRes $+ C.sourceList inputs .| C.mapM (liftIO . f) .| C.foldMap S.singleton++parallelFold+ :: Ord b => Handler IO a -> [a] -> (a -> IO b) -> Int -> IO (Set b)+parallelFold hnd inputs f n = runConduitRes $+ C.sourceList inputs+ .| parMapM handler n (liftIO . f)+ .| C.foldMap S.singleton+ where+ handler = mapHandler liftIO hnd++main :: IO ()+main = runTests "conduit" $+ [ genStreamTests "sink" sequentialSink parallelSink+ , genStreamTests "fold" sequentialFold parallelFold+ ]
+ tests/IOTest.hs view
@@ -0,0 +1,47 @@+{-# LANGUAGE BangPatterns #-}++import Control.Monad (void)+import Data.Foldable (forM_, foldlM)+import Data.Set (Set)+import qualified Data.Set as S++import BroadcastChan+import BroadcastChan.Test++sequentialSink :: Foldable f => f a -> (a -> IO b) -> IO ()+sequentialSink set f = forM_ set (void . f)++parallelSink+ :: Foldable f => Handler IO a -> f a -> (a -> IO b) -> Int -> IO ()+parallelSink hnd input f n =+ parMapM_ hnd n (void . f) input++sequentialFold :: (Foldable f, Ord b) => f a -> (a -> IO b) -> IO (Set b)+sequentialFold input f = foldlM foldFun S.empty input+ where+ foldFun bs a = (\b -> S.insert b bs) <$> f a++parallelFold+ :: (Foldable f, Ord b)+ => Handler IO a -> f a -> (a -> IO b) -> Int -> IO (Set b)+parallelFold hnd input f n =+ parFoldMap hnd n f foldFun S.empty input+ where+ foldFun :: Ord b => Set b -> b -> Set b+ foldFun s b = S.insert b s++parallelFoldM+ :: (Foldable f, Ord b)+ => Handler IO a -> f a -> (a -> IO b) -> Int -> IO (Set b)+parallelFoldM hnd input f n =+ parFoldMapM hnd n f foldFun S.empty input+ where+ foldFun :: (Ord b, Monad m) => Set b -> b -> m (Set b)+ foldFun !z b = return $ S.insert b z++main :: IO ()+main = runTests "parallel-io" $+ [ genStreamTests "sink" sequentialSink parallelSink+ , genStreamTests "fold" sequentialFold parallelFold+ , genStreamTests "foldM" sequentialFold parallelFoldM+ ]
+ tests/PipeTest.hs view
@@ -0,0 +1,35 @@+import Control.Foldl (purely, set)+import Data.Set (Set)+import Pipes+import qualified Pipes.Prelude as P+import Pipes.Safe (runSafeT)++import BroadcastChan.Pipes+import BroadcastChan.Test++sequentialSink :: [a] -> (a -> IO b) -> IO ()+sequentialSink inputs f = runSafeT . runEffect $+ each inputs >-> P.mapM_ (liftIO . void. f)++parallelSink :: Handler IO a -> [a] -> (a -> IO b) -> Int -> IO ()+parallelSink hnd inputs f n =+ runSafeT . runEffect $ parMapM_ handler n (void . f) $ each inputs+ where+ handler = mapHandler liftIO hnd++sequentialFold :: Ord b => [a] -> (a -> IO b) -> IO (Set b)+sequentialFold inputs f = runSafeT $ purely P.fold set $+ each inputs >-> P.mapM (liftIO . f)++parallelFold+ :: Ord b => Handler IO a -> [a] -> (a -> IO b) -> Int -> IO (Set b)+parallelFold hnd inputs f n =+ runSafeT . purely P.fold set . parMapM handler n f $ each inputs+ where+ handler = mapHandler liftIO hnd++main :: IO ()+main = runTests "pipes" $+ [ genStreamTests "sink" sequentialSink parallelSink+ , genStreamTests "fold" sequentialFold parallelFold+ ]