From 2733a09a40336211fa0278d192f5ca5d112437b9 Mon Sep 17 00:00:00 2001 From: Evgeny Poberezkin <2769109+epoberezkin@users.noreply.github.com> Date: Sat, 30 Apr 2022 09:36:52 +0100 Subject: [PATCH] limit concurrency when resubscribing, method to resubscribe if not subscribed (#368) --- src/Simplex/Messaging/Agent.hs | 12 +++++++++++- src/Simplex/Messaging/Agent/Client.hs | 17 +++++++++++------ src/Simplex/Messaging/Agent/Env/SQLite.hs | 2 ++ src/Simplex/Messaging/Server.hs | 2 +- 4 files changed, 25 insertions(+), 8 deletions(-) diff --git a/src/Simplex/Messaging/Agent.hs b/src/Simplex/Messaging/Agent.hs index 9ee45f362..ee67611a7 100644 --- a/src/Simplex/Messaging/Agent.hs +++ b/src/Simplex/Messaging/Agent.hs @@ -43,6 +43,7 @@ module Simplex.Messaging.Agent acceptContact, rejectContact, subscribeConnection, + resubscribeConnection, sendMessage, ackMessage, suspendConnection, @@ -139,6 +140,9 @@ rejectContact c = withAgentEnv c .: rejectContact' c subscribeConnection :: AgentErrorMonad m => AgentClient -> ConnId -> m () subscribeConnection c = withAgentEnv c . subscribeConnection' c +resubscribeConnection :: AgentErrorMonad m => AgentClient -> ConnId -> m () +resubscribeConnection c = withAgentEnv c . resubscribeConnection' c + -- | Send message to the connection (SEND command) sendMessage :: AgentErrorMonad m => AgentClient -> ConnId -> MsgBody -> m AgentMsgId sendMessage c = withAgentEnv c .: sendMessage' c @@ -355,6 +359,12 @@ subscribeConnection' c connId = SomeConn _ (RcvConnection _ rq) -> subscribeQueue c rq connId SomeConn _ (ContactConnection _ rq) -> subscribeQueue c rq connId +resubscribeConnection' :: forall m. AgentMonad m => AgentClient -> ConnId -> m () +resubscribeConnection' c connId = + unlessM + (atomically $ hasActiveSubscription c connId) + (subscribeConnection' c connId) + -- | Send message to the connection (SEND command) in Reader monad sendMessage' :: forall m. AgentMonad m => AgentClient -> ConnId -> MsgBody -> m AgentMsgId sendMessage' c connId msg = @@ -398,7 +408,7 @@ resumeMsgDelivery c connId sq@SndQueue {server, sndId} = do withStore (`getPendingMsgs` connId) >>= queuePendingMsgs c connId sq where - queueDelivering qKey = atomically $ isJust <$> TM.lookup qKey (smpQueueMsgDeliveries c) + queueDelivering qKey = atomically $ TM.member qKey (smpQueueMsgDeliveries c) connQueued = atomically $ isJust <$> TM.lookupInsert connId True (connMsgsQueued c) queuePendingMsgs :: AgentMonad m => AgentClient -> ConnId -> SndQueue -> [InternalId] -> m () diff --git a/src/Simplex/Messaging/Agent/Client.hs b/src/Simplex/Messaging/Agent/Client.hs index 21128b2ce..2d842e3c9 100644 --- a/src/Simplex/Messaging/Agent/Client.hs +++ b/src/Simplex/Messaging/Agent/Client.hs @@ -38,6 +38,7 @@ module Simplex.Messaging.Agent.Client deleteQueue, logServer, removeSubscription, + hasActiveSubscription, ) where @@ -55,7 +56,7 @@ import qualified Data.ByteString.Char8 as B import Data.List.NonEmpty (NonEmpty) import Data.Map.Strict (Map) import qualified Data.Map.Strict as M -import Data.Maybe (catMaybes, isNothing) +import Data.Maybe (catMaybes) import Data.Text.Encoding import Data.Word (Word16) import Simplex.Messaging.Agent.Env.SQLite @@ -75,7 +76,7 @@ import qualified Simplex.Messaging.TMap as TM import Simplex.Messaging.Util (bshow, catchAll_, ifM, liftEitherError, liftError, tryError, unlessM, whenM) import Simplex.Messaging.Version import System.Timeout (timeout) -import UnliftIO (async, forConcurrently) +import UnliftIO (async, pooledForConcurrentlyN) import qualified UnliftIO.Exception as E import UnliftIO.STM @@ -200,14 +201,15 @@ getSMPServerClient c@AgentClient {active, smpClients, msgQ} srv = do reconnectClient `catchError` const loop reconnectClient :: m () - reconnectClient = + reconnectClient = do + n <- asks $ resubscriptionConcurrency . config withAgentLock c . withClient c srv $ \smp -> do cs <- atomically $ mapM readTVar =<< TM.lookup srv (pendingSubscrSrvrs c) - conns <- forConcurrently (maybe [] M.toList cs) $ \sub@(connId, _) -> + conns <- pooledForConcurrentlyN n (maybe [] M.toList cs) $ \sub@(connId, _) -> ifM - (atomically $ isNothing <$> TM.lookup connId (subscrConns c)) - (subscribe_ smp sub `catchError` handleError connId) + (atomically $ hasActiveSubscription c connId) (pure $ Just connId) + (subscribe_ smp sub `catchError` handleError connId) liftIO . unless (null conns) . notifySub "" . UP srv $ catMaybes conns where subscribe_ :: SMPClient -> (ConnId, RcvQueue) -> ExceptT ProtocolClientError IO (Maybe ConnId) @@ -408,6 +410,9 @@ addSubscription c rq@RcvQueue {server} connId = atomically $ do addSubs_ (subscrSrvrs c) rq connId removePendingSubscription c server connId +hasActiveSubscription :: AgentClient -> ConnId -> STM Bool +hasActiveSubscription c connId = TM.member connId (subscrConns c) + addPendingSubscription :: AgentClient -> RcvQueue -> ConnId -> STM () addPendingSubscription = addSubs_ . pendingSubscrSrvrs diff --git a/src/Simplex/Messaging/Agent/Env/SQLite.hs b/src/Simplex/Messaging/Agent/Env/SQLite.hs index bd1207426..8dba33087 100644 --- a/src/Simplex/Messaging/Agent/Env/SQLite.hs +++ b/src/Simplex/Messaging/Agent/Env/SQLite.hs @@ -49,6 +49,7 @@ data AgentConfig = AgentConfig ntfCfg :: ProtocolClientConfig, reconnectInterval :: RetryInterval, helloTimeout :: NominalDiffTime, + resubscriptionConcurrency :: Int, caCertificateFile :: FilePath, privateKeyFile :: FilePath, certificateFile :: FilePath @@ -78,6 +79,7 @@ defaultAgentConfig = ntfCfg = defaultClientConfig {defaultTransport = ("443", transport @TLS)}, reconnectInterval = defaultReconnectInterval, helloTimeout = 2 * nominalDay, + resubscriptionConcurrency = 16, -- CA certificate private key is not needed for initialization -- ! we do not generate these caCertificateFile = "/etc/opt/simplex-agent/ca.crt", diff --git a/src/Simplex/Messaging/Server.hs b/src/Simplex/Messaging/Server.hs index c2ce106ed..5a246bd92 100644 --- a/src/Simplex/Messaging/Server.hs +++ b/src/Simplex/Messaging/Server.hs @@ -343,7 +343,7 @@ client clnt@Client {subscriptions, ntfSubscriptions, rcvQ, sndQ} Server {subscri subscribeNotifications :: m (Transmission BrokerMsg) subscribeNotifications = atomically $ do - whenM (isNothing <$> TM.lookup queueId ntfSubscriptions) $ do + unlessM (TM.member queueId ntfSubscriptions) $ do writeTBQueue ntfSubscribedQ (queueId, clnt) TM.insert queueId () ntfSubscriptions pure ok