packages feed

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 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