From 35767d74c51e9bbc6b61c693df6abdfd65f03059 Mon Sep 17 00:00:00 2001 From: shum Date: Wed, 30 Sep 2026 11:30:42 +0000 Subject: [PATCH] smp: remove unsent and timed-out proxy requests --- src/Simplex/Messaging/Client.hs | 20 +++++++++++++------- tests/SMPProxyTests.hs | 2 +- 2 files changed, 14 insertions(+), 8 deletions(-) diff --git a/src/Simplex/Messaging/Client.hs b/src/Simplex/Messaging/Client.hs index a1c5210d6..320d0f130 100644 --- a/src/Simplex/Messaging/Client.hs +++ b/src/Simplex/Messaging/Client.hs @@ -152,7 +152,7 @@ import Data.List (find, isSuffixOf) import Data.List.NonEmpty (NonEmpty (..)) import qualified Data.List.NonEmpty as L import qualified Data.Map.Strict as M -import Data.Maybe (catMaybes, fromMaybe) +import Data.Maybe (catMaybes, fromMaybe, isNothing) import Data.Text (Text) import qualified Data.Text as T import Data.Time.Clock (UTCTime (..), diffUTCTime, getCurrentTime) @@ -1359,20 +1359,22 @@ sendProtocolCommand c nm = sendProtocolCommand_ c nm Nothing Nothing -- -- Please note: if nonce is passed it is also used as a correlation ID sendProtocolCommand_ :: forall v err msg. Protocol v err msg => ProtocolClient v err msg -> NetworkRequestMode -> Maybe C.CbNonce -> Maybe Int -> Maybe C.APrivateAuthKey -> EntityId -> ProtoCommand msg -> ExceptT (ProtocolClientError err) IO msg -sendProtocolCommand_ c@ProtocolClient {client_ = PClient {sndQ}, thParams = THandleParams {blockSize, serviceAuth}} nm nonce_ tOut pKey entId cmd = +sendProtocolCommand_ c@ProtocolClient {client_ = PClient {sndQ, sentCommands}, thParams = THandleParams {blockSize, serviceAuth}} nm nonce_ tOut pKey entId cmd = ExceptT $ uncurry sendRecv =<< mkTransmission_ c nonce_ (entId, pKey, cmd) where -- two separate "atomically" needed to avoid blocking sendRecv :: Either TransportError SentRawTransmission -> Request err msg -> IO (Either (ProtocolClientError err) msg) - sendRecv t_ r = case t_ of - Left e -> pure . Left $ PCETransportError e + sendRecv t_ r@Request {corrId} = case t_ of + Left e -> notSent e Right t - | B.length s > blockSize - 2 -> pure . Left $ PCETransportError TELargeMsg + | B.length s > blockSize - 2 -> notSent TELargeMsg | otherwise -> do nonBlockingWriteTBQueue sndQ (Just r, s) response <$> getResponse c nm tOut r where s = tEncodeBatch1 serviceAuth t + where + notSent e = Left (PCETransportError e) <$ atomically (TM.delete corrId sentCommands) nonBlockingWriteTBQueue :: TBQueue a -> a -> IO () nonBlockingWriteTBQueue q x = do @@ -1380,7 +1382,7 @@ nonBlockingWriteTBQueue q x = do unless sent $ void $ forkIO $ atomically $ writeTBQueue q x getResponse :: ProtocolClient v err msg -> NetworkRequestMode -> Maybe Int -> Request err msg -> IO (Response err msg) -getResponse ProtocolClient {client_ = PClient {tcpTimeout, timeoutErrorCount}} nm tOut Request {entityId, pending, responseVar} = do +getResponse ProtocolClient {client_ = PClient {tcpTimeout, timeoutErrorCount, sentCommands, msgQ}} nm tOut Request {corrId, entityId, pending, responseVar} = do r <- fromMaybe (netTimeoutInt tcpTimeout nm) tOut `timeout` atomically (takeTMVar responseVar) response <- atomically $ do writeTVar pending False @@ -1389,7 +1391,11 @@ getResponse ProtocolClient {client_ = PClient {tcpTimeout, timeoutErrorCount}} n -- See `processMsg`. ((r <|>) <$> tryTakeTMVar responseVar) >>= \case Just r' -> writeTVar timeoutErrorCount 0 $> r' - Nothing -> modifyTVar' timeoutErrorCount (+ 1) $> Left PCEResponseTimeout + Nothing -> do + modifyTVar' timeoutErrorCount (+ 1) + -- a late response is delivered to msgQ, without msgQ it is only logged + when (isNothing msgQ) $ TM.delete corrId sentCommands + pure $ Left PCEResponseTimeout pure Response {entityId, response} mkTransmission :: Protocol v err msg => ProtocolClient v err msg -> ClientCommand msg -> IO (PCTransmission err msg) diff --git a/tests/SMPProxyTests.hs b/tests/SMPProxyTests.hs index 29ac4a0bd..4ab300c60 100644 --- a/tests/SMPProxyTests.hs +++ b/tests/SMPProxyTests.hs @@ -66,7 +66,7 @@ smpProxyTests = do testProxyReconnectAfterRelayRestart xit "must drop a stuck relay session after forward timeouts" $ \_ -> testProxyForwardTimeoutStuckSession - xit "does not keep oversized forwarded command" $ \_ -> + it "does not keep oversized forwarded command" $ \_ -> testForwardOversizedNotKept describe "agent client reconnection" $ do it "reconnects after a connect is cancelled mid-flight" $ \_ ->