http2-client 0.6.0.0 → 0.7.0.0
raw patch · 7 files changed
+188/−192 lines, 7 filesdep +deepseqdep +stmdep −connectiondep ~asyncdep ~bytestringdep ~containersPVP ok
version bump matches the API change (PVP)
Dependencies added: deepseq, stm
Dependencies removed: connection
Dependency ranges changed: async, bytestring, containers, http2, network, time, tls
API changes (from Hackage documentation)
- Network.HTTP2.Client.RawConnection: RawHttp2Connection :: (ByteString -> IO ()) -> (Int -> IO ByteString) -> IO () -> RawHttp2Connection
+ Network.HTTP2.Client.RawConnection: RawHttp2Connection :: ([ByteString] -> IO ()) -> (Int -> IO ByteString) -> IO () -> RawHttp2Connection
- Network.HTTP2.Client.RawConnection: [_sendRaw] :: RawHttp2Connection -> ByteString -> IO ()
+ Network.HTTP2.Client.RawConnection: [_sendRaw] :: RawHttp2Connection -> [ByteString] -> IO ()
Files
- app/Main.hs +1/−2
- http2-client.cabal +10/−9
- src/Network/HTTP2/Client.hs +55/−39
- src/Network/HTTP2/Client/Channels.hs +6/−108
- src/Network/HTTP2/Client/FrameConnection.hs +21/−14
- src/Network/HTTP2/Client/Helpers.hs +6/−3
- src/Network/HTTP2/Client/RawConnection.hs +89/−17
app/Main.hs view
@@ -221,7 +221,6 @@ runHttp2Client frameConn _encoderBufsize _decoderBufsize conf (throwTo parentThread) ignoreFallbackHandler withConn $ \conn -> do- linkAsyncs conn _addCredit (_incomingFlowControl conn) _initialWindowKick _ <- forkIO $ forever $ do updated <- _updateWindow $ _incomingFlowControl conn@@ -262,7 +261,7 @@ tlsParams = TLS.ClientParams { TLS.clientWantSessionResume = Nothing , TLS.clientUseMaxFragmentLength = Nothing- , TLS.clientServerIdentification = ("127.0.0.1", "")+ , TLS.clientServerIdentification = (_host, ByteString.pack $ show _port) , TLS.clientUseServerNameIndication = True , TLS.clientShared = def , TLS.clientHooks = def { TLS.onServerCertificate = \_ _ _ _ -> return []
http2-client.cabal view
@@ -1,5 +1,5 @@ name: http2-client-version: 0.6.0.0+version: 0.7.0.0 synopsis: A native HTTP2 client library. description: Please read the README.md at the homepage. homepage: https://github.com/lucasdicioccio/http2-client@@ -22,14 +22,15 @@ other-modules: Network.HTTP2.Client.Channels , Network.HTTP2.Client.Dispatch build-depends: base >= 4.7 && < 5- , async- , bytestring- , connection- , containers- , http2- , network- , time- , tls+ , async >= 2.1 && < 3+ , bytestring >= 0.10 && < 1+ , containers >= 0.5 && < 1+ , deepseq >= 1.4 && < 2+ , http2 >= 1.6 && < 2+ , network >= 2.6 && < 3+ , stm >= 2.4 && < 3+ , time >= 1.8 && < 2+ , tls >= 1.4 && < 2 default-language: Haskell2010 executable http2-client-exe
src/Network/HTTP2/Client.hs view
@@ -46,7 +46,7 @@ import Control.Concurrent.Async (Async, async, race, withAsync, link) import Control.Exception (bracket, throwIO, SomeException, catch)-import Control.Concurrent.MVar (newMVar, takeMVar, putMVar)+import Control.Concurrent.MVar (newEmptyMVar, newMVar, putMVar, takeMVar, tryPutMVar) import Control.Concurrent (threadDelay) import Control.Monad (forever, void, when, forM_) import Data.ByteString (ByteString)@@ -177,7 +177,8 @@ , _asyncs :: !Http2ClientAsyncs -- ^ Asynchronous operations threads. , _close :: IO ()- -- ^ Closes the network connection.+ -- ^ Immediately stop processing incoming frames and closes the network+ -- connection. } data InitHttp2Client = InitHttp2Client {@@ -189,6 +190,9 @@ , _initOutgoingFlowControl :: OutgoingFlowControl , _initPaylodSplitter :: IO PayloadSplitter , _initClose :: IO ()+ -- ^ Immediately closes the connection.+ , _initStop :: IO Bool+ -- ^ Stops receiving frames. } -- | Set of Async threads running an Http2Client.@@ -308,8 +312,9 @@ -- -- This function is slightly safer than 'startHttp2Client' because it uses -- 'Control.Concurrent.Async.withAsync' instead of--- 'Control.Concurrent.Async.async' however, this with-pattern takes the--- control of the thread and can be annoying at times.+-- 'Control.Concurrent.Async.async'; plus this function calls 'linkAsyncs' to+-- make sure that a network error kills the controlling thread. However, this+-- with-pattern takes the control of the thread and can be annoying at times. runHttp2Client :: Http2FrameConnection -- ^ A frame connection.@@ -332,17 +337,20 @@ withAsync incomingLoop $ \aIncoming -> do settsIO <- _initSettings initClient initSettings withAsync settsIO $ \aSettings -> do- mainHandler $ Http2Client {+ let client = Http2Client { _settings = _initSettings initClient , _ping = _initPing initClient , _goaway = _initGoaway initClient- , _close = _initClose initClient+ , _close =+ _initStop initClient >> _initClose initClient , _startStream = _initStartStream initClient , _incomingFlowControl = _initIncomingFlowControl initClient , _outgoingFlowControl = _initOutgoingFlowControl initClient , _payloadSplitter = _initPaylodSplitter initClient , _asyncs = Http2ClientAsyncs aSettings aIncoming }+ linkAsyncs client+ mainHandler client -- | Starts a new Http2Client around a frame connection. --@@ -372,7 +380,8 @@ _settings = _initSettings initClient , _ping = _initPing initClient , _goaway = _initGoaway initClient- , _close = _initClose initClient+ , _close =+ _initStop initClient >> _initClose initClient , _startStream = _initStartStream initClient , _incomingFlowControl = _initIncomingFlowControl initClient , _outgoingFlowControl = _initOutgoingFlowControl initClient@@ -404,7 +413,7 @@ (_initOutgoingFlowControl,windowUpdatesChan) <- newOutgoingFlowControl dispatchControl 0 dispatchHPACK <- newDispatchHPACKIO decoderBufSize- let incomingLoop = dispatchLoop conn dispatch dispatchControl windowUpdatesChan _initIncomingFlowControl dispatchHPACK+ (incomingLoop,endIncomingLoop) <- dispatchLoop conn dispatch dispatchControl windowUpdatesChan _initIncomingFlowControl dispatchHPACK {- Setup for client-initiated streams. -} conccurentStreams <- newIORef 0@@ -454,7 +463,10 @@ let _initPaylodSplitter = settingsPayloadSplitter <$> readSettings dispatchControl + let _initStop = endIncomingLoop+ let _initClose = closeConnection conn+ return (incomingLoop, InitHttp2Client{..}) initializeStream@@ -518,39 +530,43 @@ -> Chan (FrameHeader, FramePayload) -> IncomingFlowControl -> DispatchHPACK- -> IO ()+ -> IO (IO (), IO Bool) dispatchLoop conn d dc windowUpdatesChan inFlowControl dh = do let getNextFrame = next conn- delayException . forever $ do- frame <- getNextFrame- dispatchFramesStep frame d- whenFrame (hasStreamId 0) frame $ \got ->- dispatchControlFramesStep windowUpdatesChan got dc- whenFrame (hasTypeId [FrameData]) frame $ \got ->- creditDataFramesStep d inFlowControl got- whenFrame (hasTypeId [FrameWindowUpdate]) frame $ \got -> do- updateWindowsStep d got- whenFrame (hasTypeId [FramePushPromise, FrameHeaders]) frame $ \got -> do- let hpackLoop (FinishedWithHeaders curFh sId mkNewHdrs) = do- newHdrs <- mkNewHdrs- chan <- fmap _streamStateEvents <$> lookupStreamState d sId- let msg = StreamHeadersEvent curFh newHdrs- maybe (return ()) (flip writeChan msg) chan- hpackLoop (FinishedWithPushPromise curFh parentSid newSid mkNewHdrs) = do- newHdrs <- mkNewHdrs- chan <- fmap _streamStateEvents <$> lookupStreamState d parentSid- let msg = StreamPushPromiseEvent curFh newSid newHdrs- maybe (return ()) (flip writeChan msg) chan- hpackLoop (WaitContinuation act) =- getNextFrame >>= act >>= hpackLoop- hpackLoop (FailedHeaders curFh sId err) = do- chan <- fmap _streamStateEvents <$> lookupStreamState d sId- let msg = StreamErrorEvent curFh err- maybe (return ()) (flip writeChan msg) chan- hpackLoop (dispatchHPACKFramesStep got dh)- whenFrame (hasTypeId [FrameRSTStream]) frame $ \got -> do- handleRSTStep d got- finalizeFramesStep frame d+ let go = delayException . forever $ do+ frame <- getNextFrame+ dispatchFramesStep frame d+ whenFrame (hasStreamId 0) frame $ \got ->+ dispatchControlFramesStep windowUpdatesChan got dc+ whenFrame (hasTypeId [FrameData]) frame $ \got ->+ creditDataFramesStep d inFlowControl got+ whenFrame (hasTypeId [FrameWindowUpdate]) frame $ \got -> do+ updateWindowsStep d got+ whenFrame (hasTypeId [FramePushPromise, FrameHeaders]) frame $ \got -> do+ let hpackLoop (FinishedWithHeaders curFh sId mkNewHdrs) = do+ newHdrs <- mkNewHdrs+ chan <- fmap _streamStateEvents <$> lookupStreamState d sId+ let msg = StreamHeadersEvent curFh newHdrs+ maybe (return ()) (flip writeChan msg) chan+ hpackLoop (FinishedWithPushPromise curFh parentSid newSid mkNewHdrs) = do+ newHdrs <- mkNewHdrs+ chan <- fmap _streamStateEvents <$> lookupStreamState d parentSid+ let msg = StreamPushPromiseEvent curFh newSid newHdrs+ maybe (return ()) (flip writeChan msg) chan+ hpackLoop (WaitContinuation act) =+ getNextFrame >>= act >>= hpackLoop+ hpackLoop (FailedHeaders curFh sId err) = do+ chan <- fmap _streamStateEvents <$> lookupStreamState d sId+ let msg = StreamErrorEvent curFh err+ maybe (return ()) (flip writeChan msg) chan+ hpackLoop (dispatchHPACKFramesStep got dh)+ whenFrame (hasTypeId [FrameRSTStream]) frame $ \got -> do+ handleRSTStep d got+ finalizeFramesStep frame d+ end <- newEmptyMVar+ let run = void $ race go (takeMVar end)+ let stop = tryPutMVar end ()+ return (run, stop) handleRSTStep :: Dispatch
src/Network/HTTP2/Client/Channels.hs view
@@ -1,89 +1,28 @@ module Network.HTTP2.Client.Channels ( FramesChan- , HeadersChan- , PushPromisesChan- , waitHeaders- , waitHeadersWithStreamId- , waitFrame- , waitFrameWithStreamId- , waitPushPromiseWithParentStreamId- , waitFrameWithTypeId- , waitFrameWithTypeIdForStreamId- , isPingReply- , isSettingsReply , hasStreamId , hasTypeId , whenFrame , whenFrameElse+ -- re-exports , module Control.Concurrent.Chan ) where -import Control.Concurrent.Chan (Chan, readChan, newChan, dupChan, writeChan)+import Control.Concurrent.Chan (Chan, readChan, newChan, writeChan) import Control.Exception (Exception, throwIO)-import Data.ByteString (ByteString)-import Network.HPACK as HPACK-import Network.HTTP2 as HTTP2+import Network.HTTP2 (StreamId, FrameHeader, FramePayload, FrameTypeId, framePayloadToFrameTypeId, streamId) type FramesChan e = Chan (FrameHeader, Either e FramePayload) -type HeadersChanContent = (FrameHeader, StreamId, Either ErrorCode HeaderList)--type HeadersChan = Chan HeadersChanContent--type PushPromisesChanContent = (StreamId, StreamId, HeaderList)--type PushPromisesChan = Chan PushPromisesChanContent--waitFrameWithStreamId- :: Exception e- => StreamId- -> FramesChan e- -> IO (FrameHeader, FramePayload)-waitFrameWithStreamId sid = waitFrame (\h _ -> streamId h == sid)--waitFrameWithTypeId- :: (Exception e)- => [FrameTypeId]- -> FramesChan e- -> IO (FrameHeader, FramePayload)-waitFrameWithTypeId tids = waitFrame (\_ p -> HTTP2.framePayloadToFrameTypeId p `elem` tids)--waitFrameWithTypeIdForStreamId- :: (Exception e)- => StreamId- -> [FrameTypeId]- -> FramesChan e- -> IO (FrameHeader, FramePayload)-waitFrameWithTypeIdForStreamId sid tids =- waitFrame (\h p -> streamId h == sid && HTTP2.framePayloadToFrameTypeId p `elem` tids)--waitFrame- :: Exception e- => (FrameHeader -> FramePayload -> Bool)- -> FramesChan e- -> IO (FrameHeader, FramePayload)-waitFrame test chan =- loop- where- loop = do- (fHead, fPayload) <- readChan chan- dat <- either throwIO pure fPayload- if test fHead dat- then return (fHead, dat)- else loop- whenFrame :: Exception e => (FrameHeader -> FramePayload -> Bool) -> (FrameHeader, Either e FramePayload) -> ((FrameHeader, FramePayload) -> IO ()) -> IO ()-whenFrame test (fHead, fPayload) handle = do- dat <- either throwIO pure fPayload- if test fHead dat- then handle (fHead, dat)- else pure ()+whenFrame test frame handle = do+ whenFrameElse test frame handle (const $ pure ()) whenFrameElse :: Exception e@@ -102,45 +41,4 @@ hasStreamId sid h _ = streamId h == sid hasTypeId :: [FrameTypeId] -> FrameHeader -> FramePayload -> Bool-hasTypeId tids _ p = HTTP2.framePayloadToFrameTypeId p `elem` tids--isPingReply :: ByteString -> FrameHeader -> FramePayload -> Bool-isPingReply datSent _ (PingFrame datRcv) = datSent == datRcv-isPingReply _ _ _ = False--isSettingsReply :: FrameHeader -> FramePayload -> Bool-isSettingsReply fh (SettingsFrame _) = HTTP2.testAck (flags fh)-isSettingsReply _ _ = False--waitHeadersWithStreamId- :: StreamId- -> HeadersChan- -> IO HeadersChanContent-waitHeadersWithStreamId sid =- waitHeaders (\_ s _ -> s == sid)--waitHeaders- :: (FrameHeader -> StreamId -> Either ErrorCode HeaderList -> Bool)- -> HeadersChan- -> IO HeadersChanContent-waitHeaders test chan =- loop- where- loop = do- tuple@(fH, sId, hdrs) <- readChan chan- if test fH sId hdrs- then return tuple- else loop--waitPushPromiseWithParentStreamId- :: StreamId- -> PushPromisesChan- -> IO PushPromisesChanContent-waitPushPromiseWithParentStreamId sid chan =- loop- where- loop = do- tuple@(parentSid,_,_) <- readChan chan- if parentSid == sid- then return tuple- else loop+hasTypeId tids _ p = framePayloadToFrameTypeId p `elem` tids
src/Network/HTTP2/Client/FrameConnection.hs view
@@ -15,9 +15,10 @@ , closeConnection ) where +import Control.DeepSeq (deepseq) import Control.Exception (bracket) import Control.Concurrent.MVar (newMVar, takeMVar, putMVar)-import Control.Monad (void)+import Control.Monad ((>=>), void) import Network.HTTP2 (FrameHeader(..), FrameFlags, FramePayload, HTTP2Error, encodeInfo, decodeFramePayload) import qualified Network.HTTP2 as HTTP2 import Network.Socket (HostName, PortNumber)@@ -68,15 +69,11 @@ next :: Http2FrameConnection -> IO (FrameHeader, Either HTTP2Error FramePayload) next = _nextHeaderAndFrame . _serverStream --- | Creates a new 'Http2FrameConnection' to a given host for a frame-to-frame communication.-newHttp2FrameConnection :: HostName- -> PortNumber- -> Maybe TLS.ClientParams- -> IO Http2FrameConnection-newHttp2FrameConnection host port params = do- -- Spawns an HTTP2 connection.- http2conn <- newRawHttp2Connection host port params-+-- | Adds framing around a 'RawHttp2Connection'.+frameHttp2RawConnection+ :: RawHttp2Connection+ -> IO Http2FrameConnection+frameHttp2RawConnection http2conn = do -- Prepare a local mutex, this mutex should never escape the -- function's scope. Else it might lead to bugs (e.g., -- https://ro-che.info/articles/2014-07-30-bracket ) @@ -87,13 +84,15 @@ -- Define handlers. let makeClientStream streamID = - let putFrame modifyFF frame = do+ let putFrame modifyFF frame = let info = encodeInfo modifyFF streamID- _sendRaw http2conn $- HTTP2.encodeFrame info frame+ in HTTP2.encodeFrame info frame putFrames f = writeProtect . void $ do xs <- f- traverse (uncurry putFrame) xs+ let ys = fmap (uncurry putFrame) xs+ -- Force evaluation of frames serialization whilst+ -- write-protected to avoid out-of-order errrors.+ deepseq ys (_sendRaw http2conn ys) in Http2FrameClientStream putFrames streamID nextServerFrameChunk = Http2ServerStream $ do@@ -108,3 +107,11 @@ gtfo = _close http2conn return $ Http2FrameConnection makeClientStream nextServerFrameChunk gtfo++-- | Creates a new 'Http2FrameConnection' to a given host for a frame-to-frame communication.+newHttp2FrameConnection :: HostName+ -> PortNumber+ -> Maybe TLS.ClientParams+ -> IO Http2FrameConnection+newHttp2FrameConnection host port params = do+ frameHttp2RawConnection =<< newRawHttp2Connection host port params
src/Network/HTTP2/Client/Helpers.hs view
@@ -122,9 +122,12 @@ waitStream stream streamFlowControl ppHandler = do ev <- _waitEvent stream case ev of- StreamHeadersEvent _ hdrs -> do- (dfrms,trls) <- waitDataFrames []- return (Right hdrs, reverse dfrms, trls)+ StreamHeadersEvent fH hdrs+ | HTTP2.testEndStream (HTTP2.flags fH) -> do+ return (Right hdrs, [], Nothing)+ | otherwise -> do+ (dfrms,trls) <- waitDataFrames []+ return (Right hdrs, reverse dfrms, trls) StreamPushPromiseEvent _ ppSid ppHdrs -> do _handlePushPromise stream ppSid ppHdrs ppHandler waitStream stream streamFlowControl ppHandler
src/Network/HTTP2/Client/RawConnection.hs view
@@ -7,16 +7,22 @@ , newRawHttp2Connection ) where +import Control.Monad (forever, when)+import Control.Concurrent.Async (Async, async, cancel, pollSTM)+import Control.Concurrent.STM (STM, atomically, retry, throwSTM)+import Control.Concurrent.STM.TVar (TVar, modifyTVar', newTVarIO, readTVar, writeTVar) import Data.ByteString (ByteString)-import Network.Connection (connectTo, initConnectionContext, ConnectionParams(..), TLSSettings(..), connectionPut, connectionGetExact, connectionClose)+import qualified Data.ByteString as ByteString+import Data.ByteString.Lazy (fromChunks)+import Data.Monoid ((<>)) import qualified Network.HTTP2 as HTTP2-import Network.Socket (HostName, PortNumber)+import Network.Socket hiding (recv)+import Network.Socket.ByteString import qualified Network.TLS as TLS - -- TODO: catch connection errrors data RawHttp2Connection = RawHttp2Connection {- _sendRaw :: ByteString -> IO ()+ _sendRaw :: [ByteString] -> IO () -- ^ Function to send raw data to the server. , _nextRaw :: Int -> IO ByteString -- ^ Function to block reading a datachunk of a given size from the server.@@ -35,22 +41,88 @@ -- overwritten to always return ["h2", "h2-17"]. -> IO RawHttp2Connection newRawHttp2Connection host port mparams = do- -- Connects to SSL.- ctx <- initConnectionContext- conn <- connectTo ctx connParams+ -- Connects to TCP.+ let hints = defaultHints { addrFlags = [AI_NUMERICSERV], addrSocketType = Stream }+ addr:_ <- getAddrInfo (Just hints) (Just host) (Just $ show port)+ skt <- socket (addrFamily addr) (addrSocketType addr) (addrProtocol addr)+ connect skt (addrAddress addr) - -- Define raw byte-stream handlers.- let putRaw dat = connectionPut conn dat- let getRaw amount = connectionGetExact conn amount- let doClose = connectionClose conn+ -- Prepare structure with abstract API.+ conn <- maybe (plainTextRaw skt) (tlsRaw skt) mparams -- Initializes the HTTP2 stream.- putRaw HTTP2.connectionPreface+ _sendRaw conn [HTTP2.connectionPreface] - return $ RawHttp2Connection putRaw getRaw doClose+ return conn++plainTextRaw :: Socket -> IO RawHttp2Connection+plainTextRaw skt = do+ (b,putRaw) <- startWriteWorker (sendMany skt)+ (a,getRaw) <- startReadWorker (recv skt)+ let doClose = cancel a >> cancel b >> close skt+ return $ RawHttp2Connection (atomically . putRaw) (atomically . getRaw) doClose++tlsRaw :: Socket -> TLS.ClientParams -> IO RawHttp2Connection+tlsRaw skt params = do+ -- Connects to SSL+ tlsContext <- TLS.contextNew skt (modifyParams params)+ TLS.handshake tlsContext++ (b,putRaw) <- startWriteWorker (TLS.sendData tlsContext . fromChunks)+ (a,getRaw) <- startReadWorker (const $ TLS.recvData tlsContext)+ let doClose = cancel a >> cancel b >> TLS.bye tlsContext >> TLS.contextClose tlsContext++ return $ RawHttp2Connection (atomically . putRaw) (atomically . getRaw) doClose where- overwriteALPNHook params = (TLS.clientHooks params) {- TLS.onSuggestALPN = return $ Just [ "h2", "h2-17" ]+ modifyParams prms = prms {+ TLS.clientHooks = (TLS.clientHooks prms) {+ TLS.onSuggestALPN = return $ Just [ "h2", "h2-17" ]+ } }- modifyParams params = params { TLS.clientHooks = overwriteALPNHook params }- connParams = ConnectionParams host port (TLSSettings . modifyParams <$> mparams) Nothing++startWriteWorker+ :: ([ByteString] -> IO ())+ -> IO (Async (), [ByteString] -> STM ())+startWriteWorker sendChunks = do+ outQ <- newTVarIO []+ let putRaw chunks = modifyTVar' outQ (\xs -> xs ++ chunks)+ b <- async $ writeWorkerLoop outQ sendChunks+ return (b, putRaw)++writeWorkerLoop :: TVar [ByteString] -> ([ByteString] -> IO ()) -> IO ()+writeWorkerLoop outQ sendChunks = forever $ do+ xs <- atomically $ do+ chunks <- readTVar outQ+ when (null chunks) retry+ writeTVar outQ []+ return chunks+ sendChunks xs++startReadWorker+ :: (Int -> IO ByteString)+ -> IO (Async (), (Int -> STM ByteString))+startReadWorker get = do+ buf <- newTVarIO ""+ a <- async $ readWorkerLoop buf get+ return $ (a, getRawWorker a buf)++readWorkerLoop :: TVar ByteString -> (Int -> IO ByteString) -> IO ()+readWorkerLoop buf next = forever $ do+ dat <- next 4096+ atomically $ modifyTVar' buf (\bs -> (bs <> dat))++getRawWorker :: Async () -> TVar ByteString -> Int -> STM ByteString+getRawWorker a buf amount = do+ -- Verifies if the STM is alive, if dead, we re-throw the original+ -- exception.+ asyncStatus <- pollSTM a+ case asyncStatus of+ (Just (Left e)) -> throwSTM e+ _ -> return ()+ -- Read data consume, if there's enough, retry otherwise.+ dat <- readTVar buf+ if amount > ByteString.length dat+ then retry+ else do+ writeTVar buf (ByteString.drop amount dat)+ return $ ByteString.take amount dat