nanomsg-haskell (empty) → 0.1.0
raw patch · 9 files changed
+1228/−0 lines, 9 filesdep +QuickCheckdep +basedep +bytestringsetup-changed
Dependencies added: QuickCheck, base, bytestring, criterion, nanomsg-haskell, test-framework, test-framework-quickcheck2, test-framework-th, zeromq3-haskell
Files
- AUTHORS +3/−0
- LICENSE +20/−0
- README.md +97/−0
- Setup.hs +2/−0
- benchmarks/SendMessages.hs +25/−0
- benchmarks/Zmq.hs +44/−0
- nanomsg-haskell.cabal +77/−0
- src/Nanomsg.hsc +754/−0
- tests/Properties.hs +206/−0
+ AUTHORS view
@@ -0,0 +1,3 @@+Full list of copyright-holders:+Ivar Nymoen+
+ LICENSE view
@@ -0,0 +1,20 @@++Copyright (c) 2013 the nanomsg-haskell authors++Permission is hereby granted, free of charge, to any person obtaining a copy+of this software and associated documentation files (the "Software"), to deal+in the Software without restriction, including without limitation the rights+to use, copy, modify, merge, publish, distribute, sublicense, and/or sell+copies of the Software, and to permit persons to whom the Software is+furnished to do so, subject to the following conditions:++The above copyright notice and this permission notice shall be included in+all copies or substantial portions of the Software.++THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR+IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,+FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE+AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER+LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,+OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN+THE SOFTWARE.
+ README.md view
@@ -0,0 +1,97 @@+# nanomsg-haskell++This is a Haskell binding for the nanomsg library: <http://nanomsg.org/>.++There's support for [blocking](http://hackage.haskell.org/packages/archive/base/latest/doc/html/Control-Concurrent.html#v:threadWaitRead) send and recv, a non-blocking receive,+and for all the socket types and the functions you need to wire them up and+tear them down again.++Most socket options are available through accessor and mutator+functions. Sockets are typed, transports are not.+++## Building++You would normally make sure the nanomsg library is on your system and then+install from Hackage using "cabal update && cabal install nanomsg-haskell",+but can build from the repository following these steps:++ 1. Build and install nanomsg (and zmq3, if you are building benchmarks)+ 1. git clone https://github.com/ivarnymoen/nanomsg-haskell+ 1. cd nanomsg-haskell && cabal sandbox init+ 1. cabal install --dependencies-only [--enable-tests] [--enable-benchmarks]+ 1. cabal configure [--enable-tests] [--enable-benchmarks]+ 1. cabal build+ 1. [cabal test]+++## Contributing++Just submit a pull request for small stuff, but please get in touch+beforehand before sinking a lot of effort into major/API changes.++Remember adding your name to the AUTHORS file.+++## Usage++Simple Pub/sub example:++Server:++ module Main where++ import Nanomsg+ import qualified Data.ByteString.Char8 as C+ import Control.Monad (mapM_)+ import Control.Concurrent (threadDelay)++ main :: IO ()+ main = do+ withSocket Pub $ \s -> do+ _ <- bind s "tcp://*:5560"+ mapM_ (\num -> sendNumber s num) (cycle [1..1000000])+ where+ sendNumber s number = do+ threadDelay 1000 -- let's conserve some cycles+ let numAsString = show number+ send s (C.pack numAsString)++Client:++ module Main where++ import Nanomsg+ import qualified Data.ByteString.Char8 as C+ import Control.Monad (forever)++ main :: IO ()+ main = do+ withSocket Sub $ \s -> do+ _ <- connect s "tcp://localhost:5560"+ subscribe s $ C.pack ""+ forever $ do+ msg <- recv s+ C.putStrLn msg++Nonblocking client:++ module Main where++ import Nanomsg+ import qualified Data.ByteString.Char8 as C+ import Control.Monad (forever)+ import Control.Concurrent (threadDelay)++ main :: IO ()+ main =+ withSocket Sub $ \s -> do+ _ <- connect s "tcp://localhost:5560"+ subscribe s $ C.pack ""+ forever $ do+ threadDelay 700 -- let's conserve some cycles+ msg <- recv' s+ C.putStrLn $ case msg of+ Nothing -> C.pack "No message"+ Just s -> s+
+ Setup.hs view
@@ -0,0 +1,2 @@+import Distribution.Simple+main = defaultMain
+ benchmarks/SendMessages.hs view
@@ -0,0 +1,25 @@+module Main where++import Nanomsg+import Criterion.Main+import qualified Data.ByteString.Char8 as C+import Control.Monad (replicateM_)++pair :: Int -> Int -> IO ()+pair size count = do+ sender <- socket Pair+ _ <- bind sender "inproc://pairtest"+ recipient <- socket Pair+ _ <- connect recipient "inproc://pairtest"+ let msg = C.pack $ replicate size 'a'+ replicateM_ count (send sender msg >> recv recipient)+ close sender+ close recipient+ return ()++main :: IO ()+main = defaultMain+ [ bench "40 bytes x 10k messages" $ nfIO $ pair 40 10000+ , bench "20k bytes x 20 messages" $ nfIO $ pair 20000 20+ ]+
+ benchmarks/Zmq.hs view
@@ -0,0 +1,44 @@+module Main where++import qualified Nanomsg as N+import qualified System.ZMQ3.Monadic as Z+import Criterion.Main+import qualified Data.ByteString.Char8 as C+import Control.Monad (replicateM_)++nPair :: Int -> Int -> String -> String -> IO ()+nPair size count bindString connString = do+ s1 <- N.socket N.Pair+ _ <- N.bind s1 bindString+ s2 <- N.socket N.Pair+ _ <- N.connect s2 connString+ let msg = C.pack $ replicate size 'a'+ replicateM_ count (N.send s1 msg >> N.recv s2 >>= N.send s2 >> N.recv s1)+ N.close s1+ N.close s2+ return ()++zPair :: Int -> Int -> String -> String -> IO ()+zPair size count bindString connString = Z.runZMQ $ do+ s1 <- Z.socket Z.Pair+ _ <- Z.bind s1 bindString+ s2 <- Z.socket Z.Pair+ _ <- Z.connect s2 connString+ let msg = C.pack $ replicate size 'a'+ replicateM_ count (Z.send s1 [] msg >> Z.receive s2 >>= Z.send s2 [] >> Z.receive s1)+ Z.close s1+ Z.close s2+ return ()++main :: IO ()+main = defaultMain+ [ bench "nanomsg-haskell: 40 bytes x 1k messages, roundtrip, tcp" $ nfIO $ nPair 40 1000 "tcp://*:5566" "tcp://localhost:5566"+ , bench "zeromq3-haskell: 40 bytes x 1k messages, roundtrip, tcp" $ nfIO $ zPair 40 1000 "tcp://*:5566" "tcp://localhost:5566"+ , bench "nanomsg-haskell: 20k bytes x 20 messages, roundtrip, tcp" $ nfIO $ nPair 20000 20 "tcp://*:5566" "tcp://localhost:5566"+ , bench "zeromq3-haskell: 20k bytes x 20 messages, roundtrip, tcp" $ nfIO $ zPair 20000 20 "tcp://*:5566" "tcp://localhost:5566"+ , bench "nanomsg-haskell: 40 bytes x 1k messages, roundtrip, inproc" $ nfIO $ nPair 40 1000 "inproc://bench" "inproc://bench"+ , bench "zeromq3-haskell: 40 bytes x 1k messages, roundtrip, inproc" $ nfIO $ zPair 40 1000 "inproc://bench" "inproc://bench"+ , bench "nanomsg-haskell: 20k bytes x 20 messages, roundtrip, inproc" $ nfIO $ nPair 20000 20 "inproc://bench" "inproc://bench"+ , bench "zeromq3-haskell: 20k bytes x 20 messages, roundtrip, inproc" $ nfIO $ zPair 20000 20 "inproc://bench" "inproc://bench"+ ]+
+ nanomsg-haskell.cabal view
@@ -0,0 +1,77 @@+name: nanomsg-haskell+version: 0.1.0+synopsis:+ Bindings to the nanomsg library+description:+ This is a Haskell binding for the nanomsg library: <http://nanomsg.org/>.++ There's support for blocking send and recv, a non-blocking receive,+ and for all the socket types and the functions you need to wire+ them up and tear them down again.++ Most sockets options are available through accessor and mutator+ functions. Sockets are typed, transports are not.++homepage: https://github.com/ivarnymoen/nanomsg-haskell+license: MIT+license-file: LICENSE+author: See AUTHORS file+maintainer: <ivar.nymoen@gmail.com>+copyright: Copyright (c) 2013 the nanomsg-haskell authors+category: Network+build-type: Simple+cabal-version: >=1.8+extra-source-files: AUTHORS, README.md++library+ hs-source-dirs: src+ ghc-options: -O2 -Wall+ exposed-modules: Nanomsg+ extensions: ForeignFunctionInterface, DeriveDataTypeable+ includes: nanomsg/nn.h+ extra-libraries: nanomsg+ build-depends:+ base == 4.*,+ bytestring >= 0.9.0 && < 0.11++test-suite tests+ type: exitcode-stdio-1.0+ hs-source-dirs: tests+ main-is: Properties.hs+ ghc-options: -O2 -Wall -threaded+ build-depends:+ base == 4.*,+ bytestring >= 0.9.0 && < 0.11,+ nanomsg-haskell,+ QuickCheck,+ test-framework,+ test-framework-quickcheck2,+ test-framework-th++source-repository head+ type: git+ location: https://github.com/ivarnymoen/nanomsg-haskell++benchmark send-messages+ type: exitcode-stdio-1.0+ main-is: SendMessages.hs+ ghc-options: -O2 -Wall -threaded+ hs-source-dirs: benchmarks+ build-depends:+ base == 4.*,+ bytestring >= 0.9.0 && < 0.11,+ nanomsg-haskell,+ criterion++benchmark vs-zeromq-bindings+ type: exitcode-stdio-1.0+ main-is: Zmq.hs+ ghc-options: -O2 -Wall -threaded+ hs-source-dirs: benchmarks+ build-depends:+ base == 4.*,+ bytestring >= 0.9.0 && < 0.11,+ nanomsg-haskell,+ zeromq3-haskell,+ criterion+
+ src/Nanomsg.hsc view
@@ -0,0 +1,754 @@+{-# LANGUAGE ForeignFunctionInterface, DeriveDataTypeable #-}+-- |+-- Module: Nanomsg+-- Copyright: (c) 2013 Ivar Nymoen+-- License: MIT+-- Stability: experimental+--+-- This is a Haskell binding for the nanomsg library: <http://nanomsg.org/>.+--+-- There's support for blocking send and recv, a non-blocking receive,+-- and for all the socket types and the functions you need to wire+-- them up and tear them down again.+--+-- Most socket options are available through accessor and mutator+-- functions. Sockets are typed, transports are not.+--+-- Socket type documentation is adapted or quoted verbatim from the+-- nanomsg manual. Please refer to nanomsg.org for information on+-- how to use the library.+module Nanomsg+ (+ -- * Socket types+ Pair(..)+ , Req(..)+ , Rep(..)+ , Pub(..)+ , Sub(..)+ , Surveyor(..)+ , Respondent(..)+ , Push(..)+ , Pull(..)+ , Bus(..)+ -- * Other types+ , Socket+ , Endpoint+ , NNException+ -- * Functions+ , socket+ , withSocket+ , bind+ , connect+ , send+ , recv+ , recv'+ , subscribe+ , unsubscribe+ , shutdown+ , close+ , term+ -- ** Socket option accessors and mutators+ , linger+ , setLinger+ , sndBuf+ , setSndBuf+ , rcvBuf+ , setRcvBuf+ , reconnectInterval+ , setReconnectInterval+ , reconnectIntervalMax+ , setReconnectIntervalMax+ , sndPrio+ , setSndPrio+ , ipv4Only+ , setIpv4Only+ , requestResendInterval+ , setRequestResendInterval+ , tcpNoDelay+ , setTcpNoDelay+ ) where++#include "nanomsg/nn.h"+#include "nanomsg/pair.h"+#include "nanomsg/reqrep.h"+#include "nanomsg/pubsub.h"+#include "nanomsg/survey.h"+#include "nanomsg/pipeline.h"+#include "nanomsg/bus.h"+#include "nanomsg/tcp.h"++import Data.ByteString (ByteString)+-- import qualified Data.ByteString.Lazy as L+import qualified Data.ByteString.Char8 as C+import qualified Data.ByteString.Unsafe as U+import Foreign (peek, poke, alloca)+import Foreign.Ptr+import Foreign.C.Types+import Foreign.C.String+import Foreign.Storable (sizeOf)+import Control.Applicative ( (<$>) )+import Control.Exception.Base (bracket)+import Control.Exception (Exception, throwIO)+import Data.Typeable (Typeable)+import Control.Monad (void)+import Text.Printf (printf)+import Control.Concurrent (threadWaitRead, threadWaitWrite)+import System.Posix.Types (Fd(..))+++-- * Data and typedefs++-- | Socket for communication with exactly one peer. Each+-- party can send messages at any time. If the peer is not+-- available or send buffer is full, subsequent calls to+-- 'send' will block until it’s possible to send the message.+data Pair = Pair++-- | Used to implement a client application that sends requests+-- and receives replies. The socket will resend requests automatically+-- if there's no reply within a given time. The default timeout+-- is 1 minute.+--+-- See also 'setRequestResendInterval'.+data Req = Req++-- | Used to implement a stateless worker that receives requests+-- and sends replies.+data Rep = Rep++-- | This socket is used to distribute messages to multiple destinations.+-- Can not receive.+data Pub = Pub++-- | Receives messages from the publisher. Only messages that the socket is+-- subscribed to are received. When the socket is created there are no+-- subscriptions and thus no messages will be received.+--+-- Send is not defined on this socket. The socket can be connected+-- to at most one peer.+--+-- See also 'subscribe' and 'unsubscribe'.+data Sub = Sub++-- | Surveyor and respondent are used to broadcast a survey to multiple+-- locations and gather the responses.+--+-- This socket is used to send the survey. The survey is delivered to all+-- the connected respondents. Once the query is sent, the socket can be used+-- to receive the responses.+--+-- When the survey deadline expires, receive will return ETIMEDOUT error.+--+-- See also 'setSurveyorDeadline'+data Surveyor = Surveyor++-- | Used to respond to a survey. Survey is received using receive function,+-- response is sent using send function. This socket can be connected to+-- at most one peer.+data Respondent = Respondent++-- | Push and Pull sockets fair queue messages from one processing step, load+-- balancing them among instances of the next processing step.+--+-- This socket is used to send messages to a cluster of load-balanced nodes.+--+-- Receive operation is not implemented on this socket type.+data Push = Push++-- | This socket is used to receive a message from a cluster of nodes.+--+-- Send operation is not implemented on this socket type.+data Pull = Pull++-- | Broadcasts messages from any node to all other nodes in the topology.+-- The socket should never receives messages that it sent itself.+--+-- This pattern scales only to local level (within a single machine or+-- within a single LAN). Trying to scale it further can result in overloading+-- individual nodes with messages.+data Bus = Bus++-- | Endpoint identifier. Created by 'connect' or 'bind'.+--+-- Close connections using 'shutdown'.+data Endpoint = Endpoint CInt+ deriving (Eq, Show)++-- | Sockets are created by 'socket' and connections are established with 'connect' or 'bind'.+--+-- Free sockets using 'close'.+data Socket a = Socket a CInt+ deriving (Eq, Show)++-- | Typeclass used by all sockets, to extract their C type.+class Protocol a where+ -- | Returns the C enum value for each type. E.g. Pair => #const NN_PAIR+ protocolId :: a -> CInt++instance Protocol Pair where+ protocolId Pair = #const NN_PAIR++instance Protocol Req where+ protocolId Req = #const NN_REQ++instance Protocol Rep where+ protocolId Rep = #const NN_REP++instance Protocol Pub where+ protocolId Pub = #const NN_PUB++instance Protocol Sub where+ protocolId Sub = #const NN_SUB++instance Protocol Surveyor where+ protocolId Surveyor = #const NN_SURVEYOR++instance Protocol Respondent where+ protocolId Respondent = #const NN_RESPONDENT++instance Protocol Push where+ protocolId Push = #const NN_PUSH++instance Protocol Pull where+ protocolId Pull = #const NN_PULL++instance Protocol Bus where+ protocolId Bus = #const NN_BUS+++-- | Typeclass restricting which sockets can use the send function.+class SendType a+instance SendType Pair+instance SendType Req+instance SendType Rep+instance SendType Pub+instance SendType Surveyor+instance SendType Respondent+instance SendType Push+instance SendType Bus++-- | Typeclass for sockets that implement recv+class RecvType a+instance RecvType Pair+instance RecvType Req+instance RecvType Rep+instance RecvType Sub+instance RecvType Surveyor+instance RecvType Respondent+instance RecvType Pull+instance RecvType Bus++-- | Sub socket functionality+class SubscriberType a+instance SubscriberType Sub++-- | Surveyor socket functionality+class SurvType a+instance SurvType Surveyor++-- | Req socket functionality+class ReqType a+instance ReqType Req+++-- * Error handling+--+-- Reimplementing some of Foreign.C.Error here, to substitute nanomsg's errno+-- and strerror functions for the posix ones.++-- | Pretty much any error condition throws this exception.+data NNException = NNException String+ deriving (Eq, Show, Typeable)++instance Exception NNException++mkErrorString :: String -> IO String+mkErrorString loc = do+ errNo <- c_nn_errno+ errCString <- c_nn_strerror errNo+ errString <- peekCString errCString+ return $ printf "nanomsg-haskell error at %s. Errno %d: %s" loc (fromIntegral errNo :: Int) errString++throwErrno :: String -> IO a+throwErrno loc = do+ s <- mkErrorString loc+ throwIO $ NNException s++throwErrnoIf :: (a -> Bool) -> String -> IO a -> IO a+throwErrnoIf p loc action = do+ res <- action+ if p res then throwErrno loc else return res++throwErrnoIf_ :: (a -> Bool) -> String -> IO a -> IO ()+throwErrnoIf_ p loc action = void $ throwErrnoIf p loc action++throwErrnoIfMinus1 :: (Eq a, Num a) => String -> IO a -> IO a+throwErrnoIfMinus1 = throwErrnoIf (== -1)++throwErrnoIfMinus1_ :: (Eq a, Num a) => String -> IO a -> IO ()+throwErrnoIfMinus1_ = throwErrnoIf_ (== -1)++throwErrnoIfRetry :: (a -> Bool) -> String -> IO a -> IO a+throwErrnoIfRetry p loc f = do+ res <- f+ if p res+ then do+ err <- c_nn_errno+ if err == (#const EAGAIN) || err == (#const EINTR)+ then throwErrnoIfRetry p loc f+ else throwErrno loc+ else return res++throwErrnoIfRetry_ :: (a -> Bool) -> String -> IO a -> IO ()+throwErrnoIfRetry_ p loc f = void $ throwErrnoIfRetry p loc f++throwErrnoIfMinus1Retry :: (Eq a, Num a) => String -> IO a -> IO a+throwErrnoIfMinus1Retry = throwErrnoIfRetry (== -1)++throwErrnoIfMinus1Retry_ :: (Eq a, Num a) => String -> IO a -> IO ()+throwErrnoIfMinus1Retry_ = throwErrnoIfRetry_ (== -1)++throwErrnoIfRetryMayBlock :: (a -> Bool) -> String -> IO a -> IO b -> IO a+throwErrnoIfRetryMayBlock p loc f on_block = do+ res <- f+ if p res+ then do+ err <- c_nn_errno+ if err `elem` [ (#const EAGAIN), (#const EINTR), (#const EWOULDBLOCK) ]+ then do+ void on_block+ throwErrnoIfRetryMayBlock p loc f on_block+ else throwErrno loc+ else return res++throwErrnoIfRetryMayBlock_ :: (a -> Bool) -> String -> IO a -> IO b -> IO ()+throwErrnoIfRetryMayBlock_ p loc f on_block = void $ throwErrnoIfRetryMayBlock p loc f on_block++throwErrnoIfMinus1RetryMayBlock :: (Eq a, Num a) => String -> IO a -> IO b -> IO a+throwErrnoIfMinus1RetryMayBlock = throwErrnoIfRetryMayBlock (== -1)++throwErrnoIfMinus1RetryMayBlock_ :: (Eq a, Num a) => String -> IO a -> IO b -> IO ()+throwErrnoIfMinus1RetryMayBlock_ = throwErrnoIfRetryMayBlock_ (== -1)+++-- * FFI functions++-- NN_EXPORT int nn_socket (int domain, int protocol);+foreign import ccall unsafe "nn.h nn_socket"+ c_nn_socket :: CInt -> CInt -> IO CInt++-- NN_EXPORT int nn_bind (int s, const char *addr);+foreign import ccall unsafe "nn.h nn_bind"+ c_nn_bind :: CInt -> CString -> IO CInt++-- NN_EXPORT int nn_connect (int s, const char *addr);+foreign import ccall unsafe "nn.h nn_connect"+ c_nn_connect :: CInt -> CString -> IO CInt++-- NN_EXPORT int nn_shutdown (int s, int how);+foreign import ccall unsafe "nn.h nn_shutdown"+ c_nn_shutdown :: CInt -> CInt -> IO CInt++-- NN_EXPORT int nn_send (int s, const void *buf, size_t len, int flags);+foreign import ccall unsafe "nn.h nn_send"+ c_nn_send :: CInt -> CString -> CInt -> CInt -> IO CInt++-- NN_EXPORT int nn_recv (int s, void *buf, size_t len, int flags);+foreign import ccall unsafe "nn.h nn_recv"+ c_nn_recv_foreignbuf :: CInt -> Ptr CString -> CInt -> CInt -> IO CInt++-- NN_EXPORT int nn_freemsg (void *msg);+foreign import ccall unsafe "nn.h nn_freemsg"+ c_nn_freemsg :: Ptr CChar -> IO CInt++-- NN_EXPORT int nn_close (int s);+foreign import ccall unsafe "nn.h nn_close"+ c_nn_close :: CInt -> IO CInt++-- NN_EXPORT void nn_term (void);+foreign import ccall unsafe "nn.h nn_term"+ c_nn_term :: IO ()++-- NN_EXPORT int nn_setsockopt (int s, int level, int option, const void *optval, size_t optvallen);+foreign import ccall unsafe "nn.h nn_setsockopt"+ c_nn_setsockopt :: CInt -> CInt -> CInt -> Ptr a -> CInt -> IO CInt++-- NN_EXPORT int nn_getsockopt (int s, int level, int option, void *optval, size_t *optvallen);+foreign import ccall unsafe "nn.h nn_getsockopt"+ c_nn_getsockopt :: CInt -> CInt -> CInt -> Ptr a -> Ptr CInt -> IO CInt++-- /* Resolves system errors and native errors to human-readable string. */+-- NN_EXPORT const char *nn_strerror (int errnum);+foreign import ccall unsafe "nn.h nn_strerror"+ c_nn_strerror :: CInt -> IO CString++-- /* This function retrieves the errno as it is known to the library. */+-- /* The goal of this function is to make the code 100% portable, including */+-- /* where the library is compiled with certain CRT library (on Windows) and */+-- /* linked to an application that uses different CRT library. */+-- NN_EXPORT int nn_errno (void);+foreign import ccall unsafe "nn.h nn_errno"+ c_nn_errno :: IO CInt++{-++Unbound FFI functions:++NN_EXPORT int nn_sendmsg (int s, const struct nn_msghdr *msghdr, int flags);+NN_EXPORT int nn_recvmsg (int s, struct nn_msghdr *msghdr, int flags);++NN_EXPORT void *nn_allocmsg (size_t size, int type);+-}++-- * General functions++-- | Creates a socket. Connections are formed using 'bind' or 'connect'.+--+-- See also: 'close'.+socket :: (Protocol a) => a -> IO (Socket a)+socket t = do+ sid <- throwErrnoIfMinus1 "socket" $ c_nn_socket (#const AF_SP) (protocolId t)+ return $ Socket t sid++-- | Creates a socket and runs your action with it.+--+-- E.g. collecting 10 messages:+--+-- > withSocket Sub $ \sub -> do+-- > _ <- connect sub "tcp://localhost:5560"+-- > subscribe sub (C.pack "")+-- > replicateM 10 (recv sub)+--+-- Ensures the socket is closed when your action is done.+withSocket :: (Protocol a) => a -> (Socket a -> IO b) -> IO b+withSocket t = bracket (socket t) close++-- | Binds the socket to a local interface.+--+-- See the nanomsg documentation for specifics on transports.+-- Note that host names do not work for tcp. Some examples are:+--+-- > bind sock "tcp://*:5560"+-- > bind sock "tcp://eth0:5560"+-- > bind sock "tcp://127.0.0.1:5560"+-- > bind sock "inproc://test"+-- > bind sock "ipc:///tmp/test.ipc"+--+-- This function returns an 'Endpoint', which can be supplied+-- to 'shutdown' to remove a connection.+--+-- See also: 'connect', 'shutdown'.+bind :: Socket a -> String -> IO Endpoint+bind (Socket _ sid) addr =+ withCString addr $ \adr -> do+ epid <- throwErrnoIfMinus1 "bind" $ c_nn_bind sid adr+ return $ Endpoint epid++-- | Connects the socket to an endpoint.+--+-- e.g. :+--+-- > connect sock "tcp://localhost:5560"+-- > connect sock "inproc://test"+--+-- See also: 'bind', 'shutdown'.+connect :: Socket a -> String -> IO Endpoint+connect (Socket _ sid) addr =+ withCString addr $ \adr -> do+ epid <- throwErrnoIfMinus1 "connect" $ c_nn_connect sid adr+ return $ Endpoint epid++-- | Removes an endpoint from a socket.+--+-- See also: 'bind', 'connect'.+shutdown :: Socket a -> Endpoint -> IO ()+shutdown (Socket _ sid) (Endpoint eid) =+ throwErrnoIfMinus1_ "shutdown" $ c_nn_shutdown sid eid++-- | Blocking function for sending a message+--+-- See also: 'recv', 'recv''.+send :: (SendType a, Protocol a) => Socket a -> ByteString -> IO ()+send (Socket t sid) string =+ U.unsafeUseAsCStringLen string $ \(ptr, len) ->+ throwErrnoIfMinus1RetryMayBlock_+ "send"+ (c_nn_send sid ptr (fromIntegral len) (#const NN_DONTWAIT))+ (getOptionFd (Socket t sid) (#const NN_SNDFD) >>= threadWaitWrite)++-- | Blocking receive.+recv :: (RecvType a, Protocol a) => Socket a -> IO ByteString+recv (Socket t sid) =+ alloca $ \ptr -> do+ len <- throwErrnoIfMinus1RetryMayBlock+ "recv"+ (c_nn_recv_foreignbuf sid ptr (#const NN_MSG) (#const NN_DONTWAIT))+ (getOptionFd (Socket t sid) (#const NN_RCVFD) >>= threadWaitRead)+ buf <- peek ptr+ str <- C.packCStringLen (buf, fromIntegral len)+ throwErrnoIfMinus1_ "recv freeing message buffer" $ c_nn_freemsg buf+ return str++-- | Nonblocking receive function.+recv' :: (RecvType a, Protocol a) => Socket a -> IO (Maybe ByteString)+recv' (Socket _ sid) =+ alloca $ \ptr -> do+ len <- c_nn_recv_foreignbuf sid ptr (#const NN_MSG) (#const NN_DONTWAIT)+ if len >= 0+ then do+ buf <- peek ptr+ str <- C.packCStringLen (buf, fromIntegral len)+ throwErrnoIfMinus1_ "recv' freeing message buffer" $ c_nn_freemsg buf+ return $ Just str+ else do+ errno <- c_nn_errno+ if errno == (#const EAGAIN) || errno == (#const EINTR)+ then return Nothing+ else throwErrno "recv'"++-- | Subscribe to a given subject string.+subscribe :: (SubscriberType a, Protocol a) => Socket a -> ByteString -> IO ()+subscribe (Socket t sid) string =+ setOption (Socket t sid) (protocolId t) (#const NN_SUB_SUBSCRIBE) (StringOption string)++-- | Unsubscribes from a subject.+unsubscribe :: (SubscriberType a, Protocol a) => Socket a -> ByteString -> IO ()+unsubscribe (Socket t sid) string =+ setOption (Socket t sid) (protocolId t) (#const NN_SUB_UNSUBSCRIBE) (StringOption string)++-- | Closes the socket. Any buffered inbound messages that were not yet+-- received by the application will be discarded. The library will try to+-- deliver any outstanding outbound messages for the time specified by+-- NN_LINGER socket option. The call will block in the meantime.+close :: Socket a -> IO ()+close (Socket _ sid) =+ throwErrnoIfMinus1Retry_ "close" $ c_nn_close sid++-- | Switches nanomsg into shutdown modus and interrupts any waiting+-- function calls.+term :: IO ()+term = c_nn_term+++-- * Socket option accessors and mutators++-- not sure if this beats having setOptionInt and setOptionString..+data SocketOption = IntOption Int | StringOption ByteString+ deriving (Show)++-- Used for setting a socket option.+setOption :: Socket a -> CInt -> CInt -> SocketOption -> IO ()++setOption (Socket _ sid) level option (IntOption val) =+ alloca $ \ptr -> do+ poke ptr (fromIntegral val :: CInt)+ let cintSize = fromIntegral $ sizeOf (fromIntegral val :: CInt) :: CInt+ throwErrnoIfMinus1_ "setOption (int)" $ c_nn_setsockopt sid level option ptr cintSize++setOption (Socket _ sid) level option (StringOption str) =+ throwErrnoIfMinus1_ "setOption (string)" <$> U.unsafeUseAsCStringLen str $+ \(ptr, len) -> c_nn_setsockopt sid level option ptr (fromIntegral len)++-- Reads a socket option.+getOption :: Socket a -> CInt -> CInt -> IO CInt+getOption (Socket _ sid) level option =+ alloca $ \ptr ->+ alloca $ \sizePtr -> do+ let a = 1 :: CInt+ let cintSize = fromIntegral $ sizeOf a+ poke sizePtr cintSize+ throwErrnoIfMinus1_ "getOption" $ c_nn_getsockopt sid level option (ptr :: Ptr CInt) sizePtr+ value <- peek ptr+ size <- peek sizePtr+ if cintSize /= size then throwErrno "getOption: output size not as expected" else return value++-- Retrieves a nanomsg file descriptor for polling ready status.+getOptionFd :: Socket a -> CInt -> IO Fd+getOptionFd (Socket _ sid) option =+ alloca $ \ptr ->+ alloca $ \sizePtr -> do+ let a = 1 :: Fd+ let fdSize = fromIntegral $ sizeOf a+ poke sizePtr fdSize+ throwErrnoIfMinus1_ "getOptionFd" $ c_nn_getsockopt sid (#const NN_SOL_SOCKET) option (ptr :: Ptr Fd) sizePtr+ value <- peek ptr+ size <- peek sizePtr+ if fdSize /= size then throwErrno "getOptionFd: output size not as expected" else return value++-- | Specifies how long should the socket try to send pending outbound+-- messages after close has been called, in milliseconds.+--+-- Negative value means infinite linger. Default value is 1000 (1 second).+linger :: Socket a -> IO Int+linger s =+ fromIntegral <$> getOption s (#const NN_SOL_SOCKET) (#const NN_LINGER)++-- | Specifies how long should the socket try to send pending outbound+-- messages after close has been called, in milliseconds.+--+-- Negative value means infinite linger. Default value is 1000 (1 second).+setLinger :: Socket a -> Int -> IO ()+setLinger s val =+ setOption s (#const NN_SOL_SOCKET) (#const NN_LINGER) (IntOption val)++-- | Size of the send buffer, in bytes. To prevent blocking for messages+-- larger than the buffer, exactly one message may be buffered in addition+-- to the data in the send buffer.+--+-- Default value is 128kB.+sndBuf :: Socket a -> IO Int+sndBuf s =+ fromIntegral <$> getOption s (#const NN_SOL_SOCKET) (#const NN_SNDBUF)++-- | Size of the send buffer, in bytes. To prevent blocking for messages+-- larger than the buffer, exactly one message may be buffered in addition+-- to the data in the send buffer.+--+-- Default value is 128kB.+setSndBuf :: Socket a -> Int -> IO ()+setSndBuf s val =+ setOption s (#const NN_SOL_SOCKET) (#const NN_SNDBUF) (IntOption val)++-- | Size of the receive buffer, in bytes. To prevent blocking for messages+-- larger than the buffer, exactly one message may be buffered in addition+-- to the data in the receive buffer.+--+-- Default value is 128kB.+rcvBuf :: Socket a -> IO Int+rcvBuf s =+ fromIntegral <$> getOption s (#const NN_SOL_SOCKET) (#const NN_RCVBUF)++-- | Size of the receive buffer, in bytes. To prevent blocking for messages+-- larger than the buffer, exactly one message may be buffered in addition+-- to the data in the receive buffer.+--+-- Default value is 128kB.+setRcvBuf :: Socket a -> Int -> IO ()+setRcvBuf s val =+ setOption s (#const NN_SOL_SOCKET) (#const NN_RCVBUF) (IntOption val)++-- Think I'll just skip these. There's recv' for nonblocking receive, and+-- adding a return value to send seems awkward.+--sendTimeout+--recvTimeout++-- | For connection-based transports such as TCP, this option specifies+-- how long to wait, in milliseconds, when connection is broken before+-- trying to re-establish it.+--+-- Note that actual reconnect interval may be randomised to some extent+-- to prevent severe reconnection storms.+--+-- Default value is 100 (0.1 second).+reconnectInterval :: Socket a -> IO Int+reconnectInterval s =+ fromIntegral <$> getOption s (#const NN_SOL_SOCKET) (#const NN_RECONNECT_IVL)++-- | For connection-based transports such as TCP, this option specifies+-- how long to wait, in milliseconds, when connection is broken before+-- trying to re-establish it.+--+-- Note that actual reconnect interval may be randomised to some extent+-- to prevent severe reconnection storms.+--+-- Default value is 100 (0.1 second).+setReconnectInterval :: Socket a -> Int -> IO ()+setReconnectInterval s val =+ setOption s (#const NN_SOL_SOCKET) (#const NN_RECONNECT_IVL) (IntOption val)++-- | This option is to be used only in addition to NN_RECONNECT_IVL option.+-- It specifies maximum reconnection interval. On each reconnect attempt,+-- the previous interval is doubled until NN_RECONNECT_IVL_MAX is reached.+--+-- Value of zero means that no exponential backoff is performed and reconnect+-- interval is based only on NN_RECONNECT_IVL. If NN_RECONNECT_IVL_MAX is+-- less than NN_RECONNECT_IVL, it is ignored.+--+-- Default value is 0.+reconnectIntervalMax :: Socket a -> IO Int+reconnectIntervalMax s =+ fromIntegral <$> getOption s (#const NN_SOL_SOCKET) (#const NN_RECONNECT_IVL_MAX)++-- | This option is to be used only in addition to NN_RECONNECT_IVL option.+-- It specifies maximum reconnection interval. On each reconnect attempt,+-- the previous interval is doubled until NN_RECONNECT_IVL_MAX is reached.+--+-- Value of zero means that no exponential backoff is performed and reconnect+-- interval is based only on NN_RECONNECT_IVL. If NN_RECONNECT_IVL_MAX is+-- less than NN_RECONNECT_IVL, it is ignored.+--+-- Default value is 0.+setReconnectIntervalMax :: Socket a -> Int -> IO ()+setReconnectIntervalMax s val =+ setOption s (#const NN_SOL_SOCKET) (#const NN_RECONNECT_IVL_MAX) (IntOption val)++-- | Sets outbound priority for endpoints subsequently added to the socket.+-- This option has no effect on socket types that send messages to all the+-- peers. However, if the socket type sends each message to a single peer+-- (or a limited set of peers), peers with high priority take precedence over+-- peers with low priority.+--+-- Highest priority is 1, lowest priority is 16. Default value is 8.+sndPrio :: Socket a -> IO Int+sndPrio s =+ fromIntegral <$> getOption s (#const NN_SOL_SOCKET) (#const NN_SNDPRIO)++-- | Sets outbound priority for endpoints subsequently added to the socket.+-- This option has no effect on socket types that send messages to all the+-- peers. However, if the socket type sends each message to a single peer+-- (or a limited set of peers), peers with high priority take precedence over+-- peers with low priority.+--+-- Highest priority is 1, lowest priority is 16. Default value is 8.+setSndPrio :: Socket a -> Int -> IO ()+setSndPrio s val =+ setOption s (#const NN_SOL_SOCKET) (#const NN_SNDPRIO) (IntOption val)++-- | If set to 1, only IPv4 addresses are used. If set to 0, both IPv4+-- and IPv6 addresses are used.+--+-- Default value is 1.+ipv4Only :: Socket a -> IO Int+ipv4Only s =+ fromIntegral <$> getOption s (#const NN_SOL_SOCKET) (#const NN_IPV4ONLY)++-- | If set to 1, only IPv4 addresses are used. If set to 0, both IPv4+-- and IPv6 addresses are used.+--+-- Default value is 1.+setIpv4Only :: Socket a -> Int -> IO ()+setIpv4Only s val =+ setOption s (#const NN_SOL_SOCKET) (#const NN_IPV4ONLY) (IntOption val)++-- | This option is defined on the full REQ socket. If reply is not received+-- in specified amount of milliseconds, the request will be automatically+-- resent.+--+-- Default value is 60000 (1 minute).+requestResendInterval :: (ReqType a) => Socket a -> IO Int+requestResendInterval s =+ fromIntegral <$> getOption s (#const NN_REQ) (#const NN_REQ_RESEND_IVL)++-- | This option is defined on the full REQ socket. If reply is not received+-- in specified amount of milliseconds, the request will be automatically+-- resent.+--+-- Default value is 60000 (1 minute).+setRequestResendInterval :: (ReqType a) => Socket a -> Int -> IO ()+setRequestResendInterval s val =+ setOption s (#const NN_REQ) (#const NN_REQ_RESEND_IVL) (IntOption val)++-- | This option, when set to 1, disables Nagle's algorithm.+--+-- Default value is 0.+tcpNoDelay :: Socket a -> IO Int+tcpNoDelay s =+ fromIntegral <$> getOption s (#const NN_TCP) (#const NN_TCP_NODELAY)++-- | This option, when set to 1, disables Nagle's algorithm.+--+-- Default value is 0.+setTcpNoDelay :: Socket a -> Int -> IO ()+setTcpNoDelay s val =+ setOption s (#const NN_TCP) (#const NN_TCP_NODELAY) (IntOption val)+
+ tests/Properties.hs view
@@ -0,0 +1,206 @@+{-# LANGUAGE TemplateHaskell #-}+{-# OPTIONS_GHC -fno-warn-orphans #-}+module Main where++import Nanomsg+import Test.Framework.TH (defaultMainGenerator)+import Test.Framework.Providers.QuickCheck2 (testProperty)+import Test.QuickCheck+import Test.QuickCheck.Monadic+import Data.ByteString (ByteString)+import qualified Data.ByteString.Char8 as C+import Control.Concurrent (threadDelay)+import Control.Applicative ( (<$>) )+import Data.Maybe (catMaybes)++instance Arbitrary ByteString where+ arbitrary = C.pack <$> arbitrary++-- dummy test+prop_reverse :: [Int] -> Bool+prop_reverse xs =+ xs == reverse (reverse xs)++-- test Pub and Sub sockets+prop_PubSub :: Property+prop_PubSub = monadicIO $ do+ msgs <- pick arbitrary+ pre $ not (null msgs)+ res <- run $ do+ pub <- socket Pub+ ep1 <- bind pub "inproc://pubsub"+ sub1 <- socket Sub+ ep2 <- connect sub1 "inproc://pubsub"+ subscribe sub1 $ C.pack ""+ sub2 <- socket Sub+ ep3 <- connect sub2 "inproc://pubsub"+ subscribe sub2 $ C.pack ""+ threadDelay 1000+ r <- mapM (sendMsg pub sub1 sub2) msgs+ unsubscribe sub2 $ C.pack ""+ unsubscribe sub1 $ C.pack ""+ shutdown sub2 ep3+ shutdown sub1 ep2+ shutdown pub ep1+ close pub+ close sub1+ close sub2+ threadDelay 1000+ return r+ assert $ and res+ where+ sendMsg pub sub1 sub2 msg = do+ send pub msg+ send pub msg+ a <- recv sub1+ b <- recv sub1+ c <- recv sub2+ d <- recv sub2+ return $ a == msg && b == msg && c == msg && d == msg++-- test Pair sockets+prop_Pair :: Property+prop_Pair = monadicIO $ do+ msgs <- pick arbitrary+ pre $ not (null msgs)+ res <- run $ do+ s1 <- socket Pair+ _ <- bind s1 "inproc://pair"+ s2 <- socket Pair+ _ <- connect s2 "inproc://pair"+ threadDelay 1000+ -- Send message from s1 to s2, then back from s2 to s1, then make sure it hasn't changed+ r <- mapM (\m -> send s1 m >> recv s2 >>= send s2 >> recv s1 >>= return . (== m)) msgs+ close s1+ close s2+ threadDelay 1000+ return r+ assert $ and res++-- test Pipeline (Push & Pull) sockets+prop_Pipeline :: Property+prop_Pipeline = monadicIO $ do+ msgs <- pick arbitrary+ pre $ not (null msgs)+ res <- run $ do+ push <- socket Push+ _ <- bind push "inproc://pipeline"+ pull1 <- socket Pull+ pull2 <- socket Pull+ _ <- connect pull1 "inproc://pipeline"+ _ <- connect pull2 "inproc://pipeline"+ threadDelay 1000+ r <- mapM (testSockets push pull1 pull2) msgs+ close push+ close pull1+ close pull2+ threadDelay 1000+ return r+ assert $ and res+ where+ testSockets push pull1 pull2 msg = do+ send push msg+ send push msg+ send push msg+ threadDelay 1000+ a <- recv' pull1+ b <- recv' pull1+ c <- recv' pull1+ d <- recv' pull2+ e <- recv' pull2+ f <- recv' pull2+ let xs = catMaybes [a, b, c, d, e, f]+ return $ all (== msg) xs && (length xs == 3)++-- test Req and Rep sockets+prop_ReqRep :: Property+prop_ReqRep = monadicIO $ do+ msgs <- pick arbitrary+ pre $ not (null msgs)+ res <- run $ do+ req <- socket Req+ _ <- bind req "inproc://reqrep"+ rep <- socket Rep+ _ <- connect rep "inproc://reqrep"+ threadDelay 1000+ r <- mapM (\m -> send req m >> recv rep >>= send rep >> recv req >>= return . (== m)) msgs+ close req+ close rep+ threadDelay 1000+ return r+ assert $ and res++-- test Bus socket+prop_Bus :: Property+prop_Bus = monadicIO $ do+ msgs <- pick arbitrary+ pre $ not (null msgs)+ res <- run $ do+ -- Probably not how you're supposed to connect Bus nodes..+ b1 <- socket Bus+ _ <- bind b1 "inproc://bus1"+ b2 <- socket Bus+ _ <- connect b2 "inproc://bus1"+ _ <- bind b2 "inproc://bus2"+ b3 <- socket Bus+ _ <- connect b3 "inproc://bus2"+ _ <- bind b3 "inproc://bus3"+ _ <- connect b1 "inproc://bus3"+ threadDelay 1000+ r <- mapM (testSockets b1 b2 b3) msgs+ close b1+ close b2+ close b3+ threadDelay 1000+ return r+ assert $ and res+ where+ testSockets b1 b2 b3 msg = do+ send b1 msg+ a <- recv b2+ b <- recv b3+ send b2 msg+ c <- recv b1+ d <- recv b3+ send b3 msg+ e <- recv b1+ f <- recv b2+ return $ all (== msg) [a, b, c, d, e, f]++prop_TestOptions :: Property+prop_TestOptions = monadicIO $ do+ res <- run $ do+ req <- socket Req+ _ <- bind req "tcp://*:5560"+ threadDelay 1000+ setTcpNoDelay req 1+ v1 <- tcpNoDelay req+ setTcpNoDelay req 0+ v2 <- tcpNoDelay req+ setRequestResendInterval req 30000+ v3 <- requestResendInterval req+ setIpv4Only req 0+ v4 <- ipv4Only req+ setIpv4Only req 1+ v5 <- ipv4Only req+ setSndPrio req 7+ v6 <- sndPrio req+ setReconnectInterval req 50+ v7 <- reconnectInterval req+ setReconnectIntervalMax req 400+ v8 <- reconnectIntervalMax req+ setRcvBuf req 200000+ v9 <- rcvBuf req+ setSndBuf req 150000+ v10 <- sndBuf req+ setLinger req 500+ v11 <- linger req+ close req+ threadDelay 1000+ return [v1 == 1, v2 == 0, v3 == 30000, v4 == 0, v5 == 1, v6 == 7,+ v7 == 50, v8 == 400, v9 == 200000, v10 == 150000, v11 == 500]+ assert $ and res++main :: IO ()+main = $defaultMainGenerator+