network-run 0.5.0 → 0.6.0
raw patch · 16 files changed
+1415/−78 lines, 16 filesdep +directorydep +hspecdep +network-rundep ~bytestringdep ~networknew-component:exe:tcp-clientnew-component:exe:tcp-servernew-component:exe:udp-clientnew-component:exe:udp-serverPVP ok
version bump matches the API change (PVP)
Dependencies added: directory, hspec, network-run
Dependency ranges changed: bytestring, network
API changes (from Hackage documentation)
- Network.Run.TCP.Timeout: openServerSocket :: AddrInfo -> IO Socket
- Network.Run.TCP.Timeout: openServerSocketWithOptions :: [(SocketOption, Int)] -> AddrInfo -> IO Socket
- Network.Run.TCP.Timeout: openServerSocketWithOpts :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket
+ Network.Run.TCP: ServerSettings :: (Maybe SockAddr -> SomeException -> IO ()) -> Int -> Int -> ServerSettings
+ Network.Run.TCP: [settingsAcceptRetryDelay] :: ServerSettings -> Int
+ Network.Run.TCP: [settingsGracefulCloseTimeout] :: ServerSettings -> Int
+ Network.Run.TCP: [settingsOnException] :: ServerSettings -> Maybe SockAddr -> SomeException -> IO ()
+ Network.Run.TCP: data ServerSettings
+ Network.Run.TCP: defaultServerSettings :: ServerSettings
+ Network.Run.TCP: runTCPServerWithSettings :: ServerSettings -> Maybe HostName -> ServiceName -> (Socket -> IO a) -> IO a
+ Network.Run.TCP: runTCPServerWithSocketAndSettings :: ServerSettings -> Socket -> (Socket -> IO a) -> IO a
+ Network.Run.TCP.Timeout: ServerSettings :: (Maybe SockAddr -> SomeException -> IO ()) -> Int -> Int -> ServerSettings
+ Network.Run.TCP.Timeout: [settingsAcceptRetryDelay] :: ServerSettings -> Int
+ Network.Run.TCP.Timeout: [settingsGracefulCloseTimeout] :: ServerSettings -> Int
+ Network.Run.TCP.Timeout: [settingsOnException] :: ServerSettings -> Maybe SockAddr -> SomeException -> IO ()
+ Network.Run.TCP.Timeout: data ServerSettings
+ Network.Run.TCP.Timeout: defaultServerSettings :: ServerSettings
+ Network.Run.TCP.Timeout: openTCPServerSocket :: AddrInfo -> IO Socket
+ Network.Run.TCP.Timeout: openTCPServerSocketWithOptions :: [(SocketOption, Int)] -> AddrInfo -> IO Socket
+ Network.Run.TCP.Timeout: openTCPServerSocketWithOpts :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket
+ Network.Run.TCP.Timeout: resolve :: SocketType -> Maybe HostName -> ServiceName -> [AddrInfoFlag] -> (NonEmpty AddrInfo -> AddrInfo) -> IO AddrInfo
+ Network.Run.TCP.Timeout: runTCPServerWithSettings :: ServerSettings -> Int -> Maybe HostName -> ServiceName -> TimeoutServer () -> IO ()
+ Network.Run.TCP.Timeout: runTCPServerWithSocketAndSettings :: ServerSettings -> Int -> Socket -> TimeoutServer () -> IO ()
+ Network.Run.UDP: ServerSettings :: (Maybe SockAddr -> SomeException -> IO ()) -> Int -> Int -> ServerSettings
+ Network.Run.UDP: [settingsAcceptRetryDelay] :: ServerSettings -> Int
+ Network.Run.UDP: [settingsGracefulCloseTimeout] :: ServerSettings -> Int
+ Network.Run.UDP: [settingsOnException] :: ServerSettings -> Maybe SockAddr -> SomeException -> IO ()
+ Network.Run.UDP: data ServerSettings
+ Network.Run.UDP: defaultServerSettings :: ServerSettings
+ Network.Run.UDP: openServerSocket :: AddrInfo -> IO Socket
+ Network.Run.UDP: openServerSocketWithOptions :: [(SocketOption, Int)] -> AddrInfo -> IO Socket
+ Network.Run.UDP: openServerSocketWithOpts :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket
+ Network.Run.UDP: resolve :: SocketType -> Maybe HostName -> ServiceName -> [AddrInfoFlag] -> (NonEmpty AddrInfo -> AddrInfo) -> IO AddrInfo
+ Network.Run.UDP: runUDPServerForkWithSettings :: ServerSettings -> [HostName] -> ServiceName -> (Socket -> ByteString -> IO ()) -> IO ()
Files
- CHANGELOG.md +40/−1
- Network/Run/Core.hs +191/−32
- Network/Run/TCP.hs +38/−12
- Network/Run/TCP/Timeout.hs +46/−16
- Network/Run/UDP.hs +46/−13
- examples/tcpClient.hs +14/−0
- examples/tcpServer.hs +15/−0
- examples/udpClient.hs +24/−0
- examples/udpServer.hs +11/−0
- network-run.cabal +92/−4
- test/CoreSpec.hs +196/−0
- test/Helper.hs +168/−0
- test/Spec.hs +1/−0
- test/TCPSpec.hs +295/−0
- test/TimeoutSpec.hs +96/−0
- test/UDPSpec.hs +142/−0
CHANGELOG.md view
@@ -1,9 +1,48 @@ # Revision history for network-run +## 0.6.0++* Fixing a bug that the `[(SocketOption, Int)]` variants+ (`openClientSocketWithOptions`, `openServerSocketWithOptions` and+ `openTCPServerSocketWithOptions`) passed an `Int`, that is eight+ bytes, to `setsockopt`. BSD rejects this with `EINVAL`, so the+ documented `(IPv6Only, 0)` for a dual stack socket did not work+ there. The value is converted to `CInt` now.+* A test suite.+* Breaking change: Network.Run.TCP.Timeout no longer exports+ `openServerSocket`, `openServerSocketWithOptions` and+ `openServerSocketWithOpts`. They do not call `listen`, so+ `runTCPServerWithSocket` cannot accept on the resulting socket.+ Use `openTCPServerSocket` and friends instead.+* New API: `ServerSettings`, `defaultServerSettings` and the+ `runTCPServerWithSettings`, `runTCPServerWithSocketAndSettings` and+ `runUDPServerForkWithSettings` variants. `settingsOnException`+ receives exceptions which the library catches instead of propagating,+ `settingsGracefulCloseTimeout` controls `gracefulClose` and+ `settingsAcceptRetryDelay` controls the `accept` retry interval.+* New API: Network.Run.TCP.Timeout now exports `resolve`,+ `openTCPServerSocket`, `openTCPServerSocketWithOptions` and+ `openTCPServerSocketWithOpts`.+* `accept` no longer terminates the server on transient errors.+ `ECONNABORTED` and `EINTR` are retried immediately, and+ `EMFILE`/`ENFILE` are retried after a short delay and passed to+ `settingsOnException`.+* `runUDPServerFork` is now exception safe. A failure of `getAddrInfo`,+ `openServerSocket` or `connect` no longer leaks a socket nor kills+ the server; such a datagram is dropped and passed to+ `settingsOnException` instead. An unknown address family is dropped+ rather than calling `error`.+* Fixing a bug that `runUDPServerFork` labels every forked thread with+ the first host name.+* Documenting IPV6_V6ONLY and the single address family of+ `runTCPServer`.+* Network.Run.UDP exports `openServerSocket`,+ `openServerSocketWithOptions` and `openServerSocketWithOpts`.+ ## 0.5.0 * Fixing a bug that TimeoutServer is not killed.-* Breaking change: the signatures of Timeout.runTCPServer and +* Breaking change: the signatures of Timeout.runTCPServer and Timeout.runTCPServerWithSocket are changed. ## 0.4.3
Network/Run/Core.hs view
@@ -1,4 +1,5 @@ {-# LANGUAGE CPP #-}+{-# LANGUAGE RecordWildCards #-} module Network.Run.Core ( resolve,@@ -12,17 +13,29 @@ openTCPServerSocket, openTCPServerSocketWithOptions, openTCPServerSocketWithOpts,- gclose, labelMe,+ safeAccept,+ safeAcceptWith,+ ServerSettings (..),+ defaultServerSettings,+ gcloseWith,+ forkWith,+ forkConnection,+ forkDatagram,+ report, ) where -import Data.List.NonEmpty (NonEmpty)-import Control.Arrow+import Control.Arrow hiding (loop) import Control.Concurrent import qualified Control.Exception as E-import Control.Monad (when)+import Control.Monad (void, when)+import Data.List.NonEmpty (NonEmpty)+import Foreign.C.Error (Errno (..), eCONNABORTED)+import Foreign.C.Types (CInt) import GHC.Conc.Sync+import GHC.IO.Exception (IOErrorType (Interrupted), ioe_errno) import Network.Socket+import System.IO.Error (ioeGetErrorType, isFullError) resolve :: SocketType@@ -40,11 +53,6 @@ , addrFlags = flags } -#if !MIN_VERSION_network(3,1,2)-openSocket :: AddrInfo -> IO Socket-openSocket addr = socket (addrFamily addr) (addrSocketType addr) (addrProtocol addr)-#endif- -- | This is the same as -- -- @@@ -58,10 +66,10 @@ -- The options are set before 'connect'. This is equivalent to -- -- @--- 'openClientSocketWithOpts' . 'map' ('second' 'SockOptValue')+-- 'openClientSocketWithOpts' . 'map' ('second' ('SockOptValue' . 'fromIntegral' \@Int \@CInt)) -- @ openClientSocketWithOptions :: [(SocketOption, Int)] -> AddrInfo -> IO Socket-openClientSocketWithOptions = openClientSocketWithOpts . map (second SockOptValue)+openClientSocketWithOptions = openClientSocketWithOpts . map (second sockOptInt) -- | Open a client socket with the given options --@@ -70,7 +78,8 @@ -- ('Network.Socket.StructLinger'). -- -- The options are set before 'connect'.-openClientSocketWithOpts :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket+openClientSocketWithOpts+ :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket openClientSocketWithOpts opts addr = E.bracketOnError (openSocket addr) close $ \sock -> do mapM_ (uncurry $ setSockOptValue sock) opts connect sock $ addrAddress addr@@ -91,19 +100,32 @@ -- This is equivalent to -- -- @--- 'openServerSocketWithOpts' . 'map' ('second' 'SockOptValue')+-- 'openServerSocketWithOpts' . 'map' ('second' ('SockOptValue' . 'fromIntegral' \@Int \@CInt)) -- @ openServerSocketWithOptions :: [(SocketOption, Int)] -> AddrInfo -> IO Socket-openServerSocketWithOptions = openServerSocketWithOpts . map (second SockOptValue)+openServerSocketWithOptions = openServerSocketWithOpts . map (second sockOptInt) -- | Open socket for server use, and set the provided options before binding. -- -- In addition to the given options, the socket is configured to -- -- * allow reuse of local addresses (SO_REUSEADDR)+-- * accept IPv6 only, rejecting IPv4-mapped addresses, if the address+-- family is 'AF_INET6' (IPV6_V6ONLY) -- * automatically be closed during a successful @execve@ (FD_CLOEXEC) -- * bind to the address specified-openServerSocketWithOpts :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket+--+-- Because IPV6_V6ONLY is in effect, a socket bound to @::@ does not+-- accept IPv4 connections. To serve both families, open one socket per+-- address and run a server on each of them with+-- 'Network.Run.TCP.runTCPServerWithSocket'.+--+-- The given options are set after the ones above, so @(IPv6Only, 0)@+-- can be passed to ask for a dual stack socket. Note that OpenBSD+-- always makes IPv6 sockets IPv6 only; the option is not set there and+-- cannot be cleared.+openServerSocketWithOpts+ :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket openServerSocketWithOpts opts addr = E.bracketOnError (openSocket addr) close $ \sock -> do setSocketOption sock ReuseAddr 1 #if !defined(openbsd_HOST_OS)@@ -114,7 +136,7 @@ bind sock $ addrAddress addr return sock --- | Open TCP socket for server use+-- | Open TCP socket for server use. -- -- This is the same as: --@@ -126,36 +148,173 @@ -- | Open socket for server use, and set the provided options before binding. --+-- This is 'openServerSocketWithOpts' followed by 'listen' with a queue+-- length of 1024. See 'openServerSocketWithOpts' for the options which+-- are set in addition to the given ones.+-- -- This is equivalent to -- -- @--- 'openTCPServerSocketWithOpts' . 'map' ('second' 'SockOptValue')+-- 'openTCPServerSocketWithOpts' . 'map' ('second' ('SockOptValue' . 'fromIntegral' \@Int \@CInt)) -- @ openTCPServerSocketWithOptions :: [(SocketOption, Int)] -> AddrInfo -> IO Socket-openTCPServerSocketWithOptions = openTCPServerSocketWithOpts . map (second SockOptValue)+openTCPServerSocketWithOptions = openTCPServerSocketWithOpts . map (second sockOptInt) --- | Open socket for server use, and set the provided options before binding.------ In addition to the given options, the socket is configured to+-- | Open socket for server use, and set the provided options before+-- binding. ----- * allow reuse of local addresses (SO_REUSEADDR)--- * automatically be closed during a successful @execve@ (FD_CLOEXEC)--- * bind to the address specified--- * listen with queue length with 1024-openTCPServerSocketWithOpts :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket+-- This is 'openServerSocketWithOpts' followed by 'listen' with a queue+-- length of 1024. See 'openServerSocketWithOpts' for the options which+-- are set in addition to the given ones.+openTCPServerSocketWithOpts+ :: [(SocketOption, SockOptValue)] -> AddrInfo -> IO Socket openTCPServerSocketWithOpts opts addr = do sock <- openServerSocketWithOpts opts addr listen sock 1024 return sock -gclose :: Socket -> IO ()-#if MIN_VERSION_network(3,1,1)-gclose sock = gracefulClose sock 5000-#else-gclose = close-#endif+-- | An 'Int' option as a 'SockOptValue'.+--+-- The value must be converted to 'CInt' first. A socket option is an+-- @int@ in C, and 'SockOptValue' passes @sizeof@ of the value it is+-- given, so an 'Int' asks the kernel to read eight bytes. Linux+-- ignores the extra ones, but BSD rejects the call with @EINVAL@.+sockOptInt :: Int -> SockOptValue+sockOptInt = SockOptValue . (fromIntegral :: Int -> CInt) labelMe :: String -> IO () labelMe name = do tid <- myThreadId labelThread tid name++----------------------------------------------------------------++-- | Settings for servers.+--+-- Fields which do not apply to a given server (for instance the+-- graceful close timeout for a UDP server) are ignored.+data ServerSettings = ServerSettings+ { settingsOnException :: Maybe SockAddr -> E.SomeException -> IO ()+ -- ^ Called when an exception is caught by the library instead of+ -- being propagated. The 'SockAddr' is 'Just' the peer when the+ -- exception can be attributed to one. Exceptions thrown by this+ -- action itself are discarded, so it must not be relied on for+ -- anything but reporting. The default does nothing.+ , settingsGracefulCloseTimeout :: Int+ -- ^ Milliseconds 'gracefulClose' waits for the peer's FIN after a+ -- connection handler returns. Zero or less uses 'close' instead,+ -- which releases the file descriptor immediately. The default is+ -- 5000.+ , settingsAcceptRetryDelay :: Int+ -- ^ Microseconds to wait before retrying 'accept' after running+ -- out of file descriptors. The default is 100000.+ }++-- | Default settings. 'settingsOnException' does nothing, so the+-- behaviour is the same as before this type was introduced.+defaultServerSettings :: ServerSettings+defaultServerSettings =+ ServerSettings+ { settingsOnException = \_ _ -> return ()+ , settingsGracefulCloseTimeout = 5000+ , settingsAcceptRetryDelay = 100000+ }++-- | Calling 'settingsOnException', never letting it throw. A hook+-- must not be able to break a finalizer.+report :: ServerSettings -> Maybe SockAddr -> E.SomeException -> IO ()+report ServerSettings{..} mpeer se =+ settingsOnException mpeer se `E.catch` ignore+ where+ ignore :: E.SomeException -> IO ()+ ignore e+ | Just (E.SomeAsyncException _) <- E.fromException e = E.throwIO e+ | otherwise = return ()++-- | Closing a connected socket according to the settings.+gcloseWith :: ServerSettings -> Socket -> IO ()+gcloseWith ServerSettings{..} sock+ | settingsGracefulCloseTimeout <= 0 = close sock+ | otherwise = gracefulClose sock settingsGracefulCloseTimeout++----------------------------------------------------------------++-- | Accepting a connection, retrying on transient errors.+--+-- 'accept' fails routinely for reasons which do not mean that the+-- listening socket is broken: the peer may reset the connection before+-- it is accepted (@ECONNABORTED@), or the process or the system may+-- have run out of file descriptors (@EMFILE@\/@ENFILE@). Letting+-- these escape would terminate the accept loop, so they are retried+-- here. Errors which do suggest a broken listening socket (@EBADF@,+-- @EINVAL@, ...) are re-thrown, which is also how a closed socket+-- stops the loop.+--+-- Running out of file descriptors is passed to 'settingsOnException'+-- since a server which keeps hitting it is effectively out of service.+-- @ECONNABORTED@ and @EINTR@ are not, being routine.+--+-- This function is interruptible: a blocked or sleeping retry still+-- receives asynchronous exceptions, so the server remains killable.+safeAccept :: ServerSettings -> Socket -> IO (Socket, SockAddr)+safeAccept set sock = safeAcceptWith set $ accept sock++-- | 'safeAccept' with the accepting action passed explicitly. The+-- error paths described above cannot be provoked on a real listening+-- socket, so the test suite reaches them through this.+safeAcceptWith+ :: ServerSettings -> IO (Socket, SockAddr) -> IO (Socket, SockAddr)+safeAcceptWith set@ServerSettings{..} accept' = loop+ where+ loop = do+ ex <- E.try accept'+ case ex of+ Right r -> return r+ Left e+ -- No descriptor is available at the moment. Retrying+ -- at once would spin, since the listening socket stays+ -- readable.+ | isFullError e -> do+ report set Nothing $ E.toException e+ threadDelay settingsAcceptRetryDelay+ loop+ -- These cost nothing; retry immediately.+ | ioeGetErrorType e == Interrupted -> loop+ | ioe_errno e == Just connAborted -> loop+ | otherwise -> E.throwIO e++ Errno connAborted = eCONNABORTED++----------------------------------------------------------------++-- | Forking a thread for an accepted socket. An exception which+-- escapes the action is reported, and the socket is closed by the+-- given closer in either case.+forkWith+ :: ServerSettings+ -> (Socket -> IO ())+ -> Socket+ -> SockAddr+ -> IO a+ -> IO ()+forkWith set closer sock peer action = void $ forkFinally action finish+ where+ -- The socket must be closed even if the hook throws, which it does+ -- when an asynchronous exception arrives while it is running.+ finish er = reporting er `E.finally` closing++ reporting (Right _) = return ()+ reporting (Left se) = report set (Just peer) se++ closing = closer sock `E.catch` onCloseError++ onCloseError :: E.IOException -> IO ()+ onCloseError e = report set (Just peer) $ E.toException e++-- | 'forkWith' closing the socket gracefully. For TCP.+forkConnection :: ServerSettings -> Socket -> SockAddr -> IO a -> IO ()+forkConnection set = forkWith set (gcloseWith set)++-- | 'forkWith' closing the socket immediately. For UDP.+forkDatagram :: ServerSettings -> Socket -> SockAddr -> IO a -> IO ()+forkDatagram set = forkWith set close
Network/Run/TCP.hs view
@@ -1,11 +1,14 @@-{-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE RecordWildCards #-} -- | Simple functions to run TCP clients and servers. module Network.Run.TCP ( -- * Server runTCPServer,+ runTCPServerWithSettings, runTCPServerWithSocket,+ runTCPServerWithSocketAndSettings,+ ServerSettings (..),+ defaultServerSettings, openTCPServerSocket, openTCPServerSocketWithOptions, openTCPServerSocketWithOpts,@@ -23,9 +26,8 @@ openClientSocketWithOpts, ) where -import Control.Concurrent (forkFinally) import qualified Control.Exception as E-import Control.Monad (forever, void)+import Control.Monad (forever) import Data.List.NonEmpty (NonEmpty) import qualified Data.List.NonEmpty as NE import Network.Socket@@ -34,24 +36,48 @@ ---------------------------------------------------------------- --- | Running a TCP server with an accepted socket and its peer name.+-- | Running a TCP server with an accepted socket.+--+-- Only the first address returned for @mhost@ is used, so a server+-- created by this function listens on a single address family. Use+-- 'runTCPServerWithSocket' with one socket per address to serve both+-- IPv4 and IPv6. runTCPServer :: Maybe HostName -> ServiceName -> (Socket -> IO a) -> IO a-runTCPServer mhost port server = do+runTCPServer = runTCPServerWithSettings defaultServerSettings++-- | Running a TCP server with the given settings.+runTCPServerWithSettings+ :: ServerSettings+ -> Maybe HostName+ -> ServiceName+ -> (Socket -> IO a)+ -> IO a+runTCPServerWithSettings set mhost port server = do addr <- resolve Stream mhost port [AI_PASSIVE] NE.head E.bracket (openTCPServerSocket addr) close $ \sock ->- runTCPServerWithSocket sock server+ runTCPServerWithSocketAndSettings set sock server --- | Running a TCP client with a connected socket for a given listen--- socket.+-- | Running a TCP server on a given listen socket. runTCPServerWithSocket :: Socket+ -- ^ A listening socket created by 'openTCPServerSocket'. -> (Socket -> IO a) -- ^ Called for each incoming connection, in a new thread -> IO a-runTCPServerWithSocket sock server = forever $- E.bracketOnError (accept sock) (close . fst) $- \(conn, _peer) ->- void $ forkFinally (labelMe "TCP server" >> server conn) (const $ gclose conn)+runTCPServerWithSocket = runTCPServerWithSocketAndSettings defaultServerSettings++-- | Running a TCP server on a given listen socket with the given+-- settings.+runTCPServerWithSocketAndSettings+ :: ServerSettings+ -> Socket+ -- ^ A listening socket created by 'openTCPServerSocket'.+ -> (Socket -> IO a)+ -- ^ Called for each incoming connection, in a new thread+ -> IO a+runTCPServerWithSocketAndSettings set sock server = forever $+ E.bracketOnError (safeAccept set sock) (close . fst) $ \(conn, peer) ->+ forkConnection set conn peer (labelMe "TCP server" >> server conn) ----------------------------------------------------------------
Network/Run/TCP/Timeout.hs view
@@ -1,20 +1,22 @@-{-# LANGUAGE OverloadedStrings #-}---- | Simple functions to run TCP clients and servers.+-- | Simple functions to run TCP servers. module Network.Run.TCP.Timeout ( runTCPServer,+ runTCPServerWithSettings, TimeoutServer,+ ServerSettings (..),+ defaultServerSettings,+ resolve, -- * Generalized API runTCPServerWithSocket,- openServerSocket,- openServerSocketWithOptions,- openServerSocketWithOpts,+ runTCPServerWithSocketAndSettings,+ openTCPServerSocket,+ openTCPServerSocketWithOptions,+ openTCPServerSocketWithOpts, ) where -import Control.Concurrent (forkFinally) import qualified Control.Exception as E-import Control.Monad (forever, void)+import Control.Monad (forever) import qualified Data.List.NonEmpty as NE import Network.Socket import qualified System.TimeManager as T@@ -31,7 +33,12 @@ -- ^ A connected socket -> IO a --- | Running a TCP server with a connected socket.+-- | Running a TCP server, resolving and binding the address itself.+--+-- Only the first address returned for @mhost@ is used, so a server+-- created by this function listens on a single address family. Use+-- 'runTCPServerWithSocket' with one socket per address to serve both+-- IPv4 and IPv6. runTCPServer :: Int -- ^ Timeout in second.@@ -39,23 +46,46 @@ -> ServiceName -> TimeoutServer () -> IO ()-runTCPServer tm mhost port server = do+runTCPServer = runTCPServerWithSettings defaultServerSettings++-- | Running a TCP server with the given settings.+runTCPServerWithSettings+ :: ServerSettings+ -> Int+ -- ^ Timeout in second.+ -> Maybe HostName+ -> ServiceName+ -> TimeoutServer ()+ -> IO ()+runTCPServerWithSettings set tm mhost port server = do addr <- resolve Stream mhost port [AI_PASSIVE] NE.head E.bracket (openTCPServerSocket addr) close $ \sock ->- runTCPServerWithSocket tm sock server+ runTCPServerWithSocketAndSettings set tm sock server --- | Running a TCP client with a connected socket for a given listen--- socket.+-- | Running a TCP server on a given listen socket. runTCPServerWithSocket :: Int -- ^ Timeout in second. -> Socket+ -- ^ A listening socket created by 'openTCPServerSocket'. -> TimeoutServer () -> IO ()-runTCPServerWithSocket tm sock server = do+runTCPServerWithSocket = runTCPServerWithSocketAndSettings defaultServerSettings++-- | Running a TCP server on a given listen socket with the given+-- settings.+runTCPServerWithSocketAndSettings+ :: ServerSettings+ -> Int+ -- ^ Timeout in second.+ -> Socket+ -- ^ A listening socket created by 'openTCPServerSocket'.+ -> TimeoutServer ()+ -> IO ()+runTCPServerWithSocketAndSettings set tm sock server = T.withManager (tm * 1000000) $ \mgr -> forever $- E.bracketOnError (accept sock) (close . fst) $ \(conn, _peer) ->- void $ forkFinally (runServer mgr conn) (const $ gclose conn)+ E.bracketOnError (safeAccept set sock) (close . fst) $ \(conn, peer) ->+ forkConnection set conn peer (runServer mgr conn) where runServer mgr conn = do labelMe "TCP timeout server"
Network/Run/UDP.hs view
@@ -3,11 +3,18 @@ runUDPClient, runUDPServer, runUDPServerFork,+ runUDPServerForkWithSettings,+ ServerSettings (..),+ defaultServerSettings,+ openServerSocket,+ openServerSocketWithOptions,+ openServerSocketWithOpts,+ resolve, ) where -import Control.Concurrent (forkFinally, forkIO)+import Control.Concurrent (forkIO) import qualified Control.Exception as E-import Control.Monad (forever, void)+import Control.Monad (forever) import Data.ByteString (ByteString) import qualified Data.List.NonEmpty as NE import Network.Socket@@ -42,26 +49,52 @@ -- This approach is fragile due to NAT rebidings. runUDPServerFork :: [HostName] -> ServiceName -> (Socket -> ByteString -> IO ()) -> IO ()-runUDPServerFork [] _ _ = return ()-runUDPServerFork (h : hs) port server = do+runUDPServerFork = runUDPServerForkWithSettings defaultServerSettings++-- | 'runUDPServerFork' with the given settings.+runUDPServerForkWithSettings+ :: ServerSettings+ -> [HostName]+ -> ServiceName+ -> (Socket -> ByteString -> IO ())+ -> IO ()+runUDPServerForkWithSettings _ [] _ _ = return ()+runUDPServerForkWithSettings set (h : hs) port server = do mapM_ (forkIO . run) hs run h where run host = do- labelMe $ "UDP server for " ++ h+ labelMe $ "UDP server for " ++ host runUDPServer (Just host) port $ \lsock -> forever $ do+ -- An error from 'recvFrom' means that the listening socket+ -- itself is gone, so it is left to propagate as before. (bs0, peeraddr) <- recvFrom lsock 2048- let family = case peeraddr of- SockAddrInet{} -> AF_INET- SockAddrInet6{} -> AF_INET6- _ -> error "family"- hints =+ -- Everything below is per-datagram work. A failure here+ -- must not take the entire server down.+ dispatch peeraddr bs0 `E.catch` onDispatchError peeraddr++ onDispatchError peeraddr e =+ report set (Just peeraddr) $ E.toException (e :: E.IOException)++ dispatch peeraddr bs0 = case familyOf peeraddr of+ -- Neither IPv4 nor IPv6. Just drop the datagram.+ Nothing -> return ()+ Just family -> do+ let hints = defaultHints { addrSocketType = Datagram , addrFamily = family , addrFlags = [AI_PASSIVE] } addr <- NE.head <$> getAddrInfo (Just hints) Nothing (Just port)- s <- openServerSocket addr- connect s peeraddr- void $ forkFinally (labelMe "UDP server" >> server s bs0) (\_ -> close s)+ -- If 'connect' throws, the socket is closed here. On+ -- success it is owned by the new thread and is closed by+ -- its finalizer.+ E.bracketOnError (openServerSocket addr) close $ \s -> do+ connect s peeraddr+ forkDatagram set s peeraddr $+ labelMe "UDP server" >> server s bs0++ familyOf SockAddrInet{} = Just AF_INET+ familyOf SockAddrInet6{} = Just AF_INET6+ familyOf _ = Nothing
+ examples/tcpClient.hs view
@@ -0,0 +1,14 @@+{-# LANGUAGE OverloadedStrings #-}++module Main (main) where++import qualified Data.ByteString.Char8 as C+import Network.Run.TCP (runTCPClient)+import Network.Socket.ByteString (recv, sendAll)++main :: IO ()+main = runTCPClient "127.0.0.1" "3000" $ \s -> do+ sendAll s "Hello, world!"+ msg <- recv s 1024+ putStr "Received: "+ C.putStrLn msg
+ examples/tcpServer.hs view
@@ -0,0 +1,15 @@+module Main (main) where++import Control.Monad (unless)+import qualified Data.ByteString as S+import Network.Run.TCP (runTCPServer)+import Network.Socket.ByteString (recv, sendAll)++main :: IO ()+main = runTCPServer Nothing "3000" talk+ where+ talk s = do+ msg <- recv s 1024+ unless (S.null msg) $ do+ sendAll s msg+ talk s
+ examples/udpClient.hs view
@@ -0,0 +1,24 @@+{-# LANGUAGE OverloadedStrings #-}++-- Echo client program+module Main (main) where++import qualified Data.ByteString.Char8 as C+import Network.Run.UDP (runUDPClient)+import Network.Socket+import Network.Socket.ByteString (recvFrom, sendTo)++main :: IO ()+main = runUDPClient "127.0.0.1" "3000" $ \sock sockAddr -> do+ -- Initially the local port is 0+ my1 <- getSocketName sock+ putStrLn $ "My sock addr " ++ show my1+ putStrLn $ "Peer sock addr " ++ show sockAddr+ _ <- sendTo sock "Hello, world!" sockAddr+ -- After sendTo, the local port is implicitly bound+ my2 <- getSocketName sock+ putStrLn $ "My sock addr " ++ show my2+ (msg, peer) <- recvFrom sock 1024+ putStrLn $ "Peer sock addr " ++ show peer+ putStr "Received: "+ C.putStrLn msg
+ examples/udpServer.hs view
@@ -0,0 +1,11 @@+module Main (main) where++import Control.Monad (forever, unless, void)+import qualified Data.ByteString as S+import Network.Run.UDP (runUDPServer)+import Network.Socket.ByteString (recvFrom, sendTo)++main :: IO ()+main = runUDPServer (Just "127.0.0.1") "3000" $ \sock -> forever $ do+ (msg, peer) <- recvFrom sock 2048+ unless (S.null msg) $ void $ sendTo sock msg peer
network-run.cabal view
@@ -1,6 +1,6 @@-cabal-version: >=1.10+cabal-version: 2.0 name: network-run-version: 0.5.0+version: 0.6.0 license: BSD3 license-file: LICENSE maintainer: kazu@iij.ad.jp@@ -9,12 +9,17 @@ description: Simple functions to run network clients and servers. category: Network build-type: Simple-extra-source-files: CHANGELOG.md+extra-doc-files: CHANGELOG.md source-repository head type: git location: https://github.com/kazu-yamamoto/network-run +flag examples+ description: Build the example programs+ default: False+ manual: True+ library exposed-modules: Network.Run.TCP@@ -25,6 +30,89 @@ default-language: Haskell2010 build-depends: base >=4 && <5,+ bytestring >=0.10 && <0.13,+ network >=3.2.4 && <3.3,+ time-manager >=0.2 && <0.4++executable tcp-client+ main-is: tcpClient.hs+ hs-source-dirs: examples+ default-language: Haskell2010+ ghc-options: -Wall+ build-depends:+ base >=4 && <5,+ bytestring >=0.10 && <0.13,+ network >=3.2.4 && <3.3,+ network-run++ if !flag(examples)+ buildable: False++executable tcp-server+ main-is: tcpServer.hs+ hs-source-dirs: examples+ default-language: Haskell2010+ ghc-options: -Wall+ build-depends:+ base >=4 && <5,+ bytestring >=0.10 && <0.13,+ network >=3.2.4 && <3.3,+ network-run++ if !flag(examples)+ buildable: False++executable udp-client+ main-is: udpClient.hs+ hs-source-dirs: examples+ default-language: Haskell2010+ ghc-options: -Wall+ build-depends:+ base >=4 && <5,+ bytestring >=0.10 && <0.13,+ network >=3.2.4 && <3.3,+ network-run++ if !flag(examples)+ buildable: False++executable udp-server+ main-is: udpServer.hs+ hs-source-dirs: examples+ default-language: Haskell2010+ ghc-options: -Wall+ build-depends:+ base >=4 && <5,+ bytestring >=0.10 && <0.13,+ network >=3.2.4 && <3.3,+ network-run++ if !flag(examples)+ buildable: False++test-suite spec+ type: exitcode-stdio-1.0+ main-is: Spec.hs+ build-tool-depends: hspec-discover:hspec-discover+ hs-source-dirs: test .+ other-modules:+ CoreSpec+ Helper+ TCPSpec+ TimeoutSpec+ UDPSpec+ Network.Run.Core+ Network.Run.TCP+ Network.Run.TCP.Timeout+ Network.Run.UDP++ default-language: Haskell2010+ ghc-options: -Wall -threaded -rtsopts -with-rtsopts=-N+ build-depends:+ base >=4 && <5, bytestring,+ directory,+ hspec, network >=3.2.4,- time-manager >=0.2 && <0.4+ time-manager >=0.2 && <0.4,+ network-run
+ test/CoreSpec.hs view
@@ -0,0 +1,196 @@+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE ScopedTypeVariables #-}++-- | White box tests for the internals of "Network.Run.Core".+--+-- The transient 'accept' errors and the closer failures handled there+-- cannot be provoked through a real socket, so they are injected.+module CoreSpec (spec) where++import Control.Concurrent+import qualified Control.Exception as E+import Control.Monad (void)+import Data.IORef+import Foreign.C.Error (Errno (..), eCONNABORTED)+import GHC.IO.Exception (+ IOErrorType (Interrupted, OtherError),+ IOException (..),+ )+import Network.Socket+import System.IO.Error (fullErrorType, illegalOperationErrorType, mkIOError)+import System.Timeout (timeout)+import Test.Hspec++import Network.Run.Core++import Helper++spec :: Spec+spec = do+ describe "safeAcceptWith" $ do+ it "retries after running out of file descriptors" $ limited $ do+ (set0, getReports) <- collecting+ let set = set0{settingsAcceptRetryDelay = 50000}+ withFakeAccept [emfile, emfile] $ \(accept', count) -> do+ ((_, peer), ms) <- elapsed $ safeAcceptWith set accept'+ peer `shouldBe` fakePeer+ count `shouldReturn` 3+ -- Two retries of 50ms each.+ ms `shouldSatisfy` (>= 90)+ reports <- getReports+ map fst reports `shouldBe` [Nothing, Nothing]++ it "retries EINTR at once, without reporting it" $ limited $ do+ (set0, getReports) <- collecting+ -- A delay which would be obvious if it were taken.+ let set = set0{settingsAcceptRetryDelay = 5000000}+ withFakeAccept [eintr, eintr] $ \(accept', count) -> do+ (_, ms) <- elapsed $ safeAcceptWith set accept'+ count `shouldReturn` 3+ ms `shouldSatisfy` (< 1000)+ getReports `shouldReturn` []++ it "retries ECONNABORTED at once, without reporting it" $ limited $ do+ (set0, getReports) <- collecting+ let set = set0{settingsAcceptRetryDelay = 5000000}+ withFakeAccept [aborted] $ \(accept', count) -> do+ (_, ms) <- elapsed $ safeAcceptWith set accept'+ count `shouldReturn` 2+ ms `shouldSatisfy` (< 1000)+ getReports `shouldReturn` []++ it "rethrows an error of the listening socket" $ limited $ do+ (set, getReports) <- collecting+ withFakeAccept [bad] $ \(accept', count) -> do+ safeAcceptWith set accept'+ `shouldThrow` (\e -> ioeGetErrorType' e == ioeGetErrorType' bad)+ count `shouldReturn` 1+ getReports `shouldReturn` []++ it "is still killable while waiting to retry" $ limited $ do+ let set = defaultServerSettings{settingsAcceptRetryDelay = 5000000}+ withFakeAccept (repeat emfile) $ \(accept', _) -> do+ done <- newEmptyMVar+ tid <- forkFinally (void $ safeAcceptWith set accept') (putMVar done)+ threadDelay 100000+ killThread tid+ r <- timeout 1000000 $ takeMVar done+ case r of+ Just (Left _) -> return ()+ Just (Right _) -> expectationFailure "accept returned"+ Nothing -> expectationFailure "the retry was not interruptible"++ describe "report" $ do+ it "swallows a synchronous exception of the hook" $ do+ let set =+ defaultServerSettings+ { settingsOnException = \_ _ -> E.throwIO $ userError "hook"+ }+ report set Nothing (E.toException $ userError "boom")+ `shouldReturn` ()++ it "rethrows an asynchronous exception of the hook" $ do+ let set =+ defaultServerSettings+ { settingsOnException = \_ _ -> E.throwIO E.ThreadKilled+ }+ report set Nothing (E.toException $ userError "boom")+ `shouldThrow` (== E.ThreadKilled)++ describe "forkWith" $ do+ it "reports an exception which escapes the action" $ limited $ do+ (set, getReports) <- collecting+ withDummySocket $ \sock -> do+ closed <- newEmptyMVar+ forkWith set (\_ -> putMVar closed ()) sock fakePeer $+ E.throwIO $+ userError "boom"+ takeMVar closed `shouldReturn` ()+ reports <- waitFor 1 getReports+ map fst reports `shouldBe` [Just fakePeer]++ it "closes the socket even when the hook throws" $ limited $ do+ let set =+ defaultServerSettings+ { settingsOnException = \_ _ -> E.throwIO $ userError "hook"+ }+ withDummySocket $ \sock -> do+ closed <- newEmptyMVar+ forkWith set (\_ -> putMVar closed ()) sock fakePeer $+ E.throwIO $+ userError "boom"+ r <- timeout 1000000 $ takeMVar closed+ r `shouldBe` Just ()++ it "reports a failure of the closer" $ limited $ do+ (set, getReports) <- collecting+ withDummySocket $ \sock ->+ forkWith set (\_ -> ioError $ userError "close failed") sock fakePeer $+ return ()+ reports <- waitFor 1 getReports+ case reports of+ [(mpeer, desc)] -> do+ mpeer `shouldBe` Just fakePeer+ desc `shouldContain` "close failed"+ _ -> expectationFailure $ "unexpected reports: " ++ show reports++ describe "gcloseWith" $ do+ it "waits for the FIN of the peer if the timeout is positive" $+ limited $+ withHeldConnection $ \sock -> do+ let set = defaultServerSettings{settingsGracefulCloseTimeout = 500}+ (_, ms) <- elapsed $ gcloseWith set sock+ ms `shouldSatisfy` (>= 300)++ it "closes at once if the timeout is not positive" $+ limited $+ withHeldConnection $ \sock -> do+ let set = defaultServerSettings{settingsGracefulCloseTimeout = 0}+ (_, ms) <- elapsed $ gcloseWith set sock+ ms `shouldSatisfy` (< 300)++----------------------------------------------------------------++fakePeer :: SockAddr+fakePeer = SockAddrInet 12345 $ tupleToHostAddress (127, 0, 0, 1)++-- | An 'accept' which fails with the given errors before succeeding,+-- together with the number of times it has been called.+withFakeAccept+ :: [IOError] -> ((IO (Socket, SockAddr), IO Int) -> IO a) -> IO a+withFakeAccept errs body = withDummySocket $ \sock -> do+ ref <- newIORef errs+ cnt <- newIORef (0 :: Int)+ let accept' = do+ atomicModifyIORef' cnt $ \n -> (n + 1, ())+ me <- atomicModifyIORef' ref $ \es -> case es of+ [] -> ([], Nothing)+ e : rest -> (rest, Just e)+ case me of+ Just e -> E.throwIO e+ Nothing -> return (sock, fakePeer)+ body (accept', readIORef cnt)++-- | A socket which is never connected, standing in for an accepted one.+withDummySocket :: (Socket -> IO a) -> IO a+withDummySocket = E.bracket (socket AF_INET Stream defaultProtocol) close++-- | The server side of a connection whose peer stays open and silent,+-- so that a graceful close has to wait for its timeout.+withHeldConnection :: (Socket -> IO a) -> IO a+withHeldConnection body = withListenSocket $ \lsock port -> do+ var <- newEmptyMVar+ withServerThread (accept lsock >>= putMVar var . fst) $+ client port $+ \_held -> takeMVar var >>= body++emfile, eintr, aborted, bad :: IOError+emfile = mkIOError fullErrorType "accept" Nothing Nothing+eintr = emfile{ioe_type = Interrupted}+aborted = emfile{ioe_type = OtherError, ioe_errno = Just connAborted}+ where+ Errno connAborted = eCONNABORTED+bad = mkIOError illegalOperationErrorType "accept" Nothing Nothing++ioeGetErrorType' :: IOError -> IOErrorType+ioeGetErrorType' = ioe_type
+ test/Helper.hs view
@@ -0,0 +1,168 @@+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE ScopedTypeVariables #-}++-- | Helpers shared by the specs.+--+-- Two rules keep these tests reliable:+--+-- * A server is always given port @0@ and its real port is read back+-- with 'getSocketName', so nothing depends on a fixed port and the+-- suite can be run concurrently with anything else.+--+-- * Every test body is wrapped in 'limited'. The library is full of+-- 'forever' loops, so a regression must fail the suite instead of+-- hanging it.+module Helper where++import Control.Concurrent+import qualified Control.Exception as E+import Control.Monad+import Data.ByteString (ByteString)+import qualified Data.ByteString as BS+import Data.IORef+import qualified Data.List.NonEmpty as NE+import GHC.Clock (getMonotonicTimeNSec)+import Network.Socket+import Network.Socket.ByteString+import System.Directory (doesDirectoryExist, getDirectoryContents)+import System.Timeout (timeout)++import Network.Run.Core+import Network.Run.TCP (runTCPClient, runTCPServerWithSocketAndSettings)++----------------------------------------------------------------++-- | Failing instead of hanging.+limited :: IO a -> IO a+limited action = do+ ma <- timeout (10 * 1000000) action+ case ma of+ Nothing -> E.throwIO $ userError "the test did not finish in time"+ Just a -> return a++loopback :: HostName+loopback = "127.0.0.1"++portOf :: SockAddr -> PortNumber+portOf (SockAddrInet p _) = p+portOf (SockAddrInet6 p _ _ _) = p+portOf sa = error $ "portOf: " ++ show sa++-- | Elapsed milliseconds of an action.+elapsed :: IO a -> IO (a, Int)+elapsed action = do+ t0 <- getMonotonicTimeNSec+ a <- action+ t1 <- getMonotonicTimeNSec+ return (a, fromIntegral ((t1 - t0) `div` 1000000))++ignoreAny :: IO a -> IO ()+ignoreAny action = void action `E.catch` \(_ :: E.SomeException) -> return ()++----------------------------------------------------------------++-- | A TCP listening socket on an ephemeral port of the loopback.+withListenSocket :: (Socket -> PortNumber -> IO a) -> IO a+withListenSocket body = do+ addr <- resolve Stream (Just loopback) "0" [AI_PASSIVE] NE.head+ E.bracket (openTCPServerSocket addr) close $ \lsock -> do+ port <- portOf <$> getSocketName lsock+ body lsock port++-- | Running a server thread while the body runs.+--+-- A server which dies on its own, for instance because its port was+-- taken between the moment it was found free and 'bind', would+-- otherwise show up as the body waiting for an answer which is never+-- coming. Its exception is thrown to the caller instead, so that the+-- failure says what actually happened.+withServerThread :: IO () -> IO a -> IO a+withServerThread server body = do+ caller <- myThreadId+ stopping <- newIORef False+ let died (Right ()) = return ()+ died (Left e) = do+ stop <- readIORef stopping+ unless stop $ E.throwTo caller $ ServerDied e+ stopServer tid = writeIORef stopping True >> killThread tid+ E.bracket (forkFinally server died) stopServer $ \_ -> body++-- | A server thread which died on its own.+newtype ServerDied = ServerDied E.SomeException++instance Show ServerDied where+ show (ServerDied e) = "the server thread died: " ++ show e++instance E.Exception ServerDied++-- | Running a TCP server on an ephemeral port while the body runs.+withTCPServer+ :: ServerSettings -> (Socket -> IO ()) -> (PortNumber -> IO a) -> IO a+withTCPServer set server body = withListenSocket $ \lsock port ->+ withServerThread (void $ runTCPServerWithSocketAndSettings set lsock server) $+ body port++client :: PortNumber -> (Socket -> IO a) -> IO a+client port = runTCPClient loopback (show port)++-- | One request and one response on a fresh connection.+request :: PortNumber -> ByteString -> IO ByteString+request port bs = client port $ \sock -> sendAll sock bs >> recv sock 1024++echo :: Socket -> IO ()+echo sock = loop+ where+ loop = do+ bs <- recv sock 1024+ unless (BS.null bs) $ sendAll sock bs >> loop++----------------------------------------------------------------++-- | What 'settingsOnException' was called with.+type Report = (Maybe SockAddr, String)++-- | Settings which record every reported exception.+collecting :: IO (ServerSettings, IO [Report])+collecting = do+ ref <- newIORef []+ let set =+ defaultServerSettings+ { settingsOnException = \mpeer se ->+ atomicModifyIORef' ref $ \rs -> (rs ++ [(mpeer, show se)], ())+ }+ return (set, readIORef ref)++-- | Waiting until at least @n@ items are available, since reporting+-- happens in another thread.+waitFor :: Int -> IO [a] -> IO [a]+waitFor n getter = go (300 :: Int)+ where+ go 0 = getter+ go k = do+ xs <- getter+ if length xs >= n+ then return xs+ else threadDelay 10000 >> go (k - 1)++----------------------------------------------------------------++-- | Whether this machine can open an IPv6 socket at all.+hasIPv6 :: IO Bool+hasIPv6 = do+ er <- E.try go+ return $ either (const False) (const True) (er :: Either E.IOException ())+ where+ go = do+ addr <- resolve Stream (Just "::1") "0" [AI_PASSIVE] NE.head+ E.bracket (openServerSocket addr) close $ \_ -> return ()++-- | The number of open file descriptors, where the platform shows them.+openFds :: IO (Maybe Int)+openFds = go ["/proc/self/fd", "/dev/fd"]+ where+ go [] = return Nothing+ go (d : ds) = do+ exist <- doesDirectoryExist d+ if exist+ then Just . length <$> getDirectoryContents d+ else go ds
+ test/Spec.hs view
@@ -0,0 +1,1 @@+{-# OPTIONS_GHC -F -pgmF hspec-discover #-}
+ test/TCPSpec.hs view
@@ -0,0 +1,295 @@+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE ScopedTypeVariables #-}++module TCPSpec (spec) where++import Control.Concurrent+import qualified Control.Exception as E+import Control.Monad+import Data.ByteString (ByteString)+import Data.IORef+import Data.List (nub)+import qualified Data.List.NonEmpty as NE+import GHC.IO.Exception (IOErrorType (InvalidArgument))+import Network.Socket+import Network.Socket.ByteString+import System.IO.Error (ioeGetErrorType)+import System.Info (os)+import System.Timeout (timeout)+import Test.Hspec++import Network.Run.Core (openServerSocket)+import Network.Run.TCP++import Helper++spec :: Spec+spec = do+ describe "runTCPServer" $ do+ it "serves a connection" $+ limited $+ withTCPServer defaultServerSettings echo $ \port ->+ request port "hello" `shouldReturn` "hello"++ it "resolves and binds the address itself" $ limited $ do+ -- The port cannot be chosen in advance here, so an+ -- ephemeral one is looked up first.+ port <- freeTCPPort+ withServerThread (runTCPServer (Just loopback) (show port) echo) $ do+ threadDelay 200000+ request port "hello" `shouldReturn` "hello"++ it "keeps serving after a handler throws" $ limited $ do+ (set, getReports) <- collecting+ ref <- newIORef (0 :: Int)+ withTCPServer set (failFirst ref) $ \port -> do+ ignoreAny $ request port "hello"+ request port "hello" `shouldReturn` "hello"+ reports <- waitFor 1 getReports+ length reports `shouldBe` 1++ it "reports the peer of a failed handler" $ limited $ do+ (set, getReports) <- collecting+ withTCPServer set (\_ -> E.throwIO $ userError "boom") $ \port -> do+ peer <- client port $ \sock -> do+ sendAll sock "hello"+ getSocketName sock+ reports <- waitFor 1 getReports+ map fst reports `shouldBe` [Just peer]++ it "keeps serving when the exception hook itself throws" $ limited $ do+ let set =+ defaultServerSettings+ { settingsOnException = \_ _ -> E.throwIO $ userError "hook"+ }+ ref <- newIORef (0 :: Int)+ withTCPServer set (failFirst ref) $ \port -> do+ ignoreAny $ request port "hello"+ request port "hello" `shouldReturn` "hello"++ it "stops when the listening socket is closed" $ limited $ do+ done <- newEmptyMVar+ withListenSocket $ \lsock _ -> do+ void $+ forkFinally (runTCPServerWithSocket lsock echo) (putMVar done)+ threadDelay 100000+ close lsock+ r <- timeout 2000000 $ takeMVar done+ case r of+ Just (Left _) -> return ()+ Just (Right _) -> expectationFailure "the accept loop returned"+ Nothing -> expectationFailure "the accept loop did not stop"++ it "listens on a single address family" $ limited $ do+ port <- freeTCPPort+ let hints =+ defaultHints+ { addrSocketType = Stream+ , addrFlags = [AI_PASSIVE]+ }+ addrs <- getAddrInfo (Just hints) (Just "localhost") (Just $ show port)+ let families = NE.toList $ NE.map addrFamily addrs+ if length (nub families) < 2+ then pendingWith "localhost has a single address family here"+ else do+ -- Only the first address is used, so the other+ -- family is not served at all.+ let (served, unserved)+ | addrFamily (NE.head addrs) == AF_INET6 =+ ("::1", "127.0.0.1")+ | otherwise = ("127.0.0.1", "::1")+ withServerThread (runTCPServer (Just "localhost") (show port) echo) $ do+ threadDelay 200000+ echoOn served port `shouldReturn` "hello"+ echoOn unserved port `shouldThrow` anyIOException++ it "drains a connection when the graceful close timeout is positive" $ limited $ do+ -- The handler returns while the request it never read is+ -- still queued. 'gracefulClose' sends FIN and drains it.+ gate <- newEmptyMVar+ let set = defaultServerSettings{settingsGracefulCloseTimeout = 500}+ withTCPServer set (\_ -> takeMVar gate) $ \port ->+ client port $ \sock -> do+ sendAll sock "hello"+ threadDelay 200000+ putMVar gate ()+ recv sock 1024 `shouldReturn` ""++ it "resets a connection when the graceful close timeout is not positive" $ limited $ do+ -- The same, with 'close' instead: unread data in the queue+ -- makes the kernel answer with RST.+ gate <- newEmptyMVar+ let set = defaultServerSettings{settingsGracefulCloseTimeout = 0}+ withTCPServer set (\_ -> takeMVar gate) $ \port ->+ client port $ \sock -> do+ sendAll sock "hello"+ threadDelay 200000+ putMVar gate ()+ recv sock 1024 `shouldThrow` anyIOException++ it "serves many connections concurrently" $+ limited $+ withTCPServer defaultServerSettings echo $ \port -> do+ vars <- replicateM 50 newEmptyMVar+ forM_ vars $ \var -> forkIO $ do+ r <-+ (Just <$> request port "hello")+ `E.catch` \(_ :: E.SomeException) -> return Nothing+ putMVar var r+ rs <- mapM takeMVar vars+ rs `shouldBe` replicate 50 (Just "hello")++ describe "openTCPServerSocket" $ do+ it "listens, unlike openServerSocket" $ limited $ do+ addr <- resolve Stream (Just loopback) "0" [AI_PASSIVE] NE.head+ -- Probing with 'connect' instead would be slow: BSD drops+ -- the SYN sent to a socket which is bound but does not+ -- listen, where Linux answers with RST.+ E.bracket (openServerSocket addr) close $ \sock ->+ accept sock `shouldThrow` invalidArgument+ withTCPServer defaultServerSettings echo $ \port ->+ request port "hello" `shouldReturn` "hello"++ it "sets ReuseAddr" $ limited $ withListenSocket $ \lsock _ ->+ getSocketOption lsock ReuseAddr `shouldNotReturn` 0++ it "sets close-on-exec" $ limited $ withListenSocket $ \lsock _ ->+ withFdSocket lsock getCloseOnExec `shouldReturn` True++ it "sets a composite option" $ limited $ do+ addr <- resolve Stream (Just loopback) "0" [AI_PASSIVE] NE.head+ let opts = [(Linger, SockOptValue $ StructLinger 1 0)]+ E.bracket (openTCPServerSocketWithOpts opts addr) close $ \lsock -> do+ StructLinger onoff _ <- getSockOpt lsock Linger+ onoff `shouldBe` 1++ it "makes an IPv6 socket IPv6 only" $ limited $ onIPv6 $ do+ addr <- resolve Stream (Just "::") "0" [AI_PASSIVE] NE.head+ E.bracket (openTCPServerSocket addr) close $ \lsock -> do+ getSocketOption lsock IPv6Only `shouldNotReturn` 0+ port <- portOf <$> getSocketName lsock+ withServerThread (void $ runTCPServerWithSocket lsock echo) $ do+ -- The IPv4 loopback must not reach it.+ r <- E.try $ request port "hello"+ case r :: Either E.IOException ByteString of+ Left _ -> return ()+ Right _ -> expectationFailure "IPv4 reached an IPv6 only socket"++ it "can be asked for a dual stack socket" $+ limited $+ onIPv6 $+ if os == "openbsd"+ then pendingWith "OpenBSD always makes IPv6 sockets IPv6 only"+ else do+ addr <- resolve Stream (Just "::") "0" [AI_PASSIVE] NE.head+ E.bracket+ (openTCPServerSocketWithOptions [(IPv6Only, 0)] addr)+ close+ $ \lsock -> do+ port <- portOf <$> getSocketName lsock+ withServerThread+ (void $ runTCPServerWithSocket lsock echo)+ $ request port "hello" `shouldReturn` "hello"++ describe "runTCPClient" $ do+ it "opens the socket with settingsOpenClientSocket" $ limited $ do+ ref <- newIORef (0 :: Int)+ withTCPServer defaultServerSettings echo $ \port -> do+ let set =+ defaultSettings+ { settingsOpenClientSocket = \addr -> do+ atomicModifyIORef' ref $ \n -> (n + 1, ())+ openClientSocketWithOptions [(NoDelay, 1)] addr+ }+ r <- runTCPClientWithSettings set loopback (show port) $ \sock -> do+ sendAll sock "hello"+ recv sock 1024+ r `shouldBe` "hello"+ readIORef ref `shouldReturn` 1++ it "connects to what settingsSelectAddrInfo chose" $+ limited $+ withTCPServer defaultServerSettings echo $ \port -> do+ -- Resolving a port where nothing listens, and then+ -- replacing the address with the real one. Only a+ -- client which honours the selection can connect.+ let real = SockAddrInet port $ tupleToHostAddress (127, 0, 0, 1)+ set =+ defaultSettings+ { settingsSelectAddrInfo = \ais ->+ (NE.head ais){addrAddress = real}+ }+ wrong = show (port + 1)+ r <- runTCPClientWithSettings set loopback wrong $ \sock -> do+ sendAll sock "hello"+ recv sock 1024+ r `shouldBe` "hello"++ it "closes the socket when the action throws" $ limited $ do+ ref <- newIORef Nothing+ withTCPServer defaultServerSettings echo $ \port -> do+ let set =+ defaultSettings+ { settingsOpenClientSocket = \addr -> do+ sock <- openClientSocket addr+ writeIORef ref $ Just sock+ return sock+ }+ action = runTCPClientWithSettings set loopback (show port) $+ \_ -> E.throwIO $ userError "boom"+ (action :: IO ()) `shouldThrow` anyIOException+ msock <- readIORef ref+ case msock of+ Nothing -> expectationFailure "the socket was never opened"+ Just sock -> do+ r <- E.try $ getSocketOption sock ReuseAddr+ case r :: Either E.IOException Int of+ Left _ -> return ()+ Right _ -> expectationFailure "the client socket is still open"++ describe "resolve" $ do+ it "asks for a passive wildcard address" $ do+ addr <- resolve Stream Nothing "0" [AI_PASSIVE] NE.head+ addrSocketType addr `shouldBe` Stream+ case addrAddress addr of+ SockAddrInet _ ha -> ha `shouldBe` tupleToHostAddress (0, 0, 0, 0)+ SockAddrInet6 _ _ ha _ ->+ ha `shouldBe` tupleToHostAddress6 (0, 0, 0, 0, 0, 0, 0, 0)+ sa -> expectationFailure $ "unexpected address: " ++ show sa++ it "uses the given socket type" $ do+ addr <- resolve Datagram (Just loopback) "0" [] NE.head+ addrSocketType addr `shouldBe` Datagram++ it "returns what the selector chose" $ do+ let mark = SockAddrInet 1 $ tupleToHostAddress (1, 2, 3, 4)+ addr <- resolve Stream (Just loopback) "0" [] $ \ais ->+ (NE.head ais){addrAddress = mark}+ addrAddress addr `shouldBe` mark++----------------------------------------------------------------++-- | 'accept' on a socket which is not listening fails with @EINVAL@.+invalidArgument :: Selector E.IOException+invalidArgument e = ioeGetErrorType e == InvalidArgument++-- | A handler which fails the first connection and echoes the rest.+failFirst :: IORef Int -> Socket -> IO ()+failFirst ref sock = do+ n <- atomicModifyIORef' ref $ \n -> (n + 1, n)+ if n == 0 then E.throwIO (userError "boom") else echo sock++-- | One request and one response to the given host.+echoOn :: HostName -> PortNumber -> IO ByteString+echoOn host port = runTCPClient host (show port) $ \sock -> do+ sendAll sock "hello"+ recv sock 1024++-- | A TCP port which is free at the time of the call.+freeTCPPort :: IO PortNumber+freeTCPPort = withListenSocket $ \_ port -> return port++onIPv6 :: IO () -> IO ()+onIPv6 action = do+ ok <- hasIPv6+ if ok then action else pendingWith "no IPv6 on this machine"
+ test/TimeoutSpec.hs view
@@ -0,0 +1,96 @@+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE ScopedTypeVariables #-}++module TimeoutSpec (spec) where++import Control.Concurrent+import qualified Control.Exception as E+import Control.Monad+import Data.IORef+import Network.Socket+import Network.Socket.ByteString+import qualified System.TimeManager as T+import System.Timeout (timeout)+import Test.Hspec++import Network.Run.Core (ServerSettings (..), defaultServerSettings)+import qualified Network.Run.TCP.Timeout as Timeout++import Helper++spec :: Spec+spec = do+ describe "runTCPServer" $ do+ it "serves a connection" $+ limited $+ withTimeoutServer defaultServerSettings 2 echoServer $ \port ->+ request port "hello" `shouldReturn` "hello"++ it "kills a handler which exceeds the timeout" $ limited $ do+ let server _ _ sock = do+ threadDelay 8000000+ sendAll sock "late"+ withTimeoutServer defaultServerSettings 1 server $ \port ->+ client port $ \sock -> do+ sendAll sock "hello"+ (bs, ms) <- elapsed $ recv sock 1024+ -- The handler was killed, so the connection is+ -- closed rather than answered.+ bs `shouldBe` ""+ ms `shouldSatisfy` (< 5000)++ it "keeps a handler which tickles alive" $ limited $ do+ -- The handler lives 2.4 seconds, longer than the timeout,+ -- but it tickles every 300ms. The margin between the two+ -- is what a loaded machine may eat without the test+ -- becoming a lie, so it is kept wide.+ let server _ th sock = do+ replicateM_ 8 $ threadDelay 300000 >> T.tickle th+ sendAll sock "ok"+ withTimeoutServer defaultServerSettings 2 server $ \port ->+ client port $ \sock -> do+ sendAll sock "hello"+ recv sock 1024 `shouldReturn` "ok"++ it "keeps serving after a handler throws" $ limited $ do+ (set, getReports) <- collecting+ ref <- newIORef (0 :: Int)+ let server _ _ sock = do+ n <- atomicModifyIORef' ref $ \n -> (n + 1, n)+ if n == 0 then E.throwIO (userError "boom") else echo sock+ withTimeoutServer set 2 server $ \port -> do+ ignoreAny $ request port "hello"+ request port "hello" `shouldReturn` "hello"+ reports <- waitFor 1 getReports+ map fst reports `shouldSatisfy` all (/= Nothing)++ it "stops when the listening socket is closed" $ limited $ do+ done <- newEmptyMVar+ withListenSocket $ \lsock _ -> do+ void $+ forkFinally+ (Timeout.runTCPServerWithSocket 2 lsock echoServer)+ (putMVar done)+ threadDelay 100000+ close lsock+ r <- timeout 2000000 $ takeMVar done+ case r of+ Just (Left _) -> return ()+ Just (Right _) -> expectationFailure "the accept loop returned"+ Nothing -> expectationFailure "the accept loop did not stop"++----------------------------------------------------------------++echoServer :: Timeout.TimeoutServer ()+echoServer _ _ sock = echo sock++withTimeoutServer+ :: ServerSettings+ -> Int+ -> Timeout.TimeoutServer ()+ -> (PortNumber -> IO a)+ -> IO a+withTimeoutServer set tm server body = withListenSocket $ \lsock port ->+ withServerThread+ (Timeout.runTCPServerWithSocketAndSettings set tm lsock server)+ (body port)
+ test/UDPSpec.hs view
@@ -0,0 +1,142 @@+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE ScopedTypeVariables #-}++module UDPSpec (spec) where++import Control.Concurrent+import qualified Control.Exception as E+import Control.Monad+import Data.ByteString (ByteString)+import Data.IORef+import qualified Data.List.NonEmpty as NE+import Data.Maybe (isJust)+import Network.Socket+import Network.Socket.ByteString+import System.Timeout (timeout)+import Test.Hspec++import Network.Run.UDP++import Helper++spec :: Spec+spec = do+ describe "runUDPServer" $+ it "serves a datagram" $+ limited $ do+ port <- freeUDPPort+ let server sock = forever $ do+ (bs, peer) <- recvFrom sock 2048+ void $ sendTo sock bs peer+ withServerThread (runUDPServer (Just loopback) port server) $+ udpRequest loopback port "hello" `shouldReturn` Just "hello"++ describe "runUDPServerFork" $ do+ it "serves a datagram" $+ limited $+ withUDPServerFork defaultServerSettings [loopback] echoDatagram $ \port ->+ udpRequest loopback port "hello" `shouldReturn` Just "hello"++ it "returns at once when no host is given" $ limited $ do+ r <- timeout 1000000 $ runUDPServerFork [] "0" $ \_ _ -> return ()+ r `shouldBe` Just ()++ it "keeps serving after a handler throws" $ limited $ do+ (set, getReports) <- collecting+ ref <- newIORef (0 :: Int)+ let server sock bs = do+ n <- atomicModifyIORef' ref $ \n -> (n + 1, n)+ if n == 0+ then ioError $ userError "boom"+ else sendAll sock bs+ withUDPServerFork set [loopback] server $ \port -> do+ udpRequest loopback port "hello" `shouldReturn` Just "hello"+ reports <- waitFor 1 getReports+ case reports of+ [] -> expectationFailure "the failure was not reported"+ (mpeer, _) : _ -> mpeer `shouldSatisfy` isJust++ it "serves every given host" $ limited $ do+ ok <- hasIPv6+ if not ok+ then pendingWith "no IPv6 on this machine"+ else withUDPServerFork+ defaultServerSettings+ [loopback, "::1"]+ echoDatagram+ $ \port -> do+ udpRequest loopback port "v4" `shouldReturn` Just "v4"+ udpRequest "::1" port "v6" `shouldReturn` Just "v6"++ it "does not leak sockets" $+ limited $+ withUDPServerFork defaultServerSettings [loopback] echoDatagram $ \port -> do+ -- Warming up first, so that one-off descriptors of the+ -- runtime are not counted.+ replicateM_ 5 $ udpRequest loopback port "warm"+ threadDelay 200000+ mbefore <- openFds+ replicateM_ 40 $ udpRequest loopback port "hello"+ threadDelay 300000+ mafter <- openFds+ case (mbefore, mafter) of+ (Just n0, Just n1) ->+ n1 - n0 `shouldSatisfy` (< 10)+ _ -> pendingWith "file descriptors are not observable here"++ describe "runUDPClient" $+ it "closes the socket when the action throws" $+ limited $ do+ port <- freeUDPPort+ ref <- newIORef Nothing+ let action = runUDPClient loopback port $ \sock _ -> do+ writeIORef ref $ Just sock+ E.throwIO $ userError "boom"+ (action :: IO ()) `shouldThrow` anyIOException+ msock <- readIORef ref+ case msock of+ Nothing -> expectationFailure "the socket was never opened"+ Just sock -> do+ r <- E.try $ getSocketOption sock ReuseAddr+ case r :: Either E.IOException Int of+ Left _ -> return ()+ Right _ -> expectationFailure "the client socket is still open"++----------------------------------------------------------------++echoDatagram :: Socket -> ByteString -> IO ()+echoDatagram sock bs = sendAll sock bs++-- | A UDP port which is free at the time of the call. Unlike a TCP+-- server, a UDP server here binds the port itself, so it cannot be+-- discovered afterwards.+freeUDPPort :: IO ServiceName+freeUDPPort = do+ addr <- resolve Datagram (Just loopback) "0" [AI_PASSIVE] NE.head+ E.bracket (openServerSocket addr) close $ \sock ->+ show . portOf <$> getSocketName sock++withUDPServerFork+ :: ServerSettings+ -> [HostName]+ -> (Socket -> ByteString -> IO ())+ -> (ServiceName -> IO a)+ -> IO a+withUDPServerFork set hosts server body = do+ port <- freeUDPPort+ withServerThread (runUDPServerForkWithSettings set hosts port server) $ do+ threadDelay 200000+ body port++-- | Sending a datagram until a reply comes back. UDP may drop it, and+-- the server may not have bound its port yet.+udpRequest :: HostName -> ServiceName -> ByteString -> IO (Maybe ByteString)+udpRequest host port bs = runUDPClient host port $ \sock server -> go (20 :: Int) sock server+ where+ go 0 _ _ = return Nothing+ go n sock server = do+ void $ sendTo sock bs server+ mr <- timeout 200000 $ fst <$> recvFrom sock 2048+ case mr of+ Just r -> return $ Just r+ Nothing -> go (n - 1) sock server