Z-IO 0.6.0.0 → 0.6.1.0
raw patch · 12 files changed
+203/−208 lines, 12 filesPVP: major bump suggested
API removals or changes: PVP suggests a major version bump
API changes (from Hackage documentation)
- Z.IO.BIO: sinkBuilderToOutput :: (Output o, HasCallStack) => o -> IO (Sink (Builder ()))
- Z.IO.BIO: sinkToOutput :: HasCallStack => Output o => o -> IO (Sink Bytes)
- Z.IO.BIO: sourceFromInput :: (HasCallStack, Input i) => i -> IO (Source Bytes)
- Z.IO.BIO: sourceJSONFromInput :: (HasCallStack, Input i, JSON a) => i -> IO (Source a)
- Z.IO.BIO: sourceParseChunksBufferedInput :: (HasCallStack, Print e) => ParseChunks IO Bytes e a -> BufferedInput -> Source a
- Z.IO.BIO: sourceParseChunksInput :: (Print e, Input i, HasCallStack) => ParseChunks IO Bytes e a -> i -> IO (Source a)
- Z.IO.BIO: sourceParserBufferInput :: HasCallStack => Parser a -> BufferedInput -> Source a
- Z.IO.BIO: sourceParserInput :: (Input i, HasCallStack) => Parser a -> i -> IO (Source a)
- Z.IO.BIO: sourceTextFromInput :: (HasCallStack, Input i) => i -> IO (Source Text)
+ Z.IO.BIO: sourceFromIO :: HasCallStack => IO (Maybe a) -> Source a
+ Z.IO.BIO: sourceParseChunksFromBuffered :: (HasCallStack, Print e) => ParseChunks IO Bytes e a -> BufferedInput -> Source a
+ Z.IO.BIO: sourceParserFromBuffered :: HasCallStack => Parser a -> BufferedInput -> Source a
+ Z.IO.Network.TCP: startServerLoop :: HasCallStack => Int -> (UVManager -> Resource UVStream) -> (Ptr UVHandle -> IO ()) -> (FD -> (UVStream -> IO ()) -> IO ()) -> (UVStream -> IO ()) -> IO ()
+ Z.IO.UV.FFI: hs_uv_try_write :: Ptr UVHandle -> Ptr Word8 -> Int -> IO Int
- Z.IO.BIO: (>!>) :: BIO a b -> (HasCallStack => b -> IO c) -> BIO a c
+ Z.IO.BIO: (>!>) :: HasCallStack => BIO a b -> (b -> IO c) -> BIO a c
- Z.IO.BIO: (>|>) :: BIO a b -> BIO b c -> BIO a c
+ Z.IO.BIO: (>|>) :: HasCallStack => BIO a b -> BIO b c -> BIO a c
- Z.IO.BIO: appendSource :: Source a -> Source a -> IO (Source a)
+ Z.IO.BIO: appendSource :: HasCallStack => Source a -> Source a -> IO (Source a)
- Z.IO.BIO: concatSource :: [Source a] -> IO (Source a)
+ Z.IO.BIO: concatSource :: HasCallStack => [Source a] -> IO (Source a)
- Z.IO.BIO: fuseSink :: [Sink out] -> Sink out
+ Z.IO.BIO: fuseSink :: HasCallStack => [Sink out] -> Sink out
- Z.IO.BIO: joinSink :: Sink out -> Sink out -> Sink out
+ Z.IO.BIO: joinSink :: HasCallStack => Sink out -> Sink out -> Sink out
- Z.IO.BIO: newBase64Decoder :: IO (BIO Bytes Bytes)
+ Z.IO.BIO: newBase64Decoder :: HasCallStack => IO (BIO Bytes Bytes)
- Z.IO.BIO: zipBIO :: BIO a b -> BIO a c -> IO (BIO a (b, c))
+ Z.IO.BIO: zipBIO :: HasCallStack => BIO a b -> BIO a c -> IO (BIO a (b, c))
- Z.IO.BIO: zipSource :: Source a -> Source b -> IO (Source (a, b))
+ Z.IO.BIO: zipSource :: HasCallStack => Source a -> Source b -> IO (Source (a, b))
Files
- ChangeLog.md +6/−0
- README.md +1/−1
- Z-IO.cabal +6/−2
- Z/IO/BIO.hs +55/−96
- Z/IO/Network/IPC.hs +12/−81
- Z/IO/Network/TCP.hs +41/−25
- Z/IO/Network/UDP.hs +7/−3
- Z/IO/UV/FFI.hsc +1/−0
- Z/IO/UV/UVStream.hs +21/−0
- cbits/hs_uv_stream.c +5/−0
- include/hs_uv.h +1/−0
- test/Z/IO/Network/IPCSpec.hs +47/−0
ChangeLog.md view
@@ -1,5 +1,11 @@ # Revision history for Z-IO +## 0.6.1.0 -- 2020-02-09++* Fix a bug in `newMagicSplitter/newLineSplitter` code.+* Remove `sourceFromInput` and related functions to reduce API surface, use `newBufferedInput` with `sourceFromBuffered` instead.+* Refactor server loop to allow more code sharing between `Z.IO.Network.TCP` and `Z.IO.Network.IPC`.+ ## 0.6.0.0 -- 2020-02-04 * FileSystem: replace `DEFAULT_MODE` with `DEFAULT_FILE_MODE` & `DEFAULT_DIR_MODE`.
README.md view
@@ -6,7 +6,7 @@ [](https://github.com/ZHaskell/z-io/actions) [](https://github.com/ZHaskell/z-io/actions) -This package is part of [Z.Haskell](https://github.com/ZHaskell/Z) project, provides basic IO operations:+This package is part of [Z.Haskell](https://z.haskell.world) project, provides basic IO operations: * IO resource management, resource pool * File system operations
Z-IO.cabal view
@@ -1,6 +1,6 @@ cabal-version: 2.4 name: Z-IO-version: 0.6.0.0+version: 0.6.1.0 synopsis: Simple and high performance IO toolkit for Haskell description: Simple and high performance IO toolkit for Haskell, including file system, network, ipc and more!@@ -379,10 +379,14 @@ Z.IO.FileSystemSpec Z.IO.Network.UDPSpec Z.IO.Network.TCPSpec- -- Z.IO.Network.IPCSpec+ Z.IO.Network.IPCSpec Z.IO.LowResTimerSpec Z.IO.ProcessSpec Z.IO.ResourceSpec+++ if os(linux)+ other-modules: Z.IO.Network.IPCSpec ghc-options: -threaded default-language: Haskell2010
Z/IO/BIO.hs view
@@ -59,21 +59,20 @@ -- * Make new BIO , pureBIO, ioBIO -- ** Source+ , sourceFromIO , sourceFromList , initSourceFromFile- , sourceFromBuffered, sourceFromInput- , sourceTextFromBuffered, sourceTextFromInput- , sourceJSONFromBuffered, sourceJSONFromInput- , sourceParserBufferInput, sourceParserInput- , sourceParseChunksBufferedInput, sourceParseChunksInput+ , sourceFromBuffered+ , sourceTextFromBuffered+ , sourceJSONFromBuffered+ , sourceParserFromBuffered+ , sourceParseChunksFromBuffered -- ** Sink+ , sinkToIO , sinkToList+ , initSinkToFile , sinkToBuffered , sinkBuilderToBuffered- , sinkToOutput- , initSinkToFile- , sinkBuilderToOutput- , sinkToIO -- ** Bytes specific , newParserNode, newReChunk, newUTF8Decoder, newMagicSplitter, newLineSplitter , newBase64Encoder, newBase64Decoder@@ -177,7 +176,7 @@ infixl 3 >~> -- | Connect two 'BIO' nodes, feed left one's output to right one's input.-(>|>) :: BIO a b -> BIO b c -> BIO a c+(>|>) :: HasCallStack => BIO a b -> BIO b c -> BIO a c {-# INLINE (>|>) #-} BIO pushA pullA >|> BIO pushB pullB = BIO push_ pull_ where@@ -200,7 +199,7 @@ (>~>) = flip fmap -- | Connect BIO to an effectful function.-(>!>) :: BIO a b -> (HasCallStack => b -> IO c) -> BIO a c+(>!>) :: HasCallStack => BIO a b -> (b -> IO c) -> BIO a c {-# INLINE (>!>) #-} (>!>) BIO{..} f = BIO push_ pull_ where@@ -212,21 +211,21 @@ _ -> return Nothing -- | Connect two 'BIO' source, after first reach EOF, draw element from second.-appendSource :: Source a -> Source a -> IO (Source a)+appendSource :: HasCallStack => Source a -> Source a -> IO (Source a) {-# INLINE appendSource #-} b1 `appendSource` b2 = concatSource [b1, b2] -- | Fuse two 'BIO' sinks, i.e. everything written to the fused sink will be written to left and right sink. -- -- Flush result 'BIO' will effectively flush both sink.-joinSink :: Sink out -> Sink out -> Sink out+joinSink :: HasCallStack => Sink out -> Sink out -> Sink out {-# INLINE joinSink #-} b1 `joinSink` b2 = fuseSink [b1, b2] -- | Fuse a list of 'BIO' sinks, everything written to the fused sink will be written to every sink in the list. -- -- Flush result 'BIO' will effectively flush every sink in the list.-fuseSink :: [Sink out] -> Sink out+fuseSink :: HasCallStack => [Sink out] -> Sink out {-# INLINABLE fuseSink #-} fuseSink ss = BIO push_ pull_ where@@ -234,7 +233,7 @@ pull_ = mapM_ pull ss >> return Nothing -- | Connect list of 'BIO' sources, after one reach EOF, draw element from next.-concatSource :: [Source a] -> IO (Source a)+concatSource :: HasCallStack => [Source a] -> IO (Source a) {-# INLINABLE concatSource #-} concatSource ss0 = newIORef ss0 >>= \ ref -> return (BIO{ pull = loop ref}) where@@ -249,7 +248,7 @@ _ -> writeIORef ref rest >> loop ref -- | Zip two 'BIO' source into one, reach EOF when either one reached EOF.-zipSource :: Source a -> Source b -> IO (Source (a,b))+zipSource :: HasCallStack => Source a -> Source b -> IO (Source (a,b)) {-# INLINABLE zipSource #-} zipSource (BIO _ pullA) (BIO _ pullB) = do finRef <- newIORef False@@ -269,7 +268,7 @@ -- | Zip two 'BIO' node into one, reach EOF when either one reached EOF. -- -- The output item number should match, unmatched output will be discarded.-zipBIO :: BIO a b -> BIO a c -> IO (BIO a (b, c))+zipBIO :: HasCallStack => BIO a b -> BIO a c -> IO (BIO a (b, c)) {-# INLINABLE zipBIO #-} zipBIO (BIO pushA pullA) (BIO pushB pullB) = do finRef <- newIORef False@@ -368,7 +367,7 @@ case r of Just _ -> loop f _ -> return () --- | Wrap a stream computation into a pure interface.+-- | Wrap 'runBlock' into a pure interface. -- -- You can wrap a stateful BIO computation(including the creation of 'BIO' node), -- when you can guarantee a computation is pure, e.g. compressing, decoding, etc.@@ -409,7 +408,7 @@ Just _ -> loop _ -> return () --- | Wrap a stream computation into a pure interface.+-- | Wrap 'runBlocks' into a pure interface. -- -- Similar to 'unsafeRunBlock', but with a list of input blocks. unsafeRunBlocks :: HasCallStack => IO (BIO inp out) -> [inp] -> [out]@@ -442,6 +441,11 @@ readBuffer i >>= \ x -> if V.null x then return Nothing else return (Just x)} +-- | Turn a `IO` action into 'Source'+sourceFromIO :: HasCallStack => IO (Maybe a) -> Source a+{-# INLINABLE sourceFromIO #-}+sourceFromIO io = BIO{ pull = io }+ -- | Turn a UTF8 encoded 'BufferedInput' into 'BIO' source, map EOF to Nothing. -- sourceTextFromBuffered :: HasCallStack => BufferedInput -> Source T.Text@@ -455,17 +459,17 @@ -- Throw 'OtherError' with name "EJSON" if JSON value is not parsed or converted. sourceJSONFromBuffered :: forall a. (JSON.JSON a, HasCallStack) => BufferedInput -> Source a {-# INLINABLE sourceJSONFromBuffered #-}-sourceJSONFromBuffered = sourceParseChunksBufferedInput JSON.decodeChunks+sourceJSONFromBuffered = sourceParseChunksFromBuffered JSON.decodeChunks -- | Turn buffered input device into a packet source, throw 'OtherError' with name @EPARSE@ if parsing fail.-sourceParserBufferInput :: HasCallStack => P.Parser a -> BufferedInput -> Source a-{-# INLINABLE sourceParserBufferInput #-}-sourceParserBufferInput p = sourceParseChunksBufferedInput (P.parseChunks p)+sourceParserFromBuffered :: HasCallStack => P.Parser a -> BufferedInput -> Source a+{-# INLINABLE sourceParserFromBuffered #-}+sourceParserFromBuffered p = sourceParseChunksFromBuffered (P.parseChunks p) -- | Turn buffered input device into a packet source, throw 'OtherError' with name @EPARSE@ if parsing fail.-sourceParseChunksBufferedInput :: (HasCallStack, T.Print e) => P.ParseChunks IO V.Bytes e a -> BufferedInput -> Source a-{-# INLINABLE sourceParseChunksBufferedInput #-}-sourceParseChunksBufferedInput cp bi = BIO{ pull = do+sourceParseChunksFromBuffered :: (HasCallStack, T.Print e) => P.ParseChunks IO V.Bytes e a -> BufferedInput -> Source a+{-# INLINABLE sourceParseChunksFromBuffered #-}+sourceParseChunksFromBuffered cp bi = BIO{ pull = do bs <- readBuffer bi if V.null bs then return Nothing@@ -475,39 +479,12 @@ case r of Right v -> return (Just v) Left e -> throwOtherError "EPARSE" (T.toText e) } --- | Turn an input device into a 'V.Bytes' source.-sourceFromInput :: (HasCallStack, Input i) => i -> IO (Source V.Bytes)-{-# INLINABLE sourceFromInput #-}-sourceFromInput i = sourceFromBuffered <$> newBufferedInput i---- | Turn an input device into a 'T.Text' source.-sourceTextFromInput :: (HasCallStack, Input i) => i -> IO (Source T.Text)-{-# INLINABLE sourceTextFromInput #-}-sourceTextFromInput i = sourceTextFromBuffered <$> newBufferedInput i---- | Turn an input device into a 'JSON' source.------ Throw 'OtherError' with name "EJSON" if JSON value is not parsed or converted.-sourceJSONFromInput :: (HasCallStack, Input i, JSON.JSON a) => i -> IO (Source a)-sourceJSONFromInput i = sourceJSONFromBuffered <$> newBufferedInput i-{-# INLINABLE sourceJSONFromInput #-}- -- | Turn a file into a 'V.Bytes' source. initSourceFromFile :: HasCallStack => CBytes -> Resource (Source V.Bytes) {-# INLINABLE initSourceFromFile #-} initSourceFromFile p = do f <- FS.initFile p FS.O_RDONLY FS.DEFAULT_FILE_MODE- liftIO (sourceFromInput f)---- | Turn input device into a packet source.-sourceParserInput :: (Input i, HasCallStack) => P.Parser a -> i -> IO (Source a)-{-# INLINABLE sourceParserInput #-}-sourceParserInput p i = sourceParserBufferInput p <$> newBufferedInput i---- | Turn input device into a packet source.-sourceParseChunksInput :: (T.Print e, Input i, HasCallStack) => P.ParseChunks IO V.Bytes e a -> i -> IO (Source a)-{-# INLINABLE sourceParseChunksInput #-}-sourceParseChunksInput p i = sourceParseChunksBufferedInput p <$> newBufferedInput i+ liftIO (sourceFromBuffered <$> newBufferedInput f) -------------------------------------------------------------------------------- -- Sink@@ -529,17 +506,6 @@ push_ inp = writeBuilder bo inp >> pure Nothing pull_ = flushBuffer bo >> pure Nothing --- | Turn an 'Output' into 'V,Bytes' sink.------ 'push' will write input to buffer, and 'pull'_ will flush buffer.-sinkToOutput :: HasCallStack => Output o => o -> IO (Sink V.Bytes)-{-# INLINABLE sinkToOutput #-}-sinkToOutput o =- newBufferedOutput o >>= \ bo -> return (BIO (push_ bo) (pull_ bo))- where- push_ bo inp = writeBuffer bo inp >> pure Nothing- pull_ bo = flushBuffer bo >> pure Nothing- -- | Turn a file into a 'V.Bytes' sink. -- -- Note the file will be opened in @'FS.O_APPEND' .|. 'FS.O_CREAT' .|. 'FS.O_WRONLY'@ mode,@@ -548,22 +514,11 @@ {-# INLINABLE initSinkToFile #-} initSinkToFile p = do f <- FS.initFile p (FS.O_APPEND .|. FS.O_CREAT .|. FS.O_WRONLY) FS.DEFAULT_FILE_MODE- liftIO (sinkToOutput f)---- | Turn an 'Output' into 'B.Builder' sink.------ 'push' will write input to buffer, and 'pull'_ will flush buffer.-sinkBuilderToOutput :: (Output o, HasCallStack) => o -> IO (Sink (B.Builder ()))-{-# INLINABLE sinkBuilderToOutput #-}-sinkBuilderToOutput o =- newBufferedOutput o >>= \ bo -> return (BIO (push_ bo) (pull_ bo))- where- push_ bo inp = writeBuilder bo inp >> pure Nothing- pull_ bo = flushBuffer bo >> pure Nothing+ liftIO (sinkToBuffered <$> newBufferedOutput f) --- | Turn an 'Output' into 'BIO' sink.+-- | Turn an `IO` action into 'BIO' sink. ----- 'push' will write input to buffer then perform flush, tend to degrade performance.+-- 'push' will call `IO` action with input chunk, `pull` has no effect. sinkToIO :: HasCallStack => (a -> IO ()) -> Sink a {-# INLINABLE sinkToIO #-} sinkToIO f = BIO push_ pull_@@ -682,11 +637,11 @@ -- -- If there're invalid UTF8 bytes, an 'OtherError' with name 'EINVALIDUTF8' will be thrown.` ----- Note this node is supposed to be used with preprocess node such as compressor, decoder, etc. where bytes--- boundary cannot be controlled, UTF8 decoder will concat trailing bytes from last block to next one.--- Use this node directly with 'sourceFromBuffered' \/ 'sourceFromInput' will not be as efficient as directly use--- 'sourceTextFromBuffered' \/ 'sourceTextFromInput', because 'BufferedInput' provides push back capability,--- trailing bytes can be pushde back to reading buffer and returned with next block input together.+-- Note this node is supposed to be used with preprocess node such as decompressor, parser, etc.+-- where bytes boundary cannot be controlled, UTF8 decoder will concat trailing bytes from last block to next one.+-- Use this node directly with 'sourceFromBuffered' will not be as efficient as directly use+-- 'sourceTextFromBuffered', because 'BufferedInput' provides push back capability,+-- trailing bytes can be pushed back to reading buffer then returned with next block input together. -- newUTF8Decoder :: HasCallStack => IO (BIO V.Bytes T.Text) {-# INLINABLE newUTF8Decoder #-}@@ -731,24 +686,28 @@ where push_ trailingRef bs = do trailing <- readIORef trailingRef- case V.elemIndex magic bs of+ let chunk = trailing `V.append` bs+ case V.elemIndex magic chunk of Just i -> do- let (!line, !rest) = V.splitAt (i+1) bs- !line' = trailing `V.append` line+ let (line, rest) = V.splitAt (i+1) chunk writeIORef trailingRef rest- return (Just line')- Nothing -> do- let !chunk = trailing `V.append` bs+ return (Just line)+ _ -> do writeIORef trailingRef chunk return Nothing pull_ trailingRef = do- trailing <- readIORef trailingRef- if V.null trailing+ chunk <- readIORef trailingRef+ if V.null chunk then return Nothing- else do- writeIORef trailingRef V.empty- return (Just trailing)+ else case V.elemIndex magic chunk of+ Just i -> do+ let (line, rest) = V.splitAt (i+1) chunk+ writeIORef trailingRef rest+ return (Just line)+ _ -> do+ writeIORef trailingRef V.empty+ return (Just chunk) -- | Make a new stream splitter based on linefeed(@\r\n@ or @\n@). --@@ -774,7 +733,7 @@ return (re >~> base64Encode) -- | Make a new base64 decoder node.-newBase64Decoder :: IO (BIO V.Bytes V.Bytes)+newBase64Decoder :: HasCallStack => IO (BIO V.Bytes V.Bytes) {-# INLINABLE newBase64Decoder #-} newBase64Decoder = do re <- newReChunk 4
Z/IO/Network/IPC.hs view
@@ -32,21 +32,18 @@ , initIPCStream ) where -import Control.Concurrent.MVar import Control.Monad import Control.Monad.IO.Class-import Data.Primitive.PrimArray-import Foreign.Ptr import GHC.Generics import Z.Data.CBytes import Z.Data.Text.Print (Print) import Z.Data.JSON (JSON) import Z.IO.Exception import Z.IO.Resource+import Z.IO.Network.TCP (startServerLoop) import Z.IO.UV.FFI import Z.IO.UV.Manager import Z.IO.UV.UVStream-import Data.Coerce -------------------------------------------------------------------------------- @@ -109,84 +106,18 @@ -- run in a seperated haskell thread, -- will be closed upon exception or worker finishes. -> IO ()-startIPCServer IPCServerConfig{..} ipcServerWorker = do- let backLog = max ipcListenBacklog 128- serverUVManager <- getUVManager- withResource (initIPCStream serverUVManager) $ \ (UVStream serverHandle serverSlot _ _) -> do+startIPCServer IPCServerConfig{..} = startServerLoop+ (max ipcListenBacklog 128)+ initIPCStream+ (\ serverHandle -> withCBytesUnsafe ipcListenName $ \ name_p -> do- throwUVIfMinus_ (uv_pipe_bind serverHandle name_p)- bracket- (do check <- throwOOMIfNull $ hs_uv_check_alloc- throwUVIfMinus_ (hs_uv_check_init check serverHandle)- return check)- hs_uv_check_close $- \ check -> do---- The buffer passing of accept is a litte complicated here, to get maximum performance,--- we do batch accepting. i.e. recv multiple client inside libuv's event loop:------ we poke uvmanager's buffer table as a Ptr Word8, with byte size (backLog*sizeof(FD))--- inside libuv event loop, we cast the buffer back to int32_t* pointer.--- each accept callback push a new socket fd to the buffer, and increase a counter(buffer_size_table).--- backLog should be large enough(>128), so under windows we can't possibly filled it up within one--- uv_run, under unix we hacked uv internal to provide a stop and resume function, when backLog is--- reached, we will stop receiving.------ once back to haskell side, we read all accepted sockets and fork worker threads.--- if backLog is reached, we resume receiving from haskell side.------ Step 1.--- we allocate a buffer to hold accepted FDs, pass it just like a normal reading buffer.--- then we can start listening.- acceptBuf <- newPinnedPrimArray backLog- let acceptBufPtr = coerce (mutablePrimArrayContents acceptBuf :: Ptr FD)-- withUVManager' serverUVManager $ do- -- We use buffersize as accepted fd count(count backwards)- pokeBufferTable serverUVManager serverSlot acceptBufPtr (backLog-1)- throwUVIfMinus_ (hs_uv_listen serverHandle (fromIntegral backLog))--- Step 2.--- we start a uv_check_t for given uv_stream_t, with predefined checking callback--- see hs_accept_check_cb in hs_uv_stream.c- throwUVIfMinus_ $ hs_uv_accept_check_start check-- m <- getBlockMVar serverUVManager serverSlot- forever $ do- -- wait until accept some FDs- _ <- takeMVar m--- Step 3.--- After uv loop finishes, if we got some FDs, copy the FD buffer, fetch accepted FDs and fork worker threads.-- -- we shouldn't receive asycn exceptions here otherwise accepted FDs are not closed- mask_$ do- -- we lock uv manager here in case of next uv_run overwrite current accept buffer- acceptBufCopy <- withUVManager' serverUVManager $ do- _ <- tryTakeMVar m- acceptCountDown <- peekBufferSizeTable serverUVManager serverSlot- pokeBufferSizeTable serverUVManager serverSlot (backLog-1)-- -- if acceptCountDown count to -1, we should resume on haskell side- when (acceptCountDown == -1) (hs_uv_listen_resume serverHandle)-- -- copy accepted FDs- let acceptCount = backLog - 1 - acceptCountDown- acceptBuf' <- newPrimArray acceptCount- copyMutablePrimArray acceptBuf' 0 acceptBuf (acceptCountDown+1) acceptCount- unsafeFreezePrimArray acceptBuf'-- -- fork worker thread- forM_ [0..sizeofPrimArray acceptBufCopy-1] $ \ i -> do- let fd = indexPrimArray acceptBufCopy i- if fd < 0- -- minus fd indicate a server error and we should close server- then throwUVIfMinus_ (return fd)- -- It's important to use the worker thread's mananger instead of server's one!- else void . forkBa $ do- uvm <- getUVManager- withResource (initUVStream (\ loop hdl -> do- throwUVIfMinus_ (uv_pipe_init loop hdl 0)- throwUVIfMinus_ (uv_pipe_open hdl fd)) uvm) $ \ uvs -> do- ipcServerWorker uvs+ throwUVIfMinus_ (uv_pipe_bind serverHandle name_p))+ ( \ fd worker -> void . forkBa $ do+ uvm <- getUVManager+ withResource (initUVStream (\ loop hdl -> do+ throwUVIfMinus_ (uv_pipe_init loop hdl 0)+ throwUVIfMinus_ (uv_pipe_open hdl fd)) uvm) $ \ uvs -> do+ worker uvs) --------------------------------------------------------------------------------
Z/IO/Network/TCP.hs view
@@ -27,12 +27,13 @@ , helloWorld , echo -- * Internal helper+ , startServerLoop , setTCPNoDelay , setTCPKeepAlive , initTCPStream ) where -import Control.Concurrent.MVar+import Control.Concurrent import Control.Monad import Control.Monad.IO.Class import Data.Primitive.PrimArray@@ -106,13 +107,13 @@ defaultTCPServerConfig :: TCPServerConfig defaultTCPServerConfig = TCPServerConfig (SocketAddrIPv4 ipv4Any 8888)- 128+ 256 True 30 --- | Start a server+-- | Start a TCP server ----- Fork new worker thread upon a new connection.+-- Fork new worker threads upon a new connection. -- startTCPServer :: HasCallStack => TCPServerConfig@@ -120,13 +121,39 @@ -- run in a seperated haskell thread, -- will be closed upon exception or worker finishes. -> IO ()-startTCPServer TCPServerConfig{..} tcpServerWorker = do- let backLog = max tcpListenBacklog 128+startTCPServer TCPServerConfig{..} = startServerLoop+ (max tcpListenBacklog 128)+ initTCPStream+ -- bind is safe without withUVManager+ (\ serverHandle -> withSocketAddrUnsafe tcpListenAddr $ \ addrPtr -> do+ throwUVIfMinus_ (uv_tcp_bind serverHandle addrPtr 0))+ (\ fd worker -> void . forkBa $ do+ -- It's important to use the worker thread's mananger instead of server's one!+ uvm <- getUVManager+ withResource (initUVStream (\ loop hdl -> do+ throwUVIfMinus_ (uv_tcp_init loop hdl)+ throwUVIfMinus_ (uv_tcp_open hdl fd)) uvm) $ \ uvs -> do+ -- safe without withUVManager+ when tcpServerWorkerNoDelay . throwUVIfMinus_ $+ uv_tcp_nodelay (uvsHandle uvs) 1+ when (tcpServerWorkerKeepAlive > 0) . throwUVIfMinus_ $+ uv_tcp_keepalive (uvsHandle uvs) 1 tcpServerWorkerKeepAlive+ worker uvs)++-- | Start a server loop with different kind of @uv_stream@s, such as tcp or pipe.+--+startServerLoop :: HasCallStack+ => Int -- ^ backLog+ -> (UVManager -> Resource UVStream) -- ^ uv_tream_t initializer+ -> (Ptr UVHandle -> IO ()) -- ^ bind function+ -> (FD -> (UVStream -> IO ()) -> IO ()) -- ^ thread spawner+ -> (UVStream -> IO ()) -- ^ worker+ -> IO ()+{-# INLINABLE startServerLoop #-}+startServerLoop backLog initStream bind spawn worker = do serverUVManager <- getUVManager- withResource (initTCPStream serverUVManager) $ \ (UVStream serverHandle serverSlot _ _) -> do- -- bind is safe without withUVManager- withSocketAddrUnsafe tcpListenAddr $ \ addrPtr -> do- throwUVIfMinus_ (uv_tcp_bind serverHandle addrPtr 0)+ withResource (initStream serverUVManager) $ \ (UVStream serverHandle serverSlot _ _) -> do+ bind serverHandle bracket (do check <- throwOOMIfNull $ hs_uv_check_alloc throwUVIfMinus_ (hs_uv_check_init check serverHandle)@@ -136,14 +163,14 @@ -- The buffer passing of accept is a litte complicated here, to get maximum performance, -- we do batch accepting. i.e. recv multiple client inside libuv's event loop: ----- we poke uvmanager's buffer table as a Ptr Word8, with byte size (backLog*sizeof(FD))+-- We poke uvmanager's buffer table like a normal Ptr Word8, with byte size (backLog*sizeof(FD)) -- inside libuv event loop, we cast the buffer back to int32_t* pointer. -- each accept callback push a new socket fd to the buffer, and increase a counter(buffer_size_table). -- backLog should be large enough(>128), so under windows we can't possibly filled it up within one -- uv_run, under unix we hacked uv internal to provide a stop and resume function, when backLog is -- reached, we will stop receiving. ----- once back to haskell side, we read all accepted sockets and fork worker threads.+-- Once back to haskell side, we read all accepted sockets and fork worker threads. -- if backLog is reached, we resume receiving from haskell side. -- -- Step 1.@@ -187,24 +214,13 @@ copyMutablePrimArray acceptBuf' 0 acceptBuf (acceptCountDown+1) acceptCount unsafeFreezePrimArray acceptBuf' - -- fork worker thread+ -- looping to fork worker threads forM_ [0..sizeofPrimArray acceptBufCopy-1] $ \ i -> do let fd = indexPrimArray acceptBufCopy i if fd < 0 -- minus fd indicate a server error and we should close server then throwUVIfMinus_ (return fd)- -- It's important to use the worker thread's mananger instead of server's one!- else void . forkBa $ do- uvm <- getUVManager- withResource (initUVStream (\ loop hdl -> do- throwUVIfMinus_ (uv_tcp_init loop hdl)- throwUVIfMinus_ (uv_tcp_open hdl fd)) uvm) $ \ uvs -> do- -- safe without withUVManager- when tcpServerWorkerNoDelay . throwUVIfMinus_ $- uv_tcp_nodelay (uvsHandle uvs) 1- when (tcpServerWorkerKeepAlive > 0) . throwUVIfMinus_ $- uv_tcp_keepalive (uvsHandle uvs) 1 tcpServerWorkerKeepAlive- tcpServerWorker uvs+ else spawn fd worker --------------------------------------------------------------------------------
Z/IO/Network/UDP.hs view
@@ -81,10 +81,10 @@ import Z.IO.Exception import Z.IO.Resource --- | UDP socket client.+-- | UDP socket. -- -- UDP is not a sequential protocol, thus not an instance of 'Input\/Output'.--- Message are received or sent individually, UDP socket client is NOT thread safe!+-- Message are received or sent individually, UDP socket is NOT thread safe! -- Use 'MVar' 'UDP' in multiple threads. -- data UDP = UDP@@ -365,7 +365,11 @@ -- | Recv UDP message within a loop ----- Loop receiving can be faster since it can reuse receiving buffer.+-- Loop receiving can be faster since it can reuse receiving buffer. Unlike TCP server+-- from "Z.IO.Network.TCP"server, UDP worker function is called on current haskell thread+-- instead of a forked one, if you have heavy computations to do within the worker function,+-- consider using 'forkBa', or a producer-consumer architecture+-- recvUDPLoop :: HasCallStack => UDPRecvConfig -> UDP
Z/IO/UV/FFI.hsc view
@@ -141,6 +141,7 @@ foreign import ccall unsafe hs_uv_read_start :: Ptr UVHandle -> IO CInt foreign import ccall unsafe uv_read_stop :: Ptr UVHandle -> IO CInt foreign import ccall unsafe hs_uv_write :: Ptr UVHandle -> Ptr Word8 -> Int -> IO UVSlotUnsafe+foreign import ccall unsafe hs_uv_try_write :: Ptr UVHandle -> Ptr Word8 -> Int -> IO Int foreign import ccall unsafe hs_uv_shutdown :: Ptr UVHandle -> IO UVSlotUnsafe foreign import ccall unsafe hs_uv_accept_check_start :: Ptr UVHandle -> IO CInt
Z/IO/UV/UVStream.hs view
@@ -161,6 +161,27 @@ -- OS will guarantee writing TTY and socket will not -- hang forever anyway. throwUVIfMinus_ (uninterruptibleMask_ $ takeMVar m)+ {- wait for https://github.com/libuv/libuv/pull/2874+ -- attempt blocking write first+ r <- hs_uv_try_write hdl buf len+ if | r == len -> return ()+ | r < 0 && r /= fromIntegral UV_EAGAIN -> throwUV r+ | otherwise -> do+ m <- withUVManager' uvm $ do+ reqSlot <- if r > 0+ then getUVSlot uvm (hs_uv_write hdl (buf `plusPtr` r) (len - r))+ else getUVSlot uvm (hs_uv_write hdl buf len)+ m <- getBlockMVar uvm reqSlot+ _ <- tryTakeMVar m+ return m+ -- we can't cancel uv_write_t with current libuv,+ -- otherwise disaster will happen if buffer got collected.+ -- so we have to turn to uninterruptibleMask_'s help.+ -- i.e. writing UVStream is an uninterruptible operation.+ -- OS will guarantee writing TTY and socket will not+ -- hang forever anyway.+ throwUVIfMinus_ (uninterruptibleMask_ $ takeMVar m)+ -} --------------------------------------------------------------------------------
cbits/hs_uv_stream.c view
@@ -93,6 +93,11 @@ } else return slot; } +HsInt hs_uv_try_write(uv_stream_t* handle, char* buf, HsInt buf_siz){+ uv_buf_t buf_t = { .base = buf, .len = (size_t)buf_siz };+ return (HsInt)uv_try_write(handle, &buf_t, 1);+}+ void hs_shutdown_cb(uv_shutdown_t* req, int status){ HsInt slot = (HsInt)req->data; uv_loop_t* loop = req->handle->loop;
include/hs_uv.h view
@@ -268,6 +268,7 @@ void hs_uv_listen_resume(uv_stream_t* server); int hs_uv_read_start(uv_stream_t* handle); HsInt hs_uv_write(uv_stream_t* handle, char* buf, HsInt buf_size);+HsInt hs_uv_try_write(uv_stream_t* handle, char* buf, HsInt buf_size); int hs_uv_accept_check_start(uv_check_t* check); HsInt hs_uv_shutdown(uv_stream_t* handle);
+ test/Z/IO/Network/IPCSpec.hs view
@@ -0,0 +1,47 @@+{-# LANGUAGE OverloadedStrings #-}++module Z.IO.Network.IPCSpec where++import Control.Concurrent+import Control.Monad+import Data.Bits+import Data.List as List+import Foreign.Marshal.Array+import Foreign.Ptr+import Test.HUnit+import Test.Hspec+import Z.Data.Vector as V+import Z.Data.Vector.Base as V+import Z.IO.Buffered+import Z.IO.Exception+import Z.IO.FileSystem (mkdtemp)+import Z.IO.Network+import Z.IO.Resource++spec :: Spec+spec = describe "IPC operations" $ do+ it "roundtrip test" $ do+ let testMsg = V.cycleN 256 "abc"+ longMsg = V.cycleN 2048 "abcdefg"+ tmpDir <- mkdtemp "z-io-test"+ let addr = tmpDir <> "socket-file"++ serverThread <- forkIO $ startIPCServer defaultIPCServerConfig{ ipcListenName = addr } echo++ threadDelay 1000000 -- 1s++ replicateM_ 10 . forkIO $+ withResource (initIPCClient defaultIPCClientConfig{ipcTargetName = addr}) $ \ ipc -> do+ i <- newBufferedInput ipc+ o <- newBufferedOutput ipc++ writeBuffer o testMsg >> flushBuffer o+ testMsg' <- readAll' i+ testMsg' @=? testMsg++ writeBuffer o longMsg >> flushBuffer o+ longMsg' <- readAll' i+ longMsg' @=? longMsg++ threadDelay 5000000 -- 5s+ killThread serverThread