zeromq-haskell 0.5.0 → 0.6.0
raw patch · 13 files changed
+129/−146 lines, 13 filesPVP ok
version bump matches the API change (PVP)
API changes (from Hackage documentation)
- System.ZMQ: with :: Size -> (Context -> IO a) -> IO a
+ System.ZMQ: withContext :: Size -> (Context -> IO a) -> IO a
+ System.ZMQ: withSocket :: SType a => Context -> a -> (Socket a -> IO b) -> IO b
Files
- AUTHORS +1/−0
- README.md +5/−3
- src/System/ZMQ.hs +30/−23
- src/System/ZMQ/Internal.hs +8/−8
- test/display.hs +11/−11
- test/perf/local_lat.hs +5/−6
- test/perf/local_thr.hs +10/−11
- test/perf/remote_lat.hs +8/−9
- test/perf/remote_thr.hs +5/−6
- test/poll.hs +7/−7
- test/prompt.hs +6/−6
- test/queue.hs +32/−55
- zeromq-haskell.cabal +1/−1
AUTHORS view
@@ -4,4 +4,5 @@ Ville Tirronen added support for ZMG_SNDMORE Jeremy Fitzhardinge integrated with GHC's I/O manager Bryan O'Sullivan added resource wrappers, socket finalizer+Ben Lever fixed and tweaked some of the test examples
README.md view
@@ -3,10 +3,12 @@ Current status -------------- -Version 0.5.0 - This software currently has *beta* status, i.e. it had-seen limited testing. Changes to its API may still happen.+Version 0.6.0 - This software currently has *beta* status, i.e. it had+seen limited testing. This version renames "with" to "withContext" and+introduces a new "withSocket" resource wrapper. The API is otherwise+identical to 0.5.0 -This software requires zeromq version 2.1.0.+This software requires zeromq version 2.1.x. Installation ------------
src/System/ZMQ.hs view
@@ -38,9 +38,8 @@ Up(..), Down(..), - with,- socket,- close,+ withContext,+ withSocket, setOption, getOption, System.ZMQ.subscribe,@@ -56,7 +55,9 @@ -- * Low-level functions init,- term+ term,+ socket,+ close, ) where @@ -64,7 +65,7 @@ import Control.Applicative import Control.Exception import Control.Monad (unless, when)-import Data.IORef (readIORef, writeIORef)+import Data.IORef (atomicModifyIORef) import Data.Int import Data.Maybe import System.ZMQ.Base@@ -308,30 +309,36 @@ -- | Run an action with a 0MQ context. The 'Context' supplied to your -- action will /not/ be valid after the action either returns or -- throws an exception.-with :: Size -> (Context -> IO a) -> IO a-with ioThreads act =+withContext :: Size -> (Context -> IO a) -> IO a+withContext ioThreads act = bracket (throwErrnoIfNull "c_zmq_init" $ c_zmq_init (fromIntegral ioThreads)) (throwErrnoIfMinus1_ "c_zmq_term" . c_zmq_term) (act . Context) --- | Create a new 0MQ socket within the given context.+-- | Run an action with a 0MQ socket. The socket will be closed after running+-- the supplied action even if an error occurs. The socket supplied to your+-- action will /not/ be valid after the action terminates.+withSocket :: SType a => Context -> a -> (Socket a -> IO b) -> IO b+withSocket c t = bracket (socket c t) close++-- | Create a new 0MQ socket within the given context. 'withSocket' provides+-- automatic socket closing and may be safer to use. socket :: SType a => Context -> a -> IO (Socket a) socket (Context c) t = do let zt = typeVal . zmqSocketType $ t s <- throwErrnoIfNull "socket" (c_zmq_socket c zt) sock@(Socket _ status) <- mkSocket s addFinalizer sock $ do- alive <- readIORef status- unless alive $ c_zmq_close s >> return ()+ alive <- atomicModifyIORef status (\b -> (False, b))+ when alive $ c_zmq_close s >> return () -- socket has not been closed yet return sock --- | Close a 0MQ socket.+-- | Close a 0MQ socket. 'withSocket' provides automatic socket closing and may+-- be safer to use. close :: Socket a -> IO ()-close sock@(Socket _ status) = withSocket "close" sock $ \s -> do- alive <- readIORef status- when alive $ do- writeIORef status False- throwErrnoIfMinus1_ "close" . c_zmq_close $ s+close sock@(Socket _ status) = onSocket "close" sock $ \s -> do+ alive <- atomicModifyIORef status (\b -> (False, b))+ when alive $ throwErrnoIfMinus1_ "close" . c_zmq_close $ s -- | Set the given option on the socket. Please note that there are -- certain combatibility constraints w.r.t the socket type (cf. man@@ -393,18 +400,18 @@ -- | Bind the socket to the given address (zmq_bind) bind :: Socket a -> String -> IO ()-bind sock str = withSocket "bind" sock $+bind sock str = onSocket "bind" sock $ throwErrnoIfMinus1_ "bind" . withCString str . c_zmq_bind -- | Connect the socket to the given address (zmq_connect). connect :: Socket a -> String -> IO ()-connect sock str = withSocket "connect" sock $+connect sock str = onSocket "connect" sock $ throwErrnoIfMinus1_ "connect" . withCString str . c_zmq_connect -- | Send the given 'SB.ByteString' over the socket (zmq_send). send :: Socket a -> SB.ByteString -> [Flag] -> IO () send sock val fls = bracket (messageOf val) messageClose $ \m ->- withSocket "send" sock $ \s ->+ onSocket "send" sock $ \s -> retry "send" (waitWrite sock) $ c_zmq_send s (msgPtr m) (combine (NoBlock : fls)) @@ -413,14 +420,14 @@ -- (Lazy.toChunks lbs)) flags@ but may be more efficient. send' :: Socket a -> LB.ByteString -> [Flag] -> IO () send' sock val fls = bracket (messageOfLazy val) messageClose $ \m ->- withSocket "send'" sock $ \s ->+ onSocket "send'" sock $ \s -> retry "send'" (waitWrite sock) $ c_zmq_send s (msgPtr m) (combine (NoBlock : fls)) -- | Receive a 'ByteString' from socket (zmq_recv). receive :: Socket a -> [Flag] -> IO (SB.ByteString) receive sock fls = bracket messageInit messageClose $ \m ->- withSocket "receive" sock $ \s -> do+ onSocket "receive" sock $ \s -> do retry "receive" (waitRead sock) $ c_zmq_recv_unsafe s (msgPtr m) (combine (NoBlock : fls)) data_ptr <- c_zmq_msg_data (msgPtr m)@@ -495,8 +502,8 @@ -- Please note that this call never returns. device :: Device -> Socket a -> Socket b -> IO () device device' insock outsock =- withSocket "device" insock $ \insocket ->- withSocket "device" outsock $ \outsocket ->+ onSocket "device" insock $ \insocket ->+ onSocket "device" outsock $ \outsocket -> throwErrnoIfMinus1Retry_ "device" $ c_zmq_device (fromDevice device') insocket outsocket where
src/System/ZMQ/Internal.hs view
@@ -19,7 +19,7 @@ , toZMQFlag , combine , mkSocket- , withSocket+ , onSocket ) where import Control.Applicative@@ -65,9 +65,9 @@ -- internal helpers: -withSocket :: String -> Socket a -> (ZMQSocket -> IO b) -> IO b-withSocket _func (Socket sock _state) act = act sock-{-# INLINE withSocket #-}+onSocket :: String -> Socket a -> (ZMQSocket -> IO b) -> IO b+onSocket _func (Socket sock _state) act = act sock+{-# INLINE onSocket #-} mkSocket :: ZMQSocket -> IO (Socket a) mkSocket s = Socket s <$> newIORef True@@ -110,14 +110,14 @@ return (Message ptr) setIntOpt :: (Storable b, Integral b) => Socket a -> ZMQOption -> b -> IO ()-setIntOpt sock (ZMQOption o) i = withSocket "setIntOpt" sock $ \s ->+setIntOpt sock (ZMQOption o) i = onSocket "setIntOpt" sock $ \s -> throwErrnoIfMinus1_ "setIntOpt" $ with i $ \ptr -> c_zmq_setsockopt s (fromIntegral o) (castPtr ptr) (fromIntegral . sizeOf $ i) setStrOpt :: Socket a -> ZMQOption -> String -> IO ()-setStrOpt sock (ZMQOption o) str = withSocket "setStrOpt" sock $ \s ->+setStrOpt sock (ZMQOption o) str = onSocket "setStrOpt" sock $ \s -> throwErrnoIfMinus1_ "setStrOpt" $ withCStringLen str $ \(cstr, len) -> c_zmq_setsockopt s (fromIntegral o) (castPtr cstr)@@ -127,7 +127,7 @@ getBoolOpt s o = ((1 :: Int64) ==) <$> getIntOpt s o getIntOpt :: (Storable b, Integral b) => Socket a -> ZMQOption -> IO b-getIntOpt sock (ZMQOption o) = withSocket "getIntOpt" sock $ \s -> do+getIntOpt sock (ZMQOption o) = onSocket "getIntOpt" sock $ \s -> do let i = 0 bracket (new i) free $ \iptr -> bracket (new (fromIntegral . sizeOf $ i :: CSize)) free $ \jptr -> do@@ -136,7 +136,7 @@ peek iptr getStrOpt :: Socket a -> ZMQOption -> IO String-getStrOpt sock (ZMQOption o) = withSocket "getStrOpt" sock $ \s ->+getStrOpt sock (ZMQOption o) = onSocket "getStrOpt" sock $ \s -> bracket (mallocBytes 255) free $ \bPtr -> bracket (new (255 :: CSize)) free $ \sPtr -> do throwErrnoIfMinus1_ "getStrOpt" $
test/display.hs view
@@ -2,22 +2,22 @@ import System.IO import System.Exit import System.Environment+import Control.Exception import qualified System.ZMQ as ZMQ import qualified Data.ByteString as SB main :: IO () main = do args <- getArgs- when (length args /= 1) $ do- hPutStrLn stderr "usage: display <address>"+ when (length args < 1) $ do+ hPutStrLn stderr "usage: display <address> [<address>, ...]" exitFailure- let addr = head args- ZMQ.with 1 $ \c -> do- s <- ZMQ.socket c ZMQ.Sub- ZMQ.subscribe s ""- ZMQ.connect s addr- forever $ do- line <- ZMQ.receive s []- SB.putStrLn line- hFlush stdout+ ZMQ.withContext 1 $ \c ->+ ZMQ.withSocket c ZMQ.Sub $ \s -> do+ ZMQ.subscribe s ""+ mapM (ZMQ.connect s) args+ forever $ do+ line <- ZMQ.receive s []+ SB.putStrLn line+ hFlush stdout
test/perf/local_lat.hs view
@@ -14,12 +14,11 @@ let bindTo = args !! 0 size = read $ args !! 1 rounds = read $ args !! 2- ZMQ.with 1 $ \c -> do- s <- ZMQ.socket c ZMQ.Rep- ZMQ.bind s bindTo- loop s rounds size- ZMQ.close s- where+ ZMQ.withContext 1 $ \c ->+ ZMQ.withSocket c ZMQ.Rep $ \s -> do+ ZMQ.bind s bindTo+ loop s rounds size+ where loop s r sz = unless (r <= 0) $ do msg <- ZMQ.receive s [] when (SB.length msg /= sz) $
test/perf/local_thr.hs view
@@ -16,17 +16,16 @@ let bindTo = args !! 0 size = read $ args !! 1 :: Int count = read $ args !! 2 :: Int- ZMQ.with 1 $ \c -> do- s <- ZMQ.socket c ZMQ.Sub- ZMQ.subscribe s ""- ZMQ.bind s bindTo- receive s size- start <- getCurrentTime- loop s (count - 1) size- end <- getCurrentTime- printStat start end size count- ZMQ.close s- where+ ZMQ.withContext 1 $ \c ->+ ZMQ.withSocket c ZMQ.Sub $ \s -> do+ ZMQ.subscribe s ""+ ZMQ.bind s bindTo+ receive s size+ start <- getCurrentTime+ loop s (count - 1) size+ end <- getCurrentTime+ printStat start end size count+ where receive s sz = do msg <- ZMQ.receive s [] when (SB.length msg /= sz) $
test/perf/remote_lat.hs view
@@ -16,15 +16,14 @@ size = read $ args !! 1 rounds = read $ args !! 2 message = SB.replicate size 0x65- ZMQ.with 1 $ \c -> do- s <- ZMQ.socket c ZMQ.Req- ZMQ.connect s connTo- start <- getCurrentTime- loop s rounds message- end <- getCurrentTime- print (diffUTCTime end start)- ZMQ.close s- where+ ZMQ.withContext 1 $ \c ->+ ZMQ.withSocket c ZMQ.Req $ \s -> do+ ZMQ.connect s connTo+ start <- getCurrentTime+ loop s rounds message+ end <- getCurrentTime+ print (diffUTCTime end start)+ where loop s r msg = unless (r <= 0) $ do ZMQ.send s msg [] msg' <- ZMQ.receive s []
test/perf/remote_thr.hs view
@@ -16,12 +16,11 @@ size = read $ args !! 1 count = read $ args !! 2 message = SB.replicate size 0x65- ZMQ.with 1 $ \c -> do- s <- ZMQ.socket c ZMQ.Pub- ZMQ.connect s connTo- replicateM_ count $ ZMQ.send s message []- threadDelay 10000000- ZMQ.close s+ ZMQ.withContext 1 $ \c ->+ ZMQ.withSocket c ZMQ.Pub $ \s -> do+ ZMQ.connect s connTo+ replicateM_ count $ ZMQ.send s message []+ threadDelay 10000000 usage :: String usage = "usage: remote_thr <connect-to> <message-size> <message-count>"
test/poll.hs view
@@ -12,13 +12,13 @@ when (length args /= 1) $ do hPutStrLn stderr usage exitFailure- ZMQ.with 1 $ \c -> do- s <- ZMQ.socket c ZMQ.Rep- let bindTo = head args- toPoll = [ZMQ.S s ZMQ.In]- ZMQ.bind s bindTo- forever $- ZMQ.poll toPoll 1000000 >>= receive+ ZMQ.withContext 1 $ \c ->+ ZMQ.withSocket c ZMQ.Rep $ \s -> do+ let bindTo = head args+ toPoll = [ZMQ.S s ZMQ.In]+ ZMQ.bind s bindTo+ forever $+ ZMQ.poll toPoll 1000000 >>= receive where receive [] = return () receive ((ZMQ.S s e):ss) = do
test/prompt.hs view
@@ -16,10 +16,10 @@ exitFailure let addr = args !! 0 name = SB.append (SB.fromString $ args !! 1) ": "- ZMQ.with 1 $ \c -> do- s <- ZMQ.socket c ZMQ.Pub- ZMQ.connect s addr- forever $ do- line <- SB.fromString <$> getLine- ZMQ.send s (SB.append name line) []+ ZMQ.withContext 1 $ \c ->+ ZMQ.withSocket c ZMQ.Pub $ \s -> do+ ZMQ.bind s addr+ forever $ do+ line <- SB.fromString <$> getLine+ ZMQ.send s (SB.append name line) []
test/queue.hs view
@@ -4,39 +4,31 @@ -- -- ghc --make -threaded queue.hs -import Control.Concurrent (forkOS, threadDelay)+import Control.Concurrent (forkIO, threadDelay) import Control.Concurrent.MVar (MVar, newEmptyMVar, putMVar, takeMVar)- import Control.Monad (forever, forM_, replicateM, replicateM_)- import qualified Data.ByteString.Char8 as SB- import qualified System.ZMQ as ZMQ main :: IO ()-main = do- context <- ZMQ.init 1-+main = ZMQ.withContext 1 $ \context -> do lock <- newEmptyMVar-- _ <- forkOS $ launchQueue context lock- _ <- takeMVar lock+ _ <- forkIO $ launchQueue context lock+ _ <- takeMVar lock forM_ [0..numWorkers] $ \i ->- forkOS $ launchWorker context i+ forkIO $ launchWorker context i locks <- replicateM numClients newEmptyMVar forM_ (zip [0..numClients] locks) $ \(i, lock') ->- forkOS $ launchClient context i lock'+ forkIO $ launchClient context i lock' -- Wait untill all clients signal completion forM_ locks takeMVar - -- We can't clean up the context since our queue device is still running,- -- and can't be killed...- -- ZMQ.term context+ -- our queue device is still running, and can't be killed... - where+ where numWorkers :: Int numWorkers = 5 @@ -56,48 +48,33 @@ delay = 1000000 launchQueue :: ZMQ.Context -> MVar () -> IO ()- launchQueue context lock = do- workers <- ZMQ.socket context ZMQ.Xreq- ZMQ.bind workers workersAddress-- clients <- ZMQ.socket context ZMQ.Xrep- ZMQ.bind clients clientsAddress-- putMVar lock ()-- ZMQ.device ZMQ.Queue clients workers-- -- This isn't reached- ZMQ.close workers- ZMQ.close clients+ launchQueue context lock =+ ZMQ.withSocket context ZMQ.Xreq $ \workers ->+ ZMQ.withSocket context ZMQ.Xrep $ \clients -> do+ ZMQ.bind workers workersAddress+ ZMQ.bind clients clientsAddress+ putMVar lock ()+ ZMQ.device ZMQ.Queue clients workers launchWorker :: ZMQ.Context -> Int -> IO ()- launchWorker context i = do- socket <- ZMQ.socket context ZMQ.Rep- ZMQ.connect socket workersAddress-- forever $ do- request <- ZMQ.receive socket []- putStrLn $- "Message received in worker " ++ show i ++ ": " ++- SB.unpack request- threadDelay delay -- Do some 'work'- ZMQ.send socket message []- putStrLn $ "Reply sent in worker " ++ show i-- -- This isn't reached- ZMQ.close socket+ launchWorker context i =+ ZMQ.withSocket context ZMQ.Rep $ \socket -> do+ ZMQ.connect socket workersAddress+ forever $ do+ request <- ZMQ.receive socket []+ putStrLn $+ "Message received in worker " ++ show i ++ ": " +++ SB.unpack request+ threadDelay delay -- Do some 'work'+ ZMQ.send socket message []+ putStrLn $ "Reply sent in worker " ++ show i launchClient :: ZMQ.Context -> Int -> MVar () -> IO () launchClient context i lock = do- socket <- ZMQ.socket context ZMQ.Req- ZMQ.connect socket clientsAddress-- putStrLn $ "Sending message in client " ++ show i- ZMQ.send socket (SB.pack $ show i) []- _ <- ZMQ.receive socket []- putStrLn $ "Reply received in client " ++ show i-- ZMQ.close socket-+ ZMQ.withSocket context ZMQ.Req $ \socket -> do+ ZMQ.connect socket clientsAddress+ putStrLn $ "Sending message in client " ++ show i+ ZMQ.send socket (SB.pack $ show i) []+ _ <- ZMQ.receive socket []+ putStrLn $ "Reply received in client " ++ show i putMVar lock ()
zeromq-haskell.cabal view
@@ -1,5 +1,5 @@ name: zeromq-haskell-version: 0.5.0+version: 0.6.0 synopsis: bindings to zeromq description: Bindings to zeromq (http://zeromq.org) category: System, FFI