pipes-network 0.0.2 → 0.1.0
raw patch · 11 files changed
+1430/−156 lines, 11 filesdep +HUnitdep +pipesdep +pipes-networkdep −pipes-coredep ~basedep ~bytestringdep ~transformersnew-uploaderPVP ok
version bump matches the API change (PVP)
Dependencies added: HUnit, pipes, pipes-network, pipes-safe, test-framework, test-framework-hunit
Dependencies removed: pipes-core
Dependency ranges changed: base, bytestring, transformers
API changes (from Hackage documentation)
- Control.Pipe.Network: ClientSettings :: Int -> String -> ClientSettings
- Control.Pipe.Network: ServerSettings :: Int -> Maybe String -> ServerSettings
- Control.Pipe.Network: clientHost :: ClientSettings -> String
- Control.Pipe.Network: clientPort :: ClientSettings -> Int
- Control.Pipe.Network: data ClientSettings
- Control.Pipe.Network: data ServerSettings
- Control.Pipe.Network: runTCPClient :: MonadIO m => ClientSettings -> Application m r -> IO r
- Control.Pipe.Network: runTCPServer :: MonadIO m => ServerSettings -> Application m r -> IO r
- Control.Pipe.Network: serverHost :: ServerSettings -> Maybe String
- Control.Pipe.Network: serverPort :: ServerSettings -> Int
- Control.Pipe.Network: socketReader :: MonadIO m => Socket -> Pipe () ByteString m ()
- Control.Pipe.Network: socketWriter :: MonadIO m => Socket -> Pipe ByteString Void m r
- Control.Pipe.Network: type Application m r = Pipe () ByteString m () -> Pipe ByteString Void m () -> IO r
+ Control.Proxy.TCP: Host :: HostName -> HostPreference
+ Control.Proxy.TCP: HostAny :: HostPreference
+ Control.Proxy.TCP: HostIPv4 :: HostPreference
+ Control.Proxy.TCP: HostIPv6 :: HostPreference
+ Control.Proxy.TCP: Timeout :: String -> Timeout
+ Control.Proxy.TCP: accept :: Socket -> ((Socket, SockAddr) -> IO b) -> IO b
+ Control.Proxy.TCP: acceptFork :: Socket -> ((Socket, SockAddr) -> IO ()) -> IO ThreadId
+ Control.Proxy.TCP: bindSock :: HostPreference -> ServiceName -> IO (Socket, SockAddr)
+ Control.Proxy.TCP: connect :: HostName -> ServiceName -> ((Socket, SockAddr) -> IO r) -> IO r
+ Control.Proxy.TCP: connectSock :: HostName -> ServiceName -> IO (Socket, SockAddr)
+ Control.Proxy.TCP: data HostPreference
+ Control.Proxy.TCP: data Timeout
+ Control.Proxy.TCP: listen :: HostPreference -> ServiceName -> ((Socket, SockAddr) -> IO r) -> IO r
+ Control.Proxy.TCP: nsocketReadS :: Proxy p => Socket -> Int -> Server p Int ByteString IO ()
+ Control.Proxy.TCP: nsocketReadTimeoutS :: Proxy p => Int -> Socket -> Int -> Server (EitherP Timeout p) Int ByteString IO ()
+ Control.Proxy.TCP: serve :: HostPreference -> ServiceName -> ((Socket, SockAddr) -> IO r) -> IO r
+ Control.Proxy.TCP: serveFork :: HostPreference -> ServiceName -> ((Socket, SockAddr) -> IO ()) -> IO ()
+ Control.Proxy.TCP: socketReadS :: Proxy p => Int -> Socket -> () -> Producer p ByteString IO ()
+ Control.Proxy.TCP: socketReadTimeoutS :: Proxy p => Int -> Int -> Socket -> () -> Producer (EitherP Timeout p) ByteString IO ()
+ Control.Proxy.TCP: socketWriteD :: Proxy p => Socket -> x -> p x ByteString x ByteString IO r
+ Control.Proxy.TCP: socketWriteTimeoutD :: Proxy p => Int -> Socket -> x -> (EitherP Timeout p) x ByteString x ByteString IO r
+ Control.Proxy.TCP.Safe: Host :: HostName -> HostPreference
+ Control.Proxy.TCP.Safe: HostAny :: HostPreference
+ Control.Proxy.TCP.Safe: HostIPv4 :: HostPreference
+ Control.Proxy.TCP.Safe: HostIPv6 :: HostPreference
+ Control.Proxy.TCP.Safe: Timeout :: String -> Timeout
+ Control.Proxy.TCP.Safe: accept :: (Proxy p, Monad m) => (forall x. SafeIO x -> m x) -> Socket -> ((Socket, SockAddr) -> ExceptionP p a' a b' b m r) -> ExceptionP p a' a b' b m r
+ Control.Proxy.TCP.Safe: acceptFork :: (Proxy p, Monad m) => (forall x. SafeIO x -> m x) -> Socket -> ((Socket, SockAddr) -> IO ()) -> ExceptionP p a' a b' b m ThreadId
+ Control.Proxy.TCP.Safe: connect :: (Proxy p, Monad m) => (forall x. SafeIO x -> m x) -> HostName -> ServiceName -> ((Socket, SockAddr) -> ExceptionP p a' a b' b m r) -> ExceptionP p a' a b' b m r
+ Control.Proxy.TCP.Safe: connectReadS :: Proxy p => Maybe Int -> Int -> HostName -> ServiceName -> () -> Producer (ExceptionP p) ByteString SafeIO ()
+ Control.Proxy.TCP.Safe: connectWriteD :: Proxy p => Maybe Int -> HostName -> ServiceName -> x -> (ExceptionP p) x ByteString x ByteString SafeIO ()
+ Control.Proxy.TCP.Safe: data HostPreference
+ Control.Proxy.TCP.Safe: data Timeout
+ Control.Proxy.TCP.Safe: listen :: (Proxy p, Monad m) => (forall x. SafeIO x -> m x) -> HostPreference -> ServiceName -> ((Socket, SockAddr) -> ExceptionP p a' a b' b m r) -> ExceptionP p a' a b' b m r
+ Control.Proxy.TCP.Safe: nsocketReadS :: Proxy p => Maybe Int -> Socket -> Int -> Server (ExceptionP p) Int ByteString SafeIO ()
+ Control.Proxy.TCP.Safe: serve :: (Proxy p, Monad m) => (forall x. SafeIO x -> m x) -> HostPreference -> ServiceName -> ((Socket, SockAddr) -> ExceptionP p a' a b' b m r) -> ExceptionP p a' a b' b m r
+ Control.Proxy.TCP.Safe: serveFork :: (Proxy p, Monad m) => (forall x. SafeIO x -> m x) -> HostPreference -> ServiceName -> ((Socket, SockAddr) -> IO ()) -> ExceptionP p a' a b' b m r
+ Control.Proxy.TCP.Safe: serveReadS :: Proxy p => Maybe Int -> Int -> HostPreference -> ServiceName -> () -> Producer (ExceptionP p) ByteString SafeIO ()
+ Control.Proxy.TCP.Safe: serveWriteD :: Proxy p => Maybe Int -> HostPreference -> ServiceName -> x -> (ExceptionP p) x ByteString x ByteString SafeIO ()
+ Control.Proxy.TCP.Safe: socketReadS :: Proxy p => Maybe Int -> Int -> Socket -> () -> Producer (ExceptionP p) ByteString SafeIO ()
+ Control.Proxy.TCP.Safe: socketWriteD :: Proxy p => Maybe Int -> Socket -> x -> (ExceptionP p) x ByteString x ByteString SafeIO r
+ Control.Proxy.TCP.Safe.Sync: Receive :: Int -> Request t
+ Control.Proxy.TCP.Safe.Sync: Received :: ByteString -> Response
+ Control.Proxy.TCP.Safe.Sync: Send :: t -> Request t
+ Control.Proxy.TCP.Safe.Sync: Sent :: Response
+ Control.Proxy.TCP.Safe.Sync: data Request t
+ Control.Proxy.TCP.Safe.Sync: data Response
+ Control.Proxy.TCP.Safe.Sync: socketSyncProxy :: Proxy p => Maybe Int -> Socket -> Request a' -> (ExceptionP p) a' ByteString (Request a') Response SafeIO ()
+ Control.Proxy.TCP.Safe.Sync: socketSyncServer :: Proxy p => Maybe Int -> Socket -> Request ByteString -> Server (ExceptionP p) (Request ByteString) Response SafeIO ()
+ Control.Proxy.TCP.Safe.Sync: syncDelimit :: (Monad m, Proxy p) => Int -> ByteString -> b' -> p (Request b') Response b' ByteString m r
+ Control.Proxy.TCP.Sync: Receive :: Int -> Request t
+ Control.Proxy.TCP.Sync: Received :: ByteString -> Response
+ Control.Proxy.TCP.Sync: Send :: t -> Request t
+ Control.Proxy.TCP.Sync: Sent :: Response
+ Control.Proxy.TCP.Sync: data Request t
+ Control.Proxy.TCP.Sync: data Response
+ Control.Proxy.TCP.Sync: instance Eq Response
+ Control.Proxy.TCP.Sync: instance Eq t => Eq (Request t)
+ Control.Proxy.TCP.Sync: instance Read Response
+ Control.Proxy.TCP.Sync: instance Read t => Read (Request t)
+ Control.Proxy.TCP.Sync: instance Show Response
+ Control.Proxy.TCP.Sync: instance Show t => Show (Request t)
+ Control.Proxy.TCP.Sync: socketSyncProxy :: Proxy p => Socket -> Request a' -> p a' ByteString (Request a') Response IO ()
+ Control.Proxy.TCP.Sync: socketSyncProxyTimeout :: Proxy p => Int -> Socket -> Request a' -> (EitherP Timeout p) a' ByteString (Request a') Response IO ()
+ Control.Proxy.TCP.Sync: socketSyncServer :: Proxy p => Socket -> Request ByteString -> Server p (Request ByteString) Response IO ()
+ Control.Proxy.TCP.Sync: socketSyncServerTimeout :: Proxy p => Int -> Socket -> Request ByteString -> Server (EitherP Timeout p) (Request ByteString) Response IO ()
+ Control.Proxy.TCP.Sync: syncDelimit :: (Monad m, Proxy p) => Int -> ByteString -> b' -> p (Request b') Response b' ByteString m r
Files
- AUTHORS +3/−0
- Control/Pipe/Network.hs +0/−130
- LICENSE +3/−2
- README.md +13/−0
- pipes-network.cabal +62/−24
- src/Control/Proxy/Network/Internal.hs +59/−0
- src/Control/Proxy/TCP.hs +385/−0
- src/Control/Proxy/TCP/Safe.hs +435/−0
- src/Control/Proxy/TCP/Safe/Sync.hs +121/−0
- src/Control/Proxy/TCP/Sync.hs +203/−0
- tests/Simple.hs +146/−0
+ AUTHORS view
@@ -0,0 +1,3 @@+Renzo Carbonara+Gabriel Gonzalez+Paolo Capriotti
− Control/Pipe/Network.hs
@@ -1,130 +0,0 @@-{-# LANGUAGE ScopedTypeVariables #-}-module Control.Pipe.Network (- Application,- socketReader,- socketWriter,- ServerSettings(..),- runTCPServer,- ClientSettings(..),- runTCPClient,- ) where--import qualified Network.Socket as NS-import Network.Socket (Socket)-import Network.Socket.ByteString (sendAll, recv)-import Data.ByteString (ByteString)-import qualified Data.ByteString as B-import Control.Concurrent (forkIO)-import qualified Control.Exception as E-import Control.Monad (forever, unless)-import Control.Monad.IO.Class-import Control.Monad.Trans.Class-import Control.Pipe---- adapted from conduit---- | Stream data from the socket.-socketReader :: MonadIO m => Socket -> Pipe () ByteString m ()-socketReader socket = go- where- go = do- bs <- lift . liftIO $ recv socket 4096- unless (B.null bs) $- yield bs >> go---- | Stream data to the socket.-socketWriter :: MonadIO m => Socket -> Pipe ByteString Void m r-socketWriter socket = forever $ await >>= lift . liftIO . sendAll socket---- | A simple TCP application. It takes two arguments: the 'Producer' to read--- input data from, and the 'Consumer' to send output data to.-type Application m r = Pipe () ByteString m ()- -> Pipe ByteString Void m ()- -> IO r---- | Settings for a TCP server. It takes a port to listen on, and an optional--- hostname to bind to.-data ServerSettings = ServerSettings- { serverPort :: Int- , serverHost :: Maybe String -- ^ 'Nothing' indicates no preference- }---- | Run an @Application@ with the given settings. This function will create a--- new listening socket, accept connections on it, and spawn a new thread for--- each connection.-runTCPServer :: MonadIO m => ServerSettings -> Application m r -> IO r-runTCPServer (ServerSettings port host) app = E.bracket- (bindPort host port)- NS.sClose- (forever . serve)- where- serve lsocket = do- (socket, _addr) <- NS.accept lsocket- forkIO $ do- E.finally- (app (socketReader socket) (socketWriter socket))- (NS.sClose socket)- return ()---- | Settings for a TCP client, specifying how to connect to the server.-data ClientSettings = ClientSettings- { clientPort :: Int- , clientHost :: String- }---- | Run an 'Application' by connecting to the specified server.-runTCPClient :: MonadIO m => ClientSettings -> Application m r -> IO r-runTCPClient (ClientSettings port host) app = E.bracket- (getSocket host port)- NS.sClose- (\s -> app (socketReader s) (socketWriter s))---- | Attempt to connect to the given host/port.-getSocket :: String -> Int -> IO NS.Socket-getSocket host' port' = do- let hints = NS.defaultHints {- NS.addrFlags = [NS.AI_ADDRCONFIG]- , NS.addrSocketType = NS.Stream- }- (addr:_) <- NS.getAddrInfo (Just hints) (Just host') (Just $ show port')- E.bracketOnError- (NS.socket (NS.addrFamily addr)- (NS.addrSocketType addr)- (NS.addrProtocol addr))- NS.sClose- (\sock -> NS.connect sock (NS.addrAddress addr) >> return sock)---- | Attempt to bind a listening @Socket@ on the given host/port. If no host is--- given, will use the first address available.-bindPort :: Maybe String -> Int -> IO Socket-bindPort host p = do- let hints = NS.defaultHints- { NS.addrFlags =- [ NS.AI_PASSIVE- , NS.AI_NUMERICSERV- , NS.AI_NUMERICHOST- ]- , NS.addrSocketType = NS.Stream- }- port = Just . show $ p- addrs <- NS.getAddrInfo (Just hints) host port- let- tryAddrs (addr1:rest@(_:_)) = E.catch- (theBody addr1)- (\(_ :: E.IOException) -> tryAddrs rest)- tryAddrs (addr1:[]) = theBody addr1- tryAddrs _ = error "bindPort: addrs is empty"- theBody addr =- E.bracketOnError- (NS.socket- (NS.addrFamily addr)- (NS.addrSocketType addr)- (NS.addrProtocol addr))- NS.sClose- (\sock -> do- NS.setSocketOption sock NS.ReuseAddr 1- NS.bindSocket sock (NS.addrAddress addr)- NS.listen sock NS.maxListenQueue- return sock- )- tryAddrs addrs
LICENSE view
@@ -1,4 +1,5 @@-Copyright (c)2012, Paolo Capriotti+Copyright (c) 2012-2013, Renzo Carbonara+Copyright (c) 2012-2012, Paolo Capriotti All rights reserved. @@ -13,7 +14,7 @@ disclaimer in the documentation and/or other materials provided with the distribution. - * Neither the name of Paolo Capriotti nor the names of other+ * Neither the name of Renzo Carbonara nor the names of other contributors may be used to endorse or promote products derived from this software without specific prior written permission.
+ README.md view
@@ -0,0 +1,13 @@+# pipes-network++Utilities to deal with sockets using the **pipes** and **pipes-safe**+libraries.++Check the source or rendered Haddocks for extensive documentation.++This code is licensed under the terms of the so called **3-clause BSD+license**. Read the file named ``LICENSE`` found in this same directory+for details.++See the ``AUTHORS`` file to learn about the people involved in this+effort.
pipes-network.cabal view
@@ -1,27 +1,65 @@-Name: pipes-network-Version: 0.0.2-License: BSD3-License-file: LICENSE-Author: Paolo Capriotti-Maintainer: p.capriotti@gmail.com-Stability: Experimental-Homepage: https://github.com/pcapriotti/pipes-extra-Category: Control, Enumerator, Network-Build-type: Simple-Synopsis: Utilities to deal with sockets.-Description: Utilities to deal with sockets.-Cabal-version: >=1.8+name: pipes-network+version: 0.1.0+license: BSD3+license-file: LICENSE+copyright: Copyright (c) Renzo Carbonara 2012-2013, Paolo Capriotti 2012-2012.+author: Renzo Carbonara+maintainer: renzocarbonaraλgmail.com+stability: Experimental+tested-with: GHC == 7.4.1+homepage: https://github.com/k0001/pipes-network+bug-reports: https://github.com/k0001/pipes-network/issues+category: Control, Pipes, Network+build-type: Simple+synopsis: Use network sockets together with the pipes library.+cabal-version: >=1.8+extra-source-files: README.md AUTHORS+description:+ Use network sockets together with the @pipes@ library.+ .+ This package is organized using the following namespaces:+ .+ * "Control.Proxy.TCP" exports @pipes@ proxies and functions to deal with TCP+ connections. Such proxies don't acquire nor release new resources within a+ proxy pipeline.+ .+ * "Control.Proxy.TCP.Safe" exports @pipes-safe@ proxies and functions to deal+ with TCP connections. Such proxies may safely acquire and release resources+ within a pipeline, using the facilities provided by the @pipes-safe@ package. -Source-Repository head- Type: git- Location: https://github.com/pcapriotti/pipes-extra+source-repository head+ type: git+ location: git://github.com/k0001/pipes-network.git -Library- Build-Depends:+library+ hs-source-dirs: src+ build-depends: base (== 4.*),- transformers (>= 0.2 && < 0.4),- pipes-core (== 0.1.*),- bytestring (== 0.9.*),- network- Exposed-Modules:- Control.Pipe.Network+ bytestring (>= 0.9.2.1),+ network,+ pipes (==3.1.*),+ pipes-safe (>=1.0),+ transformers (>= 0.2 && < 0.4)+ exposed-modules:+ Control.Proxy.Network.Internal+ Control.Proxy.TCP+ Control.Proxy.TCP.Sync+ Control.Proxy.TCP.Safe+ Control.Proxy.TCP.Safe.Sync+ ghc-options: -Wall -fno-warn-unused-do-bind++test-suite simple+ hs-source-dirs: tests+ main-is: Simple.hs+ type: exitcode-stdio-1.0+ build-depends:+ base < 5,+ bytestring,+ HUnit,+ network,+ pipes-network,+ pipes,+ pipes-safe,+ transformers,+ test-framework,+ test-framework-hunit
+ src/Control/Proxy/Network/Internal.hs view
@@ -0,0 +1,59 @@+{-# LANGUAGE DeriveDataTypeable #-}+{-# OPTIONS_HADDOCK hide #-}++-- | This module doesn't belong to this namespace, but really, I don't+-- know where it belongs. Suggestions welcome.++-- Some code in this file was adapted from the @network-conduit@ library by+-- Michael Snoyman. Copyright (c) 2011. See its licensing terms (BSD3) at:+-- https://github.com/snoyberg/conduit/blob/master/network-conduit/LICENSE++module Control.Proxy.Network.Internal (+ HostPreference(..),+ hpHostName,+ Timeout(..)+ ) where++import qualified Control.Exception as E+import Data.String (IsString (fromString))+import Data.Typeable (Typeable)+import qualified Network.Socket as NS++-- | Preferred host to bind.+data HostPreference+ = HostAny -- ^Any avaiable host.+ | HostIPv4 -- ^Any avaiable IPv4 host.+ | HostIPv6 -- ^Any avaiable IPv6 host.+ | Host NS.HostName -- ^An explicit host name.+ deriving (Eq, Ord, Show, Read)+++-- | The following special values are recognized:+--+-- * @*@ means 'HostAny'+--+-- * @*4@ means 'HostIPv4'+--+-- * @*6@ means 'HostIPv6'+--+-- * Any other string is 'Host'+instance IsString HostPreference where+ fromString "*" = HostAny+ fromString "*4" = HostIPv4+ fromString "*6" = HostIPv6+ fromString s = Host s+++-- | Extract the 'NS.HostName' from a 'Host' preference, or 'Nothing' otherwise.+hpHostName:: HostPreference -> Maybe NS.HostName+hpHostName (Host s) = Just s+hpHostName _ = Nothing++--------------------------------------------------------------------------------++-- |Exception thrown when a timeout has elapsed.+data Timeout+ = Timeout String -- ^Timeouted with an additional explanatory message.+ deriving (Eq, Show, Typeable)++instance E.Exception Timeout where
+ src/Control/Proxy/TCP.hs view
@@ -0,0 +1,385 @@+-- | This module exports functions that allow you to safely use 'NS.Socket'+-- resources acquired and released outside a 'P.Proxy' pipeline.+--+-- Instead, if want to safely acquire and release resources within the+-- pipeline itself, then you should use the functions exported by+-- "Control.Proxy.TCP.Safe".++-- Some code in this file was adapted from the @network-conduit@ library by+-- Michael Snoyman. Copyright (c) 2011. See its licensing terms (BSD3) at:+-- https://github.com/snoyberg/conduit/blob/master/network-conduit/LICENSE++module Control.Proxy.TCP (+ -- * Server side+ -- $server-side+ serve,+ serveFork,+ -- ** Listening+ listen,+ -- ** Accepting+ accept,+ acceptFork,++ -- * Client side+ -- $client-side+ connect,++ -- * Socket streams+ -- $socket-streaming+ socketReadS,+ nsocketReadS,+ socketWriteD,+ -- ** Timeouts+ -- $socket-streaming-timeout+ socketReadTimeoutS,+ nsocketReadTimeoutS,+ socketWriteTimeoutD,+ -- * Low level support+ bindSock,+ connectSock,+ -- * Exports+ HostPreference(..),+ Timeout(..)+ ) where++import Control.Concurrent (ThreadId, forkIO)+import qualified Control.Exception as E+import Control.Monad+import Control.Monad.Trans.Class+import qualified Control.Proxy as P+import qualified Control.Proxy.Trans.Either as PE+import Control.Proxy.Network.Internal+import qualified Data.ByteString as B+import Data.Monoid+import Data.List (partition)+import qualified Network.Socket as NS+import Network.Socket.ByteString (recv, sendAll)+import System.Timeout (timeout)++--------------------------------------------------------------------------------++-- $client-side+--+-- The following functions allow you to obtain and use 'NS.Socket's useful to+-- the client side of a TCP connection.+--+-- Here's how you could run a TCP client:+--+-- > connect "www.example.org" "80" $ \(connectionSocket, remoteAddr) -> do+-- > putStrLn $ "Connection established to " ++ show remoteAddr+-- > -- now you may use connectionSocket as you please within this scope,+-- > -- possibly with any of the socketReadS, nsocketReadS or socketWriteD+-- > -- proxies explained below.++-- | Connect to a TCP server and use the connection.+--+-- The connection socket is closed when done or in case of exceptions.+--+-- If you prefer to acquire and close the socket yourself, then use+-- 'connectSock' and the 'NS.sClose' function from "Network.Socket" instead.+connect+ :: NS.HostName -- ^Server hostname.+ -> NS.ServiceName -- ^Server service port.+ -> ((NS.Socket, NS.SockAddr) -> IO r)+ -- ^Computation taking the communication socket+ -- and the server address.+ -> IO r+connect host port = E.bracket (connectSock host port) (NS.sClose . fst)++--------------------------------------------------------------------------------++-- $server-side+--+-- The following functions allow you to obtain and use 'NS.Socket's useful to+-- the server side of a TCP connection.+--+-- Here's how you could run a TCP server that handles in different threads each+-- incoming connection to port @8000@ at address @127.0.0.1@:+--+-- > listen (Host "127.0.0.1") "8000" $ \(listeningSocket, listeningAddr) -> do+-- > putStrLn $ "Listening for incoming connections at " ++ show listeningAddr+-- > forever . acceptFork listeningSocket $ \(connectionSocket, remoteAddr) -> do+-- > putStrLn $ "Connection established from " ++ show remoteAddr+-- > -- now you may use connectionSocket as you please within this scope,+-- > -- possibly with any of the socketReadS, nsocketReadS or socketWriteD+-- > -- proxies explained below.+--+-- If you keep reading you'll discover there are different ways to achieve+-- the same, some ways more general than others. The above one was just an+-- example using a pretty general approach, you are encouraged to use simpler+-- approaches such as 'serve' if those suit your needs.++-- | Bind a TCP listening socket and use it.+--+-- The listening socket is closed when done or in case of exceptions.+--+-- If you prefer to acquire and close the socket yourself, then use+-- 'bindSock' and the 'NS.listen' and 'NS.sClose' functions from+-- "Network.Socket" instead.+--+-- Note: 'N.maxListenQueue' is tipically 128, which is too small for high+-- performance servers. So, we use the maximum between 'N.maxListenQueue' and+-- 2048 as the default size of the listening queue.+listen+ :: HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> ((NS.Socket, NS.SockAddr) -> IO r)+ -- ^Computation taking the listening socket and+ -- the address it's bound to.+ -> IO r+listen hp port = E.bracket listen' (NS.sClose . fst)+ where+ listen' = do x@(bsock,_) <- bindSock hp port+ NS.listen bsock $ max 2048 NS.maxListenQueue+ return x++-- | Start a TCP server that sequentially accepts and uses each incoming+-- connection.+--+-- Both the listening and connection sockets are closed when done or in case of+-- exceptions.+--+-- Note: You don't need to use 'listen' nor 'accept' if you use this function.+serve+ :: HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> ((NS.Socket, NS.SockAddr) -> IO r)+ -- ^Computation to run once an incoming+ -- connection is accepted. Takes the connection socket+ -- and remote end address.+ -> IO r+serve hp port k = do+ listen hp port $ \(lsock,_) -> do+ forever $ accept lsock k++-- | Start a TCP server that accepts incoming connections and uses them+-- concurrently in different threads.+--+-- The listening and connection sockets are closed when done or in case of+-- exceptions.+--+-- Note: You don't need to use 'listen' nor 'acceptFork' if you use this+-- function.+serveFork+ :: HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> ((NS.Socket, NS.SockAddr) -> IO ())+ -- ^Computation to run in a different thread+ -- once an incoming connection is accepted. Takes the+ -- connection socket and remote end address.+ -> IO ()+serveFork hp port k = do+ listen hp port $ \(lsock,_) -> do+ forever $ acceptFork lsock k++-- | Accept a single incoming connection and use it.+--+-- The connection socket is closed when done or in case of exceptions.+accept+ :: NS.Socket -- ^Listening and bound socket.+ -> ((NS.Socket, NS.SockAddr) -> IO b)+ -- ^Computation to run once an incoming+ -- connection is accepted. Takes the connection socket+ -- and remote end address.+ -> IO b+accept lsock k = do+ conn@(csock,_) <- NS.accept lsock+ E.finally (k conn) (NS.sClose csock)+{-# INLINABLE accept #-}++-- | Accept a single incoming connection and use it in a different thread.+--+-- The connection socket is closed when done or in case of exceptions.+acceptFork+ :: NS.Socket -- ^Listening and bound socket.+ -> ((NS.Socket, NS.SockAddr) -> IO ())+ -- ^Computation to run in a different thread+ -- once an incoming connection is accepted. Takes the+ -- connection socket and remote end address.+ -> IO ThreadId+acceptFork lsock f = do+ client@(csock,_) <- NS.accept lsock+ forkIO $ E.finally (f client) (NS.sClose csock)+{-# INLINABLE acceptFork #-}++--------------------------------------------------------------------------------++-- $socket-streaming+--+-- Once you have a connected 'NS.Socket', you can use the following 'P.Proxy's+-- to interact with the other connection end using streams.++-- | Receives bytes from the remote end sends them downstream.+--+-- Less than the specified maximum number of bytes might be received at once.+--+-- If the remote peer closes its side of the connection, this proxy returns.+socketReadS+ :: P.Proxy p+ => Int -- ^Maximum number of bytes to receive at once.+ -> NS.Socket -- ^Connected socket.+ -> () -> P.Producer p B.ByteString IO ()+socketReadS nbytes sock () = P.runIdentityP loop where+ loop = do+ bs <- lift $ recv sock nbytes+ unless (B.null bs) $ P.respond bs >> loop+{-# INLINABLE socketReadS #-}++-- | Just like 'socketReadS', except each request from downstream specifies the+-- maximum number of bytes to receive.+nsocketReadS+ :: P.Proxy p+ => NS.Socket -- ^Connected socket.+ -> Int -> P.Server p Int B.ByteString IO ()+nsocketReadS sock = P.runIdentityK loop where+ loop nbytes = do+ bs <- lift $ recv sock nbytes+ unless (B.null bs) $ P.respond bs >>= loop+{-# INLINABLE nsocketReadS #-}++-- | Sends to the remote end the bytes received from upstream, then forwards+-- such same bytes downstream.+--+-- Requests from downstream are forwarded upstream.+socketWriteD+ :: P.Proxy p+ => NS.Socket -- ^Connected socket.+ -> x -> p x B.ByteString x B.ByteString IO r+socketWriteD sock = P.runIdentityK loop where+ loop x = do+ a <- P.request x+ lift $ sendAll sock a+ P.respond a >>= loop+{-# INLINABLE socketWriteD #-}++--------------------------------------------------------------------------------++-- $socket-streaming-timeout+--+-- These proxies behave like the similarly named ones above, except support for+-- timing out the interaction with the remote end is added.++-- | Like 'socketReadS', except it throws a 'Timeout' exception in the+-- 'PE.EitherP' proxy transformer if receiving data from the remote end takes+-- more time than specified.+socketReadTimeoutS+ :: P.Proxy p+ => Int -- ^Timeout in microseconds (1/10^6 seconds).+ -> Int -- ^Maximum number of bytes to receive at once.+ -> NS.Socket -- ^Connected socket.+ -> () -> P.Producer (PE.EitherP Timeout p) B.ByteString IO ()+socketReadTimeoutS wait nbytes sock () = loop where+ loop = do+ mbs <- lift . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> PE.throw ex+ Just bs -> unless (B.null bs) $ P.respond bs >> loop+ ex = Timeout $ "recv: " <> show wait <> " microseconds."+{-# INLINABLE socketReadTimeoutS #-}++-- | Like 'nsocketReadS', except it throws a 'Timeout' exception in the+-- 'PE.EitherP' proxy transformer if receiving data from the remote end takes+-- more time than specified.+nsocketReadTimeoutS+ :: P.Proxy p+ => Int -- ^Timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> Int -> P.Server (PE.EitherP Timeout p) Int B.ByteString IO ()+nsocketReadTimeoutS wait sock = loop where+ loop nbytes = do+ mbs <- lift . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> PE.throw ex+ Just bs -> unless (B.null bs) $ P.respond bs >>= loop+ ex = Timeout $ "recv: " <> show wait <> " microseconds."+{-# INLINABLE nsocketReadTimeoutS #-}++-- | Like 'socketWriteD', except it throws a 'Timeout' exception in the+-- 'PE.EitherP' proxy transformer if sending data to the remote end takes+-- more time than specified.+socketWriteTimeoutD+ :: P.Proxy p+ => Int -- ^Timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> x -> (PE.EitherP Timeout p) x B.ByteString x B.ByteString IO r+socketWriteTimeoutD wait sock = loop where+ loop x = do+ a <- P.request x+ mbs <- lift . timeout wait $ sendAll sock a+ case mbs of+ Nothing -> PE.throw ex+ Just () -> P.respond a >>= loop+ ex = Timeout $ "recv: " <> show wait <> " microseconds."+{-# INLINABLE socketWriteTimeoutD #-}++--------------------------------------------------------------------------------++-- | Obtain a 'NS.Socket' connected to the given host and TCP service port.+--+-- The obtained 'NS.Socket' should be closed manually using 'NS.sClose' when+-- it's not needed anymore, otherwise you risk having the socket open for much+-- longer than needed.+--+-- Prefer to use 'connect' if you will be using the socket within a limited+-- scope and would like it to be closed immediately after its usage or in case+-- of exceptions.+connectSock :: NS.HostName -> NS.ServiceName -> IO (NS.Socket, NS.SockAddr)+connectSock host port = do+ (addr:_) <- NS.getAddrInfo (Just hints) (Just host) (Just port)+ E.bracketOnError (newSocket addr) NS.sClose $ \sock -> do+ let sockAddr = NS.addrAddress addr+ NS.connect sock sockAddr+ return (sock, sockAddr)+ where+ hints = NS.defaultHints { NS.addrFlags = [NS.AI_ADDRCONFIG]+ , NS.addrSocketType = NS.Stream }++-- | Obtain a 'NS.Socket' bound to the given host name and TCP service port.+--+-- The obtained 'NS.Socket' should be closed manually using 'NS.sClose' when+-- it's not needed anymore.+--+-- Prefer to use 'listen' if you will be listening on this socket and using it+-- within a limited scope, and would like it to be closed immediately after its+-- usage or in case of exceptions.+bindSock :: HostPreference -> NS.ServiceName -> IO (NS.Socket, NS.SockAddr)+bindSock hp port = do+ addrs <- NS.getAddrInfo (Just hints) (hpHostName hp) (Just port)+ let addrs' = case hp of+ HostIPv4 -> prioritize isIPv4addr addrs+ HostIPv6 -> prioritize isIPv6addr addrs+ _ -> addrs+ tryAddrs addrs'+ where+ hints = NS.defaultHints { NS.addrFlags = [NS.AI_PASSIVE]+ , NS.addrSocketType = NS.Stream }++ tryAddrs [] = error "listen: no addresses available"+ tryAddrs [x] = useAddr x+ tryAddrs (x:xs) = E.catch (useAddr x)+ (\e -> let _ = e :: E.IOException in tryAddrs xs)++ useAddr addr = E.bracketOnError (newSocket addr) NS.sClose $ \sock -> do+ let sockAddr = NS.addrAddress addr+ NS.setSocketOption sock NS.NoDelay 1+ NS.setSocketOption sock NS.ReuseAddr 1+ NS.bindSocket sock sockAddr+ return (sock, sockAddr)++--------------------------------------------------------------------------------++-- Misc++newSocket :: NS.AddrInfo -> IO NS.Socket+newSocket addr = NS.socket (NS.addrFamily addr)+ (NS.addrSocketType addr)+ (NS.addrProtocol addr)++isIPv4addr, isIPv6addr :: NS.AddrInfo -> Bool+isIPv4addr x = NS.addrFamily x == NS.AF_INET+isIPv6addr x = NS.addrFamily x == NS.AF_INET6++-- | Move the elements that match the predicate closer to the head of the list.+-- Preserve relative order.+prioritize :: (a -> Bool) -> [a] -> [a]+prioritize p = uncurry (++) . partition p
+ src/Control/Proxy/TCP/Safe.hs view
@@ -0,0 +1,435 @@+{-# LANGUAGE Rank2Types #-}++-- | This module exports functions that allow you to safely use 'NS.Socket'+-- resources within a 'P.Proxy' pipeline, possibly acquiring and releasing such+-- resources within the pipeline itself, using the facilities provided by+-- 'P.ExceptionP' from the @pipes-safe@ library.+--+-- Instead, if just want to use resources already acquired or released outside+-- the pipeline, then you could use the simpler functions exported by+-- "Control.Proxy.TCP".++module Control.Proxy.TCP.Safe (+ -- * Server side+ -- $server-side+ serve,+ serveFork,+ -- ** Listening+ listen,+ -- ** Accepting+ accept,+ acceptFork,+ -- ** Streaming+ -- $server-streaming+ serveReadS,+ serveWriteD,++ -- * Client side+ -- $client-side+ connect,+ -- ** Streaming+ -- $client-streaming+ connectReadS,+ connectWriteD,++ -- * Socket streams+ -- $socket-streaming+ socketReadS,+ nsocketReadS,+ socketWriteD,++ -- * Exports+ HostPreference(..),+ Timeout(..)+ ) where++import Control.Concurrent (forkIO, ThreadId)+import qualified Control.Exception as E+import Control.Monad+import qualified Control.Proxy as P+import Control.Proxy.Network.Internal+import qualified Control.Proxy.TCP as T+import qualified Control.Proxy.Safe as P+import qualified Data.ByteString as B+import Data.Monoid+import qualified Network.Socket as NS+import Network.Socket.ByteString (sendAll, recv)+import System.Timeout (timeout)++--------------------------------------------------------------------------------++-- $client-side+--+-- The following functions allow you to obtain and use 'NS.Socket's useful to+-- the client side of a TCP connection.+--+-- Here's how you could run a TCP client:+--+-- > connect id "www.example.org" "80" $ \(connectionSocket, remoteAddr) -> do+-- > tryIO . putStrLn $ "Connection established to " ++ show remoteAddr+-- > -- now you may use connectionSocket as you please within this scope,+-- > -- possibly with any of the socketReadS, nsocketReadS or socketWriteD+-- > -- proxies explained below.+--+-- You might instead prefer the simpler but less general solutions offered by+-- 'connectReadS' and 'connectWriteD', so check those too.++-- | Connect to a TCP server and use the connection.+--+-- The connection socket is closed when done or in case of exceptions.+--+-- If you prefer to acquire close the socket yourself, then use+-- 'T.connectSock' and the 'NS.sClose' from "Network.Socket" instead.+connect+ :: (P.Proxy p, Monad m)+ => (forall x. P.SafeIO x -> m x) -- ^Monad morphism.+ -> NS.HostName -- ^Server hostname.+ -> NS.ServiceName -- ^Server service port.+ -> ((NS.Socket, NS.SockAddr) -> P.ExceptionP p a' a b' b m r)+ -- ^Computation taking the+ -- communication socket and the server+ -- address.+ -> P.ExceptionP p a' a b' b m r+connect morph host port =+ P.bracket morph (T.connectSock host port) (NS.sClose . fst)++--------------------------------------------------------------------------------++-- $client-streaming+--+-- The following proxies allow you to easily connect to a TCP server and+-- immediately interact with it using streams, all at once, instead of+-- having to perform the individual steps separately.++-- | Connect to a TCP server and send downstream the bytes received from the+-- remote end.+--+-- If an optional timeout is given and receiveing data from the remote end takes+-- more time that such timeout, then throw a 'Timeout' exception in the+-- 'P.ExceptionP' proxy transformer.+--+-- The connection socket is closed when done or in case of exceptions.+--+-- Using this proxy you can write straightforward code like the following, which+-- prints whatever is received from a single TCP connection to a given server+-- listening locally on port 9000:+--+-- >>> runSafeIO . runProxy . runEitherK $ connectReadS Nothing 4096 "127.0.0.1" "9000" >-> tryK printD+connectReadS+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> Int -- ^Maximum number of bytes to receive at once.+ -> NS.HostName -- ^Server host name.+ -> NS.ServiceName -- ^Server service port.+ -> () -> P.Producer (P.ExceptionP p) B.ByteString P.SafeIO ()+connectReadS mwait nbytes host port () = do+ connect id host port $ \(csock,_) -> do+ socketReadS mwait nbytes csock ()++-- | Connects to a TCP server, sends to the remote end the bytes received from+-- upstream, then forwards such same bytes downstream.+--+-- Requests from downstream are forwarded upstream.+--+-- If an optional timeout is given and sending data to the remote end takes+-- more time that such timeout, then throw a 'Timeout' exception in the+-- 'P.ExceptionP' proxy transformer.+--+-- The connection socket is closed when done or in case of exceptions.+--+-- Using this proxy you can write straightforward code like the following, which+-- greets a TCP client listening locally at port 9000:+--+-- >>> :set -XOverloadedStrings+-- >>> runSafeIO . runProxy . runEitherK $ fromListS ["He","llo\r\n"] >-> connectWriteD Nothing "127.0.0.1" "9000"+connectWriteD+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> NS.HostName -- ^Server host name.+ -> NS.ServiceName -- ^Server service port.+ -> x -> (P.ExceptionP p) x B.ByteString x B.ByteString P.SafeIO ()+connectWriteD mwait hp port x = do+ connect id hp port $ \(csock,_) ->+ socketWriteD mwait csock x++--------------------------------------------------------------------------------++-- $server-side+--+-- The following functions allow you to obtain and use 'NS.Socket's useful to+-- the server side of a TCP connection.+--+-- Here's how you could run a TCP server that handles in different threads each+-- incoming connection to port @8000@ at address @127.0.0.1@:+--+-- > listen id (Host "127.0.0.1") "8000" $ \(listeningSocket, listeningAddr) -> do+-- > tryIO . putStrLn $ "Listening for incoming connections at " ++ show listeningAddr+-- > forever . acceptFork id listeningSocket $ \(connectionSocket, remoteAddr) -> do+-- > putStrLn $ "Connection established from " ++ show remoteAddr+-- > -- now you may use connectionSocket as you please within this scope,+-- > -- possibly with any of the socketReadS, nsocketReadS or socketWriteD+-- > -- proxies explained below.+--+-- If you keep reading you'll discover there are different ways to achieve+-- the same, some ways more general than others. The above one was just an+-- example using a pretty general approach, you are encouraged to use simpler+-- approaches such as 'serve' or 'serveReadS' if those suit your needs.++-- | Bind a TCP listening socket and use it.+--+-- The listening socket is closed when done or in case of exceptions.+--+-- If you prefer to acquire and close the socket yourself, then use+-- 'T.bindSock' and the 'NS.listen' and 'NS.sClose' functions from+-- "Network.Socket" instead.+--+-- Note: 'N.maxListenQueue' is tipically 128, which is too small for high+-- performance servers. So, we use the maximum between 'N.maxListenQueue' and+-- 2048 as the default size of the listening queue.+listen+ :: (P.Proxy p, Monad m)+ => (forall x. P.SafeIO x -> m x) -- ^Monad morphism.+ -> HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> ((NS.Socket, NS.SockAddr) -> P.ExceptionP p a' a b' b m r)+ -- ^Computation taking the listening+ -- socket and the address it's bound to.+ -> P.ExceptionP p a' a b' b m r+listen morph hp port = P.bracket morph listen' (NS.sClose . fst)+ where+ listen' = do x@(bsock,_) <- T.bindSock hp port+ NS.listen bsock $ max 2048 NS.maxListenQueue+ return x++-- | Start a TCP server that sequentially accepts and uses each incoming+-- connection.+--+-- Both the listening and connection sockets are closed when done or in case of+-- exceptions.+--+-- Note: You don't need to use 'listen' nor 'accept' if you use this function.+serve+ :: (P.Proxy p, Monad m)+ => (forall x. P.SafeIO x -> m x) -- ^Monad morphism.+ -> HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> ((NS.Socket, NS.SockAddr) -> P.ExceptionP p a' a b' b m r)+ -- ^Computation to run once an incoming+ -- connection is accepted. Takes the+ -- connection socket and remote end address.+ -> P.ExceptionP p a' a b' b m r+serve morph hp port k = do+ listen morph hp port $ \(lsock,_) -> do+ forever $ accept morph lsock k++-- | Start a TCP server that accepts incoming connections and uses them+-- concurrently in different threads.+--+-- The listening and connection sockets are closed when done or in case of+-- exceptions.+--+-- Note: You don't need to use 'listen' nor 'acceptFork' if you use this+-- function.+serveFork+ :: (P.Proxy p, Monad m)+ => (forall x. P.SafeIO x -> m x) -- ^Monad morphism.+ -> HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> ((NS.Socket, NS.SockAddr) -> IO ())+ -- ^Computation to run in a different thread+ -- once an incoming connection is accepted.+ -- Takes the connection socket and remote end+ -- address.+ -> P.ExceptionP p a' a b' b m r+serveFork morph hp port k = do+ listen morph hp port $ \(lsock,_) -> do+ forever $ acceptFork morph lsock k++-- | Accept a single incoming connection and use it.+--+-- The connection socket is closed when done or in case of exceptions.+accept+ :: (P.Proxy p, Monad m)+ => (forall x. P.SafeIO x -> m x) -- ^Monad morphism.+ -> NS.Socket -- ^Listening and bound socket.+ -> ((NS.Socket, NS.SockAddr) -> P.ExceptionP p a' a b' b m r)+ -- ^Computation to run once an incoming+ -- connection is accepted. Takes the+ -- connection socket and remote end address.+ -> P.ExceptionP p a' a b' b m r+accept morph lsock k = do+ conn@(csock,_) <- P.hoist morph . P.tryIO $ NS.accept lsock+ P.finally morph (NS.sClose csock) (k conn)+{-# INLINABLE accept #-}++-- | Accept a single incoming connection and use it in a different thread.+--+-- The connection socket is closed when done or in case of exceptions.+acceptFork+ :: (P.Proxy p, Monad m)+ => (forall x. P.SafeIO x -> m x) -- ^Monad morphism.+ -> NS.Socket -- ^Listening and bound socket.+ -> ((NS.Socket, NS.SockAddr) -> IO ())+ -- ^Computation to run in a different thread+ -- once an incoming connection is accepted.+ -- Takes the connection socket and remote end+ -- address.+ -> P.ExceptionP p a' a b' b m ThreadId+acceptFork morph lsock f = P.hoist morph . P.tryIO $ do+ client@(csock,_) <- NS.accept lsock+ forkIO $ E.finally (f client) (NS.sClose csock)+{-# INLINABLE acceptFork #-}++--------------------------------------------------------------------------------++-- $server-streaming+--+-- The following proxies allow you to easily run a TCP server and immediately+-- interact with incoming connections using streams, all at once, instead of+-- having to perform the individual steps separately.++-- | Binds a listening socket, accepts a single connection and sends downstream+-- any bytes received from the remote end.+--+-- If an optional timeout is given and receiveing data from the remote end takes+-- more time that such timeout, then throw a 'Timeout' exception in the+-- 'P.ExceptionP' proxy transformer.+--+-- Less than the specified maximum number of bytes might be received at once.+--+-- If the remote peer closes its side of the connection, this proxy returns.+--+-- Both the listening and connection sockets are closed when done or in case of+-- exceptions.+--+-- Using this proxy you can write straightforward code like the following, which+-- prints whatever is received from a single TCP connection to port 9000:+--+-- >>> :set -XOverloadedStrings+-- >>> runSafeIO . runProxy . runEitherK $ serveReadS Nothing 4096 "127.0.0.1" "9000" >-> tryK printD+serveReadS+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> Int -- ^Maximum number of bytes to receive at once.+ -> HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> () -> P.Producer (P.ExceptionP p) B.ByteString P.SafeIO ()+serveReadS mwait nbytes hp port () = do+ listen id hp port $ \(lsock,_) -> do+ accept id lsock $ \(csock,_) -> do+ socketReadS mwait nbytes csock ()++-- | Binds a listening socket, accepts a single connection, sends to the remote+-- end the bytes received from upstream, then forwards such sames bytes+-- downstream.+--+-- Requests from downstream are forwarded upstream.+--+-- If an optional timeout is given and sending data to the remote end takes+-- more time that such timeout, then throw a 'Timeout' exception in the+-- 'P.ExceptionP' proxy transformer.+--+-- Both the listening and connection sockets are closed when done or in case of+-- exceptions.+--+-- Using this proxy you can write straightforward code like the following, which+-- greets a TCP client connecting to port 9000:+--+-- >>> :set -XOverloadedStrings+-- >>> runSafeIO . runProxy . runEitherK $ fromListS ["He","llo\r\n"] >-> serveWriteD Nothing "127.0.0.1" "9000"+serveWriteD+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> HostPreference -- ^Preferred host to bind.+ -> NS.ServiceName -- ^Service port to bind.+ -> x -> (P.ExceptionP p) x B.ByteString x B.ByteString P.SafeIO ()+serveWriteD mwait hp port x = do+ listen id hp port $ \(lsock,_) -> do+ accept id lsock $ \(csock,_) -> do+ socketWriteD mwait csock x++--------------------------------------------------------------------------------++-- $socket-streaming+--+-- Once you have a connected 'NS.Socket', you can use the following 'P.Proxy's+-- to interact with the other connection end using streams.++-- | Receives bytes from the remote end and sends them downstream.+--+-- If an optional timeout is given and receiveing data from the remote end takes+-- more time that such timeout, then throw a 'Timeout' exception in the+-- 'P.ExceptionP' proxy transformer.+--+-- Less than the specified maximum number of bytes might be received at once.+--+-- If the remote peer closes its side of the connection, this proxy returns.+socketReadS+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> Int -- ^Maximum number of bytes to receive at once.+ -> NS.Socket -- ^Connected socket.+ -> () -> P.Producer (P.ExceptionP p) B.ByteString P.SafeIO ()+socketReadS Nothing nbytes sock () = loop where+ loop = do+ bs <- P.tryIO $ recv sock nbytes+ unless (B.null bs) $ P.respond bs >> loop+socketReadS (Just wait) nbytes sock () = loop where+ loop = do+ mbs <- P.tryIO . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> P.throw ex+ Just bs -> unless (B.null bs) $ P.respond bs >> loop+ ex = Timeout $ "recv: " <> show wait <> " microseconds."+{-# INLINABLE socketReadS #-}++-- | Just like 'socketReadS', except each request from downstream specifies the+-- maximum number of bytes to receive.+nsocketReadS+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> Int -> P.Server (P.ExceptionP p) Int B.ByteString P.SafeIO ()+nsocketReadS Nothing sock = loop where+ loop nbytes = do+ bs <- P.tryIO $ recv sock nbytes+ unless (B.null bs) $ P.respond bs >>= loop+nsocketReadS (Just wait) sock = loop where+ loop nbytes = do+ mbs <- P.tryIO . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> P.throw ex+ Just bs -> unless (B.null bs) $ P.respond bs >>= loop+ ex = Timeout $ "recv: " <> show wait <> " microseconds."+{-# INLINABLE nsocketReadS #-}++-- | Sends to the remote end the bytes received from upstream, then forwards+-- such same bytes downstream.+--+-- If an optional timeout is given and sending data to the remote end takes+-- more time that such timeout, then throw a 'Timeout' exception in the+-- 'P.ExceptionP' proxy transformer.+--+-- Requests from downstream are forwarded upstream.+socketWriteD+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> x -> (P.ExceptionP p) x B.ByteString x B.ByteString P.SafeIO r+socketWriteD Nothing sock = loop where+ loop x = do+ a <- P.request x+ P.tryIO $ sendAll sock a+ P.respond a >>= loop+socketWriteD (Just wait) sock = loop where+ loop x = do+ a <- P.request x+ m <- P.tryIO . timeout wait $ sendAll sock a+ case m of+ Nothing -> P.throw ex+ Just () -> P.respond a >>= loop+ ex = Timeout $ "sendAll: " <> show wait <> " microseconds."+{-# INLINABLE socketWriteD #-}+++
+ src/Control/Proxy/TCP/Safe/Sync.hs view
@@ -0,0 +1,121 @@+-- | This module exports 'P.Proxy's that allow implementing synchronous RPC-like+-- communication with a remote end by using a simple protocol on their+-- downstream interface.+--+-- As opposed to the similar proxies found in "Control.Proxy.TCP.Sync",+-- these use the exception handling facilities provided by 'P.ExceptionP'.+--+-- You may prefer the more general proxies from+-- "Control.Proxy.TCP.Safe".++module Control.Proxy.TCP.Safe.Sync (+ -- * Socket proxies+ socketSyncServer,+ socketSyncProxy,+ -- * RPC support+ syncDelimit,+ -- * Protocol+ Request(..),+ Response(..),+ ) where++import Control.Monad+import qualified Control.Proxy as P+import Control.Proxy.TCP.Sync+ (Request(..), Response(..), syncDelimit)+import Control.Proxy.Network.Internal+import qualified Control.Proxy.Safe as P+import qualified Data.ByteString as B+import Data.Monoid+import qualified Network.Socket as NS+import Network.Socket.ByteString (recv, sendAll)+import System.Timeout (timeout)+++-- | 'P.Server' able to send and receive bytes through a 'NS.Socket'.+--+-- If downstream requests @'Send' bytes@, then such @bytes@ are sent to the+-- remote end and then this proxy responds 'Sent' downstream.+--+-- If downstream requests @'Receive' num@, then at most @num@ bytes are received+-- from the remote end. This proxy then responds downstream such received+-- bytes as @'Received' bytes@. Less than the specified maximum number of bytes+-- might be received at once.+--+-- If an optional timeout is given and interactions with the remote end+-- take more time that such timeout, then throw a 'Timeout' exception in+-- the 'P.ExceptionP' proxy transformer.+--+-- If the remote peer closes its side of the connection, this proxy returns.+socketSyncServer+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> Request B.ByteString+ -> P.Server (P.ExceptionP p) (Request B.ByteString) Response P.SafeIO()+socketSyncServer Nothing sock = loop where+ loop (Send bs) = do+ P.tryIO $ sendAll sock bs+ P.respond Sent >>= loop+ loop (Receive nbytes) = do+ bs <- P.tryIO $ recv sock nbytes+ unless (B.null bs) $ P.respond (Received bs) >>= loop+socketSyncServer (Just wait) sock = loop where+ loop (Send bs) = do+ m <- P.tryIO . timeout wait $ sendAll sock bs+ case m of+ Nothing -> P.throw $ ex "sendAll"+ Just () -> P.respond Sent >>= loop+ loop (Receive nbytes) = do+ mbs <- P.tryIO . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> P.throw $ ex "recv"+ Just bs -> unless (B.null bs) $ P.respond (Received bs) >>= loop+ ex s = Timeout $ s <> ": " <> show wait <> " microseconds."+{-# INLINABLE socketSyncServer #-}++-- | 'P.Proxy' able to send and receive bytes through a 'NS.Socket'.+--+-- If downstream requests @'Send' a'@, then such @a'@ request is forwarded+-- upstream, which in return responds a 'B.ByteString' that this proxy sends to+-- the remote end. After sending to the remote end, this proxy responds 'Sent'+-- downstream.+--+-- If downstream requests @'Receive' num@, then at most @num@ bytes are received+-- from the remote end. This proxy then responds downstream such received+-- bytes as @'Received' bytes@. Less than the specified maximum number of bytes+-- might be received at once.+--+-- If an optional timeout is given and interactions with the remote end+-- take more time that such timeout, then throw a 'Timeout' exception in+-- the 'P.ExceptionP' proxy transformer.+--+-- If the remote peer closes its side of the connection, this proxy returns.+socketSyncProxy+ :: P.Proxy p+ => Maybe Int -- ^Optional timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> Request a'+ -> (P.ExceptionP p) a' B.ByteString (Request a') Response P.SafeIO ()+socketSyncProxy Nothing sock = loop where+ loop (Send a') = do+ P.request a' >>= P.tryIO . sendAll sock+ P.respond Sent >>= loop+ loop (Receive nbytes) = do+ bs <- P.tryIO $ recv sock nbytes+ unless (B.null bs) $ P.respond (Received bs) >>= loop+socketSyncProxy (Just wait) sock = loop where+ loop (Send a') = do+ bs <- P.request a'+ m <- P.tryIO . timeout wait $ sendAll sock bs+ case m of+ Nothing -> P.throw $ ex "sendAll"+ Just () -> P.respond Sent >>= loop+ loop (Receive nbytes) = do+ mbs <- P.tryIO . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> P.throw $ ex "recv"+ Just bs -> unless (B.null bs) $ P.respond (Received bs) >>= loop+ ex s = Timeout $ s <> ": " <> show wait <> " microseconds."+{-# INLINABLE socketSyncProxy #-}+
+ src/Control/Proxy/TCP/Sync.hs view
@@ -0,0 +1,203 @@+{-# LANGUAGE OverloadedStrings #-}++-- | This module exports 'P.Proxy's that allow implementing synchronous RPC-like+-- communication with a remote end by using a simple protocol on their+-- downstream interface.+--+-- As opposed to the similar proxies found in+-- "Control.Proxy.TCP.Safe.Sync", these don't use the exception handling+-- facilities provided by 'P.ExceptionP'.+--+-- You may prefer the more general and efficient proxies from+-- "Control.Proxy.TCP".++module Control.Proxy.TCP.Sync (+ -- * Socket proxies+ socketSyncServer,+ socketSyncProxy,+ -- ** Timeouts+ -- $timeouts+ socketSyncServerTimeout,+ socketSyncProxyTimeout,+ -- * RPC support+ syncDelimit,+ -- * Protocol types+ Request(..),+ Response(..),+ ) where++import Control.Monad+import Control.Monad.Trans.Class+import qualified Control.Proxy as P+import Control.Proxy.Network.Internal+import qualified Control.Proxy.Trans.Either as PE+import qualified Data.ByteString.Char8 as B+import Data.Monoid+import qualified Network.Socket as NS+import Network.Socket.ByteString (recv, sendAll)+import System.Timeout (timeout)+++-- | A request made to one of the @socketSync*@ proxies.+data Request t = Send t | Receive Int+ deriving (Eq, Read, Show)++-- | A response received from one of the @socketSync*@ proxies.+data Response = Sent | Received B.ByteString+ deriving (Eq, Read, Show)++--------------------------------------------------------------------------------++-- | 'P.Server' able to send and receive bytes through a 'NS.Socket'.+--+-- If downstream requests @'Send' bytes@, then such @bytes@ are sent to the+-- remote end and then this proxy responds 'Sent' downstream.+--+-- If downstream requests @'Receive' num@, then at most @num@ bytes are received+-- from the remote end. This proxy then responds downstream such received+-- bytes as @'Received' bytes@. Less than the specified maximum number of bytes+-- might be received at once.+--+-- If the remote peer closes its side of the connection, this proxy returns.+socketSyncServer+ :: P.Proxy p+ => NS.Socket -- ^Connected socket.+ -> Request B.ByteString+ -> P.Server p (Request B.ByteString) Response IO ()+socketSyncServer sock = P.runIdentityK loop where+ loop (Send bs) = do+ lift $ sendAll sock bs+ P.respond Sent >>= loop+ loop (Receive nbytes) = do+ bs <- lift $ recv sock nbytes+ unless (B.null bs) $ P.respond (Received bs) >>= loop+{-# INLINABLE socketSyncServer #-}++-- | 'P.Proxy' able to send and receive bytes through a 'NS.Socket'.+--+-- If downstream requests @'Send' a'@, then such @a'@ request is forwarded+-- upstream, which in return responds a 'B.ByteString' that this proxy sends to+-- the remote end. After sending to the remote end, this proxy responds 'Sent'+-- downstream.+--+-- If downstream requests @'Receive' num@, then at most @num@ bytes are received+-- from the remote end. This proxy then responds downstream such received+-- bytes as @'Received' bytes@. Less than the specified maximum number of bytes+-- might be received at once.+--+-- If the remote peer closes its side of the connection, this proxy returns.+socketSyncProxy+ :: P.Proxy p+ => NS.Socket -- ^Connected socket.+ -> Request a'+ -> p a' B.ByteString (Request a') Response IO ()+socketSyncProxy sock = P.runIdentityK loop where+ loop (Send a') = do+ P.request a' >>= lift . sendAll sock+ P.respond Sent >>= loop+ loop (Receive nbytes) = do+ bs <- lift $ recv sock nbytes+ unless (B.null bs) $ P.respond (Received bs) >>= loop+{-# INLINABLE socketSyncProxy #-}++--------------------------------------------------------------------------------++-- $timeouts+--+-- These proxies behave like the similarly named ones above, except support for+-- timing out the interaction with the remote end is added.++-- | Like 'socketSyncServer', except it throws a 'Timeout' exception in the+-- 'PE.EitherP' proxy transformer if interacting with the remote end takes+-- more time than specified.+socketSyncServerTimeout+ :: P.Proxy p+ => Int -- ^Timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> Request B.ByteString+ -> P.Server (PE.EitherP Timeout p) (Request B.ByteString) Response IO ()+socketSyncServerTimeout wait sock = loop where+ loop (Send bs) = do+ m <- lift . timeout wait $ sendAll sock bs+ case m of+ Nothing -> PE.throw $ ex "sendAll"+ Just () -> P.respond Sent >>= loop+ loop (Receive nbytes) = do+ mbs <- lift . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> PE.throw $ ex "recv"+ Just bs -> unless (B.null bs) $ P.respond (Received bs) >>= loop+ ex s = Timeout $ s <> ": " <> show wait <> " microseconds."+{-# INLINABLE socketSyncServerTimeout #-}++-- | Like 'socketSyncProxy', except it throws a 'Timeout' exception in the+-- 'PE.EitherP' proxy transformer if interacting with the remote end takes+-- more time than specified.+socketSyncProxyTimeout+ :: P.Proxy p+ => Int -- ^Timeout in microseconds (1/10^6 seconds).+ -> NS.Socket -- ^Connected socket.+ -> Request a'+ -> (PE.EitherP Timeout p) a' B.ByteString (Request a') Response IO ()+socketSyncProxyTimeout wait sock = loop where+ loop (Send a') = do+ bs <- P.request a'+ m <- lift . timeout wait $ sendAll sock bs+ case m of+ Nothing -> PE.throw $ ex "sendAll"+ Just () -> P.respond Sent >>= loop+ loop (Receive nbytes) = do+ mbs <- lift . timeout wait $ recv sock nbytes+ case mbs of+ Nothing -> PE.throw $ ex "recv"+ Just bs -> unless (B.null bs) $ P.respond (Received bs) >>= loop+ ex s = Timeout $ s <> ": " <> show wait <> " microseconds."+{-# INLINABLE socketSyncProxyTimeout #-}++--------------------------------------------------------------------------------++-- | When used together with one of the @socketSync*@ proxies upstream, this+-- proxy sends a single 'B.ByteString' to the remote end and then repeatedly+-- receives bytes from the remote end until the given delimiter is found.+-- Finally, a single 'B.ByteString' up to the given delimiter (inclusive) is+-- sent downstream and then the whole process is repeated.+--+-- This proxy works cooperatively with any @socketSync*@ proxy immediately+-- upstream, so read their documentation to understand the purpose of the+-- @b'@ value received from downstream.+--+-- For example, if you'd like to convert a 'NS.Socket' into an synchronous+-- line-oriented RPC client implemented as a 'P.Server' in which RPC calls are+-- received via the downstream interface and RPC responses are sent downstream,+-- then you could use this proxy as:+--+-- > socketSyncServer ... >-> syncDelimit 4096 "\r\n"+--+-- Otherwise, if you'd like to convert a 'NS.Socket' into an synchronous+-- line-oriented RPC client implemented as a 'P.Proxy' in which RPC calls are+-- received via the upstream interface and RPC responses are sent downstream,+-- then you could use this proxy as:+--+-- > socketSyncProxy ... >-> syncDelimit 4096 "\r\n"+syncDelimit+ :: (Monad m, P.Proxy p)+ => Int -- ^Maximum number of bytes to receive at once.+ -> B.ByteString -- ^Delimiting bytes.+ -> b'-> p (Request b') Response b' B.ByteString m r+syncDelimit nbytes delim b' =+ -- XXX this implementation might be inefficient.+ P.runIdentityP $ use =<< more mempty (Send b')+ where+ more buf req = do+ a <- P.request req+ case a of+ Received bs -> return (buf <> bs)+ Sent -> more buf (Receive nbytes)+ use buf = do+ let (pre,suf) = B.breakSubstring delim buf+ case B.length suf of+ 0 -> use =<< more buf (Receive nbytes)+ _ -> do b'2 <- P.respond (pre <> delim)+ use =<< more (B.drop (B.length delim) suf) (Send b'2)+{-# INLINABLE syncDelimit #-}+
+ tests/Simple.hs view
@@ -0,0 +1,146 @@+{-# LANGUAGE OverloadedStrings #-}+{-# OPTIONS_GHC -fno-warn-unused-do-bind -fno-warn-missing-signatures #-}++module Main where++import Control.Concurrent (forkIO, threadDelay)+import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)+import qualified Control.Exception as E+import qualified Data.ByteString.Char8 as B+import qualified Network.Socket as NS+import Test.Framework (Test, defaultMain, testGroup)+import Test.Framework.Providers.HUnit (testCase)+import Test.HUnit (Assertion, (@=?))+import Control.Proxy ((>->))+import qualified Control.Proxy as P+import qualified Control.Proxy.Safe as P+import qualified Control.Proxy.TCP as T+import qualified Control.Proxy.TCP.Safe as T'++host1 = "127.0.0.1" :: NS.HostName+host1p = T.Host host1 :: T.HostPreference+ports = fmap show [14000..14010] :: [NS.ServiceName]+msg1 = take 1000 $ cycle ["Hell","o\r\n"] :: [B.ByteString]+msg1b = B.concat msg1 :: B.ByteString+++-- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- --+-- The following 4 IO actions are used throughout the various tests as the+-- default implementations for reading/writing a server/client. They themselves+-- are also tested below.++-- tested by 'test_listen_accept_socketWriteD_then_connect_socketReadD'+connectAndRead :: NS.HostName -> NS.ServiceName -> IO [B.ByteString]+connectAndRead host port = do+ T.connect host port $ \(csock, _caddr) -> do+ let p = P.raiseK (T.socketReadS 4096 csock) >-> P.toListD+ fmap snd $ P.runWriterT . P.runProxy $ p+-- let p = P.raiseK (T'.connectReadS Nothing 4096 host port) >-> P.toListD+-- (eex,out) <- P.trySafeIO . P.runWriterT . P.runProxy .P.runEitherK $ p+-- case eex of+-- Left ex -> E.throw ex+-- Right () -> return out++-- tested by 'test_listen_accept_socketReadS_then_connect_socketWriteD'+connectAndWrite :: NS.HostName -> NS.ServiceName -> [B.ByteString] -> IO ()+connectAndWrite host port msg = do+ T.connect host port $ \(csock, _caddr) -> do+ P.runProxy $ P.fromListS msg >-> T.socketWriteD csock+-- let p = P.fromListS msg >-> T'.connectWriteD Nothing host1 port+-- P.runSafeIO . P.runProxy .P.runEitherK $ p++-- tested by 'test_listen_accept_socketWriteD_then_connect_socketReadD'+serveOnceAndRead :: T.HostPreference -> NS.ServiceName -> IO [B.ByteString]+serveOnceAndRead hp port = do+ T.listen hp port $ \(lsock, _laddr) -> do+ T.accept lsock $ \(csock, _caddr) -> do+ let p = P.raiseK (T.socketReadS 4096 csock) >-> P.toListD+ fmap snd $ P.runWriterT . P.runProxy $ p+-- let p = P.raiseK (T'.serveReadS Nothing 4096 host1p port) >-> P.toListD+-- (eex,out) <- P.trySafeIO . P.runWriterT . P.runProxy .P.runEitherK $ p+-- case eex of+-- Left ex -> E.throw ex+-- Right () -> return out++-- tested by 'test_listen_accept_socketWriteD_then_connect_socketReadD'+serveOnceAndWrite :: T.HostPreference -> NS.ServiceName -> [B.ByteString] -> IO ()+serveOnceAndWrite hp port msg = do+ T.listen hp port $ \(lsock, _laddr) -> do+ T.accept lsock $ \(csock, _caddr) -> do+ P.runProxy $ P.fromListS msg >-> T.socketWriteD csock+-- let p = P.fromListS msg >-> T'.serveWriteD Nothing host1p port+-- P.runSafeIO . P.runProxy .P.runEitherK $ p+-- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- -- --++-- Note: In all the tests below we wait a bit before starting the+-- client, hoping that by then the server has already started.+-- Yes, I know, it's not the best approach. Hopefully it will be enough.+waitTime :: Int -- in microseconds (1e6)+waitTime = 200000++test_listen_accept_socketReadS_then_connect_socketWriteD :: Assertion+test_listen_accept_socketReadS_then_connect_socketWriteD = do+ let port = ports !! 0+ mvout <- newEmptyMVar+ forkIO $ putMVar mvout =<< serveOnceAndRead host1p port+ threadDelay waitTime+ connectAndWrite host1 port msg1+ out <- takeMVar mvout+ B.concat out @=? msg1b++test_listen_accept_socketWriteD_then_connect_socketReadD :: Assertion+test_listen_accept_socketWriteD_then_connect_socketReadD = do+ let port = ports !! 1+ forkIO $ serveOnceAndWrite host1p port msg1+ threadDelay waitTime+ out <- connectAndRead host1 port+ B.concat out @=? msg1b++test_safe_serveWriteD :: Assertion+test_safe_serveWriteD = do+ let port = ports !! 2+ serveOnceAndWrite' = do+ let p = P.fromListS msg1 >-> T'.serveWriteD Nothing host1p port+ P.runSafeIO . P.runProxy .P.runEitherK $ p+ forkIO serveOnceAndWrite'+ threadDelay waitTime+ out <- connectAndRead host1 port+ B.concat out @=? msg1b++test_safe_serveReadS :: Assertion+test_safe_serveReadS = do+ let port = ports !! 3+ serveOnceAndRead' = do+ let p = P.raiseK (T'.serveReadS Nothing 4096 host1p port) >-> P.toListD+ (eex,out) <- P.trySafeIO . P.runWriterT . P.runProxy .P.runEitherK $ p+ case eex of+ Left ex -> E.throw ex+ Right () -> return out+ mvout <- newEmptyMVar+ forkIO $ putMVar mvout =<< serveOnceAndRead'+ threadDelay waitTime+ connectAndWrite host1 port msg1+ out <- takeMVar mvout+ B.concat out @=? msg1b+++tests :: [Test]+tests =+ [ testGroup "TCP"+ [ testGroup "{listen*accept,connect}*{socketReadS,socketWriteD}"+ [ testCase "test_listen_accept_socketReadS_then_connect_socketWriteD"+ test_listen_accept_socketReadS_then_connect_socketWriteD+ , testCase "test_listen_accept_socketWriteD_then_connect_socketReadD"+ test_listen_accept_socketWriteD_then_connect_socketReadD+ ]+ ]+ , testGroup "TCP.Safe"+ [ testGroup "{serve,connect}{WriteD,ReadS}"+ [ testCase "test_safe_serveWriteD" test_safe_serveWriteD+ , testCase "test_safe_serveReadS" test_safe_serveReadS+ ]+ ]+ ]++main :: IO ()+main = NS.withSocketsDo $ defaultMain tests