quic 0.3.6 → 0.3.7
raw patch · 11 files changed
+325/−73 lines, 11 filesPVP: major bump suggested
API removals or changes: PVP suggests a major version bump
API changes (from Hackage documentation)
+ Network.QUIC.Internal: OpenedStreams :: StreamId -> IntSet -> OpenedStreams
+ Network.QUIC.Internal: [openedNext] :: OpenedStreams -> StreamId
+ Network.QUIC.Internal: [openedSkipped] :: OpenedStreams -> IntSet
+ Network.QUIC.Internal: [peerOpened] :: Connection -> IORef OpenedStreams
+ Network.QUIC.Internal: [peerUniOpened] :: Connection -> IORef OpenedStreams
+ Network.QUIC.Internal: claimPeerStream :: Connection -> StreamId -> IO Bool
+ Network.QUIC.Internal: data OpenedStreams
+ Network.QUIC.Internal: newOpenedStreams :: Concurrency -> OpenedStreams
- Network.QUIC.Internal: Connection :: ConnState -> DebugLogger -> QLogger -> Hooks -> ~Send -> ~Recv -> RecvQ -> DatagramQ -> IORef Socket -> (CID -> StatelessResetToken) -> IORef (Map Word64 (Weak ThreadId)) -> ThreadId -> Rate -> IORef RoleInfo -> IORef VersionInfo -> VersionInfo -> Parameters -> IORef CIDDB -> IORef Parameters -> TVar CIDDB -> IORef PeerInfo -> InputQ -> CryptoQ -> OutputQ -> Rate -> Shared -> IORef Int -> IORef (IO ()) -> IORef PacketNumber -> IORef StreamTable -> TVar Concurrency -> TVar Concurrency -> IORef Concurrency -> IORef Concurrency -> TVar TxFlow -> IORef RxFlow -> TVar MigrationState -> IORef Bool -> IORef Microseconds -> IORef Int -> IORef Int -> TVar Bool -> Array EncryptionLevel (TVar [ReceivedPacket]) -> IOArray EncryptionLevel Cipher -> IOArray EncryptionLevel Coder -> IOArray Bool Coder1RTT -> IOArray EncryptionLevel Protector -> IORef (Bool, PacketNumber) -> IORef Negotiated -> IORef AuthCIDs -> IORef AuthCIDs -> Buffer -> SizedBuffer -> Buffer -> IORef (IO ()) -> LDCC -> Connection
+ Network.QUIC.Internal: Connection :: ConnState -> DebugLogger -> QLogger -> Hooks -> ~Send -> ~Recv -> RecvQ -> DatagramQ -> IORef Socket -> (CID -> StatelessResetToken) -> IORef (Map Word64 (Weak ThreadId)) -> ThreadId -> Rate -> IORef RoleInfo -> IORef VersionInfo -> VersionInfo -> Parameters -> IORef CIDDB -> IORef Parameters -> TVar CIDDB -> IORef PeerInfo -> InputQ -> CryptoQ -> OutputQ -> Rate -> Shared -> IORef Int -> IORef (IO ()) -> IORef PacketNumber -> IORef StreamTable -> TVar Concurrency -> TVar Concurrency -> IORef Concurrency -> IORef Concurrency -> IORef OpenedStreams -> IORef OpenedStreams -> TVar TxFlow -> IORef RxFlow -> TVar MigrationState -> IORef Bool -> IORef Microseconds -> IORef Int -> IORef Int -> TVar Bool -> Array EncryptionLevel (TVar [ReceivedPacket]) -> IOArray EncryptionLevel Cipher -> IOArray EncryptionLevel Coder -> IOArray Bool Coder1RTT -> IOArray EncryptionLevel Protector -> IORef (Bool, PacketNumber) -> IORef Negotiated -> IORef AuthCIDs -> IORef AuthCIDs -> Buffer -> SizedBuffer -> Buffer -> IORef (IO ()) -> LDCC -> Connection
Files
- ChangeLog.md +38/−0
- Network/QUIC/Connection/Crypto.hs +0/−3
- Network/QUIC/Connection/StreamTable.hs +23/−0
- Network/QUIC/Connection/Types.hs +23/−0
- Network/QUIC/Receiver.hs +64/−44
- Network/QUIC/Server/Reader.hs +13/−2
- Network/QUIC/Server/Run.hs +11/−6
- quic.cabal +9/−1
- test/Config.hs +87/−15
- test/HandshakeSpec.hs +8/−1
- test/IOSpec.hs +49/−1
ChangeLog.md view
@@ -1,5 +1,43 @@ # ChangeLog +## 0.3.7++* Don't open a closed stream again for a late copy of its data. A STREAM+ frame for a stream no longer in the table opened it anew, and after the+ stream was closed what arrives is a copy of data already received, sent+ again because the packet carrying it was taken for lost while it was+ only late. The new stream starts from the initial window, 256K, so a+ copy from past that point was called a flow control error and the+ connection closed with FLOW_CONTROL_ERROR; a copy from within the+ window was worse, handing the application a closed stream as a new one.+ [#118](https://github.com/kazu-yamamoto/quic/pull/118)++* Report a server that could not be started. `run` and `runWithSockets`+ created their sockets inside a handler whose logger discards what it is+ given, so a failure to bind was swallowed, `onServerReady` was never+ reached, and `run` returned as if all were well. **This changes what+ callers see**: a `run` that cannot bind now raises where it used to+ return quietly.+ [#119](https://github.com/kazu-yamamoto/quic/pull/119)++* Don't let one datagram take the server's dispatcher down for good. An+ exception anywhere in the dispatcher loop ended it, and nothing+ restarts it -- the socket stays bound, so the server went on looking+ like a server and answering nothing, with the failure logged to the+ same discarding logger. Decode and dispatch are guarded per datagram+ now. No input was found that raises; this is the blast radius being+ closed, not a known hole.+ [#120](https://github.com/kazu-yamamoto/quic/pull/120)++* Tests only: the IOSpec relay no longer connects its sockets, since a+ connected UDP socket turns an ICMP port-unreachable into ECONNREFUSED+ on the next operation and the peers there come and go with every test.+ And qlog is behind a `qlog` flag, off by default, which makes its own+ directory when it is on -- the directory used to be the CI's job, so a+ fresh clone failed 33 examples with `openFile: does not exist`.+ [#117](https://github.com/kazu-yamamoto/quic/pull/117),+ [#121](https://github.com/kazu-yamamoto/quic/pull/121)+ ## 0.3.6 A security fix and three for a stalled handshake.
Network/QUIC/Connection/Crypto.hs view
@@ -25,7 +25,6 @@ ) where import Control.Concurrent.STM-import Network.TLS.Extra.Cipher import Network.TLS.QUIC import Network.QUIC.Connection.Misc@@ -141,7 +140,6 @@ Secret sN1 = nextSecret ver cipher $ Secret sN secN1 = (ClientTrafficSecret cN1, ServerTrafficSecret sN1) - genNiteCoder :: Bool -> Version -> Cipher -> TrafficSecrets a -> IO (Coder, Protector) genNiteCoder cli ver cipher (ClientTrafficSecret c, ServerTrafficSecret s) = do@@ -185,7 +183,6 @@ rxPayloadIV = initialVector ver cipher rxSecret rxHeaderKey = headerProtectionKey ver cipher rxSecret unp = protectionMask cipher rxHeaderKey- genNiteCoder1RTT :: Bool -> Version -> Cipher -> TrafficSecrets a -> Coder -> IO Coder
Network/QUIC/Connection/StreamTable.hs view
@@ -2,6 +2,7 @@ module Network.QUIC.Connection.StreamTable ( createStream,+ claimPeerStream, findStream, addStream, delStream,@@ -11,6 +12,8 @@ getCryptoStream, ) where +import qualified Data.IntSet as IntSet+ import Network.QUIC.Connection.Misc import Network.QUIC.Connection.Queue import Network.QUIC.Connection.Types@@ -25,6 +28,26 @@ strm <- addStream conn sid putInput conn $ InpStream strm return strm++-- | Recording that a peer-initiated stream is opened. 'False' if it was+-- opened before, which, for one no longer in the stream table, means that it+-- has been closed.+--+-- The skipped ids never outnumber the streams we allow the peer, as a frame+-- for a stream past that limit is refused before it gets here.+claimPeerStream :: Connection -> StreamId -> IO Bool+claimPeerStream conn sid = atomicModifyIORef' ref claim+ where+ ref+ | isUnidirectional sid = peerUniOpened conn+ | otherwise = peerOpened conn+ claim os@OpenedStreams{..}+ | sid >= openedNext =+ let skipped = IntSet.fromDistinctAscList [openedNext, openedNext + 4 .. sid - 4]+ in (OpenedStreams (sid + 4) (openedSkipped <> skipped), True)+ | sid `IntSet.member` openedSkipped =+ (os{openedSkipped = IntSet.delete sid openedSkipped}, True)+ | otherwise = (os, False) findStream :: Connection -> StreamId -> IO (Maybe Stream) findStream Connection{..} sid = lookupStream sid <$> readIORef streamTable
Network/QUIC/Connection/Types.hs view
@@ -12,6 +12,8 @@ import Data.ByteString.Internal import Data.IntMap.Strict (IntMap) import qualified Data.IntMap.Strict as IntMap+import Data.IntSet (IntSet)+import qualified Data.IntSet as IntSet import Data.Map.Strict (Map) import qualified Data.Map.Strict as Map import Data.X509 (CertificateChain)@@ -200,6 +202,23 @@ | otherwise = if bidi then 1 else 3 maxStreams = StreamIdBase n +-- | The peer-initiated streams of one type that have been opened so far.+--+-- Every id below 'openedNext' has been opened, save those in+-- 'openedSkipped': the peer used a higher-numbered stream first, which opens+-- these too (RFC 9000 Sec 3.2), but nothing has arrived for them yet. A+-- stream that has been opened and is no longer in the stream table has been+-- closed, and a frame for it -- typically a retransmission that crossed our+-- acknowledgement -- must not open it again.+data OpenedStreams = OpenedStreams+ { openedNext :: StreamId+ , openedSkipped :: IntSet+ }+ deriving (Show)++newOpenedStreams :: Concurrency -> OpenedStreams+newOpenedStreams conc = OpenedStreams (currentStream conc) IntSet.empty+ ---------------------------------------------------------------- type Send = Buffer -> Int -> IO ()@@ -272,6 +291,8 @@ , myUniStreamId :: TVar Concurrency -- C:2 S:3 , peerStreamId :: IORef Concurrency -- C:1 S:0 , peerUniStreamId :: IORef Concurrency -- C:3 S:2+ , peerOpened :: IORef OpenedStreams -- C:1 S:0+ , peerUniOpened :: IORef OpenedStreams -- C:3 S:2 , flowTx :: TVar TxFlow , flowRx :: IORef RxFlow , migrationState :: TVar MigrationState@@ -370,6 +391,8 @@ myUniStreamId <- newTVarIO (newConcurrency rl Unidirectional 0) peerStreamId <- newIORef peerConcurrency peerUniStreamId <- newIORef peerUniConcurrency+ peerOpened <- newIORef (newOpenedStreams peerConcurrency)+ peerUniOpened <- newIORef (newOpenedStreams peerUniConcurrency) flowTx <- newTVarIO (newTxFlow 0) -- limit is set in Handshake flowRx <- newIORef (newRxFlow $ initialMaxData myParameters) migrationState <- newTVarIO NonMigration
Network/QUIC/Receiver.hs view
@@ -234,6 +234,16 @@ closeConnection conn StreamStateError emsg streamNotCreatedYet _ _ _ = return () +-- | Opening a stream for a STREAM frame that found none, unless the stream+-- was open once and has been closed since. 'guardStream' has already+-- refused one of ours that we have not created yet.+openStream :: Connection -> StreamId -> IO (Maybe Stream)+openStream conn sid+ | isInitiated conn sid = return Nothing+ | otherwise = do+ new <- claimPeerStream conn sid+ if new then Just <$> createStream conn sid else return Nothing+ processFrame :: Connection -> EncryptionLevel -> Frame -> IO () processFrame _ _ Padding{} = return () processFrame conn lvl Ping = do@@ -316,28 +326,33 @@ closeConnection conn StreamStateError "send-only stream" mstrm <- findStream conn sid guardStream conn sid mstrm- strm <- maybe (createStream conn sid) return mstrm- let len = BS.length dat- rx = RxStreamData dat off len fin- fc <- putRxStreamData strm rx- case fc of- -- FLOW CONTROL: MAX_STREAM_DATA: recv: rejecting if over my limit- OverLimit ->- closeConnection conn FlowControlError "Flow control error for stream in 0-RTT"- -- Not a flow control error: the peer is inside its window, it is- -- just spending it in more pieces than we will hold. Rate control- -- answers with InternalError too.- TooFragmented ->- closeConnection conn QUIC.InternalError "Too many stream fragments"- Duplicated -> return ()- Reassembled -> do- ok' <- checkRxMaxData conn len- -- FLOW CONTROL: MAX_DATA: send: respecting peer's limit- unless ok' $- closeConnection- conn- FlowControlError- "Flow control error for connection in 0-RTT"+ -- Nothing when the stream has been closed. What arrives for it now+ -- is a copy of something already received, sent again because our+ -- acknowledgement crossed it. Opening the stream anew would hold+ -- the copy to the initial window and call it a flow control error.+ mstrm' <- maybe (openStream conn sid) (return . Just) mstrm+ forM_ mstrm' $ \strm -> do+ let len = BS.length dat+ rx = RxStreamData dat off len fin+ fc <- putRxStreamData strm rx+ case fc of+ -- FLOW CONTROL: MAX_STREAM_DATA: recv: rejecting if over my limit+ OverLimit ->+ closeConnection conn FlowControlError "Flow control error for stream in 0-RTT"+ -- Not a flow control error: the peer is inside its window, it is+ -- just spending it in more pieces than we will hold. Rate control+ -- answers with InternalError too.+ TooFragmented ->+ closeConnection conn QUIC.InternalError "Too many stream fragments"+ Duplicated -> return ()+ Reassembled -> do+ ok' <- checkRxMaxData conn len+ -- FLOW CONTROL: MAX_DATA: send: respecting peer's limit+ unless ok' $+ closeConnection+ conn+ FlowControlError+ "Flow control error for connection in 0-RTT" processFrame conn RTT1Level (StreamF sid _ [""] False) = do -- FLOW CONTROL: MAX_STREAMS: recv: rejecting if over my limit ok <- checkRxMaxStreams conn sid@@ -355,28 +370,33 @@ closeConnection conn StreamStateError "send-only stream" mstrm <- findStream conn sid guardStream conn sid mstrm- strm <- maybe (createStream conn sid) return mstrm- let len = BS.length dat- rx = RxStreamData dat off len fin- fc <- putRxStreamData strm rx- case fc of- -- FLOW CONTROL: MAX_STREAM_DATA: recv: rejecting if over my limit- OverLimit ->- closeConnection conn FlowControlError "Flow control error for stream in 1-RTT"- -- Not a flow control error: the peer is inside its window, it is- -- just spending it in more pieces than we will hold. Rate control- -- answers with InternalError too.- TooFragmented ->- closeConnection conn QUIC.InternalError "Too many stream fragments"- Duplicated -> return ()- Reassembled -> do- ok' <- checkRxMaxData conn len- -- FLOW CONTROL: MAX_DATA: send: respecting peer's limit- unless ok' $- closeConnection- conn- FlowControlError- "Flow control error for connection in 1-RTT"+ -- Nothing when the stream has been closed. What arrives for it now+ -- is a copy of something already received, sent again because our+ -- acknowledgement crossed it. Opening the stream anew would hold+ -- the copy to the initial window and call it a flow control error.+ mstrm' <- maybe (openStream conn sid) (return . Just) mstrm+ forM_ mstrm' $ \strm -> do+ let len = BS.length dat+ rx = RxStreamData dat off len fin+ fc <- putRxStreamData strm rx+ case fc of+ -- FLOW CONTROL: MAX_STREAM_DATA: recv: rejecting if over my limit+ OverLimit ->+ closeConnection conn FlowControlError "Flow control error for stream in 1-RTT"+ -- Not a flow control error: the peer is inside its window, it is+ -- just spending it in more pieces than we will hold. Rate control+ -- answers with InternalError too.+ TooFragmented ->+ closeConnection conn QUIC.InternalError "Too many stream fragments"+ Duplicated -> return ()+ Reassembled -> do+ ok' <- checkRxMaxData conn len+ -- FLOW CONTROL: MAX_DATA: send: respecting peer's limit+ unless ok' $+ closeConnection+ conn+ FlowControlError+ "Flow control error for connection in 1-RTT" processFrame conn lvl (MaxData n) = do when (lvl == InitialLevel || lvl == HandshakeLevel) $ closeConnection conn ProtocolViolation "MAX_DATA in Initial or Handshake"
Network/QUIC/Server/Reader.hs view
@@ -186,10 +186,21 @@ let send' b = void $ NSB.sendTo mysock b peersa -- cf: greaseQuicBit $ getMyParameters conn quicBit = greaseQuicBit $ scParameters conf- cpckts <- decodeCryptPackets bs (not quicBit) let bytes = BS.length bs switch = dispatch d conf forkConnection logAction mysock peersa send' bytes now- mapM_ switch cpckts+ -- One datagram's failure is one datagram's problem.+ --+ -- Letting it out of the loop ends the dispatcher, and nothing+ -- restarts it. The socket stays bound, because 'run' holds it+ -- for as long as the server lives, so the server goes on looking+ -- like a server and answers nothing for the rest of its life --+ -- with the handler above logging it to a logger that discards+ -- what it is given. Whatever a peer could find that raises in+ -- here would be one datagram against the whole server, from any+ -- address, before any handshake.+ handleLogUnit logAction $ do+ cpckts <- decodeCryptPackets bs (not quicBit)+ mapM_ switch cpckts loop logAction _msg = return ()
Network/QUIC/Server/Run.hs view
@@ -40,10 +40,16 @@ -- The action is executed with a new connection -- in a new lightweight thread. run :: ServerConfig -> (Connection -> IO ()) -> IO ()-run conf server = handleLogUnit debugLog $ do+run conf server = do labelMe "QUIC run" stvar <- newTVarIO Running- E.bracket (setup stvar) teardown $ \(_, _, _) -> do+ -- Outside handleLogUnit on purpose. If the addresses cannot be bound+ -- there is no server, and that is the caller's business: swallowing it+ -- returned from 'run' as if all were well, having never reached+ -- onServerReady, and left anyone waiting on that hook waiting for good.+ -- An IOSpec run wedged for two and a half days that way, on a port the+ -- previous test had not finished releasing.+ E.bracket (setup stvar) teardown $ \(_, _, _) -> handleLogUnit debugLog $ do onServerReady $ scHooks conf atomically $ do st <- readTVar stvar@@ -53,7 +59,6 @@ setup stvar = do dispatch <- newDispatch conf let forkConn acc = void $ forkIO (runServer conf server dispatch stvar acc)- -- fixme: the case where sockets cannot be created. ssas <- mapM serverSocket $ scAddresses conf tids <- mapM (runDispatcher dispatch conf stvar forkConn) ssas return (dispatch, tids, ssas)@@ -66,10 +71,11 @@ -- The action is executed with a new connection -- in a new lightweight thread. runWithSockets :: [NS.Socket] -> ServerConfig -> (Connection -> IO ()) -> IO ()-runWithSockets ssas conf server = handleLogUnit debugLog $ do+runWithSockets ssas conf server = do labelMe "QUIC runWithSockets" stvar <- newTVarIO Running- E.bracket (setup stvar) teardown $ \(_, _) -> do+ -- As in 'run'.+ E.bracket (setup stvar) teardown $ \(_, _) -> handleLogUnit debugLog $ do onServerReady $ scHooks conf atomically $ do st <- readTVar stvar@@ -79,7 +85,6 @@ setup stvar = do dispatch <- newDispatch conf let forkConn acc = void $ forkIO (runServer conf server dispatch stvar acc)- -- fixme: the case where sockets cannot be created. tids <- mapM (runDispatcher dispatch conf stvar forkConn) ssas return (dispatch, tids) teardown (dispatch, tids) = do
quic.cabal view
@@ -1,6 +1,6 @@ cabal-version: 2.0 name: quic-version: 0.3.6+version: 0.3.7 license: BSD3 license-file: LICENSE maintainer: kazu@iij.ad.jp@@ -24,6 +24,10 @@ description: Development commands default: False +flag qlog+ description: Write qlog from the test suite+ default: False+ library exposed-modules: Network.QUIC@@ -240,6 +244,10 @@ default-language: Haskell2010 default-extensions: Strict StrictData ghc-options: -Wall -threaded -rtsopts++ if flag(qlog)+ cpp-options: -DQLOG+ build-depends: base >=4.9 && <5, QuickCheck,
test/Config.hs view
@@ -1,3 +1,4 @@+{-# LANGUAGE CPP #-} {-# LANGUAGE OverloadedStrings #-} module Config (@@ -6,6 +7,7 @@ testClientConfig, testClientConfigR, setServerQlog,+ prepareQlog, setClientQlog, withPipe, withPipeStray,@@ -16,15 +18,18 @@ import Control.Concurrent import qualified Control.Exception as E import Control.Monad+import Data.Bits ((.&.)) import Data.ByteString (ByteString) import qualified Data.ByteString as BS-import Data.Bits ((.&.)) import Data.IORef import qualified Data.List as L import qualified Data.List.NonEmpty as NE import Network.Socket import Network.Socket.ByteString import Network.TLS hiding (Version)+#ifdef QLOG+import System.Directory (createDirectoryIfMissing)+#endif import Network.QUIC.Client import Network.QUIC.Internal@@ -102,6 +107,7 @@ } } +#ifdef QLOG -- | Write qlog for the connections that go through 'withPipe'. -- -- These are the tests that lose packets on purpose, and so the ones that@@ -111,16 +117,46 @@ -- fired. Two stalls found at a rate of one run in a few hundred were read -- straight off these traces, and neither would have been diagnosable -- without them. CI keeps the directory when a job fails.+-- | Where qlog goes when the @qlog@ flag is on.+qlogDir :: FilePath+qlogDir = "qlog"+#endif++-- | Create 'qlogDir' if the tests are going to write into it.+--+-- The directory used to be the CI's job, and a checkout without it -- a fresh+-- clone, or a @git clean@ -- failed 33 examples with @openFile: does not+-- exist@, which says nothing about qlog. It is the test suite's directory,+-- so the test suite makes it.+prepareQlog :: IO ()+#ifdef QLOG+prepareQlog = createDirectoryIfMissing True qlogDir+#else+prepareQlog = return ()+#endif+ setServerQlog :: ServerConfig -> ServerConfig-setServerQlog sc = sc{scQLog = Just "qlog"}+#ifdef QLOG+setServerQlog sc = sc{scQLog = Just qlogDir}+#else+setServerQlog sc = sc+#endif setClientQlog :: ClientConfig -> ClientConfig-setClientQlog cc = cc{ccQLog = Just "qlog"}+#ifdef QLOG+setClientQlog cc = cc{ccQLog = Just qlogDir}+#else+setClientQlog cc = cc+#endif data Scenario = Randomly Int | DropClientPacket [Int] | DropServerPacket [Int]+ | -- | Hold the n-th datagram from the client back for a second.+ DelayClientPacket Int+ | -- | Hold the n-th datagram from the server back for a second.+ DelayServerPacket Int withPipe :: Scenario -> IO () -> IO () withPipe = withPipeWith False@@ -149,7 +185,21 @@ setSocketOption sockC ReuseAddr 1 setSocketOption sockS ReuseAddr 1 bind sockC saC- connect sockS saS+ -- Bound, not connected. A connected UDP socket turns an ICMP+ -- port-unreachable from the peer into ECONNREFUSED on the next+ -- operation, and the peers here come and go with every test, so+ -- the relay was being killed by a reply to something it had sent+ -- to an address that had just closed. It surfaced on the+ -- receive, not the send:+ --+ -- DIED C->S: Network.Socket.recvBuf: does not exist+ -- (Connection refused)+ --+ -- and since these are forkIO threads with nobody watching, the+ -- relay simply stopped. The client then sent Initial packets+ -- until the idle timeout and heard nothing.+ addrAny <- resolve "0"+ bind sockS $ addrAddress addrAny when stray $ E.bracket (openSocket addrC) close $ \sock -> void $ sendTo sock (BS.pack [0x40, 1, 2, 3]) saC@@ -160,9 +210,13 @@ -- the bracket has just closed. That surfaces as "threadWait: -- invalid argument (Bad file descriptor)" from a thread nobody is -- watching, and buries whatever the test was really failing on.- E.bracket (startRelay sockC sockS irefC irefS) stopRelay $ \_ -> body+ E.bracket (startRelay saS sockC sockS irefC irefS) stopRelay $ \_ -> body where- startRelay sockC sockS irefC irefS = do+ startRelay saS sockC sockS irefC irefS = do+ -- Where to send what comes back from the server. Set once the client+ -- has introduced itself, which is before the server can have anything+ -- to say about it.+ peerRef <- newIORef Nothing -- from client tid0 <- forkIO $ do -- Wait for the client to introduce itself, and take the first@@ -181,23 +235,34 @@ -- A client always opens with a long header; a leftover from an -- established connection is a short one. That tells them apart. (bs, saO) <- waitForClientHello sockC- connect sockC saO+ writeIORef peerRef $ Just saO n0 <- atomicModifyIORef' irefC $ \x -> (x + 1, x) dropPacket0 <- shouldDrop scenario True n0- unless dropPacket0 $ void $ send sockS bs+ unless dropPacket0 $ void $ sendTo sockS bs saS forever $ do- bs1 <- recv sockC 2048- n <- atomicModifyIORef' irefC $ \x -> (x + 1, x)- dropPacket <- shouldDrop scenario True n- let isCC = BS.length bs1 < 200- when (isCC || not dropPacket) $ void $ send sockS bs1+ (bs1, sa) <- recvFrom sockC 2048+ -- Only from the client we latched onto. The connect this+ -- replaces did that in the kernel; doing it here keeps+ -- leftovers from the connection that just closed from being+ -- counted, which would shift the drop indices.+ when (sa == saO) $ do+ n <- atomicModifyIORef' irefC $ \x -> (x + 1, x)+ dropPacket <- shouldDrop scenario True n+ let isCC = BS.length bs1 < 200+ when (isCC || not dropPacket) $+ delayIf (shouldDelay scenario True n) $+ void $+ sendTo sockS bs1 saS -- from server tid1 <- forkIO $ forever $ do- bs <- recv sockS 2048+ (bs, _) <- recvFrom sockS 2048 n <- atomicModifyIORef' irefS $ \x -> (x + 1, x) dropPacket <- shouldDrop scenario False n let isCC = BS.length bs < 200- when (isCC || not dropPacket) $ void $ send sockC bs+ when (isCC || not dropPacket) $ do+ mpeer <- readIORef peerRef+ forM_ mpeer $ \sa ->+ delayIf (shouldDelay scenario False n) $ void $ sendTo sockC bs sa return (tid0, tid1) stopRelay (tid0, tid1) = killThread tid0 >> killThread tid1 waitForClientHello sockC = do@@ -222,6 +287,13 @@ shouldDrop (DropServerPacket ns) fromC pn | fromC = return False | otherwise = return (pn `elem` ns)+ shouldDrop _ _ _ = return False+ shouldDelay (DelayClientPacket k) fromC pn = fromC && pn == k+ shouldDelay (DelayServerPacket k) fromC pn = not fromC && pn == k+ shouldDelay _ _ _ = False+ -- The packets that follow go ahead of the one held back.+ delayIf True send_ = void $ forkIO $ threadDelay 1000000 >> send_+ delayIf False send_ = send_ chooseALPN :: Version -> [ByteString] -> IO ByteString chooseALPN _ver protos = return $ case mh3idx of
test/HandshakeSpec.hs view
@@ -8,6 +8,7 @@ import qualified Data.ByteString as BS import Network.TLS (Group (..), HandshakeMode13 (..)) import qualified Network.TLS as TLS+import qualified System.Timeout as Timeout import Test.Hspec import Network.QUIC@@ -30,7 +31,13 @@ { onServerReady = putMVar var () } }- let waitS = takeMVar var :: IO ()+ -- With a timeout: 'run' reports a server it could not start, but it is+ -- reported to whoever forked it, and that is not us. Without this the+ -- take never returns and the whole suite stops rather than failing.+ let waitS =+ Timeout.timeout 5000000 (takeMVar var) >>= \r -> case r of+ Just () -> return ()+ Nothing -> expectationFailure "server never became ready" describe "handshake" $ do it "can handshake in the normal case" $ do let cc = testClientConfig
test/IOSpec.hs view
@@ -19,6 +19,7 @@ spec :: Spec spec = do+ runIO prepareQlog sc0 <- runIO makeTestServerConfigR var <- runIO newEmptyMVar let sc =@@ -29,7 +30,13 @@ } } let cc = setClientQlog testClientConfigR- let waitS = takeMVar var :: IO ()+ -- With a timeout: 'run' reports a server it could not start, but it is+ -- reported to whoever forked it, and that is not us. Without this the+ -- take never returns and the whole suite stops rather than failing.+ let waitS =+ Timeout.timeout 5000000 (takeMVar var) >>= \r -> case r of+ Just () -> return ()+ Nothing -> expectationFailure "server never became ready" describe "send & recv" $ do it "can exchange data on random dropping" $ do withPipe (Randomly 20) $ testSendRecv cc sc waitS 1000@@ -93,6 +100,11 @@ -- number of bytes sent by the RESET_STREAM sender. it "sends RESET_STREAM with the bytes sent as final size" $ do withPipe (DropClientPacket []) $ testResetStreamFinalSize cc sc waitS+ describe "closed stream" $ do+ it "ignores a late copy of the data it received on a stream it opened" $ do+ 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 "port handover" $ do it "ignores a leftover datagram from the connection that just closed" $ withPipeStray (Randomly 20) $ testSendRecv cc sc waitS 20@@ -147,6 +159,42 @@ consumeBytes strm (BS.length request) sendStream strm payload takeMVar doneVar++-- | One end has read a stream to its end and closed it while a packet for+-- it is still on the way. That packet was taken for lost and its data sent+-- again, so what arrives late is a copy, well past the initial window of+-- 256K. It must not open the stream anew: the new stream holds the copy to+-- that window and calls it a flow control error, or else hands the+-- 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+ where+ (upLen, downLen)+ | upload = (1000000, 10)+ | otherwise = (10, 1000000)++ client = do+ waitS+ C.run cc $ \conn -> do+ exchange conn+ -- The copy arrives while we wait.+ threadDelay 1500000+ exchange conn+ exchange conn = do+ strm <- stream conn+ sendStream strm (BS.replicate upLen 0)+ shutdownStream strm+ consumeBytes strm downLen+ assertEndOfStream strm+ closeStream strm++ server = run sc $ \conn -> forever $ do+ strm <- acceptStream conn+ consumeBytes strm upLen+ assertEndOfStream strm+ sendStream strm (BS.replicate downLen 0)+ closeStream strm testRecvStreamClientStopFirst :: C.ClientConfig -> ServerConfig -> IO () -> IO ()