From 84d91bd445ea6567e47ad982442d5fef169d9f95 Mon Sep 17 00:00:00 2001 From: shum Date: Wed, 30 Sep 2026 13:45:07 +0000 Subject: [PATCH] bench: measure proxy mesh servers separately --- bench/MemBench.hs | 311 +++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 306 insertions(+), 5 deletions(-) diff --git a/bench/MemBench.hs b/bench/MemBench.hs index 69e2ca2ce..85c672976 100644 --- a/bench/MemBench.hs +++ b/bench/MemBench.hs @@ -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