packages feed

conduit-algorithms 0.0.8.1 → 0.0.8.2

raw patch · 9 files changed

+104/−49 lines, 9 filesPVP: major bump suggested

API removals or changes: PVP suggests a major version bump

API changes (from Hackage documentation)

+ Data.Conduit.Algorithms.Async: withPossiblyCompressedFile :: (MonadUnliftIO m, MonadResource m, MonadThrow m) => FilePath -> (ConduitT () ByteString m () -> m a) -> m a
- Data.Conduit.Algorithms: mergeC :: (Ord a, Monad m) => [Source m a] -> Source m a
+ Data.Conduit.Algorithms: mergeC :: (Ord a, Monad m) => [ConduitT () a m ()] -> ConduitT () a m ()
- Data.Conduit.Algorithms: mergeC2 :: (Ord a, Monad m) => Source m a -> Source m a -> Source m a
+ Data.Conduit.Algorithms: mergeC2 :: (Ord a, Monad m) => ConduitT () a m () -> ConduitT () a m () -> ConduitT () a m ()
- Data.Conduit.Algorithms: removeRepeatsC :: (Eq a, Monad m) => Conduit a m a
+ Data.Conduit.Algorithms: removeRepeatsC :: (Eq a, Monad m) => ConduitT a a m ()
- Data.Conduit.Algorithms: uniqueC :: (Ord a, Monad m) => Conduit a m a
+ Data.Conduit.Algorithms: uniqueC :: (Ord a, Monad m) => ConduitT a a m ()
- Data.Conduit.Algorithms: uniqueOnC :: (Ord b, Monad m) => (a -> b) -> Conduit a m a
+ Data.Conduit.Algorithms: uniqueOnC :: (Ord b, Monad m) => (a -> b) -> ConduitT a a m ()
- Data.Conduit.Algorithms.Async: asyncBzip2From :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> Source m ByteString
+ Data.Conduit.Algorithms.Async: asyncBzip2From :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> ConduitT () ByteString m ()
- Data.Conduit.Algorithms.Async: asyncBzip2FromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> Source m ByteString
+ Data.Conduit.Algorithms.Async: asyncBzip2FromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> ConduitT () ByteString m ()
- Data.Conduit.Algorithms.Async: asyncBzip2To :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> Sink ByteString m ()
+ Data.Conduit.Algorithms.Async: asyncBzip2To :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> ConduitT ByteString Void m ()
- Data.Conduit.Algorithms.Async: asyncBzip2ToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> Sink ByteString m ()
+ Data.Conduit.Algorithms.Async: asyncBzip2ToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> ConduitT ByteString Void m ()
- Data.Conduit.Algorithms.Async: asyncGzipFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> Source m ByteString
+ Data.Conduit.Algorithms.Async: asyncGzipFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> ConduitT () ByteString m ()
- Data.Conduit.Algorithms.Async: asyncGzipFromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> Source m ByteString
+ Data.Conduit.Algorithms.Async: asyncGzipFromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> ConduitT () ByteString m ()
- Data.Conduit.Algorithms.Async: asyncGzipTo :: forall m. (MonadIO m, MonadUnliftIO m) => Handle -> Sink ByteString m ()
+ Data.Conduit.Algorithms.Async: asyncGzipTo :: forall m. (MonadIO m, MonadUnliftIO m) => Handle -> ConduitT ByteString Void m ()
- Data.Conduit.Algorithms.Async: asyncGzipToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> Sink ByteString m ()
+ Data.Conduit.Algorithms.Async: asyncGzipToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> ConduitT ByteString Void m ()
- Data.Conduit.Algorithms.Async: asyncMapC :: forall a m b. (MonadIO m, NFData b) => Int -> (a -> b) -> Conduit a m b
+ Data.Conduit.Algorithms.Async: asyncMapC :: forall a m b. (MonadIO m, NFData b) => Int -> (a -> b) -> ConduitT a b m ()
- Data.Conduit.Algorithms.Async: asyncMapEitherC :: forall a m b e. (MonadIO m, NFData b, NFData e, MonadError e m) => Int -> (a -> Either e b) -> Conduit a m b
+ Data.Conduit.Algorithms.Async: asyncMapEitherC :: forall a m b e. (MonadIO m, NFData b, NFData e, MonadError e m) => Int -> (a -> Either e b) -> ConduitT a b m ()
- Data.Conduit.Algorithms.Async: asyncXzFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m, MonadThrow m) => Handle -> Source m ByteString
+ Data.Conduit.Algorithms.Async: asyncXzFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m, MonadThrow m) => Handle -> ConduitT () ByteString m ()
- Data.Conduit.Algorithms.Async: asyncXzFromFile :: forall m. (MonadResource m, MonadUnliftIO m, MonadThrow m) => FilePath -> Source m ByteString
+ Data.Conduit.Algorithms.Async: asyncXzFromFile :: forall m. (MonadResource m, MonadUnliftIO m, MonadThrow m) => FilePath -> ConduitT () ByteString m ()
- Data.Conduit.Algorithms.Async: asyncXzTo :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> Sink ByteString m ()
+ Data.Conduit.Algorithms.Async: asyncXzTo :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> ConduitT ByteString Void m ()
- Data.Conduit.Algorithms.Async: asyncXzToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> Sink ByteString m ()
+ Data.Conduit.Algorithms.Async: asyncXzToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> ConduitT ByteString Void m ()
- Data.Conduit.Algorithms.Async: conduitPossiblyCompressedFile :: (MonadUnliftIO m, MonadResource m, MonadThrow m) => FilePath -> Source m ByteString
+ Data.Conduit.Algorithms.Async: conduitPossiblyCompressedFile :: (MonadUnliftIO m, MonadResource m, MonadThrow m) => FilePath -> ConduitT () ByteString m ()
- Data.Conduit.Algorithms.Async: conduitPossiblyCompressedToFile :: (MonadUnliftIO m, MonadResource m) => FilePath -> Sink ByteString m ()
+ Data.Conduit.Algorithms.Async: conduitPossiblyCompressedToFile :: (MonadUnliftIO m, MonadResource m) => FilePath -> ConduitT ByteString Void m ()
- Data.Conduit.Algorithms.Async: unorderedAsyncMapC :: forall a m b. (MonadIO m, NFData b) => Int -> (a -> b) -> Conduit a m b
+ Data.Conduit.Algorithms.Async: unorderedAsyncMapC :: forall a m b. (MonadIO m, NFData b) => Int -> (a -> b) -> ConduitT a b m ()
- Data.Conduit.Algorithms.Async.ByteString: asyncFilterLinesC :: MonadIO m => Int -> (ByteString -> Bool) -> Conduit ByteString m ByteString
+ Data.Conduit.Algorithms.Async.ByteString: asyncFilterLinesC :: MonadIO m => Int -> (ByteString -> Bool) -> ConduitT ByteString ByteString m ()
- Data.Conduit.Algorithms.Async.ByteString: asyncMapLineGroupsC :: (MonadIO m, NFData a) => Int -> ([ByteString] -> a) -> Conduit ByteString m a
+ Data.Conduit.Algorithms.Async.ByteString: asyncMapLineGroupsC :: (MonadIO m, NFData a) => Int -> ([ByteString] -> a) -> ConduitT ByteString a m ()
- Data.Conduit.Algorithms.Storable: writeStorableV :: forall m a. (MonadIO m, Monad m, Storable a) => Conduit (Vector a) m ByteString
+ Data.Conduit.Algorithms.Storable: writeStorableV :: forall m a. (MonadIO m, Monad m, Storable a) => ConduitT (Vector a) ByteString m ()
- Data.Conduit.Algorithms.Utils: awaitJust :: Monad m => (a -> Conduit a m b) -> Conduit a m b
+ Data.Conduit.Algorithms.Utils: awaitJust :: Monad m => (a -> ConduitT a b m ()) -> ConduitT a b m ()
- Data.Conduit.Algorithms.Utils: enumerateC :: Monad m => Conduit a m (Int, a)
+ Data.Conduit.Algorithms.Utils: enumerateC :: Monad m => ConduitT a (Int, a) m ()
- Data.Conduit.Algorithms.Utils: groupC :: (Monad m) => Int -> Conduit a m [a]
+ Data.Conduit.Algorithms.Utils: groupC :: (Monad m) => Int -> ConduitT a [a] m ()

Files

ChangeLog view
@@ -1,3 +1,6 @@+Version 0.0.8.2 2018-09-27 by luispedro+	* Add withPossiblyCompressedFile function for prompt resource deallocation+ Version 0.0.8.1 2018-05-29 by luispedro 	* faster mergeC (use a priority queue instead of an ordered list) 
Data/Conduit/Algorithms.hs view
@@ -34,7 +34,7 @@ -- it uses a 'Data.Set' to store previously seen elements. Thus, memory usage -- is O(N) and time is O(N log N). If the input is sorted, you can use -- 'removeRepeatsC'-uniqueOnC :: (Ord b, Monad m) => (a -> b) -> C.Conduit a m a+uniqueOnC :: (Ord b, Monad m) => (a -> b) -> C.ConduitT a a m () uniqueOnC f = checkU (S.empty :: S.Set b)     where         checkU cur = awaitJust $ \val ->@@ -46,7 +46,7 @@ -- | Unique conduit -- -- See 'uniqueOnC' and 'removeRepeatsC'-uniqueC :: (Ord a, Monad m) => C.Conduit a m a+uniqueC :: (Ord a, Monad m) => C.ConduitT a a m () uniqueC = uniqueOnC id  -- | Removes repeated elements@@ -58,7 +58,7 @@ -- is equivalent to @[0, 1, 2, 0]@ -- -- See 'uniqueC' and 'uniqueOnC'-removeRepeatsC :: (Eq a, Monad m) => C.Conduit a m a+removeRepeatsC :: (Eq a, Monad m) => C.ConduitT a a m () removeRepeatsC = awaitJust removeRepeatsC'     where         removeRepeatsC' prev = C.await >>= \case@@ -76,7 +76,7 @@ -- all elements in sorted order. -- -- See 'mergeC2'-mergeC :: (Ord a, Monad m) => [C.Source m a] -> C.Source m a+mergeC :: (Ord a, Monad m) => [C.ConduitT () a m ()] -> C.ConduitT () a m () mergeC [a] = a mergeC [a,b] = mergeC2 a b mergeC cs = CI.ConduitT $ \rest -> let@@ -86,19 +86,18 @@             _ -> error "This situation should have been impossible (mergeC/go)"         norm1insert :: (Monad m, Ord o) => PQ.MinPQueue o (CI.Pipe () i o () m ()) -> CI.Pipe () i o () m () -> CI.Pipe () i o () m (PQ.MinPQueue o (CI.Pipe () i o () m ()))         norm1insert q c@(CI.HaveOutput _ v) = return (PQ.insert v c q)-        norm1insert q c@CI.Done{} = return q+        norm1insert q CI.Done{} = return q         norm1insert q (CI.PipeM p) = lift p >>= norm1insert q         norm1insert q (CI.NeedInput _ next) = norm1insert q (next ())         norm1insert q (CI.Leftover next ()) = norm1insert q next     in do         let st = map (($ CI.Done) . CI.unConduitT) cs-        init <- foldM norm1insert PQ.empty st-        go init+        go =<< foldM norm1insert PQ.empty st  -- | Take two sorted sources and merge them. -- -- See 'mergeC'-mergeC2 :: (Ord a, Monad m) => C.Source m a -> C.Source m a -> C.Source m a+mergeC2 :: (Ord a, Monad m) => C.ConduitT () a m () -> C.ConduitT () a m () -> C.ConduitT () a m () mergeC2 (CI.ConduitT s1) (CI.ConduitT s2) = CI.ConduitT $ \rest -> let         go right@(CI.HaveOutput s1' v1) left@(CI.HaveOutput s2' v2)             | v1 <= v2 = CI.HaveOutput (go s1' left) v1
Data/Conduit/Algorithms/Async.hs view
@@ -11,6 +11,7 @@ module Data.Conduit.Algorithms.Async     ( conduitPossiblyCompressedFile     , conduitPossiblyCompressedToFile+    , withPossiblyCompressedFile     , asyncMapC     , asyncMapEitherC     , asyncGzipTo@@ -50,13 +51,14 @@ import           Data.Foldable (toList) import           Control.Monad.IO.Class (MonadIO, liftIO) import           Control.Monad.Error.Class (MonadError(..))-import           Control.Monad.IO.Unlift (MonadUnliftIO)+import           Control.Monad.IO.Unlift (MonadUnliftIO, withRunInIO) import           Control.Monad.Trans.Resource (MonadResource) import           Control.Monad.Catch (MonadThrow) import           Control.Exception (evaluate, displayException) import           Control.DeepSeq (NFData, force) import           System.IO.Error (mkIOError, userErrorType) import           System.IO+import qualified System.IO as IO import           Data.List (isSuffixOf) import           Data.Conduit.Algorithms.Utils (awaitJust) @@ -83,7 +85,7 @@ asyncMapC :: forall a m b . (MonadIO m, NFData b) =>                     Int -- ^ Maximum number of worker threads                     -> (a -> b) -- ^ Function to execute-                    -> C.Conduit a m b+                    -> C.ConduitT a b m () asyncMapC = asyncMapCHelper True  -- | A version of 'asyncMapC' which can reorder results in the stream@@ -95,17 +97,17 @@ unorderedAsyncMapC :: forall a m b . (MonadIO m, NFData b) =>                     Int -- ^ Maximum number of worker threads                     -> (a -> b) -- ^ Function to execute-                    -> C.Conduit a m b+                    -> C.ConduitT a b m () unorderedAsyncMapC = asyncMapCHelper False  asyncMapCHelper  :: forall a m b . (MonadIO m, NFData b) =>                     Bool                     -> Int -- ^ Maximum number of worker threads                     -> (a -> b) -- ^ Function to execute-                    -> C.Conduit a m b+                    -> C.ConduitT a b m () asyncMapCHelper isSynchronous maxThreads f = initLoop (0 :: Int) (Seq.empty :: Seq.Seq (A.Async b))     where-        initLoop :: Int -> Seq.Seq (A.Async b) -> C.Conduit a m b+        initLoop :: Int -> Seq.Seq (A.Async b) -> C.ConduitT a b m ()         initLoop size q             | size == maxThreads = loop q             | otherwise = C.await >>= \case@@ -117,7 +119,7 @@         sched v = liftIO . A.async . evaluate . force $ f v          -- | yield all-        yAll :: Seq.Seq (A.Async b) -> C.Conduit a m b+        yAll :: Seq.Seq (A.Async b) -> C.ConduitT a b m ()         yAll q             | Seq.null q = return ()             | otherwise = do@@ -125,7 +127,7 @@                 C.yield r                 yAll q' -        loop :: Seq.Seq (A.Async b) -> C.Conduit a m b+        loop :: Seq.Seq (A.Async b) -> C.ConduitT a b m ()         loop q = C.await >>= \case                 Nothing -> yAll q                 Just v -> do@@ -151,7 +153,7 @@ -- processing elements after the first error). -- -- See 'asyncMapC'-asyncMapEitherC :: forall a m b e . (MonadIO m, NFData b, NFData e, MonadError e m) => Int -> (a -> Either e b) -> C.Conduit a m b+asyncMapEitherC :: forall a m b e . (MonadIO m, NFData b, NFData e, MonadError e m) => Int -> (a -> Either e b) -> C.ConduitT a b m () asyncMapEitherC maxThreads f = asyncMapC maxThreads f .| (C.awaitForever $ \case                                 Right v -> C.yield v                                 Left err -> throwError err)@@ -164,7 +166,7 @@ -- the chunkSize parameter is a hint, not an exact element. In particular, -- larger chunks are not split up and smaller chunks can be yielded too. bsConcatTo :: MonadIO m => Int -- ^ chunk hint-                            -> C.Conduit B.ByteString m [B.ByteString]+                            -> C.ConduitT B.ByteString [B.ByteString] m () bsConcatTo chunkSize = awaitJust start     where         start v@@ -176,9 +178,9 @@                 | B.length v + s > chunkSize -> C.yield chunks >> start v                 | otherwise -> continue (v:chunks) (s + B.length v) -untilNothing :: forall m i. (Monad m) => C.Conduit (Maybe i) m i-untilNothing = C.await >>= \case-    Just (Just val) -> do+untilNothing :: forall m i. (Monad m) => C.ConduitT (Maybe i) i m ()+untilNothing = awaitJust $ \case+    Just val -> do         C.yield val         untilNothing     _ -> return ()@@ -187,7 +189,7 @@ -- writes the results to `h`. -- -- See also 'asyncGzipToFile'-asyncGzipTo :: forall m. (MonadIO m, MonadUnliftIO m) => Handle -> C.Sink B.ByteString m ()+asyncGzipTo :: forall m. (MonadIO m, MonadUnliftIO m) => Handle -> C.ConduitT B.ByteString C.Void m () asyncGzipTo h = do     let drain q = liftIO . C.runConduit $                 CA.sourceTBQueue q@@ -202,7 +204,7 @@ -- performed in a separate thread. -- -- See also 'asyncGzipTo'-asyncGzipToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.Sink B.ByteString m ()+asyncGzipToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.ConduitT B.ByteString C.Void m () asyncGzipToFile fname = C.bracketP     (openFile fname WriteMode)     hClose@@ -213,7 +215,7 @@ -- will probably be left at an undefined position in the file. -- -- See also 'asyncGzipFromFile'-asyncGzipFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.Source m B.ByteString+asyncGzipFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.ConduitT () B.ByteString m () asyncGzipFrom h = do     let prod q = liftIO $ do                     C.runConduit $@@ -230,7 +232,7 @@ -- separate thread. -- -- See also 'asyncGzipFrom'-asyncGzipFromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.Source m B.ByteString+asyncGzipFromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.ConduitT () B.ByteString m () asyncGzipFromFile fname = C.bracketP     (openFile fname ReadMode)     hClose@@ -240,7 +242,7 @@ -- writes the results to `h`. -- -- See also 'asyncBzip2ToFile'-asyncBzip2To :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.Sink B.ByteString m ()+asyncBzip2To :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.ConduitT B.ByteString C.Void m () asyncBzip2To h = do     let drain q = C.runConduit $                 CA.sourceTBQueue q@@ -255,7 +257,7 @@ -- performed in a separate thread. -- -- See also 'asyncBzip2To'-asyncBzip2ToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.Sink B.ByteString m ()+asyncBzip2ToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.ConduitT B.ByteString C.Void m () asyncBzip2ToFile fname = C.bracketP     (openFile fname WriteMode)     hClose@@ -266,7 +268,7 @@ -- will probably be left at an undefined position in the file. -- -- See also 'asyncBzip2FromFile'-asyncBzip2From :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.Source m B.ByteString+asyncBzip2From :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.ConduitT () B.ByteString m () asyncBzip2From h = do     let prod q = do                     C.runConduit $@@ -281,7 +283,7 @@ -- separate thread. -- -- See also 'asyncBzip2From'-asyncBzip2FromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.Source m B.ByteString+asyncBzip2FromFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.ConduitT () B.ByteString m () asyncBzip2FromFile fname = C.bracketP     (openFile fname ReadMode)     hClose@@ -291,7 +293,7 @@ -- writes the results to `h`. -- -- See also 'asyncXzToFile'-asyncXzTo :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.Sink B.ByteString m ()+asyncXzTo :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m) => Handle -> C.ConduitT B.ByteString C.Void m () asyncXzTo h = do     let drain q = C.runConduit $                 CA.sourceTBQueue q@@ -306,7 +308,7 @@ -- performed in a separate thread. -- -- See also 'asyncXzTo'-asyncXzToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.Sink B.ByteString m ()+asyncXzToFile :: forall m. (MonadResource m, MonadUnliftIO m) => FilePath -> C.ConduitT B.ByteString C.Void m () asyncXzToFile fname = C.bracketP     (openFile fname WriteMode)     hClose@@ -317,7 +319,7 @@ -- will probably be left at an undefined position in the file. -- -- See also 'asyncXzFromFile'-asyncXzFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m, MonadThrow m) => Handle -> C.Source m B.ByteString+asyncXzFrom :: forall m. (MonadIO m, MonadResource m, MonadUnliftIO m, MonadThrow m) => Handle -> C.ConduitT () B.ByteString m () asyncXzFrom h = do     let oneGBmembuffer = Just $ 1024 ^ (3 :: Integer)         prod q = do@@ -334,17 +336,50 @@ -- separate thread. -- -- See also 'asyncXzFrom'-asyncXzFromFile :: forall m. (MonadResource m, MonadUnliftIO m, MonadThrow m) => FilePath -> C.Source m B.ByteString+asyncXzFromFile :: forall m. (MonadResource m, MonadUnliftIO m, MonadThrow m) => FilePath -> C.ConduitT () B.ByteString m () asyncXzFromFile fname = C.bracketP     (openFile fname ReadMode)     hClose     asyncXzFrom ++-- | If the filename indicates a supported compressed file (gzip, xz, and, on+-- Unix, bzip2), then it reads it and uncompresses it.+--+-- Usage+--+-- @+--+--      withPossiblyCompressedFile fname $ \src ->+--          src .| mySink+-- @+--+-- Unlike 'conduitPossiblyCompressedFile', this ensures that the file is closed+-- even if the conduit terminates early.+--+-- On Windows, attempting to read from a bzip2 file, results in 'error'.+withPossiblyCompressedFile :: (MonadUnliftIO m, MonadResource m, MonadThrow m) => FilePath -> (C.ConduitT () B.ByteString m () -> m a) -> m a+withPossiblyCompressedFile fname inner = withRunInIO $ \run -> do+    IO.withBinaryFile fname IO.ReadMode $+        run . inner . withPossiblyCompressedFile' fname++withPossiblyCompressedFile' :: (MonadUnliftIO m, MonadResource m, MonadThrow m) => FilePath -> Handle -> C.ConduitT () B.ByteString m ()+withPossiblyCompressedFile' fname+    | ".gz" `isSuffixOf` fname = asyncGzipFrom+    | ".xz" `isSuffixOf` fname = asyncXzFrom+    | ".bz2" `isSuffixOf` fname = asyncBzip2From+    | otherwise = C.sourceHandle++ -- | If the filename indicates a gzipped file (or, on Unix, also a bz2 file), -- then it reads it and uncompresses it. --+--+-- To ensure that the file is closed even if the downstream finishes early,+-- consider using 'withPossiblyCompressedFile'.+-- -- On Windows, attempting to read from a bzip2 file, results in 'error'.-conduitPossiblyCompressedFile :: (MonadUnliftIO m, MonadResource m, MonadThrow m) => FilePath -> C.Source m B.ByteString+conduitPossiblyCompressedFile :: (MonadUnliftIO m, MonadResource m, MonadThrow m) => FilePath -> C.ConduitT () B.ByteString m () conduitPossiblyCompressedFile fname     | ".gz" `isSuffixOf` fname = asyncGzipFromFile fname     | ".xz" `isSuffixOf` fname = asyncXzFromFile fname@@ -355,7 +390,7 @@ -- then it compresses and write with the algorithm matching the filename -- -- On Windows, attempting to write to a bzip2 file, results in 'error'.-conduitPossiblyCompressedToFile :: (MonadUnliftIO m, MonadResource m) => FilePath -> C.Sink B.ByteString m ()+conduitPossiblyCompressedToFile :: (MonadUnliftIO m, MonadResource m) => FilePath -> C.ConduitT B.ByteString C.Void m () conduitPossiblyCompressedToFile fname     | ".gz" `isSuffixOf` fname = asyncGzipToFile fname     | ".xz" `isSuffixOf` fname = asyncXzToFile fname
Data/Conduit/Algorithms/Async/ByteString.hs view
@@ -44,14 +44,14 @@ -- Only Unix-style ASCII lines are supported (splitting at Bytes with value -- 10, i.e., \\n). When Windows lines (\\r\\n) are passed to this function, this -- results in each element having an extra \\r at the end.-asyncMapLineGroupsC :: (MonadIO m, NFData a) => Int -> ([B.ByteString] -> a) -> C.Conduit B.ByteString m a+asyncMapLineGroupsC :: (MonadIO m, NFData a) => Int -> ([B.ByteString] -> a) -> C.ConduitT B.ByteString a m () asyncMapLineGroupsC nthreads f = breakAtLineBoundary .| CAlg.asyncMapC nthreads (f . asLines)     where         asLines :: BL.ByteString -> [B.ByteString]         asLines = fmap BL.toStrict . BL.split 10          -- The purpose is to break input blocks at a line boundary-        breakAtLineBoundary :: Monad m => C.Conduit B.ByteString m BL.ByteString+        breakAtLineBoundary :: Monad m => C.ConduitT B.ByteString BL.ByteString m ()         breakAtLineBoundary = continue BL.empty         continue prev = C.await >>= \case                     Nothing -> unless (BL.null prev) $@@ -72,7 +72,7 @@ --      CB.lines .| CL.filer f -- @ ---asyncFilterLinesC :: MonadIO m => Int -> (B.ByteString -> Bool) -> C.Conduit B.ByteString m B.ByteString+asyncFilterLinesC :: MonadIO m => Int -> (B.ByteString -> Bool) -> C.ConduitT B.ByteString B.ByteString m () asyncFilterLinesC n f = asyncMapLineGroupsC n (filter f) .| CL.concat {-# INLINE asyncFilterLinesC #-} 
Data/Conduit/Algorithms/Storable.hs view
@@ -4,7 +4,7 @@ License     : MIT Maintainer  : luis@luispedro.org -Higher level async processing interfaces.+Read/write Storable vectors -} {-# LANGUAGE FlexibleContexts, ScopedTypeVariables #-} module Data.Conduit.Algorithms.Storable@@ -33,7 +33,7 @@ -- This uses the same format as in-memory -- -- See 'readStorableV'-writeStorableV :: forall m a. (MonadIO m, Monad m, Storable a) => C.Conduit (VS.Vector a) m B.ByteString+writeStorableV :: forall m a. (MonadIO m, Monad m, Storable a) => C.ConduitT (VS.Vector a) B.ByteString m () writeStorableV = CL.mapM (liftIO. encodeStorable')     where         encodeStorable' :: Storable a => VS.Vector a -> IO B.ByteString@@ -43,8 +43,12 @@  -- | read a Storable vector ----- This expects the same format as the in-memory vector+-- This expects the same format as the in-memory vector. --+-- This will break up the incoming data into vectors of the given size. The+-- last vector may be smaller if there is not enough data. Any unconsumed Bytes+-- will be leftover for the next conduit in the pipeline.+-- -- See 'writeStorableV' readStorableV :: forall m a. (MonadIO m, Storable a) => Int -> C.ConduitM B.ByteString (VS.Vector a) m () readStorableV nelems = CC.chunksOfE blockBytes .| parseBlocks@@ -54,7 +58,7 @@         a' = undefined  -        parseBlocks :: MonadIO m => C.Conduit B.ByteString m (VS.Vector a)+        parseBlocks :: MonadIO m => C.ConduitT B.ByteString (VS.Vector a) m ()         parseBlocks = C.awaitForever $ \bs -> do             let (n,rest) = B.length bs `divMod` sizeOf a'             r <- liftIO $ do
Data/Conduit/Algorithms/Utils.hs view
@@ -28,12 +28,12 @@ -- -- This is a simple utility adapted from -- http://neilmitchell.blogspot.de/2015/07/thoughts-on-conduits.html-awaitJust :: Monad m => (a -> C.Conduit a m b) -> C.Conduit a m b+awaitJust :: Monad m => (a -> C.ConduitT a b m ()) -> C.ConduitT a b m () awaitJust f = C.await >>= maybe (return ()) f {-# INLINE awaitJust #-}  -- | Conduit analogue to Python's enumerate function-enumerateC :: Monad m => C.Conduit a m (Int, a)+enumerateC :: Monad m => C.ConduitT a (Int, a) m () enumerateC = enumerateC' 0     where         enumerateC' !i = awaitJust $ \v -> do@@ -53,7 +53,7 @@ -- results in @[ [0,1,2], [3,4,5], [6,7,8], [9, 10] ]@ -- -- This function is deprecated; use 'Data.Conduit.List.chunksOf'-groupC :: (Monad m) => Int -> C.Conduit a m [a]+groupC :: (Monad m) => Int -> C.ConduitT a [a] m () groupC n = loop n []     where         loop 0 ps = C.yield (reverse ps) >> loop n []
README.md view
@@ -7,7 +7,8 @@ Some conduit-based algorithms.  Much of this code was originally part of [NGLess](http://ngless.embl.de) and-has been in production use for years. However, it can be of generic use.+has been in production use for years. However, it was spun of from that project+as it can be of generic use.  License: MIT 
conduit-algorithms.cabal view
@@ -1,11 +1,11 @@--- This file has been generated from package.yaml by hpack version 0.21.2.+-- This file has been generated from package.yaml by hpack version 0.28.2. -- -- see: https://github.com/sol/hpack ----- hash: c74053b52f17daacb3df2612ca8127e9e24189040509592d94be61a097cfaf0e+-- hash: 0fe3779c7b8153d9edb2c1061a1ea1ce629a6e6474a0c02fd622a91b5c237a31  name:           conduit-algorithms-version:        0.0.8.1+version:        0.0.8.2 synopsis:       Conduit-based algorithms description:    Algorithms on Conduits, including higher level asynchronous processing and some other utilities. category:       Conduit@@ -17,7 +17,6 @@ license-file:   COPYING build-type:     Simple cabal-version:  >= 1.10- extra-source-files:     ChangeLog     README.md
tests/Tests.hs view
@@ -22,6 +22,7 @@ import           Control.Exception (catch, ErrorCall) import           Control.Monad (forM_) import           Control.Monad.Trans.Resource.Internal (ResourceT)+import qualified Control.Monad.Trans.Resource as R import           Control.Monad.IO.Unlift (MonadUnliftIO)  import qualified Data.Conduit.Algorithms as CAlg@@ -202,6 +203,19 @@             .| CL.map (read . B8.unpack)     removeFile testingFileNameGZ     removeFile testingFileNameGZ2++case_withPossiblyCompressedFile :: IO ()+case_withPossiblyCompressedFile = do+    let testdata = [0 :: Int .. 12]+    C.runConduitRes $+        CC.yieldMany testdata+            .| CL.map (B8.pack . (\n -> show n ++ "\n"))+            .| CAlg.asyncGzipToFile testingFileNameGZ+    back <- R.runResourceT $+                CAlg.withPossiblyCompressedFile testingFileNameGZ $ \src ->+                    C.runConduit (src .| CB.lines .| CL.map (read . B8.unpack) .| CC.sinkList)+    removeFile testingFileNameGZ+    back @?= testdata  case_async_bzip2_to_from :: IO () case_async_bzip2_to_from = do