mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-06 14:17:20 +00:00
bench: add production-shaped memory phase
This commit is contained in:
+138
-1
@@ -23,7 +23,7 @@
|
||||
--
|
||||
-- Single-server phases (server on testPort):
|
||||
-- plain | svc | svcrace | ntf | conc | svcsubs | getp | stuck | certchurn | link | ntfexp
|
||||
-- ntfloop | ntfdeliver | subslice | conns | load | cpsave -- PostgreSQL build only
|
||||
-- ntfloop | ntfdeliver | subslice | conns | load | prodmix | cpsave -- PostgreSQL build only
|
||||
-- tlsstall | tlshalf | tlschurn | tlspartial -- TLS/TCP stack
|
||||
--
|
||||
-- Two-server phases (proxy on testPort, lagged destination relay on testPort2):
|
||||
@@ -48,6 +48,7 @@ import Data.ByteString.Char8 (ByteString)
|
||||
import Data.Int (Int64)
|
||||
import Data.Foldable (toList)
|
||||
import Data.List.NonEmpty (NonEmpty (..), fromList)
|
||||
import qualified Data.Map.Strict as M
|
||||
import Data.Maybe (fromMaybe)
|
||||
import Data.Time.Clock (diffUTCTime, getCurrentTime)
|
||||
import qualified Data.X509.Validation as XV
|
||||
@@ -685,6 +686,140 @@ runLoadClient g = do
|
||||
n <- readTVarIO ops
|
||||
putStrLn $ "DONE " <> show n <> " " <> show (realToFrac (diffUTCTime end start) :: Double)
|
||||
|
||||
-- Production-shaped server memory: C connections (PROD_CONNS, default 1000) each subscribing
|
||||
-- Q queues (PROD_QUEUES, default 20) in batches of up to 100 per block, as the agent resubscribes,
|
||||
-- plus PROD_NEW subscriptions per connection created with NEW; then PROD_SEC seconds (default 60)
|
||||
-- of traffic where PROD_SENDERS connections (default 16) send to random queues and every recipient
|
||||
-- acknowledges its messages. The client runs in a child process (phase prodclient); this process
|
||||
-- reports the server's live heap, memory in use and RSS when idle and at peak under traffic.
|
||||
runProdMix :: IO ()
|
||||
runProdMix = do
|
||||
base <- liveBytesMiB
|
||||
rssBase <- rssMiB
|
||||
exe <- getExecutablePath
|
||||
(Just hIn, Just hOut, _, ph) <- createProcess (proc exe ["prodclient", "0"]) {std_in = CreatePipe, std_out = CreatePipe}
|
||||
hSetBuffering hIn LineBuffering
|
||||
(conns, subs) <- awaitLine hOut "READY" >>= \case
|
||||
[c, n] -> pure (read c, read n) :: IO (Int, Int)
|
||||
_ -> fail "bad READY line"
|
||||
idle <- liveBytesMiB
|
||||
idleStats <- getRTSStats
|
||||
idleRss <- rssMiB
|
||||
printf "prodmix: IDLE connections=%d subscriptions=%d live=%.1f MiB (%.2f KiB/conn) mem_in_use=%.1f MiB rss=%.1f MiB (base live=%.1f rss=%.1f)\n"
|
||||
conns subs (idle - base) ((idle - base) * 1024 / fromIntegral conns) (mib $ gcdetails_mem_in_use_bytes $ gc idleStats) idleRss base rssBase
|
||||
hPutStrLn hIn "GO"
|
||||
peak <- newTVarIO (0 :: Double, 0 :: Double, 0 :: Double)
|
||||
done <- newEmptyTMVarIO
|
||||
let sampler = do
|
||||
threadDelay 2000000
|
||||
st <- getRTSStats
|
||||
r <- rssMiB
|
||||
atomically $ modifyTVar' peak $ \(m, l, rs) -> (max m (mib $ gcdetails_mem_in_use_bytes $ gc st), max l (mib $ gcdetails_live_bytes $ gc st), max rs r)
|
||||
atomically (isEmptyTMVar done) >>= (`when` sampler)
|
||||
msgs <- withAsync sampler $ \_ -> do
|
||||
[n] <- awaitLine hOut "DONE"
|
||||
atomically $ putTMVar done ()
|
||||
pure (read n :: Int)
|
||||
(pkMem, pkLive, pkRss) <- readTVarIO peak
|
||||
end <- getRTSStats
|
||||
after <- liveBytesMiB
|
||||
printf "prodmix: TRAFFIC messages=%d peak_mem_in_use=%.1f MiB peak_live_after_gc=%.1f MiB peak_rss=%.1f MiB live_after=%.1f MiB major_gcs=%d gc_cpu=%.1fs mutator_cpu=%.1fs\n"
|
||||
msgs pkMem pkLive pkRss (after - base) (major_gcs end - major_gcs idleStats)
|
||||
(fromIntegral (gc_cpu_ns end - gc_cpu_ns idleStats) / 1e9 :: Double) (fromIntegral (mutator_cpu_ns end - mutator_cpu_ns idleStats) / 1e9 :: Double)
|
||||
hClose hIn
|
||||
void $ waitForProcess ph
|
||||
where
|
||||
mib :: Word64 -> Double
|
||||
mib b = fromIntegral b / (1024 * 1024)
|
||||
awaitLine :: Handle -> String -> IO [String]
|
||||
awaitLine h tag = hGetLine h >>= \l -> case words l of
|
||||
(t : rest) | t == tag -> pure rest
|
||||
_ -> awaitLine h tag
|
||||
|
||||
-- resident set size of this process, from /proc
|
||||
rssMiB :: IO Double
|
||||
rssMiB = do
|
||||
ls <- lines <$> readFile "/proc/self/status"
|
||||
pure $ case [w | l <- ls, ("VmRSS:" : w : _) <- [words l]] of
|
||||
(kb : _) -> read kb / 1024
|
||||
_ -> 0
|
||||
|
||||
runProdClient :: TVar ChaChaDRG -> IO ()
|
||||
runProdClient g = do
|
||||
hSetBuffering stdout LineBuffering
|
||||
let env k d = fromMaybe d . (>>= readMaybe) <$> lookupEnv k
|
||||
conns <- env "PROD_CONNS" 1000
|
||||
nq <- env "PROD_QUEUES" 20
|
||||
newPer <- env "PROD_NEW" 2
|
||||
secs <- env "PROD_SEC" 60
|
||||
senders <- env "PROD_SENDERS" 16
|
||||
let creators = 8 :: Int
|
||||
total = conns * nq
|
||||
-- queues are created on a few connections first, as if by earlier sessions
|
||||
created <- forConcurrently ([0 .. creators - 1] :: [Int]) $ \w -> benchClient $ \h ->
|
||||
forM [i | i <- [0 .. total - 1], i `mod` creators == w] $ \i -> do
|
||||
(rPub, rKey, dhPub) <- genKeys g
|
||||
Resp _ _ (Ids rId sId _) <- signSendRecv h rKey (B.pack $ "c" <> show i, NoEntity, New0 rPub dhPub)
|
||||
pure (rId, sId, rKey)
|
||||
let qs = concat created
|
||||
perConn = chunksOf nq qs
|
||||
subscribed <- newTVarIO (0 :: Int)
|
||||
go <- newEmptyTMVarIO
|
||||
stop <- newEmptyTMVarIO
|
||||
sIdsVar <- newTVarIO ([] :: [SenderId])
|
||||
received <- newTVarIO (0 :: Int)
|
||||
let recipient cqs = benchClient $ \h -> do
|
||||
forM_ (chunksOf 100 cqs) $ \b -> do
|
||||
let sub (rId, _, rKey) = Right $ signTransmission h rKey Nothing (B.pack "s", rId, SUB)
|
||||
void $ tPut h (fromList $ map sub b)
|
||||
awaitResponses h (length b)
|
||||
newQs <- forM ([1 .. newPer] :: [Int]) $ \i -> do
|
||||
(rPub, rKey, dhPub) <- genKeys g
|
||||
Resp _ _ (Ids rId sId _) <- signSendRecv h rKey (B.pack $ "n" <> show i, NoEntity, New rPub dhPub)
|
||||
pure (rId, sId, rKey)
|
||||
let keys = M.fromList [(rId, rKey) | (rId, _, rKey) <- cqs ++ newQs]
|
||||
atomically $ do
|
||||
modifyTVar' sIdsVar ([sId | (_, sId, _) <- cqs ++ newQs] ++)
|
||||
modifyTVar' subscribed (+ (length cqs + length newQs))
|
||||
-- acknowledge every delivered message until the run is released
|
||||
withAsync (forever $ tGetClient h >>= mapM_ (ack h keys)) $ \_ -> atomically $ readTMVar stop
|
||||
ack h keys = \case
|
||||
(_, rId, Right (MSG RcvMessage {msgId})) | Just rKey <- M.lookup rId keys -> do
|
||||
let t = signTransmission h rKey Nothing (B.pack "a", rId, ACK msgId)
|
||||
void $ tPut h [Right t]
|
||||
atomically $ modifyTVar' received (+ 1)
|
||||
_ -> pure ()
|
||||
sender w = benchClient $ \h -> do
|
||||
atomically $ readTMVar go
|
||||
sIds <- readTVarIO sIdsVar
|
||||
let n = length sIds
|
||||
loop i t0 = do
|
||||
t <- getCurrentTime
|
||||
when (diffUTCTime t t0 < fromIntegral (secs :: Int)) $ do
|
||||
let sId = sIds !! ((i * 7919 + w * 104729) `mod` n)
|
||||
void $ sendRecv h (Nothing, B.pack (show i), sId, _SEND "hello")
|
||||
loop (i + 1) t0
|
||||
getCurrentTime >>= loop 0
|
||||
withAsync (forConcurrently_ perConn recipient) $ \rs -> do
|
||||
atomically $ readTVar subscribed >>= \n -> when (n < total + conns * newPer) retry
|
||||
putStrLn $ "READY " <> show conns <> " " <> show (total + conns * newPer)
|
||||
"GO" <- getLine
|
||||
atomically $ putTMVar go ()
|
||||
forConcurrently_ ([1 .. senders] :: [Int]) sender
|
||||
n <- readTVarIO received
|
||||
putStrLn $ "DONE " <> show n
|
||||
void (E.try getLine :: IO (Either E.IOException String))
|
||||
atomically $ putTMVar stop ()
|
||||
wait rs
|
||||
where
|
||||
chunksOf k xs = case splitAt k xs of
|
||||
(c, []) -> [c | not (null c)]
|
||||
(c, rest) -> c : chunksOf k rest
|
||||
awaitResponses :: H -> Int -> IO ()
|
||||
awaitResponses h n = when (n > 0) $ do
|
||||
rs :: NonEmpty (Transmission (Either ErrorType BrokerMsg)) <- tGetClient h
|
||||
awaitResponses h (n - length rs)
|
||||
|
||||
-- Client side of subslice: creates the subscriptions, prints "READY <subscribed>", and holds the
|
||||
-- connection until stdin is closed.
|
||||
runSubClient :: TVar ChaChaDRG -> Int -> IO ()
|
||||
@@ -1411,6 +1546,8 @@ main = do
|
||||
else if phase == "connclient" then runConnClient iters
|
||||
else if phase == "load" then postgressBracket benchDBConnectInfo $ withSmpServerConfigOn (transport @TLS) benchPgCfg benchPort $ \_ -> runLoad
|
||||
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
|
||||
#endif
|
||||
else withSmpServerConfigOn (transport @TLS) srvCfg testPort $ \_ -> settle leakDiagSec $ do
|
||||
threadDelay 250000
|
||||
|
||||
Reference in New Issue
Block a user