packages feed

quic 0.3.8 → 0.3.9

raw patch · 11 files changed

+484/−56 lines, 11 filesPVP: major bump suggested

API removals or changes: PVP suggests a major version bump

API changes (from Hackage documentation)

+ Network.QUIC.Internal: [closedStreams] :: Concurrency -> Int
+ Network.QUIC.Internal: markReleased :: Stream -> IO Bool
- Network.QUIC.Internal: Concurrency :: StreamId -> StreamIdBase -> Concurrency
+ Network.QUIC.Internal: Concurrency :: StreamId -> StreamIdBase -> Int -> Concurrency

Files

ChangeLog.md view
@@ -1,5 +1,79 @@ # ChangeLog +## 0.3.9++A stream limit the peer could walk past, two ways for a connection to+hang, and the stream states.++* Keep the peer's open streams within `initial_max_streams`.  MAX_STREAMS+  counts streams cumulatively (RFC 9000 Sec 4.6), so the limit should go+  up by one for each of the peer's streams we are done with.  It went up+  by the highest stream the peer had opened plus the initial number,+  every time one closed and the peer was near its limit -- a whole new+  window for one closed stream.  A peer that keeps its streams open could+  have any number open at once: against 0.3.7 with a limit of 64, a+  client whose server closed one stream in ten held 3132 open after three+  seconds, and on an HTTP/3 server each of those is a handler thread.+  The limit for unidirectional streams was taken from+  `initialMaxStreamsBidi` as well, and `closeStream` on a unidirectional+  stream the peer opened sent a FIN on a stream with no sending side, so+  the peer closed the connection with STREAM_STATE_ERROR and those+  streams could not be closed at all.+  [#124](https://github.com/kazu-yamamoto/quic/pull/124)++* Count what cannot be read against the anti-amplification limit.  RFC+  9000 Sec 8.1 says to count all the payload bytes received, "including+  datagrams that contain packets that are discarded"; they were counted+  only for a packet that decrypted.  A server that stops being able to+  read the peer therefore stops earning the credit it needs to answer,+  and its sender waits for good.  Reached by way of compatible version+  negotiation: the server answers in a version of its own and replaces+  its Initial keys, its first flight is lost, the client retransmits in+  the version it started with, and the server -- having sent exactly+  three times what it read -- can neither read those nor send again.  The+  handshake never finishes.+  [#128](https://github.com/kazu-yamamoto/quic/pull/128)++* Open the stream a RESET_STREAM arrives for.  RFC 9000 Sec 3.2 has the+  receiving part of a peer's stream created by the first STREAM,+  STREAM_DATA_BLOCKED or RESET_STREAM frame for it; only a STREAM frame+  created it, and a RESET_STREAM that found none was dropped.  It arrives+  first whenever the peer resets a stream it has just sent on, since the+  sender empties the queue the RESET_STREAM is on before the one the data+  is on -- not a race but the order it works in.  The stream was then+  opened by the data that came after, with nothing to say it had been+  reset and no FIN to end it, and `recvStream` waited on it until the+  idle timeout.  A stream never created is never counted, so each one+  lost this way took a unit of MAX_STREAMS credit with it for good.+  [#129](https://github.com/kazu-yamamoto/quic/pull/129)++* Set the sending part closed on STOP_SENDING, not on RESET_STREAM.  The+  two end opposite directions and were the wrong way round.  A+  RESET_STREAM from the peer closed our sending part as well as the+  receiving one, so a reply to what the peer sent before the reset could+  not go out; a STOP_SENDING left our sending part open, so `sendStream`+  went on working after we had answered with RESET_STREAM.  **This+  changes what callers see**: `sendStream` on a stream the peer stopped+  now raises `StreamIsClosed`, where it used to succeed.  `closeStream`+  and `resetStream` also end the stream for its reader whatever else they+  do -- with the sending part already closed they skipped it, and+  `recvStream` blocked for good on a stream that had left the table and+  could receive nothing more.  And a MAX_STREAMS that does not raise the+  limit is ignored (RFC 9000 Sec 4.6); one reordered or retransmitted+  after a newer one lowered the limit on the streams we may open.+  [#125](https://github.com/kazu-yamamoto/quic/pull/125)++* Tests only: the IOSpec ports moved out of the range the kernel hands+  out for a socket bound to port 0, which is 49152 to 65535 on macOS and+  32768 to 60999 on Linux.  Inside it, the relay's own socket was now and+  then given the very port the test server wanted, and "server never+  became ready" followed -- about one run in three hundred.  Each test+  also waits for its server to stop before the next one starts;+  `killThread` returns once the exception is delivered, not once the+  thread is done with it, so the server outlived it and went on serving.+  [#126](https://github.com/kazu-yamamoto/quic/pull/126),+  [#127](https://github.com/kazu-yamamoto/quic/pull/127)+ ## 0.3.8  * Tell a stream that was reset from one that ended.  After a
Network/QUIC/Connection/Stream.hs view
@@ -58,8 +58,18 @@ setTxUniMaxStreams :: Connection -> Int -> IO () setTxUniMaxStreams Connection{..} = set myUniStreamId +-- | Raising the limit on the streams I may open.+--+-- A MAX_STREAMS that does not raise it is ignored (RFC 9000, section 4.6).+-- One of them reordered on the way, or retransmitted after a newer one has+-- gone out, would otherwise lower the limit.  'setTxMaxData' and+-- 'setTxMaxStreamData' guard the limits on data the same way. set :: TVar Concurrency -> Int -> IO ()-set tvar mx = atomically $ modifyTVar tvar $ \c -> c{maxStreams = StreamIdBase mx}+set tvar mx = atomically $ modifyTVar' tvar raise+  where+    raise conc@Concurrency{..}+        | fromStreamIdBase maxStreams < mx = conc{maxStreams = StreamIdBase mx}+        | otherwise = conc  updatePeerStreamId :: Connection -> StreamId -> IO () updatePeerStreamId conn sid = do@@ -101,6 +111,21 @@         3 -> readTVarIO myUniStreamId         _ -> E.throwIO MustNotReached +-- | Counting one of the peer's streams as done with, and answering the new+--   limit for MAX_STREAMS if it is time to announce one.+--+-- The limit is the initial one plus the number of the peer's streams we are+-- done with, so that no more than the initial number are open at once+-- (RFC 9000, section 4.6).  It is announced once it has gone up by half the+-- initial number, rather than for every stream.+--+-- It used to be the highest stream the peer had opened plus the initial+-- number, whenever a stream was closed and the peer was close to its limit.+-- A peer keeping its streams open was then given a whole new window every+-- time we closed one of them, and could keep any number open: one closing a+-- tenth of what it opened held over three thousand open within seconds,+-- against a limit of 64.  The initial number was taken from the limit for+-- bidirectional streams for unidirectional ones too. checkStreamIdRoom :: Connection -> Direction -> IO (Maybe Int) checkStreamIdRoom conn dir = do     let ref@@ -108,12 +133,15 @@             | otherwise = peerUniStreamId conn     atomicModifyIORef' ref checkConc   where+    params = getMyParameters conn+    initialStreams+        | dir == Bidirectional = initialMaxStreamsBidi params+        | otherwise = initialMaxStreamsUni params     checkConc conc@Concurrency{..} =         let StreamIdBase base = maxStreams-            initialStreams = initialMaxStreamsBidi $ getMyParameters conn-            cbase = currentStream !>>. 2-         in if base - cbase < (initialStreams !>>. 3)-                then-                    let base' = cbase + initialStreams-                     in (conc{maxStreams = StreamIdBase base'}, Just base')-                else (conc, Nothing)+            closed = closedStreams + 1+            base' = initialStreams + closed+            conc' = conc{closedStreams = closed}+         in if base' - base >= max 1 (initialStreams `div` 2)+                then (conc'{maxStreams = StreamIdBase base'}, Just base')+                else (conc', Nothing)
Network/QUIC/Connection/Types.hs view
@@ -190,6 +190,8 @@ data Concurrency = Concurrency     { currentStream :: StreamId     , maxStreams :: StreamIdBase+    , closedStreams :: Int+    -- ^ For the peer's streams: how many of them we are done with     }     deriving (Show) @@ -201,6 +203,7 @@         | rl == Client = if bidi then 0 else 2         | otherwise = if bidi then 1 else 3     maxStreams = StreamIdBase n+    closedStreams = 0  -- | The peer-initiated streams of one type that have been opened so far. --
Network/QUIC/IO.hs view
@@ -142,31 +142,55 @@ closeStream :: Stream -> IO () closeStream s = do     let conn = streamConnection s-    let sid = streamId s+        sid = streamId s+        -- A unidirectional stream of the peer's has no sending side here.+        -- Sending a FIN on it anyway made the peer close the connection with+        -- STREAM_STATE_ERROR, so that one could not be closed at all.+        receiveOnly =+            (isClient conn && isServerInitiatedUnidirectional sid)+                || (isServer conn && isClientInitiatedUnidirectional sid)     closed <- isConnectionClosed conn     sclosed <- isTxStreamClosed s-    unless (closed || sclosed) $ do+    unless (receiveOnly || closed || sclosed) $ do         setTxStreamClosed s-        setRxStreamClosed s         putSendStreamQ conn $ TxStreamData s [] 0 True         waitFinTx s+    -- Outside the guard above: whether or not the FIN goes out, we are done+    -- with the stream, and 'recvStream' has to say so rather than block.+    setRxStreamClosed s     delStream conn s-    when-        ( (isClient conn && isServerInitiatedBidirectional sid)-            || (isServer conn && isClientInitiatedBidirectional sid)-        )-        $ do-            -- FLOW CONTROL: MAX_STREAMS: recv: announcing my limit properly-            checkMaxStreams conn Bidirectional-    when-        ( (isClient conn && isServerInitiatedUnidirectional sid)-            || (isServer conn && isClientInitiatedUnidirectional sid)-        )-        $ do-            -- FLOW CONTROL: MAX_STREAMS: recv: announcing my limit properly-            checkMaxStreams conn Unidirectional+    releaseStream s++-- | Counting a stream of the peer's as done with, once, and announcing a+--   new MAX_STREAMS if that is due.+--+-- Only when the application is done with it: closing it or resetting it.  A+-- RESET_STREAM from the peer takes the stream out of the table, but whatever+-- is serving it may be at work yet, and counting it then would let a peer+-- that opens streams and resets them straight away have any number served+-- at once.+releaseStream :: Stream -> IO ()+releaseStream s = do+    first <- markReleased s+    when first $ do+        when+            ( (isClient conn && isServerInitiatedBidirectional sid)+                || (isServer conn && isClientInitiatedBidirectional sid)+            )+            $ do+                -- FLOW CONTROL: MAX_STREAMS: recv: announcing my limit properly+                checkMaxStreams Bidirectional+        when+            ( (isClient conn && isServerInitiatedUnidirectional sid)+                || (isServer conn && isClientInitiatedUnidirectional sid)+            )+            $ do+                -- FLOW CONTROL: MAX_STREAMS: recv: announcing my limit properly+                checkMaxStreams Unidirectional   where-    checkMaxStreams conn dir = do+    conn = streamConnection s+    sid = streamId s+    checkMaxStreams dir = do         mx <- checkStreamIdRoom conn dir         case mx of             Nothing -> return ()@@ -219,11 +243,16 @@     unless sclosed $ do         finalSize <- getTxStreamFinalSize s         setTxStreamClosed s-        setRxStreamClosed s         lvl <- getEncryptionLevel conn         let frame = ResetStream sid aerr finalSize         putOutput conn $ OutControl lvl [frame]+    -- Outside the guard above: whether or not the RESET_STREAM goes out, we+    -- are done with the stream, and 'recvStream' has to say so rather than+    -- block.  The peer's STOP_SENDING has closed the sending part already+    -- when an application resets a stream because 'sendStream' failed.+    setRxStreamClosed s     delStream conn s+    releaseStream s  -- | Asking the peer to stop sending. --   This sends STOP_SENDING to the peer@@ -248,19 +277,22 @@     -- Determine the send level from connection readiness state rather than     -- the encryptionLevel TVar, which is not set to RTT0Level during 0-RTT.     ready1rtt <- isConnection1RTTReady conn-    lvl <- if ready1rtt-        then return RTT1Level-        else do-            ready0rtt <- isConnection0RTTReady conn-            if ready0rtt-                then return RTT0Level-                else E.throwIO $ ConnectionIsClosed "Cannot send DATAGRAM"+    lvl <-+        if ready1rtt+            then return RTT1Level+            else do+                ready0rtt <- isConnection0RTTReady conn+                if ready0rtt+                    then return RTT0Level+                    else E.throwIO $ ConnectionIsClosed "Cannot send DATAGRAM"     limitBytes <- maxDatagramFrameSize <$> getPeerParameters conn     when (limitBytes == 0) $-      E.throwIO $ ConnectionIsClosed "DATAGRAM not supported by peer"+        E.throwIO $+            ConnectionIsClosed "DATAGRAM not supported by peer"     let frameOverhead = 1 + BS.length (encodeInt (fromIntegral $ BS.length dat))     when (BS.length dat + frameOverhead > limitBytes) $-        E.throwIO $ ConnectionIsClosed "DATAGRAM size violation"+        E.throwIO $+            ConnectionIsClosed "DATAGRAM size violation"     let frame = Datagram False dat     putOutput conn $ OutControl lvl [frame] 
Network/QUIC/Receiver.hs view
@@ -156,12 +156,22 @@     let CryptPacket hdr crypt = rpCryptPacket rpkt         lvl = rpEncryptionLevel rpkt         tim = rpTimeRecevied rpkt+    -- Before the decryption below.  RFC 9000, section 8.1: "servers MUST+    -- count all of the payload bytes received in datagrams that are uniquely+    -- attributed to a single connection.  This includes datagrams that+    -- contain packets that are discarded."+    --+    -- Counted only for what it could read, a server that stops being able to+    -- read the peer stops earning the credit it needs to answer, and the+    -- anti-amplification limit holds its sender for good.  Whatever the peer+    -- retransmits from then on is discarded and buys nothing, so the+    -- connection does not recover; it hangs until the idle timeout.+    pathInfo <- getPathInfo conn+    addPathRxBytes pathInfo $ rpReceivedBytes rpkt     mplain <- decryptCrypt conn crypt lvl     case mplain of         Just plain@Plain{..} -> do             addRxBytes conn $ rpReceivedBytes rpkt-            pathInfo <- getPathInfo conn-            addPathRxBytes pathInfo $ rpReceivedBytes rpkt             when (isIllegalReservedBits plainMarks || isNoFrames plainMarks) $                 closeConnection conn ProtocolViolation "Non 0 RR bits or no frames"             when (isUnknownFrame plainMarks) $@@ -271,7 +281,27 @@         closeConnection conn ProtocolViolation "RESET_STREAM"     when (isSendOnly conn sid) $         closeConnection conn StreamStateError "Received in a send-only stream"-    mstrm <- findStream conn sid+    updatePeerStreamId conn sid+    -- FLOW CONTROL: MAX_STREAMS: recv: rejecting if over my limit+    ok <- checkRxMaxStreams conn sid+    unless ok $ closeConnection conn StreamLimitError "stream id is too large"+    mstrm0 <- findStream conn sid+    guardStream conn sid mstrm0+    -- RFC 9000 Sec 3.2: "The receiving part of a stream initiated by a peer+    -- ... is created when the first STREAM, STREAM_DATA_BLOCKED, or+    -- RESET_STREAM frame is received for that stream."  Only a STREAM frame+    -- created it here, and a RESET_STREAM that arrived before one was+    -- dropped without a word.+    --+    -- It arrives first whenever the peer resets a stream it has just sent+    -- on: the sender empties the queue the RESET_STREAM is on before the one+    -- the data is on, so the reset overtakes it.  The peer's application had+    -- then said its piece and been ignored, and ours waited on 'recvStream'+    -- for a stream that would never end, until the idle timeout.+    --+    -- 'openStream' answers Nothing for a stream that has been closed since,+    -- so a late copy does not open one anew.+    mstrm <- maybe (openStream conn sid) (return . Just) mstrm0     onResetStreamReceived2 (connHooks conn) mstrm aerr finlen     case mstrm of         Nothing -> return ()@@ -280,7 +310,11 @@             -- Before the pseudo FIN below, so that whoever reads it can             -- tell it was not a real one.             setResetReceived strm aerr-            setTxStreamClosed strm+            -- RESET_STREAM ends the peer's sending part alone+            -- (RFC 9000, section 3.2).  Our own sending part is untouched,+            -- so that a reply to what the peer sent before the reset can+            -- still go out.  It used to be closed here as well, and+            -- 'sendStream' then threw StreamIsClosed.             setRxStreamClosed strm             delStream conn strm processFrame conn lvl (StopSending sid err) = do@@ -296,6 +330,10 @@             Just strm -> do                 finalSize <- getTxStreamFinalSize strm                 sendFrames conn lvl [ResetStream sid err finalSize]+                -- Our sending part is reset now (RFC 9000, section 3.5), so+                -- no more STREAM frames may go out on it.  Without this,+                -- 'sendStream' went on working after the RESET_STREAM.+                setTxStreamClosed strm processFrame _ _ (CryptoF _ "") = return () processFrame conn lvl (CryptoF off cdat) = do     when (lvl == RTT0Level) $
Network/QUIC/Stream.hs view
@@ -22,6 +22,7 @@     setRxStreamClosed,     resetReceived,     setResetReceived,+    markReleased,     readStreamFlowTx,     addTxStreamData,     setTxMaxStreamData,
Network/QUIC/Stream/Misc.hs view
@@ -11,6 +11,7 @@     setRxStreamClosed,     resetReceived,     setResetReceived,+    markReleased,     --     readStreamFlowTx,     addTxStreamData,@@ -82,6 +83,10 @@  setResetReceived :: Stream -> ApplicationProtocolError -> IO () setResetReceived Stream{..} aerr = writeIORef streamResetRx $ Just aerr++-- | Marking the stream as done with; 'True' only the first time.+markReleased :: Stream -> IO Bool+markReleased Stream{..} = atomicModifyIORef' streamReleased $ \done -> (True, not done)  ---------------------------------------------------------------- 
Network/QUIC/Stream/Types.hs view
@@ -43,6 +43,8 @@     , streamSyncFinTx :: MVar ()     , streamResetRx :: IORef (Maybe ApplicationProtocolError)     -- ^ The error code of a RESET_STREAM from the peer+    , streamReleased :: IORef Bool+    -- ^ Whether we are done with it and have counted it so     }  instance Show Stream where@@ -59,6 +61,7 @@     streamReass     <- newIORef (0, Skew.empty)     streamSyncFinTx <- newEmptyMVar     streamResetRx   <- newIORef Nothing+    streamReleased  <- newIORef False     return Stream{..} {- FOURMOLU_ENABLE -} 
quic.cabal view
@@ -1,6 +1,6 @@ cabal-version:      2.0 name:               quic-version:            0.3.8+version:            0.3.9 license:            BSD3 license-file:       LICENSE maintainer:         kazu@iij.ad.jp
test/Config.hs view
@@ -50,7 +50,7 @@ testServerConfig =     defaultServerConfig         { -- Don't use "0.0.0.0" and "::" for Windows (UDP dispatching bug)-          scAddresses = [("127.0.0.1", 50003)]+          scAddresses = [("127.0.0.1", 15003)]         , scParameters =             (scParameters defaultServerConfig)                 { maxIdleTimeout = Milliseconds 10000@@ -74,7 +74,7 @@ testServerConfigR =     defaultServerConfig         { -- Don't use "0.0.0.0" and "::" for Windows (UDP dispatching bug)-          scAddresses = [("127.0.0.1", 50003)]+          scAddresses = [("127.0.0.1", 15003)]         , scParameters =             (scParameters defaultServerConfig)                 { maxIdleTimeout = Milliseconds 10000@@ -85,7 +85,7 @@ testClientConfig =     defaultClientConfig         { ccServerName = "127.0.0.1"-        , ccPortName = "50003"+        , ccPortName = "15003"         , ccValidate = False         , ccDebugLog = True         , ccParameters =@@ -98,7 +98,7 @@ testClientConfigR =     defaultClientConfig         { ccServerName = "127.0.0.1"-        , ccPortName = "50002"+        , ccPortName = "15002"         , ccValidate = False         , ccDebugLog = True         , ccParameters =@@ -174,9 +174,20 @@  withPipeWith :: Bool -> Scenario -> IO () -> IO () withPipeWith stray scenario body = do-    addrC <- resolve "50002"+    -- Ports outside the range the kernel hands out for a socket bound to+    -- port 0, which is 49152 to 65535 on macOS and 32768 to 60999 on Linux.+    -- Inside it, the relay's own socket for the server side, bound to+    -- 127.0.0.1:0 below, is now and then given the very port the test+    -- server wants, and the server cannot bind it:+    --+    --   Network.Socket.bind: resource busy (Address already in use)+    --+    -- 'run' reports that to whoever forked it, which is nobody, so+    -- onServerReady never fires and the test fails with "server never+    -- became ready".  One run in a few hundred, on ports 50002 and 50003.+    addrC <- resolve "15002"     let saC = addrAddress addrC-    addrS <- resolve "50003"+    addrS <- resolve "15003"     let saS = addrAddress addrS     irefC <- newIORef 0     irefS <- newIORef 0
test/IOSpec.hs view
@@ -7,6 +7,7 @@ import qualified Control.Exception as E import Control.Monad import qualified Data.ByteString as BS+import Data.IORef import qualified System.Timeout as Timeout import Test.Hspec @@ -37,6 +38,9 @@             Timeout.timeout 5000000 (takeMVar var) >>= \r -> case r of                 Just () -> return ()                 Nothing -> expectationFailure "server never became ready"+    describe "handshake" $ do+        it "can exchange data when the server's first flight is lost" $ do+            withPipe (DropServerPacket [0, 1, 2]) $ testSendRecv cc sc waitS 20     describe "send & recv" $ do         it "can exchange data on random dropping" $ do             withPipe (Randomly 20) $ testSendRecv cc sc waitS 1000@@ -107,9 +111,26 @@             withPipe (DelayServerPacket 300) $ testLateCopy False cc sc waitS         it "ignores a late copy of the data it received on a stream the peer opened" $ do             withPipe (DelayClientPacket 300) $ testLateCopy True cc sc waitS+    describe "stream limits" $ do+        it "keeps the peer's open bidirectional streams within the limit" $ do+            withPipe (DropClientPacket []) $ testOpenStreams cc sc waitS+        it "raises the limit on unidirectional streams by their own number" $ do+            withPipe (DropClientPacket []) $ testUniStreams cc sc waitS+    describe "stream state" $ do+        it "can still send after the peer's RESET_STREAM" $ do+            withPipe (DropClientPacket []) $ testSendAfterReset cc sc waitS+        it "cannot send after the peer's STOP_SENDING" $ do+            withPipe (DropClientPacket []) $ testStopSending cc sc waitS+        it "ends the stream for the reader when it is reset" $ do+            withPipe (DropClientPacket []) $ testRecvAfterReset cc sc waitS+        it "accepts a stream the peer only reset" $ do+            withPipe (DropClientPacket []) $ testResetOnly cc sc waitS+        it "accepts a reset that overtook the data it followed" $ do+            withPipe (DropClientPacket []) $ testResetOvertakes cc sc waitS     describe "port handover" $ do         it "ignores a leftover datagram from the connection that just closed" $-            withPipeStray (Randomly 20) $ testSendRecv cc sc waitS 20+            withPipeStray (Randomly 20) $+                testSendRecv cc sc waitS 20     describe "concurrency" $ do         it "can handle multiple clients" $ do             withPipe (Randomly 20) $ testMultiSendRecv cc sc waitS 500@@ -138,7 +159,7 @@         payload = BS.replicate 1234 0         hooks = (ccHooks cc0){onResetStreamReceived2 = record finalSizeVar}         cc = cc0{ccHooks = hooks}-    E.bracket (forkIO $ server request payload doneVar) killThread $ \_ ->+    withAsync (server request payload doneVar) $ \_ ->         client cc request payload finalSizeVar doneVar   where     aerr = ApplicationProtocolError 0@@ -170,7 +191,7 @@ --   application a stream that was never opened. testLateCopy :: Bool -> C.ClientConfig -> ServerConfig -> IO () -> IO () testLateCopy upload cc sc waitS =-    E.bracket (forkIO server) killThread $ \_ -> client+    withAsync server $ \_ -> client   where     (upLen, downLen)         | upload = (1000000, 10)@@ -204,7 +225,7 @@ testResetReceived :: C.ClientConfig -> ServerConfig -> IO () -> IO () testResetReceived cc sc waitS = do     resultVar <- newEmptyMVar-    E.bracket (forkIO $ server resultVar) killThread $ \_ -> client resultVar+    withAsync (server resultVar) $ \_ -> client resultVar   where     aerr = ApplicationProtocolError 7 @@ -235,11 +256,71 @@         putMVar resultVar (r0, r1)         threadDelay 1000000 +-- | The server closes one stream in ten and keeps the others open; the+-- client opens streams for as long as it is let.  No more than+-- initial_max_streams_bidi may be open at once (RFC 9000, section 4.6).+--+-- The limit used to be raised to the highest stream opened plus the initial+-- number whenever a stream was closed, so that the client was given a whole+-- new window each time: over three thousand were held open within seconds.+testOpenStreams :: C.ClientConfig -> ServerConfig -> IO () -> IO ()+testOpenStreams cc sc waitS = do+    maxOpen <- newIORef (0 :: Int)+    withAsync (server maxOpen) $ \_ -> do+        client+        readIORef maxOpen+            >>= (`shouldSatisfy` (<= initialMaxStreamsBidi (scParameters sc)))+  where+    client = do+        waitS+        C.run cc $ \conn -> do+            _ <- Timeout.timeout 1500000 $ forever $ do+                strm <- stream conn+                sendStream strm "x"+            return ()+    server maxOpen = run sc $ \conn -> do+        openRef <- newIORef (0 :: Int)+        forM_ [0 :: Int ..] $ \n -> do+            strm <- acceptStream conn+            _ <- recvStream strm 1+            if n `mod` 10 == 0+                then closeStream strm+                else do+                    o <- atomicModifyIORef' openRef $ \x -> (x + 1, x + 1)+                    atomicModifyIORef' maxOpen $ \m -> (max m o, ())++-- | The server closes the first unidirectional stream and keeps the rest.+-- The client may then have initial_max_streams_uni plus one.  The limit was+-- raised by initial_max_streams_bidi instead, 64 against 3.+testUniStreams :: C.ClientConfig -> ServerConfig -> IO () -> IO ()+testUniStreams cc sc waitS = do+    opened <- newIORef (0 :: Int)+    withAsync server $ \_ -> do+        client opened+        readIORef opened+            >>= (`shouldSatisfy` (<= initialMaxStreamsUni (scParameters sc) + 1))+  where+    client opened = do+        waitS+        C.run cc $ \conn -> do+            _ <- Timeout.timeout 1500000 $ forever $ do+                strm <- unidirectionalStream conn+                sendStream strm "x"+                modifyIORef' opened (+ 1)+            return ()+    server = run sc $ \conn -> do+        strm0 <- acceptStream conn+        _ <- recvStream strm0 1+        closeStream strm0+        forever $ do+            strm <- acceptStream conn+            recvStream strm 1+ testRecvStreamClientStopFirst     :: C.ClientConfig -> ServerConfig -> IO () -> IO () testRecvStreamClientStopFirst cc sc waitS = do     mvar <- newEmptyMVar-    E.bracket (forkIO $ server mvar) killThread $ \_ -> client mvar+    withAsync (server mvar) $ \_ -> client mvar   where     aerr = ApplicationProtocolError 0 @@ -265,7 +346,7 @@     :: C.ClientConfig -> ServerConfig -> IO () -> IO () testRecvStreamServerStopFirst cc sc waitS = do     mvar <- newEmptyMVar-    E.bracket (forkIO $ server mvar) killThread $ \_ -> client mvar+    withAsync (server mvar) $ \_ -> client mvar   where     aerr = ApplicationProtocolError 0 @@ -288,7 +369,7 @@ testSendRecv :: C.ClientConfig -> ServerConfig -> IO () -> Int -> IO () testSendRecv cc sc waitS times = do     mvar <- newEmptyMVar-    E.bracket (forkIO $ server mvar) killThread $ \_ -> client mvar+    withAsync (server mvar) $ \_ -> client mvar   where     client mvar = do         waitS@@ -307,7 +388,7 @@ testMultiSendRecv :: C.ClientConfig -> ServerConfig -> IO () -> Int -> IO () testMultiSendRecv cc sc waitS times = do     mvars <- replicateM concurrency newEmptyMVar-    E.bracket (forkIO $ server mvars) killThread $ \_ -> client mvars+    withAsync (server mvars) $ \_ -> client mvars   where     concurrency = 10     chunklen = 12345@@ -344,7 +425,7 @@  testAbort :: C.ClientConfig -> ServerConfig -> IO () -> IO () testAbort cc sc waitS = do-    E.bracket (forkIO server) killThread $ \_ ->+    withAsync server $ \_ ->         client `shouldThrow` appErr   where     client = do@@ -364,3 +445,155 @@                         (ApplicationProtocolError 1)                         "testing abortConnection"             loop conn++-- | RESET_STREAM ends the peer's sending part alone (RFC 9000, section 3.2).+--   What we have to send on a bidirectional stream must still go out.+testSendAfterReset :: C.ClientConfig -> ServerConfig -> IO () -> IO ()+testSendAfterReset cc sc waitS = do+    res <- newEmptyMVar+    withAsync (server res) $ \_ -> do+        client+        takeMVar res >>= (`shouldBe` True)+  where+    client = do+        waitS+        C.run cc $ \conn -> do+            strm <- stream conn+            sendStream strm "ping"+            -- The reply says the "ping" has arrived, so that the+            -- RESET_STREAM cannot overtake it and be dropped as a frame+            -- for a stream that does not exist yet.+            _ <- recvStream strm 1024+            resetStream strm (ApplicationProtocolError 0)+            threadDelay 300000+    server res = run sc $ \conn -> do+        strm <- acceptStream conn+        _ <- recvStream strm 1024+        sendStream strm "ack"+        -- An empty ByteString: the pseudo FIN that the RESET_STREAM left.+        let recvEOF = do+                bs <- recvStream strm 1024+                unless (bs == "") recvEOF+        recvEOF+        sent <-+            (True <$ sendStream strm "pong") `E.catch` \e -> case e of+                StreamIsClosed -> return False+                _ -> E.throwIO e+        putMVar res sent++-- | STOP_SENDING makes us reset our sending part (RFC 9000, section 3.5),+--   so nothing more may go out on it.+testStopSending :: C.ClientConfig -> ServerConfig -> IO () -> IO ()+testStopSending cc sc waitS =+    withAsync server $ \_ -> client+  where+    client = do+        waitS+        C.run cc $ \conn -> do+            strm <- stream conn+            sendStream strm "ping"+            stopped <-+                ( do+                    _ <- Timeout.timeout 2000000 $ forever $ do+                        threadDelay 20000+                        sendStream strm "x"+                    return False+                )+                    `E.catch` \e -> case e of+                        StreamIsClosed -> return True+                        _ -> E.throwIO e+            stopped `shouldBe` True+    server = run sc $ \conn -> do+        strm <- acceptStream conn+        _ <- recvStream strm 1024+        stopStream strm (ApplicationProtocolError 0)+        threadDelay 3000000++-- | 'resetStream' leaves nothing to wait for.  An application that resets a+--   stream because 'sendStream' failed does so with the sending part closed+--   by the peer's STOP_SENDING already, and 'recvStream' has to end there+--   rather than block: what the peer sends afterwards is dropped anyway,+--   the stream being out of the table.+testRecvAfterReset :: C.ClientConfig -> ServerConfig -> IO () -> IO ()+testRecvAfterReset cc sc waitS =+    withAsync server $ \_ -> client+  where+    client = do+        waitS+        C.run cc $ \conn -> do+            strm <- stream conn+            sendStream strm "ping"+            _ <-+                ( do+                    _ <- Timeout.timeout 2000000 $ forever $ do+                        threadDelay 20000+                        sendStream strm "x"+                    return ()+                )+                    `E.catch` \e -> case e of+                        StreamIsClosed -> return ()+                        _ -> E.throwIO e+            resetStream strm (ApplicationProtocolError 0)+            mbs <- Timeout.timeout 1000000 $ recvStream strm 1024+            mbs `shouldBe` Just ""+    server = run sc $ \conn -> do+        strm <- acceptStream conn+        _ <- recvStream strm 1024+        stopStream strm (ApplicationProtocolError 0)+        threadDelay 3000000++-- | RFC 9000, section 3.2: the receiving part of a stream the peer opened is+--   created when the first STREAM, STREAM_DATA_BLOCKED or RESET_STREAM frame+--   for it arrives.  Here the client never sends on the stream, so the+--   RESET_STREAM is the first the server hears of it.+testResetOnly :: C.ClientConfig -> ServerConfig -> IO () -> IO ()+testResetOnly cc sc waitS = do+    got <- newEmptyMVar+    withAsync (server got) $ \_ -> do+        client+        -- With a timeout: without the stream, 'acceptStream' on the server+        -- never returns, and the test would hang rather than fail.+        r <- Timeout.timeout 1000000 $ takeMVar got+        r `shouldBe` Just (Just aerr)+  where+    aerr = ApplicationProtocolError 7+    client = do+        waitS+        C.run cc $ \conn -> do+            strm <- stream conn+            resetStream strm aerr+            threadDelay 300000+    server got = run sc $ \conn -> do+        strm <- acceptStream conn+        bs <- recvStream strm 1024+        bs `shouldBe` ""+        resetReceived strm >>= putMVar got++-- | The same, for the reset that follows data on the stream.  It overtakes+--   the data: the sender empties the queue the RESET_STREAM is on before the+--   one the STREAM frame is on.+testResetOvertakes :: C.ClientConfig -> ServerConfig -> IO () -> IO ()+testResetOvertakes cc sc waitS = do+    got <- newEmptyMVar+    withAsync (server got) $ \_ -> do+        client+        -- With a timeout: without the stream, 'acceptStream' on the server+        -- never returns, and the test would hang rather than fail.+        r <- Timeout.timeout 1000000 $ takeMVar got+        r `shouldBe` Just (Just aerr)+  where+    aerr = ApplicationProtocolError 7+    client = do+        waitS+        C.run cc $ \conn -> do+            strm <- stream conn+            sendStream strm "ping"+            resetStream strm aerr+            threadDelay 300000+    server got = run sc $ \conn -> do+        strm <- acceptStream conn+        let recvEOF = do+                bs <- recvStream strm 1024+                unless (bs == "") recvEOF+        recvEOF+        resetReceived strm >>= putMVar got