stm-conduit 2.1.0 → 2.1.1
raw patch · 2 files changed
+60/−7 lines, 2 filesdep +asyncdep +monad-control
Dependencies added: async, monad-control
Files
- Data/Conduit/Async.hs +50/−0
- stm-conduit.cabal +10/−7
+ Data/Conduit/Async.hs view
@@ -0,0 +1,50 @@+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE RankNTypes #-}++-- | * Introduction+--+-- Contains a combinator for concurrently joining a producer and a consumer,+-- such that the producer may continue to produce (up to the queue size) as+-- the consumer is concurrently consuming.+module Data.Conduit.Async where++import Control.Concurrent.Async+import Control.Concurrent.STM+import Control.Monad.IO.Class+import Control.Monad.Trans.Control+import Data.Conduit+import Data.Conduit.List+import Prelude hiding (mapM_)++-- | Concurrently join the producer and consumer, using a bounded queue of the+-- given size. The producer will block when the queue is full, if it is+-- producing faster than the consumers is taking from it. Likewise, if the+-- consumer races ahead, it will block until more input is available.+--+-- Exceptions are properly managed and propagated between the two sides, so+-- the net effect should be equivalent to not using buffer at all, save for+-- the concurrent interleaving of effects.+buffer :: (MonadBaseControl IO m, MonadIO m)+ => Int -> Producer m a -> Consumer a m b -> m b+buffer size input output = do+ chan <- liftIO $ newTBQueueIO size+ control $ \runInIO ->+ withAsync (runInIO $ input $$ mapM_ (send chan)) $ \input' ->+ withAsync (runInIO $ recv chan $$ output) $ \output' -> do+ link2 input' output'+ wait output'+ where+ send chan = liftIO . atomically . writeTBQueue chan . Just++ recv chan = do+ mx <- liftIO $ atomically $ readTBQueue chan+ case mx of+ Nothing -> return ()+ Just x -> yield x >> recv chan++-- | An operator form of 'buffer'. In general you should be able to replace+-- any use of 'Data.Conduit.$$' with '$$&' and suddenly reap the benefit of+-- concurrency, if your conduits were spending time waiting on each other.+($$&) :: (MonadIO m, MonadBaseControl IO m)+ => Producer m a -> Consumer a m b -> m b+($$&) = buffer 64
stm-conduit.cabal view
@@ -1,5 +1,5 @@ Name: stm-conduit-Version: 2.1.0+Version: 2.1.1 Synopsis: Introduces conduits to channels, and promotes using conduits concurrently. Description: Provides two simple conduit wrappers around STM@@ -17,16 +17,19 @@ Library exposed-modules:+ Data.Conduit.Async Data.Conduit.TMChan Data.Conduit.TQueue build-depends:- base == 4.*- , transformers >= 0.2 && <= 0.4- , stm == 2.4.*- , stm-chans >= 2.0 && < 3.1- , conduit == 1.0.*- , resourcet >= 0.3 && < 0.5+ base == 4.*+ , transformers >= 0.2 && <= 0.4+ , stm == 2.4.*+ , stm-chans >= 2.0 && < 3.1+ , conduit == 1.0.*+ , resourcet >= 0.3 && < 0.5+ , async >= 2.0.1+ , monad-control >= 0.3.2 ghc-options: -Wall -fwarn-tabs -fwarn-unused-imports