mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-05 23:07:22 +00:00
smp-server: limit concurrent resolver requests
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -167,6 +167,8 @@ iniFileContent cfgPath logPath opts host basicAuth controlPortPwds =
|
||||
\# resolver_auth: basic <username>:<password>\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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
+16
-1
@@ -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 =
|
||||
|
||||
Reference in New Issue
Block a user