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 +74/−0
- Network/QUIC/Connection/Stream.hs +36/−8
- Network/QUIC/Connection/Types.hs +3/−0
- Network/QUIC/IO.hs +60/−28
- Network/QUIC/Receiver.hs +42/−4
- Network/QUIC/Stream.hs +1/−0
- Network/QUIC/Stream/Misc.hs +5/−0
- Network/QUIC/Stream/Types.hs +3/−0
- quic.cabal +1/−1
- test/Config.hs +17/−6
- test/IOSpec.hs +242/−9
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