From 83eb9a55fdc08b3d356bcbf634198f8a34243d44 Mon Sep 17 00:00:00 2001 From: shum Date: Fri, 2 Oct 2026 11:03:32 +0000 Subject: [PATCH] smp-server: remove deleted queues from store maps --- .../Messaging/Server/QueueStore/Postgres.hs | 15 ++-- .../Messaging/Server/QueueStore/STM.hs | 4 +- .../Messaging/Server/StoreLog/ReadWrite.hs | 13 +--- tests/CoreTests/MsgStoreTests.hs | 78 ++++++++++++++++++- tests/CoreTests/StoreLogTests.hs | 12 --- 5 files changed, 90 insertions(+), 32 deletions(-) diff --git a/src/Simplex/Messaging/Server/QueueStore/Postgres.hs b/src/Simplex/Messaging/Server/QueueStore/Postgres.hs index 39f02131c..80a49a64f 100644 --- a/src/Simplex/Messaging/Server/QueueStore/Postgres.hs +++ b/src/Simplex/Messaging/Server/QueueStore/Postgres.hs @@ -212,21 +212,21 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where loadRcvQueue = do (rId, qRec) <- loadQueue " WHERE recipient_id = ?" liftIO $ cacheQueue rId qRec $ \_ -> pure () -- recipient map already checked, not caching sender ref - loadSndQueue = loadSndQueue_ " WHERE sender_id = ?" - loadLinkQueue = loadSndQueue_ " WHERE link_id = ?" + loadSndQueue = loadSndQueue_ senders " WHERE sender_id = ?" + loadLinkQueue = loadSndQueue_ links " WHERE link_id = ?" loadNtfQueue = do (rId, qRec) <- loadQueue " WHERE notifier_id = ?" liftIO $ TM.lookupIO rId queues -- checking recipient map first, not creating lock in map, not caching queue >>= maybe (mkQ False rId qRec) pure - loadSndQueue_ condition = do + loadSndQueue_ refs condition = do (rId, qRec) <- loadQueue condition liftIO $ TM.lookupIO rId queues -- checking recipient map first - >>= maybe (cacheQueue rId qRec cacheSender) (atomically (cacheSender rId) $>) + >>= maybe (cacheQueue rId qRec $ cacheRef refs) (atomically (cacheRef refs rId) $>) loadQueueNoCache cond = mask $ loadQueue cond >>= liftIO . uncurry (mkQ True) mask = E.uninterruptibleMask_ . runExceptT - cacheSender rId = TM.insert qId rId senders + cacheRef refs rId = TM.insert qId rId refs loadQueue condition = withDB "getQueue_" st $ \db -> firstRow rowToQueueRec AUTH $ DB.query db (queueRecQuery <> condition <> " AND deleted_at IS NULL") (Only qId) @@ -329,9 +329,10 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where deleteQueueLinkData :: PostgresQueueStore q -> q -> IO (Either ErrorType ()) deleteQueueLinkData st sq = withQueueRec sq "deleteQueueLinkData" $ \q -> case queueData q of - Just _ -> do + Just (lnkId, _) -> do assertUpdated $ withDB' "deleteQueueLinkData" st $ \db -> DB.execute db "UPDATE msg_queues SET link_id = NULL, fixed_data = NULL, user_data = NULL WHERE recipient_id = ? AND deleted_at IS NULL" (Only rId) + when (useCache st) $ atomically $ TM.delete lnkId $ links st atomically $ writeTVar (queueRec sq) $ Just q {queueData = Nothing} withLog "deleteQueueLinkData" st (`logDeleteLink` rId) _ -> throwE AUTH @@ -457,7 +458,9 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where DB.execute db "UPDATE msg_queues SET deleted_at = ? WHERE recipient_id = ? AND deleted_at IS NULL" (ts, rId) atomically $ writeTVar qr Nothing when (useCache st) $ do + atomically $ TM.delete rId $ queues st atomically $ TM.delete (senderId q) $ senders st + forM_ (queueData q) $ \(lnkId, _) -> atomically $ TM.delete lnkId $ links st forM_ (notifier q) $ \NtfCreds {notifierId} -> do atomically $ TM.delete notifierId $ notifiers st atomically $ TM.delete notifierId $ notifierLocks st diff --git a/src/Simplex/Messaging/Server/QueueStore/STM.hs b/src/Simplex/Messaging/Server/QueueStore/STM.hs index 1e17b9051..4c32505e6 100644 --- a/src/Simplex/Messaging/Server/QueueStore/STM.hs +++ b/src/Simplex/Messaging/Server/QueueStore/STM.hs @@ -277,9 +277,11 @@ instance StoreQueueClass q => QueueStoreClass q (STMQueueStore q) where where rId = recipientId sq qr = queueRec sq - delete q@QueueRec {senderId, rcvServiceId} = do + delete q@QueueRec {senderId, queueData, rcvServiceId} = do writeTVar qr Nothing + TM.delete rId $ queues st TM.delete senderId $ senders st + forM_ queueData $ \(lnkId, _) -> TM.delete lnkId $ links st mapM_ (removeServiceQueue st serviceRcvQueues rId) rcvServiceId mapM_ (removeNotifier st) $ notifier q pure q diff --git a/src/Simplex/Messaging/Server/StoreLog/ReadWrite.hs b/src/Simplex/Messaging/Server/StoreLog/ReadWrite.hs index 6bf327f70..f5536d66f 100644 --- a/src/Simplex/Messaging/Server/StoreLog/ReadWrite.hs +++ b/src/Simplex/Messaging/Server/StoreLog/ReadWrite.hs @@ -14,9 +14,6 @@ module Simplex.Messaging.Server.StoreLog.ReadWrite import Control.Concurrent.STM import Control.Logger.Simple -import Control.Monad -import Control.Monad.IO.Class -import Control.Monad.Trans.Except import qualified Data.ByteString.Char8 as B import qualified Data.Text as T import Data.Text.Encoding (decodeLatin1) @@ -26,7 +23,7 @@ import Simplex.Messaging.Server.QueueStore (QueueRec, ServiceRec (..)) import Simplex.Messaging.Server.QueueStore.STM (STMQueueStore (..), STMService (..)) import Simplex.Messaging.Server.QueueStore.Types import Simplex.Messaging.Server.StoreLog -import Simplex.Messaging.Util (tshow) +import Simplex.Messaging.Util (tshow, ($>>=)) import System.IO writeQueueStore :: forall q. StoreQueueClass q => StoreLog 'WriteMode -> STMQueueStore q -> IO () @@ -68,13 +65,7 @@ readQueueStore tty mkQ f st = readLogLines tty f $ \_ -> processLine printError :: String -> IO () printError e = B.putStrLn $ "Error parsing log: " <> B.pack e <> " - " <> s withQueue :: forall a. RecipientId -> T.Text -> (q -> IO (Either ErrorType a)) -> IO () - withQueue qId op a = runExceptT go >>= qError qId op - where - go = do - q <- ExceptT $ getQueue_ st (\_ -> mkQ) SRecipient qId - liftIO (readTVarIO $ queueRec q) >>= \case - Nothing -> logWarn $ logPfx qId op <> "already deleted" - Just _ -> void $ ExceptT $ a q + withQueue qId op a = (getQueue_ st (\_ -> mkQ) SRecipient qId $>>= a) >>= qError qId op qError qId op = \case Left e -> logError $ logPfx qId op <> tshow e Right _ -> pure () diff --git a/tests/CoreTests/MsgStoreTests.hs b/tests/CoreTests/MsgStoreTests.hs index 05878cc14..a05082ffd 100644 --- a/tests/CoreTests/MsgStoreTests.hs +++ b/tests/CoreTests/MsgStoreTests.hs @@ -28,13 +28,14 @@ import Data.ByteString.Char8 (ByteString) import qualified Data.ByteString.Char8 as B import Data.Int (Int64) import Data.List (isPrefixOf, isSuffixOf) +import qualified Data.Map.Strict as M import Data.Maybe (fromJust) import Data.Time.Clock (addUTCTime) import Data.Time.Clock.System (SystemTime (..), getSystemTime) import SMPClient (testStoreLogFile, testStoreMsgsDir, testStoreMsgsDir2, testStoreMsgsFile, testStoreMsgsFile2) import Simplex.Messaging.Crypto (pattern MaxLenBS) import qualified Simplex.Messaging.Crypto as C -import Simplex.Messaging.Protocol (EntityId (..), ErrorType, LinkId, Message (..), QueueLinkData, RecipientId, SParty (..), noMsgFlags) +import Simplex.Messaging.Protocol (EncDataBytes (..), EntityId (..), ErrorType (..), LinkId, Message (..), NotifierId, QueueLinkData, RecipientId, SParty (..), SenderId, noMsgFlags) import Simplex.Messaging.Server (exportMessages, importMessages, printMessageStats) import Simplex.Messaging.Server.Env.STM (MsgStore (..), journalMsgStoreDepth, readWriteQueueStore) import Simplex.Messaging.Server.Expiration (ExpirationConfig (..), expireBeforeEpoch) @@ -43,7 +44,10 @@ import Simplex.Messaging.Server.MsgStore.STM import Simplex.Messaging.Server.MsgStore.Types import Simplex.Messaging.Server.QueueStore import Simplex.Messaging.Server.QueueStore.QueueInfo +import Simplex.Messaging.Server.QueueStore.STM (STMQueueStore (..)) +import Simplex.Messaging.Server.QueueStore.Types import Simplex.Messaging.Server.StoreLog (closeStoreLog, logCreateQueue) +import Simplex.Messaging.TMap (TMap) import System.Directory (copyFile, createDirectoryIfMissing, listDirectory, removeFile, renameFile) import System.FilePath (()) import System.IO (IOMode (..), withFile) @@ -57,7 +61,6 @@ import Simplex.Messaging.Agent.Store.Postgres.Common import Simplex.Messaging.Agent.Store.Shared (MigrationConfirmation (..)) import Simplex.Messaging.Server.MsgStore.Postgres import Simplex.Messaging.Server.QueueStore.Postgres -import Simplex.Messaging.Server.QueueStore.Types import SMPClient (postgressBracket, testServerDBConnectInfo, testStoreDBOpts) #endif @@ -68,12 +71,14 @@ msgStoreTests = do someMsgStoreTests journalMsgStoreTests it "should export and import journal store" testExportImportStore + it "should remove deleted queues from queue store maps" $ testDeleteQueueMaps stmQueueMapSizes #if defined(dbServerPostgres) around_ (postgressBracket testServerDBConnectInfo) $ do around (withMsgStore $ testJournalStoreCfg $ PQStoreCfg testPostgresStoreCfg) $ describe "Postgres+journal message store" $ do someMsgStoreTests journalMsgStoreTests + it "should remove deleted queues from queue cache maps" $ testDeleteQueueMaps postgresQueueMapSizes around (withMsgStore testPostgresStoreConfig) $ describe "Postgres-only message store" $ do someMsgStoreTests @@ -184,6 +189,18 @@ testNewQueueRecData g qm queueData = do where rndId = atomically $ EntityId <$> C.randomBytes 24 g +testNtfCreds :: TVar ChaChaDRG -> IO NtfCreds +testNtfCreds g = do + (notifierKey, _) <- atomically $ C.generateAuthKeyPair C.SX25519 g + (k, pk) <- atomically $ C.generateKeyPair @'C.X25519 g + pure + NtfCreds + { notifierId = EntityId "ijkl", + notifierKey, + rcvNtfDhSecret = C.dh' k pk, + ntfServiceId = Nothing + } + testGetQueue :: MsgStoreClass s => s -> IO () testGetQueue ms = do g <- C.newRandom @@ -319,7 +336,64 @@ testExportImportStore ms = do exportMessages False (StoreMemory stmStore) testStoreMsgsFile False (B.sort <$> B.readFile testStoreMsgsFile `shouldReturn`) =<< (B.sort <$> B.readFile (testStoreMsgsFile2 <> ".bak")) +-- sizes of queues, senders, links and notifiers maps +type QueueMapSizes = (Int, Int, Int, Int) + +stmQueueMapSizes :: JournalMsgStore 'QSMemory -> IO QueueMapSizes +stmQueueMapSizes ms = queueMapSizes queues senders links notifiers + where + STMQueueStore {queues, senders, links, notifiers} = stmQueueStore ms + +queueMapSizes :: TMap RecipientId q -> TMap SenderId RecipientId -> TMap LinkId RecipientId -> TMap NotifierId RecipientId -> IO QueueMapSizes +queueMapSizes qs ss ls ns = (,,,) <$> size qs <*> size ss <*> size ls <*> size ns + where + size :: TMap k v -> IO Int + size = fmap M.size . readTVarIO + +testDeleteQueueMaps :: forall s. MsgStoreClass s => (s -> IO QueueMapSizes) -> s -> IO () +testDeleteQueueMaps mapSizes ms = do + g <- C.newRandom + ntfCreds <- testNtfCreds g + let qd = (EncDataBytes "fixed data", EncDataBytes "user data") + newLinkId = atomically $ EntityId <$> C.randomBytes 24 g + lnkId1 <- newLinkId + lnkId2 <- newLinkId + lnkId3 <- newLinkId + (rId1, qr1) <- testNewQueueRec g QMMessaging + (rId2, qr2) <- testNewQueueRecData g QMContact (Just (lnkId1, qd)) + (rId3, qr3) <- testNewQueueRec g QMMessaging + (rId4, qr4) <- testNewQueueRec g QMMessaging + let rIds = [rId1, rId2, rId3, rId4] :: [RecipientId] + sIds = map senderId [qr1, qr2, qr3, qr4] + lnkIds = [lnkId1, lnkId2, lnkId3] :: [LinkId] + mapSizes ms `shouldReturn` (0, 0, 0, 0) + runRight_ $ do + q1 <- ExceptT $ addQueue ms rId1 qr1 {notifier = Just ntfCreds} + q2 <- ExceptT $ addQueue ms rId2 qr2 + q3 <- ExceptT $ addQueue ms rId3 qr3 + q4 <- ExceptT $ addQueue ms rId4 qr4 + ExceptT $ addQueueLinkData (queueStore ms) q3 lnkId2 qd + ExceptT $ addQueueLinkData (queueStore ms) q4 lnkId3 qd + forM_ sIds $ void . ExceptT . getQueue ms SSender + forM_ lnkIds $ void . ExceptT . getQueue ms SSenderLink + liftIO $ mapSizes ms `shouldReturn` (4, 4, 3, 1) + ExceptT $ deleteQueueLinkData (queueStore ms) q3 + liftIO $ mapSizes ms `shouldReturn` (4, 4, 2, 1) + forM_ ([q1, q2, q3, q4] :: [StoreQueue s]) $ void . ExceptT . deleteQueue ms + mapSizes ms `shouldReturn` (0, 0, 0, 0) + forM_ rIds $ \rId -> getQueue ms SRecipient rId >>= expectAuth + forM_ sIds $ \sId -> getQueue ms SSender sId >>= expectAuth + forM_ lnkIds $ \lnkId -> getQueue ms SSenderLink lnkId >>= expectAuth + mapSizes ms `shouldReturn` (0, 0, 0, 0) + where + expectAuth = either (`shouldBe` AUTH) (\_ -> expectationFailure "deleted queue is still found") + #if defined(dbServerPostgres) +postgresQueueMapSizes :: JournalMsgStore 'QSPostgres -> IO QueueMapSizes +postgresQueueMapSizes ms = queueMapSizes queues senders links notifiers + where + PostgresQueueStore {queues, senders, links, notifiers} = postgresQueueStore ms + testUpdateMessageCounts :: PostgresMsgStore -> IO () testUpdateMessageCounts ms = do g <- C.newRandom diff --git a/tests/CoreTests/StoreLogTests.hs b/tests/CoreTests/StoreLogTests.hs index 81a0535d5..ac10231e4 100644 --- a/tests/CoreTests/StoreLogTests.hs +++ b/tests/CoreTests/StoreLogTests.hs @@ -48,18 +48,6 @@ import Simplex.Messaging.Server.Main testPublicAuthKey :: C.APublicAuthKey testPublicAuthKey = C.APublicAuthKey C.SEd25519 (C.publicKey "MC4CAQAwBQYDK2VwBCIEIDfEfevydXXfKajz3sRkcQ7RPvfWUPoq6pu1TYHV1DEe") -testNtfCreds :: TVar ChaChaDRG -> IO NtfCreds -testNtfCreds g = do - (notifierKey, _) <- atomically $ C.generateAuthKeyPair C.SX25519 g - (k, pk) <- atomically $ C.generateKeyPair @'C.X25519 g - pure - NtfCreds - { notifierId = EntityId "ijkl", - notifierKey, - rcvNtfDhSecret = C.dh' k pk, - ntfServiceId = Nothing - } - data StoreLogTestCase r s = SLTC {name :: String, saved :: [r], state :: s, compacted :: [r]} type SMPStoreLogTestCase = StoreLogTestCase StoreLogRecord (M.Map RecipientId QueueRec)