From b260bbb80d9218505dde7e2e32c9a28a4e73b978 Mon Sep 17 00:00:00 2001 From: Evgeny Poberezkin Date: Wed, 23 Oct 2024 15:03:40 +0100 Subject: [PATCH] process all queues more efficiently --- src/Simplex/Messaging/Server.hs | 7 ++-- .../Messaging/Server/MsgStore/Journal.hs | 32 +++++++++---------- 2 files changed, 20 insertions(+), 19 deletions(-) diff --git a/src/Simplex/Messaging/Server.hs b/src/Simplex/Messaging/Server.hs index bed242e7a..084e75303 100644 --- a/src/Simplex/Messaging/Server.hs +++ b/src/Simplex/Messaging/Server.hs @@ -378,9 +378,10 @@ smpServer started cfg@ServerConfig {transports, transportConfig = tCfg} attachHT liftIO $ forever $ do threadDelay' interval old <- expireBeforeEpoch expCfg - void $ withActiveMsgQueues ms (\_ -> expireQueueMsgs stats old) 0 + deleted <- withActiveMsgQueues ms (\_ -> expireQueueMsgs stats old) 0 + logInfo $ "STORE: expireMessagesThread, expired " <> tshow deleted <> " messages" where - expireQueueMsgs stats old q acc = + expireQueueMsgs stats old q !acc = runExceptT (deleteExpiredMsgs q True old) >>= \case Right deleted -> (acc + deleted) <$ atomicModifyIORef'_ (msgExpired stats) (+ deleted) Left _ -> pure 0 @@ -1741,7 +1742,7 @@ exportMessages ms f drainMsgs = do total <- liftIO $ withFile f WriteMode $ \h -> withAllMsgQueues ms (saveQueueMsgs h) 0 logInfo $ "messages saved: " <> tshow total where - saveQueueMsgs h rId q acc = getQueueMessages drainMsgs q >>= \msgs -> (acc + length msgs) <$ BLD.hPutBuilder h (encodeMessages rId msgs) + saveQueueMsgs h rId q !acc = getQueueMessages drainMsgs q >>= \msgs -> (acc + length msgs) <$ BLD.hPutBuilder h (encodeMessages rId msgs) encodeMessages rId = mconcat . map (\msg -> BLD.byteString (strEncode $ MLRv3 rId msg) <> BLD.char8 '\n') processServerMessages :: M MessageStats diff --git a/src/Simplex/Messaging/Server/MsgStore/Journal.hs b/src/Simplex/Messaging/Server/MsgStore/Journal.hs index b6549abe3..cd0fe3923 100644 --- a/src/Simplex/Messaging/Server/MsgStore/Journal.hs +++ b/src/Simplex/Messaging/Server/MsgStore/Journal.hs @@ -220,31 +220,28 @@ instance MsgStoreClass JournalMsgStore where processStore = do closeMsgStore ms lock <- createLockIO -- the same lock is used for all queues - dirs <- zip [0..] <$> listQueueDirs 0 ("", storePath) - let count = length dirs - res' <- foldM (processQueue lock count) res dirs - progress count count + (!count, !res') <- foldQueues 0 (processQueue lock) (0, res) ("", storePath) + progress count putStrLn "" pure res' JournalStoreConfig {storePath, pathParts} = config - processQueue queueLock count acc (i :: Int, (queueId, dir)) = do - when (i `mod` 100 == 0) $ progress i count - let statePath = dir queueLogFileName <> "." <> queueId <> logFileExt - queue = JMQueue {queueDirectory = dir, queueLock, statePath} - q <- openMsgQueue ms queue + processQueue queueLock (!i :: Int, !acc) (queueId, dir) = do + when (i `mod` 100 == 0) $ progress i + let statePath = msgQueueStatePath dir queueId + q <- openMsgQueue ms JMQueue {queueDirectory = dir, queueLock, statePath} acc' <- case strDecode $ B.pack queueId of Right rId -> action rId q acc Left e -> do putStrLn ("Error: message queue directory " <> dir <> " is invalid: " <> e) exitFailure closeMsgQueue q - pure acc' - progress i count = do - putStr $ "Processed: " <> show i <> "/" <> show count <> " queues\r" + pure (i + 1, acc') + progress i = do + putStr $ "Processed: " <> show i <> " queues\r" IO.hFlush stdout - listQueueDirs depth (queueId, path) - | depth == pathParts - 1 = listDirs - | otherwise = fmap concat . mapM (listQueueDirs (depth + 1)) =<< listDirs + foldQueues depth f acc (queueId, path) = do + let f' = if depth == pathParts - 1 then f else foldQueues (depth + 1) f + listDirs >>= foldM f' acc where listDirs = fmap catMaybes . mapM queuePath =<< listDirectory path queuePath dir = do @@ -271,7 +268,7 @@ instance MsgStoreClass JournalMsgStore where newQ = do queueLock <- atomically $ getMapLock queueLocks rId let dir = msgQueueDirectory ms rId - statePath = dir (queueLogFileName <> "." <> B.unpack (strEncode rId) <> logFileExt) + statePath = msgQueueStatePath dir $ B.unpack (strEncode rId) queue = JMQueue {queueDirectory = dir, queueLock, statePath} q <- ifM (doesDirectoryExist dir) (openMsgQueue ms queue) (createQ queue) atomically $ TM.insert rId q msgQueues @@ -448,6 +445,9 @@ msgQueueDirectory JournalMsgStore {config = JournalStoreConfig {storePath, pathP let (seg, s') = B.splitAt 2 s in seg : splitSegments (n - 1) s' +msgQueueStatePath :: FilePath -> String -> FilePath +msgQueueStatePath dir queueId = dir (queueLogFileName <> "." <> queueId <> logFileExt) + createNewJournal :: FilePath -> ByteString -> IO Handle createNewJournal dir journalId = do let path = journalFilePath dir journalId -- TODO retry if file exists