From 706211dc24f4e1920d8a10a5c3548ec696ddc1d2 Mon Sep 17 00:00:00 2001 From: shum Date: Wed, 30 Sep 2026 12:03:25 +0000 Subject: [PATCH] bench: add server-only memory and load phases --- bench/MemBench.hs | 411 ++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 401 insertions(+), 10 deletions(-) diff --git a/bench/MemBench.hs b/bench/MemBench.hs index 1f2daf548..9a14002a5 100644 --- a/bench/MemBench.hs +++ b/bench/MemBench.hs @@ -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 ", 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