packages feed

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 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 ()