From 6db79808aa31b1c719c9d3c3d73e6181711a08f7 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Tue, 18 Mar 2025 09:40:22 +0000 Subject: [PATCH] smp server: use COPY to import store log to postgres db, improve concurrency and error handling (#1487) * smp server: use COPY to import store log to postgres db * compact queues when importing to postgres * mempty * version * handle errors while expiring, mask async exceptions while getting queue * whitespace * version --- simplexmq.cabal | 2 +- src/Simplex/Messaging/Server.hs | 12 +-- src/Simplex/Messaging/Server/Main.hs | 14 ++-- .../Messaging/Server/QueueStore/Postgres.hs | 75 +++++++++++++------ .../Messaging/Server/QueueStore/STM.hs | 1 - src/Simplex/Messaging/Server/StoreLog.hs | 5 +- 6 files changed, 74 insertions(+), 35 deletions(-) diff --git a/simplexmq.cabal b/simplexmq.cabal index 5cd28a1f1..6bbf363bb 100644 --- a/simplexmq.cabal +++ b/simplexmq.cabal @@ -1,7 +1,7 @@ cabal-version: 1.12 name: simplexmq -version: 6.3.0.805 +version: 6.3.0.8 synopsis: SimpleXMQ message broker description: This package includes <./docs/Simplex-Messaging-Server.html server>, <./docs/Simplex-Messaging-Client.html client> and diff --git a/src/Simplex/Messaging/Server.hs b/src/Simplex/Messaging/Server.hs index d816118d5..00dc7cba7 100644 --- a/src/Simplex/Messaging/Server.hs +++ b/src/Simplex/Messaging/Server.hs @@ -399,15 +399,17 @@ smpServer started cfg@ServerConfig {transports, transportConfig = tCfg, startOpt expire :: forall s. MsgStoreClass s => s -> ServerStats -> Int64 -> IO () expire ms stats interval = do threadDelay' interval + logInfo "Started expiring messages..." n <- compactQueues @(StoreQueue s) $ queueStore ms when (n > 0) $ logInfo $ "Removed " <> tshow n <> " old deleted queues from the database." old <- expireBeforeEpoch expCfg now <- systemSeconds <$> getSystemTime - msgStats@MessageStats {storedMsgsCount = stored, expiredMsgsCount = expired} <- - withAllMsgQueues False "idleDeleteExpiredMsgs" ms $ expireQueueMsgs now ms old - atomicWriteIORef (msgCount stats) stored - atomicModifyIORef'_ (msgExpired stats) (+ expired) - printMessageStats "STORE: messages" msgStats + tryAny (withAllMsgQueues False "idleDeleteExpiredMsgs" ms $ expireQueueMsgs now ms old) >>= \case + Right msgStats@MessageStats {storedMsgsCount = stored, expiredMsgsCount = expired} -> do + atomicWriteIORef (msgCount stats) stored + atomicModifyIORef'_ (msgExpired stats) (+ expired) + printMessageStats "STORE: messages" msgStats + Left e -> logError $ "STORE: withAllMsgQueues, error expiring messages, " <> tshow e expireQueueMsgs now ms old q = do (expired_, stored) <- idleDeleteExpiredMsgs now ms q old pure MessageStats {storedMsgsCount = stored, expiredMsgsCount = fromMaybe 0 expired_, storedQueues = 1} diff --git a/src/Simplex/Messaging/Server/Main.hs b/src/Simplex/Messaging/Server/Main.hs index 14ea6fe9e..6c930ba0b 100644 --- a/src/Simplex/Messaging/Server/Main.hs +++ b/src/Simplex/Messaging/Server/Main.hs @@ -71,7 +71,7 @@ import Simplex.Messaging.Server.MsgStore.Types (QSType (..)) import Simplex.Messaging.Server.MsgStore.Journal (postgresQueueStore) import Simplex.Messaging.Server.QueueStore.Postgres (batchInsertQueues, foldQueueRecs) import Simplex.Messaging.Server.QueueStore.Types -import Simplex.Messaging.Server.StoreLog (logCreateQueue, openWriteStoreLog) +import Simplex.Messaging.Server.StoreLog (closeStoreLog, logCreateQueue, openWriteStoreLog) import System.Directory (renameFile) #endif @@ -176,10 +176,11 @@ smpServerCLI_ generateSite serveStaticFiles attachStaticFiles cfgPath logPath = | otherwise -> do storeLogFile <- getRequiredStoreLogFile ini confirmOrExit - ("WARNING: store log file " <> storeLogFile <> " will be imported to PostrgreSQL database: " <> B.unpack connstr <> ", schema: " <> B.unpack schema) + ("WARNING: store log file " <> storeLogFile <> " will be compacted and imported to PostrgreSQL database: " <> B.unpack connstr <> ", schema: " <> B.unpack schema) "Queue records not imported" ms <- newJournalMsgStore MQStoreCfg - readQueueStore True (mkQueue ms) storeLogFile (queueStore ms) + sl <- readWriteQueueStore True (mkQueue ms) storeLogFile (queueStore ms) + closeStoreLog sl queues <- readTVarIO $ loadedQueues $ stmQueueStore ms let storeCfg = PostgresStoreCfg {dbOpts = dbOpts {createSchema = True}, dbStoreLogPath = Nothing, confirmMigrations = MCConsole, deletedTTL = iniDeletedTTL ini} ps <- newJournalMsgStore $ PQStoreCfg storeCfg @@ -187,10 +188,13 @@ smpServerCLI_ generateSite serveStaticFiles attachStaticFiles cfgPath logPath = renameFile storeLogFile $ storeLogFile <> ".bak" putStrLn $ "Import completed: " <> show qCnt <> " queues" putStrLn $ case readStoreType ini of - Right (ASType SQSMemory SMSMemory) -> "store_messages set to `memory`.\nImport messages to journal to use PostgreSQL database for queues (`smp-server journal import`)" - Right (ASType SQSMemory SMSJournal) -> "store_queues set to `memory`, update it to `database` in INI file" + Right (ASType SQSMemory SMSMemory) -> setToDbStr <> "\nstore_messages set to `memory`, import messages to journal to use PostgreSQL database for queues (`smp-server journal import`)" + Right (ASType SQSMemory SMSJournal) -> setToDbStr Right (ASType SQSPostgres SMSJournal) -> "store_queues set to `database`, start the server." Left e -> e <> ", configure storage correctly" + where + setToDbStr :: String + setToDbStr = "store_queues set to `memory`, update it to `database` in INI file" SCExport | schemaExists && storeLogExists -> exitConfigureQueueStore connstr schema | not schemaExists -> do diff --git a/src/Simplex/Messaging/Server/QueueStore/Postgres.hs b/src/Simplex/Messaging/Server/QueueStore/Postgres.hs index 3c9bbbcd8..3062e2313 100644 --- a/src/Simplex/Messaging/Server/QueueStore/Postgres.hs +++ b/src/Simplex/Messaging/Server/QueueStore/Postgres.hs @@ -32,18 +32,25 @@ import Control.Monad import Control.Monad.Except import Control.Monad.IO.Class import Control.Monad.Trans.Except +import Data.ByteString.Builder (Builder) +import qualified Data.ByteString.Builder as BB +import Data.ByteString.Char8 (ByteString) +import qualified Data.ByteString.Char8 as B +import qualified Data.ByteString.Lazy as LB import Data.Bitraversable (bimapM) import Data.Either (fromRight) import Data.Functor (($>)) import Data.Int (Int64) +import Data.List (intersperse) import qualified Data.Map.Strict as M import Data.Maybe (catMaybes) import qualified Data.Text as T import Data.Time.Clock.System (SystemTime (..), getSystemTime) import Database.PostgreSQL.Simple (Binary (..), Only (..), Query, SqlError) import qualified Database.PostgreSQL.Simple as DB +import qualified Database.PostgreSQL.Simple.Copy as DB import Database.PostgreSQL.Simple.FromField (FromField (..)) -import Database.PostgreSQL.Simple.ToField (ToField (..)) +import Database.PostgreSQL.Simple.ToField (Action (..), ToField (..)) import Database.PostgreSQL.Simple.Errors (ConstraintViolation (..), constraintViolation) import Database.PostgreSQL.Simple.SqlQQ (sql) import GHC.IO (catchAny) @@ -160,7 +167,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where loadSndQueue = loadQueue " WHERE sender_id = ?" $ \rId -> TM.insert qId rId senders loadNtfQueue = loadQueue " WHERE notifier_id = ?" $ \_ -> pure () -- do NOT cache ref - ntf subscriptions are rare loadQueue condition insertRef = - runExceptT $ do + E.uninterruptibleMask_ $ runExceptT $ do (rId, qRec) <- withDB "getQueue_" st $ \db -> firstRow rowToQueueRec AUTH $ DB.query db (queueRecQuery <> condition <> " AND deleted_at IS NULL") (Only qId) @@ -182,7 +189,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where secureQueue :: PostgresQueueStore q -> q -> SndPublicAuthKey -> IO (Either ErrorType ()) secureQueue st sq sKey = - withQueueDB sq "secureQueue" $ \q -> do + withQueueRec sq "secureQueue" $ \q -> do verify q assertUpdated $ withDB' "secureQueue" st $ \db -> DB.execute db "UPDATE msg_queues SET sender_key = ? WHERE recipient_id = ? AND deleted_at IS NULL" (sKey, rId) @@ -196,7 +203,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where addQueueNotifier :: PostgresQueueStore q -> q -> NtfCreds -> IO (Either ErrorType (Maybe NotifierId)) addQueueNotifier st sq ntfCreds@NtfCreds {notifierId = nId, notifierKey, rcvNtfDhSecret} = - withQueueDB sq "addQueueNotifier" $ \q -> + withQueueRec sq "addQueueNotifier" $ \q -> ExceptT $ withLockMap (notifierLocks st) nId "addQueueNotifier" $ ifM (TM.memberIO nId notifiers) (pure $ Left DUPLICATE_) $ runExceptT $ do assertUpdated $ withDB "addQueueNotifier" st $ \db -> @@ -223,7 +230,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where deleteQueueNotifier :: PostgresQueueStore q -> q -> IO (Either ErrorType (Maybe NotifierId)) deleteQueueNotifier st sq = - withQueueDB sq "deleteQueueNotifier" $ \q -> + withQueueRec sq "deleteQueueNotifier" $ \q -> ExceptT $ fmap sequence $ forM (notifier q) $ \NtfCreds {notifierId = nId} -> withLockMap (notifierLocks st) nId "deleteQueueNotifier" $ runExceptT $ do assertUpdated $ withDB' "deleteQueueNotifier" st update @@ -260,7 +267,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where updateQueueTime :: PostgresQueueStore q -> q -> RoundedSystemTime -> IO (Either ErrorType QueueRec) updateQueueTime st sq t = - withQueueDB sq "updateQueueTime" $ \q@QueueRec {updatedAt} -> + withQueueRec sq "updateQueueTime" $ \q@QueueRec {updatedAt} -> if updatedAt == Just t then pure q else do @@ -295,21 +302,19 @@ batchInsertQueues tty queues toStore = do qs <- catMaybes <$> mapM (\(rId, q) -> (rId,) <$$> readTVarIO (queueRec q)) (M.assocs queues) putStrLn $ "Importing " <> show (length qs) <> " queues..." let st = dbStore toStore - (qCnt, count) <- foldM (processChunk st) (0, 0) $ toChunks 1000000 qs + count <- + withConnection st $ \db -> do + DB.copy_ db "COPY msg_queues (recipient_id, recipient_key, rcv_dh_secret, sender_id, sender_key, snd_secure, notifier_id, notifier_key, rcv_ntf_dh_secret, status, updated_at) FROM STDIN WITH (FORMAT CSV)" + mapM_ (putQueue db) (zip [1..] qs) + DB.putCopyEnd db + Only qCnt : _ <- withConnection st (`DB.query_` "SELECT count(*) FROM msg_queues") putStrLn $ progress count pure qCnt where - processChunk st (qCnt, i) qs = do - qCnt' <- withConnection st $ \db -> DB.executeMany db insertQueueQuery $ map queueRecToRow qs - let i' = i + length qs - when tty $ putStr (progress i' <> "\r") >> hFlush stdout - pure (qCnt + qCnt', i') + putQueue db (i :: Int, q) = do + DB.putCopyData db $ queueRecToText q + when (tty && i `mod` 100000 == 0) $ putStr (progress i <> "\r") >> hFlush stdout progress i = "Imported: " <> show i <> " queues" - toChunks :: Int -> [a] -> [[a]] - toChunks _ [] = [] - toChunks n xs = - let (ys, xs') = splitAt n xs - in ys : toChunks n xs' insertQueueQuery :: Query insertQueueQuery = @@ -349,6 +354,34 @@ queueRecToRow :: (RecipientId, QueueRec) -> QueueRecRow queueRecToRow (rId, QueueRec {recipientKey, rcvDhSecret, senderId, senderKey, sndSecure, notifier = n, status, updatedAt}) = (rId, recipientKey, rcvDhSecret, senderId, senderKey, sndSecure, notifierId <$> n, notifierKey <$> n, rcvNtfDhSecret <$> n, status, updatedAt) +queueRecToText :: (RecipientId, QueueRec) -> ByteString +queueRecToText (rId, QueueRec {recipientKey, rcvDhSecret, senderId, senderKey, sndSecure, notifier = n, status, updatedAt}) = + LB.toStrict $ BB.toLazyByteString $ mconcat tabFields <> BB.char7 '\n' + where + tabFields = BB.char7 ',' `intersperse` fields + fields = + [ renderField (toField rId), + renderField (toField recipientKey), + renderField (toField rcvDhSecret), + renderField (toField senderId), + nullable senderKey, + renderField (toField sndSecure), + nullable (notifierId <$> n), + nullable (notifierKey <$> n), + nullable (rcvNtfDhSecret <$> n), + BB.char7 '"' <> renderField (toField status) <> BB.char7 '"', + nullable updatedAt + ] + nullable :: ToField a => Maybe a -> Builder + nullable = maybe mempty (renderField . toField) + renderField :: Action -> Builder + renderField = \case + Plain bld -> bld + Escape s -> BB.byteString s + EscapeByteA s -> BB.string7 "\\x" <> BB.byteStringHex s + EscapeIdentifier s -> BB.byteString s -- Not used in COPY data + Many as -> mconcat (map renderField as) + rowToQueueRec :: QueueRecRow -> (RecipientId, QueueRec) rowToQueueRec (rId, recipientKey, rcvDhSecret, senderId, senderKey, sndSecure, notifierId_, notifierKey_, rcvNtfDhSecret_, status, updatedAt) = let notifier = NtfCreds <$> notifierId_ <*> notifierKey_ <*> rcvNtfDhSecret_ @@ -356,14 +389,14 @@ rowToQueueRec (rId, recipientKey, rcvDhSecret, senderId, senderKey, sndSecure, n setStatusDB :: StoreQueueClass q => String -> PostgresQueueStore q -> q -> ServerEntityStatus -> ExceptT ErrorType IO () -> IO (Either ErrorType ()) setStatusDB op st sq status writeLog = - withQueueDB sq op $ \q -> do + withQueueRec sq op $ \q -> do assertUpdated $ withDB' op st $ \db -> DB.execute db "UPDATE msg_queues SET status = ? WHERE recipient_id = ? AND deleted_at IS NULL" (status, recipientId sq) atomically $ writeTVar (queueRec sq) $ Just q {status} writeLog -withQueueDB :: StoreQueueClass q => q -> String -> (QueueRec -> ExceptT ErrorType IO a) -> IO (Either ErrorType a) -withQueueDB sq op action = +withQueueRec :: StoreQueueClass q => q -> String -> (QueueRec -> ExceptT ErrorType IO a) -> IO (Either ErrorType a) +withQueueRec sq op action = withQueueLock sq op $ E.uninterruptibleMask_ $ runExceptT $ ExceptT (readQueueRecIO $ queueRec sq) >>= action assertUpdated :: ExceptT ErrorType IO Int64 -> ExceptT ErrorType IO () @@ -379,7 +412,7 @@ withDB op st action = logErr :: E.SomeException -> IO (Either ErrorType a) logErr e = logError ("STORE: " <> T.pack err) $> Left (STORE err) where - err = op <> ", withLog, " <> show e + err = op <> ", withDB, " <> show e withLog :: MonadIO m => String -> PostgresQueueStore q -> (StoreLog 'WriteMode -> IO ()) -> m () withLog op PostgresQueueStore {dbStoreLog} action = diff --git a/src/Simplex/Messaging/Server/QueueStore/STM.hs b/src/Simplex/Messaging/Server/QueueStore/STM.hs index c119beb23..8a360c3a0 100644 --- a/src/Simplex/Messaging/Server/QueueStore/STM.hs +++ b/src/Simplex/Messaging/Server/QueueStore/STM.hs @@ -18,7 +18,6 @@ module Simplex.Messaging.Server.QueueStore.STM ( STMQueueStore (..), setStoreLog, withLog', - withQueueRec, readQueueRecIO, setStatus, ) diff --git a/src/Simplex/Messaging/Server/StoreLog.hs b/src/Simplex/Messaging/Server/StoreLog.hs index 1e0565a61..1cc8ebd6c 100644 --- a/src/Simplex/Messaging/Server/StoreLog.hs +++ b/src/Simplex/Messaging/Server/StoreLog.hs @@ -259,8 +259,9 @@ removeStoreLogBackups f = do times2 = take (length times1 - minOldBackups) times1 -- keep 3 backups older than 24 hours toDelete = filter (< old) times2 -- remove all backups older than 21 day mapM_ (removeFile . backupPath) toDelete - putStrLn $ "Removed " <> show (length toDelete) <> " backups:" - mapM_ (putStrLn . backupPath) toDelete + when (length toDelete > 0) $ do + putStrLn $ "Removed " <> show (length toDelete) <> " backups:" + mapM_ (putStrLn . backupPath) toDelete where backupPathTime :: FilePath -> Maybe UTCTime backupPathTime = iso8601ParseM <=< stripPrefix backupPathPfx