diff --git a/bench/MemBench.hs b/bench/MemBench.hs index 7be41611c..69e2ca2ce 100644 --- a/bench/MemBench.hs +++ b/bench/MemBench.hs @@ -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 ", 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