mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-05 20:57:16 +00:00
smp: remove unsent and timed-out proxy requests
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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" $ \_ ->
|
||||
|
||||
Reference in New Issue
Block a user