packages feed

http2 5.4.6 → 5.4.7

raw patch · 15 files changed

+1203/−124 lines, 15 filesPVP ok

version bump matches the API change (PVP)

API changes (from Hackage documentation)

Files

ChangeLog.md view
@@ -1,5 +1,57 @@ # 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
Network/HPACK/HeaderBlock/Decode.hs view
@@ -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,8 +56,12 @@ --     '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@@ -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) $ E.throwIO IllegalHeaderName+                            when (isJust mx) $ malformed IllegalHeaderName                             -- unknown-                            when (isMaxTokenIx tokenIx) $ E.throwIO IllegalHeaderName+                            when (isMaxTokenIx tokenIx) $ malformed IllegalHeaderName                             unsafeWrite arr tokenIx (Just v)                             pseudo                         else do                             -- 0-Length Headers Leak - CVE-2019-9516-                            when (tokenKey == "") $ E.throwIO IllegalHeaderName+                            when (tokenKey == "") $ malformed IllegalHeaderName                             when (isMaxTokenIx tokenIx && B8.any isUpper (original tokenKey)) $-                                E.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 = E.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 $ E.throwIO IllegalHeaderName+                        when isPseudo $ malformed IllegalHeaderName                         -- 0-Length Headers Leak - CVE-2019-9516-                        when (tokenKey == "") $ E.throwIO IllegalHeaderName+                        when (tokenKey == "") $ malformed IllegalHeaderName                         when (isMaxTokenIx tokenIx && B8.any isUpper (original tokenKey)) $-                            E.throwIO IllegalHeaderName+                            malformed IllegalHeaderName                         unsafeWrite arr tokenIx (Just v)                         if isCookieTokenIx tokenIx                             then normal (n + 1) builder (cookie << v)@@ -192,6 +204,22 @@                                     tvs = (tokenCookie, v) : tvs0                                 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
Network/HPACK/Huffman/Decode.hs view
@@ -59,11 +59,26 @@     -> 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)@@ -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/Table/Dynamic.hs view
@@ -222,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@@ -309,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 =
Network/HTTP2/Client/Run.hs view
@@ -89,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.@@ -117,6 +121,26 @@         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@@ -149,10 +173,27 @@     runSender = frameSender ctx conf     runClientReceiver = do         labelMe "H2 ClientReceiver"-        er <- race runReceiver runClient-        case er of-            Right r -> return r-            Left err -> E.throwIO err+        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
Network/HTTP2/H2/Context.hs view
@@ -7,6 +7,7 @@ import Control.Concurrent.STM 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@@ -91,6 +92,13 @@     -- ^ 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 -} @@ -156,6 +164,8 @@     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@@ -289,6 +299,31 @@     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@@ -369,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@@ -384,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@@ -395,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
@@ -115,10 +115,17 @@ -- -- 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+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.+-- 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@@ -132,18 +139,21 @@ hpackDecode illegal hdrblk sid Context{..} =     decodeTokenHeader decodeDynamicTable hdrblk `E.catch` handl   where-    -- Connection errors, both of them, even though a malformed message is a-    -- stream error by RFC 9113 section 8.1.1.  Either way the field block was-    -- abandoned part-way through, so 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.-    ---    -- A malformed message caught /after/ a complete decode is a different-    -- matter, and 'hpackDecodeHeader' reports those as stream errors.+    -- 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 $ ConnectionErrorIsSent ProtocolError sid illegal+        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 $ ConnectionErrorIsSent CompressionError sid msg
Network/HTTP2/H2/Receiver.hs view
@@ -296,7 +296,7 @@  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@@ -311,13 +311,12 @@     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  -- Transition (process2)-processState (Open _ (HasBody tbl@(_, reqvt))) ctx@Context{..} strm@Stream{streamInput, streamRxQ} _streamId = do+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@@ -336,8 +335,7 @@     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  -- Transition (process4)@@ -356,6 +354,21 @@     -- Idle     setStreamState ctx strm s +-- | 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 -}@@ -477,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@@ -507,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)@@ -524,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@@ -618,9 +646,9 @@     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@@ -638,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 == ""@@ -649,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@@ -743,9 +800,20 @@ stream FrameContinuation FrameHeader{streamId} _ _ _ _ =     E.throwIO $         ConnectionErrorIsSent ProtocolError streamId "continue frame cannot come here"+-- 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)
Network/HTTP2/H2/Sender.hs view
@@ -89,23 +89,29 @@     ctx@Context{outputQ, controlQ, encodeDynamicTable, outputBufferLimit}     Config{..} = do         labelMe "H2 sender"+        -- 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 0 `E.catch` return+        loop parked 0 `E.catch` return       where         -----------------------------------------------------------------        loop :: Offset -> IO E.SomeException-        loop off = do+        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 off+                    x <- atomically $ dequeue parked off                     case x of-                        C ctl -> flushN off >> control ctl >> loop 0-                        O out -> outputAndSync out off >>= flushIfNecessary >>= loop-                        Flush -> flushN off >> loop 0+                        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.@@ -122,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          ----------------------------------------------------------------@@ -169,10 +189,10 @@         --         -- Both the stream window and the connection window are open.         -----------------------------------------------------------------        outputAndSync :: Output -> Offset -> IO Offset+        outputAndSync :: TVar [Output] -> Output -> Offset -> IO Offset         -- "handler" catches an asynchronous exception and         -- re-throws it.-        outputAndSync out@(Output strm otyp sync) off = E.handle (handler strm off) $ do+        outputAndSync parked out@(Output strm otyp sync) off = E.handle (handler strm off) $ do             state <- readStreamState strm             if isHalfClosedLocal state                 then do@@ -207,10 +227,15 @@                         return off                     _ -> do                         sws <- getStreamWindowSize strm-                        cws <- getConnectionWindowSize ctx -- not 0+                        cws <- getConnectionWindowSize ctx                         let lim = min cws sws                         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
Network/HTTP2/H2/Window.hs view
@@ -30,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@@ -68,14 +71,22 @@ ---------------------------------------------------------------- -- 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@@ -92,7 +103,7 @@ -- charge them and give them straight back. informIgnoredData :: Context -> StreamId -> Int -> IO () informIgnoredData _ _ 0 = return ()-informIgnoredData Context{controlQ, rxFlow} sid len = do+informIgnoredData ctx@Context{rxFlow} sid len = do     ok <- atomicModifyIORef' rxFlow $ checkRxLimit len     unless ok $         E.throwIO $@@ -100,6 +111,14 @@                 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]@@ -107,22 +126,37 @@ -- 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/Worker.hs view
@@ -44,7 +44,6 @@         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'}@@ -126,9 +125,10 @@     push _ [] n = return (n :: Int)     push tvar (pp : pps) n = do         T.forkManaged threadManager "H2 server push" $ do-            (newstrm, lc) <- promise pp `E.finally` increment tvar-            let Response rsp = promiseResponse pp-            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@@ -138,9 +138,12 @@     -- 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.+    -- 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-        (pid, newstrm) <- makePushStream ctx pstrm+        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 =@@ -171,12 +174,12 @@  ---------------------------------------------------------------- -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  ---------------------------------------------------------------- 
http2.cabal view
@@ -1,6 +1,6 @@ cabal-version:      2.0 name:               http2-version:            5.4.6+version:            5.4.7 license:            BSD3 license-file:       LICENSE maintainer:         Kazu Yamamoto <kazu@iij.ad.jp>
test/HPACK/DecodeSpec.hs view
@@ -4,9 +4,12 @@  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@@ -64,6 +67,40 @@                             , 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@@ -78,6 +115,33 @@                                 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.
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/HTTP2/ServerSpec.hs view
@@ -233,6 +233,44 @@                 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.@@ -254,6 +292,124 @@                                     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@@ -368,6 +524,32 @@             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@@ -418,10 +600,89 @@         -- 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@@ -433,6 +694,8 @@                 [("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) []@@ -440,6 +703,10 @@         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) []@@ -498,6 +765,9 @@ 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@@ -513,6 +783,15 @@ 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@@ -885,6 +1164,336 @@                     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)@@ -922,6 +1531,17 @@                 | 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))