bench: measure proxy mesh servers separately

This commit is contained in:
shum
2026-09-30 23:55:05 +00:00
parent 37ee78f8c3
commit 84d91bd445
+306 -5
View File
@@ -36,7 +36,7 @@
module Main (main) where
import Control.Concurrent (threadDelay)
import Control.Concurrent.Async (concurrently_, forConcurrently, forConcurrently_, mapConcurrently_, wait, withAsync)
import Control.Concurrent.Async (async, concurrently_, forConcurrently, forConcurrently_, mapConcurrently_, wait, withAsync)
import Control.Logger.Simple (LogConfig (..), LogLevel (..), setLogLevel, withGlobalLogging)
import qualified Control.Exception as E
import Control.Concurrent.STM
@@ -49,6 +49,9 @@ import Data.Int (Int64)
import Data.Foldable (toList)
import Data.List.NonEmpty (NonEmpty (..), fromList)
import qualified Data.Map.Strict as M
import qualified Data.List.NonEmpty as L
import Data.Functor ((<&>))
import GHC.Conc (listThreads)
import Data.Maybe (fromMaybe)
import Data.Time.Clock (diffUTCTime, getCurrentTime)
import qualified Data.X509.Validation as XV
@@ -60,12 +63,12 @@ import SMPClient
import Simplex.Messaging.Client
import qualified Simplex.Messaging.Crypto as C
import Simplex.Messaging.Protocol
import Simplex.Messaging.Client.Agent (SMPClientAgentConfig (msgQSize))
import Simplex.Messaging.Server.Env.STM (AStoreType (..), ServerConfig (controlPort, controlPortAdminAuth, maxJournalMsgCount, msgQueueQuota, notificationExpiration, ntfDeliveryInterval, serverClientConcurrency, smpAgentCfg, storeNtfsFile))
import Simplex.Messaging.Client.Agent (SMPClientAgentConfig (msgQSize, persistErrorInterval))
import Simplex.Messaging.Server.Env.STM (AStoreType (..), ServerConfig (controlPort, controlPortAdminAuth, maxJournalMsgCount, msgQueueQuota, notificationExpiration, ntfDeliveryInterval, serverClientConcurrency, smpAgentCfg, storeNtfsFile, tbqSize))
import Simplex.Messaging.Server.Expiration (ExpirationConfig (..))
import Simplex.Messaging.Server.MsgStore.Types (SMSType (..), SQSType (..))
import Simplex.Messaging.Transport
import Simplex.Messaging.Transport.Client (TransportClientConfig (..), defaultTransportClientConfig, runTransportClient)
import Simplex.Messaging.Transport.Client (TransportClientConfig (..), TransportHost (..), defaultTransportClientConfig, runTransportClient)
import Simplex.Messaging.Transport.Credentials (genCredentials, tlsCredentials)
import Simplex.Messaging.Version (mkVersionRange)
import System.Environment (getArgs, lookupEnv, setEnv)
@@ -1548,6 +1551,8 @@ main = do
else if phase == "loadclient" then runLoadClient g
else if phase == "prodmix" then postgressBracket benchDBConnectInfo $ withSmpServerConfigOn (transport @TLS) benchPgCfg benchPort $ \_ -> runProdMix
else if phase == "prodclient" then runProdClient g
else if phase == "srvchild" then runSrvChild (args !! 1)
else if phase == "mesh" then postgressBracket benchDBConnectInfo $ runMesh g iters
#endif
else withSmpServerConfigOn (transport @TLS) srvCfg testPort $ \_ -> settle leakDiagSec $ do
threadDelay 250000
@@ -1632,8 +1637,304 @@ runPfwdBig g iters cp = do
withCheckpoints "pfwdbig" iters cp $ \_ ->
void $ runExceptT $ sendProtocolCommand pc NRMInteractive Nothing (EntityId prSessionId) pfwd
#if defined(dbServerPostgres)
-- Proxy mesh with each server in its own process, so the live bytes of the proxy and of the relay
-- are measured separately and without the bench clients. Both use the pure PostgreSQL store and the
-- production queue sizes (tbqSize 128, client_concurrency 32 unless MESH_CONC is set, quota 128).
-- Ports and database follow BENCHID, so runs do not need the test lock.
--
-- MESH selects the scenario, `iters` its size:
-- sessions PRXY to `iters` aliases of the relay (same host and port, distinct second host), so the
-- proxy opens `iters` relay sessions and the relay holds `iters` proxy connections
-- inflight MESH_CONNS client connections each with MESH_K concurrent PFWDs; the relay stops
-- answering (MESH_RELAY=drop, default) or stops reading (MESH_RELAY=stall)
-- rate MESH_CONNS connections sending `iters` PFWD/s in total for MESH_SEC seconds with
-- MESH_LAG_MS one-way relay lag
-- throughput MESH_CONNS connections x MESH_K concurrent PFWDs back to back for MESH_SEC seconds, then
-- the same concurrency sending SEND directly to the relay
-- slowreader `iters` raw connections sending PFWD without ever reading a response
runMesh :: TVar ChaChaDRG -> Int -> IO ()
runMesh g n = do
hSetBuffering stdout LineBuffering
scen <- fromMaybe "inflight" <$> lookupEnv "MESH"
withSrvChild "relay" $ \relay -> withSrvChild "proxy" $ \proxy -> do
let measureBoth tag = do
p <- childMeasure proxy
r <- childMeasure relay
printf "mesh %s: proxy %s | relay %s\n" (tag :: String) (showM p) (showM r)
pure (p, r)
case scen of
"sessions" -> meshSessions g n measureBoth
"inflight" -> meshInflight g relay measureBoth
"rate" -> meshRate g n relay measureBoth
"throughput" -> meshThroughput g measureBoth
"slowreader" -> meshSlowReader g n measureBoth
_ -> error $ "unknown MESH scenario: " <> scen
data ChildMem = ChildMem {cmLive :: Double, cmFrag :: Double, cmInUse :: Double, cmThreads :: Int}
showM :: ChildMem -> String
showM ChildMem {cmLive, cmFrag, cmInUse, cmThreads} = printf "live=%.1f frag=%.1f in_use=%.1f threads=%d" cmLive cmFrag cmInUse cmThreads
deltaPer :: String -> Int -> (ChildMem, ChildMem) -> (ChildMem, ChildMem) -> IO ()
deltaPer unit k (p0, r0) (p1, r1) =
printf "mesh SUMMARY per %s (n=%d): proxy %+.2f KiB live, %+.2f KiB in_use, %+.2f threads | relay %+.2f KiB live, %+.2f KiB in_use, %+.2f threads\n"
unit k (per cmLive p0 p1) (per cmInUse p0 p1) (perT p0 p1) (per cmLive r0 r1) (per cmInUse r0 r1) (perT r0 r1)
where
per f a b = (f b - f a) * 1024 / fromIntegral (max 1 k)
perT a b = fromIntegral (cmThreads b - cmThreads a) / fromIntegral (max 1 k) :: Double
data SrvChild = SrvChild {chIn :: Handle, chOut :: Handle}
withSrvChild :: String -> (SrvChild -> IO a) -> IO a
withSrvChild role action = do
exe <- getExecutablePath
rts <- maybe ["+RTS", "-N", "-F1.2", "-A16m", "-I0.01", "-Iw15", "-T", "-RTS"] words <$> lookupEnv ("MESH_RTS_" <> role)
(Just hIn, Just hOut, _, ph) <- createProcess (proc exe (["srvchild", role] <> rts)) {std_in = CreatePipe, std_out = CreatePipe}
hSetBuffering hIn LineBuffering
let c = SrvChild hIn hOut
(childReply c >> action c) `E.finally` (hClose hIn >> waitForProcess ph)
childReply :: SrvChild -> IO [String]
childReply c@SrvChild {chOut} = hGetLine chOut >>= \l -> case words l of
"R" : ws -> pure ws
_ -> childReply c
childCmd :: SrvChild -> String -> IO [String]
childCmd c@SrvChild {chIn} cmd = hPutStrLn chIn cmd >> childReply c
childMeasure :: SrvChild -> IO ChildMem
childMeasure c = childCmd c "m" >>= \case
[l, f, u, t] -> pure $ ChildMem (read l) (read f) (read u) (read t)
r -> fail $ "bad child reply: " <> unwords r
-- a server of the mesh; replies to each stdin command with one line starting with "R"
runSrvChild :: String -> IO ()
runSrvChild role = do
hSetBuffering stdout LineBuffering
case role of
"proxy" -> withSmpServerConfigOn (transport @TLS) meshProxyCfg benchPort $ \_ -> serve
"relay" -> withSmpServerConfigOn (transport @LagTLS) meshRelayCfg meshRelayPort $ \_ -> serve
_ -> error $ "unknown server role: " <> role
where
serve = putStrLn "R ready" >> loop
loop = (E.try getLine :: IO (Either E.IOException String)) >>= either (const $ pure ()) (\l -> cmd (words l) >> loop)
cmd = \case
["m"] -> do
live <- liveBytesMiB
frag <- fragmentationMiB
inUse <- (/ (1024 * 1024)) . fromIntegral . gcdetails_mem_in_use_bytes . gc <$> getRTSStats
ts <- length <$> listThreads
printf "R %.3f %.3f %.3f %d\n" live frag (inUse :: Double) ts
["lag", r, s] -> setLag (read r * 1000) (read s * 1000) >> putStrLn "R ok"
["drop", b] -> setDropSnd (b == "1") >> putStrLn "R ok"
["clear"] -> clearLag >> putStrLn "R ok"
c -> putStrLn $ "R unknown " <> unwords c
meshSrvCfg :: Int -> ServerConfig s -> ServerConfig s
meshSrvCfg conc c = c {tbqSize = 128, msgQueueQuota = 128, serverClientConcurrency = conc}
meshConc :: Int
meshConc = unsafePerformIO $ fromMaybe 32 . (>>= readMaybe) <$> lookupEnv "MESH_CONC"
{-# NOINLINE meshConc #-}
meshProxyCfg :: AServerConfig
meshProxyCfg = updateCfg (meshPgCfg "smp_server" $ proxyCfgMS (ASType SQSPostgres SMSPostgres)) $ \c ->
c {smpAgentCfg = (smpAgentCfg c) {persistErrorInterval = 30}}
meshRelayCfg :: AServerConfig
meshRelayCfg = meshPgCfg "smp_server2" $ cfgMS (ASType SQSPostgres SMSPostgres)
meshPgCfg :: String -> AServerConfig -> AServerConfig
meshPgCfg sch = \case
ASrvCfg SQSPostgres SMSPostgres c -> ASrvCfg SQSPostgres SMSPostgres (meshSrvCfg meshConc c) {serverStoreCfg = SSCDatabase (meshStoreCfg sch)}
c -> c
-- the BENCHID database, one schema per server
meshStoreCfg :: String -> PostgresStoreCfg
meshStoreCfg sch = PostgresStoreCfg {dbOpts = testStoreDBOpts {connstr, schema = B.pack sch}, dbStoreLogPath = Nothing, confirmMigrations = MCYesUp, deletedTTL = 86400}
where
connstr = B.pack $ "postgresql://" <> connectUser benchDBConnectInfo <> "@/" <> connectDatabase benchDBConnectInfo
meshRelayPort :: ServiceName
meshRelayPort = show $ 15002 + 100 * benchId
meshProxySrv :: SMPServer
meshProxySrv = SMPServer testHost benchPort testKeyHash
meshRelaySrv :: SMPServer
meshRelaySrv = SMPServer testHost2 meshRelayPort testKeyHash
envInt :: String -> Int -> IO Int
envInt name def = fromMaybe def . (>>= readMaybe) <$> lookupEnv name
-- sender IDs of fresh unsubscribed queues on the relay, created directly
relayQueues :: TVar ChaChaDRG -> Int -> IO [SenderId]
relayQueues g k = do
ts <- getCurrentTime
rc <- getProtocolClient g NRMInteractive (98, meshRelaySrv, Nothing) benchClientCfg [] Nothing ts (\_ -> pure ()) >>= either (fail . show) pure
qs <- forM ([1 .. k] :: [Int]) $ \_ -> do
(rPub, rKey) <- atomically $ C.generateAuthKeyPair C.SEd25519 g
(dhPub, _ :: C.PrivateKeyX25519) <- atomically $ C.generateKeyPair g
QIK {sndId} <- runExceptT' $ createSMPQueue rc NRMInteractive Nothing (rPub, rKey) dhPub Nothing SMOnlyCreate (QRMessaging Nothing) Nothing
pure sndId
closeProtocolClient rc
pure qs
meshProxyClient :: TVar ChaChaDRG -> Int64 -> IO SMPClient
meshProxyClient g n = do
ts <- getCurrentTime
getProtocolClient g NRMInteractive (n, meshProxySrv, Nothing) benchClientCfg [] Nothing ts (\_ -> pure ()) >>= either (fail . show) pure
-- client connections to the proxy, each with its own PRXY session reference
proxyClients :: TVar ChaChaDRG -> Int -> IO [(SMPClient, ProxiedRelay)]
proxyClients g k = forM ([1 .. k] :: [Int]) $ \i -> do
pc <- meshProxyClient g (fromIntegral i)
sess <- runExceptT' $ connectSMPProxiedRelay pc NRMInteractive meshRelaySrv Nothing
pure (pc, sess)
forwardOnce :: SMPClient -> ProxiedRelay -> SenderId -> IO Bool
forwardOnce pc sess sId =
runExceptT (proxySMPMessage pc NRMInteractive sess Nothing sId noMsgFlags "hello") <&> \case
Right (Right ()) -> True
Left (PCEProtocolError QUOTA) -> True
_ -> False
meshSessions :: TVar ChaChaDRG -> Int -> (String -> IO (ChildMem, ChildMem)) -> IO ()
meshSessions g n measureBoth = do
pc <- meshProxyClient g 1
base <- measureBoth "before"
forM_ ([1 .. n] :: [Int]) $ \i -> runExceptT (connectSMPProxiedRelay pc NRMInteractive (alias i) Nothing) >>= \case
Right _ -> pure ()
Left e -> fail $ "PRXY " <> show i <> ": " <> show e
cur <- measureBoth (show n <> " relay sessions")
deltaPer "relay session" n base cur
where
alias i = SMPServer (L.head testHost2 :| [THDomainName $ "alias" <> show i <> ".invalid"]) meshRelayPort testKeyHash
meshInflight :: TVar ChaChaDRG -> SrvChild -> (String -> IO (ChildMem, ChildMem)) -> IO ()
meshInflight g relay measureBoth = do
conns <- envInt "MESH_CONNS" 16
k <- envInt "MESH_K" 32
mode <- fromMaybe "drop" <$> lookupEnv "MESH_RELAY"
sIds <- relayQueues g 1
pcs <- proxyClients g conns
forM_ pcs $ \(pc, sess) -> forwardOnce pc sess (head sIds)
base <- measureBoth "before"
void $ childCmd relay $ if mode == "stall" then "lag 60000 0" else "drop 1"
as <- forM pcs $ \(pc, sess) -> forM ([1 .. k] :: [Int]) $ \_ -> async $ forwardOnce pc sess (head sIds)
threadDelay 8000000
cur <- measureBoth $ printf "%d x %d forwards in flight, relay %s" conns k mode
deltaPer "in-flight forward" (conns * min k meshConc) base cur
mapM_ (mapM_ wait) as
void $ childCmd relay "clear"
threadDelay 45000000
void $ measureBoth "45s after the clients timed out"
meshRate :: TVar ChaChaDRG -> Int -> SrvChild -> (String -> IO (ChildMem, ChildMem)) -> IO ()
meshRate g rate relay measureBoth = do
conns <- envInt "MESH_CONNS" 64
secs <- envInt "MESH_SEC" 30
lagMs <- envInt "MESH_LAG_MS" 100
sIds <- relayQueues g 256
pcs <- proxyClients g conns
forM_ pcs $ \(pc, sess) -> forwardOnce pc sess (head sIds)
base <- measureBoth "before"
void $ childCmd relay $ "lag " <> show lagMs <> " " <> show lagMs
ok <- newTVarIO (0 :: Int)
failed <- newTVarIO (0 :: Int)
let perConn = fromIntegral rate / fromIntegral conns :: Double
intervalUs = round (1000000 / perConn) :: Int
sender i (pc, sess) = forM_ ([1 .. round (perConn * fromIntegral secs)] :: [Int]) $ \j -> do
void $ async $ forwardOnce pc sess (sIds !! ((i * 7 + j) `mod` length sIds)) >>= \r -> atomically $ modifyTVar' (if r then ok else failed) (+ 1)
threadDelay intervalUs
sampler = forever $ do
threadDelay 5000000
(o, f) <- (,) <$> readTVarIO ok <*> readTVarIO failed
void $ measureBoth $ printf "rate=%d/s lag=%dms ok=%d failed=%d" rate lagMs o f
cur <- withAsync sampler $ \_ -> do
mapConcurrently_ (uncurry sender) (zip [0 ..] pcs)
measureBoth "end of sending"
deltaPer "forward/s" rate base cur
void $ childCmd relay "clear"
threadDelay 35000000
(o, f) <- (,) <$> readTVarIO ok <*> readTVarIO failed
printf "mesh rate: ok=%d failed=%d\n" o f
meshThroughput :: TVar ChaChaDRG -> (String -> IO (ChildMem, ChildMem)) -> IO ()
meshThroughput g measureBoth = do
conns <- envInt "MESH_CONNS" 16
k <- envInt "MESH_K" 32
secs <- envInt "MESH_SEC" 20
sIds <- relayQueues g 1024
pcs <- proxyClients g conns
forM_ pcs $ \(pc, sess) -> forwardOnce pc sess (head sIds)
void $ measureBoth "before"
let run label go = do
cnt <- newTVarIO (0 :: Int)
t0 <- getCurrentTime
let worker w = loop (0 :: Int)
where
loop j = do
t <- getCurrentTime
when (diffUTCTime t t0 < fromIntegral secs) $ do
r <- go w (sIds !! ((w * 31 + j) `mod` length sIds))
when r $ atomically $ modifyTVar' cnt (+ 1)
loop (j + 1)
forConcurrently_ ([0 .. conns * k - 1] :: [Int]) worker
c <- readTVarIO cnt
printf "mesh throughput %s: %d ok in %ds = %.0f/s (%d connections x %d concurrent)\n" (label :: String) c secs (fromIntegral c / fromIntegral secs :: Double) conns k
run "via proxy" $ \w sId -> let (pc, sess) = pcs !! (w `mod` conns) in forwardOnce pc sess sId
void $ measureBoth "after proxy run"
ts <- getCurrentTime
rcs <- forM ([1 .. conns] :: [Int]) $ \i -> getProtocolClient g NRMInteractive (fromIntegral i, meshRelaySrv, Nothing) benchClientCfg [] Nothing ts (\_ -> pure ()) >>= either (fail . show) pure
run "direct SEND" $ \w sId ->
runExceptT (sendSMPMessage (rcs !! (w `mod` conns)) NRMInteractive Nothing sId noMsgFlags "hello") <&> \case
Right () -> True
Left (PCEProtocolError QUOTA) -> True
_ -> False
-- An unauthenticated client that sends valid PFWDs and never reads: the proxy forwards them, and the
-- responses accumulate in its sndQ, then in the forked command threads, then in its rcvQ.
meshSlowReader :: TVar ChaChaDRG -> Int -> (String -> IO (ChildMem, ChildMem)) -> IO ()
meshSlowReader g n measureBoth = do
secs <- envInt "MESH_SEC" 30
sIds <- relayQueues g 1
(pc, sess) <- head <$> proxyClients g 1
_ <- forwardOnce pc sess (head sIds)
base <- measureBoth "before"
sent <- newTVarIO (0 :: Int)
let attacker = benchClient $ \h -> do
t <- pfwdTransmission g (thParams' h) sess (head sIds)
forever $ tPut1 h (Nothing, t) >> atomically (modifyTVar' sent (+ 1))
thParams' THandle {params} = params
-- not awaited: closing TLS sends an alert, which blocks as the proxy no longer reads, so the
-- connections stay until the process exits
void $ async $ forConcurrently_ ([1 .. n] :: [Int]) $ \_ -> attacker
threadDelay $ secs * 1000000
s <- readTVarIO sent
cur <- measureBoth $ printf "%d non-reading clients, %d PFWD accepted by TCP" n s
deltaPer "non-reading client" n base cur
-- the PFWD a client builds in proxySMPCommand, as one transmission for a raw connection
pfwdTransmission :: TVar ChaChaDRG -> THandleParams SMPVersion 'TClient -> ProxiedRelay -> SenderId -> IO ByteString
pfwdTransmission g params ProxiedRelay {prSessionId, prVersion, prServerKey} sId = do
let serverThAuth = (\ta -> ta {peerServerPubKey = prServerKey}) <$> thAuth params
serverThParams = smpTHParamsSetVersion prVersion params {sessionId = prSessionId, thAuth = serverThAuth}
(cmdPubKey, cmdPrivKey) <- atomically $ C.generateKeyPair @'C.X25519 g
nonce@(C.CbNonce corrId) <- atomically $ C.randomCbNonce g
let TransmissionForAuth {tToSend} = encodeTransmissionForAuth serverThParams (CorrId corrId, sId, Cmd SSender $ SEND noMsgFlags "hello")
b <- case batchTransmissions serverThParams [Right (Nothing, tToSend)] of
TBTransmission s _ : _ -> pure s
TBTransmissions s _ _ : _ -> pure s
_ -> fail "pfwdTransmission: batch"
et <- either (fail . show) (pure . EncTransmission) $ C.cbEncrypt (C.dh' prServerKey cmdPrivKey) nonce b paddedProxiedTLength
let TransmissionForAuth {tToSend = t} = encodeTransmissionForAuth params (CorrId corrId, EntityId prSessionId, Cmd SProxiedClient $ PFWD prVersion cmdPubKey et)
pure t
#endif
proxyPhases :: [String]
proxyPhases = ["proxyfwd", "proxytmo", "proxychurn", "subtmo", "conclimit", "fastfwd", "msgqfill", "pfwdbig"]
proxyPhases =["proxyfwd", "proxytmo", "proxychurn", "subtmo", "conclimit", "fastfwd", "msgqfill", "pfwdbig"]
-- proxy topologies use the test databases; create them for the run and drop them after
pgBracket :: IO a -> IO a