http2 5.3.11 → 5.4.7
raw patch · 38 files changed
Files
- ChangeLog.md +177/−0
- Network/HPACK/HeaderBlock/Decode.hs +51/−23
- Network/HPACK/HeaderBlock/Encode.hs +20/−47
- Network/HPACK/HeaderBlock/Integer.hs +36/−6
- Network/HPACK/Huffman/Decode.hs +25/−10
- Network/HPACK/Huffman/Encode.hs +2/−2
- Network/HPACK/Table/Dynamic.hs +52/−26
- Network/HPACK/Types.hs +5/−2
- Network/HTTP2/Client.hs +14/−1
- Network/HTTP2/Client/Internal.hs +1/−0
- Network/HTTP2/Client/Run.hs +103/−19
- Network/HTTP2/Frame/Decode.hs +103/−39
- Network/HTTP2/H2/Config.hs +1/−0
- Network/HTTP2/H2/Context.hs +235/−43
- Network/HTTP2/H2/HPACK.hs +75/−5
- Network/HTTP2/H2/OutBodyIface.hs +130/−0
- Network/HTTP2/H2/Receiver.hs +416/−146
- Network/HTTP2/H2/Sender.hs +183/−61
- Network/HTTP2/H2/Settings.hs +3/−1
- Network/HTTP2/H2/Stream.hs +13/−77
- Network/HTTP2/H2/StreamTable.hs +25/−11
- Network/HTTP2/H2/Sync.hs +49/−12
- Network/HTTP2/H2/Types.hs +62/−44
- Network/HTTP2/H2/Window.hs +79/−22
- Network/HTTP2/Server.hs +12/−1
- Network/HTTP2/Server/Internal.hs +2/−0
- Network/HTTP2/Server/Run.hs +17/−7
- Network/HTTP2/Server/Worker.hs +75/−27
- bench-hpack/Main.hs +1/−1
- http2.cabal +11/−9
- test-hpack/HPACKDecode.hs +2/−2
- test/HPACK/DecodeSpec.hs +117/−0
- test/HPACK/EncodeSpec.hs +28/−0
- test/HPACK/HuffmanSpec.hs +4/−0
- test/HPACK/IntegerSpec.hs +32/−0
- test/HTTP2/ClientSpec.hs +8/−2
- test/HTTP2/FrameSpec.hs +55/−0
- test/HTTP2/ServerSpec.hs +1696/−419
ChangeLog.md view
@@ -1,5 +1,182 @@ # ChangeLog for http2 +## 5.4.7++* A valid request could close the whole connection, with every other+ stream on it:+ - a Huffman-coded field value longer than 4096 octets, such as a long+ token in `authorization`, was taken for a truncated block+ [#201](https://github.com/kazu-yamamoto/http2/pull/201);+ - so was a header block with no fields, which is what empty trailers+ are sent as [#202](https://github.com/kazu-yamamoto/http2/pull/202);+ - a malformed field (an upper-case name, a pseudo-header out of place,+ more than 200 fields) left the rest of its block undecoded, and the+ HPACK tables out of step. The block is now decoded to the end and the+ message refused with RST_STREAM(PROTOCOL_ERROR) on its stream alone+ (RFC 9113, section 8.1.1).+ [#211](https://github.com/kazu-yamamoto/http2/pull/211)+* Flow control lost octets, so that a long-lived connection could stall:+ - the padding of DATA frames was charged to both windows and never+ given back [#204](https://github.com/kazu-yamamoto/http2/pull/204);+ - DATA refused on a stream in the wrong state, or ignored on a stream we+ had reset, was not charged to the connection window, though the peer+ had charged it [#205](https://github.com/kazu-yamamoto/http2/pull/205);+ - on a client, the rest of a response that `processResponse` did not+ read to the end, or that it threw on, was never given back, and its+ stream held a slot of the server's SETTINGS_MAX_CONCURRENT_STREAMS.+ Such a stream is now reset with CANCEL.+ [#207](https://github.com/kazu-yamamoto/http2/pull/207)+ - on a client, the DATA of a push nobody asked for was held against the+ connection window for good. Pushes now give it back as it arrives.+ [#208](https://github.com/kazu-yamamoto/http2/pull/208)+* A padded body that matched its content-length was reset as malformed:+ the padding was counted into its length.+ [#203](https://github.com/kazu-yamamoto/http2/pull/203)+* A response carrying a push waited for ever when the client announced+ SETTINGS_MAX_CONCURRENT_STREAMS of 0, the way to refuse pushes. A push+ there is no room for is now not made.+ [#206](https://github.com/kazu-yamamoto/http2/pull/206)+* GOAWAY:+ - the last stream identifier of the server's GOAWAY left out streams+ whose handlers were still running, so a client could send again a+ request that had been acted on+ [#209](https://github.com/kazu-yamamoto/http2/pull/209);+ - a GOAWAY with NO_ERROR closed the connection at once, failing every+ stream in flight. Streams up to its last stream identifier now go+ on, those above it fail with `ConnectionIsClosed`, no new stream is+ opened, and the connection closes once nothing is left; on a client,+ the client function is let finish.+ [#210](https://github.com/kazu-yamamoto/http2/pull/210)+* With the connection window shut, nothing went out at all, though only+ DATA is flow-controlled: not the response to a request with no body,+ not RST_STREAM. DATA now waits for the window on its own.+ [#212](https://github.com/kazu-yamamoto/http2/pull/212)++## 5.4.6++* Security: a regression in 5.4.5. Since stream errors reset the stream+ rather than the connection, a peer could have the server reset streams+ for it -- with a PRIORITY on a stream depending on itself, DATA on a+ half-closed stream, and the like -- and so free concurrency slots while+ the handlers went on running, without ever sending RST_STREAM itself+ (MadeYouReset, CVE-2025-8671). Resets we send because of the peer now+ count against `rstRateLimit` with the peer's own.+ [#190](https://github.com/kazu-yamamoto/http2/pull/190)+* Security: a PRIORITY frame for a stream that was never opened created+ the stream and took a concurrency slot for good, so 64 PRIORITY frames+ were enough to have every later request refused.+ [#195](https://github.com/kazu-yamamoto/http2/pull/195)+* Security: a SETTINGS_INITIAL_WINDOW_SIZE that overflowed a stream's+ window stopped the sender without a word, leaving the connection open+ and silent. It is now a connection error of type FLOW_CONTROL_ERROR,+ and any failure of the sender closes the connection.+ [#196](https://github.com/kazu-yamamoto/http2/pull/196)+* The HPACK dynamic table lost entries, or had the encoder send the wrong+ one (index 61 of the static table), once it held as many entries as it+ has room for -- which a small or odd SETTINGS_HEADER_TABLE_SIZE from the+ peer makes easy. Headers were silently wrong on both sides.+ [#192](https://github.com/kazu-yamamoto/http2/pull/192)+* A Huffman-coded string of 16K or more was corrupted by the encoder: the+ length's fourth octet overwrote the start of the code.+ [#188](https://github.com/kazu-yamamoto/http2/pull/188)+* Header blocks and trailers larger than a frame are sent and received as+ HEADERS and CONTINUATION frames, and the header blocks of streams that+ are already reset are still decoded, so that the HPACK tables stay in+ step. Thanks to Edsko de Vries.+ [#187](https://github.com/kazu-yamamoto/http2/pull/187)+ [#189](https://github.com/kazu-yamamoto/http2/pull/189)+* A race between the receiver and the sender lost a stream's half-closed+ state, so that it was never removed from the stream table: with both+ ends streaming, a client ran out of streams and a server refused every+ new one.+ [#193](https://github.com/kazu-yamamoto/http2/pull/193)+* A client no longer rejects a response that has no content but a+ non-zero content-length, as responses to HEAD and 304 responses do.+ [#194](https://github.com/kazu-yamamoto/http2/pull/194)+* A client request that failed before it was queued -- a `requestFile` for+ a file that cannot be opened, say -- made every later request on the+ connection wait for ever.+ [#198](https://github.com/kazu-yamamoto/http2/pull/198)+* Server push: a PUSH_PROMISE could come after the response it belongs+ to, and pushed streams were never closed, so a connection stopped after+ 64 pushes.+ [#199](https://github.com/kazu-yamamoto/http2/pull/199)+* An upload through `runIO` larger than the stream's window was cut short+ with END_STREAM after the first window's worth.+ [#200](https://github.com/kazu-yamamoto/http2/pull/200)+* GHC 9.12 and later, with `-O`, miscompile a value holding a+ never-returning streaming body into one with no body+ ([GHC #27857](https://gitlab.haskell.org/ghc/ghc/-/work_items/27857)).+ The test suite works around it.+ [#197](https://github.com/kazu-yamamoto/http2/pull/197)++## 5.4.5++* Security: frame payload decoders read their fixed-size fields without+ checking that the payload holds them, so a truncated frame, or padding+ covering a field, read past the end of the buffer -- and an empty payload+ is the shared empty `ByteString`, whose pointer is null. An+ unauthenticated peer could segfault the process with 33 bytes.+ [#182](https://github.com/kazu-yamamoto/http2/pull/182)+* Security: HPACK integer decoding overflowed `Int` silently, so a long+ enough encoding decoded to whatever value the sender aimed at and two+ different byte strings could decode to the same header. Integers are now+ bounded and over-long encodings are a decoding error, as RFC 7541+ section 5.1 requires.+ [#181](https://github.com/kazu-yamamoto/http2/pull/181)+* A RST_STREAM gave a stream's concurrency slot back twice, so a peer could+ walk `SETTINGS_MAX_CONCURRENT_STREAMS` upwards and hold open as many+ streams as it liked.+ [#178](https://github.com/kazu-yamamoto/http2/pull/178)+* A stream reset while its response was still being produced left the+ worker blocked until the timeout manager killed it, one thread per reset+ stream.+ [#179](https://github.com/kazu-yamamoto/http2/pull/179)+* Stream errors now reset the stream and the connection carries on, as+ RFC 9113 section 5.4.2 requires. A field block abandoned part-way is+ still a connection error, since the HPACK tables have diverged by then.+ [#183](https://github.com/kazu-yamamoto/http2/pull/183)+* A stream over `SETTINGS_MAX_CONCURRENT_STREAMS` is refused with+ RST_STREAM(REFUSED_STREAM) rather than ending the connection.+ [#184](https://github.com/kazu-yamamoto/http2/pull/184)+* `DecodeError` has a new constructor, `TooLargeInteger`. Strictly this is+ a breaking change -- an exhaustive match on `DecodeError` no longer+ compiles -- but it ships as a patch version on purpose: no package on+ Hackage names any constructor of that type, while a minor bump would+ shut out every dependant carrying a `< 5.5` bound, these security fixes+ along with it.+* A malformed request now reaches a client as `StreamResetIsReceived` on+ the stream it concerns, where it used to arrive as+ `ConnectionErrorIsReceived` on the connection.++## 5.4.4++* Improvements for dealing with RST_STREAM+ [#172](https://github.com/kazu-yamamoto/http2/pull/172)++## 5.4.3++* auxSendInformational: gate usage with CPP to http-semantics >= 0.4.1+ [#170](https://github.com/kazu-yamamoto/http2/pull/170)++## 5.4.2++* Support informational (1xx) responses, e.g. 103 Early Hints. Servers can send+ them via `auxSendInformational`; clients can observe them via the new+ `confOnInformational` callback in `Config`.+ [#168](https://github.com/kazu-yamamoto/http2/pull/168)++## 5.4.1++* Ensure sender notices when receiver has terminated.+ [#167](https://github.com/kazu-yamamoto/http2/pull/167)++## 5.4.0++* Providing `defaultConfig`.+* Except the item above, this version is identical to v5.3.11 which+ includes breaking changes and is thus deprecated.+ ## 5.3.11 * Implementing `auxSendPing` for client.
Network/HPACK/HeaderBlock/Decode.hs view
@@ -14,7 +14,7 @@ decodeSimple, -- testing ) where -import Control.Exception (catch, throwIO)+import qualified Control.Exception as E import Data.Array.Base (unsafeRead, unsafeWrite) import qualified Data.Array.IO as IOA import qualified Data.Array.Unsafe as Unsafe@@ -37,7 +37,6 @@ -- -- * Headers are decoded as is. -- * 'DecodeError' would be thrown if the HPACK format is broken.--- * 'BufferOverrun' will be thrown if the temporary buffer for Huffman decoding is too small. decodeHeader :: DynamicTable -> ByteString@@ -57,15 +56,19 @@ -- 'IllegalHeaderName' is thrown. -- * If a header key contains capital letters, -- 'IllegalHeaderName' is thrown.+-- * If the number of header fields is too large,+-- 'TooLargeHeader' is thrown.+-- * 'IllegalHeaderName' and 'TooLargeHeader' are thrown only once the+-- whole block has been decoded, so that the dynamic table is up to+-- date: the message is malformed, not the block. -- * 'DecodeError' would be thrown if the HPACK format is broken.--- * 'BufferOverrun' will be thrown if the temporary buffer for Huffman decoding is too small. decodeTokenHeader :: DynamicTable -> ByteString -- ^ An HPACK format -> IO TokenHeaderTable decodeTokenHeader dyntbl inp =- decodeHPACK dyntbl inp (decodeSophisticated (toTokenHeader dyntbl)) `catch` \BufferOverrun -> throwIO HeaderBlockTruncated+ decodeHPACK dyntbl inp (decodeSophisticated (toTokenHeader dyntbl)) `E.catch` \BufferOverrun -> E.throwIO HeaderBlockTruncated decodeHPACK :: DynamicTable@@ -75,20 +78,27 @@ decodeHPACK dyntbl inp dec = withReadBuffer inp chkChange where chkChange rbuf = do- w <- read8 rbuf- if isTableSizeUpdate w- then do- tableSizeUpdate dyntbl w rbuf- chkChange rbuf+ -- A block can be empty, or hold nothing but table size updates:+ -- no fields, which is what empty trailers are sent as. Reading+ -- on regardless threw 'BufferOverrun', reported as a truncated+ -- block, so the connection was closed over them.+ leftover <- remainingSize rbuf+ if leftover < 1+ then dec rbuf else do- ff rbuf (-1)- dec rbuf+ w <- read8 rbuf+ if isTableSizeUpdate w+ then do+ tableSizeUpdate dyntbl w rbuf+ chkChange rbuf+ else do+ ff rbuf (-1)+ dec rbuf -- | Converting to '[Header]'. -- -- * Headers are decoded as is. -- * 'DecodeError' would be thrown if the HPACK format is broken.--- * 'BufferOverrun' will be thrown if the temporary buffer for Huffman decoding is too small. decodeSimple :: (Word8 -> ReadBuffer -> IO TokenHeader) -> ReadBuffer@@ -124,8 +134,10 @@ -- 'IllegalHeaderName' is thrown. -- * If the number of header fields is too large, -- 'TooLargeHeader' is thrown+-- * 'IllegalHeaderName' and 'TooLargeHeader' are thrown only once the+-- whole block has been decoded, so that the dynamic table is up to+-- date: the message is malformed, not the block. -- * 'DecodeError' would be thrown if the HPACK format is broken.--- * 'BufferOverrun' will be thrown if the temporary buffer for Huffman decoding is too small. decodeSophisticated :: (Word8 -> ReadBuffer -> IO TokenHeader) -> ReadBuffer@@ -150,34 +162,34 @@ then do mx <- unsafeRead arr tokenIx -- duplicated- when (isJust mx) $ throwIO IllegalHeaderName+ when (isJust mx) $ malformed IllegalHeaderName -- unknown- when (isMaxTokenIx tokenIx) $ throwIO IllegalHeaderName+ when (isMaxTokenIx tokenIx) $ malformed IllegalHeaderName unsafeWrite arr tokenIx (Just v) pseudo else do -- 0-Length Headers Leak - CVE-2019-9516- when (tokenKey == "") $ throwIO IllegalHeaderName+ when (tokenKey == "") $ malformed IllegalHeaderName when (isMaxTokenIx tokenIx && B8.any isUpper (original tokenKey)) $- throwIO IllegalHeaderName+ malformed IllegalHeaderName unsafeWrite arr tokenIx (Just v) if isCookieTokenIx tokenIx then normal 0 empty (empty << v) else normal 0 (empty << tv) empty else return [] normal n builder cookie- | n > headerLimit = throwIO TooLargeHeader+ | n > headerLimit = malformed TooLargeHeader | otherwise = do leftover <- remainingSize rbuf if leftover >= 1 then do w <- read8 rbuf tv@(Token{..}, v) <- decTokenHeader w rbuf- when isPseudo $ throwIO IllegalHeaderName+ when isPseudo $ malformed IllegalHeaderName -- 0-Length Headers Leak - CVE-2019-9516- when (tokenKey == "") $ throwIO IllegalHeaderName+ when (tokenKey == "") $ malformed IllegalHeaderName when (isMaxTokenIx tokenIx && B8.any isUpper (original tokenKey)) $- throwIO IllegalHeaderName+ malformed IllegalHeaderName unsafeWrite arr tokenIx (Just v) if isCookieTokenIx tokenIx then normal (n + 1) builder (cookie << v)@@ -193,11 +205,27 @@ unsafeWrite arr cookieTokenIx (Just v) return tvs + -- A field that makes the message malformed, as opposed to the block.+ -- The rest of the block is decoded all the same, and only then is the+ -- error thrown: every field of it may change the dynamic table, and one+ -- left undecoded leaves our table out of step with the peer's encoder,+ -- so that nothing after it on the connection decodes. So decoded, a+ -- malformed message can be refused on its own (RFC 9113, section 8.1.1:+ -- a stream error), rather than with the connection.+ malformed :: DecodeError -> IO a+ malformed err = skipRest >> E.throwIO err+ skipRest = do+ leftover <- remainingSize rbuf+ when (leftover >= 1) $ do+ w <- read8 rbuf+ _ <- decTokenHeader w rbuf+ skipRest+ toTokenHeader :: DynamicTable -> Word8 -> ReadBuffer -> IO TokenHeader toTokenHeader dyntbl w rbuf | w `testBit` 7 = indexed dyntbl w rbuf | w `testBit` 6 = incrementalIndexing dyntbl w rbuf- | w `testBit` 5 = throwIO IllegalTableSizeUpdate+ | w `testBit` 5 = E.throwIO IllegalTableSizeUpdate | w `testBit` 4 = neverIndexing dyntbl w rbuf | otherwise = withoutIndexing dyntbl w rbuf @@ -206,7 +234,7 @@ let w' = mask5 w siz <- decodeI 5 w' rbuf suitable <- isSuitableSize siz dyntbl- unless suitable $ throwIO TooLargeTableSize+ unless suitable $ E.throwIO TooLargeTableSize renewDynamicTable siz dyntbl ----------------------------------------------------------------
Network/HPACK/HeaderBlock/Encode.hs view
@@ -7,7 +7,6 @@ encodeS, ) where -import Control.Exception (bracket, throwIO) import qualified Control.Exception as E import qualified Data.ByteString as BS import Data.ByteString.Internal (create)@@ -68,13 +67,13 @@ -> TokenHeaderList -> IO ByteString -- ^ An HPACK format-encodeHeader' stgy siz dyntbl hs = bracket (mallocBytes siz) free enc+encodeHeader' stgy siz dyntbl hs = E.bracket (mallocBytes siz) free enc where enc buf = do (hs', len) <- encodeTokenHeader buf siz stgy True dyntbl hs case hs' of [] -> create len $ \p -> copyBytes p buf len- _ -> throwIO BufferOverrun+ _ -> E.throwIO BufferOverrun ---------------------------------------------------------------- @@ -298,49 +297,23 @@ -> IO ByteString encodeString h bs = withWriteBuffer 4096 $ \wbuf -> encStr wbuf h bs -{--N+ 1 2 3 <- bytes-8 254 382 16638-7 126 254 16510-6 62 190 16446-5 30 158 16414-4 14 142 16398-3 6 134 16390-2 2 130 16386-1 0 128 16384--}-+-- | The number of octets 'encodeI' produces for @l@ with an N-bit prefix.+--+-- 'encodeS' reserves this much before it knows the Huffman-coded length, and+-- moves the code if the guess was wrong, so it has to be exact. It used to+-- stop at three octets, which is enough only up to 2^N - 1 + 2^14 - 1:+-- a Huffman-coded string of 16K or more needs four, and 'encodeI' then wrote+-- the last of them over the first octet of the code.+--+-- >>> map (integerLength 7) [126, 127, 254, 255, 16510, 16511]+-- [1,2,2,3,3,4] {-# INLINE integerLength #-} integerLength :: Int -> Int -> Int-integerLength 8 l- | l <= 254 = 1- | l <= 382 = 2- | otherwise = 3-integerLength 7 l- | l <= 126 = 1- | l <= 254 = 2- | otherwise = 3-integerLength 6 l- | l <= 62 = 1- | l <= 190 = 2- | otherwise = 3-integerLength 5 l- | l <= 30 = 1- | l <= 158 = 2- | otherwise = 3-integerLength 4 l- | l <= 14 = 1- | l <= 142 = 2- | otherwise = 3-integerLength 3 l- | l <= 6 = 1- | l <= 134 = 2- | otherwise = 3-integerLength 2 l- | l <= 2 = 1- | l <= 130 = 2- | otherwise = 3-integerLength _ l- | l <= 0 = 1- | l <= 128 = 2- | otherwise = 3+integerLength n l+ | l < p = 1+ | otherwise = go 2 (l - p)+ where+ p = (1 `shiftL` n) - 1+ go k r+ | r < 128 = k+ | otherwise = go (k + 1) (r `shiftR` 7)
Network/HPACK/HeaderBlock/Integer.hs view
@@ -3,13 +3,16 @@ encodeInteger, decodeI, decodeInteger,+ integerLimit, ) where +import qualified Control.Exception as E import Data.Array (Array, listArray) import Data.Array.Base (unsafeAt) import Network.ByteOrder import Imports+import Network.HPACK.Types (DecodeError (..)) -- $setup -- >>> import qualified Data.ByteString as BS@@ -127,9 +130,36 @@ p = powerArray `unsafeAt` (n - 1) i = fromIntegral w decode :: Int -> Int -> IO Int- decode m j = do- b <- fromIntegral <$> read8 rbuf- let j' = j + (b .&. 0x7f) * 2 ^ m- m' = m + 7- cont = b `testBit` 7- if cont then decode m' j' else return j'+ decode m j+ -- Checked before the shift rather than after: shifting an 'Int' by a+ -- word width or more is not defined to give zero, and the value would+ -- have wrapped long before there were anything to notice.+ | m > maxShift = E.throwIO TooLargeInteger+ | otherwise = do+ b <- fromIntegral <$> read8 rbuf+ let d = b .&. 0x7f+ -- d * 2^m > integerLimit - j, without evaluating the product.+ when (d > (integerLimit - j) `shiftR` m) $ E.throwIO TooLargeInteger+ let j' = j + (d `shiftL` m)+ if b `testBit` 7 then decode (m + 7) j' else return j'++-- | The largest integer 'decodeI' will return.+--+-- HPACK's integer encoding carries no bound of its own, so a decoder has to+-- impose one. RFC 7541, section 5.1: "Integer encodings that exceed+-- implementation limits -- in value or octet length -- MUST be treated as+-- decoding errors."+--+-- 2^30 - 1 is far above anything HTTP\/2 can ask for -- a frame payload is at+-- most 2^24 - 1 octets, so no length or index comes near it -- and it still+-- fits in an 'Int' on a platform where that is 32 bits wide.+--+-- >>> integerLimit+-- 1073741823+integerLimit :: Int+integerLimit = 1073741823++-- | The largest shift that can carry a continuation octet into+-- 'integerLimit'; past it every further octet is an overflow.+maxShift :: Int+maxShift = 28
Network/HPACK/Huffman/Decode.hs view
@@ -9,7 +9,7 @@ GCBuffer, ) where -import Control.Exception (throwIO)+import qualified Control.Exception as E import Data.Array (Array, listArray) import Data.Array.Base (unsafeAt) import qualified Data.ByteString as BS@@ -59,26 +59,41 @@ -> Int -- ^ The target length -> IO ByteString-decodeH gcbuf bufsiz rbuf len = withForeignPtr gcbuf $ \buf -> do- wbuf <- newWriteBuffer buf bufsiz- decH wbuf rbuf len- toByteString wbuf+decodeH gcbuf bufsiz rbuf len+ -- The working space is only a cache. A value that may not fit gets a+ -- buffer of its own: running out of room part-way used to throw+ -- 'BufferOverrun', which the header block decoder reported as a+ -- truncated block, so a valid field longer than the working space+ -- (a long Huffman-coded @authorization@, say) closed the connection.+ | maxDecodedLength len > bufsiz =+ withWriteBuffer (maxDecodedLength len) $ \wbuf -> decH wbuf rbuf len+ | otherwise = withForeignPtr gcbuf $ \buf -> do+ wbuf <- newWriteBuffer buf bufsiz+ decH wbuf rbuf len+ toByteString wbuf +-- | The longest a Huffman-coded string of this many octets can decode to.+--+-- The shortest code is 5 bits long (RFC 7541, Appendix B), so each+-- decoded octet takes at least 5 of the input's bits.+maxDecodedLength :: Int -> Int+maxDecodedLength len = len * 8 `div` 5+ -- | Low devel Huffman decoding in a write buffer. decH :: WriteBuffer -> ReadBuffer -> Int -> IO () decH wbuf rbuf len = go len (way256 `unsafeAt` 0) where go 0 way0 = case way0 of- WayStep Nothing _ -> throwIO IllegalEos+ WayStep Nothing _ -> E.throwIO IllegalEos WayStep (Just i) _ | i <= 8 -> return ()- | otherwise -> throwIO TooLongEos+ | otherwise -> E.throwIO TooLongEos go n way0 = do w <- read8 rbuf way <- doit way0 w go (n - 1) way doit way w = case next way w of- EndOfString -> throwIO EosInTheMiddle+ EndOfString -> E.throwIO EosInTheMiddle Forward n -> return $ way256 `unsafeAt` fromIntegral n GoBack n v -> do write8 wbuf v@@ -88,9 +103,9 @@ write8 wbuf v2 return $ way256 `unsafeAt` fromIntegral n --- | Huffman decoding with a temporary buffer whose size is 4096.+-- | Huffman decoding. decodeHuffman :: ByteString -> IO ByteString-decodeHuffman bs = withWriteBuffer 4096 $ \wbuf ->+decodeHuffman bs = withWriteBuffer (max 1 $ maxDecodedLength $ BS.length bs) $ \wbuf -> withReadBuffer bs $ \rbuf -> decH wbuf rbuf $ BS.length bs ----------------------------------------------------------------
Network/HPACK/Huffman/Encode.hs view
@@ -6,7 +6,7 @@ encodeHuffman, ) where -import Control.Exception (throwIO)+import qualified Control.Exception as E import Data.Array.Base (unsafeAt) import Data.Array.IArray (listArray) import Data.Array.Unboxed (UArray)@@ -75,7 +75,7 @@ off' = off - len {-# INLINE write #-} write p w = do- when (p >= limit) $ throwIO BufferOverrun+ when (p >= limit) $ E.throwIO BufferOverrun let w8 = fromIntegral (w `shiftR` shiftForWrite) :: Word8 poke p w8 let p' = p `plusPtr` 1
Network/HPACK/Table/Dynamic.hs view
@@ -24,7 +24,7 @@ getRevIndex, ) where -import Control.Exception (throwIO)+import qualified Control.Exception as E import Data.Array.Base (unsafeRead, unsafeWrite) import Data.Array.IO (IOArray, newArray) import qualified Data.ByteString.Char8 as BS@@ -43,7 +43,7 @@ {-# INLINE toIndexedEntry #-} toIndexedEntry :: DynamicTable -> Index -> IO Entry toIndexedEntry dyntbl idx- | idx <= 0 = throwIO $ IndexOverrun idx+ | idx <= 0 = E.throwIO $ IndexOverrun idx | idx <= staticTableSize = return $ toStaticEntry idx | otherwise = toDynamicEntry dyntbl idx @@ -55,7 +55,13 @@ maxN <- readIORef maxNumOfEntries off <- readIORef offset x <- adj maxN (didx - off)- return $ x + staticTableSize+ -- Entries sit at off+1 .. off+n, so the relative position is 1 .. n.+ -- When the ring is full, n is maxN and the oldest entry is at off+maxN,+ -- which is off itself: the modulus makes that 0 rather than maxN, and 0+ -- is index 61 of the static table. 'toDynamicEntry', going the other+ -- way, lands on the right slot either way.+ let x' = if x == 0 then maxN else x+ return $ x' + staticTableSize ---------------------------------------------------------------- @@ -121,7 +127,7 @@ {-# INLINE adj #-} adj :: Int -> Int -> IO Int adj maxN x- | maxN == 0 = throwIO TooSmallTableSize+ | maxN == 0 = E.throwIO TooSmallTableSize | otherwise = let ret = (x + maxN) `mod` maxN in return ret@@ -216,7 +222,8 @@ :: Size -- ^ The dynamic table size -> Size- -- ^ The size of temporary buffer for Huffman decoding+ -- ^ The size of temporary buffer for Huffman decoding.+ -- A longer value is decoded in a buffer of its own. -> IO DynamicTable newDynamicTableForDecoding maxsiz huftmpsiz = do lim <- newIORef maxsiz@@ -303,7 +310,8 @@ :: Size -- ^ The dynamic table size -> Size- -- ^ The size of temporary buffer for Huffman+ -- ^ The size of temporary buffer for Huffman decoding.+ -- A longer value is decoded in a buffer of its own. -> (DynamicTable -> IO a) -> IO a withDynamicTableForDecoding maxsiz huftmpsiz action =@@ -312,18 +320,49 @@ ---------------------------------------------------------------- -- | Inserting 'Entry' to 'DynamicTable'.+--+-- Entries are evicted first and the new one added after, as RFC 7541+-- section 4.4 has it: "Before a new entry is added to the dynamic table,+-- entries are evicted from the end of the dynamic table until the size of+-- the dynamic table is less than or equal to (maximum size - new entry+-- size) or until the table is empty." An entry larger than the table+-- empties it and is not added.+--+-- The order matters to the ring. It has room for maxNumbers entries, and+-- the table holds that many whenever they are all close to the 32-octet+-- minimum -- at a size of 40 or 100, say. Added first, the new entry+-- landed on the oldest one's slot, and the eviction that followed read that+-- slot back and took out the new entry instead, leaving a dummy. After+-- evicting there is always a free slot: every entry is 32 octets or more,+-- so the entries left and the new one come to at most maxNumbers. insertEntry :: Entry -> DynamicTable -> IO () insertEntry e dyntbl@DynamicTable{..} = do- -- Theoretically speaking, dropping entries by adjustTableSize- -- should be first. However, non-used slots always exist since the- -- size of dynamic table calculated via the minimum entry size (32- -- bytes). To simply adjustTableSize, insertFront is called first.- insertFront e dyntbl- es <- adjustTableSize dyntbl+ es <- evictFor (entrySize e) dyntbl+ -- Before the new entry goes in: the reverse index is keyed by name and+ -- value, so an evicted entry equal to the new one would take its+ -- mapping out with it. case codeInfo of CIE (EncodeInfo rev _) -> deleteRevIndexList es rev _ -> return ()+ maxdsize <- readIORef maxDynamicTableSize+ when (entrySize e <= maxdsize) $ insertFront e dyntbl +-- | Evicting entries until one of the given size fits, or the table is+-- empty.+evictFor :: Size -> DynamicTable -> IO [Entry]+evictFor siz dyntbl@DynamicTable{..} = evict []+ where+ evict :: [Entry] -> IO [Entry]+ evict es = do+ n <- readIORef numOfEntries+ dsize <- readIORef dynamicTableSize+ maxdsize <- readIORef maxDynamicTableSize+ if n == 0 || dsize + siz <= maxdsize+ then return es+ else do+ e <- removeEnd dyntbl+ evict (e : es)+ insertFront :: Entry -> DynamicTable -> IO () insertFront e DynamicTable{..} = do maxN <- readIORef maxNumOfEntries@@ -345,19 +384,6 @@ CIE (EncodeInfo rev _) -> insertRevIndex e (DIndex i) rev _ -> return () -adjustTableSize :: DynamicTable -> IO [Entry]-adjustTableSize dyntbl@DynamicTable{..} = adjust []- where- adjust :: [Entry] -> IO [Entry]- adjust es = do- dsize <- readIORef dynamicTableSize- maxdsize <- readIORef maxDynamicTableSize- if dsize <= maxdsize- then return es- else do- e <- removeEnd dyntbl- adjust (e : es)- ---------------------------------------------------------------- -- Used in copyEntries.@@ -402,7 +428,7 @@ maxN <- readIORef maxNumOfEntries off <- readIORef offset n <- readIORef numOfEntries- when (idx > n + staticTableSize) $ throwIO $ IndexOverrun idx+ when (idx > n + staticTableSize) $ E.throwIO $ IndexOverrun idx didx <- adj maxN (idx + off - staticTableSize) table <- readIORef circularTable unsafeRead table didx
Network/HPACK/Types.hs view
@@ -20,7 +20,7 @@ BufferOverrun (..), ) where -import Control.Exception as E+import qualified Control.Exception as E import Network.ByteOrder (Buffer, BufferOverrun (..), BufferSize) import Imports@@ -76,6 +76,9 @@ IllegalEos | -- | Eos of huffman string is more than 7 bits TooLongEos+ | -- | An integer is encoded above the limit this decoder accepts,+ -- or in more octets than reaching that limit can take+ TooLargeInteger | -- | A peer set the dynamic table size less than 32 TooSmallTableSize | -- | A peer tried to change the dynamic table size over the limit@@ -87,4 +90,4 @@ | TooLargeHeader deriving (Eq, Show) -instance Exception DecodeError+instance E.Exception DecodeError
Network/HTTP2/Client.hs view
@@ -71,7 +71,18 @@ rstRateLimit, -- * Common configuration- Config (..),+ Config,+ defaultConfig,+ confWriteBuffer,+ confBufferSize,+ confSendAll,+ confReadN,+ confPositionReadMaker,+ confTimeoutManager,+ confMySockAddr,+ confPeerSockAddr,+ confReadNTimeout,+ confOnInformational, allocSimpleConfig, allocSimpleConfig', freeSimpleConfig,@@ -79,6 +90,7 @@ -- * Error HTTP2Error (..),+ StreamTerminated (..), ReasonPhrase, ErrorCode ( ErrorCode,@@ -104,3 +116,4 @@ import Network.HTTP2.Client.Run import Network.HTTP2.Frame import Network.HTTP2.H2 hiding (authority, scheme)+import Network.HTTP2.H2.OutBodyIface
Network/HTTP2/Client/Internal.hs view
@@ -1,6 +1,7 @@ module Network.HTTP2.Client.Internal ( Request (..), Response (..),+ Config (..), ClientConfig (..), Settings (..), Aux (..),
Network/HTTP2/Client/Run.hs view
@@ -7,7 +7,7 @@ import Control.Concurrent import Control.Concurrent.Async import Control.Concurrent.STM-import Control.Exception+import qualified Control.Exception as E import qualified Data.ByteString.UTF8 as UTF8 import Data.IORef import Data.IP (IPv6)@@ -22,6 +22,7 @@ import Imports import Network.HTTP2.Frame import Network.HTTP2.H2+import Network.HTTP2.H2.OutBodyIface -- | Client configuration data ClientConfig = ClientConfig@@ -88,15 +89,19 @@ False "Haskell!" -- 8 bytes }- clientCore ctx req processResponse = do+ clientCore ctx req processResponse = counted ctx $ do (strm, moutobj) <- makeStream ctx scheme authority req case moutobj of Nothing -> return () Just outobj -> sendRequest conf ctx strm outobj False rsp <- getResponse strm- x <- processResponse rsp- adjustRxWindow ctx strm- return x+ processResponse rsp `E.finally` doneWithStream ctx strm+ -- After a GOAWAY, the connection lasts as long as a request is in+ -- 'activeRequests' ('drained').+ counted ctx =+ E.bracket_+ (atomically $ modifyTVar' (activeRequests ctx) (+ 1))+ (atomically $ modifyTVar' (activeRequests ctx) (subtract 1)) runClient ctx = client (clientCore ctx) $ aux ctx -- | Launching a receiver and a sender.@@ -116,11 +121,31 @@ action $ ClientIO confMySockAddr confPeerSockAddr putR get putB create runH2 conf ctx runClient +-- | Called once 'processResponse' is done with a stream, however it ended.+--+-- A response it did not read to the end left its stream open. The server+-- went on sending the rest of the body, which nobody would read: it held+-- the stream's slot of the server's SETTINGS_MAX_CONCURRENT_STREAMS, and+-- the octets that came in were never given back to the connection window.+-- Enough such requests, or ones whose 'processResponse' threw, and new+-- requests waited for a slot, or the connection stalled. So such a stream+-- is reset (CANCEL), and what was left of its body is given back.+doneWithStream :: Context -> Stream -> IO ()+doneWithStream ctx strm = do+ cancelled <-+ closeIfReceiving ctx strm $ ResetByMe $ E.toException CancelledStream+ if cancelled+ then do+ enqueueControl (controlQ ctx) $+ CFrames Nothing [resetFrame Cancel $ streamNumber strm]+ giveBackUnread ctx strm+ else adjustRxWindow ctx strm+ getResponse :: Stream -> IO Response getResponse strm = do mRsp <- takeMVar $ streamInput strm case mRsp of- Left err -> throwIO err+ Left err -> E.throwIO err Right rsp -> return $ Response rsp setup :: ClientConfig -> Config -> IO Context@@ -140,7 +165,7 @@ runH2 :: Config -> Context -> IO a -> IO a runH2 conf ctx runClient = do- T.stopAfter mgr (try runAll >>= closureClient conf ctx) $ \res ->+ T.stopAfter mgr (E.try runAll >>= closureClient conf ctx) $ \res -> closeAllStreams (oddStreamTable ctx) (evenStreamTable ctx) res where mgr = threadManager ctx@@ -148,13 +173,55 @@ runSender = frameSender ctx conf runClientReceiver = do labelMe "H2 ClientReceiver"- er <- race runReceiver runClient- case er of- Right r -> return r- -- never reached because runReceiver throws an exception to exit.- Left () -> throwIO ConnectionIsClosed- runAll = snd <$> concurrently runSender runClientReceiver+ withAsync runReceiver $ \ar ->+ withAsync runClient $ \ac -> do+ er <- waitEither ar ac+ case er of+ Right r -> return r+ Left err -> do+ goingaway <- isJust <$> readTVarIO (peerGoAway ctx)+ case E.fromException err of+ -- The connection has run its course after the+ -- server's GOAWAY. The client function is let+ -- finish, rather than killed with it: every+ -- request it makes from here is refused, and+ -- one still waiting is failed now.+ Just ConnectionIsClosed+ | goingaway -> do+ closeAllStreams+ (oddStreamTable ctx)+ (evenStreamTable ctx)+ (Just err)+ wait ac+ _ -> E.throwIO err + -- When 'runClientReceiver' terminates, it is important we give the sender+ -- a chance to terminate cleanly also (it's possible the client terminated+ -- but there are still some messages in the queue to be sent).+ --+ -- If the client terminated successfully, we ignore any other errors in the+ -- sender (indeed, any exception here might simply be that the background+ -- threads were cancelled /because/ the client terminated).+ --+ -- If the sender terminates first, it failed, and no request can go out+ -- any more: the client is stopped and the sender's error reported, rather+ -- than the client left waiting on a connection nothing sends on.+ runAll =+ withAsync runSender $ \as ->+ withAsync runClientReceiver $ \ac -> do+ r <- waitEither as ac+ case r of+ Right x -> wait as >> return x+ Left e -> do+ -- The sender also finishes, normally, as soon as the+ -- receiver is done and the queues are empty, and may+ -- get there before the client side is seen to. Only+ -- with the receiver still running did it fail.+ done <- readTVarIO $ receiverDone ctx+ case done of+ Just _ -> wait ac+ Nothing -> E.throwIO e+ makeStream :: Context -> Scheme@@ -190,12 +257,13 @@ req' = req{outObjHeaders = hdr2} -- FLOW CONTROL: SETTINGS_MAX_CONCURRENT_STREAMS: send: respecting peer's limit (_sid, newstrm) <- openOddStreamWait ctx+ writeIORef (streamRequestMethod newstrm) $ Just method return (newstrm, Just req') sendRequest :: Config -> Context -> Stream -> OutObj -> Bool -> IO () sendRequest Config{..} ctx@Context{..} strm OutObj{..} io = do let sid = streamNumber strm- (mnext, mtbq) <- case outObjBody of+ (mnext, mtbq) <- (`E.onException` abandon sid) $ case outObjBody of OutBodyNone -> return (Nothing, Nothing) OutBodyFile (FileSpec path fileoff bytecount) -> do (pread, sentinel) <- confPositionReadMaker path@@ -216,11 +284,11 @@ let ot = OHeader outObjHeaders mnext outObjTrailers if io then do- let out = makeOutputIO ctx strm ot- pushOutput sid out+ let out = makeOutputIO ctx strm mtbq ot+ pushOutput sid out `E.onException` abandon sid else do (pop, out) <- makeOutput strm ot- pushOutput sid out+ pushOutput sid out `E.onException` abandon sid lc <- newLoopCheck strm mtbq T.forkManaged threadManager label $ syncWithSender' ctx pop lc where@@ -230,16 +298,32 @@ check (sidOK == sid) writeTVar outputQStreamID (sid + 2) enqueueOutputSTM outputQ out+ -- The request failed before it was queued -- the file of a+ -- 'requestFile' could not be opened, say, or the thread was killed while+ -- waiting for its turn. Its stream id was taken but nothing went out on+ -- it, and requests go out in stream id order: 'pushOutput' waits for+ -- 'outputQStreamID' to reach its own id. Left as it was, that turn never+ -- came, so every later request waited for ever, and the stream held its+ -- concurrency slot. So the stream is taken out of the table, and a+ -- thread passes its turn on once it arrives; the id goes unused, which a+ -- later, higher one closes implicitly (RFC 9113, section 5.1.1).+ abandon sid = do+ closed ctx strm Killed+ T.forkManaged threadManager ("H2 skipping stream " ++ show sid) $+ atomically $ do+ sidOK <- readTVar outputQStreamID+ check (sidOK == sid)+ writeTVar outputQStreamID (sid + 2) sendStreaming :: Context -> Stream -> (OutBodyIface -> IO ()) -> IO (TBQueue StreamingChunk)-sendStreaming Context{..} strm strmbdy = do+sendStreaming ctx@Context{..} strm strmbdy = do tbq <- newTBQueueIO 10 -- fixme: hard coding: 10 T.forkManagedUnmask threadManager label $ \unmask ->- withOutBodyIface tbq unmask strmbdy+ withOutBodyIface ctx strm tbq unmask strmbdy return tbq where label = "H2 request streaming sender for stream " ++ show (streamNumber strm)
Network/HTTP2/Frame/Decode.hs view
@@ -1,3 +1,4 @@+{-# LANGUAGE NamedFieldPuns #-} {-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE RecordWildCards #-} @@ -23,7 +24,7 @@ decodeContinuationFrame, ) where -import Control.Exception (Exception)+import qualified Control.Exception as E import Data.Array (Array, listArray, (!)) import qualified Data.ByteString as BS import Foreign.Ptr (Ptr, plusPtr)@@ -38,7 +39,7 @@ data FrameDecodeError = FrameDecodeError ErrorCode StreamId ShortByteString deriving (Eq, Show) -instance Exception FrameDecodeError+instance E.Exception FrameDecodeError ---------------------------------------------------------------- @@ -92,6 +93,10 @@ Left $ FrameDecodeError ProtocolError streamId "cannot used in non-zero stream" | otherwise = checkType typ where+ checkType FrameData+ | testPadded flags && payloadLength < 1 =+ Left $+ FrameDecodeError FrameSizeError streamId "insufficient payload for Pad Length" checkType FrameHeaders | testPadded flags && payloadLength < 1 = Left $@@ -142,6 +147,18 @@ ProtocolError streamId "push promise must be used with an odd stream identifier"+ | testPadded flags && payloadLength < 5 =+ Left $+ FrameDecodeError+ FrameSizeError+ streamId+ "insufficient payload for Pad Length and promised stream id"+ | not (testPadded flags) && payloadLength < 4 =+ Left $+ FrameDecodeError+ FrameSizeError+ streamId+ "insufficient payload for promised stream id" checkType FramePing | payloadLength /= 8 = Left $@@ -206,39 +223,52 @@ decodeFramePayload :: FrameType -> FramePayloadDecoder decodeFramePayload ftyp | ftyp > maxFrameType = checkFrameSize $ decodeUnknownFrame ftyp-decodeFramePayload ftyp = checkFrameSize decoder- where- decoder = payloadDecoders ! ftyp+decodeFramePayload ftyp = payloadDecoders ! ftyp -- each one checks its own size ---------------------------------------------------------------- -- | Frame payload decoder for DATA frame. decodeDataFrame :: FramePayloadDecoder-decodeDataFrame header _bs = decodeWithPadding header _bs DataFrame+decodeDataFrame = checkFrameSize $ \header bs ->+ decodeWithPadding header bs $ Right . DataFrame -- | Frame payload decoder for HEADERS frame. decodeHeadersFrame :: FramePayloadDecoder-decodeHeadersFrame header _bs = decodeWithPadding header _bs $ \bs' ->- if hasPriority- then- let (bs0, bs1) = BS.splitAt 5 bs'- p = priority bs0- in HeadersFrame (Just p) bs1- else HeadersFrame Nothing bs'- where- hasPriority = testPriority $ flags header+decodeHeadersFrame = checkFrameSize $ \header@FrameHeader{streamId} bs ->+ decodeWithPadding header bs $ \bs' ->+ if testPriority $ flags header+ then+ -- The header check knows the payload is long enough to hold+ -- the priority fields, but not that the padding leaves them+ -- there: Pad Length may cover the lot.+ if BS.length bs' < 5+ then+ Left $+ FrameDecodeError+ FrameSizeError+ streamId+ "no room for priority fields"+ else+ let (bs0, bs1) = BS.splitAt 5 bs'+ in Right $ HeadersFrame (Just (priority bs0)) bs1+ else Right $ HeadersFrame Nothing bs' -- | Frame payload decoder for PRIORITY frame. decodePriorityFrame :: FramePayloadDecoder-decodePriorityFrame _ bs = Right $ PriorityFrame $ priority bs+decodePriorityFrame = checkFrameSize $ requireBytes 5 $ \_ bs ->+ Right $ PriorityFrame $ priority bs -- | Frame payload decoder for RST_STREAM frame. decodeRSTStreamFrame :: FramePayloadDecoder-decodeRSTStreamFrame _ bs = Right $ RSTStreamFrame $ toErrorCode $ N.word32 bs+decodeRSTStreamFrame = checkFrameSize $ requireBytes 4 $ \_ bs ->+ Right $ RSTStreamFrame $ toErrorCode $ N.word32 bs -- | Frame payload decoder for SETTINGS frame. decodeSettingsFrame :: FramePayloadDecoder-decodeSettingsFrame FrameHeader{..} (PS fptr off _)+decodeSettingsFrame = checkFrameSize decodeSettingsFrame'++decodeSettingsFrame' :: FramePayloadDecoder+decodeSettingsFrame' FrameHeader{..} (PS fptr off _) | num > 10 = Left $ FrameDecodeError EnhanceYourCalm streamId "Settings is too large" | otherwise = Right $ SettingsFrame alist@@ -258,20 +288,31 @@ -- | Frame payload decoder for PUSH_PROMISE frame. decodePushPromiseFrame :: FramePayloadDecoder-decodePushPromiseFrame header _bs = decodeWithPadding header _bs $ \bs' ->- let (bs0, bs1) = BS.splitAt 4 bs'- sid = streamIdentifier (N.word32 bs0)- in PushPromiseFrame sid bs1+decodePushPromiseFrame = checkFrameSize $ \header@FrameHeader{streamId} bs ->+ decodeWithPadding header bs $ \bs' ->+ -- As in HEADERS: the padding may cover the promised stream id.+ if BS.length bs' < 4+ then+ Left $+ FrameDecodeError+ FrameSizeError+ streamId+ "no room for the promised stream id"+ else+ let (bs0, bs1) = BS.splitAt 4 bs'+ sid = streamIdentifier (N.word32 bs0)+ in Right $ PushPromiseFrame sid bs1 -- | Frame payload decoder for PING frame. decodePingFrame :: FramePayloadDecoder-decodePingFrame _ _bs = Right $ PingFrame bs- where- bs = BS.copy _bs+decodePingFrame = checkFrameSize $ \_ _bs -> Right $ PingFrame $ BS.copy _bs -- | Frame payload decoder for GOAWAY frame. decodeGoAwayFrame :: FramePayloadDecoder-decodeGoAwayFrame _ _bs = Right $ GoAwayFrame sid ecid bs2+decodeGoAwayFrame = checkFrameSize $ requireBytes 8 decodeGoAwayFrame'++decodeGoAwayFrame' :: FramePayloadDecoder+decodeGoAwayFrame' _ _bs = Right $ GoAwayFrame sid ecid bs2 where bs = BS.copy _bs (bs0, bs1') = BS.splitAt 4 bs@@ -281,7 +322,10 @@ -- | Frame payload decoder for WINDOW_UPDATE frame. decodeWindowUpdateFrame :: FramePayloadDecoder-decodeWindowUpdateFrame FrameHeader{..} bs+decodeWindowUpdateFrame = checkFrameSize $ requireBytes 4 decodeWindowUpdateFrame'++decodeWindowUpdateFrame' :: FramePayloadDecoder+decodeWindowUpdateFrame' FrameHeader{..} bs | wsi == 0 = Left $ FrameDecodeError ProtocolError streamId "window update must not be 0" | otherwise = Right $ WindowUpdateFrame wsi@@ -290,9 +334,7 @@ -- | Frame payload decoder for CONTINUATION frame. decodeContinuationFrame :: FramePayloadDecoder-decodeContinuationFrame _ _bs = Right $ ContinuationFrame bs- where- bs = BS.copy _bs+decodeContinuationFrame = checkFrameSize $ \_ _bs -> Right $ ContinuationFrame $ BS.copy _bs decodeUnknownFrame :: FrameType -> FramePayloadDecoder decodeUnknownFrame typ _ _bs = Right $ UnknownFrame typ bs@@ -307,6 +349,22 @@ Left $ FrameDecodeError FrameSizeError streamId "payload is too short" | otherwise = func header body +-- | Require the payload to actually hold the fixed fields about to be read+-- from it.+--+-- The reads below sit at fixed offsets and never consult the length of the+-- 'ByteString' they read from, so a payload shorter than the field runs off+-- the end of the buffer -- and an empty one is the shared empty+-- 'ByteString', whose pointer is null. 'checkFrameHeader' pins these lengths+-- down, but it is a separate function that a caller of the decoders is free+-- not to have used, and 'checkFrameSize' only compares the payload against+-- the length the frame header claims, which may itself be wrong.+requireBytes :: Int -> FramePayloadDecoder -> FramePayloadDecoder+requireBytes n func header@FrameHeader{streamId} body+ | BS.length body < n =+ Left $ FrameDecodeError FrameSizeError streamId "payload is too short"+ | otherwise = func header body+ -- | Helper function to pull off the padding if its there, and will -- eat up the trailing padding automatically. Calls the decoder func -- passed in with the length of the unpadded portion between the@@ -314,17 +372,23 @@ decodeWithPadding :: FrameHeader -> ByteString- -> (ByteString -> FramePayload)+ -> (ByteString -> Either FrameDecodeError FramePayload) -> Either FrameDecodeError FramePayload decodeWithPadding FrameHeader{..} bs body- | padded =- let (w8, rest) = fromMaybe (error "decodeWithPadding") $ BS.uncons bs'- padlen = intFromWord8 w8- bodylen = payloadLength - padlen - 1- in if bodylen < 0- then Left $ FrameDecodeError ProtocolError streamId "padding is not enough"- else Right . body $ BS.take bodylen rest- | otherwise = Right $ body bs'+ | padded = case BS.uncons bs' of+ -- The header checks rule this out for every frame type that can be+ -- padded, but the type does not, and the reply to a payload with no+ -- room for its Pad Length is an error, never a crash.+ Nothing ->+ Left $+ FrameDecodeError FrameSizeError streamId "insufficient payload for Pad Length"+ Just (w8, rest)+ | bodylen < 0 ->+ Left $ FrameDecodeError ProtocolError streamId "padding is not enough"+ | otherwise -> body $ BS.take bodylen rest+ where+ bodylen = payloadLength - intFromWord8 w8 - 1+ | otherwise = body bs' where bs' = BS.copy bs padded = testPadded flags
Network/HTTP2/H2/Config.hs view
@@ -33,6 +33,7 @@ confMySockAddr <- getSocketName s confPeerSockAddr <- getPeerName s let confReadNTimeout = False+ let confOnInformational = \_ _ -> return () return Config{..} -- | Deallocating the resource of the simple configuration.
Network/HTTP2/H2/Context.hs view
@@ -5,9 +5,9 @@ module Network.HTTP2.H2.Context where import Control.Concurrent.STM-import Control.Exception import qualified Control.Exception as E import Data.IORef+import qualified Data.IntMap.Strict as IntMap import Network.Control import Network.Socket (SockAddr) import qualified System.ThreadManager as T@@ -64,11 +64,7 @@ , peerSettings :: IORef Settings , oddStreamTable :: TVar OddStreamTable , evenStreamTable :: TVar EvenStreamTable- , continued :: IORef (Maybe StreamId)- -- ^ RFC 9113 says "Other frames (from any stream) MUST NOT- -- occur between the HEADERS frame and any CONTINUATION- -- frames that might follow". This field is used to implement- -- this requirement.+ , continued :: IORef (Maybe HeaderContinuation) , myStreamId :: TVar StreamId , peerStreamId :: IORef StreamId , peerLastStreamId :: IORef StreamId@@ -90,11 +86,38 @@ , mySockAddr :: SockAddr , peerSockAddr :: SockAddr , threadManager :: T.ThreadManager- , receiverDone :: TVar Bool+ , receiverDone :: TVar (Maybe E.SomeException) , workersDone :: STM Bool+ , informationalCallback :: StreamId -> TokenHeaderTable -> IO ()+ -- ^ Client only: called when a 1xx informational response (e.g. 103 Early+ -- Hints) is received, ahead of the final response. Copied from+ -- 'confOnInformational'; no-op by default.+ , peerGoAway :: TVar (Maybe StreamId)+ -- ^ The last stream identifier of the peer's GOAWAY(NO_ERROR), once+ -- one has come: see 'goingAway'.+ , activeRequests :: TVar Int+ -- ^ Client only: requests whose 'processResponse' has not returned.+ -- A response can be complete, its stream gone from the table, and+ -- its body still being read. } {- FOURMOLU_ENABLE -} +-- | Header/trailer continuation+--+-- RFC 9113 says "Other frames (from any stream) MUST NOT occur between the+-- HEADERS frame and any CONTINUATION frames that might follow". This is used to+-- implement this requirement.+--+-- It also accumulates the fragments of the block. These are connection-level+-- state: the block must be decoded even if its stream is reset before the+-- block is complete, since it may modify the dynamic table.+data HeaderContinuation = HeaderContinuation+ { hcStreamId :: StreamId+ , hcBlock :: PartialHeaderBlock+ , hcEndOfStream :: Bool+ -- ^ END_STREAM, from the HEADERS frame that started the block+ }+ ---------------------------------------------------------------- {- FOURMOLU_DISABLE -}@@ -138,8 +161,11 @@ let mySockAddr = confMySockAddr let peerSockAddr = confPeerSockAddr threadManager <- T.newThreadManager timmgr- receiverDone <- newTVarIO False+ receiverDone <- newTVarIO Nothing+ let informationalCallback = confOnInformational let workersDone = fromMaybe (T.isAllGone threadManager) mdone+ peerGoAway <- newTVarIO Nothing+ activeRequests <- newTVarIO 0 return Context{..} where role = case roleInfo of@@ -189,47 +215,115 @@ {-# INLINE setStreamState #-} setStreamState :: Context -> Stream -> StreamState -> IO ()-setStreamState _ Stream{streamState} newState = do- oldState <- readIORef streamState+setStreamState _ Stream{streamNumber, streamState} newState = atomically $ do+ oldState <- readTVar streamState+ informReplaced streamNumber oldState newState+ writeTVar streamState newState++-- | Replacing the open state of a stream as the receiver moves it on, from+-- headers to body.+--+-- The receiver reads a stream's state, works out the next one from the+-- frame, and writes it back -- in a transaction of its own. In between, the+-- sender may have half-closed the stream on our side ('halfClosedLocal',+-- which records it as @Open (Just cc) _@) or closed it. Writing the whole+-- state back undid that: the half-close was lost, the peer's END_STREAM then+-- took the stream to half-closed (remote) rather than closed, and it stayed+-- in the stream table, holding its concurrency slot for good. With both ends+-- streaming at once -- gRPC-style -- a client ran out of streams and a+-- server refused every new one.+--+-- So only the open state is replaced, keeping whatever the sender recorded+-- about our side, and a stream that is no longer open is left alone.+setOpenState :: Context -> Stream -> OpenState -> IO ()+setOpenState _ Stream{streamNumber, streamState} o = atomically $ do+ oldState <- readTVar streamState+ case oldState of+ Open hcl _ -> do+ let newState = Open hcl o+ informReplaced streamNumber oldState newState+ writeTVar streamState newState+ _otherwise -> return ()++-- | Inform consumers of any streams that we close+informReplaced :: StreamId -> StreamState -> StreamState -> STM ()+informReplaced streamNumber oldState newState = case (oldState, newState) of (Open _ (Body q _ _ _), Open _ (Body q' _ _ _)) | q == q' -> -- The stream stays open with the same body; nothing to do return ()+ (Open _ (Body q _ _ _), Closed cc) ->+ writeTQueue q $ Left $ E.toException $ closedCodeToError streamNumber cc (Open _ (Body q _ _ _), _) ->- -- The stream is either closed, or is open with a /new/ body- -- We need to close the old queue so that any reads from it won't block- atomically $ writeTQueue q $ Left $ toException ConnectionIsClosed+ -- The stream is opened with a /new/ body+ writeTQueue q $ Left $ E.toException ConnectionIsClosed _otherwise -> -- The stream wasn't open to start with; nothing to do return ()- writeIORef streamState newState +-- | Opening an idle stream.+--+-- Only an idle one: the receiver checks that the stream is idle and then+-- opens it, and a client's request stream stays idle while the request is+-- being sent -- so by the time it opens the stream for the response's+-- HEADERS, the sender may already have half-closed it ('halfClosedLocal'+-- turns an idle stream into @Open (Just cc) JustOpened@). Opening it over+-- that lost the half-close, with the same result as described at+-- 'setOpenState'. opened :: Context -> Stream -> IO ()-opened ctx strm = setStreamState ctx strm (Open Nothing JustOpened)+opened _ Stream{streamState} = atomically $ modifyTVar' streamState open+ where+ open Idle = Open Nothing JustOpened+ open st = st halfClosedRemote :: Context -> Stream -> IO () halfClosedRemote ctx stream@Stream{streamState} = do- closingCode <- atomicModifyIORef streamState closeHalf+ closingCode <- atomically $ stateTVar streamState closeHalf traverse_ (closed ctx stream) closingCode where- closeHalf :: StreamState -> (StreamState, Maybe ClosedCode)- closeHalf x@(Closed _) = (x, Nothing)- closeHalf (Open (Just cc) _) = (Closed cc, Just cc)- closeHalf _ = (HalfClosedRemote, Nothing)+ closeHalf :: StreamState -> (Maybe ClosedCode, StreamState)+ closeHalf x@(Closed _) = (Nothing, x)+ closeHalf (Open (Just cc) _) = (Just cc, Closed cc)+ closeHalf _ = (Nothing, HalfClosedRemote) halfClosedLocal :: Context -> Stream -> ClosedCode -> IO () halfClosedLocal ctx stream@Stream{streamState} cc = do- shouldFinalize <- atomicModifyIORef streamState closeHalf+ shouldFinalize <- atomically $ stateTVar streamState closeHalf when shouldFinalize $ closed ctx stream cc where- closeHalf :: StreamState -> (StreamState, Bool)- closeHalf x@(Closed _) = (x, False)- closeHalf HalfClosedRemote = (Closed cc, True)- closeHalf (Open Nothing o) = (Open (Just cc) o, False)- closeHalf _ = (Open (Just cc) JustOpened, False)+ closeHalf :: StreamState -> (Bool, StreamState)+ closeHalf x@(Closed _) = (False, x)+ closeHalf HalfClosedRemote = (True, Closed cc)+ closeHalf (Open Nothing o) = (False, Open (Just cc) o)+ closeHalf _ = (False, Open (Just cc) JustOpened) +-- | Closing a stream whose response is still coming in, and saying+-- whether it was.+--+-- Decided and done in one transaction, so that a stream the peer finishes+-- meanwhile is not taken for one still open. A response whose END_STREAM+-- the receiver has already queued, and not yet recorded in the state, can+-- still be taken for one coming in; the reset that follows is one RFC 9113+-- section 5.1 has the peer ignore ("for a short period after a DATA or+-- HEADERS frame containing an END_STREAM flag is sent").+closeIfReceiving :: Context -> Stream -> ClosedCode -> IO Bool+closeIfReceiving ctx strm@Stream{streamNumber, streamState} cc = do+ receiving <- atomically $ do+ st <- readTVar streamState+ case st of+ -- END_STREAM came with the headers.+ Open _ (NoBody _) -> return False+ Open{} -> do+ informReplaced streamNumber st (Closed cc)+ writeTVar streamState (Closed cc)+ return True+ _otherwise -> return False+ -- Out of the stream table, giving its concurrency slot back.+ when receiving $ closed ctx strm cc+ return receiving+ closed :: Context -> Stream -> ClosedCode -> IO () closed ctx@Context{oddStreamTable, evenStreamTable} strm@Stream{streamNumber} cc = do if isServerInitiated streamNumber@@ -237,19 +331,20 @@ else deleteOdd oddStreamTable streamNumber err setStreamState ctx strm (Closed cc) -- anyway where- err :: SomeException- err = toException (closedCodeToError streamNumber cc)+ err :: E.SomeException+ err = E.toException (closedCodeToError streamNumber cc) ---------------------------------------------------------------- -- From peer -- Server+--+-- Note that this does not apply SETTINGS_MAX_CONCURRENT_STREAMS. A stream+-- over the limit still has to be admitted this far, because its field block+-- has to be decoded before it can be refused; 'checkOddConcurrency' does the+-- refusing once that has happened. openOddStreamCheck :: Context -> StreamId -> FrameType -> IO Stream openOddStreamCheck ctx@Context{oddStreamTable, peerSettings, mySettings} sid ftyp = do- -- My SETTINGS_MAX_CONCURRENT_STREAMS- when (ftyp == FrameHeaders) $ do- conc <- getOddConcurrency oddStreamTable- checkMyConcurrency sid mySettings (conc + 1) txws <- initialWindowSize <$> readIORef peerSettings let rxws = initialWindowSize mySettings newstrm <- newOddStream sid txws rxws@@ -268,6 +363,25 @@ newstrm <- newEvenStream sid txws rxws insertEvenCache evenStreamTable method path newstrm +-- | Refuse a peer-initiated stream that puts us over the limit we advertised+-- in SETTINGS_MAX_CONCURRENT_STREAMS.+--+-- Checked once the stream's field block has been decoded, rather than when+-- its HEADERS frame arrived. A block has to be decoded whatever becomes of+-- its stream -- RFC 9113 section 10.5.1, "The field block MUST be processed+-- to ensure a consistent connection state" -- and refusing at arrival meant+-- throwing before the frame's payload had even been read, which left nothing+-- to do but drop the connection. From here the throw lands inside the+-- receiver's per-frame reset handler, so the answer is+-- RST_STREAM(REFUSED_STREAM) and the connection carries on, which is what+-- section 5.1.2 asks for and what section 8.7 lets the peer retry against.+--+-- The stream is in the table by the time we get here, so it counts itself.+checkOddConcurrency :: Context -> StreamId -> IO ()+checkOddConcurrency Context{oddStreamTable, mySettings} sid = do+ conc <- getOddConcurrency oddStreamTable+ checkMyConcurrency sid mySettings conc+ checkMyConcurrency :: StreamId -> Settings -> Int -> IO () checkMyConcurrency sid settings conc = do@@ -290,13 +404,16 @@ let rxws = initialWindowSize mySettings case mMaxConc of Nothing -> do- sid <- atomically $ getMyNewStreamId ctx+ sid <- atomically $ do+ refuseIfGoingAway ctx+ getMyNewStreamId ctx txws <- initialWindowSize <$> readIORef peerSettings newstrm <- newOddStream sid txws rxws insertOdd oddStreamTable sid newstrm return (sid, newstrm) Just maxConc -> do sid <- atomically $ do+ refuseIfGoingAway ctx waitIncOdd oddStreamTable maxConc getMyNewStreamId ctx txws <- initialWindowSize <$> readIORef peerSettings@@ -305,8 +422,21 @@ return (sid, newstrm) -- Server-openEvenStreamWait :: Context -> IO (StreamId, Stream)-openEvenStreamWait ctx@Context{..} = do++-- | Opening a stream for a push, if the peer's+-- SETTINGS_MAX_CONCURRENT_STREAMS leaves room for one.+--+-- Not waiting for room: the response the push belongs to waits for its+-- PUSH_PROMISE to go out, and with no room -- a peer can announce 0 to+-- refuse pushes (RFC 9113, section 8.4) -- it waited for ever. A push is+-- only ever an offer, so one there is no room for is not made.+openEvenStreamTry :: Context -> IO (Maybe (StreamId, Stream))+openEvenStreamTry ctx@Context{..} = do+ goingaway <- isJust <$> readTVarIO peerGoAway+ if goingaway then return Nothing else openEvenStreamTry' ctx++openEvenStreamTry' :: Context -> IO (Maybe (StreamId, Stream))+openEvenStreamTry' ctx@Context{..} = do -- Peer SETTINGS_MAX_CONCURRENT_STREAMS mMaxConc <- maxConcurrentStreams <$> readIORef peerSettings let rxws = initialWindowSize mySettings@@ -316,12 +446,74 @@ txws <- initialWindowSize <$> readIORef peerSettings newstrm <- newEvenStream sid txws rxws insertEven evenStreamTable sid newstrm- return (sid, newstrm)+ return $ Just (sid, newstrm) Just maxConc -> do- sid <- atomically $ do- waitIncEven evenStreamTable maxConc- getMyNewStreamId ctx- txws <- initialWindowSize <$> readIORef peerSettings- newstrm <- newEvenStream sid txws rxws- insertEven' evenStreamTable sid newstrm- return (sid, newstrm)+ msid <- atomically $ do+ let open = do+ waitIncEven evenStreamTable maxConc+ Just <$> getMyNewStreamId ctx+ open `orElse` return Nothing+ forM msid $ \sid -> do+ txws <- initialWindowSize <$> readIORef peerSettings+ newstrm <- newEvenStream sid txws rxws+ insertEven' evenStreamTable sid newstrm+ return (sid, newstrm)++----------------------------------------------------------------+-- GOAWAY from the peer++-- | No new stream once the peer has sent GOAWAY: "Receivers of a GOAWAY+-- frame MUST NOT open additional streams on the connection" (RFC 9113,+-- section 6.8). A request that has not got a stream yet is answered as+-- one on a stream above the last stream identifier would be.+refuseIfGoingAway :: Context -> STM ()+refuseIfGoingAway Context{peerGoAway} = do+ goingaway <- isJust <$> readTVar peerGoAway+ when goingaway $ throwSTM ConnectionIsClosed++-- | Taking the peer's GOAWAY(NO_ERROR) on board.+--+-- The peer has said it will go no further than the last stream identifier,+-- not that it is going at once. Streams of ours above it will not be+-- processed, and are closed with 'ConnectionIsClosed', which a client may+-- retry elsewhere; no stream is opened from now on. The others go on+-- until they are done, and then so is the connection: 'drained' tells+-- when. It used to be closed right away, cutting off every stream in+-- flight, even those the peer had promised to finish.+--+-- A later GOAWAY can lower the last stream identifier, not raise it.+-- Answers whether this is the first.+goingAway :: Context -> StreamId -> IO Bool+goingAway ctx@Context{peerGoAway, oddStreamTable, evenStreamTable} lastSid = do+ first <- atomically $ do+ old <- readTVar peerGoAway+ writeTVar peerGoAway $ Just $ maybe lastSid (min lastSid) old+ return $ isNothing old+ mine <-+ if isClient ctx+ then getOddStreams oddStreamTable+ else getEvenStreams evenStreamTable+ let (_, above) = IntMap.split lastSid mine+ forM_ above $ \strm -> closed ctx strm Finished+ return first++-- | Waiting until nothing is left in flight after the peer's GOAWAY:+-- neither a stream of the peer's, nor one of ours the peer is to process,+-- nor -- on a client -- a response still being read. Streams of ours above+-- the last stream identifier are not waited for; those that got their+-- identifier too late to be closed by 'goingAway' are closed with the+-- connection.+drained :: Context -> STM ()+drained ctx@Context{peerGoAway, oddStreamTable, evenStreamTable, activeRequests} = do+ mlast <- readTVar peerGoAway+ case mlast of+ Nothing -> retry+ Just lastSid -> do+ odds <- oddTable <$> readTVar oddStreamTable+ evens <- evenTable <$> readTVar evenStreamTable+ active <- readTVar activeRequests+ let (mine, theirs)+ | isClient ctx = (odds, evens)+ | otherwise = (evens, odds)+ noneOfMine = maybe True ((> lastSid) . fst) $ IntMap.lookupMin mine+ check $ IntMap.null theirs && noneOfMine && active == 0
Network/HTTP2/H2/HPACK.hs view
@@ -3,14 +3,20 @@ module Network.HTTP2.H2.HPACK ( hpackEncodeHeader,- hpackEncodeHeaderLoop,+ hpackEncodeHeaderRest, hpackDecodeHeader, hpackDecodeTrailer,+ hpackDiscardHeader, just, fixHeaders, ) where import qualified Control.Exception as E+import qualified Data.ByteString as BS+import Data.ByteString.Internal (create)+import qualified Data.ByteString.Lazy as BS.Lazy+import Foreign.Marshal.Alloc (free, mallocBytes)+import Foreign.Marshal.Utils (copyBytes) import Network.ByteOrder import Network.HTTP.Semantics import Network.HTTP.Types@@ -68,25 +74,89 @@ hpackEncodeHeaderLoop Context{..} buf siz hs = encodeTokenHeader buf siz strategy False encodeDynamicTable hs +-- | Encode the rest of a header block whose start 'hpackEncodeHeader' wrote+--+-- For a block that did not fit where it was being written. Grows the buffer+-- as needed: a header that does not fit is retried with a larger buffer (the+-- encoder does not modify the dynamic table for a header it could not write).+hpackEncodeHeaderRest+ :: Context+ -> BufferSize+ -- ^ Initial buffer size+ -> TokenHeaderList+ -> IO BS.Lazy.ByteString+hpackEncodeHeaderRest ctx = go []+ where+ go acc _ [] = return $ BS.Lazy.fromChunks (reverse acc)+ go acc siz ths = do+ (chunk, ths') <- E.bracket (mallocBytes siz) free $ \buf -> do+ (ths', len) <- hpackEncodeHeaderLoop ctx buf siz ths+ chunk <- create len $ \p -> copyBytes p buf len+ return (chunk, ths')+ if BS.null chunk+ then go acc (siz * 2) ths -- no progress: grow+ else go (chunk : acc) siz ths'+ ---------------------------------------------------------------- hpackDecodeHeader :: HeaderBlockFragment -> StreamId -> Context -> IO TokenHeaderTable hpackDecodeHeader hdrblk sid ctx = do- tbl@(_, vt) <- hpackDecodeTrailer hdrblk sid ctx+ tbl@(_, vt) <- hpackDecode "illegal header" hdrblk sid ctx if isClient ctx || checkRequestHeader vt then return tbl else E.throwIO $ StreamErrorIsSent ProtocolError sid "illegal header" hpackDecodeTrailer :: HeaderBlockFragment -> StreamId -> Context -> IO TokenHeaderTable-hpackDecodeTrailer hdrblk sid Context{..} = decodeTokenHeader decodeDynamicTable hdrblk `E.catch` handl+hpackDecodeTrailer = hpackDecode "illegal trailer"++-- | Decode a field block for a stream we no longer have, and discard the result+--+-- The block must still be decoded: it may modify the dynamic table.+hpackDiscardHeader :: HeaderBlockFragment -> StreamId -> Context -> IO ()+hpackDiscardHeader hdrblk sid ctx =+ void (hpackDecode "illegal header" hdrblk sid ctx) `E.catch` ignore where+ -- Decoded to the end, so there is nothing to refuse: the message goes+ -- nowhere anyway.+ ignore StreamErrorIsSent{} = return ()+ ignore e = E.throwIO e++-- | Decode a field block, reporting a block we could not get through as a+-- connection error, and a malformed message in a block decoded to the end+-- as a stream error.+--+-- The first argument says which kind of block it was, since the peer reads+-- this in the GOAWAY and "illegal trailer" about a request's headers is a+-- confusing thing to be told.+hpackDecode+ :: ReasonPhrase+ -> HeaderBlockFragment+ -> StreamId+ -> Context+ -> IO TokenHeaderTable+hpackDecode illegal hdrblk sid Context{..} =+ decodeTokenHeader decodeDynamicTable hdrblk `E.catch` handl+ where+ -- A malformed message: 'decodeTokenHeader' says so only once it has+ -- decoded the whole block, so the dynamic table is up to date and the+ -- connection can go on. A stream error, as RFC 9113 section 8.1.1 has+ -- it. This used to be a connection error, as the block was abandoned+ -- at the malformed field; one request with an upper-case field name+ -- closed the connection, every other stream on it with it. handl IllegalHeaderName =- E.throwIO $ StreamErrorIsSent ProtocolError sid "illegal trailer"+ E.throwIO $ StreamErrorIsSent ProtocolError sid illegal+ handl TooLargeHeader =+ E.throwIO $ StreamErrorIsSent ProtocolError sid "too many fields"+ -- A block we could not get through: our dynamic table now holds the+ -- entries decoded before the throw and nothing after them -- no longer+ -- what the peer's encoder believes we have. Section 10.5.1: "The field+ -- block MUST be processed to ensure a consistent connection state,+ -- unless the connection is closed." We did not, so it must be. handl e = do let msg = fromString $ show e- E.throwIO $ StreamErrorIsSent CompressionError sid msg+ E.throwIO $ ConnectionErrorIsSent CompressionError sid msg {-# INLINE checkRequestHeader #-} checkRequestHeader :: ValueTable -> Bool
+ Network/HTTP2/H2/OutBodyIface.hs view
@@ -0,0 +1,130 @@+{-# LANGUAGE DeriveAnyClass #-}+{-# LANGUAGE DerivingStrategies #-}+{-# LANGUAGE NamedFieldPuns #-}+{-# LANGUAGE RankNTypes #-}++module Network.HTTP2.H2.OutBodyIface (+ StreamTerminated (..),+ withOutBodyIface,+) where++import Control.Concurrent.STM+import Control.Exception+import Network.HTTP.Semantics+import Network.HTTP.Semantics.IO+import Network.HTTP2.H2.Context+import Network.HTTP2.H2.Sync+import Network.HTTP2.H2.Types++----------------------------------------------------------------++data StreamTerminated+ = StreamPushedFinal+ | StreamCancelled+ | StreamOutOfScope+ | StreamRemoteReset ClosedCode+ deriving (Show)+ deriving anyclass (Exception)++----------------------------------------------------------------++withOutBodyIface+ :: Context+ -> Stream+ -> TBQueue StreamingChunk+ -> (forall a. IO a -> IO a)+ -> (OutBodyIface -> IO r)+ -> IO r+withOutBodyIface ctx@Context{outputQ} strm tbq unmask k = do+ terminated <- newTVarIO Nothing+ let checkNotTerminated :: STM ()+ checkNotTerminated = do+ mTerminated <- readTVar terminated+ maybe (return ()) throwSTM mTerminated++ -- Check if the peer is still listening for messages+ --+ -- It is important to call 'checkNotClosed' prior to enqueuing stream+ -- chunks to ensure that 'writeTBQueue' will not block indefinitely+ -- (because nothing is consuming elements from the queue anymore).+ --+ -- Assumes 'checkNotTerminated'.+ checkNotClosed :: STM ()+ checkNotClosed = do+ mClosed <- getIsClosed+ case mClosed of+ Just code ->+ -- When the stream is closed, but /we/ did not close it (or+ -- 'checkNotTerminated' would have thrown an exception), it+ -- must mean that our peer send us a RST_STREAM, indicating+ -- that they do not want to receive any further messages.+ throwSTM $ StreamRemoteReset code+ _otherwise ->+ return ()++ getIsClosed :: STM (Maybe ClosedCode)+ getIsClosed = do+ st <- readTVar (streamState strm)+ case st of+ Closed code -> return $ Just code+ _otherwise -> return Nothing++ cancelAfterFinish :: Maybe SomeException -> STM ()+ cancelAfterFinish mErr =+ writeTQueue outputQ $ makeOutputIO ctx strm Nothing (OReset mErr)++ iface :: OutBodyIface+ iface =+ OutBodyIface+ { outBodyUnmask = unmask+ , outBodyPush = \b -> atomically $ do+ checkNotTerminated+ checkNotClosed+ writeTBQueue tbq $ StreamingBuilder b NotEndOfStream+ , outBodyPushFinal = \b -> atomically $ do+ checkNotTerminated+ checkNotClosed+ writeTVar terminated (Just StreamPushedFinal)+ writeTBQueue tbq $ StreamingBuilder b (EndOfStream Nothing)+ writeTBQueue tbq $ StreamingFinished Nothing+ , outBodyFlush = atomically $ do+ checkNotTerminated+ checkNotClosed+ writeTBQueue tbq StreamingFlush+ , outBodyCancel = \mErr -> atomically $ do+ mTerminated <- readTVar terminated+ mClosed <- getIsClosed+ case (mClosed, mTerminated) of+ (Nothing, Nothing) -> do+ writeTVar terminated (Just StreamCancelled)+ writeTBQueue tbq $ StreamingCancelled mErr+ (Nothing, Just StreamCancelled) ->+ -- Already cancelled+ return ()+ (Nothing, Just _) -> do+ -- We finished streaming (that is, sending messages to the peer),+ -- but we must still be able to cancel the stream entirely+ -- (that is, tell the peer that we no longer want to /receive/ messages: RST_STREAM)+ writeTVar terminated (Just StreamCancelled)+ cancelAfterFinish mErr+ (Just _code, _) ->+ -- Peer already closed+ return ()+ }++ finished :: IO ()+ finished = atomically $ do+ mTerminated <- readTVar terminated+ mClosed <- getIsClosed+ case (mClosed, mTerminated) of+ (Nothing, Nothing) -> do+ writeTVar terminated (Just StreamOutOfScope)+ writeTBQueue tbq $ StreamingFinished Nothing+ (Nothing, Just _) ->+ -- We already terminated+ return ()+ (Just _code, _) ->+ -- Peer already closed+ return ()++ k iface `finally` finished
Network/HTTP2/H2/Receiver.hs view
@@ -19,8 +19,10 @@ import qualified Data.ByteString.Short as Short import qualified Data.ByteString.UTF8 as UTF8 import Data.IORef+import Data.Void import Network.Control import Network.HTTP.Semantics+import qualified System.IO.Error as E import qualified System.ThreadManager as T import Imports hiding (delete, insert)@@ -45,14 +47,20 @@ ---------------------------------------------------------------- -frameReceiver :: Context -> Config -> IO ()-frameReceiver ctx conf@Config{..} =- (switch `E.catch` handler)- `E.finally` atomically- (writeTVar (receiverDone ctx) True)+frameReceiver :: Context -> Config -> IO E.SomeException+frameReceiver ctx@Context{receiverDone} conf@Config{..} =+ E.mask $ \unmask -> do+ -- This catches an asynchronous exception.+ -- It is re-thrown by "runH2"+ mErr <- E.try $ unmask switch+ case mErr of+ Left err -> do+ atomically $ writeTVar receiverDone $ Just err+ return err+ Right x -> do+ absurd x -- We only terminate due to exceptions where- handler ConnectionIsClosed = return ()- handler e = E.throwIO e+ switch :: IO Void switch = do labelMe "H2 receiver" tid <- myThreadId@@ -60,13 +68,16 @@ then loop1 else- void $- T.withHandle (threadManager ctx) (E.throwTo tid ConnectionIsTimeout) loop2+ T.withHandle (threadManager ctx) (E.throwTo tid ConnectionIsTimeout) loop2++ loop1 :: IO Void loop1 = do hd <- confReadN frameHeaderLength -- throwing an exception on timeout when (BS.null hd) $ E.throwIO ConnectionIsClosed processFrame ctx conf $ decodeFrameHeader hd loop1++ loop2 :: T.Handle -> IO Void loop2 th = do -- If 'confReadN' is timeouted, 'ConnectionIsTimeout' is thrown -- to destroy the thread trees.@@ -89,13 +100,13 @@ | isServer ctx = E.throwIO $ ConnectionErrorIsSent ProtocolError streamId "push promise is not allowed"-processFrame Context{..} Config{..} (ftyp, FrameHeader{payloadLength, streamId})+processFrame Context{..} conf (ftyp, FrameHeader{payloadLength, streamId}) | ftyp > maxFrameType = do mx <- readIORef continued case mx of Nothing -> do -- ignoring unknown frame- void $ confReadN payloadLength+ void $ readPayload conf payloadLength Just _ -> E.throwIO $ ConnectionErrorIsSent ProtocolError streamId "unknown frame" processFrame ctx@Context{..} conf typhdr@(ftyp, header) = do -- My SETTINGS_MAX_FRAME_SIZE@@ -115,51 +126,182 @@ ---------------------------------------------------------------- +-- | Read a frame payload in full.+--+-- 'confReadN' answers with an empty string at end of input, so a payload that+-- comes back short means the peer hung up in the middle of the frame. Saying+-- so here keeps every decoder below from being handed fewer bytes than the+-- frame header promised it.+readPayload :: Config -> Int -> IO ByteString+readPayload Config{..} len = do+ bs <- confReadN len+ when (BS.length bs /= len) $ E.throwIO ConnectionIsClosed+ return bs+ controlOrStream :: Context -> Config -> FrameType -> FrameHeader -> IO ()-controlOrStream ctx@Context{..} Config{..} ftyp header@FrameHeader{streamId, payloadLength}+controlOrStream ctx@Context{..} conf ftyp header@FrameHeader{flags, streamId, payloadLength} | isControl streamId = do- bs <- confReadN payloadLength+ bs <- readPayload conf payloadLength control ftyp header bs ctx | ftyp == FramePushPromise = do- bs <- confReadN payloadLength- push header bs ctx+ bs <- readPayload conf payloadLength+ -- A promised stream can be refused over concurrency too, and by the+ -- time 'push' gets that far it has decoded the field block, so the+ -- same reasoning as 'resettable' applies: reset the promised stream+ -- and read on. There is no 'Stream' to close -- it was refused+ -- before one was made -- so this resets by identifier alone.+ push header bs ctx `E.catch` resetPromised | otherwise = do- checkContinued+ mcont <- checkContinued mstrm <- getStream ctx ftyp streamId- bs <- confReadN payloadLength- case mstrm of- Just strm -> do- state0 <- readStreamState strm- state <- stream ftyp header bs ctx state0 strm- resetContinued- set <- processState state ctx strm streamId- when set setContinued- Nothing- | ftyp == FramePriority -> do- -- for h2spec only- PriorityFrame newpri <- guardIt $ decodePriorityFrame header bs- checkPriority newpri streamId- | otherwise -> return ()+ bs <- readPayload conf payloadLength+ case mcont of+ Just hc -> continuation hc bs mstrm+ Nothing ->+ case mstrm of+ Just strm -> resettable strm $ do+ state0 <- readStreamState strm+ state <- stream ftyp header bs ctx state0 strm+ processState state ctx strm streamId+ Nothing+ | ftyp == FramePriority -> do+ -- for h2spec only+ PriorityFrame newpri <- guardIt $ decodePriorityFrame header bs+ checkPriority newpri streamId+ | ftyp == FrameData ->+ -- Dropped, but still paid for.+ informIgnoredData ctx streamId payloadLength+ | ftyp == FrameHeaders -> do+ HeadersFrame _ frag <- guardIt $ decodeHeadersFrame header bs+ if testEndHeader flags+ then hpackDiscardHeader frag streamId ctx+ else startHeaderBlock ctx streamId (testEndStream flags) frag+ | otherwise -> return () where- setContinued = writeIORef continued $ Just streamId resetContinued = writeIORef continued Nothing+ resetPromised (StreamErrorIsSent err sid _msg) =+ enqueueControl controlQ $ CFrames Nothing [resetFrame err sid]+ resetPromised e = E.throwIO e+ -- Answer a stream error by resetting that stream and reading on, which is+ -- what RFC 9113 section 5.4.2 asks for: "an error related to a specific+ -- stream that does not affect processing of other streams".+ --+ -- Safe only here, after the payload has been read and any field block in+ -- it decoded, so that the connection sits at a frame boundary and the+ -- HPACK tables still agree with the peer's. Where neither holds -- a+ -- field block abandoned part-way, a stream refused before its payload was+ -- read -- the error is raised as a connection error where it is detected,+ -- and travels straight past this handler.+ resettable strm act = act `E.catch` reset+ where+ reset e@(StreamErrorIsSent err sid _msg) = do+ -- MadeYouReset: CVE-2025-8671. A reset we send because of+ -- what the peer sent frees the stream's concurrency slot just as+ -- one the peer sends does, while a handler already launched for+ -- the stream runs on. Counted with the peer's own resets, or a+ -- peer that never sends RST_STREAM -- a PRIORITY on the stream+ -- depending on itself is enough -- could keep any number of+ -- handlers running past SETTINGS_MAX_CONCURRENT_STREAMS.+ --+ -- REFUSED_STREAM is left out: it launches nothing, and it is the+ -- answer section 8.7 means a peer to be able to retry.+ when (err /= RefusedStream) $ do+ rate <- getRate rstRate+ when (rate > rstRateLimit mySettings) $+ E.throwIO $+ ConnectionErrorIsSent EnhanceYourCalm sid "too many stream errors"+ resetContinued+ -- 'closed' hands the exception to whoever is reading the stream+ -- and takes it out of the stream table.+ closed ctx strm $ ResetByMe $ E.toException e+ enqueueControl controlQ $ CFrames Nothing [resetFrame err sid]+ reset e = E.throwIO e++ checkContinued :: IO (Maybe HeaderContinuation) checkContinued = do mx <- readIORef continued case mx of Nothing -> return ()- Just sid- | sid == streamId && ftyp == FrameContinuation -> return ()+ Just hc+ | hcStreamId hc == streamId && ftyp == FrameContinuation -> return () | otherwise -> E.throwIO $ ConnectionErrorIsSent ProtocolError streamId "continuation frame must follow"+ return mx + continuation+ :: HeaderContinuation -> HeaderBlockFragment -> Maybe Stream -> IO ()+ continuation hc frag mstrm+ | frag == "" && not (testEndHeader flags) = do+ -- Empty Frame Flooding - CVE-2019-9518+ rate <- getRate emptyFrameRate+ when (rate > emptyFrameRateLimit mySettings) $+ E.throwIO $+ ConnectionErrorIsSent EnhanceYourCalm streamId "too many empty continuation"+ | otherwise = do+ phb' <- addFragment streamId frag (hcBlock hc)+ if testEndHeader flags+ then completeBlock hc (completeHeaderBlock phb') mstrm+ else writeIORef continued $ Just hc{hcBlock = phb'}++ completeBlock+ :: HeaderContinuation -> HeaderBlockFragment -> Maybe Stream -> IO ()+ completeBlock hc blk mstrm = do+ resetContinued+ case mstrm of+ Just strm -> resettable strm $ do+ state0 <- readStreamState strm+ case state0 of+ Open hcl JustOpened -> do+ tbl <- hpackDecodeHeader blk streamId ctx+ state <- onResponseHeaders ctx streamId hcl (hcEndOfStream hc) tbl+ processState state ctx strm streamId+ Open _ (Body q _ _ tlr) -> do+ state <- onTrailers ctx streamId blk q tlr+ processState state ctx strm streamId+ _otherwise ->+ -- The block began on a stream that was open for+ -- headers or trailers, and the sender has since+ -- closed it -- a reset crossing the block. There is+ -- nothing to deliver, but the block still has to go+ -- through the decoder, since it may have changed the+ -- dynamic table. Handing it to 'stream' as a+ -- CONTINUATION would get it refused as one that+ -- cannot come here, closing the connection.+ hpackDiscardHeader blk streamId ctx+ Nothing ->+ hpackDiscardHeader blk streamId ctx++-- | Is this a response that is defined to have no content?+--+-- RFC 9113, section 8.1.1: "A response that is defined to have no content,+-- as described in Section 6.4.1 of [HTTP], can have a non-zero+-- content-length header field, even though no content is included in DATA+-- frames." Those are the responses to HEAD, 204 and 304, and 2xx to+-- CONNECT; the content-length of one to HEAD, in particular, is that of the+-- content a GET would have had. Checking it against the content that+-- arrived made every such response to HEAD a stream error.+hasNoContent :: Context -> Stream -> ValueTable -> IO Bool+hasNoContent ctx Stream{streamRequestMethod} vt+ | isServer ctx = return False+ | otherwise = do+ mmethod <- readIORef streamRequestMethod+ let status = getFieldValue tokenStatus vt+ return $+ mmethod == Just "HEAD"+ || status `elem` [Just "204", Just "304"]+ || (mmethod == Just "CONNECT" && maybe False ("2" `BS.isPrefixOf`) status)+ ---------------------------------------------------------------- -processState :: StreamState -> Context -> Stream -> StreamId -> IO Bool+processState :: StreamState -> Context -> Stream -> StreamId -> IO () -- Transition (process1)-processState (Open _ (NoBody tbl@(_, reqvt))) ctx@Context{..} strm@Stream{streamInput} streamId = do+processState (Open _ (NoBody tbl@(_, reqvt))) ctx strm@Stream{streamInput} streamId = do+ -- My SETTINGS_MAX_CONCURRENT_STREAMS+ when (isServer ctx) $ checkOddConcurrency ctx streamId+ noContent <- hasNoContent ctx strm reqvt let mcl = fst <$> (getFieldValue tokenContentLength reqvt >>= C8.readInt)- when (just mcl (/= (0 :: Int))) $+ when (not noContent && just mcl (/= (0 :: Int))) $ E.throwIO $ StreamErrorIsSent ProtocolError@@ -169,52 +311,64 @@ let inpObj = InpObj tbl (Just 0) (return (mempty, True)) tlr if isServer ctx then do- let ServerInfo{..} = toServerInfo roleInfo- launch ctx strm inpObj+ launchHandler ctx strm inpObj else putMVar streamInput $ Right inpObj halfClosedRemote ctx strm- return False -- Transition (process2)-processState (Open hcl (HasBody tbl@(_, reqvt))) ctx@Context{..} strm@Stream{streamInput, streamRxQ} _streamId = do- let mcl = fst <$> (getFieldValue tokenContentLength reqvt >>= C8.readInt)+processState (Open _ (HasBody tbl@(_, reqvt))) ctx strm@Stream{streamInput, streamRxQ} _streamId = do+ -- My SETTINGS_MAX_CONCURRENT_STREAMS+ when (isServer ctx) $ checkOddConcurrency ctx _streamId+ noContent <- hasNoContent ctx strm reqvt+ let mcl+ -- Its content-length describes content it does not have.+ | noContent = Just 0+ | otherwise = fst <$> (getFieldValue tokenContentLength reqvt >>= C8.readInt) bodyLength <- newIORef 0 tlr <- newIORef Nothing q <- newTQueueIO writeIORef streamRxQ $ Just q- setStreamState ctx strm $ Open hcl (Body q mcl bodyLength tlr)+ setOpenState ctx strm $ Body q mcl bodyLength tlr -- FLOW CONTROL: WINDOW_UPDATE 0: recv: announcing my limit properly -- FLOW CONTROL: WINDOW_UPDATE: recv: announcing my limit properly bodySource <- mkSource q $ informWindowUpdate ctx strm let inpObj = InpObj tbl mcl (readSource bodySource) tlr if isServer ctx then do- let ServerInfo{..} = toServerInfo roleInfo- launch ctx strm inpObj+ launchHandler ctx strm inpObj else putMVar streamInput $ Right inpObj- return False --- Transition (process3)-processState s@(Open _ Continued{}) ctx strm _streamId = do- setStreamState ctx strm s- return True- -- Transition (process4) processState HalfClosedRemote ctx strm _streamId = do halfClosedRemote ctx strm- return False -- Transition (process5) processState (Closed cc) ctx strm _streamId = do closed ctx strm cc- return False -- Transition (process6)+processState (Open _ o) ctx strm _streamId =+ -- Open JustOpened, Open Body. Not the whole state: see 'setOpenState'.+ setOpenState ctx strm o processState s ctx strm _streamId = do- -- Idle, Open Body, Closed+ -- Idle setStreamState ctx strm s- return False +-- | Handing a request to the server's handler.+--+-- The stream counts towards the last stream identifier of our GOAWAY from+-- here on: from now its request "might have been processed" (RFC 9113,+-- section 6.8), and a client is free to retry, on another connection, any+-- request on a stream above it. It used to be counted once the handler+-- had returned, so a GOAWAY sent while handlers were running left their+-- streams out, and a client could send again a request -- a POST, say --+-- that had been acted on.+launchHandler :: Context -> Stream -> InpObj -> IO ()+launchHandler ctx@Context{roleInfo} strm inpObj = do+ modifyPeerLastStreamId ctx $ streamNumber strm+ let ServerInfo{..} = toServerInfo roleInfo+ launch ctx strm inpObj+ ---------------------------------------------------------------- {- FOURMOLU_DISABLE -}@@ -253,7 +407,14 @@ csid <- getPeerStreamID ctx if streamId <= csid -- consider the stream closed then- if ftyp `elem` [FrameWindowUpdate, FrameRSTStream, FramePriority]+ -- RFC 9113 section 5.1: "An endpoint MUST ignore frames that+ -- it receives on closed streams after it has sent a+ -- RST_STREAM frame." DATA is in that list because resetting a+ -- stream mid-body leaves whatever the peer already put on the+ -- wire still to arrive. HEADERS is not: that would be reuse+ -- of a stream identifier, which section 5.1.1 makes a+ -- connection error.+ if ftyp `elem` [FrameData, FrameWindowUpdate, FrameRSTStream, FramePriority] then return Nothing -- will be ignored else E.throwIO $@@ -270,9 +431,21 @@ `BS.append` C8.pack (show ftyp) ) E.throwIO $ ConnectionErrorIsSent ProtocolError streamId errmsg- when (ftyp == FrameHeaders) $ setPeerStreamID ctx streamId- -- FLOW CONTROL: SETTINGS_MAX_CONCURRENT_STREAMS: recv: rejecting if over my limit- Just <$> openOddStreamCheck ctx streamId ftyp+ if ftyp == FramePriority+ then+ -- PRIORITY does not open a stream (RFC 9113, section+ -- 5.1): it is checked and dropped like one for a+ -- stream we do not have. It used to create the+ -- stream, taking a concurrency slot that nothing ever+ -- gave back, since no HEADERS need follow: a peer+ -- could fill SETTINGS_MAX_CONCURRENT_STREAMS with+ -- PRIORITY frames alone and have every request after+ -- them refused.+ return Nothing+ else do+ setPeerStreamID ctx streamId+ -- FLOW CONTROL: SETTINGS_MAX_CONCURRENT_STREAMS: recv: rejecting if over my limit+ Just <$> openOddStreamCheck ctx streamId ftyp | otherwise = -- We received a frame from the server on an unknown stream -- (likely a previously created and then subsequently reset stream).@@ -317,10 +490,19 @@ if rate > pingRateLimit mySettings then E.throwIO $ ConnectionErrorIsSent EnhanceYourCalm streamId "too many ping" else sendPing ctx True bs-control FrameGoAway header bs _ = do+control FrameGoAway header bs ctx = do GoAwayFrame sid err msg <- guardIt $ decodeGoAwayFrame header bs if err == NoError- then E.throwIO ConnectionIsClosed+ then do+ first <- goingAway ctx sid+ -- The receiver goes on reading for the streams left, and stops+ -- the way it does when the peer closes the connection, once+ -- they are done.+ when first $ do+ receiver <- myThreadId+ T.forkManaged (threadManager ctx) "H2 draining after GOAWAY" $ do+ atomically $ drained ctx+ E.throwTo receiver ConnectionIsClosed else E.throwIO $ ConnectionErrorIsReceived err sid $ Short.toShort msg control FrameWindowUpdate header bs ctx = do WindowUpdateFrame n <- guardIt $ decodeWindowUpdateFrame header bs@@ -347,7 +529,9 @@ ProtocolError streamId "wrong header fragment for push promise"- (_, vt) <- hpackDecodeHeader frag streamId ctx+ -- A malformed promised request is an error on the promised stream,+ -- which 'resetPromised' resets.+ (_, vt) <- hpackDecodeHeader frag streamId ctx `E.catch` onPromised sid let ClientInfo{..} = toClientInfo $ roleInfo ctx when ( getFieldValue tokenAuthority vt == Just (UTF8.fromString authority)@@ -364,6 +548,10 @@ ---------------------------------------------------------------- +onPromised :: StreamId -> HTTP2Error -> IO a+onPromised sid (StreamErrorIsSent err _ msg) = E.throwIO $ StreamErrorIsSent err sid msg+onPromised _ e = E.throwIO e+ {-# INLINE guardIt #-} guardIt :: Either FrameDecodeError a -> IO a guardIt x = case x of@@ -379,6 +567,27 @@ where dep = streamDependency p +-- | Handle a decoded response HEADERS section. On the client, a 1xx+-- informational response (e.g. 103 Early Hints) is delivered to the+-- informational callback and the stream keeps waiting for the final response;+-- otherwise the headers become the (final) response.+onResponseHeaders+ :: Context+ -> StreamId+ -> Maybe ClosedCode+ -> Bool+ -> TokenHeaderTable+ -> IO StreamState+onResponseHeaders ctx streamId hcl endOfStream tbl+ | endOfStream = return $ Open hcl (NoBody tbl)+ | role ctx == Client && isInformational = do+ informationalCallback ctx streamId tbl+ return $ Open hcl JustOpened+ | otherwise = return $ Open hcl (HasBody tbl)+ where+ isInformational =+ maybe False ("1" `BS.isPrefixOf`) $ getFieldValue tokenStatus (snd tbl)+ stream :: FrameType -> FrameHeader@@ -408,41 +617,38 @@ if endOfHeader then do tbl <- hpackDecodeHeader frag streamId ctx- return $- if endOfStream- then -- turned into HalfClosedRemote in processState- Open hcl (NoBody tbl)- else Open hcl (HasBody tbl)+ onResponseHeaders ctx streamId hcl endOfStream tbl else do- let siz = BS.length frag- return $ Open hcl $ Continued [frag] siz 1 endOfStream+ startHeaderBlock ctx streamId endOfStream frag+ return s -- Transition (stream2)-stream FrameHeaders header@FrameHeader{flags, streamId} bs ctx (Open _ (Body q _ _ tlr)) _ = do+stream FrameHeaders header@FrameHeader{flags, streamId} bs ctx s@(Open _ (Body q _ _ tlr)) _ = do HeadersFrame _ frag <- guardIt $ decodeHeadersFrame header bs let endOfStream = testEndStream flags -- checking frag == "" is not necessary if endOfStream then do- tbl <- hpackDecodeTrailer frag streamId ctx- writeIORef tlr (Just tbl)- atomically $ writeTQueue q $ Right (mempty, True)- return HalfClosedRemote- else -- we don't support continuation here.+ if testEndHeader flags+ then onTrailers ctx streamId frag q tlr+ else do+ startHeaderBlock ctx streamId endOfStream frag+ return s+ else E.throwIO $ ConnectionErrorIsSent ProtocolError streamId- "continuation in trailer is not supported"+ "trailers without END_STREAM" -- Transition (stream4) stream FrameData header@FrameHeader{flags, payloadLength, streamId} bs- Context{emptyFrameRate, rxFlow, mySettings}+ ctx@Context{emptyFrameRate, rxFlow, mySettings} s@(Open _ (Body q mcl bodyLength _))- Stream{..} = do+ strm@Stream{..} = do DataFrame body <- guardIt $ decodeDataFrame header bs -- FLOW CONTROL: WINDOW_UPDATE 0: recv: rejecting if over my limit okc <- atomicModifyIORef' rxFlow $ checkRxLimit payloadLength@@ -460,8 +666,22 @@ EnhanceYourCalm streamId "exceeds stream flow-control limit"+ -- A push gives the connection window everything back now, padding+ -- and all ('connectionCreditedOnArrival'); its stream window, and+ -- any other stream's windows, as below.+ when (connectionCreditedOnArrival streamId) $+ giveBackConnectionWindow ctx payloadLength+ -- The padding is charged to both windows, as it must be, but it+ -- never reaches the reader, whose reading is what gives octets+ -- back ('readSource'). Left there, each padded frame shrank the+ -- peer's windows for good, until the connection stalled. It is+ -- done with as soon as it arrives, so it goes straight back.+ informWindowUpdate ctx strm $ payloadLength - BS.length body len0 <- readIORef bodyLength- let len = len0 + payloadLength+ -- The content itself: 'payloadLength', which flow control goes by,+ -- also counts the padding, and so made a padded body look longer+ -- than its content-length.+ let len = len0 + BS.length body endOfStream = testEndStream flags -- Empty Frame Flooding - CVE-2019-9518 if body == ""@@ -471,7 +691,22 @@ E.throwIO $ ConnectionErrorIsSent EnhanceYourCalm streamId "too many empty data" else do writeIORef bodyLength len- atomically $ writeTQueue q $ Right (body, endOfStream)+ -- Not for a stream closed since its state was read, by+ -- a reader that has done with it ('giveBackUnread') or a+ -- reset of ours: nothing would ever read it, or give it+ -- back to the connection window. Checked in the same+ -- transaction, so that a closer that has emptied the+ -- queue finds nothing put in after it.+ queued <- atomically $ do+ st <- readTVar streamState+ if isClosed st+ then return False+ else do+ writeTQueue q $ Right (body, endOfStream)+ return True+ unless (queued || connectionCreditedOnArrival streamId) $+ giveBackConnectionWindow ctx $+ BS.length body if endOfStream then do case mcl of@@ -488,39 +723,6 @@ return HalfClosedRemote else return s --- Transition (stream5)-stream FrameContinuation FrameHeader{flags, streamId} frag ctx s@(Open hcl (Continued rfrags siz n endOfStream)) _ = do- let endOfHeader = testEndHeader flags- if frag == "" && not endOfHeader- then do- -- Empty Frame Flooding - CVE-2019-9518- rate <- getRate $ emptyFrameRate ctx- if rate > emptyFrameRateLimit (mySettings ctx)- then- E.throwIO $- ConnectionErrorIsSent EnhanceYourCalm streamId "too many empty continuation"- else return s- else do- let rfrags' = frag : rfrags- siz' = siz + BS.length frag- n' = n + 1- when (siz' > headerFragmentLimit) $- E.throwIO $- ConnectionErrorIsSent EnhanceYourCalm streamId "Header is too big"- when (n' > continuationLimit) $- E.throwIO $- ConnectionErrorIsSent EnhanceYourCalm streamId "Header is too fragmented"- if endOfHeader- then do- let hdrblk = BS.concat $ reverse rfrags'- tbl <- hpackDecodeHeader hdrblk streamId ctx- return $- if endOfStream- then -- turned into HalfClosedRemote in processState- Open hcl (NoBody tbl)- else Open hcl (HasBody tbl)- else return $ Open hcl $ Continued rfrags' siz' n' endOfStream- -- (No state transition) stream FrameWindowUpdate header bs _ s strm = do WindowUpdateFrame n <- guardIt $ decodeWindowUpdateFrame header bs@@ -560,19 +762,32 @@ -- > Either endpoint can send a RST_STREAM frame from this state, causing it -- > to transition immediately to "closed". --- -- This justifies the two non-error cases, below. (Section 8.1 of the spec+ -- This justifies the non-error cases, below. (Section 8.1 of the spec -- is also relevant, but it is less explicit about the /either endpoint/ -- part.)+ --+ -- The error code the peer sent does not enter into it. Receiving a+ -- RST_STREAM closes that stream and nothing else, whatever the reason+ -- given; the code is for whoever is reading the stream, and reaches them+ -- as 'StreamResetIsReceived' by way of 'closed' above. Ending the whole+ -- connection over it would punish every other stream on the connection+ -- for a peer's complaint about one. case s of- Open _ _- | isNonCritical err ->- -- Open /or/ half-closed (local)- return (Closed cc)- HalfClosedRemote- | isNonCritical err ->- return (Closed cc)- _otherwise -> do- E.throwIO $ StreamErrorIsReceived err streamId+ -- Open /or/ half-closed (local)+ Open _ _ -> return (Closed cc)+ HalfClosedRemote -> return (Closed cc)+ Reserved -> return (Closed cc)+ Closed _ -> return (Closed cc)+ -- Only an idle stream is left, which a PRIORITY frame can have+ -- created. Section 5.1 again, on "idle": "Receiving any frame other+ -- than HEADERS or PRIORITY on a stream in this state MUST be treated+ -- as a connection error (Section 5.4.1) of type PROTOCOL_ERROR."+ Idle ->+ E.throwIO $+ ConnectionErrorIsSent+ ProtocolError+ streamId+ "rst_stream on an idle stream" -- (No state transition) stream FramePriority header bs _ s Stream{streamNumber} = do -- ignore@@ -585,15 +800,20 @@ stream FrameContinuation FrameHeader{streamId} _ _ _ _ = E.throwIO $ ConnectionErrorIsSent ProtocolError streamId "continue frame cannot come here"-stream _ FrameHeader{streamId} _ _ (Open _ Continued{}) _ =- E.throwIO $- ConnectionErrorIsSent- ProtocolError- streamId- "an illegal frame follows header/continuation frames"+-- DATA that is not taken in, below, is still paid for. RFC 9113 section+-- 6.9: "A receiver that receives a flow-controlled frame MUST always account+-- for its contribution against the connection flow-control window, unless+-- the receiver treats this as a connection error." The peer charged it to+-- the connection window before sending it, and unless it is charged and+-- given back here as well, the peer's view of that window shrinks for good.+-- -- Ignore frames to streams we have just reset, per section 5.1.+stream FrameData FrameHeader{payloadLength, streamId} _ ctx st@(Closed (ResetByMe _)) _ = do+ informIgnoredData ctx streamId payloadLength+ return st stream _ _ _ _ st@(Closed (ResetByMe _)) _ = return st-stream FrameData FrameHeader{streamId} _ _ _ _ =+stream FrameData FrameHeader{payloadLength, streamId} _ ctx _ _ = do+ informIgnoredData ctx streamId payloadLength E.throwIO $ StreamErrorIsSent StreamClosed streamId $ fromString ("illegal data frame for " ++ show streamId)@@ -602,17 +822,6 @@ StreamErrorIsSent ProtocolError streamId $ fromString ("illegal frame " ++ show x ++ " for " ++ show streamId) -{- FOURMOLU_DISABLE -}--- Although some stream errors indicate misbehaving peers, such as--- FLOW_CONTROL_ERROR, not all errors do. We will close the connection only--- for critical errors.-isNonCritical :: ErrorCode -> Bool-isNonCritical NoError = True-isNonCritical Cancel = True-isNonCritical InternalError = True-isNonCritical _ = False-{- FOURMOLU_ENABLE -}- ---------------------------------------------------------------- -- | Type for input streaming.@@ -687,12 +896,7 @@ return $ goawayFrame sid err msg sendGoaway :: Config -> ByteString -> IO ()-sendGoaway Config{..} frame = confSendAll frame `E.catch` ignore--ignore :: E.SomeException -> IO ()-ignore (E.SomeException e)- | isAsyncException e = E.throwIO e- | otherwise = return ()+sendGoaway Config{..} frame = confSendAll frame `E.catchIOError` \_ -> return () ---------------------------------------------------------------- @@ -700,3 +904,69 @@ sendPing Context{..} ack bs = enqueueControl controlQ $ CFrames Nothing [frame] where frame = pingFrame ack bs++----------------------------------------------------------------++-- | Deliver a complete trailer block: the body ends with it.+onTrailers+ :: Context+ -> StreamId+ -> HeaderBlockFragment+ -> TQueue (Either E.SomeException (ByteString, Bool))+ -> IORef (Maybe TokenHeaderTable)+ -> IO StreamState+onTrailers ctx streamId blk q tlr = do+ tbl <- hpackDecodeTrailer blk streamId ctx+ writeIORef tlr (Just tbl)+ atomically $ writeTQueue q $ Right (mempty, True)+ return HalfClosedRemote++-- | Start accumulating a header block that does not fit in a single frame+startHeaderBlock+ :: Context+ -> StreamId+ -> Bool+ -- ^ END_STREAM, from the HEADERS frame+ -> HeaderBlockFragment+ -- ^ The fragment in the HEADERS frame+ -> IO ()+startHeaderBlock Context{continued} streamId endOfStream frag =+ writeIORef continued . Just $+ HeaderContinuation+ { hcStreamId = streamId+ , hcBlock = newPartialHeaderBlock frag+ , hcEndOfStream = endOfStream+ }++newPartialHeaderBlock :: HeaderBlockFragment -> PartialHeaderBlock+newPartialHeaderBlock frag =+ PartialHeaderBlock+ { phbFragments = [frag]+ , phbTotalSize = BS.length frag+ , phbNumFrames = 1+ }++addFragment+ :: StreamId+ -- ^ Used for error messages only+ -> HeaderBlockFragment+ -> PartialHeaderBlock+ -> IO PartialHeaderBlock+addFragment streamId frag phb = do+ when (phbTotalSize phb' > headerFragmentLimit) $+ E.throwIO $+ ConnectionErrorIsSent EnhanceYourCalm streamId "Header is too big"+ when (phbNumFrames phb' > continuationLimit) $+ E.throwIO $+ ConnectionErrorIsSent EnhanceYourCalm streamId "Header is too fragmented"+ return phb'+ where+ phb' =+ PartialHeaderBlock+ { phbFragments = frag : phbFragments phb+ , phbTotalSize = phbTotalSize phb + BS.length frag+ , phbNumFrames = phbNumFrames phb + 1+ }++completeHeaderBlock :: PartialHeaderBlock -> HeaderBlockFragment+completeHeaderBlock = BS.concat . reverse . phbFragments
Network/HTTP2/H2/Sender.hs view
@@ -10,13 +10,14 @@ import Control.Concurrent.STM import qualified Control.Exception as E+import qualified Data.ByteString as BS+import qualified Data.ByteString.Lazy as BS.Lazy import Data.IORef (modifyIORef', readIORef, writeIORef) import Data.IntMap.Strict (IntMap)-import Foreign.Ptr (minusPtr, plusPtr)+import Foreign.Ptr (castPtr, minusPtr, plusPtr) import Network.ByteOrder import Network.HTTP.Semantics.Client import Network.HTTP.Semantics.IO-import System.ThreadManager import Imports import Network.HPACK (setLimitForEncoding, toTokenHeaderTable)@@ -57,40 +58,60 @@ where updateAllStreamTxFlow :: WindowSize -> IntMap Stream -> IO () updateAllStreamTxFlow siz strms =- forM_ strms $ \strm -> increaseStreamWindowSize strm siz+ forM_ strms $ \strm -> increaseStreamWindowSize strm siz `E.catch` connectionError+ -- RFC 9113, section 6.9.2: "An endpoint MUST treat a change to+ -- SETTINGS_INITIAL_WINDOW_SIZE that causes any flow-control window to+ -- exceed the maximum size as a connection error of type+ -- FLOW_CONTROL_ERROR." The same overflow from a WINDOW_UPDATE is a stream+ -- error, which is what 'increaseStreamWindowSize' raises.+ connectionError (StreamErrorIsSent err sid msg) =+ E.throwIO $ ConnectionErrorIsSent err sid msg+ connectionError e = E.throwIO e -checkDone :: Context -> Int -> IO Bool+checkDone :: Context -> Int -> IO (Maybe E.SomeException) checkDone Context{..} 0 = atomically $ do isEmptyC <- isEmptyTQueue controlQ isEmptyO <- isEmptyTQueue outputQ if not isEmptyC || not isEmptyO then- return False+ return Nothing else do- gone <- isAllGone threadManager- unless gone retry- done <- readTVar receiverDone- unless done retry- return True-checkDone _ _ = return False+ recv <- readTVar receiverDone+ case recv of+ Just done ->+ return $ Just done+ _otherwise ->+ retry+checkDone _ _ = return Nothing -frameSender :: Context -> Config -> IO ()+frameSender :: Context -> Config -> IO E.SomeException frameSender ctx@Context{outputQ, controlQ, encodeDynamicTable, outputBufferLimit} Config{..} = do labelMe "H2 sender"- loop 0+ -- DATA that has to wait for the connection window, in the order it+ -- came: see 'dequeue'.+ parked <- newTVarIO []+ -- This catches an asynchronous exception.+ -- It is re-thrown by "runH2"+ loop parked 0 `E.catch` return where ----------------------------------------------------------------- loop :: Offset -> IO ()- loop off = do- done <- checkDone ctx off- unless done $ do- x <- atomically $ dequeue off- case x of- C ctl -> flushN off >> control ctl >> loop 0- O out -> outputAndSync out off >>= flushIfNecessary >>= loop- Flush -> flushN off >> loop 0+ loop :: TVar [Output] -> Offset -> IO E.SomeException+ loop parked off = do+ mDone <- checkDone ctx off+ case mDone of+ Just done ->+ return done+ Nothing -> do+ x <- atomically $ dequeue parked off+ case x of+ C ctl -> flushN off >> control ctl >> loop parked 0+ O out ->+ outputAndSync parked out off+ >>= flushIfNecessary+ >>= loop parked+ Flush -> flushN off >> loop parked 0 -- Flush the connection buffer to the socket, where the first 'n' bytes of -- the buffer are filled.@@ -107,17 +128,31 @@ flushN off return 0 - dequeue :: Offset -> STM Switch- dequeue off = do+ -- Only DATA is flow-controlled (RFC 9113, section 6.9), so only DATA+ -- waits for the connection window. The whole output queue used to:+ -- with the window shut, no HEADERS, PUSH_PROMISE or RST_STREAM went+ -- out either, on any stream -- not even the response to a request+ -- that has no body, or the reset that would have freed some of the+ -- window. Now outputs are taken as they come, and DATA there is no+ -- connection window for is parked ('outputAndSync') and goes out,+ -- first and in order, once there is.+ dequeue :: TVar [Output] -> Offset -> STM Switch+ dequeue parked off = do isEmptyC <- isEmptyTQueue controlQ if isEmptyC then do- -- FLOW CONTROL: WINDOW_UPDATE 0: send: respecting peer's limit- waitConnectionWindowSize ctx- isEmptyO <- isEmptyTQueue outputQ- if isEmptyO- then if off /= 0 then return Flush else retry- else O <$> readTQueue outputQ+ ps <- readTVar parked+ cws <- connectionWindowSizeSTM ctx+ case ps of+ -- FLOW CONTROL: WINDOW_UPDATE 0: send: respecting peer's limit+ p : rest | cws > 0 -> do+ writeTVar parked rest+ return $ O p+ _ -> do+ isEmptyO <- isEmptyTQueue outputQ+ if isEmptyO+ then if off /= 0 then return Flush else retry+ else O <$> readTQueue outputQ else C <$> readTQueue controlQ ----------------------------------------------------------------@@ -154,24 +189,74 @@ -- -- Both the stream window and the connection window are open. ----------------------------------------------------------------- outputAndSync :: Output -> Offset -> IO Offset- outputAndSync out@(Output strm otyp sync) off = E.handle (\e -> resetStream strm InternalError e >> return off) $ do+ outputAndSync :: TVar [Output] -> Output -> Offset -> IO Offset+ -- "handler" catches an asynchronous exception and+ -- re-throws it.+ outputAndSync parked out@(Output strm otyp sync) off = E.handle (handler strm off) $ do state <- readStreamState strm if isHalfClosedLocal state- then return off+ then do+ case otyp of+ OReset mErr+ | not (isClosed state) ->+ -- RST_STREAM is the only frame we can still send+ -- after half-closing+ resetStreamWith strm mErr+ _otherwise ->+ return ()+ -- Nothing more can go out on this stream, but whoever+ -- enqueued this output is waiting in 'syncWithSender'' to+ -- be told so. Dropping the notification parked that+ -- thread on an MVar nothing would ever fill, until the+ -- timeout manager killed it -- one stranded worker per+ -- stream the peer resets while a response is in flight.+ sync Nothing+ return off else case otyp of OHeader hdr mnext tlrmkr -> do (off', mout') <- outputHeader strm hdr mnext tlrmkr sync off sync mout' return off'+ OInformational hdr -> do+ off' <- outputInformational strm hdr off+ sync Nothing+ return off'+ OReset mErr -> do+ resetStreamWith strm mErr+ sync Nothing+ return off _ -> do sws <- getStreamWindowSize strm- cws <- getConnectionWindowSize ctx -- not 0+ cws <- getConnectionWindowSize ctx let lim = min cws sws- (off', mout') <- output out off lim- sync mout'- return off'+ case otyp of+ ONext{}+ | cws <= 0 -> do+ -- To wait for the connection window,+ -- without holding up anything else.+ atomically $ modifyTVar' parked (++ [out])+ return off+ | lim <= 0 -> do+ -- No room for any of the body: the+ -- window was shut after this was queued+ -- (a SETTINGS_INITIAL_WINDOW_SIZE+ -- decrease, say). Filling a DATA frame+ -- into no room reads 0 octets of a file,+ -- which is taken for its end; handed back+ -- instead, it is queued again once the+ -- window opens.+ sync $ Just out+ return off+ _ -> do+ (off', mout') <- output out off lim+ sync mout'+ return off' + ----------------------------------------------------------------+ handler strm off e = do+ resetStream strm InternalError e+ return off+ resetStream :: Stream -> ErrorCode -> E.SomeException -> IO () resetStream strm err e | isAsyncException e = E.throwIO e@@ -180,6 +265,12 @@ let rst = resetFrame err $ streamNumber strm enqueueControl controlQ $ CFrames Nothing [rst] + resetStreamWith :: Stream -> Maybe E.SomeException -> IO ()+ resetStreamWith strm (Just err) =+ resetStream strm InternalError err+ resetStreamWith strm Nothing =+ resetStream strm Cancel (E.toException CancelledStream)+ ---------------------------------------------------------------- outputHeader :: Stream@@ -208,6 +299,21 @@ return (off, Just out') ----------------------------------------------------------------+ -- Emit an informational (1xx) HEADERS section. Unlike 'outputHeader',+ -- this never sets END_STREAM and never half-closes the stream, so the+ -- final response can still be sent afterwards.+ outputInformational+ :: Stream+ -> [Header]+ -> Offset+ -> IO Offset+ outputInformational strm hdr off0 = do+ let sid = streamNumber strm+ (ths, _) <- toTokenHeaderTable $ fixHeaders hdr+ off' <- headerContinue sid ths False {- not endOfStream -} off0+ flushIfNecessary off'++ ---------------------------------------------------------------- output :: Output -> Offset -> WindowSize -> IO (Offset, Maybe Output) output out@(Output strm (ONext curr tlrmkr) _) off0 lim = do -- Data frame payload@@ -217,7 +323,10 @@ datBufSiz = buflim - payloadOff curr datBuf (min datBufSiz lim) >>= \case Next datPayloadLen reqflush mnext -> do- NextTrailersMaker tlrmkr' <- runTrailersMaker tlrmkr datBuf datPayloadLen+ tm <- runTrailersMaker tlrmkr datBuf datPayloadLen+ let tlrmkr' = case tm of+ NextTrailersMaker t -> t+ _ -> defaultTrailersMaker fillDataHeader strm off0@@ -237,11 +346,7 @@ -- outputs for this stream already enqueued. Therefore, we can -- safely cancel it knowing that we won't try and send any -- more data frames on this stream.- case mErr of- Just err ->- resetStream strm InternalError err- Nothing ->- resetStream strm Cancel (E.toException CancelledStream)+ resetStreamWith strm mErr return (off0, Nothing) output (Output strm (OPush ths pid) _) off0 _lim = do -- Creating a push promise header@@ -259,22 +364,31 @@ let offkv = off0 + frameHeaderLength bufkv = confWriteBuffer `plusPtr` offkv limkv = buflim - offkv- (ths, kvlen) <- hpackEncodeHeader ctx bufkv limkv ths0- if kvlen == 0- then continue off0 ths FrameHeaders+ -- Most blocks fit where they are going: encode in place, which+ -- is one HEADERS frame and no copying.+ (rest, kvlen) <- hpackEncodeHeader ctx bufkv limkv ths0+ if null rest+ then do+ let buf = confWriteBuffer `plusPtr` off0+ fillFrameHeader FrameHeaders kvlen sid (getFlag FrameHeaders BS.Lazy.empty) buf+ return $ offkv + kvlen else do- let flag = getFlag ths- buf = confWriteBuffer `plusPtr` off0- off = offkv + kvlen- fillFrameHeader FrameHeaders kvlen sid flag buf- continue off ths FrameContinuation+ -- It did not fit. What was written is the start of the+ -- block, and the dynamic table has taken it into+ -- account, so it is kept; the rest is encoded after it,+ -- and the whole block then starts in a fresh buffer to+ -- avoid emitting a tiny HEADERS frame.+ start <- BS.packCStringLen (castPtr bufkv, kvlen)+ ths1 <- hpackEncodeHeaderRest ctx (buflim - frameHeaderLength) rest+ continue off0 (BS.Lazy.fromStrict start <> ths1) FrameHeaders where eos = if endOfStream then setEndStream else id- getFlag [] = eos $ setEndHeader defaultFlags- getFlag _ = eos defaultFlags+ getFlag ft ths =+ (if ft == FrameHeaders then eos else id) $+ if BS.Lazy.null ths then setEndHeader defaultFlags else defaultFlags - continue :: Offset -> TokenHeaderList -> FrameType -> IO Offset- continue off [] _ = return off+ continue :: Offset -> BS.Lazy.ByteString -> FrameType -> IO Offset+ continue off ths _ | BS.Lazy.null ths = return off continue off ths ft = do flushN off -- Now off is 0@@ -283,15 +397,20 @@ headerPayloadLim = buflim - frameHeaderLength (ths', kvlen') <-- hpackEncodeHeaderLoop ctx bufHeaderPayload headerPayloadLim ths- when (ths == ths') $- E.throwIO $- ConnectionErrorIsSent CompressionError sid "cannot compress the header"- let flag = getFlag ths'+ copyFragment bufHeaderPayload headerPayloadLim ths+ let flag = getFlag ft ths' off' = frameHeaderLength + kvlen' fillFrameHeader ft kvlen' sid flag confWriteBuffer continue off' ths' FrameContinuation + -- Copy as much of the block as fits; return the rest and the number of bytes copied+ copyFragment+ :: Buffer -> Int -> BS.Lazy.ByteString -> IO (BS.Lazy.ByteString, Int)+ copyFragment buf lim ths = do+ let (frag, rest) = BS.Lazy.splitAt (fromIntegral (max 0 lim)) ths+ _ <- foldM copy buf (BS.Lazy.toChunks frag)+ return (rest, fromIntegral (BS.Lazy.length frag))+ ---------------------------------------------------------------- fillDataHeader :: Stream@@ -312,7 +431,10 @@ reqflush = do let buf = confWriteBuffer `plusPtr` off (mtrailers, flag) <- do- Trailers trailers <- tlrmkr Nothing+ tm <- tlrmkr Nothing+ let trailers = case tm of+ Trailers t -> t+ _ -> [] if null trailers then return (Nothing, setEndStream defaultFlags) else return (Just trailers, defaultFlags)
Network/HTTP2/H2/Settings.hs view
@@ -29,7 +29,9 @@ , settingsRateLimit :: Int -- ^ Maximum number of settings frames allowed per second (CVE-2019-9515) , rstRateLimit :: Int- -- ^ Maximum number of reset frames allowed per second (CVE-2023-44487)+ -- ^ Maximum number of streams reset per second, whether by the peer's+ -- RST_STREAM (CVE-2023-44487) or by ours in answer to a stream error+ -- the peer caused (CVE-2025-8671) } deriving (Eq, Show)
Network/HTTP2/H2/Stream.hs view
@@ -1,18 +1,14 @@-{-# LANGUAGE DeriveAnyClass #-}-{-# LANGUAGE DerivingStrategies #-} {-# LANGUAGE NamedFieldPuns #-}-{-# LANGUAGE RankNTypes #-} module Network.HTTP2.H2.Stream where import Control.Concurrent import Control.Concurrent.STM-import Control.Exception+import qualified Control.Exception as E import Control.Monad import Data.IORef import Data.Maybe (fromMaybe) import Network.Control-import Network.HTTP.Semantics import Network.HTTP.Semantics.IO import Network.HTTP2.Frame@@ -52,37 +48,42 @@ newOddStream :: StreamId -> WindowSize -> WindowSize -> IO Stream newOddStream sid txwin rxwin = Stream sid- <$> newIORef Idle+ <$> newTVarIO Idle <*> newEmptyMVar <*> newTVarIO (newTxFlow txwin) <*> newIORef (newRxFlow rxwin) <*> newIORef Nothing+ <*> newIORef Nothing newEvenStream :: StreamId -> WindowSize -> WindowSize -> IO Stream newEvenStream sid txwin rxwin = Stream sid- <$> newIORef Reserved+ <$> newTVarIO Reserved <*> newEmptyMVar <*> newTVarIO (newTxFlow txwin) <*> newIORef (newRxFlow rxwin) <*> newIORef Nothing+ <*> newIORef Nothing ---------------------------------------------------------------- {-# INLINE readStreamState #-} readStreamState :: Stream -> IO StreamState-readStreamState Stream{streamState} = readIORef streamState+readStreamState Stream{streamState} = readTVarIO streamState ---------------------------------------------------------------- closeAllStreams- :: TVar OddStreamTable -> TVar EvenStreamTable -> Maybe SomeException -> IO ()-closeAllStreams ovar evar mErr' = do+ :: TVar OddStreamTable -> TVar EvenStreamTable -> Maybe E.SomeException -> IO ()+closeAllStreams ovar evar mErr = do ostrms <- clearOddStreamTable ovar mapM_ finalize ostrms estrms <- clearEvenStreamTable evar mapM_ finalize estrms where+ -- We treat /every/ exception, including 'ConectionIsClosed', as abnormal+ -- termination: we should only report a clean termination when we receive an+ -- explicit @END_STREAM@ frame. finalize strm = do st <- readStreamState strm void $ tryPutMVar (streamInput strm) err@@ -92,75 +93,10 @@ _otherwise -> return () - mErr :: Maybe SomeException- mErr = case mErr' of- Just e- | Just ConnectionIsClosed <- fromException e ->- Nothing- _otherwise ->- mErr'-- err :: Either SomeException a- err = Left $ fromMaybe (toException ConnectionIsClosed) mErr+ err :: Either E.SomeException a+ err = Left $ fromMaybe (E.toException ConnectionIsClosed) mErr ------------------------------------------------------------------data StreamTerminated- = StreamPushedFinal- | StreamCancelled- | StreamOutOfScope- deriving (Show)- deriving anyclass (Exception)--withOutBodyIface- :: TBQueue StreamingChunk- -> (forall a. IO a -> IO a)- -> (OutBodyIface -> IO r)- -> IO r-withOutBodyIface tbq unmask k = do- terminated <- newTVarIO Nothing- let whenNotTerminated act = do- mTerminated <- readTVar terminated- maybe act throwSTM mTerminated-- terminateWith reason act = do- mTerminated <- readTVar terminated- case mTerminated of- Just _ ->- -- Already terminated- return ()- Nothing -> do- writeTVar terminated (Just reason)- act-- iface =- OutBodyIface- { outBodyUnmask = unmask- , outBodyPush = \b ->- atomically $- whenNotTerminated $- writeTBQueue tbq $- StreamingBuilder b NotEndOfStream- , outBodyPushFinal = \b ->- atomically $ whenNotTerminated $ do- writeTVar terminated (Just StreamPushedFinal)- writeTBQueue tbq $ StreamingBuilder b (EndOfStream Nothing)- writeTBQueue tbq $ StreamingFinished Nothing- , outBodyFlush =- atomically $- whenNotTerminated $- writeTBQueue tbq StreamingFlush- , outBodyCancel =- atomically- . terminateWith StreamCancelled- . writeTBQueue tbq- . StreamingCancelled- }- finished = atomically $ do- terminateWith StreamOutOfScope $- writeTBQueue tbq $- StreamingFinished Nothing- k iface `finally` finished nextForStreaming :: TBQueue StreamingChunk
Network/HTTP2/H2/StreamTable.hs view
@@ -33,7 +33,7 @@ import Control.Concurrent import Control.Concurrent.STM-import Control.Exception+import qualified Control.Exception as E import Data.IntMap.Strict (IntMap) import qualified Data.IntMap.Strict as IntMap import Network.Control (LRUCache)@@ -78,7 +78,13 @@ let oddTable' = IntMap.insert k v oddTable in OddStreamTable oddConc oddTable' -deleteOdd :: TVar OddStreamTable -> IntMap.Key -> SomeException -> IO ()+-- | Remove a stream and give its concurrency slot back.+--+-- 'closed' can be called more than once for the same stream -- a RST_STREAM+-- carrying a non-critical error code goes through both 'stream' and+-- 'processState', each of which closes it -- so the count must follow an+-- entry that was really there, not the number of calls.+deleteOdd :: TVar OddStreamTable -> IntMap.Key -> E.SomeException -> IO () deleteOdd var k err = do mv <- atomically deleteStream case mv of@@ -88,10 +94,13 @@ deleteStream :: STM (Maybe Stream) deleteStream = do OddStreamTable{..} <- readTVar var- let oddConc' = oddConc - 1- oddTable' = IntMap.delete k oddTable- writeTVar var $ OddStreamTable oddConc' oddTable'- return $ IntMap.lookup k oddTable+ case IntMap.lookup k oddTable of+ Nothing -> return Nothing+ Just v -> do+ let oddConc' = oddConc - 1+ oddTable' = IntMap.delete k oddTable+ writeTVar var $ OddStreamTable oddConc' oddTable'+ return $ Just v lookupOdd :: TVar OddStreamTable -> IntMap.Key -> IO (Maybe Stream) lookupOdd var k = IntMap.lookup k . oddTable <$> readTVarIO var@@ -128,7 +137,9 @@ let evenTable' = IntMap.insert k v evenTable in EvenStreamTable evenConc evenTable' evenCache -deleteEven :: TVar EvenStreamTable -> IntMap.Key -> SomeException -> IO ()+-- | Remove a stream and give its concurrency slot back.+-- Idempotent, for the same reason as 'deleteOdd'.+deleteEven :: TVar EvenStreamTable -> IntMap.Key -> E.SomeException -> IO () deleteEven var k err = do mv <- atomically deleteStream case mv of@@ -138,10 +149,13 @@ deleteStream :: STM (Maybe Stream) deleteStream = do EvenStreamTable{..} <- readTVar var- let evenConc' = evenConc - 1- evenTable' = IntMap.delete k evenTable- writeTVar var $ EvenStreamTable evenConc' evenTable' evenCache- return $ IntMap.lookup k evenTable+ case IntMap.lookup k evenTable of+ Nothing -> return Nothing+ Just v -> do+ let evenConc' = evenConc - 1+ evenTable' = IntMap.delete k evenTable+ writeTVar var $ EvenStreamTable evenConc' evenTable' evenCache+ return $ Just v lookupEven :: TVar EvenStreamTable -> IntMap.Key -> IO (Maybe Stream) lookupEven var k = IntMap.lookup k . evenTable <$> readTVarIO var
Network/HTTP2/H2/Sync.hs view
@@ -1,3 +1,5 @@+{-# LANGUAGE MultiWayIf #-}+{-# LANGUAGE NamedFieldPuns #-} {-# LANGUAGE RecordWildCards #-} module Network.HTTP2.H2.Sync (@@ -15,6 +17,7 @@ import Control.Monad import Network.Control import Network.HTTP.Semantics.IO+import qualified System.ThreadManager as T import Network.HTTP2.H2.Context import Network.HTTP2.H2.Queue@@ -46,14 +49,32 @@ } return (pop, out) -makeOutputIO :: Context -> Stream -> OutputType -> Output-makeOutputIO Context{..} strm otyp = out+-- | An output for the 'runIO' interfaces, which have no thread waiting to+-- put the rest of a body back on the queue.+--+-- The rest used to go back at once, whatever the stream's window. With+-- none left, the sender filled a DATA frame into no room; a file read into+-- no room reads 0 octets, which is the end of the file, so a body larger than+-- the window went out cut short with END_STREAM. A streaming body with+-- nothing queued made the sender spin instead. So the rest goes back once+-- it can go on, the way 'syncWithSender'' does it for the other interfaces;+-- only when it has to wait is a thread used for it.+makeOutputIO+ :: Context -> Stream -> Maybe (TBQueue StreamingChunk) -> OutputType -> Output+makeOutputIO Context{..} strm mtbq otyp = out where push mout = case mout of Nothing -> return ()- -- Sender enqueues output again ignoring- -- the stream TX window.- Just ot -> enqueueOutput outputQ ot+ Just ot -> do+ now <- atomically $ (Just <$> ready) `orElse` return Nothing+ case now of+ Just True -> enqueueOutput outputQ ot+ Just False -> return ()+ Nothing ->+ T.forkManaged threadManager "H2 output waiting for its window" $ do+ ok <- atomically ready+ when ok $ enqueueOutput outputQ ot+ ready = readyToContinue strm mtbq out = Output { outputStream = strm@@ -61,9 +82,22 @@ , outputSync = push } +-- | Whether the rest of a stream's body can go on: waiting while the+-- stream's window is shut or a streaming body has nothing queued, and 'False'+-- once the stream is closed.+readyToContinue :: Stream -> Maybe (TBQueue StreamingChunk) -> STM Bool+readyToContinue Stream{streamState, streamTxFlow} mtbq = do+ state <- readTVar streamState+ case state of+ Closed{} -> return False+ _ -> do+ waitStreaming' mtbq+ waitStreamWindowSizeSTM streamTxFlow+ return True+ enqueueOutputSIO :: Context -> Stream -> OutputType -> IO () enqueueOutputSIO ctx@Context{..} strm otyp = do- let out = makeOutputIO ctx strm otyp+ let out = makeOutputIO ctx strm Nothing otyp enqueueOutput outputQ out syncWithSender' :: Context -> IO Sync -> LoopCheck -> IO ()@@ -76,7 +110,6 @@ Cont newout -> do cont <- checkLoop lc when cont $ do- -- This is justified by the precondition above enqueueOutput outputQ newout loop @@ -85,13 +118,15 @@ tovar <- newTVarIO False return $ LoopCheck- { lcTBQ = mtbq+ { lcState = streamState strm+ , lcTBQ = mtbq , lcTimeout = tovar , lcWindow = streamTxFlow strm } data LoopCheck = LoopCheck- { lcTBQ :: Maybe (TBQueue StreamingChunk)+ { lcState :: TVar StreamState+ , lcTBQ :: Maybe (TBQueue StreamingChunk) , lcTimeout :: TVar Bool , lcWindow :: TVar TxFlow }@@ -99,9 +134,11 @@ checkLoop :: LoopCheck -> IO Bool checkLoop LoopCheck{..} = atomically $ do tout <- readTVar lcTimeout- if tout- then return False- else do+ state <- readTVar lcState+ if+ | tout -> return False+ | Closed{} <- state -> return False+ | otherwise -> do waitStreaming' lcTBQ waitStreamWindowSizeSTM lcWindow return True
Network/HTTP2/H2/Types.hs view
@@ -7,13 +7,9 @@ import Control.Concurrent import Control.Concurrent.STM-import Control.Exception (- Exception,- SomeAsyncException (..),- SomeException (..),- ) import qualified Control.Exception as E import Data.IORef+import Foreign.Ptr (nullPtr) import Network.Control import Network.HTTP.Semantics.Client import Network.HTTP.Semantics.IO@@ -49,30 +45,15 @@ is labelled with the relevant case in either the function 'stream' or the function 'processState'. -> [Open JustOpened]-> |-> |-> HEADERS-> |-> | (stream1)-> |-> END_HEADERS?-> |-> ______/ \______-> / yes no \-> | |-> | [Open Continued] <--\-> | | |-> | CONTINUATION |-> | | |-> | | (stream5) |-> | | |-> | END_HEADERS? |-> | | |-> v yes / \ no |-> END_STREAM? <-------/ \-----------/-> | (process3)+> [Open JustOpened] > |+> |+> HEADERS CONTINUATION*+> |+> | (stream1)+> |+> END_STREAM?+> | > _________/ \_________ > / yes no \ > | |@@ -84,7 +65,7 @@ > | | | | > | | +---------------\ | > RST_STREAM | | | |-> | | HEADERS DATA |+> | | HEADERS CONT* DATA | > | (stream6) | | | | > | | | (stream2) | (stream4) | > | (process5) | | | |@@ -105,25 +86,29 @@ data OpenState = JustOpened- | Continued- [HeaderBlockFragment]- Int -- Total size- Int -- The number of continuation frames- Bool -- End of stream | NoBody TokenHeaderTable | HasBody TokenHeaderTable | Body- (TQueue (Either SomeException (ByteString, Bool)))+ (TQueue (Either E.SomeException (ByteString, Bool))) (Maybe Int) -- received Content-Length -- compared the body length for error checking (IORef Int) -- actual body length (IORef (Maybe TokenHeaderTable)) -- trailers +-- | Header block fragments accumulated so far.+--+-- Fragments are stored in reverse order (newest first).+data PartialHeaderBlock = PartialHeaderBlock+ { phbFragments :: [HeaderBlockFragment]+ , phbTotalSize :: Int+ , phbNumFrames :: Int+ }+ data ClosedCode = Finished | Killed | Reset ErrorCode- | ResetByMe SomeException+ | ResetByMe E.SomeException deriving (Show) -- | Used for streams which are cancelled by calling@@ -136,7 +121,7 @@ case cc of Finished -> ConnectionIsClosed Killed -> ConnectionIsTimeout- Reset err -> ConnectionErrorIsReceived err sid "Connection was reset"+ Reset err -> StreamResetIsReceived err sid ResetByMe err -> BadThingHappen err ----------------------------------------------------------------@@ -162,11 +147,14 @@ data Stream = Stream { streamNumber :: StreamId- , streamState :: IORef StreamState- , streamInput :: MVar (Either SomeException InpObj) -- Client only+ , streamState :: TVar StreamState+ , streamInput :: MVar (Either E.SomeException InpObj) -- Client only , streamTxFlow :: TVar TxFlow , streamRxFlow :: IORef RxFlow , streamRxQ :: IORef (Maybe RxQ)+ , streamRequestMethod :: IORef (Maybe ByteString)+ -- ^ Client only: the method of the request, which decides whether the+ -- response may have content at all (RFC 9110, section 6.4.1) } instance Show Stream where@@ -174,7 +162,7 @@ "Stream{id=" ++ show streamNumber ++ ",state="- ++ show (unsafePerformIO (readIORef streamState))+ ++ show (unsafePerformIO (readTVarIO streamState)) ++ "}" ----------------------------------------------------------------@@ -189,6 +177,8 @@ = OHeader [Header] (Maybe DynaNext) TrailersMaker | OPush TokenHeaderList StreamId -- associated stream id from client | ONext DynaNext TrailersMaker+ | OInformational [Header]+ | OReset (Maybe E.SomeException) data Sync = Done | Cont Output @@ -201,8 +191,12 @@ type ReasonPhrase = ShortByteString -- | The connection error or the stream error.--- Stream errors are treated as connection errors since--- there are no good recovery ways.+-- A stream error resets that stream and the connection carries on, as+-- RFC 9113 section 5.4.2 asks. One kind of trouble is a connection error+-- even though the spec calls it a stream error, because this+-- implementation cannot carry on through it: a field block abandoned+-- part-way leaves the HPACK tables disagreeing with the peer's, and+-- nothing sent afterwards would decode. -- `ErrorCode` in connection errors should be the highest stream identifier -- but in this implementation it identifies the stream that -- caused this error.@@ -212,6 +206,7 @@ | ConnectionErrorIsReceived ErrorCode StreamId ReasonPhrase | ConnectionErrorIsSent ErrorCode StreamId ReasonPhrase | StreamErrorIsReceived ErrorCode StreamId+ | StreamResetIsReceived ErrorCode StreamId | StreamErrorIsSent ErrorCode StreamId ReasonPhrase | BadThingHappen E.SomeException deriving (Show)@@ -270,10 +265,33 @@ , confPeerSockAddr :: SockAddr -- ^ This is copied into 'Aux', if exist, on server. , confReadNTimeout :: Bool+ , confOnInformational :: StreamId -> TokenHeaderTable -> IO ()+ -- ^ Client only: called when a 1xx informational response (e.g. 103 Early+ -- Hints) is received on the given stream, ahead of the final response.+ -- No-op by default.+ --+ -- @since 5.4.2 } -isAsyncException :: Exception e => e -> Bool+-- | Default config. This is just a template to modify via+-- field names. Don't use this without modifications.+defaultConfig :: Config+defaultConfig =+ Config+ { confWriteBuffer = nullPtr+ , confBufferSize = 0+ , confSendAll = \_ -> return ()+ , confReadN = \_ -> return ""+ , confPositionReadMaker = defaultPositionReadMaker+ , confTimeoutManager = T.defaultManager+ , confMySockAddr = SockAddrInet 0 0+ , confPeerSockAddr = SockAddrInet 0 0+ , confReadNTimeout = False+ , confOnInformational = \_ _ -> return ()+ }++isAsyncException :: E.Exception e => e -> Bool isAsyncException e = case E.fromException (E.toException e) of- Just (SomeAsyncException _) -> True+ Just (E.SomeAsyncException _) -> True Nothing -> False
Network/HTTP2/H2/Window.hs view
@@ -1,5 +1,6 @@ {-# LANGUAGE BangPatterns #-} {-# LANGUAGE NamedFieldPuns #-}+{-# LANGUAGE OverloadedStrings #-} module Network.HTTP2.H2.Window where @@ -29,6 +30,9 @@ w <- txWindowSize <$> readTVar streamTxFlow check (w > 0) +connectionWindowSizeSTM :: Context -> STM WindowSize+connectionWindowSizeSTM Context{txFlow} = txWindowSize <$> readTVar txFlow+ waitConnectionWindowSize :: Context -> STM () waitConnectionWindowSize Context{txFlow} = do w <- txWindowSize <$> readTVar txFlow@@ -67,39 +71,92 @@ ---------------------------------------------------------------- -- Sending window update +-- | Whether a stream's DATA is given back to the connection window as it+-- arrives, rather than as it is read.+--+-- So it is for pushes, the only streams of the server's we receive on. A+-- push is read only if a request for it comes along, and maybe never: held+-- against the connection window until then, pushes that nobody asked for+-- used it up, and the connection stalled, responses to requests and all.+-- Their stream windows still hold back what each can send unread.+connectionCreditedOnArrival :: StreamId -> Bool+connectionCreditedOnArrival = isServerInitiated+ informWindowUpdate :: Context -> Stream -> Int -> IO () informWindowUpdate _ _ 0 = return ()-informWindowUpdate Context{controlQ, rxFlow} Stream{streamNumber, streamRxFlow} len = do- mxc <- atomicModifyIORef rxFlow $ maybeOpenRxWindow len FCTWindowUpdate- forM_ mxc $ \ws -> do- let frame = windowUpdateFrame 0 ws- cframe = CFrames Nothing [frame]- enqueueControl controlQ cframe+informWindowUpdate ctx@Context{controlQ} Stream{streamNumber, streamRxFlow} len = do+ unless (connectionCreditedOnArrival streamNumber) $+ giveBackConnectionWindow ctx len mxs <- atomicModifyIORef streamRxFlow $ maybeOpenRxWindow len FCTWindowUpdate forM_ mxs $ \ws -> do let frame = windowUpdateFrame streamNumber ws cframe = CFrames Nothing [frame] enqueueControl controlQ cframe +-- | Account for a DATA frame that is being dropped.+--+-- Its stream is gone -- reset, or closed and forgotten -- so there is no+-- stream window to adjust. The peer charged these octets against the+-- connection window before sending them, though, and if we say nothing its+-- view of that window shrinks for good; enough dropped frames and the+-- connection stalls with both sides believing the other is at fault. So+-- charge them and give them straight back.+informIgnoredData :: Context -> StreamId -> Int -> IO ()+informIgnoredData _ _ 0 = return ()+informIgnoredData ctx@Context{rxFlow} sid len = do+ ok <- atomicModifyIORef' rxFlow $ checkRxLimit len+ unless ok $+ E.throwIO $+ ConnectionErrorIsSent+ EnhanceYourCalm+ sid+ "exceeds connection flow-control limit"+ giveBackConnectionWindow ctx len++-- | Give octets already charged to the connection window back to it, and to+-- it alone: for a stream that is closed, whose own window no longer+-- matters.+giveBackConnectionWindow :: Context -> Int -> IO ()+giveBackConnectionWindow _ 0 = return ()+giveBackConnectionWindow Context{controlQ, rxFlow} len = do+ mxc <- atomicModifyIORef rxFlow $ maybeOpenRxWindow len FCTWindowUpdate+ forM_ mxc $ \ws ->+ enqueueControl controlQ $ CFrames Nothing [windowUpdateFrame 0 ws]+ -- This must be called after an application is finished -- to adjust RX window. adjustRxWindow :: Context -> Stream -> IO ()-adjustRxWindow ctx stream@Stream{streamRxQ} = do+adjustRxWindow ctx stream = do+ len <- takeUnread stream+ informWindowUpdate ctx stream len++-- | Like 'adjustRxWindow', for a stream that has been closed: what was+-- left unread goes back to the connection window only.+--+-- Closed first, so that nothing is queued after this has looked: the+-- receiver does not queue DATA for a closed stream ('stream'), but gives+-- it back to the connection window itself.+giveBackUnread :: Context -> Stream -> IO ()+giveBackUnread ctx stream = do+ len <- takeUnread stream+ unless (connectionCreditedOnArrival $ streamNumber stream) $+ giveBackConnectionWindow ctx len++-- | Take what is left unread in a stream's queue, and say how many octets+-- of body it was.+takeUnread :: Stream -> IO Int+takeUnread Stream{streamRxQ} = do mq <- readIORef streamRxQ case mq of- Nothing -> return ()- Just q -> do- len <- readQ q- informWindowUpdate ctx stream len+ Nothing -> return 0+ Just q -> atomically $ loop q 0 where- readQ q = atomically $ loop 0- where- loop !total = do- meb <- tryReadTQueue q- case meb of- Just (Right (bs, _)) -> loop (total + BS.length bs)- Just le@(Left _) -> do- -- reserving HTTP2Error- writeTQueue q le- return total- _ -> return total+ loop q !total = do+ meb <- tryReadTQueue q+ case meb of+ Just (Right (bs, _)) -> loop q (total + BS.length bs)+ Just le@(Left _) -> do+ -- reserving HTTP2Error+ writeTQueue q le+ return total+ _ -> return total
Network/HTTP2/Server.hs view
@@ -51,7 +51,18 @@ rstRateLimit, -- * Common configuration- Config (..),+ Config,+ defaultConfig,+ confWriteBuffer,+ confBufferSize,+ confSendAll,+ confReadN,+ confPositionReadMaker,+ confTimeoutManager,+ confMySockAddr,+ confPeerSockAddr,+ confReadNTimeout,+ confOnInformational, allocSimpleConfig, allocSimpleConfig', freeSimpleConfig,
Network/HTTP2/Server/Internal.hs view
@@ -1,6 +1,8 @@ module Network.HTTP2.Server.Internal ( Request (..), Response (..),+ Config (..),+ ServerConfig (..), Aux (..), -- * Low level
Network/HTTP2/Server/Run.hs view
@@ -3,9 +3,8 @@ module Network.HTTP2.Server.Run where -import Control.Concurrent.Async (concurrently_)+import Control.Concurrent.Async import Control.Concurrent.STM-import qualified Control.Exception as E import Imports import Network.Control (defaultMaxData) import Network.HTTP.Semantics.IO@@ -127,11 +126,22 @@ let mgr = threadManager ctx runReceiver = frameReceiver ctx conf runSender = frameSender ctx conf- runBackgroundThreads = do- er <- E.try $ concurrently_ runReceiver runSender- case er of- Right () -> return ()- Left e -> closureServer conf ctx e+ runBackgroundThreads =+ withAsync runReceiver $ \ar ->+ withAsync runSender $ \as -> do+ r <- waitEither ar as+ e <- case r of+ -- The receiver is done; the sender finishes once it+ -- has flushed what is queued.+ Left _ -> wait as+ -- The sender finished first. Either the receiver is+ -- done too and not yet seen to be, and this is its+ -- error, or the sender failed: nothing more would go+ -- out, and leaving the receiver to run on left the+ -- connection open and silent, with no GOAWAY. Both+ -- are closed with it.+ Right e -> return e+ closureServer conf ctx e T.stopAfter mgr runBackgroundThreads $ \res -> closeAllStreams (oddStreamTable ctx) (evenStreamTable ctx) res
Network/HTTP2/Server/Worker.hs view
@@ -1,3 +1,4 @@+{-# LANGUAGE CPP #-} {-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE RecordWildCards #-} @@ -6,6 +7,7 @@ ) where import Control.Concurrent.STM+import qualified Control.Exception as E import Data.IORef import Network.HTTP.Semantics import Network.HTTP.Semantics.IO@@ -17,7 +19,12 @@ import Imports hiding (insert) import Network.HTTP2.Frame import Network.HTTP2.H2+import Network.HTTP2.H2.OutBodyIface +#if MIN_VERSION_http_semantics(0,4,1)+import qualified Data.ByteString.Char8 as C8+#endif+ ---------------------------------------------------------------- runServer :: Config -> Server -> Launch@@ -29,12 +36,14 @@ { auxTimeHandle = th , auxMySockAddr = mySockAddr , auxPeerSockAddr = peerSockAddr+#if MIN_VERSION_http_semantics(0,4,1)+ , auxSendInformational = sendInformational ctx strm+#endif } request = Request req' lc <- newLoopCheck strm Nothing server request aux $ sendResponse conf ctx lc strm request adjustRxWindow ctx strm- modifyPeerLastStreamId ctx $ streamNumber strm where label = "H2 response sender for stream " ++ show (streamNumber strm) pauseRequestBody th = req{inpObjBody = readBody'}@@ -48,6 +57,21 @@ ---------------------------------------------------------------- +#if MIN_VERSION_http_semantics(0,4,1)+-- | Send an informational (1xx) response, e.g. 103 Early Hints, on the given+-- stream ahead of the final response. This is wired into 'auxSendInformational'+-- so that a server (or WAI handler via Warp) can emit early hints. It blocks+-- until the informational HEADERS have been handed to the sender, preserving+-- ordering with respect to the final response.+sendInformational :: Context -> Stream -> Status -> ResponseHeaders -> IO ()+sendInformational ctx strm st hdrs = do+ lc <- newLoopCheck strm Nothing+ let hdr = (":status", C8.pack (show (statusCode st))) : hdrs+ syncWithSender ctx strm (OInformational hdr) lc+#endif++----------------------------------------------------------------+ -- | This function is passed to workers. -- They also pass 'Response's from a server to this function. -- This function enqueues commands for the HTTP/2 sender.@@ -101,37 +125,61 @@ push _ [] n = return (n :: Int) push tvar (pp : pps) n = do T.forkManaged threadManager "H2 server push" $ do- (pid, newstrm) <- makePushStream ctx pstrm- let scheme = fromJust $ getFieldValue tokenScheme reqvt- -- fixme: this value can be Nothing- auth =- fromJust- ( getFieldValue tokenAuthority reqvt- <|> getFieldValue tokenHost reqvt- )- path = promiseRequestPath pp- promiseRequest =- [ (tokenMethod, methodGet)- , (tokenScheme, scheme)- , (tokenAuthority, auth)- , (tokenPath, path)- ]- ot = OPush promiseRequest pid- Response rsp = promiseResponse pp- increment tvar- lc <- newLoopCheck newstrm Nothing- syncWithSender ctx newstrm ot lc- sendHeaderBody conf ctx lc newstrm rsp+ mpushed <- promise pp `E.finally` increment tvar+ forM_ mpushed $ \(newstrm, lc) -> do+ let Response rsp = promiseResponse pp+ sendHeaderBody conf ctx lc newstrm rsp push tvar pps (n + 1)+ -- Sending the PUSH_PROMISE, and only then counting the push as done:+ -- 'waiter' holds the parent's response back until every push is+ -- counted. The PUSH_PROMISE has to go out before the parent's frames+ -- (RFC 9113, section 8.4.1) -- before its END_STREAM above all, after+ -- which a PUSH_PROMISE on it is a connection error. Counted before it+ -- was queued, the parent's response could overtake it, and a client+ -- asked for the pushed resource itself before hearing of the promise.+ -- 'syncWithSender' returns once the sender has written the frame.+ -- Counted however it ends, or the parent would wait for ever. A push+ -- the peer has no room for is not made ('openEvenStreamTry').+ promise pp = do+ mstrm <- makePushStream ctx pstrm+ forM mstrm $ \(pid, newstrm) -> promiseOn pp pid newstrm+ promiseOn pp pid newstrm = do+ let scheme = fromJust $ getFieldValue tokenScheme reqvt+ -- fixme: this value can be Nothing+ auth =+ fromJust+ ( getFieldValue tokenAuthority reqvt+ <|> getFieldValue tokenHost reqvt+ )+ path = promiseRequestPath pp+ promiseRequest =+ [ (tokenMethod, methodGet)+ , (tokenScheme, scheme)+ , (tokenAuthority, auth)+ , (tokenPath, path)+ ]+ ot = OPush promiseRequest pid+ lc <- newLoopCheck newstrm Nothing+ syncWithSender ctx newstrm ot lc+ -- Reserved (local) until now. The peer sends nothing on a pushed+ -- stream, so its side is closed from here (RFC 9113, section 5.1:+ -- "half-closed (remote)" once the HEADERS go out), and the END_STREAM+ -- of the pushed response closes the stream. Left reserved, that+ -- END_STREAM only half-closed it: the stream stayed in the table+ -- holding a slot of the peer's SETTINGS_MAX_CONCURRENT_STREAMS, and+ -- once that many pushes had been made, the next waited for a slot+ -- for ever, and so did the response it belonged to.+ halfClosedRemote ctx newstrm+ return (newstrm, lc) ---------------------------------------------------------------- -makePushStream :: Context -> Stream -> IO (StreamId, Stream)+makePushStream :: Context -> Stream -> IO (Maybe (StreamId, Stream)) makePushStream ctx pstrm = do -- FLOW CONTROL: SETTINGS_MAX_CONCURRENT_STREAMS: send: respecting peer's limit- (_, newstrm) <- openEvenStreamWait ctx+ mstrm <- openEvenStreamTry ctx let pid = streamNumber pstrm- return (pid, newstrm)+ return $ (\(_, newstrm) -> (pid, newstrm)) <$> mstrm ---------------------------------------------------------------- @@ -170,10 +218,10 @@ -> Stream -> (OutBodyIface -> IO ()) -> IO (TBQueue StreamingChunk)-sendStreaming Context{..} strm strmbdy = do+sendStreaming ctx@Context{..} strm strmbdy = do tbq <- newTBQueueIO 10 -- fixme: hard coding: 10 T.forkManagedTimeout threadManager label $ \th ->- withOutBodyIface tbq id $ \iface -> do+ withOutBodyIface ctx strm tbq id $ \iface -> do let iface' = iface { outBodyPush = \b -> do
bench-hpack/Main.hs view
@@ -2,7 +2,7 @@ module Main where -import Control.Exception+import qualified Control.Exception as E import Criterion.Main import Data.ByteString (ByteString) import Network.HPACK
http2.cabal view
@@ -1,6 +1,6 @@-cabal-version: >=1.10+cabal-version: 2.0 name: http2-version: 5.3.11+version: 5.4.7 license: BSD3 license-file: LICENSE maintainer: Kazu Yamamoto <kazu@iij.ad.jp>@@ -8,7 +8,7 @@ homepage: https://github.com/kazu-yamamoto/http2 synopsis: HTTP/2 library description:- HTTP/2 library including frames, priority queues, HPACK, client and server.+ HTTP/2 library including frames, HPACK, client and server. category: Network build-type: Simple@@ -89,6 +89,7 @@ Network.HTTP2.H2.Context Network.HTTP2.H2.EncodeFrame Network.HTTP2.H2.HPACK+ Network.HTTP2.H2.OutBodyIface Network.HTTP2.H2.Queue Network.HTTP2.H2.Receiver Network.HTTP2.H2.Sender@@ -114,15 +115,15 @@ bytestring >=0.10, case-insensitive >=1.2 && <1.3, containers >=0.6,- http-semantics >= 0.3.1 && <0.4,+ http-semantics >= 0.4 && <0.5, http-types >=0.12 && <0.13, iproute >= 1.7 && < 1.8, network >=3.1, network-byte-order >=0.1.7 && <0.2, network-control >=0.1 && <0.2, stm >=2.5 && <2.6,- time-manager >=0.2 && <0.4,- unix-time >=0.4.11 && <0.5,+ time-manager >=0.3.0 && <0.4,+ unix-time >=0.4.11 && <0.6, utf8-string >=1.0 && <1.1 executable h2c-client@@ -139,7 +140,7 @@ http-types, http2, network,- network-run >= 0.5 && <0.6,+ network-run >= 0.6 && <0.7, unix-time if flag(devel)@@ -308,7 +309,8 @@ http-types, http2, network,- network-run >= 0.5 && <0.6,+ network-byte-order,+ network-run >= 0.6 && <0.7, random, typed-process @@ -327,7 +329,7 @@ hspec >=1.3, http-types, http2,- network-run >= 0.5 && <0.6,+ network-run >= 0.6 && <0.7, typed-process if flag(h2spec)
test-hpack/HPACKDecode.hs view
@@ -12,7 +12,7 @@ #if __GLASGOW_HASKELL__ < 709 import Control.Applicative ((<$>)) #endif-import Control.Exception+import qualified Control.Exception as E import Control.Monad (when) import qualified Data.ByteString.Base16 as B16 import qualified Data.ByteString.Char8 as B8@@ -66,7 +66,7 @@ case size c of Nothing -> return () Just siz -> renewDynamicTable siz dyntbl- x <- try $ decodeHeader dyntbl inp+ x <- E.try $ decodeHeader dyntbl inp case x of Left e -> return $ Just $ show (e :: DecodeError) Right hs' -> do
test/HPACK/DecodeSpec.hs view
@@ -2,8 +2,14 @@ module HPACK.DecodeSpec where +import Control.Monad (forM_)+import qualified Data.ByteString as BS+import qualified Data.ByteString.Char8 as BS8+import Data.String (fromString)+import Data.Word (Word8) import Network.HPACK import Network.HPACK.Table+import Network.HPACK.Token (tokenKey) import Test.Hspec import HPACK.HeaderBlock@@ -39,6 +45,117 @@ h1 `shouldBe` hl1 isDynamicTableEmpty etbl `shouldReturn` True isDynamicTableEmpty dtbl `shouldReturn` True+ it "keeps the newest entry when a full table evicts" $+ -- A size update to 40 leaves room for one entry. Two literals+ -- with incremental indexing, then index 62: the newest entry,+ -- "b". Inserting before evicting used to write "b" over "a",+ -- then evict the slot it had just written, so 62 came back as+ -- a dummy entry.+ withDynamicTableForDecoding 4096 4096 $ \dtbl -> do+ let blk =+ BS.pack+ [ 0x3f+ , 0x09 -- size update: 31 + 9+ , 0x40+ , 0x01+ , 0x61+ , 0x00 -- a: (incremental)+ , 0x40+ , 0x01+ , 0x62+ , 0x00 -- b: (incremental)+ , 0xbe -- indexed 62+ ]+ decodeHeader dtbl blk `shouldReturn` [("a", ""), ("b", ""), ("b", "")]+ it "decodes a Huffman-coded value longer than the Huffman buffer" $+ -- The value decodes to more than the 4096 octets of the+ -- decoder's Huffman buffer. It used to be reported as a+ -- truncated block, although the same value as a plain literal+ -- was accepted.+ withDynamicTableForEncoding 4096 $ \etbl ->+ withDynamicTableForDecoding 4096 4096 $ \dtbl ->+ forM_ [False, True] $ \huff -> do+ let hs = [("x-long", BS8.replicate 5000 'a')]+ stgy = defaultEncodeStrategy{useHuffman = huff}+ blk <- encodeHeader stgy 8192 etbl hs+ decodeHeader dtbl blk `shouldReturn` hs+ (tvs, _) <- decodeTokenHeader dtbl blk+ map (\(t, v) -> (tokenKey t, v)) tvs `shouldBe` hs+ it "decodes a block with no fields" $+ -- Empty, or only dynamic table size updates: both are valid+ -- blocks of no fields, and both used to be taken for truncated.+ withDynamicTableForDecoding 4096 4096 $ \dtbl ->+ forM_ ["", "\x20", "\x3f\xe1\x1f"] $ \blk -> do+ decodeHeader dtbl blk `shouldReturn` []+ (tvs, _) <- decodeTokenHeader dtbl blk+ tvs `shouldBe` []+ it "decodes the rest of a block with a malformed field" $+ -- The field after the malformed ones goes into the dynamic+ -- table, and the next block refers to it: index 62, the newest+ -- entry. The decoder used to stop at the malformed field, so+ -- that reference went astray.+ forM_ [illegalName, tooMany] $ \(fields, err) ->+ withDynamicTableForDecoding 4096 4096 $ \dtbl -> do+ let blk1 = fields <> incremental "x-after" "2"+ blk2 = BS.pack [0xbe]+ decodeTokenHeader dtbl blk1 `shouldThrow` (== err)+ (tvs, _) <- decodeTokenHeader dtbl blk2+ map (\(t, v) -> (tokenKey t, v)) tvs `shouldBe` [("x-after", "2")]+ it "round-trips through tables small enough to fill up" $+ -- Entries near the 32-octet minimum fill a table of these sizes+ -- to its last slot. The encoder follows the peer's+ -- SETTINGS_HEADER_TABLE_SIZE, so any of them can be asked for;+ -- the encoder used to send index 61 of the static table+ -- (www-authenticate) for an entry it had lost.+ forM_ [33, 40, 63, 64, 100, 127, 1023] $ \siz ->+ forM_ [False, True] $ \huff ->+ withDynamicTableForEncoding siz $ \etbl ->+ withDynamicTableForDecoding siz 4096 $ \dtbl ->+ forM_ smallBlocks $ \hs -> do+ let stgy = defaultEncodeStrategy{useHuffman = huff}+ blk <- encodeHeader stgy 4096 etbl hs+ decodeHeader dtbl blk `shouldReturn` hs++-- | A field name the encoder would have made lower-case.+illegalName :: (BS.ByteString, DecodeError)+illegalName = (literal "X-Upper" "1", IllegalHeaderName)++-- | One field more than the decoder takes.+tooMany :: (BS.ByteString, DecodeError)+tooMany =+ ( mconcat [literal (BS8.pack ('f' : show i)) "v" | i <- [1 .. 202 :: Int]]+ , TooLargeHeader+ )++-- | A literal field with a new name, without indexing (RFC 7541, 6.2.2).+literal :: BS.ByteString -> BS.ByteString -> BS.ByteString+literal = field 0x00++-- | A literal field with a new name, with incremental indexing (6.2.1).+incremental :: BS.ByteString -> BS.ByteString -> BS.ByteString+incremental = field 0x40++-- | Names and values shorter than 127 octets.+field :: Word8 -> BS.ByteString -> BS.ByteString -> BS.ByteString+field w k v =+ BS.pack [w, fromIntegral (BS.length k)]+ <> k+ <> BS.pack [fromIntegral (BS.length v)]+ <> v++-- | Blocks of fields close to the 32-octet minimum entry size, coming back+-- to earlier ones so that the encoder refers to what it inserted.+smallBlocks :: [[Header]]+smallBlocks =+ concat $+ replicate 3 $+ [ [("aa", "x")]+ , [("bb", "y")]+ , [("aa", "x")]+ , [("cc", ""), ("aa", "x")]+ , [("dd", "z"), ("bb", "y"), ("cc", "")]+ ]+ ++ [[(fromString ('k' : show i), "v")] | i <- [0 .. 40 :: Int]] hl1 :: [Header] hl1 =
test/HPACK/EncodeSpec.hs view
@@ -8,8 +8,12 @@ import qualified Control.Exception as E import Data.Bits import qualified Data.ByteString as BS+import qualified Data.ByteString.Char8 as C8 import Data.Maybe (fromMaybe)+import GHC.ForeignPtr (mallocPlainForeignPtrBytes)+import Network.ByteOrder (withReadBuffer, withWriteBuffer) import Network.HPACK+import Network.HPACK.Internal (decodeH, decodeS, encodeS) import Test.Hspec spec :: Spec@@ -32,6 +36,17 @@ run (Just 0) EncodeStrategy{compressionAlgo = Linear, useHuffman = False} [] it "does not use indexed fields" $ do runNotIndexed EncodeStrategy{compressionAlgo = Linear, useHuffman = False}+ describe "encodeS" $ do+ it "round-trips a Huffman-coded string whose length needs four octets" $ do+ -- 'a' is a five-bit code, so these come to 5/8 of their length+ -- when Huffman-coded: 26416 is the first to need a four-octet+ -- length with a 7-bit prefix. The fourth octet used to overwrite+ -- the first octet of the code.+ sequence_+ [ roundTripS n len+ | n <- [3, 5, 7]+ , len <- [100, 20000, 26415, 26416, 30000, 100000]+ ] run :: Maybe Int -> EncodeStrategy -> [Int] -> Expectation run msz stgy lens0 = do@@ -101,3 +116,16 @@ linearLens :: [Int] linearLens = [250,312,26,390,288,204,224,204,200,202,204,204,206,206,228,100,204,204,218,208,228,434,208,608,232,208,208,208,98,202,208,256,168,208,208,224,208,208,382,84,242,208,208,232,208,208,208,210,210,210,210,208,210,222,208,210,400,224,238,206,206,230,252,222,202,202,198,138,250,204,216,204,204,108,96,306,250,242,208,94,226,206,264,222,40,224,810,204,38,266,144,158,254,100,206,110,132,38,254,144,102,132,102,102,102,102,102,210,230,208,204,464,224,142,198,198,410,156,250,218,130,18,26,338,284,238,222,36,142,208,92,34,552,152,206,1020,288,42,490,98,40,1884,434,300,240,206,278,278,268,252,460,632,178,220,298,144,430,746,724,202,330,144,204,206,782,146,206,206,146,240,228,204,206,208,300,144,160,146,146,38,280,220,144,146,100,144,418,206,204,294,144,300,228,204,204,146,144,240,204,244,218,230,286,102,256,202,208,206,144,146,206,836,204,842,300,220,326,182,300,148,150,204,144,144,98,146,204,206,146,100,204,222,202,202,166,268,146,40,38,142,38,206,418,318,226,174,256,246,274,208,208,208,208,208,544,254,146,146,144,268,160,572,362,178,224,590,362,3150,1034,316,402,204,228,206,206,40,146,142,266,158,142,354,380,264,702,74,424,674,410,688,322,250,300,204,188,60,298,204,206,468,230,200,232,222,208,210,272,282,252,218,724,144,238,206,208,210,100,254,146,144,124,38,112,204,204,216,168,208,276,100,206,116,100,326,892,194,102,210,102,210,206,40,126,102,100,208,98,242,206,218,278,282,292,234,144,40,144,202,288,206,98,40,146,148,40,116,850,242,38,40,40,148,204,110,290,162,662,212,218,230,100,100,134,100,1026,100,2442,100,100,100,208,100,100,112,100,164,144,100,100,100,110,100,518,202,232,342,728,46,384,204,230,100,398,100,208,114,102,290,208,246,324,782,296,280,796,636,268,84,74,246,34,38,284,612,1090,332,602,378,84,24,256,204,234,26,226,654,60,206,28,160,220,238,38,204,484,206,440,308,206,246,392,314,814,714,200,244,290,258,50,94,252,572,38,284,1050,286,24,252,24,728,46,400,390,330,214,740,368,244,38,252,32,244,252,246,36,94,22,638,296,206,304,32,34,246,240,20,306,340,28,276,226,814,638,278,40,226,50,38,34,42,630,552,252,84,244,252,240,20,198,346,284,290,202,240,300,206,102,214,204,210,430,210,208,144,252,210,240,208,304,224,208,100,354,102,210,764,102,240,210,208,208,102,208,102,208,208,100,102,208,210,100,100,154,268,222,286,256,260,92,642,232,208,262,204,146,100,260,226,146,72,206,38,98,394,1090,348,2602,112,102,490,526,312,486,366,368,368,368,368,674,46,462,202,220,210,516,906,154,384,300,280,206,102,102,102,102,102,102,102,626,102,160,88,226,50,248,34,36,632,308,1124,684,450,254,252,714,60] -}++roundTripS :: Int -> Int -> Expectation+roundTripS n len = do+ let bs = C8.replicate len 'a'+ bufsiz = len * 4 + 64+ enc <- withWriteBuffer bufsiz $ \wbuf -> encodeS wbuf True id (`setBit` n) n bs+ gcbuf <- mallocPlainForeignPtrBytes bufsiz+ dec <-+ withReadBuffer enc $+ decodeS (.&. mask) (`testBit` n) n (decodeH gcbuf bufsiz)+ dec `shouldBe` bs+ where+ mask = (1 `shiftL` n) - 1
test/HPACK/HuffmanSpec.hs view
@@ -60,6 +60,10 @@ es <- encodeHuffman bs ds <- decodeHuffman es ds `shouldBe` bs+ it "decodes a string longer than 4096 octets" $ do+ let bs = BS.replicate 6000 'a' -- 3750 octets encoded+ es <- encodeHuffman bs+ decodeHuffman es `shouldReturn` bs describe "encode" $ do it "encodes" $ do mapM_ (\(x, y) -> x `shouldBeEncoded` y) testData
test/HPACK/IntegerSpec.hs view
@@ -2,6 +2,8 @@ import qualified Data.ByteString as BS import Data.Maybe (fromMaybe)+import Data.Word (Word8)+import Network.HPACK (DecodeError (..)) import Network.HPACK.Internal import Test.Hspec import Test.Hspec.QuickCheck@@ -14,8 +16,38 @@ x' <- decodeInteger n w ws x `shouldBe` x' +roundtrip7 :: BS.ByteString -> IO Int+roundtrip7 bs = do+ let (w, ws) = fromMaybe (error "roundtrip7") $ BS.uncons bs+ decodeInteger 7 w ws++-- | Decode with a 7-bit prefix that is all ones, so that the continuation+-- octets in 'ws' are what decides the value.+decode7 :: [Word8] -> IO Int+decode7 ws = decodeInteger 7 127 (BS.pack ws)+ spec :: Spec spec = do+ describe "decodeInteger" $ do+ it "rejects an encoding that runs past the limit" $ do+ r <- encodeInteger 7 integerLimit >>= roundtrip7+ r `shouldBe` integerLimit+ ws <- BS.unpack . BS.tail <$> encodeInteger 7 (integerLimit + 1)+ decode7 ws `shouldThrow` (== TooLargeInteger)++ it "rejects an encoding in more octets than the limit can take" $+ -- Continuation octets that each add nothing, so only their number+ -- is objectionable.+ decode7 (replicate 8 0x80 ++ [0x00]) `shouldThrow` (== TooLargeInteger)++ it "rejects an encoding that would wrap around" $+ -- This used to come back as 2, by overflowing 'Int' until it+ -- landed there: the same as the single octet 0x82, ":method: GET".+ -- Two byte strings decoding alike is exactly what RFC 7541+ -- section 5.1 asks a decoder to refuse.+ decode7 [0x83, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0x01]+ `shouldThrow` (== TooLargeInteger)+ describe "encode and decode" $ do prop "duality" $ dual 1 prop "duality" $ dual 2
test/HTTP2/ClientSpec.hs view
@@ -136,7 +136,13 @@ putMVar resultVar result threadDelay 10000 +-- | A malformed request is a stream error (RFC 9113 section 8.1.1), so the+-- server resets that stream and the connection carries on. The client learns+-- of it through the stream it was waiting on, as 'StreamResetIsReceived'.+--+-- This used to also admit 'ConnectionErrorIsReceived', from back when the+-- server escalated every stream error to the connection and answered one bad+-- request by hanging up on all of them. streamError :: Selector HTTP2Error-streamError StreamErrorIsReceived{} = True-streamError ConnectionErrorIsReceived{} = True+streamError StreamResetIsReceived{} = True streamError _ = False
test/HTTP2/FrameSpec.hs view
@@ -4,12 +4,67 @@ import Test.Hspec +import qualified Data.ByteString as BS import Data.ByteString.Char8 () import Data.Either import Network.HTTP2.Frame +-- | The error a decoder reports, or Nothing when it accepted the payload.+decodeError :: FrameType -> FrameHeader -> BS.ByteString -> Maybe ErrorCode+decodeError typ header body = case decodeFramePayload typ header body of+ Left (FrameDecodeError ec _ _) -> Just ec+ Right _ -> Nothing+ spec :: Spec spec = do+ describe "decodeFramePayload" $ do+ -- Each of these used to reach a peek at a fixed offset that never+ -- consulted the length of the ByteString it was reading from. An+ -- empty one is the shared empty ByteString, whose pointer is null, so+ -- the result was a segfault rather than an exception -- which is why+ -- none of this could be written as a failing assertion before.+ it "rejects a padded frame with no room for Pad Length" $ do+ let padded = FrameHeader 0 (setPadded defaultFlags) 1+ decodeError FrameData padded "" `shouldBe` Just FrameSizeError+ decodeError FramePushPromise padded "" `shouldBe` Just FrameSizeError++ it "rejects a padded HEADERS whose padding covers the priority fields" $ do+ -- Six octets is the smallest payload the header check accepts for+ -- PADDED and PRIORITY together, and a Pad Length of five leaves+ -- none of the five priority octets behind.+ let flags = setPadded $ setPriority defaultFlags+ header = FrameHeader 6 flags 1+ decodeError FrameHeaders header (BS.pack [5, 0, 0, 0, 0, 0])+ `shouldBe` Just FrameSizeError++ it "rejects a padded PUSH_PROMISE whose padding covers the promised id" $ do+ let flags = setPadded defaultFlags+ header = FrameHeader 5 flags 1+ decodeError FramePushPromise header (BS.pack [4, 0, 0, 0, 0])+ `shouldBe` Just FrameSizeError++ it "rejects a payload shorter than the frame header promised" $ do+ -- What a peer that hangs up mid-frame leaves behind.+ decodeError FramePriority (FrameHeader 5 defaultFlags 1) ""+ `shouldBe` Just FrameSizeError+ decodeError FrameRSTStream (FrameHeader 4 defaultFlags 1) ""+ `shouldBe` Just FrameSizeError+ decodeError FrameWindowUpdate (FrameHeader 4 defaultFlags 1) ""+ `shouldBe` Just FrameSizeError+ decodeError FrameSettings (FrameHeader 6 defaultFlags 0) ""+ `shouldBe` Just FrameSizeError++ it "rejects a payload too short for the fields it holds" $ do+ -- A payloadLength of zero satisfies checkFrameSize against an+ -- empty payload, but each of these still has a fixed-size field+ -- to read. GOAWAY was a segfault; the other three quietly+ -- returned whatever lay past the end of the buffer.+ let lying = FrameHeader 0 defaultFlags 1+ decodeError FrameRSTStream lying "" `shouldBe` Just FrameSizeError+ decodeError FrameWindowUpdate lying "" `shouldBe` Just FrameSizeError+ decodeError FramePriority lying "" `shouldBe` Just FrameSizeError+ decodeError FrameGoAway lying "" `shouldBe` Just FrameSizeError+ describe "encodeFrameHeader & decodeFrameHeader" $ do it "encode/decodes frames properly" $ do let header =
test/HTTP2/ServerSpec.hs view
@@ -1,423 +1,1700 @@ {-# LANGUAGE BangPatterns #-}-{-# LANGUAGE OverloadedStrings #-}-{-# LANGUAGE RankNTypes #-}-{-# LANGUAGE RecordWildCards #-}--module HTTP2.ServerSpec (spec) where--import Control.Concurrent-import Control.Concurrent.Async-import qualified Control.Exception as E-import Control.Monad-import Crypto.Hash (Context, SHA1)-import qualified Crypto.Hash as CH-import Data.ByteString (ByteString)-import qualified Data.ByteString as B-import Data.ByteString.Builder (Builder, byteString)-import qualified Data.ByteString.Char8 as C8-import Data.IORef-import Network.HTTP.Semantics-import Network.HTTP.Types-import Network.Run.TCP-import Network.Socket-import Network.Socket.ByteString-import System.IO-import System.IO.Unsafe-import System.Random-import Test.Hspec--import Network.HPACK-import Network.HPACK.Internal-import qualified Network.HTTP2.Client as C-import qualified Network.HTTP2.Client.Internal as C-import Network.HTTP2.Frame-import Network.HTTP2.Server--port :: String-port = show $ unsafePerformIO (randomPort <$> getStdGen)- where- randomPort = fst . randomR (43124 :: Int, 44320)--host :: String-host = "127.0.0.1"--spec :: Spec-spec = do- describe "server" $ do- it "handles normal cases" $- E.bracket (forkIO runServer) killThread $ \_ -> do- threadDelay 10000- runClient allocSimpleConfig-- it "should always send the connection preface first" $ do- prefaceVar <- newEmptyMVar- E.bracket (forkIO (runFakeServer prefaceVar)) killThread $ \_ -> do- threadDelay 10000- E.catch (runClient allocSlowPrefaceConfig) ignoreHTTP2Error-- preface <- takeMVar prefaceVar- preface `shouldBe` connectionPreface-- it "prevents attacks" $- E.bracket (forkIO runServer) killThread $ \_ -> do- threadDelay 10000- runAttack rapidSettings `shouldThrow` connectionError "too many settings"- runAttack rapidPing `shouldThrow` connectionError "too many ping"- runAttack rapidEmptyHeader- `shouldThrow` connectionError "too many empty headers"- runAttack rapidEmptyData `shouldThrow` connectionError "too many empty data"- runAttack rapidRst `shouldThrow` connectionError "too many rst_stream"--ignoreHTTP2Error :: C.HTTP2Error -> IO ()-ignoreHTTP2Error _ = pure ()--runServer :: IO ()-runServer = runTCPServer (Just host) port runHTTP2Server- where- runHTTP2Server s =- E.bracket- (allocSimpleConfig s 32768)- freeSimpleConfig- (\conf -> run defaultServerConfig conf server)--runFakeServer :: MVar ByteString -> IO ()-runFakeServer prefaceVar = do- runTCPServer (Just host) port $ \s -> do- ref <- newIORef Nothing-- -- send settings- sendAll s $- "\x00\x00\x12\x04\x00\x00\x00\x00\x00"- `mappend` "\x00\x03\x00\x00\x00\x80\x00\x04\x00"- `mappend` "\x01\x00\x00\x00\x05\x00\xff\xff\xff"-- -- receive preface- value <- defaultReadN s ref (B.length connectionPreface)- putMVar prefaceVar value-- -- send goaway frame- sendAll s "\x00\x00\x08\x07\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x01"-- -- wait for a few ms to make sure the client has a chance to close the- -- socket on its end- threadDelay 10000--server :: Server-server req _aux sendResponse = case requestMethod req of- Just "GET" -> case requestPath req of- Just "/" -> sendResponse responseHello []- Just "/stream" -> sendResponse responseInfinite []- Just "/push" -> do- let pp = pushPromise "/push-pp" responsePP 0- sendResponse responseHello [pp]- _ -> sendResponse response404 []- Just "POST" -> case requestPath req of- Just "/echo" -> sendResponse (responseEcho req) []- _ -> sendResponse responseHello []- _ -> sendResponse response405 []--responseHello :: Response-responseHello = responseBuilder ok200 header body- where- header = [("Content-Type", "text/plain")]- body = byteString "Hello, world!\n"--responsePP :: Response-responsePP = responseBuilder ok200 header body- where- header =- [ ("Content-Type", "text/plain")- , ("x-push", "True")- ]- body = byteString "Push\n"--responseInfinite :: Response-responseInfinite = responseStreaming ok200 header body- where- header = [("Content-Type", "text/plain")]- body :: (Builder -> IO ()) -> IO () -> IO ()- body write flush = do- let go n = write (byteString (C8.pack (show n)) `mappend` "\n") *> flush *> go (succ n)- go (0 :: Int)--response404 :: Response-response404 = responseNoBody notFound404 []--response405 :: Response-response405 = responseNoBody methodNotAllowed405 []--responseEcho :: Request -> Response-responseEcho req = setResponseTrailersMaker h2rsp maker- where- h2rsp = responseStreaming ok200 header streamingBody- header = [("Content-Type", "text/plain")]- mhx = getFieldValue (toToken "X-Tag") (snd (requestHeaders req))- streamingBody write _flush = do- loop- mt <- getRequestTrailers req- firstTrailerValue <$> mt `shouldBe` mhx- where- loop = do- bs <- getRequestBodyChunk req- when (bs /= "") $ do- void $ write $ byteString bs- loop- maker = trailersMaker (CH.hashInit :: Context SHA1)---- Strictness is important for Context.-trailersMaker :: Context SHA1 -> Maybe ByteString -> IO NextTrailersMaker-trailersMaker ctx Nothing = return $ Trailers [("X-SHA1", sha1)]- where- !sha1 = C8.pack $ show $ CH.hashFinalize ctx-trailersMaker ctx (Just bs) = return $ NextTrailersMaker $ trailersMaker ctx'- where- !ctx' = CH.hashUpdate ctx bs--runClient :: (Socket -> BufferSize -> IO Config) -> IO ()-runClient allocConfig =- runTCPClient host port runHTTP2Client- where- auth = host- cliconf = C.defaultClientConfig{C.authority = auth}- runHTTP2Client s =- E.bracket- (allocConfig s 4096)- freeSimpleConfig- (\conf -> C.run cliconf conf client)-- client :: C.Client ()- client sendRequest aux =- foldr1- concurrently_- [ client0 sendRequest aux- , client1 sendRequest aux- , client2 sendRequest aux- , client3 sendRequest aux- , client3' sendRequest aux- , client3'' sendRequest aux- , client4 sendRequest aux- , client5 sendRequest aux- ]---- delay sending preface to be able to test if it is always sent first-allocSlowPrefaceConfig :: Socket -> BufferSize -> IO Config-allocSlowPrefaceConfig s size = do- config <- allocSimpleConfig s size- pure config{confSendAll = slowPrefaceSend (confSendAll config)}- where- slowPrefaceSend :: (ByteString -> IO ()) -> ByteString -> IO ()- slowPrefaceSend orig chunk = do- when (C8.pack "PRI" `C8.isPrefixOf` chunk) $ do- threadDelay 10000- orig chunk--client0 :: C.Client ()-client0 sendRequest _aux = do- let req = C.requestNoBody methodGet "/" []- sendRequest req $ \rsp -> do- C.responseStatus rsp `shouldBe` Just ok200- fmap statusMessage (C.responseStatus rsp) `shouldBe` Just "OK"--client1 :: C.Client ()-client1 sendRequest _aux = do- let req = C.requestNoBody methodGet "/push-pp" []- sendRequest req $ \rsp -> do- C.responseStatus rsp `shouldBe` Just notFound404--client2 :: C.Client ()-client2 sendRequest _aux = do- let req = C.requestNoBody methodPut "/" []- sendRequest req $ \rsp -> do- C.responseStatus rsp `shouldBe` Just methodNotAllowed405--client3 :: C.Client ()-client3 sendRequest _aux = do- let hx = "b0870457df2b8cae06a88657a198d9b52f8e2b0a"- req0 =- C.requestFile methodPost "/echo" [("X-Tag", hx)] $- FileSpec "test/inputFile" 0 1012731- req = C.setRequestTrailersMaker req0 maker- sendRequest req $ \rsp -> do- let consumeBody = do- bs <- C.getResponseBodyChunk rsp- when (bs /= "") consumeBody- consumeBody- mt <- C.getResponseTrailers rsp- firstTrailerValue <$> mt `shouldBe` Just hx- where- !maker = trailersMaker (CH.hashInit :: Context SHA1)--client3' :: C.Client ()-client3' sendRequest _aux = do- let hx = "b0870457df2b8cae06a88657a198d9b52f8e2b0a"- req0 = C.requestStreaming methodPost "/echo" [("X-Tag", hx)] $ \write _flush -> do- let sendFile h = do- bs <- B.hGet h 1024- when (bs /= "") $ do- write $ byteString bs- sendFile h- withFile "test/inputFile" ReadMode sendFile- req = C.setRequestTrailersMaker req0 maker- sendRequest req $ \rsp -> do- let consumeBody = do- bs <- C.getResponseBodyChunk rsp- when (bs /= "") consumeBody- consumeBody- mt <- C.getResponseTrailers rsp- firstTrailerValue <$> mt `shouldBe` Just hx- where- !maker = trailersMaker (CH.hashInit :: Context SHA1)--client3'' :: C.Client ()-client3'' sendRequest _axu = do- let hx = "59f82dfddc0adf5bdf7494b8704f203a67e25d4a"- req0 = C.requestStreaming methodPost "/echo" [("X-Tag", hx)] $ \write _flush -> do- let chunk = C8.replicate (16384 * 2) 'c'- tag = C8.replicate 16 't'- -- I don't think 9 is important here, this is just what I have, the client hangs on receiving the last one- replicateM_ 9 $ write $ byteString chunk- write $ byteString tag- req = C.setRequestTrailersMaker req0 maker- sendRequest req $ \rsp -> do- let consumeBody = do- bs <- C.getResponseBodyChunk rsp- when (bs /= "") consumeBody- consumeBody- mt <- C.getResponseTrailers rsp- firstTrailerValue <$> mt `shouldBe` Just hx- where- !maker = trailersMaker (CH.hashInit :: Context SHA1)--client4 :: C.Client ()-client4 sendRequest _aux = do- let req0 = C.requestNoBody methodGet "/push" []- sendRequest req0 $ \rsp -> do- C.responseStatus rsp `shouldBe` Just ok200- let req1 = C.requestNoBody methodGet "/push-pp" []- sendRequest req1 $ \rsp -> do- C.responseStatus rsp `shouldBe` Just ok200--client5 :: C.Client ()-client5 sendRequest _aux = do- let req0 = C.requestNoBody methodGet "/stream" []- sendRequest req0 $ \rsp -> do- C.responseStatus rsp `shouldBe` Just ok200- let go n- | n > 0 = do- _ <- C.getResponseBodyChunk rsp- go (pred n)- | otherwise = pure ()- go (100 :: Int)--firstTrailerValue :: TokenHeaderTable -> FieldValue-firstTrailerValue tbl = case fst tbl of- [] -> error "firstTrailerValue"- x : _ -> snd x--runAttack :: (C.ClientIO -> IO ()) -> IO ()-runAttack attack =- runTCPClient host port runHTTP2Client- where- auth = host- cliconf = C.defaultClientConfig{C.authority = auth}- runHTTP2Client s =- E.bracket- (allocSimpleConfig s 4096)- freeSimpleConfig- (\conf -> C.runIO cliconf conf client)- client cconf = return $ do- attack cconf- threadDelay 1000000--rapidSettings :: C.ClientIO -> IO ()-rapidSettings C.ClientIO{..} = do- let einfo = EncodeInfo defaultFlags 0 Nothing- bs = encodeFrame einfo $ SettingsFrame [(SettingsEnablePush, 0)]- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs--rapidPing :: C.ClientIO -> IO ()-rapidPing C.ClientIO{..} = do- let einfo = EncodeInfo defaultFlags 0 Nothing- opaque64 = "01234567"- bs = encodeFrame einfo $ PingFrame opaque64- replicateM_ 20 $ cioWriteBytes bs--rapidEmptyHeader :: C.ClientIO -> IO ()-rapidEmptyHeader C.ClientIO{..} = do- (sid, _) <- cioCreateStream- let einfo = EncodeInfo defaultFlags sid Nothing- bs = encodeFrame einfo $ HeadersFrame Nothing ""- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs- cioWriteBytes bs--rapidEmptyData :: C.ClientIO -> IO ()-rapidEmptyData C.ClientIO{..} = do- (sid, _) <- cioCreateStream- let einfoH = EncodeInfo (setEndHeader defaultFlags) sid Nothing- hdr =- hpackEncode- [ (":scheme", "http")- , (":authority", "127.0.0.1")- , (":path", "/")- , (":method", "GET")- ]- bsH = encodeFrame einfoH $ HeadersFrame Nothing hdr- cioWriteBytes bsH- let einfoD = EncodeInfo defaultFlags sid Nothing- bsD = encodeFrame einfoD $ DataFrame ""- cioWriteBytes bsD- cioWriteBytes bsD- cioWriteBytes bsD- cioWriteBytes bsD- cioWriteBytes bsD- cioWriteBytes bsD- cioWriteBytes bsD- cioWriteBytes bsD--rapidRst :: C.ClientIO -> IO ()-rapidRst C.ClientIO{..} = do- reset- reset- reset- reset- reset- reset- reset- reset- where- reset = do- (sid, _) <- cioCreateStream- -- setEndStream for HalfClosedRemote- let einfoH = EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing- hdr =- hpackEncode- [ (":scheme", "http")- , (":authority", "127.0.0.1")- , (":path", "/")- , (":method", "GET")- ]- bsH = encodeFrame einfoH $ HeadersFrame Nothing hdr- cioWriteBytes bsH- let einfoR = EncodeInfo defaultFlags sid Nothing- -- Only (HalfClosedRemote, NoError) is accepted.- -- Otherwise, a stream error terminates the connection.- bsR = encodeFrame einfoR $ RSTStreamFrame NoError- cioWriteBytes bsR+{-# LANGUAGE NamedFieldPuns #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE RecordWildCards #-}+-- GHC 9.12 and later (still on master), with -O, compile 'responseInfinite'+-- to a response with no body: 'OutBodyNone' instead of 'OutBodyStreaming'.+-- Full laziness floats the constructor application to a top-level thunk,+-- and that thunk reaches code that switches on the pointer tag without+-- evaluating it: https://gitlab.haskell.org/ghc/ghc/-/work_items/27857+-- The "infinite" stream then ends with its HEADERS, and the MadeYouReset+-- test sometimes sees a stream closed before its PRIORITY arrives (#191).+-- 9.10 and earlier are not affected.+{-# OPTIONS_GHC -fno-full-laziness #-}++module HTTP2.ServerSpec (spec) where++import Control.Concurrent+import Control.Concurrent.Async+import qualified Control.Exception as E+import Control.Monad+import Crypto.Hash (Context, SHA1)+import qualified Crypto.Hash as CH+import Data.ByteString (ByteString)+import qualified Data.ByteString as B+import Data.ByteString.Builder (Builder, byteString)+import qualified Data.ByteString.Char8 as C8+import Data.IORef+import Data.Maybe (isJust, isNothing)+import Network.HTTP.Semantics+import Network.HTTP.Types+import Network.Run.TCP+import Network.Socket+import Network.Socket.ByteString+import System.IO+import System.IO.Unsafe+import System.Random+import System.Timeout (timeout)+import Test.Hspec++import Network.HPACK+import Network.HPACK.Internal+import qualified Network.HTTP2.Client as C+import qualified Network.HTTP2.Client.Internal as C+import Network.HTTP2.Frame+import Network.HTTP2.Server++port :: String+port = show $ unsafePerformIO (randomPort <$> getStdGen)+ where+ randomPort = fst . randomR (43124 :: Int, 44320)++host :: String+host = "127.0.0.1"++spec :: Spec+spec = do+ describe "server" $ do+ it "sends a header block and trailers larger than a frame" $+ -- Both have to go out as HEADERS and CONTINUATION frames and be+ -- put back together on receipt; the requests after them check+ -- that the two ends' HPACK tables still agree.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ r <- timeout 5000000 $ runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run C.defaultClientConfig{C.authority = host} conf $ \sendRequest _ -> do+ sendRequest (C.requestNoBody methodGet "/big" []) $ \rsp -> do+ getFieldValue (toToken "x-big") (snd (C.responseHeaders rsp))+ `shouldBe` Just bigVal+ let drain = do+ bs <- C.getResponseBodyChunk rsp+ unless (B.null bs) drain+ drain+ mt <- C.getResponseTrailers rsp+ (mt >>= getFieldValue (toToken "x-big-trailer") . snd)+ `shouldBe` Just bigVal+ -- Same connection: the HPACK state must still agree.+ forM_ [1 :: Int, 2] $ \_ ->+ sendRequest (C.requestNoBody methodGet "/" []) $ \rsp ->+ C.responseStatus rsp `shouldBe` Just ok200+ r `shouldBe` Just ()++ it "handles normal cases" $+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ runClient allocSimpleConfig++ it "delivers 103 Early Hints to the client's informational handler" $+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ hintsRef <- newIORef []+ runClientEarly hintsRef >>= (`shouldBe` Just ok200)+ hints <- readIORef hintsRef+ map (getFieldValue (toToken "link") . snd) hints+ `shouldBe` [ Just "</style.css>; rel=preload; as=style"+ , Just "</app.js>; rel=preload; as=script"+ ]++ it "should always send the connection preface first" $ do+ prefaceVar <- newEmptyMVar+ E.bracket (forkIO (runFakeServer prefaceVar)) killThread $ \_ -> do+ threadDelay 10000+ E.catch (runClient allocSlowPrefaceConfig) ignoreHTTP2Error++ preface <- takeMVar prefaceVar+ preface `shouldBe` connectionPreface++ it "refuses one stream over the limit and keeps the connection" $+ E.bracket (forkIO runServerMaxConc1) killThread $ \_ -> do+ threadDelay 10000+ -- The server announced room for one concurrent stream. Open+ -- one, reset it, then open two more: the second of those is+ -- the one over the limit.+ --+ -- Two things are on trial. That the reset gives the slot back+ -- exactly once -- decrementing the count twice, as it used to,+ -- would leave room for both. And that being over the limit+ -- costs you that stream and not the connection: no GOAWAY.+ frames <-+ rawExchange+ [ openStreamFrame 1+ , encodeFrame (EncodeInfo defaultFlags 1 Nothing) $+ RSTStreamFrame Cancel+ , openStreamFrame 3+ , openStreamFrame 5+ ]+ [(sid, ec) | (FrameRSTStream, sid, ec) <- resets frames]+ `shouldBe` [(5, RefusedStream)]+ [() | (FrameGoAway, _, _) <- resets frames] `shouldBe` []++ it "releases a worker whose stream the peer reset" $ do+ doneVar <- newEmptyMVar+ E.bracket (forkIO (runServerCancel doneVar)) killThread $ \_ -> do+ threadDelay 10000+ runAttack cancelInFlight+ timeout 1000000 (takeMVar doneVar) `shouldReturn` Just ()++ it "survives a padded HEADERS whose padding covers the priority fields" $+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ runAttack paddingOverPriority+ `shouldThrow` connectionError "no room for priority fields"++ it "resets one stream and goes on serving the connection" $+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ runStreamErrorClient++ it "limits the resets a peer can make us send (MadeYouReset)" $+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ -- Not through the client library: it would take the+ -- server's first RST_STREAM, on a stream it never opened+ -- itself, for a protocol error of its own.+ timeout 5000000 rapidStreamError+ `shouldReturn` Just (Just (EnhanceYourCalm, "too many stream errors"))++ it "gives back the slot of a stream both sides streamed on" $+ -- Room for four concurrent streams, so a few slots that are+ -- never given back stop the connection within a few thousand+ -- requests one after another. Both sides stream with flushes: the receiver used to write back a stream state+ -- it had read before the sender half-closed the stream, undoing+ -- the half-close, so the peer's END_STREAM then left the stream+ -- half-closed instead of closed and in the table for good.+ --+ -- It is a race between the receiver and the sender, so it needs+ -- them running in parallel: on one capability it hardly ever+ -- shows.+ withCapabilities 4 $+ E.bracket (forkIO runServerSmallWindow) killThread $ \_ -> do+ threadDelay 10000+ done <- newIORef (0 :: Int)+ r <- timeout 60000000 $ runTCPClient host port $ \s -> do+ -- Fifty small writes each way per request: without+ -- this, Nagle and delayed ACKs can hold each one up.+ setSocketOption s NoDelay 1+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run C.defaultClientConfig{C.authority = host} conf $ \sendRequest _ ->+ forM_ [1 .. 2000 :: Int] $ \_ -> do+ let req = C.requestStreaming methodPost "/both" [] $ \write flush ->+ replicateM_ 50 $ write (byteString (C8.replicate 50 'a')) >> flush+ sendRequest req $ \rsp -> do+ let drain n = do+ bs <- C.getResponseBodyChunk rsp+ if B.null bs then return n else drain (n + B.length bs)+ drain 0 `shouldReturn` 2500+ modifyIORef' done (+ 1)+ -- How far it got tells a hang from a slow run.+ n <- readIORef done+ when (isNothing r) $+ expectationFailure $+ "timed out after " ++ show n ++ " of 2000 requests"++ it "accepts a content-length on a response with no content" $+ -- RFC 9113, section 8.1.1: the response to HEAD, 204 and 304 can+ -- carry a non-zero content-length without content. The client+ -- used to take each of these for a malformed response.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run C.defaultClientConfig{C.authority = host} conf $ \sendRequest _ -> do+ let noContent method path =+ sendRequest (C.requestNoBody method path []) $ \rsp -> do+ C.responseStatus rsp `shouldSatisfy` isJust+ C.getResponseBodyChunk rsp `shouldReturn` ""+ noContent methodHead "/"+ noContent methodHead "/data"+ noContent methodGet "/not-modified"+ -- A response that is meant to have content still+ -- has to match its content-length.+ sendRequest (C.requestNoBody methodGet "/no-content" []) (const $ return ())+ `shouldThrow` malformedResponse++ it "does not open a stream for a PRIORITY frame" $+ -- Over a raw socket, as the client library does not send+ -- PRIORITY. The server allows 64 concurrent streams; each of+ -- these PRIORITY frames used to open one and hold its slot,+ -- so the request after them was refused.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 idlePriority `shouldReturn` Just (Just "HEADERS")++ it "closes the connection when SETTINGS overflow a stream's window" $+ -- RFC 9113, section 6.9.2: a connection error of type+ -- FLOW_CONTROL_ERROR. The overflow is found in the sender,+ -- which used to stop on it without a word, leaving the+ -- connection open and silent.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ -- A connection error: GOAWAY, with no RST_STREAM before it.+ timeout 5000000 settingsOverflow+ `shouldReturn` Just (False, Just FlowControlError)++ it "accepts empty trailers" $+ -- A HEADERS frame with END_STREAM and an empty field block ends+ -- the body with no trailer fields. The empty block used to be+ -- taken for a truncated one: COMPRESSION_ERROR, and the+ -- connection closed.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 emptyTrailers `shouldReturn` Just (Just "HEADERS")++ it "checks a padded body against its content-length" $+ -- Padding is not content (RFC 9113, section 6.1). It used to+ -- be counted into the body's length, so a padded body that+ -- matched its content-length was reset as one that did not.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 paddedBody `shouldReturn` Just (Just "DATA 4")++ it "gives the padding of a body back to the windows" $+ -- 2000 DATA frames of one octet of content and 255 of padding:+ -- about twice the stream's window. The padding used to be+ -- charged and never given back, so a peer keeping to the+ -- windows stalled, and one that did not, like this one, broke+ -- the stream's limit and had the connection closed.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 paddingWindow `shouldReturn` Just (Just "DATA 2000")++ it "gives DATA it refuses back to the connection window" $+ -- DATA on a stream the peer has half-closed is a stream error+ -- (RFC 9113, section 5.1), but it still counts against the+ -- connection window (section 6.9). It used to be left out, so+ -- the peer's view of that window shrank for good. With a+ -- window of 65535, the refused 16384 octets and the 16384 of+ -- the next request make up the half that is given back.+ E.bracket (forkIO runServerSmallConnWindow) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 refusedData `shouldReturn` Just (Just (32768, "16384"))++ it "goes on sending requests after one fails before it is queued" $+ -- The file of this requestFile does not exist, so the request+ -- fails after its stream id is taken and before it is queued.+ -- Requests are queued in stream id order, so every one after it+ -- used to wait for its turn for ever.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ r <- timeout 5000000 $ runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run C.defaultClientConfig{C.authority = host} conf $ \sendRequest _ -> do+ let missing =+ C.requestFile methodPost "/echo" [] $+ FileSpec "test/no-such-file" 0 10+ failed <- E.try $ sendRequest missing (const $ return ())+ either (const True) (const False) (failed :: Either E.SomeException ())+ `shouldBe` True+ replicateM_ 3 $+ sendRequest (C.requestNoBody methodGet "/" []) $ \rsp ->+ C.responseStatus rsp `shouldBe` Just ok200+ r `shouldBe` Just ()++ it "answers without a push when the peer has no room for one" $+ -- SETTINGS_MAX_CONCURRENT_STREAMS of 0 is how a peer can refuse+ -- pushes (RFC 9113, section 8.4). The push of /push-pp used to+ -- wait for room for ever, and the response to /push with it.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 pushNoRoom `shouldReturn` Just (Just "HEADERS")++ it "frees the stream of a response the client did not read to the end" $+ -- /endless never ends. Each request here reads one chunk of it+ -- and is done, by returning or by throwing. Its stream used to+ -- stay open, holding one of the server's 64 slots, with what the+ -- server sent never given back to the connection window: the+ -- 65th request waited for a slot for ever. Now each is reset,+ -- 70 in a burst, which takes a server allowing more resets a+ -- second than the default.+ E.bracket (forkIO runServerManyResets) killThread $ \_ -> do+ threadDelay 10000+ r <- timeout 10000000 $ runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run C.defaultClientConfig{C.authority = host} conf $ \sendRequest _ -> do+ forM_ [1 .. 70 :: Int] $ \i -> do+ let abandon rsp = do+ _ <- C.getResponseBodyChunk rsp+ when (even i) $ E.throwIO $ userError "done with it"+ r <- E.try $ sendRequest (C.requestNoBody methodGet "/endless" []) abandon+ either (\e -> const (return ()) (e :: E.IOException)) return r+ sendRequest (C.requestNoBody methodGet "/" []) $ \rsp ->+ C.responseStatus rsp `shouldBe` Just ok200+ r `shouldBe` Just ()++ it "goes on when pushes nobody asks for fill the connection window" $+ -- Each /push-big comes with a push of 20000 octets that is never+ -- asked for. With a connection window of 65535, the fourth push+ -- used to find it used up by the first three, unread, and the+ -- connection stalled, the responses to /push-big with it.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ let cconf =+ C.defaultClientConfig+ { C.authority = host+ , C.connectionWindowSize = defaultWindowSize+ }+ r <- timeout 5000000 $ runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run cconf conf $ \sendRequest _ ->+ replicateM_ 10 $+ sendRequest (C.requestNoBody methodGet "/push-big" []) $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just ok200+ let body = do+ bs <- C.getResponseBodyChunk rsp+ unless (B.null bs) body+ body+ r `shouldBe` Just ()++ it "counts streams whose handlers are running in its GOAWAY" $+ -- RFC 9113, section 6.8: the last stream identifier is the+ -- highest one that "might have been processed". Streams 1 and+ -- 3 are being answered when the connection is closed, and used+ -- to be left out until their handlers had returned, so the+ -- GOAWAY said 0: as if the client could send them again.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 goAwayLastStream `shouldReturn` Just (Just (3, ProtocolError))++ it "finishes what the server will still answer after its GOAWAY" $+ -- RFC 9113, section 6.8: a GOAWAY with NO_ERROR and last stream 1+ -- says stream 1 will still be answered, stream 3 will not. The+ -- client used to close the connection as soon as it came,+ -- failing both, and the client function with them.+ E.bracket (forkIO runGoAwayServer) killThread $ \_ -> do+ threadDelay 10000+ r <- timeout 5000000 $ E.try $ runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run C.defaultClientConfig{C.authority = host} conf $ \sendRequest _ -> do+ let get = E.try . flip sendRequest readAll . C.requestNoBody methodGet "/" $ []+ readAll rsp = do+ bs <- C.getResponseBodyChunk rsp+ if B.null bs then return "" else (bs <>) <$> readAll rsp+ (r1, r3) <- concurrently get (threadDelay 50000 >> get)+ -- By now the connection has run its course,+ -- and the client function goes on: nothing new+ -- goes out, and it is not killed either.+ threadDelay 100000+ r5 <- get+ return (status r1, status r3, status r5)+ case r of+ Just (Right rs) -> rs `shouldBe` ("hello", "closed", "closed")+ Just (Left e) -> expectationFailure $ show (e :: C.HTTP2Error)+ Nothing -> expectationFailure "timed out"++ it "answers the requests it has after the client's GOAWAY" $+ -- A client's GOAWAY speaks of the server's own streams, and+ -- the requests it has already sent are still to be answered.+ -- The server used to close the connection on it at once.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 goAwayFromClient+ `shouldReturn` Just ["HEADERS 1", "DATA 1 END_STREAM", "GOAWAY NoError"]++ it "resets a malformed request and goes on serving the connection" $+ -- An upper-case field name makes the request malformed: a+ -- stream error (RFC 9113, section 8.1.1). It used to close the+ -- connection, and the next request on it went unanswered.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 malformedRequest+ `shouldReturn` Just ["RST_STREAM 1 ProtocolError", "HEADERS 3"]++ it "answers without a body while the connection window is shut" $+ -- HEADERS are not flow-controlled (RFC 9113, section 6.9). The+ -- sender used to wait for the connection window before taking+ -- anything off its queue, so once /endless had used it up, the+ -- response to a request with no body never went out.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ timeout 5000000 shutWindow `shouldReturn` Just (Just "HEADERS 3, then DATA 1")++ it "sends a PUSH_PROMISE before the response that carries it" $+ -- /push answers with a push of /push-pp, so a request for+ -- /push-pp after it is served from the push. The server used to+ -- let the response to /push overtake the PUSH_PROMISE now and+ -- then; the client then asked the server for /push-pp itself,+ -- and got 404. One round in a few dozen did, so 200 of them --+ -- which also takes more pushes than the peer allows concurrent+ -- streams, so pushed streams that are never closed show too.+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ done <- newIORef (0 :: Int)+ r <- timeout 30000000 $ runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run C.defaultClientConfig{C.authority = host} conf $ \sendRequest _ ->+ replicateM_ 200 $ do+ -- Bodies are read to the end, so that the+ -- streams close and give their slots back.+ let drain rsp = do+ bs <- C.getResponseBodyChunk rsp+ unless (B.null bs) $ drain rsp+ sendRequest (C.requestNoBody methodGet "/push" []) $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just ok200+ drain rsp+ sendRequest (C.requestNoBody methodGet "/push-pp" []) $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just ok200+ drain rsp+ modifyIORef' done (+ 1)+ -- How far it got tells a hang (at 64, the peer's concurrency+ -- limit, if pushed streams leak) from a slow run.+ n <- readIORef done+ when (isNothing r) $+ expectationFailure $+ "timed out after " ++ show n ++ " of 200 rounds"++ it "uploads a file through runIO past the stream's window" $+ -- The server announces an 8192-octet window. runIO put the rest+ -- of a body back on the queue without waiting for the window to+ -- open; with none left, the file was read into no room, and a+ -- read of 0 octets is the end of the file, so the request ended+ -- with END_STREAM after the first window's worth.+ E.bracket (forkIO runServerSmallWindow) killThread $ \_ -> do+ threadDelay 10000+ timeout 10000000 uploadIO `shouldReturn` Just 100000++ it "prevents attacks" $+ E.bracket (forkIO runServer) killThread $ \_ -> do+ threadDelay 10000+ runAttack rapidSettings `shouldThrow` connectionError "too many settings"+ runAttack rapidPing `shouldThrow` connectionError "too many ping"+ runAttack rapidEmptyHeader+ `shouldThrow` connectionError "too many empty headers"+ runAttack rapidEmptyData `shouldThrow` connectionError "too many empty data"+ runAttack rapidRst `shouldThrow` connectionError "too many rst_stream"++ignoreHTTP2Error :: C.HTTP2Error -> IO ()+ignoreHTTP2Error _ = pure ()++runServer :: IO ()+runServer = runTCPServer (Just host) port runHTTP2Server+ where+ runHTTP2Server s =+ E.bracket+ (allocSimpleConfig s 32768)+ freeSimpleConfig+ (\conf -> run defaultServerConfig conf server)++-- | Like 'runServer', but announcing room for a single concurrent stream.+-- | Uploading 100000 octets of a file through 'C.runIO', and what the server+-- says it received.+uploadIO :: IO Int+uploadIO = runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.runIO C.defaultClientConfig{C.authority = host} conf $ \C.ClientIO{..} ->+ return $ do+ let body rsp acc = do+ bs <- C.getResponseBodyChunk rsp+ if B.null bs then return acc else body rsp (acc <> bs)+ exchange req = cioWriteRequest req >>= cioReadResponse . snd+ -- A request first, so that the server's SETTINGS -- and its+ -- small window -- are known before the upload starts.+ _ <- exchange (C.requestNoBody methodGet "/" []) >>= (`body` "")+ rsp <-+ exchange $+ C.requestFile methodPost "/count" [] $+ FileSpec "test/inputFile" 0 100000+ read . C8.unpack <$> body rsp ""++-- | Running with at least this many capabilities.+withCapabilities :: Int -> IO a -> IO a+withCapabilities n act =+ E.bracket getNumCapabilities setNumCapabilities $ \old -> do+ setNumCapabilities (max n old)+ act++-- | Room for four concurrent streams and a small window, so that WINDOW_UPDATE+-- frames go back and forth all the time.+runServerSmallWindow :: IO ()+runServerSmallWindow = runTCPServer (Just host) port runHTTP2Server+ where+ sconf =+ defaultServerConfig+ { settings =+ (settings defaultServerConfig)+ { maxConcurrentStreams = Just 4+ , initialWindowSize = 8192+ }+ }+ runHTTP2Server s = do+ setSocketOption s NoDelay 1+ E.bracket+ (allocSimpleConfig s 32768)+ freeSimpleConfig+ (\conf -> run sconf conf server)++-- | Like 'runServer', but with the connection window left at its initial+-- 65535 octets.+runServerSmallConnWindow :: IO ()+runServerSmallConnWindow = runTCPServer (Just host) port runHTTP2Server+ where+ sconf = defaultServerConfig{connectionWindowSize = defaultWindowSize}+ runHTTP2Server s =+ E.bracket+ (allocSimpleConfig s 32768)+ freeSimpleConfig+ (\conf -> run sconf conf server)++-- | Like 'runServer', but allowing a client to reset 1000 streams a second.+runServerManyResets :: IO ()+runServerManyResets = runTCPServer (Just host) port runHTTP2Server+ where+ sconf =+ defaultServerConfig+ { settings = (settings defaultServerConfig){rstRateLimit = 1000}+ }+ runHTTP2Server s =+ E.bracket+ (allocSimpleConfig s 32768)+ freeSimpleConfig+ (\conf -> run sconf conf server)++runServerMaxConc1 :: IO ()+runServerMaxConc1 = runTCPServer (Just host) port runHTTP2Server+ where+ sconf =+ defaultServerConfig+ { settings = (settings defaultServerConfig){maxConcurrentStreams = Just 1}+ }+ runHTTP2Server s =+ E.bracket+ (allocSimpleConfig s 32768)+ freeSimpleConfig+ (\conf -> run sconf conf server)++-- | A server whose handler waits long enough for a RST_STREAM to arrive+-- before it responds, and then signals that 'sendResponse' returned.+runServerCancel :: MVar () -> IO ()+runServerCancel doneVar = runTCPServer (Just host) port runHTTP2Server+ where+ runHTTP2Server s =+ E.bracket+ (allocSimpleConfig s 32768)+ freeSimpleConfig+ (\conf -> run defaultServerConfig conf cancelServer)+ cancelServer _req _aux sendResponse = do+ threadDelay 200000+ sendResponse responseHello []+ putMVar doneVar ()++runFakeServer :: MVar ByteString -> IO ()+runFakeServer prefaceVar = do+ runTCPServer (Just host) port $ \s -> do+ ref <- newIORef Nothing++ -- send settings+ sendAll s $+ "\x00\x00\x12\x04\x00\x00\x00\x00\x00"+ `mappend` "\x00\x03\x00\x00\x00\x80\x00\x04\x00"+ `mappend` "\x01\x00\x00\x00\x05\x00\xff\xff\xff"++ -- receive preface+ value <- defaultReadN s ref (B.length connectionPreface)+ putMVar prefaceVar value++ -- send goaway frame+ sendAll s "\x00\x00\x08\x07\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x01"++ -- wait for a few ms to make sure the client has a chance to close the+ -- socket on its end+ threadDelay 10000++-- | Answering two requests with a GOAWAY that leaves out the second: the+-- headers of the first response, GOAWAY(NO_ERROR) with last stream 1, and+-- the rest of the first response a little later.+runGoAwayServer :: IO ()+runGoAwayServer = runTCPServer (Just host) port $ \s -> do+ _ <- recvAll s (B.length connectionPreface)+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let awaitRequests :: Int -> IO ()+ awaitRequests 2 = return ()+ awaitRequests n = do+ mf <- recvFrame s+ case mf of+ Nothing -> return ()+ Just (FrameSettings, fh, _)+ | not (testAck (flags fh)) -> do+ sendAll s $+ encodeFrame (EncodeInfo (setAck defaultFlags) 0 Nothing) $+ SettingsFrame []+ awaitRequests n+ Just (FrameHeaders, _, _) -> awaitRequests (n + 1)+ Just _ -> awaitRequests n+ awaitRequests 0+ sendAll s $+ encodeFrame (EncodeInfo (setEndHeader defaultFlags) 1 Nothing) $+ HeadersFrame Nothing $+ hpackEncode [(":status", "200")]+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $+ GoAwayFrame 1 NoError "going"+ threadDelay 200000+ sendAll s $+ encodeFrame (EncodeInfo (setEndStream defaultFlags) 1 Nothing) $+ DataFrame "hello"+ -- Until the client closes it.+ let drain = recvFrame s >>= maybe (return ()) (const drain)+ drain++-- | How a request through the client ended: the body, or "closed" for+-- 'ConnectionIsClosed'.+status :: Either C.HTTP2Error ByteString -> ByteString+status (Right bs) = bs+status (Left C.ConnectionIsClosed) = "closed"+status (Left e) = C8.pack $ show e++-- | A request for /slow, then GOAWAY(NO_ERROR) at once. What the server+-- sends from then on until it closes the connection.+goAwayFromClient :: IO [String]+goAwayFromClient = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ sendAll s $+ encodeFrame (EncodeInfo (setEndStream $ setEndHeader defaultFlags) 1 Nothing) $+ HeadersFrame Nothing $+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/slow")+ , (":method", "GET")+ ]+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $+ GoAwayFrame 0 NoError "going"+ collect s+ where+ collect s = do+ mf <- recvFrame s+ case mf of+ Nothing -> return []+ Just (FrameHeaders, fh, _) -> (("HEADERS " ++ show (streamId fh)) :) <$> collect s+ Just (FrameData, fh, _)+ | testEndStream (flags fh) ->+ (("DATA " ++ show (streamId fh) ++ " END_STREAM") :) <$> collect s+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ (("GOAWAY " ++ show err) :) <$> collect s+ Just _ -> collect s++server :: Server+server req aux sendResponse = case requestMethod req of+ Just "GET" -> case requestPath req of+ Just "/" -> sendResponse responseHello []+ -- A moment before the answer.+ Just "/slow" -> threadDelay 200000 >> sendResponse responseHello []+ Just "/early" -> do+ auxSendInformational+ aux+ earlyHints103+ [("link", "</style.css>; rel=preload; as=style")]+ auxSendInformational+ aux+ earlyHints103+ [("link", "</app.js>; rel=preload; as=script")]+ sendResponse responseHello []+ Just "/stream" -> sendResponse responseInfinite []+ -- Like /stream, but going quietly once the client resets it.+ Just "/endless" -> sendResponse responseEndless []+ Just "/not-modified" -> sendResponse (responseNoBody notModified304 bigLength) []+ -- Says it has content, and has none: malformed.+ Just "/no-content" -> sendResponse (responseNoBody ok200 bigLength) []+ Just "/big" -> sendResponse responseBig []+ Just "/push" -> do+ let pp = pushPromise "/push-pp" responsePP 0+ sendResponse responseHello [pp]+ -- A push of 20000 octets.+ Just "/push-big" -> do+ let pp = pushPromise "/push-big-pp" responsePushBig 0+ sendResponse responseHello [pp]+ _ -> sendResponse response404 []+ Just "POST" -> case requestPath req of+ Just "/echo" -> sendResponse (responseEcho req) []+ -- How many octets of body arrived.+ Just "/count" -> do+ let count n = do+ bs <- getRequestBodyChunk req+ if B.null bs then return n else count (n + B.length bs)+ n <- count (0 :: Int)+ sendResponse (responseBuilder ok200 [] (byteString (C8.pack (show n)))) []+ Just "/both" -> do+ -- Read the body on the side, so that the response does not+ -- wait for it.+ _ <-+ forkIO $+ let d = getRequestBodyChunk req >>= \bs -> unless (B.null bs) d+ in d+ sendResponse responseBoth []+ _ -> sendResponse responseHello []+ Just "HEAD" -> case requestPath req of+ -- HEADERS, then an empty DATA frame with END_STREAM.+ Just "/data" -> sendResponse (responseBuilder ok200 bigLength mempty) []+ -- HEADERS with END_STREAM.+ _ -> sendResponse (responseNoBody ok200 bigLength) []+ _ -> sendResponse response405 []++-- | Larger than the default frame size and than the server's 32K buffer.+bigVal :: ByteString+bigVal = C8.replicate 40000 'x'++responseBig :: Response+responseBig = setResponseTrailersMaker rsp maker+ where+ rsp = responseBuilder ok200 [("x-big", bigVal)] "hello"+ maker Nothing = return $ Trailers [("x-big-trailer", bigVal)]+ maker (Just _) = return $ NextTrailersMaker maker++-- | The stream error a client raises for a malformed response, as+-- 'sendRequest' hands it on.+malformedResponse :: C.HTTP2Error -> Bool+malformedResponse (C.StreamErrorIsSent C.ProtocolError _ _) = True+malformedResponse (C.BadThingHappen se) =+ maybe False malformedResponse $ E.fromException se+malformedResponse _ = False++-- | The content-length of content that is not there.+bigLength :: ResponseHeaders+bigLength = [("content-length", "1234")]++responseHello :: Response+responseHello = responseBuilder ok200 header body+ where+ header = [("Content-Type", "text/plain")]+ body = byteString "Hello, world!\n"++earlyHints103 :: Status+earlyHints103 = mkStatus 103 "Early Hints"++responsePushBig :: Response+responsePushBig = responseBuilder ok200 [] $ byteString $ C8.replicate 20000 'p'++responsePP :: Response+responsePP = responseBuilder ok200 header body+ where+ header =+ [ ("Content-Type", "text/plain")+ , ("x-push", "True")+ ]+ body = byteString "Push\n"++-- | A streaming response that does not wait for the request body, so that+-- both ends are sending at once and either can finish first.+responseBoth :: Response+responseBoth = responseStreaming ok200 [] $ \write flush ->+ replicateM_ 50 $ write (byteString (C8.replicate 50 'b')) >> flush++responseEndless :: Response+responseEndless = responseStreaming ok200 [] body+ where+ body :: (Builder -> IO ()) -> IO () -> IO ()+ body write flush = forever (write (byteString chunk) *> flush) `E.catch` quiet+ chunk = C8.replicate 1024 'x'+ quiet :: E.SomeException -> IO ()+ quiet _ = return ()++responseInfinite :: Response+responseInfinite = responseStreaming ok200 header body+ where+ header = [("Content-Type", "text/plain")]+ body :: (Builder -> IO ()) -> IO () -> IO ()+ body write flush = do+ let go n = write (byteString (C8.pack (show n)) `mappend` "\n") *> flush *> go (succ n)+ go (0 :: Int)++response404 :: Response+response404 = responseNoBody notFound404 []++response405 :: Response+response405 = responseNoBody methodNotAllowed405 []++responseEcho :: Request -> Response+responseEcho req = setResponseTrailersMaker h2rsp maker+ where+ h2rsp = responseStreaming ok200 header streamingBody+ header = [("Content-Type", "text/plain")]+ mhx = getFieldValue (toToken "X-Tag") (snd (requestHeaders req))+ streamingBody write _flush = do+ loop+ mt <- getRequestTrailers req+ firstTrailerValue <$> mt `shouldBe` mhx+ where+ loop = do+ bs <- getRequestBodyChunk req+ when (bs /= "") $ do+ void $ write $ byteString bs+ loop+ maker = trailersMaker (CH.hashInit :: Context SHA1)++-- Strictness is important for Context.+trailersMaker :: Context SHA1 -> Maybe ByteString -> IO NextTrailersMaker+trailersMaker ctx Nothing = return $ Trailers [("X-SHA1", sha1)]+ where+ !sha1 = C8.pack $ show $ CH.hashFinalize ctx+trailersMaker ctx (Just bs) = return $ NextTrailersMaker $ trailersMaker ctx'+ where+ !ctx' = CH.hashUpdate ctx bs++-- | Request @/early@ with an informational handler installed, recording each+-- 103 Early Hints section and returning the final response status.+runClientEarly :: IORef [TokenHeaderTable] -> IO (Maybe Status)+runClientEarly hintsRef = runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf0 ->+ C.run cliconf (conf0{confOnInformational = onInformational}) $ \sendRequest _aux ->+ sendRequest (C.requestNoBody methodGet "/early" []) (return . C.responseStatus)+ where+ cliconf = C.defaultClientConfig{C.authority = host}+ onInformational _sid tbl = modifyIORef' hintsRef (++ [tbl])++runClient :: (Socket -> BufferSize -> IO Config) -> IO ()+runClient allocConfig =+ runTCPClient host port runHTTP2Client+ where+ auth = host+ cliconf = C.defaultClientConfig{C.authority = auth}+ runHTTP2Client s =+ E.bracket+ (allocConfig s 4096)+ freeSimpleConfig+ (\conf -> C.run cliconf conf client)++ client :: C.Client ()+ client sendRequest aux =+ foldr1+ concurrently_+ [ client0 sendRequest aux+ , client1 sendRequest aux+ , client2 sendRequest aux+ , client3 sendRequest aux+ , client3' sendRequest aux+ , client3'' sendRequest aux+ , client4 sendRequest aux+ , client5 sendRequest aux+ ]++-- delay sending preface to be able to test if it is always sent first+allocSlowPrefaceConfig :: Socket -> BufferSize -> IO Config+allocSlowPrefaceConfig s size = do+ config <- allocSimpleConfig s size+ pure config{confSendAll = slowPrefaceSend (confSendAll config)}+ where+ slowPrefaceSend :: (ByteString -> IO ()) -> ByteString -> IO ()+ slowPrefaceSend orig chunk = do+ when (C8.pack "PRI" `C8.isPrefixOf` chunk) $ do+ threadDelay 10000+ orig chunk++client0 :: C.Client ()+client0 sendRequest _aux = do+ let req = C.requestNoBody methodGet "/" []+ sendRequest req $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just ok200+ fmap statusMessage (C.responseStatus rsp) `shouldBe` Just "OK"++client1 :: C.Client ()+client1 sendRequest _aux = do+ let req = C.requestNoBody methodGet "/push-pp" []+ sendRequest req $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just notFound404++client2 :: C.Client ()+client2 sendRequest _aux = do+ let req = C.requestNoBody methodPut "/" []+ sendRequest req $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just methodNotAllowed405++client3 :: C.Client ()+client3 sendRequest _aux = do+ let hx = "b0870457df2b8cae06a88657a198d9b52f8e2b0a"+ req0 =+ C.requestFile methodPost "/echo" [("X-Tag", hx)] $+ FileSpec "test/inputFile" 0 1012731+ req = C.setRequestTrailersMaker req0 maker+ sendRequest req $ \rsp -> do+ let consumeBody = do+ bs <- C.getResponseBodyChunk rsp+ when (bs /= "") consumeBody+ consumeBody+ mt <- C.getResponseTrailers rsp+ firstTrailerValue <$> mt `shouldBe` Just hx+ where+ !maker = trailersMaker (CH.hashInit :: Context SHA1)++client3' :: C.Client ()+client3' sendRequest _aux = do+ let hx = "b0870457df2b8cae06a88657a198d9b52f8e2b0a"+ req0 = C.requestStreaming methodPost "/echo" [("X-Tag", hx)] $ \write _flush -> do+ let sendFile h = do+ bs <- B.hGet h 1024+ when (bs /= "") $ do+ write $ byteString bs+ sendFile h+ withFile "test/inputFile" ReadMode sendFile+ req = C.setRequestTrailersMaker req0 maker+ sendRequest req $ \rsp -> do+ let consumeBody = do+ bs <- C.getResponseBodyChunk rsp+ when (bs /= "") consumeBody+ consumeBody+ mt <- C.getResponseTrailers rsp+ firstTrailerValue <$> mt `shouldBe` Just hx+ where+ !maker = trailersMaker (CH.hashInit :: Context SHA1)++client3'' :: C.Client ()+client3'' sendRequest _axu = do+ let hx = "59f82dfddc0adf5bdf7494b8704f203a67e25d4a"+ req0 = C.requestStreaming methodPost "/echo" [("X-Tag", hx)] $ \write _flush -> do+ let chunk = C8.replicate (16384 * 2) 'c'+ tag = C8.replicate 16 't'+ -- I don't think 9 is important here, this is just what I have, the client hangs on receiving the last one+ replicateM_ 9 $ write $ byteString chunk+ write $ byteString tag+ req = C.setRequestTrailersMaker req0 maker+ sendRequest req $ \rsp -> do+ let consumeBody = do+ bs <- C.getResponseBodyChunk rsp+ when (bs /= "") consumeBody+ consumeBody+ mt <- C.getResponseTrailers rsp+ firstTrailerValue <$> mt `shouldBe` Just hx+ where+ !maker = trailersMaker (CH.hashInit :: Context SHA1)++client4 :: C.Client ()+client4 sendRequest _aux = do+ let req0 = C.requestNoBody methodGet "/push" []+ sendRequest req0 $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just ok200+ let req1 = C.requestNoBody methodGet "/push-pp" []+ sendRequest req1 $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just ok200++client5 :: C.Client ()+client5 sendRequest _aux = do+ let req0 = C.requestNoBody methodGet "/stream" []+ sendRequest req0 $ \rsp -> do+ C.responseStatus rsp `shouldBe` Just ok200+ let go n+ | n > 0 = do+ _ <- C.getResponseBodyChunk rsp+ go (pred n)+ | otherwise = pure ()+ go (100 :: Int)++firstTrailerValue :: TokenHeaderTable -> FieldValue+firstTrailerValue tbl = case fst tbl of+ [] -> error "firstTrailerValue"+ x : _ -> snd x++runAttack :: (C.ClientIO -> IO ()) -> IO ()+runAttack attack =+ runTCPClient host port runHTTP2Client+ where+ auth = host+ cliconf = C.defaultClientConfig{C.authority = auth}+ runHTTP2Client s =+ E.bracket+ (allocSimpleConfig s 4096)+ freeSimpleConfig+ (\conf -> C.runIO cliconf conf client)+ client cconf = return $ do+ attack cconf+ threadDelay 1000000++rapidSettings :: C.ClientIO -> IO ()+rapidSettings C.ClientIO{..} = do+ let einfo = EncodeInfo defaultFlags 0 Nothing+ bs = encodeFrame einfo $ SettingsFrame [(SettingsEnablePush, 0)]+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs++rapidPing :: C.ClientIO -> IO ()+rapidPing C.ClientIO{..} = do+ let einfo = EncodeInfo defaultFlags 0 Nothing+ opaque64 = "01234567"+ bs = encodeFrame einfo $ PingFrame opaque64+ replicateM_ 20 $ cioWriteBytes bs++rapidEmptyHeader :: C.ClientIO -> IO ()+rapidEmptyHeader C.ClientIO{..} = do+ (sid, _) <- cioCreateStream+ let einfo = EncodeInfo defaultFlags sid Nothing+ bs = encodeFrame einfo $ HeadersFrame Nothing ""+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs+ cioWriteBytes bs++rapidEmptyData :: C.ClientIO -> IO ()+rapidEmptyData C.ClientIO{..} = do+ (sid, _) <- cioCreateStream+ let einfoH = EncodeInfo (setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/")+ , (":method", "GET")+ ]+ bsH = encodeFrame einfoH $ HeadersFrame Nothing hdr+ cioWriteBytes bsH+ let einfoD = EncodeInfo defaultFlags sid Nothing+ bsD = encodeFrame einfoD $ DataFrame ""+ cioWriteBytes bsD+ cioWriteBytes bsD+ cioWriteBytes bsD+ cioWriteBytes bsD+ cioWriteBytes bsD+ cioWriteBytes bsD+ cioWriteBytes bsD+ cioWriteBytes bsD++rapidRst :: C.ClientIO -> IO ()+rapidRst C.ClientIO{..} = do+ reset+ reset+ reset+ reset+ reset+ reset+ reset+ reset+ where+ reset = do+ (sid, _) <- cioCreateStream+ -- setEndStream for HalfClosedRemote+ let einfoH = EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/")+ , (":method", "GET")+ ]+ bsH = encodeFrame einfoH $ HeadersFrame Nothing hdr+ cioWriteBytes bsH+ let einfoR = EncodeInfo defaultFlags sid Nothing+ -- Only (HalfClosedRemote, NoError) is accepted.+ -- Otherwise, a stream error terminates the connection.+ bsR = encodeFrame einfoR $ RSTStreamFrame NoError+ cioWriteBytes bsR++-- | MadeYouReset (CVE-2025-8671): the same churn as 'rapidRst' without a+-- single RST_STREAM from us. Each stream gets a handler that goes on+-- running, then a PRIORITY making it depend on itself, which the server+-- answers by resetting the stream -- giving its concurrency slot back while+-- the handler runs on. Those resets did not count against the limit on+-- resets, so this could be kept up for as long as the peer liked.+rapidStreamError :: IO (Maybe (ErrorCode, ByteString))+rapidStreamError = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ forM_ [1, 3 .. 15] $ \sid -> do+ let einfoH = EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/stream")+ , (":method", "GET")+ ]+ einfoP = EncodeInfo defaultFlags sid Nothing+ sendAll s $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ sendAll s $ encodeFrame einfoP $ PriorityFrame $ Priority False sid 16+ awaitGoAway s++-- | What the server says in its GOAWAY, if it sends one before closing.+awaitGoAway :: Socket -> IO (Maybe (ErrorCode, ByteString))+awaitGoAway s = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err msg) <- decodeGoAwayFrame fh p ->+ return $ Just (err, msg)+ Just _ -> awaitGoAway s++-- | A SETTINGS_INITIAL_WINDOW_SIZE that takes an open stream's window past+-- 2^31-1. What the server answers with: whether it reset the stream, and+-- the error in its GOAWAY.+settingsOverflow :: IO (Bool, Maybe ErrorCode)+settingsOverflow = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let sid = 1+ -- No END_STREAM: the stream stays open, waiting for the body.+ einfoH = EncodeInfo (setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/echo")+ , (":method", "POST")+ ]+ sendAll s $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ -- The stream's window is now the largest there is ...+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags sid Nothing) $+ WindowUpdateFrame (maxWindowSize - defaultWindowSize)+ -- ... and one more octet of initial window takes it over.+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $+ SettingsFrame [(SettingsInitialWindowSize, defaultWindowSize + 1)]+ answer s False+ where+ answer s reset = do+ mf <- recvFrame s+ case mf of+ Nothing -> return (reset, Nothing)+ Just (FrameRSTStream, _, _) -> answer s True+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ return (reset, Just err)+ Just _ -> answer s reset++-- | A request whose body ends with an empty trailer block. What the+-- server answers it with.+emptyTrailers :: IO (Maybe String)+emptyTrailers = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let sid = 1+ einfoH = EncodeInfo (setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/count")+ , (":method", "POST")+ ]+ einfoT = EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing+ sendAll s $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ sendAll s $ encodeFrame (EncodeInfo defaultFlags sid Nothing) $ DataFrame "body"+ sendAll s $ encodeFrame einfoT $ HeadersFrame Nothing ""+ answer s sid+ where+ answer s sid = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameHeaders, fh, _)+ | streamId fh == sid -> return $ Just "HEADERS"+ Just (FrameRSTStream, fh, p)+ | streamId fh == sid+ , Right (RSTStreamFrame err) <- decodeRSTStreamFrame fh p ->+ return $ Just $ "RST_STREAM " ++ show err+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ return $ Just $ "GOAWAY " ++ show err+ Just _ -> answer s sid++-- | A request with a content-length of 4, whose body comes in padded DATA+-- frames: "bo", "dy", and an empty one with END_STREAM. What the server+-- answers it with: the body of its response, which is the number of octets+-- of body that arrived.+paddedBody :: IO (Maybe String)+paddedBody = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let sid = 1+ einfoH = EncodeInfo (setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/count")+ , (":method", "POST")+ , ("content-length", "4")+ ]+ padding = C8.replicate 10 '\0'+ einfoD = EncodeInfo defaultFlags sid (Just padding)+ einfoE = EncodeInfo (setEndStream defaultFlags) sid (Just padding)+ sendAll s $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ sendAll s $ encodeFrame einfoD $ DataFrame "bo"+ sendAll s $ encodeFrame einfoD $ DataFrame "dy"+ sendAll s $ encodeFrame einfoE $ DataFrame ""+ answer s sid+ where+ answer s sid = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameData, fh, p)+ | streamId fh == sid -> return $ Just $ "DATA " ++ C8.unpack p+ Just (FrameRSTStream, fh, p)+ | streamId fh == sid+ , Right (RSTStreamFrame err) <- decodeRSTStreamFrame fh p ->+ return $ Just $ "RST_STREAM " ++ show err+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ return $ Just $ "GOAWAY " ++ show err+ Just _ -> answer s sid++-- | A request whose body is 2000 DATA frames of one octet each, padded to+-- 257 octets of payload, more than the stream's window in all. Sent+-- without waiting for WINDOW_UPDATE: every octet of padding has to have+-- been given back by the time the next frame is checked. What the server+-- answers it with: the number of octets of body that arrived.+paddingWindow :: IO (Maybe String)+paddingWindow = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let sid = 1+ einfoH = EncodeInfo (setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/count")+ , (":method", "POST")+ ]+ padding = C8.replicate 255 '\0'+ einfoD = EncodeInfo defaultFlags sid (Just padding)+ einfoE = EncodeInfo (setEndStream defaultFlags) sid Nothing+ sendAll s $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ replicateM_ 2000 $ sendAll s $ encodeFrame einfoD $ DataFrame "x"+ sendAll s $ encodeFrame einfoE $ DataFrame ""+ answer s sid+ where+ answer s sid = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameData, fh, p)+ | streamId fh == sid -> return $ Just $ "DATA " ++ C8.unpack p+ Just (FrameRSTStream, fh, p)+ | streamId fh == sid+ , Right (RSTStreamFrame err) <- decodeRSTStreamFrame fh p ->+ return $ Just $ "RST_STREAM " ++ show err+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err msg) <- decodeGoAwayFrame fh p ->+ return $ Just $ "GOAWAY " ++ show err ++ " " ++ C8.unpack msg+ Just _ -> answer s sid++-- | 16384 octets of DATA on a stream we have half-closed, then a request+-- with a body of 16384 octets. What the server gives back to the+-- connection window before answering the request, and the answer: the+-- number of octets of body that arrived.+refusedData :: IO (Maybe (Int, String))+refusedData = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let request sid method path flags =+ encodeFrame (EncodeInfo (flags $ setEndHeader defaultFlags) sid Nothing) $+ HeadersFrame Nothing $+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", path)+ , (":method", method)+ ]+ chunk = C8.replicate 16384 'x'+ -- A response that goes on for ever keeps stream 1 in the table.+ sendAll s $ request 1 "GET" "/stream" setEndStream+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 1 Nothing) $ DataFrame chunk+ sendAll s $ request 3 "POST" "/count" id+ sendAll s $+ encodeFrame (EncodeInfo (setEndStream defaultFlags) 3 Nothing) $+ DataFrame chunk+ answer s 0+ where+ answer s n = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameWindowUpdate, fh, p)+ | streamId fh == 0+ , Right (WindowUpdateFrame w) <- decodeWindowUpdateFrame fh p ->+ answer s (n + w)+ Just (FrameData, fh, p)+ | streamId fh == 3 -> return $ Just (n, C8.unpack p)+ Just (FrameGoAway, _, _) -> return Nothing+ Just _ -> answer s n++-- | A request for /push, which comes with a push, from a peer that has+-- announced room for no streams of the server's. What the server answers+-- it with first.+pushNoRoom :: IO (Maybe String)+pushNoRoom = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $+ SettingsFrame [(SettingsMaxConcurrentStreams, 0)]+ -- The server takes our SETTINGS on board once it has acknowledged them,+ -- and before it goes on to what comes next: the answer to a PING sent+ -- after them.+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $+ PingFrame "12345678"+ awaitPingAck s+ let sid = 1+ einfoH = EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/push")+ , (":method", "GET")+ ]+ sendAll s $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ answer s sid+ where+ awaitPingAck s = do+ mf <- recvFrame s+ case mf of+ Just (FramePing, fh, _) | testAck (flags fh) -> return ()+ Just _ -> awaitPingAck s+ Nothing -> return ()+ answer s sid = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameHeaders, fh, _)+ | streamId fh == sid -> return $ Just "HEADERS"+ Just (FramePushPromise, _, _) -> return $ Just "PUSH_PROMISE"+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ return $ Just $ "GOAWAY " ++ show err+ Just _ -> answer s sid++-- | Two requests that are answered for ever, then SETTINGS that are a+-- connection error. The last stream identifier and the error of the+-- server's GOAWAY.+goAwayLastStream :: IO (Maybe (StreamId, ErrorCode))+goAwayLastStream = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ forM_ [1, 3] $ \sid ->+ sendAll s $+ encodeFrame (EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing) $+ HeadersFrame Nothing $+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/endless")+ , (":method", "GET")+ ]+ -- SETTINGS_ENABLE_PUSH can only be 0 or 1.+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $+ SettingsFrame [(SettingsEnablePush, 2)]+ answer s+ where+ answer s = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame sid err _) <- decodeGoAwayFrame fh p ->+ return $ Just (sid, err)+ Just _ -> answer s++-- | A request with an upper-case field name on stream 1, then a good one+-- on stream 3. What the server sends on them, up to the answer on 3.+malformedRequest :: IO [String]+malformedRequest = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let request sid extra =+ encodeFrame (EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing) $+ HeadersFrame Nothing $+ hpackEncode $+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/")+ , (":method", "GET")+ ]+ ++ extra+ sendAll s $ request 1 [("X-Upper", "1")]+ sendAll s $ request 3 []+ collect s+ where+ collect s = do+ mf <- recvFrame s+ case mf of+ Nothing -> return []+ Just (FrameHeaders, fh, _)+ | streamId fh == 3 -> return ["HEADERS 3"]+ Just (FrameRSTStream, fh, p)+ | Right (RSTStreamFrame err) <- decodeRSTStreamFrame fh p ->+ (("RST_STREAM " ++ show (streamId fh) ++ " " ++ show err) :) <$> collect s+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ return ["GOAWAY " ++ show err]+ Just _ -> collect s++-- | /endless until the server has used up the connection window, which we+-- do not open, then a request whose answer has no body: what the server+-- sends for it. Then the windows opened a little: whether the body held+-- back goes on.+shutWindow :: IO (Maybe String)+shutWindow = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ let request sid path =+ encodeFrame (EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing) $+ HeadersFrame Nothing $+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", path)+ , (":method", "GET")+ ]+ sendAll s $ request 1 "/endless"+ used <- untilShut s 0+ if not used+ then return Nothing+ else do+ sendAll s $ request 3 "/not-modified"+ ma <- answer s+ case ma of+ Just "HEADERS 3" -> do+ forM_ [0, 1] $ \sid ->+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags sid Nothing) $+ WindowUpdateFrame 1000+ fmap ("HEADERS 3, then " ++) <$> resumed s+ _ -> return ma+ where+ resumed s = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameData, fh, _)+ | streamId fh == 1 -> return $ Just "DATA 1"+ Just _ -> resumed s+ untilShut s n+ | n >= defaultWindowSize = return True+ | otherwise = do+ mf <- recvFrame s+ case mf of+ Nothing -> return False+ Just (FrameData, fh, _) -> untilShut s (n + payloadLength fh)+ Just _ -> untilShut s n+ answer s = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameHeaders, fh, _)+ | streamId fh == 3 -> return $ Just "HEADERS 3"+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ return $ Just $ "GOAWAY " ++ show err+ Just _ -> answer s++-- | PRIORITY frames for 100 streams that are never opened, then a request.+-- What the server answers the request with.+idlePriority :: IO (Maybe String)+idlePriority = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ forM_ [3, 5 .. 201] $ \sid ->+ sendAll s $+ encodeFrame (EncodeInfo defaultFlags sid Nothing) $+ PriorityFrame $+ Priority False 0 16+ let sid = 203+ einfoH = EncodeInfo (setEndStream $ setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/")+ , (":method", "GET")+ ]+ sendAll s $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ answer s sid+ where+ answer s sid = do+ mf <- recvFrame s+ case mf of+ Nothing -> return Nothing+ Just (FrameHeaders, fh, _)+ | streamId fh == sid -> return $ Just "HEADERS"+ Just (FrameRSTStream, fh, p)+ | streamId fh == sid+ , Right (RSTStreamFrame err) <- decodeRSTStreamFrame fh p ->+ return $ Just $ "RST_STREAM " ++ show err+ Just (FrameGoAway, fh, p)+ | Right (GoAwayFrame _ err _) <- decodeGoAwayFrame fh p ->+ return $ Just $ "GOAWAY " ++ show err+ Just _ -> answer s sid++-- | Exactly so many octets off a raw connection, or fewer once it is closed.+recvAll :: Socket -> Int -> IO ByteString+recvAll s n0 = go n0 []+ where+ go 0 acc = return $ B.concat $ reverse acc+ go k acc = do+ bs <- recv s k+ if B.null bs+ then return $ B.concat $ reverse acc+ else go (k - B.length bs) (bs : acc)++-- | One frame off a raw connection, or 'Nothing' once it is closed.+recvFrame :: Socket -> IO (Maybe (FrameType, FrameHeader, ByteString))+recvFrame s = do+ mh <- recvExactly frameHeaderLength+ case mh of+ Nothing -> return Nothing+ Just h -> do+ let (ftyp, fh) = decodeFrameHeader h+ fmap (\p -> (ftyp, fh, p)) <$> recvExactly (payloadLength fh)+ where+ recvExactly n = go n []+ where+ go 0 acc = return $ Just $ B.concat $ reverse acc+ go k acc = do+ bs <- recv s k+ if B.null bs+ then return Nothing+ else go (k - B.length bs) (bs : acc)++-- | Open a stream, reset it, then open two more. The server announced room+-- for one concurrent stream, so the third one here must be refused.+--+-- Closing a stream used to give its slot back twice -- a RST_STREAM carrying a+-- non-critical error code is closed by both 'stream' and 'processState' -- so+-- the count drifted down by one on every reset and this sequence went through+-- unchallenged.+-- | A HEADERS frame opening a stream and leaving it open, so that it goes on+-- holding a concurrency slot.+--+-- Stream identifiers are written out rather than taken from+-- 'C.cioCreateStream': the limit being overrun is the one the server+-- announced, and asking for a stream the proper way would block on that same+-- limit on this side.+openStreamFrame :: StreamId -> ByteString+openStreamFrame sid = encodeFrame einfo $ HeadersFrame Nothing hdr+ where+ einfo = EncodeInfo (setEndHeader defaultFlags) sid Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/")+ , (":method", "GET")+ ]++-- | Speak raw frames to the server and collect what it says back.+rawExchange :: [ByteString] -> IO [(FrameType, StreamId, ByteString)]+rawExchange out = runTCPClient host port $ \s -> do+ sendAll s connectionPreface+ sendAll s $ encodeFrame (EncodeInfo defaultFlags 0 Nothing) $ SettingsFrame []+ mapM_ (sendAll s) out+ splitFrames <$> collect mempty s+ where+ collect acc s = do+ mbs <- timeout 300000 $ recv s 4096+ case mbs of+ Just bs | not (B.null bs) -> collect (acc `B.append` bs) s+ _ -> return acc++splitFrames :: ByteString -> [(FrameType, StreamId, ByteString)]+splitFrames bs+ | B.length bs < frameHeaderLength = []+ | otherwise =+ let (h, rest) = B.splitAt frameHeaderLength bs+ (typ, FrameHeader{payloadLength, streamId}) = decodeFrameHeader h+ (body, rest') = B.splitAt payloadLength rest+ in (typ, streamId, body) : splitFrames rest'++-- | The RST_STREAM and GOAWAY frames among them, with their error codes.+resets+ :: [(FrameType, StreamId, ByteString)] -> [(FrameType, StreamId, ErrorCode)]+resets frames =+ [ (typ, sid, ec)+ | (typ, sid, body) <- frames+ , typ == FrameRSTStream || typ == FrameGoAway+ , Just ec <- [errorCodeOf typ sid body]+ ]+ where+ errorCodeOf FrameRSTStream sid body =+ case decodeRSTStreamFrame (FrameHeader (B.length body) defaultFlags sid) body of+ Right (RSTStreamFrame ec) -> Just ec+ _ -> Nothing+ errorCodeOf FrameGoAway sid body =+ case decodeGoAwayFrame (FrameHeader (B.length body) defaultFlags sid) body of+ Right (GoAwayFrame _ ec _) -> Just ec+ _ -> Nothing+ errorCodeOf _ _ _ = Nothing++-- | Open a stream and cancel it straight away, while the server is still+-- working on the response.+--+-- The sender skips a stream that is already half-closed, and used to return+-- without telling the thread that enqueued the output. That thread sat in+-- 'syncWithSender'' on an MVar nothing would fill, so 'sendResponse' never+-- returned and the worker was only reclaimed when the timeout manager killed+-- it, seconds later.+cancelInFlight :: C.ClientIO -> IO ()+cancelInFlight C.ClientIO{..} = do+ -- setEndStream for HalfClosedRemote, so that CANCEL is accepted as a+ -- stream error rather than taken down the connection.+ let einfoH = EncodeInfo (setEndStream $ setEndHeader defaultFlags) 1 Nothing+ hdr =+ hpackEncode+ [ (":scheme", "http")+ , (":authority", "127.0.0.1")+ , (":path", "/")+ , (":method", "GET")+ ]+ cioWriteBytes $ encodeFrame einfoH $ HeadersFrame Nothing hdr+ cioWriteBytes $+ encodeFrame (EncodeInfo defaultFlags 1 Nothing) $+ RSTStreamFrame Cancel++-- | Send a malformed request, then a good one down the same connection.+--+-- RFC 9113 section 8.1.1 makes a malformed request a stream error, so the+-- server must reset that one stream and keep serving: the second request is+-- the point of the test. The whole connection used to come down with the+-- first, taking every other stream on it along.+runStreamErrorClient :: IO ()+runStreamErrorClient = runTCPClient host port $ \s ->+ E.bracket (allocSimpleConfig s 4096) freeSimpleConfig $ \conf ->+ C.run cliconf conf $ \sendRequest _aux -> do+ -- "te" may only ever be "trailers" (section 8.2.2), and unlike+ -- "connection" it is not one of the headers the sender strips.+ let bad = C.requestNoBody methodGet "/" [("te", "gzip")]+ sendRequest bad (\_ -> return ()) `shouldThrow` streamWasReset+ let good = C.requestNoBody methodGet "/" []+ sendRequest good $ \rsp ->+ C.responseStatus rsp `shouldBe` Just ok200+ where+ cliconf = C.defaultClientConfig{C.authority = host}++streamWasReset :: Selector C.HTTP2Error+streamWasReset C.StreamResetIsReceived{} = True+streamWasReset _ = False++-- | A HEADERS frame with PADDED and PRIORITY set, six octets of payload and a+-- Pad Length of five, so that the padding covers the whole of the priority+-- fields the flag promises.+--+-- Six octets is the smallest payload the frame header check accepts for those+-- two flags together, so this gets through it; the decoder then took the five+-- priority octets out of what padding had left empty, reading off the end of+-- the buffer. The empty ByteString is the shared one, whose pointer is null,+-- so what died was the process rather than the connection.+paddingOverPriority :: C.ClientIO -> IO ()+paddingOverPriority C.ClientIO{..} = do+ let flags = setPadded $ setPriority $ setEndHeader defaultFlags+ header = encodeFrameHeader FrameHeaders $ FrameHeader 6 flags 1+ payload = B.pack [5, 0, 0, 0, 0, 0] -- Pad Length 5, then the padding+ cioWriteBytes $ header `B.append` payload connectionError :: C.ReasonPhrase -> C.HTTP2Error -> Bool connectionError phrase (C.ConnectionErrorIsReceived _ _ p)