packages feed

http2 5.3.11 → 5.4.7

raw patch · 38 files changed

Files

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)