conduit-extra 1.1.11 → 1.1.12
raw patch · 4 files changed
+115/−11 lines, 4 filesdep ~asyncdep ~basePVP: major bump suggested
API removals or changes: PVP suggests a major version bump
Dependency ranges changed: async, base
API changes (from Hackage documentation)
+ Data.Conduit.Process: sourceCmdWithStreams :: String -> Producer IO ByteString -> Consumer ByteString IO a -> Consumer ByteString IO b -> IO (ExitCode, a, b)
+ Data.Conduit.Process: sourceProcessWithStreams :: CreateProcess -> Producer IO ByteString -> Consumer ByteString IO a -> Consumer ByteString IO b -> IO (ExitCode, a, b)
- Data.Conduit.Process: sourceCmdWithConsumer :: MonadIO m => String -> Consumer ByteString m a -> m (ExitCode, a)
+ Data.Conduit.Process: sourceCmdWithConsumer :: (MonadMask m, MonadIO m) => String -> Consumer ByteString m a -> m (ExitCode, a)
- Data.Conduit.Process: sourceProcessWithConsumer :: MonadIO m => CreateProcess -> Consumer ByteString m a -> m (ExitCode, a)
+ Data.Conduit.Process: sourceProcessWithConsumer :: (MonadMask m, MonadIO m) => CreateProcess -> Consumer ByteString m a -> m (ExitCode, a)
Files
- ChangeLog.md +4/−0
- Data/Conduit/Process.hs +76/−9
- conduit-extra.cabal +2/−1
- test/Data/Conduit/ProcessSpec.hs +33/−1
ChangeLog.md view
@@ -1,3 +1,7 @@+## 1.1.12++* Add sourceProcessWithStreams [#258](https://github.com/snoyberg/conduit/pull/258)+ ## 1.1.11 * `withCheckedProcessCleanup`
Data/Conduit/Process.hs view
@@ -14,6 +14,8 @@ ( -- * Functions sourceCmdWithConsumer , sourceProcessWithConsumer+ , sourceCmdWithStreams+ , sourceProcessWithStreams , withCheckedProcessCleanup -- * Reexport , module Data.Streaming.Process@@ -27,7 +29,8 @@ import Data.Conduit import Data.Conduit.Binary (sourceHandle, sinkHandle) import Data.ByteString (ByteString)-import Control.Monad.Catch (MonadMask, onException, throwM)+import Control.Concurrent.Async (runConcurrently, Concurrently(..))+import Control.Monad.Catch (MonadMask, onException, throwM, finally) instance (r ~ (), MonadIO m, i ~ ByteString) => InputSource (ConduitM i o m r) where isStdStream = (\(Just h) -> return $ sinkHandle h, Just CreatePipe)@@ -44,24 +47,85 @@ -- return a tuple of the @ExitCode@ from the process and the output collected -- from the @Consumer@. --+-- If an exception is raised by the consumer,+-- the process is terminated.+-- -- Since 1.1.2-sourceProcessWithConsumer :: MonadIO m => CreateProcess -> Consumer ByteString m a -> m (ExitCode, a)+sourceProcessWithConsumer :: (MonadMask m, MonadIO m)+ => CreateProcess+ -> Consumer ByteString m a -- ^stdout+ -> m (ExitCode, a) sourceProcessWithConsumer cp consumer = do- (ClosedStream, (source, close), ClosedStream, cph) <- streamingProcess cp- res <- source $$ consumer- close- ec <- waitForStreamingProcess cph+ (ClosedStream, (source, close), ClosedStream, sph) <- streamingProcess cp+ res <- (source $$ consumer)+ `finally` close+ `onException` liftIO (terminateStreamingProcess sph)+ ec <- waitForStreamingProcess sph return (ec, res) -- | Like @sourceProcessWithConsumer@ but providing the command to be run as -- a @String@. -- -- Since 1.1.2-sourceCmdWithConsumer :: MonadIO m => String -> Consumer ByteString m a -> m (ExitCode, a)+sourceCmdWithConsumer :: (MonadMask m, MonadIO m)+ => String -- ^command+ -> Consumer ByteString m a -- ^stdout+ -> m (ExitCode, a) sourceCmdWithConsumer cmd = sourceProcessWithConsumer (shell cmd) ++-- | Given a @CreateProcess@, run the process+-- and feed the provided @Producer@+-- to the stdin @Sink@ of the process.+-- Use the process outputs (stdout, stderr) as @Source@s+-- and feed it to the provided @Consumer@s.+-- Once the process has completed,+-- return a tuple of the @ExitCode@ from the process+-- and the results collected from the @Consumer@s.+--+-- If an exception is raised by any of the streams,+-- the process is terminated.+--+-- IO is required because the streams are run concurrently+-- using the <https://hackage.haskell.org/package/async async> package+--+-- @since 1.1.12+sourceProcessWithStreams :: CreateProcess+ -> Producer IO ByteString -- ^stdin+ -> Consumer ByteString IO a -- ^stdout+ -> Consumer ByteString IO b -- ^stderr+ -> IO (ExitCode, a, b)+sourceProcessWithStreams cp producerStdin consumerStdout consumerStderr = do+ ( (sinkStdin, closeStdin)+ , (sourceStdout, closeStdout)+ , (sourceStderr, closeStderr)+ , sph) <- streamingProcess cp+ (_, resStdout, resStderr) <-+ runConcurrently (+ (,,)+ <$> Concurrently ((producerStdin $$ sinkStdin) `finally` closeStdin)+ <*> Concurrently (sourceStdout $$ consumerStdout)+ <*> Concurrently (sourceStderr $$ consumerStderr))+ `finally` (closeStdout >> closeStderr)+ `onException` terminateStreamingProcess sph+ ec <- waitForStreamingProcess sph+ return (ec, resStdout, resStderr)++-- | Like @sourceProcessWithStreams@ but providing the command to be run as+-- a @String@.+--+-- @since 1.1.12+sourceCmdWithStreams :: String -- ^command+ -> Producer IO ByteString -- ^stdin+ -> Consumer ByteString IO a -- ^stdout+ -> Consumer ByteString IO b -- ^stderr+ -> IO (ExitCode, a, b)+sourceCmdWithStreams cmd = sourceProcessWithStreams (shell cmd)+ -- | Same as 'withCheckedProcess', but kills the child process in the case of -- an exception being thrown by the provided callback function.+--+-- @since 1.1.11 withCheckedProcessCleanup :: ( InputSource stdin , OutputSink stderr@@ -74,9 +138,12 @@ -> m b withCheckedProcessCleanup cp f = do (x, y, z, sph) <- streamingProcess cp- res <- f x y z `onException`- liftIO (terminateProcess (streamingProcessHandleRaw sph))+ res <- f x y z `onException` liftIO (terminateStreamingProcess sph) ec <- waitForStreamingProcess sph if ec == ExitSuccess then return res else throwM $ ProcessExitedUnsuccessfully cp ec+++terminateStreamingProcess :: StreamingProcessHandle -> IO ()+terminateStreamingProcess = terminateProcess . streamingProcessHandleRaw
conduit-extra.cabal view
@@ -1,5 +1,5 @@ Name: conduit-extra-Version: 1.1.11+Version: 1.1.12 Synopsis: Batteries included conduit: adapters for common libraries. Description: The conduit package itself maintains relative small dependencies. The purpose of this package is to collect commonly used utility functions wrapping other library dependencies, without depending on heavier-weight dependencies. The basic idea is that this package should only depend on haskell-platform packages and conduit.@@ -45,6 +45,7 @@ , transformers , transformers-base + , async , attoparsec >= 0.10 , blaze-builder >= 0.3 , directory
test/Data/Conduit/ProcessSpec.hs view
@@ -1,4 +1,5 @@ {-# LANGUAGE CPP #-}+{-# LANGUAGE OverloadedStrings #-} module Data.Conduit.ProcessSpec (spec, main) where import Test.Hspec@@ -9,6 +10,7 @@ import Control.Concurrent.Async (concurrently) import qualified Data.ByteString.Lazy as L import qualified Data.ByteString as S+import qualified Data.ByteString.Char8 as S8 import System.Exit import Control.Concurrent (threadDelay) @@ -45,6 +47,36 @@ `shouldReturn` (ExitFailure 11, ()) (sourceCmdWithConsumer "exit 12" CL.sinkNull) `shouldReturn` (ExitFailure 12, ())+ (sourceCmdWithStreams "exit 0" CL.sourceNull CL.sinkNull CL.sinkNull)+ `shouldReturn` (ExitSuccess, (), ())+ (sourceCmdWithStreams "exit 11" CL.sourceNull CL.sinkNull CL.sinkNull)+ `shouldReturn` (ExitFailure 11, (), ())+ (sourceCmdWithStreams "exit 12" CL.sourceNull CL.sinkNull CL.sinkNull)+ `shouldReturn` (ExitFailure 12, (), ())++ it "consumes stdout" $ do+ let mystr = "this is a test string" :: String+ sourceCmdWithStreams ("bash -c \"echo -n " ++ mystr ++ "\"")+ CL.sourceNull+ CL.consume -- stdout+ CL.consume -- stderr+ `shouldReturn` (ExitSuccess, [S8.pack mystr], [])++ it "consumes stderr" $ do+ let mystr = "this is a test string" :: String+ sourceCmdWithStreams ("bash -c \">&2 echo -n " ++ mystr ++ "\"")+ CL.sourceNull+ CL.consume -- stdout+ CL.consume -- stderr+ `shouldReturn` (ExitSuccess, [], [S8.pack mystr])++ it "feeds stdin" $ do+ let mystr = "this is a test string" :: S.ByteString+ sourceCmdWithStreams "cat"+ (mapM_ yield . L.toChunks $ L.fromStrict mystr)+ CL.consume -- stdout+ CL.consume -- stderr+ `shouldReturn` (ExitSuccess, [mystr], []) #endif it "blocking vs non-blocking" $ do (ClosedStream, ClosedStream, ClosedStream, cph) <- streamingProcess (shell "sleep 1")@@ -64,7 +96,7 @@ threadDelay 500000 loop (pred i) Just _ -> mec2 `shouldBe` Just ExitSuccess- loop 5+ loop (5 :: Int) ec <- waitForStreamingProcess cph ec `shouldBe` ExitSuccess