streaming-commons 0.1.0.2 → 0.1.1
raw patch · 4 files changed
+144/−8 lines, 4 filesdep +arraydep +asyncdep ~network
Dependencies added: array, async
Dependency ranges changed: network
Files
- Data/Streaming/Network.hs +103/−7
- Data/Streaming/Network/Internal.hs +1/−0
- streaming-commons.cabal +6/−1
- test/Data/Streaming/NetworkSpec.hs +34/−0
Data/Streaming/Network.hs view
@@ -14,6 +14,7 @@ #endif -- ** Smart constructors , serverSettingsTCP+ , serverSettingsTCPSocket , clientSettingsTCP , serverSettingsUDP , clientSettingsUDP@@ -52,10 +53,14 @@ -- * Functions -- ** General , bindPortGen+ , bindRandomPortGen , getSocketGen , acceptSafe+ , unassignedPorts+ , getUnassignedPort -- ** TCP , bindPortTCP+ , bindRandomPortTCP , getSocketTCP , safeRecv , runTCPServer@@ -64,6 +69,7 @@ , runTCPServerWithHandle -- ** UDP , bindPortUDP+ , bindRandomPortUDP , getSocketUDP #if !WINDOWS -- ** Unix@@ -89,6 +95,9 @@ import Data.Functor.Identity (Identity (Identity, runIdentity)) import Control.Concurrent (forkIO) import Control.Monad (forever)+import Data.IORef (IORef, newIORef, atomicModifyIORef)+import Data.Array.Unboxed ((!), UArray, bounds, listArray)+import System.IO.Unsafe (unsafePerformIO) #if WINDOWS import Control.Concurrent.MVar (putMVar, takeMVar, newEmptyMVar) #endif@@ -157,6 +166,62 @@ ) tryAddrs addrs' +-- | Bind to a random port number. Especially useful for writing network tests.+--+-- This will attempt 30 different port numbers before giving up and throwing an+-- exception.+--+-- Since 0.1.1+bindRandomPortGen :: SocketType -> HostPreference -> IO (Int, Socket)+bindRandomPortGen sockettype s = do+ loop 30+ where+ loop cnt | cnt <= 0 = error "Data.Streaming.Network.bindRandomPortGen: Could not get port"+ loop cnt = do+ port <- getUnassignedPort+ esocket <- try $ bindPortGen sockettype port s+ case esocket :: Either IOException Socket of+ Left _ -> loop (cnt - 1)+ Right socket -> return (port, socket)++-- | Top 10 Largest IANA unassigned port ranges with no unauthorized uses known+unassignedPortsList :: [Int]+unassignedPortsList = concat+ [ [43124..44320]+ , [28120..29166]+ , [45967..46997]+ , [28241..29117]+ , [40001..40840]+ , [29170..29998]+ , [38866..39680]+ , [43442..44122]+ , [41122..41793]+ , [35358..36000]+ ]++unassignedPorts :: UArray Int Int+unassignedPorts = listArray (unassignedPortsMin, unassignedPortsMax) unassignedPortsList++unassignedPortsMin, unassignedPortsMax :: Int+unassignedPortsMin = 1+unassignedPortsMax = length unassignedPortsList++nextUnusedPort :: IORef Int+nextUnusedPort = unsafePerformIO $ newIORef unassignedPortsMin+{-# NOINLINE nextUnusedPort #-}++-- | Get a port from the IANA list of unassigned ports.+--+-- Internally, this function uses an @IORef@ to cycle through the list of ports+getUnassignedPort :: IO Int+getUnassignedPort = do+ port <- atomicModifyIORef nextUnusedPort go+ return $! port+ where+ go i+ | i > unassignedPortsMax = (succ unassignedPortsMin, unassignedPorts ! unassignedPortsMin)+ | otherwise = (succ i, unassignedPorts ! i)+ -- | Attempt to connect to the given host/port. getSocketUDP :: String -> Int -> IO (Socket, AddrInfo) getSocketUDP = getSocketGen NS.Datagram@@ -166,6 +231,14 @@ bindPortUDP :: Int -> HostPreference -> IO Socket bindPortUDP = bindPortGen NS.Datagram +-- | Bind a random UDP port.+--+-- See 'bindRandomPortGen'+--+-- Since 0.1.1+bindRandomPortUDP :: HostPreference -> IO (Int, Socket)+bindRandomPortUDP = bindRandomPortGen NS.Datagram+ #if !WINDOWS -- | Attempt to connect to the given Unix domain socket path. getSocketUnix :: FilePath -> IO Socket@@ -250,10 +323,24 @@ serverSettingsTCP port host = ServerSettings { serverPort = port , serverHost = host+ , serverSocket = Nothing , serverAfterBind = const $ return () , serverNeedLocalAddr = False } +-- | Create a server settings that uses an already available listening socket.+-- Any port and host modifications made to this value will be ignored.+--+-- Since 0.1.1+serverSettingsTCPSocket :: Socket -> ServerSettings+serverSettingsTCPSocket lsocket = ServerSettings+ { serverPort = 0+ , serverHost = HostAny+ , serverSocket = Just lsocket+ , serverAfterBind = const $ return ()+ , serverNeedLocalAddr = False+ }+ -- | Smart constructor. clientSettingsUDP :: Int -- ^ port to connect to@@ -294,6 +381,17 @@ NS.listen sock (max 2048 NS.maxListenQueue) return sock +-- | Bind a random TCP port.+--+-- See 'bindRandomPortGen'.+--+-- Since 0.1.1+bindRandomPortTCP :: HostPreference -> IO (Int, Socket)+bindRandomPortTCP s = do+ (port, sock) <- bindRandomPortGen NS.Stream s+ NS.listen sock (max 2048 NS.maxListenQueue)+ return (port, sock)+ -- | Try to accept a connection, recovering automatically from exceptions. -- -- As reported by Kazu against Warp, "resource exhausted (Too many open files)"@@ -381,14 +479,12 @@ type ConnectionHandle = Socket -> NS.SockAddr -> Maybe NS.SockAddr -> IO () runTCPServerWithHandle :: ServerSettings -> ConnectionHandle -> IO ()-runTCPServerWithHandle (ServerSettings port host afterBind needLocalAddr) handle =- E.bracket- (bindPortTCP port host)- NS.sClose- (\socket -> do- afterBind socket- forever $ serve socket)+runTCPServerWithHandle (ServerSettings port host msocket afterBind needLocalAddr) handle =+ case msocket of+ Nothing -> E.bracket (bindPortTCP port host) NS.sClose inner+ Just lsocket -> inner lsocket where+ inner lsocket = afterBind lsocket >> forever (serve lsocket) serve lsocket = E.bracketOnError (acceptSafe lsocket) (\(socket, _) -> NS.sClose socket)
Data/Streaming/Network/Internal.hs view
@@ -21,6 +21,7 @@ data ServerSettings = ServerSettings { serverPort :: !Int , serverHost :: !HostPreference+ , serverSocket :: !(Maybe Socket) -- ^ listening socket , serverAfterBind :: !(Socket -> IO ()) , serverNeedLocalAddr :: !Bool }
streaming-commons.cabal view
@@ -1,5 +1,5 @@ name: streaming-commons-version: 0.1.0.2+version: 0.1.1 synopsis: Common lower-level functions needed by various streaming data libraries description: Provides low-dependency functionality commonly needed by various streaming data libraries, such as conduit and pipes. homepage: https://github.com/fpco/streaming-commons@@ -36,6 +36,7 @@ Data.Text.Internal.Encoding.Utf32 build-depends: base >= 4 && < 5+ , array , bytestring , directory , network@@ -62,6 +63,7 @@ ghc-options: -Wall -threaded other-modules: Data.Streaming.FileReadSpec Data.Streaming.FilesystemSpec+ Data.Streaming.NetworkSpec Data.Streaming.TextSpec Data.Streaming.ZlibSpec build-depends: base@@ -69,8 +71,11 @@ , hspec >= 1.8 , QuickCheck+ , array+ , async , bytestring , deepseq+ , network , text , zlib
+ test/Data/Streaming/NetworkSpec.hs view
@@ -0,0 +1,34 @@+{-# LANGUAGE OverloadedStrings #-}+module Data.Streaming.NetworkSpec where++import Control.Concurrent.Async (withAsync)+import Control.Exception (bracket)+import Control.Monad (forever, replicateM_)+import Data.Array.Unboxed (elems)+import qualified Data.ByteString.Char8 as S8+import Data.Char (toUpper)+import Data.Streaming.Network+import Network.Socket (sClose)+import Test.Hspec+import Test.Hspec.QuickCheck++spec :: Spec+spec = do+ describe "getUnassignedPort" $ do+ it "sanity" $ replicateM_ 100000 $ do+ port <- getUnassignedPort+ (port `elem` elems unassignedPorts) `shouldBe` True+ describe "bindRandomPortTCP" $ do+ modifyMaxSuccess (const 5) $ prop "sanity" $ \content -> bracket+ (bindRandomPortTCP "*4")+ (sClose . snd)+ $ \(port, socket) -> do+ let server ad = forever $ appRead ad >>= appWrite ad . S8.map toUpper+ client ad = do+ appWrite ad bs+ appRead ad >>= (`shouldBe` S8.map toUpper bs)+ bs+ | null content = "hello"+ | otherwise = S8.pack $ take 1000 content+ withAsync (runTCPServer (serverSettingsTCPSocket socket) server) $ \_ -> do+ runTCPClient (clientSettingsTCP port "127.0.0.1") client