diff --git a/src/Simplex/Messaging/Server.hs b/src/Simplex/Messaging/Server.hs index 0a1f18a90..83164d791 100644 --- a/src/Simplex/Messaging/Server.hs +++ b/src/Simplex/Messaging/Server.hs @@ -1497,7 +1497,7 @@ client Nothing -> incStat (rslvDisabled st) $> Nothing Just nenv -> pure (Just nenv) -- Runs on a forked thread so RSLV does not block other commands; - -- concurrency is limited by serverResolverConcurrency in forkCmd. + -- concurrency is limited per connection by serverResolverConcurrency in forkCmd and globally in resolveName. resolveNameMsg :: VersionSMP -> NamesEnv -> NameQuery -> M s BrokerMsg resolveNameMsg v nenv q = do st <- asks (rslvStats . serverStats) diff --git a/src/Simplex/Messaging/Server/Main.hs b/src/Simplex/Messaging/Server/Main.hs index 138a25fc0..83bef9db0 100644 --- a/src/Simplex/Messaging/Server/Main.hs +++ b/src/Simplex/Messaging/Server/Main.hs @@ -814,7 +814,8 @@ readNamesConfig ini { resolverEndpoint = either (error . ("[NAMES] resolver_endpoint: " <>)) id (validateUrl endpoint resolverAuth_), resolverAuth = resolverAuth_, resolverTimeoutMs = boundedIniInt 3000 100 60000 "resolver_timeout_ms", - resolverMaxResponseBytes = boundedIniInt 16000 1024 16000 "resolver_max_response_bytes" + resolverMaxResponseBytes = boundedIniInt 16000 1024 16000 "resolver_max_response_bytes", + resolverGlobalConcurrency = boundedIniInt 32 1 1000 "resolver_global_concurrency" } where enabled = fromMaybe False (iniOnOff "NAMES" "enable" ini) diff --git a/src/Simplex/Messaging/Server/Main/Init.hs b/src/Simplex/Messaging/Server/Main/Init.hs index 89df1ad93..ac947f55c 100644 --- a/src/Simplex/Messaging/Server/Main/Init.hs +++ b/src/Simplex/Messaging/Server/Main/Init.hs @@ -167,6 +167,8 @@ iniFileContent cfgPath logPath opts host basicAuth controlPortPwds = \# resolver_auth: basic :\n\ \# resolver_timeout_ms: 3000\n\ \# resolver_max_response_bytes: 16000\n\ + \# Max concurrent requests to the resolver from all connections.\n\ + \# resolver_global_concurrency: 32\n\ \# Max concurrent name resolutions per connection (forwarded RSLVs from many\n\ \# clients share one proxy connection, so this is much higher than PROXY client_concurrency).\n" <> ("# resolver_concurrency = " <> tshow defaultNameResolverConcurrency) diff --git a/src/Simplex/Messaging/Server/Names.hs b/src/Simplex/Messaging/Server/Names.hs index 601519c56..17e90c919 100644 --- a/src/Simplex/Messaging/Server/Names.hs +++ b/src/Simplex/Messaging/Server/Names.hs @@ -15,6 +15,7 @@ module Simplex.Messaging.Server.Names ) where +import Control.Concurrent.QSem (QSem, newQSem, signalQSem, waitQSem) import qualified Control.Exception as E import Control.Logger.Simple (logError) import Data.Bifunctor (first) @@ -38,19 +39,22 @@ data NamesConfig = NamesConfig { resolverEndpoint :: String, resolverAuth :: Maybe RpcAuth, resolverTimeoutMs :: Int, - resolverMaxResponseBytes :: Int + resolverMaxResponseBytes :: Int, + resolverGlobalConcurrency :: Int } deriving (Show) data NamesEnv = NamesEnv { config :: NamesConfig, - resolverEnv :: ResolverEnv + resolverEnv :: ResolverEnv, + resolverSlots :: QSem } newNamesEnv :: NamesConfig -> IO NamesEnv newNamesEnv config = do - resolverEnv <- newResolverEnv (resolverEndpoint config) (resolverAuth config) (resolverTimeoutMs config) (resolverMaxResponseBytes config) - pure NamesEnv {config, resolverEnv} + resolverEnv <- newResolverEnv (resolverEndpoint config) (resolverAuth config) (resolverTimeoutMs config) (resolverMaxResponseBytes config) (resolverGlobalConcurrency config) + resolverSlots <- newQSem (resolverGlobalConcurrency config) + pure NamesEnv {config, resolverEnv, resolverSlots} closeNamesEnv :: NamesEnv -> IO () closeNamesEnv NamesEnv {resolverEnv} = closeResolverEnv resolverEnv @@ -60,8 +64,9 @@ pingEndpoint NamesEnv {resolverEnv, config} = fromMaybe (Left ResolverTimeout) <$> timeout (resolverTimeoutMs config * 1000) (healthHttp resolverEnv) resolveName :: NamesEnv -> NameQuery -> IO (Either NameErrorType NameResponse) -resolveName env q = do - r <- E.try (timeout (resolverTimeoutMs (config env) * 1000) (fetch env q)) +resolveName env@NamesEnv {resolverSlots} q = do + -- waiting for a slot counts towards the timeout, so lookups do not queue longer than resolverTimeoutMs + r <- E.try (timeout (resolverTimeoutMs (config env) * 1000) (E.bracket_ (waitQSem resolverSlots) (signalQSem resolverSlots) (fetch env q))) case r of Right result -> pure (fromMaybe (Left (RESOLVER "timeout")) result) Left e diff --git a/src/Simplex/Messaging/Server/Names/HttpResolver.hs b/src/Simplex/Messaging/Server/Names/HttpResolver.hs index 0f272bcf4..504348fe8 100644 --- a/src/Simplex/Messaging/Server/Names/HttpResolver.hs +++ b/src/Simplex/Messaging/Server/Names/HttpResolver.hs @@ -85,9 +85,9 @@ data ResolverError | ResolverTimeout deriving (Show) -newResolverEnv :: String -> Maybe RpcAuth -> Int -> Int -> IO ResolverEnv -newResolverEnv baseUrl auth_ timeoutMs maxResponseBytes = do - manager <- HC.newManager tlsManagerSettings {managerConnCount = 10} +newResolverEnv :: String -> Maybe RpcAuth -> Int -> Int -> Int -> IO ResolverEnv +newResolverEnv baseUrl auth_ timeoutMs maxResponseBytes maxConcurrency = do + manager <- HC.newManager tlsManagerSettings {managerConnCount = maxConcurrency} pure ResolverEnv { manager, diff --git a/tests/NamesResolverServer.hs b/tests/NamesResolverServer.hs index 054d55e40..59414dc6b 100644 --- a/tests/NamesResolverServer.hs +++ b/tests/NamesResolverServer.hs @@ -61,7 +61,8 @@ testNamesConfig port = { resolverEndpoint = "http://127.0.0.1:" <> show port, resolverAuth = Nothing, resolverTimeoutMs = 1000, - resolverMaxResponseBytes = 65536 + resolverMaxResponseBytes = 65536, + resolverGlobalConcurrency = 8 } memCfg :: AServerConfig diff --git a/tests/RSLVTests.hs b/tests/RSLVTests.hs index 96bd0639f..ded51a0a5 100644 --- a/tests/RSLVTests.hs +++ b/tests/RSLVTests.hs @@ -51,6 +51,7 @@ import Simplex.Messaging.Protocol ) import qualified Simplex.Messaging.Protocol as SMP import Simplex.Messaging.Server.Env.STM (ServerConfig (..)) +import Simplex.Messaging.Server.Names (NamesConfig (..)) import Simplex.Messaging.SimplexName (SimplexDomain) import Simplex.Messaging.Transport import Simplex.Messaging.Version (mkVersionRange) @@ -106,8 +107,9 @@ rslvTests = do it "RSLV sends the 2LD as its hash" testRslvSendsTheHash it "a name with subnames is sent as text" testSubnameKeepsItsLabels it "a record naming a different name is rejected" testRslvWrongName - describe "RSLV resource use" $ + describe "RSLV resource use" $ do it "one connection has at most resolver_concurrency lookups in flight" testRslvConnectionCap + it "all connections have at most resolver_global_concurrency lookups in flight" testRslvFanOut -- | /v2/resolve answers 200, 400 or 502, so a 404 is a resolver that predates -- the route, not a name that does not exist. @@ -295,6 +297,19 @@ testRslvConnectionCap = where connCap = 4 +testRslvFanOut :: IO () +testRslvFanOut = + NRS.withResolverServerDelayed 3000 (NRS.resolveResp status200 "{}") $ \port reqs -> + withSmpServerConfigOn (transport @TLS) (withNames port memCfg) testPort $ const $ + testSMPClient @TLS $ \h1 -> testSMPClient @TLS $ \h2 -> do + sendRslvs h1 "a" 32 + sendRslvs h2 "b" 32 + threadDelay 800000 + length <$> resolvePaths reqs `shouldReturn` resolverGlobalConcurrency (NRS.testNamesConfig port) + let timedOut = replicate 32 (Right (ERR (NAME (RESOLVER "timeout")))) + recvResponses h1 32 `shouldReturn` timedOut + recvResponses h2 32 `shouldReturn` timedOut + -- | One RSLV per block, so no batch limit applies. sendRslvs :: THandleSMP TLS 'TClient -> String -> Int -> IO () sendRslvs h@THandle {params} prefix n =