mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-06 05:37:17 +00:00
bench: add server-only memory and load phases
This commit is contained in:
+401
-10
@@ -23,6 +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
|
||||
-- tlsstall | tlshalf | tlschurn | tlspartial -- TLS/TCP stack
|
||||
--
|
||||
-- Two-server phases (proxy on testPort, lagged destination relay on testPort2):
|
||||
@@ -59,7 +60,7 @@ 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, serverClientConcurrency, smpAgentCfg))
|
||||
import Simplex.Messaging.Server.Env.STM (AStoreType (..), ServerConfig (controlPort, controlPortAdminAuth, maxJournalMsgCount, msgQueueQuota, notificationExpiration, ntfDeliveryInterval, serverClientConcurrency, smpAgentCfg, storeNtfsFile))
|
||||
import Simplex.Messaging.Server.Expiration (ExpirationConfig (..))
|
||||
import Simplex.Messaging.Server.MsgStore.Types (SMSType (..), SQSType (..))
|
||||
import Simplex.Messaging.Transport
|
||||
@@ -67,7 +68,25 @@ import Simplex.Messaging.Transport.Client (TransportClientConfig (..), defaultTr
|
||||
import Simplex.Messaging.Transport.Credentials (genCredentials, tlsCredentials)
|
||||
import Simplex.Messaging.Version (mkVersionRange)
|
||||
import System.Environment (getArgs, lookupEnv, setEnv)
|
||||
import System.IO (BufferMode (..), IOMode (..), hClose, hGetLine, hPutStrLn, hSetBuffering, hSetNewlineMode, universalNewlineMode)
|
||||
import System.IO (BufferMode (..), IOMode (..), hClose, hGetLine, hPutStrLn, hSetBuffering, hSetNewlineMode, universalNewlineMode, withFile)
|
||||
#if defined(dbServerPostgres)
|
||||
import Database.PostgreSQL.Simple (ConnectInfo (..))
|
||||
import Network.Socket (ServiceName)
|
||||
import Simplex.Messaging.Agent.Store.Postgres.Options (DBOpts (..))
|
||||
import Simplex.Messaging.Agent.Store.Shared (MigrationConfirmation (..))
|
||||
import Simplex.Messaging.Server.Env.STM (ServerStoreCfg (..), serverStoreCfg)
|
||||
import Simplex.Messaging.Server.QueueStore.Postgres.Config (PostgresStoreCfg (..))
|
||||
import qualified Data.ByteString.Builder as BLD
|
||||
import Data.List (unfoldr)
|
||||
import Data.Time.Clock.System (getSystemTime)
|
||||
import Data.Word (Word64)
|
||||
import Simplex.Messaging.Encoding.String (strEncode)
|
||||
import Simplex.Messaging.Server.NtfStore (MsgNtf (..), NtfLogRecord (..))
|
||||
import System.Directory (createDirectoryIfMissing, doesFileExist, removeFile)
|
||||
import System.Environment (getExecutablePath)
|
||||
import System.IO (Handle, stdout)
|
||||
import System.Process (CreateProcess (..), StdStream (..), callProcess, createProcess, proc, waitForProcess)
|
||||
#endif
|
||||
import System.Mem (performMajorGC)
|
||||
import System.Timeout (timeout)
|
||||
import Text.Printf (printf)
|
||||
@@ -126,10 +145,14 @@ serviceSignSendRecv h pk serviceKey t = do
|
||||
pure r
|
||||
|
||||
signSendRecv_ :: forall p. PartyI p => H -> C.APrivateAuthKey -> Maybe C.PrivateKeyEd25519 -> (ByteString, EntityId, Command p) -> IO (NonEmpty (Transmission (Either ErrorType BrokerMsg)))
|
||||
signSendRecv_ h@THandle {params} (C.APrivateAuthKey a pk) serviceKey_ (corrId, qId, cmd) = do
|
||||
let TransmissionForAuth {tForAuth, tToSend} = encodeTransmissionForAuth params (CorrId corrId, qId, cmd)
|
||||
Right () <- tPut1 h (authorize tForAuth, tToSend)
|
||||
signSendRecv_ h pk serviceKey_ t = do
|
||||
Right () <- tPut1 h (signTransmission h pk serviceKey_ t)
|
||||
tGetClient h
|
||||
|
||||
signTransmission :: forall p. PartyI p => H -> C.APrivateAuthKey -> Maybe C.PrivateKeyEd25519 -> (ByteString, EntityId, Command p) -> SentRawTransmission
|
||||
signTransmission THandle {params} (C.APrivateAuthKey a pk) serviceKey_ (corrId, qId, cmd) =
|
||||
let TransmissionForAuth {tForAuth, tToSend} = encodeTransmissionForAuth params (CorrId corrId, qId, cmd)
|
||||
in (authorize tForAuth, tToSend)
|
||||
where
|
||||
authorize t = (,(`C.sign'` t) <$> serviceKey_) <$> case a of
|
||||
C.SEd25519 -> Just . TASignature . C.ASignature C.SEd25519 $ C.sign' pk t'
|
||||
@@ -171,6 +194,10 @@ liveBytesMiB = do
|
||||
s <- getRTSStats
|
||||
pure $ fromIntegral (gcdetails_live_bytes (gc s)) / (1024 * 1024)
|
||||
|
||||
-- memory in partially used blocks (mostly pinned blocks kept by a few live objects), as of the last GC
|
||||
fragmentationMiB :: IO Double
|
||||
fragmentationMiB = (/ (1024 * 1024)) . fromIntegral . gcdetails_block_fragmentation_bytes . gc <$> getRTSStats
|
||||
|
||||
report :: String -> Int -> Double -> Double -> IO ()
|
||||
report phase i base cur =
|
||||
printf "%-8s iter=%7d live=%9.1f MiB delta=%+9.1f MiB (%+.4f KiB/iter)\n"
|
||||
@@ -259,6 +286,178 @@ runNtf g iters cp =
|
||||
Resp _ _ OK <- signSendRecv recip rKey (corr, rId, ACK mId)
|
||||
pure ()
|
||||
|
||||
#if defined(dbServerPostgres)
|
||||
-- Notification delivery loop cost with N stored notifiers.
|
||||
--
|
||||
-- deliverNtfsThread (Server.hs:380) runs every ntfDeliveryInterval, 1.5s in production. Each tick
|
||||
-- takes every NtfStore key, including keys whose list is already empty, and resolves their services
|
||||
-- with one `notifier_id IN ?` query read in full (QueueStore/Postgres.hs:538).
|
||||
--
|
||||
-- Seeds N notifier queues with SQL and N stored notifications through the server's own restore
|
||||
-- file, then samples RTS counters with no client traffic. NTFLOOP_DELIVER=0 disables the loop, so
|
||||
-- the difference between the two runs is the loop's own cost.
|
||||
runNtfLoop :: Int -> IO ()
|
||||
runNtfLoop n = postgressBracket benchDBConnectInfo $ do
|
||||
deliver <- (/= Just "0") <$> lookupEnv "NTFLOOP_DELIVER"
|
||||
secs <- fromMaybe 30 . (>>= readMaybe) <$> lookupEnv "NTFLOOP_SEC"
|
||||
let productionInterval = 1500000
|
||||
disabledInterval = 86400 * 1000000
|
||||
srvCfg = updateCfg benchPgCfg $ \c ->
|
||||
c {storeNtfsFile = Just testStoreNtfsFile, ntfDeliveryInterval = if deliver then productionInterval else disabledInterval}
|
||||
createDirectoryIfMissing True "tests/tmp"
|
||||
doesFileExist testStoreNtfsFile >>= (`when` removeFile testStoreNtfsFile)
|
||||
-- the first start runs the migrations that create msg_queues
|
||||
withSmpServerConfigOn (transport @TLS) srvCfg benchPort $ \_ -> threadDelay 500000
|
||||
seedNotifierQueues n
|
||||
writeNtfsFile n
|
||||
printf "ntfloop: n=%d deliver=%s window=%ds\n" n (show deliver) secs
|
||||
withSmpServerConfigOn (transport @TLS) srvCfg benchPort $ \_ -> do
|
||||
performMajorGC
|
||||
base <- getRTSStats
|
||||
let mib :: Word64 -> Double
|
||||
mib b = fromIntegral b / (1024 * 1024)
|
||||
cpuNs s = mutator_cpu_ns s + gc_cpu_ns s
|
||||
printf "ntfloop: after restore live=%.1f MiB mem_in_use=%.1f MiB max_mem_in_use=%.1f MiB\n"
|
||||
(mib $ gcdetails_live_bytes $ gc base) (mib $ gcdetails_mem_in_use_bytes $ gc base) (mib $ max_mem_in_use_bytes base)
|
||||
peak <- newTVarIO (0 :: Word64)
|
||||
forM_ ([1 .. secs] :: [Int]) $ \t -> do
|
||||
threadDelay 1000000
|
||||
s <- getRTSStats
|
||||
atomically $ modifyTVar' peak (max $ gcdetails_mem_in_use_bytes $ gc s)
|
||||
when (t `mod` 5 == 0) $
|
||||
printf "ntfloop: t=%3ds alloc=%8.1f MiB/s cpu=%5.2f cores mem_in_use=%8.1f MiB major_gcs=%d\n"
|
||||
t
|
||||
(mib (allocated_bytes s - allocated_bytes base) / fromIntegral t)
|
||||
(fromIntegral (cpuNs s - cpuNs base) / (fromIntegral t * 1e9) :: Double)
|
||||
(mib $ gcdetails_mem_in_use_bytes $ gc s)
|
||||
(major_gcs s - major_gcs base)
|
||||
end <- getRTSStats
|
||||
pk <- readTVarIO peak
|
||||
performMajorGC
|
||||
after <- getRTSStats
|
||||
printf "ntfloop: SUMMARY n=%d deliver=%s alloc=%.1f MiB/s (%.1f MiB per 1.5s tick) cpu=%.2f cores peak_mem_in_use=%.1f MiB max_mem_in_use=%.1f MiB live_end=%.1f MiB\n"
|
||||
n
|
||||
(show deliver)
|
||||
(mib (allocated_bytes end - allocated_bytes base) / fromIntegral secs)
|
||||
(1.5 * mib (allocated_bytes end - allocated_bytes base) / fromIntegral secs)
|
||||
(fromIntegral (cpuNs end - cpuNs base) / (fromIntegral secs * 1e9) :: Double)
|
||||
(mib pk)
|
||||
(mib $ max_mem_in_use_bytes end)
|
||||
(mib $ gcdetails_live_bytes $ gc after)
|
||||
|
||||
-- Its own database, user and port, so a bench run cannot collide with a concurrent test-suite run,
|
||||
-- which binds testPort and drops and recreates test_server_db and test_server_user.
|
||||
benchDBConnectInfo :: ConnectInfo
|
||||
benchDBConnectInfo = testServerDBConnectInfo {connectUser = "mem_bench_user", connectDatabase = "mem_bench_db"}
|
||||
|
||||
benchPort :: ServiceName
|
||||
benchPort = "15001"
|
||||
|
||||
benchPgCfg :: AServerConfig
|
||||
benchPgCfg = case cfgMS (ASType SQSPostgres SMSPostgres) of
|
||||
ASrvCfg SQSPostgres SMSPostgres c -> ASrvCfg SQSPostgres SMSPostgres c {serverStoreCfg = SSCDatabase storeCfg}
|
||||
c -> c
|
||||
where
|
||||
storeCfg =
|
||||
PostgresStoreCfg
|
||||
{ dbOpts = testStoreDBOpts {connstr = B.pack $ "postgresql://" <> connectUser benchDBConnectInfo <> "@/" <> connectDatabase benchDBConnectInfo},
|
||||
dbStoreLogPath = Nothing,
|
||||
confirmMigrations = MCYesUp,
|
||||
deletedTTL = 86400
|
||||
}
|
||||
|
||||
benchClient :: (H -> IO a) -> IO a
|
||||
benchClient client = do
|
||||
Right useHost <- pure $ chooseTransportHost defaultNetworkConfig testHost
|
||||
testSMPClient_ useHost benchPort supportedClientSMPRelayVRange Nothing client
|
||||
|
||||
-- NtfStore keys left behind by delivered notifications.
|
||||
--
|
||||
-- With a subscribed ntf client, each tick flushes a notifier's list to [] (addNtfs) but keeps the
|
||||
-- key, and every key is re-sent to getQueueNtfServices on every tick until deleteExpiredNtfs
|
||||
-- removes it (hourly on this branch, never on stable). Creates `iters` notifier queues, subscribes
|
||||
-- one ntf client to all of them, delivers one notification to each, then samples RTS counters while
|
||||
-- nothing else happens. LEAKDIAG ntfStore_keys shows the retained keys.
|
||||
runNtfDeliver :: TVar ChaChaDRG -> Int -> IO ()
|
||||
runNtfDeliver g n = do
|
||||
secs <- fromMaybe 30 . (>>= readMaybe) <$> lookupEnv "NTFLOOP_SEC"
|
||||
benchClient $ \recip -> benchClient $ \sndr -> benchClient $ \nh -> do
|
||||
sIds <- forM ([1 .. n] :: [Int]) $ \i -> do
|
||||
(rPub, rKey, dhPub) <- genKeys g
|
||||
(nPub, nKey) <- atomically $ C.generateAuthKeyPair C.SEd25519 g
|
||||
(rcvNtfPubDh, _dh :: C.PrivateKeyX25519) <- atomically $ C.generateKeyPair g
|
||||
let corr = B.pack (show i)
|
||||
Resp _ _ (Ids rId sId _) <- signSendRecv recip rKey (corr, NoEntity, New0 rPub dhPub)
|
||||
Resp _ _ (NID nId _) <- signSendRecv recip rKey (corr, rId, NKEY nPub rcvNtfPubDh)
|
||||
Resp _ _ (SOK Nothing) <- signSendRecv nh nKey (corr, nId, NSUB)
|
||||
pure sId
|
||||
forM_ (zip [1 :: Int ..] sIds) $ \(i, sId) -> do
|
||||
Resp _ _ OK <- sendRecv sndr (Nothing, B.pack (show i), sId, _SEND' "hi")
|
||||
pure ()
|
||||
let awaitNtfs k = when (k > 0) $ do
|
||||
rs :: NonEmpty (Transmission (Either ErrorType BrokerMsg)) <- tGetClient nh
|
||||
awaitNtfs (k - length [() | (_, _, Right NMSG {}) <- toList rs])
|
||||
awaitNtfs n
|
||||
printf "ntfdeliver: n=%d notifications delivered, sampling %ds\n" n secs
|
||||
performMajorGC
|
||||
base <- getRTSStats
|
||||
threadDelay $ secs * 1000000
|
||||
end <- getRTSStats
|
||||
let mib :: Word64 -> Double
|
||||
mib b = fromIntegral b / (1024 * 1024)
|
||||
cpuNs st = mutator_cpu_ns st + gc_cpu_ns st
|
||||
printf "ntfdeliver: SUMMARY n=%d idle alloc=%.1f MiB/s cpu=%.2f cores\n"
|
||||
n
|
||||
(mib (allocated_bytes end - allocated_bytes base) / fromIntegral secs)
|
||||
(fromIntegral (cpuNs end - cpuNs base) / (fromIntegral secs * 1e9) :: Double)
|
||||
|
||||
-- Control-port `save` on the PostgreSQL store.
|
||||
--
|
||||
-- CPSave (Server.hs:1153) runs saveServer False, and its closeMsgStore closes the DB pool
|
||||
-- (closeDBStore drains and closes every pooled connection). The server keeps running, so the next
|
||||
-- withConnectionPool waits forever on the empty pool while holding dbSem.
|
||||
runCpSave :: TVar ChaChaDRG -> IO ()
|
||||
runCpSave g = benchClient $ \h -> do
|
||||
(rPub, rKey, dhPub) <- genKeys g
|
||||
Resp _ _ (Ids _ _ _) <- signSendRecv h rKey ("1", NoEntity, New rPub dhPub)
|
||||
putStrLn "cpsave: NEW before save: IDS"
|
||||
r <- cpCommand benchCpPort "save" 2 30000000
|
||||
printf "cpsave: save replied %s\n" (show r)
|
||||
(rPub', rKey', dhPub') <- genKeys g
|
||||
res <- timeout 10000000 $ signSendRecv h rKey' ("2", NoEntity, New rPub' dhPub')
|
||||
putStrLn $ "cpsave: NEW after save: " <> maybe "no response in 10s (DB access blocked)" (const "responded") res
|
||||
|
||||
benchCpPort :: ServiceName
|
||||
benchCpPort = "15010"
|
||||
|
||||
-- The IDs are a tag byte followed by i in 23 big-endian bytes, the same bytes the SQL builds with
|
||||
-- lpad(to_hex(i), 46, '0'), so file entries and rows match without passing IDs between them.
|
||||
seqId :: Char -> Int -> ByteString
|
||||
seqId tag i = B.cons tag $ B.replicate (23 - B.length be) '\0' <> be
|
||||
where
|
||||
be = B.pack . reverse $ unfoldr (\x -> if x == 0 then Nothing else Just (toEnum (x `mod` 256), x `div` 256)) i
|
||||
|
||||
seedNotifierQueues :: Int -> IO ()
|
||||
seedNotifierQueues n =
|
||||
callProcess "psql" ["-q", "-v", "ON_ERROR_STOP=1", "-U", "postgres", "-d", connectDatabase benchDBConnectInfo, "-c", sql]
|
||||
where
|
||||
hexId :: String -> String
|
||||
hexId tag = "decode('" <> tag <> "' || lpad(to_hex(i), 46, '0'), 'hex')"
|
||||
sql =
|
||||
"INSERT INTO smp_server.msg_queues (recipient_id, recipient_keys, rcv_dh_secret, sender_id, notifier_id, notifier_key, rcv_ntf_dh_secret, status, updated_at) SELECT "
|
||||
<> hexId "01" <> ", '\\x00', '\\x00', " <> hexId "02" <> ", " <> hexId "03" <> ", '\\x00', '\\x00', 'active', 0 FROM generate_series(1, "
|
||||
<> show n <> ") AS i"
|
||||
|
||||
-- MsgNtf field sizes match mkMessageNotification: 24-byte msg id and nonce, 128-byte padded meta plus 16-byte tag
|
||||
writeNtfsFile :: Int -> IO ()
|
||||
writeNtfsFile n = do
|
||||
ts <- getSystemTime
|
||||
let ntf = MsgNtf {ntfMsgId = B.replicate 24 'm', ntfTs = ts, ntfNonce = C.cbNonce (B.replicate 24 'n'), ntfEncMeta = B.replicate 144 'e'}
|
||||
withFile testStoreNtfsFile WriteMode $ \h ->
|
||||
forM_ ([1 .. n] :: [Int]) $ \i ->
|
||||
BLD.hPutBuilder h $ BLD.byteString (strEncode $ NLRv1 (EntityId $ seqId '\3' i) ntf) <> BLD.char8 '\n'
|
||||
#endif
|
||||
|
||||
-- reusable steps ------------------------------------------------------------
|
||||
|
||||
genKeys :: TVar ChaChaDRG -> IO (RcvPublicAuthKey, C.APrivateAuthKey, RcvPublicDhKey)
|
||||
@@ -356,6 +555,181 @@ runGet g iters cp =
|
||||
Resp _ _ OK <- signSendRecv recip rKey (corr, rId, DEL)
|
||||
pure ()
|
||||
|
||||
#if defined(dbServerPostgres)
|
||||
-- Server memory per subscription.
|
||||
--
|
||||
-- The server stores the parsed entity ID as the subscription key (Server.hs:1839, :1971, and
|
||||
-- queueSubscribers via subQ). The ID is an attoparsec slice of the decrypted ~16 KB block, so while
|
||||
-- any key from a block is subscribed, the whole block stays live. NEW subscriptions are keyed by the
|
||||
-- server-generated randomId instead.
|
||||
--
|
||||
-- The client runs in a child process (phase subclient), so the live bytes measured here are the
|
||||
-- server's alone. SUBMODE=sub (default) creates `iters` queues without subscribing, then subscribes
|
||||
-- them from a second connection, SUBBATCH per block (default 1); SUBKEEP=k keeps k subscriptions per
|
||||
-- block and deletes the other queues. SUBMODE=new subscribes with NEW; NEWSUB=0 only creates.
|
||||
runSubSlice :: Int -> IO ()
|
||||
runSubSlice iters = do
|
||||
mode <- fromMaybe "sub" <$> lookupEnv "SUBMODE"
|
||||
base <- liveBytesMiB
|
||||
fragBase <- fragmentationMiB
|
||||
report "subslice" 0 base base
|
||||
exe <- getExecutablePath
|
||||
(Just hIn, Just hOut, _, ph) <- createProcess (proc exe ["subclient", show iters]) {std_in = CreatePipe, std_out = CreatePipe}
|
||||
kept <- awaitReady hOut
|
||||
cur <- liveBytesMiB
|
||||
frag <- fragmentationMiB
|
||||
report "subslice" iters base cur
|
||||
printf "subslice: SUMMARY mode=%s subscribed=%d server_live_delta=%.1f MiB per_subscription=%.2f KiB block_fragmentation_per_subscription=%.2f KiB\n"
|
||||
mode kept (cur - base) ((cur - base) * 1024 / fromIntegral (max 1 kept)) ((frag - fragBase) * 1024 / fromIntegral (max 1 kept))
|
||||
hClose hIn
|
||||
void $ waitForProcess ph
|
||||
-- the server drops the connection's subscriptions on disconnect
|
||||
threadDelay 1000000
|
||||
after <- liveBytesMiB
|
||||
printf "subslice: after disconnect server_live_delta=%.1f MiB\n" (after - base)
|
||||
where
|
||||
awaitReady :: Handle -> IO Int
|
||||
awaitReady h = hGetLine h >>= \l -> case words l of
|
||||
["READY", k] -> pure $ read k
|
||||
_ -> awaitReady h
|
||||
|
||||
-- Server memory per idle connection. The child (phase connclient) opens `iters` SMP connections,
|
||||
-- completes the handshakes, prints READY and holds them until stdin closes, so the live bytes
|
||||
-- measured here are the server's alone.
|
||||
runConns :: Int -> IO ()
|
||||
runConns iters = do
|
||||
base <- liveBytesMiB
|
||||
report "conns" 0 base base
|
||||
exe <- getExecutablePath
|
||||
(Just hIn, Just hOut, _, ph) <- createProcess (proc exe ["connclient", show iters]) {std_in = CreatePipe, std_out = CreatePipe}
|
||||
n <- awaitReadyLine hOut
|
||||
cur <- liveBytesMiB
|
||||
report "conns" n base cur
|
||||
printf "conns: SUMMARY connections=%d server_live_delta=%.1f MiB per_connection=%.2f KiB\n" n (cur - base) ((cur - base) * 1024 / fromIntegral (max 1 n))
|
||||
hClose hIn
|
||||
void $ waitForProcess ph
|
||||
threadDelay 2000000
|
||||
after <- liveBytesMiB
|
||||
printf "conns: after disconnect server_live_delta=%.1f MiB\n" (after - base)
|
||||
|
||||
runConnClient :: Int -> IO ()
|
||||
runConnClient iters = do
|
||||
hSetBuffering stdout LineBuffering
|
||||
opened <- newTVarIO (0 :: Int)
|
||||
release <- newEmptyTMVarIO
|
||||
let hold = benchClient $ \_ -> do
|
||||
atomically $ modifyTVar' opened (+ 1)
|
||||
atomically $ readTMVar release
|
||||
withAsync (forConcurrently_ ([1 .. iters] :: [Int]) $ \_ -> hold) $ \a -> do
|
||||
atomically $ readTVar opened >>= \n -> when (n < iters) retry
|
||||
putStrLn $ "READY " <> show iters
|
||||
void (E.try getLine :: IO (Either E.IOException String))
|
||||
atomically $ putTMVar release ()
|
||||
wait a
|
||||
|
||||
awaitReadyLine :: Handle -> IO Int
|
||||
awaitReadyLine h = hGetLine h >>= \l -> case words l of
|
||||
["READY", k] -> pure $ read k
|
||||
_ -> awaitReadyLine h
|
||||
|
||||
-- Server CPU under a fixed mixed workload, to compare RTS settings. The child (phase loadclient)
|
||||
-- runs LOAD_WORKERS workers (default 16) for LOAD_SEC seconds (default 60); each worker reconnects
|
||||
-- every 50 steps so TLS handshakes are part of the load, and every third step also uses notifications.
|
||||
-- The server keeps the RTS flags given to this process; the child keeps its defaults.
|
||||
runLoad :: IO ()
|
||||
runLoad = do
|
||||
performMajorGC
|
||||
base <- getRTSStats
|
||||
exe <- getExecutablePath
|
||||
(_, Just hOut, _, ph) <- createProcess (proc exe ["loadclient", "0"]) {std_out = CreatePipe}
|
||||
(ops, secs) <- awaitDone hOut
|
||||
end <- getRTSStats
|
||||
void $ waitForProcess ph
|
||||
let cpu f = fromIntegral (f end - f base) / 1e9 :: Double
|
||||
mutCpu = cpu mutator_cpu_ns
|
||||
gcCpu = cpu gc_cpu_ns
|
||||
printf "load: SUMMARY ops=%d secs=%.1f ops_per_sec=%.1f server_cpu_ms_per_op=%.3f mutator_cpu=%.1fs gc_cpu=%.1fs gc_share=%.1f%% gcs=%d major_gcs=%d max_mem_in_use=%.1f MiB\n"
|
||||
ops secs (fromIntegral ops / secs) ((mutCpu + gcCpu) * 1000 / fromIntegral ops) mutCpu gcCpu (100 * gcCpu / (mutCpu + gcCpu))
|
||||
(gcs end - gcs base) (major_gcs end - major_gcs base) (fromIntegral (max_mem_in_use_bytes end) / (1024 * 1024) :: Double)
|
||||
where
|
||||
awaitDone :: Handle -> IO (Int, Double)
|
||||
awaitDone h = hGetLine h >>= \l -> case words l of
|
||||
["DONE", o, t] -> pure (read o, read t)
|
||||
_ -> awaitDone h
|
||||
|
||||
runLoadClient :: TVar ChaChaDRG -> IO ()
|
||||
runLoadClient g = do
|
||||
hSetBuffering stdout LineBuffering
|
||||
workers <- fromMaybe 16 . (>>= readMaybe) <$> lookupEnv "LOAD_WORKERS"
|
||||
secs <- fromMaybe 60 . (>>= readMaybe) <$> lookupEnv "LOAD_SEC"
|
||||
ops <- newTVarIO (0 :: Int)
|
||||
start <- getCurrentTime
|
||||
let deadline = fromIntegral (secs :: Int)
|
||||
worker w = do
|
||||
t <- getCurrentTime
|
||||
when (diffUTCTime t start < deadline) $ do
|
||||
benchClient $ \recip -> benchClient $ \sndr ->
|
||||
forM_ ([1 .. 50] :: [Int]) $ \i -> do
|
||||
(if i `mod` 3 == 0 then ntfStep else plainStep) g recip sndr (w * 1000000 + i)
|
||||
atomically $ modifyTVar' ops (+ 1)
|
||||
worker w
|
||||
forConcurrently_ ([1 .. workers] :: [Int]) worker
|
||||
end <- getCurrentTime
|
||||
n <- readTVarIO ops
|
||||
putStrLn $ "DONE " <> show n <> " " <> show (realToFrac (diffUTCTime end start) :: Double)
|
||||
|
||||
-- Client side of subslice: creates the subscriptions, prints "READY <subscribed>", and holds the
|
||||
-- connection until stdin is closed.
|
||||
runSubClient :: TVar ChaChaDRG -> Int -> IO ()
|
||||
runSubClient g iters = do
|
||||
hSetBuffering stdout LineBuffering
|
||||
lookupEnv "SUBMODE" >>= \case
|
||||
Just "new" -> do
|
||||
subscribe <- (/= Just "0") <$> lookupEnv "NEWSUB"
|
||||
benchClient $ \h -> do
|
||||
forM_ ([1 .. iters] :: [Int]) $ \i -> do
|
||||
(rPub, rKey, dhPub) <- genKeys g
|
||||
Resp _ _ Ids {} <- signSendRecv h rKey (B.pack $ "n" <> show i, NoEntity, (if subscribe then New else New0) rPub dhPub)
|
||||
pure ()
|
||||
ready $ if subscribe then iters else 0
|
||||
-- create without subscribing, then SUB on the same connection
|
||||
Just "newthensub" -> benchClient $ \h -> do
|
||||
forM_ ([1 .. iters] :: [Int]) $ \i -> do
|
||||
(rPub, rKey, dhPub) <- genKeys g
|
||||
Resp _ _ (Ids rId _ _) <- signSendRecv h rKey (B.pack $ "n" <> show i, NoEntity, New0 rPub dhPub)
|
||||
Resp _ _ (SOK Nothing) <- signSendRecv h rKey (B.pack $ "s" <> show i, rId, SUB)
|
||||
pure ()
|
||||
ready iters
|
||||
_ -> do
|
||||
batch <- fromMaybe 1 . (>>= readMaybe) <$> lookupEnv "SUBBATCH"
|
||||
keep <- fromMaybe batch . (>>= readMaybe) <$> lookupEnv "SUBKEEP"
|
||||
qs <- benchClient $ \h -> forM ([1 .. iters] :: [Int]) $ \i -> do
|
||||
(rPub, rKey, dhPub) <- genKeys g
|
||||
Resp _ _ (Ids rId _ _) <- signSendRecv h rKey (B.pack $ "c" <> show i, NoEntity, New0 rPub dhPub)
|
||||
pure (i, rId, rKey)
|
||||
benchClient $ \h -> do
|
||||
forM_ (chunks batch qs) $ \b -> do
|
||||
let sub (i, rId, rKey) = Right $ signTransmission h rKey Nothing (B.pack $ "s" <> show i, rId, SUB)
|
||||
void $ tPut h (fromList $ map sub b)
|
||||
awaitResponses h (length b)
|
||||
forM_ (drop keep b) $ \(i, rId, rKey) -> do
|
||||
Resp _ _ OK <- signSendRecv h rKey (B.pack $ "d" <> show i, rId, DEL)
|
||||
pure ()
|
||||
ready $ sum $ map (min keep . length) (chunks batch qs)
|
||||
where
|
||||
ready :: Int -> IO ()
|
||||
ready k = do
|
||||
putStrLn $ "READY " <> show k
|
||||
void (E.try getLine :: IO (Either E.IOException String))
|
||||
chunks n xs = case splitAt n xs of
|
||||
(c, []) -> [c]
|
||||
(c, rest) -> c : chunks n 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)
|
||||
#endif
|
||||
|
||||
-- forkDeliver blocked in SubPending: a subscriber that never reads its sndQ.
|
||||
-- Deliveries to a full sndQ fork a deliverThread that blocks forever -> threads/subs_thread grow.
|
||||
runStuck :: TVar ChaChaDRG -> Int -> Int -> IO ()
|
||||
@@ -814,20 +1188,24 @@ heldConns phase iters
|
||||
-- connection from `active` before gracefulClose (up to 5s) and before it increments `closed`,
|
||||
-- so a connection in teardown is counted in neither and shows up as leaked.
|
||||
cpSockets :: N.ServiceName -> IO [String]
|
||||
cpSockets port = do
|
||||
cpSockets port = fromMaybe [] <$> cpCommand port "sockets" 5 5000000 -- "Sockets for port N:" + accepted/closed/active/leaked
|
||||
|
||||
-- Run one admin command on the control port and read `n` reply lines, or Nothing on timeout.
|
||||
cpCommand :: N.ServiceName -> String -> Int -> Int -> IO (Maybe [String])
|
||||
cpCommand port cmd n tmo = do
|
||||
sock <- rawConnect port
|
||||
h <- N.socketToHandle sock ReadWriteMode
|
||||
hSetBuffering h LineBuffering
|
||||
hSetNewlineMode h universalNewlineMode
|
||||
r <- timeout 5000000 $ do
|
||||
r <- timeout tmo $ do
|
||||
_ <- hGetLine h -- banner line 1
|
||||
_ <- hGetLine h -- banner line 2
|
||||
hPutStrLn h "auth bench"
|
||||
_ <- hGetLine h
|
||||
hPutStrLn h "sockets"
|
||||
replicateM 5 (hGetLine h) -- "Sockets for port N:" + accepted/closed/active/leaked
|
||||
hPutStrLn h cmd
|
||||
replicateM n (hGetLine h)
|
||||
hClose h `E.catch` \(_ :: E.SomeException) -> pure ()
|
||||
pure $ fromMaybe [] r
|
||||
pure r
|
||||
|
||||
cpPort :: N.ServiceName
|
||||
cpPort = "5010"
|
||||
@@ -1012,6 +1390,19 @@ main = do
|
||||
"fastfwd" -> runFastFwd g iters cp
|
||||
"msgqfill" -> runMsgQFill g iters cp
|
||||
_ -> error $ "unknown proxy phase: " <> phase
|
||||
#if defined(dbServerPostgres)
|
||||
else if phase == "ntfloop" then runNtfLoop iters
|
||||
else if phase == "ntfdeliver" then postgressBracket benchDBConnectInfo $ withSmpServerConfigOn (transport @TLS) (updateCfg benchPgCfg $ \c -> c {ntfDeliveryInterval = 1500000}) benchPort $ \_ -> runNtfDeliver g iters
|
||||
else if phase == "cpsave" then
|
||||
postgressBracket benchDBConnectInfo $
|
||||
withSmpServerConfigOn (transport @TLS) (updateCfg benchPgCfg $ \c -> c {controlPort = Just benchCpPort, controlPortAdminAuth = Just "bench"}) benchPort $ \_ -> runCpSave g
|
||||
else if phase == "subslice" then postgressBracket benchDBConnectInfo $ withSmpServerConfigOn (transport @TLS) benchPgCfg benchPort $ \_ -> runSubSlice iters
|
||||
else if phase == "subclient" then runSubClient g iters
|
||||
else if phase == "conns" then postgressBracket benchDBConnectInfo $ withSmpServerConfigOn (transport @TLS) benchPgCfg benchPort $ \_ -> runConns iters
|
||||
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
|
||||
#endif
|
||||
else withSmpServerConfigOn (transport @TLS) srvCfg testPort $ \_ -> settle leakDiagSec $ do
|
||||
threadDelay 250000
|
||||
case phase of
|
||||
|
||||
Reference in New Issue
Block a user