From 8eade7ef1bb949e44f9391a2439b001045657a5e Mon Sep 17 00:00:00 2001 From: shum Date: Sun, 4 Oct 2026 11:20:48 +0000 Subject: [PATCH] smp-server: fix queue cache collisions and leaks --- .../Messaging/Server/MsgStore/Journal.hs | 7 +- .../Messaging/Server/MsgStore/Postgres.hs | 2 + src/Simplex/Messaging/Server/MsgStore/STM.hs | 2 + .../Messaging/Server/QueueStore/Postgres.hs | 98 +++++++++------ .../Messaging/Server/QueueStore/Types.hs | 2 + tests/CoreTests/MsgStoreTests.hs | 119 +++++++++++++++--- 6 files changed, 172 insertions(+), 58 deletions(-) diff --git a/src/Simplex/Messaging/Server/MsgStore/Journal.hs b/src/Simplex/Messaging/Server/MsgStore/Journal.hs index 185c113b7..8e9005091 100644 --- a/src/Simplex/Messaging/Server/MsgStore/Journal.hs +++ b/src/Simplex/Messaging/Server/MsgStore/Journal.hs @@ -161,6 +161,7 @@ data QStoreCfg s where data JournalQueue (s :: QSType) = JournalQueue { recipientId' :: RecipientId, queueLock :: Lock, + queueLocks' :: TMap RecipientId Lock, sharedLock :: TMVar RecipientId, -- To avoid race conditions and errors when restoring queues, -- Nothing is written to TVar when queue is deleted. @@ -300,6 +301,9 @@ instance StoreQueueClass (JournalQueue s) where withQueueLock JournalQueue {recipientId', queueLock, sharedLock} = withLockWaitShared recipientId' queueLock sharedLock {-# INLINE withQueueLock #-} + removeQueueLock :: JournalQueue s -> IO () + removeQueueLock JournalQueue {recipientId', queueLock, queueLocks'} = + atomically $ TM.lookup recipientId' queueLocks' >>= \l -> when (l == Just queueLock) $ TM.delete recipientId' queueLocks' instance QueueStoreClass (JournalQueue s) (QStore s) where type QueueStoreCfg (QStore s) = QStoreCfg s @@ -361,7 +365,7 @@ instance QueueStoreClass (JournalQueue s) (QStore s) where {-# INLINE getServiceQueueCountHash #-} makeQueue_ :: JournalMsgStore s -> RecipientId -> QueueRec -> Lock -> IO (JournalQueue s) -makeQueue_ JournalMsgStore {sharedLock} rId qr queueLock = do +makeQueue_ JournalMsgStore {queueLocks, sharedLock} rId qr queueLock = do queueRec' <- newTVarIO $ Just qr msgQueue' <- newTVarIO Nothing activeAt <- newTVarIO 0 @@ -370,6 +374,7 @@ makeQueue_ JournalMsgStore {sharedLock} rId qr queueLock = do JournalQueue { recipientId' = rId, queueLock, + queueLocks' = queueLocks, sharedLock, queueRec', msgQueue', diff --git a/src/Simplex/Messaging/Server/MsgStore/Postgres.hs b/src/Simplex/Messaging/Server/MsgStore/Postgres.hs index a79643f08..be99fa265 100644 --- a/src/Simplex/Messaging/Server/MsgStore/Postgres.hs +++ b/src/Simplex/Messaging/Server/MsgStore/Postgres.hs @@ -83,6 +83,8 @@ instance StoreQueueClass PostgresQueue where {-# INLINE queueRec #-} withQueueLock PostgresQueue {} _ = id -- TODO [messages] maybe it's just transaction? {-# INLINE withQueueLock #-} + removeQueueLock _ = pure () + {-# INLINE removeQueueLock #-} newtype DBTransaction = DBTransaction {dbConn :: DB.Connection} diff --git a/src/Simplex/Messaging/Server/MsgStore/STM.hs b/src/Simplex/Messaging/Server/MsgStore/STM.hs index f118e007c..8388efd34 100644 --- a/src/Simplex/Messaging/Server/MsgStore/STM.hs +++ b/src/Simplex/Messaging/Server/MsgStore/STM.hs @@ -63,6 +63,8 @@ instance StoreQueueClass STMQueue where {-# INLINE queueRec #-} withQueueLock _ _ = id {-# INLINE withQueueLock #-} + removeQueueLock _ = pure () + {-# INLINE removeQueueLock #-} instance MsgStoreClass STMMsgStore where type StoreMonad STMMsgStore = STM diff --git a/src/Simplex/Messaging/Server/QueueStore/Postgres.hs b/src/Simplex/Messaging/Server/QueueStore/Postgres.hs index 80a49a64f..a0bea33f0 100644 --- a/src/Simplex/Messaging/Server/QueueStore/Postgres.hs +++ b/src/Simplex/Messaging/Server/QueueStore/Postgres.hs @@ -107,7 +107,6 @@ data PostgresQueueStore q = PostgresQueueStore queues :: TMap RecipientId q, -- this map only cashes the queues that were attempted to send messages to, senders :: TMap SenderId RecipientId, - links :: TMap LinkId RecipientId, -- this map only cashes the queues that were attempted to be subscribed to, notifiers :: TMap NotifierId RecipientId, notifierLocks :: TMap NotifierId Lock, @@ -127,11 +126,10 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where dbStoreLog <- mapM (openWriteStoreLog True) dbStoreLogPath queues <- TM.emptyIO senders <- TM.emptyIO - links <- TM.emptyIO notifiers <- TM.emptyIO notifierLocks <- TM.emptyIO serviceLocks <- TM.emptyIO - pure PostgresQueueStore {dbStore, dbStoreLog, queues, senders, links, notifiers, notifierLocks, serviceLocks, deletedTTL, useCache} + pure PostgresQueueStore {dbStore, dbStoreLog, queues, senders, notifiers, notifierLocks, serviceLocks, deletedTTL, useCache} where err e = do logError $ "STORE: newQueueStore, error opening PostgreSQL database, " <> tshow e @@ -174,7 +172,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where -- and relies on unique constraints in the database to prevent duplicate IDs. addQueue_ :: PostgresQueueStore q -> (RecipientId -> QueueRec -> IO q) -> RecipientId -> QueueRec -> IO (Either ErrorType q) addQueue_ st mkQ rId qr = do - sq <- mkQ rId qr + sq <- mkQ rId qr {queueData = withoutLinkData . fst <$> queueData qr} withQueueLock sq "addQueue_" $ E.uninterruptibleMask_ $ runExceptT $ do void $ withDB "addQueue_" st $ \db -> E.try (DB.execute db insertQueueQuery $ queueRecToRow (rId, qr)) @@ -183,11 +181,10 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where atomically $ TM.insert rId sq queues atomically $ TM.insert (senderId qr) rId senders forM_ (notifier qr) $ \NtfCreds {notifierId = nId} -> atomically $ TM.insert nId rId notifiers - forM_ (queueData qr) $ \(lnkId, _) -> atomically $ TM.insert lnkId rId links withLog "addStoreQueue" st $ \s -> logCreateQueue s rId qr pure sq where - PostgresQueueStore {queues, senders, links, notifiers, useCache} = st + PostgresQueueStore {queues, senders, notifiers, useCache} = st -- Not doing duplicate checks in maps as the probability of duplicates is very low. -- It needs to be reconsidered when IDs are supplied by the users. -- hasId = anyM [TM.memberIO rId queues, TM.memberIO senderId senders, hasNotifier] @@ -198,7 +195,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where | useCache = case party of SRecipient -> getRcvQueue qId SSender -> TM.lookupIO qId senders >>= maybe (mask loadSndQueue) getRcvQueue - SSenderLink -> TM.lookupIO qId links >>= maybe (mask loadLinkQueue) getRcvQueue + SSenderLink -> mask loadLinkQueue -- loaded queue is deleted from notifiers map to reduce cache size after queue was subscribed to by ntf server SNotifier -> TM.lookupIO qId notifiers >>= maybe (mask loadNtfQueue) (getRcvQueue >=> (atomically (TM.delete qId notifiers) $>)) | otherwise = case party of @@ -207,44 +204,53 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where SSenderLink -> loadQueueNoCache " WHERE link_id = ?" SNotifier -> loadQueueNoCache " WHERE notifier_id = ?" where - PostgresQueueStore {queues, senders, links, notifiers, useCache} = st + PostgresQueueStore {queues, senders, notifiers, useCache} = st getRcvQueue rId = TM.lookupIO rId queues >>= maybe (mask loadRcvQueue) (pure . Right) loadRcvQueue = do (rId, qRec) <- loadQueue " WHERE recipient_id = ?" - liftIO $ cacheQueue rId qRec $ \_ -> pure () -- recipient map already checked, not caching sender ref - loadSndQueue = loadSndQueue_ senders " WHERE sender_id = ?" - loadLinkQueue = loadSndQueue_ links " WHERE link_id = ?" + cacheQueue rId qRec $ \_ -> pure () -- recipient map already checked, not caching sender ref + loadSndQueue = do + (rId, qRec) <- loadQueue " WHERE sender_id = ?" + -- checking recipient map first, sender ref is only cached for a queue in the map + atomically (TM.lookup rId queues >>= mapM (\sq -> sq <$ cacheSender rId)) + >>= maybe (cacheQueue rId qRec cacheSender) pure + -- link IDs are supplied by clients, they are not cached to prevent collisions with sender IDs + loadLinkQueue = do + (rId, qRec) <- loadQueue " WHERE link_id = ?" + liftIO (TM.lookupIO rId queues) >>= maybe (cacheQueue rId qRec $ \_ -> pure ()) pure 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_ refs condition = do - (rId, qRec) <- loadQueue condition - liftIO $ - TM.lookupIO rId queues -- checking recipient map first - >>= maybe (cacheQueue rId qRec $ cacheRef refs) (atomically (cacheRef refs rId) $>) loadQueueNoCache cond = mask $ loadQueue cond >>= liftIO . uncurry (mkQ True) mask = E.uninterruptibleMask_ . runExceptT - cacheRef refs rId = TM.insert qId rId refs - loadQueue condition = + cacheSender rId = TM.insert qId rId senders + loadQueue condition = loadQueueBy condition qId + loadQueueBy condition qId' = withDB "getQueue_" st $ \db -> firstRow rowToQueueRec AUTH $ - DB.query db (queueRecQuery <> condition <> " AND deleted_at IS NULL") (Only qId) + DB.query db (queueRecQuery <> condition <> " AND deleted_at IS NULL") (Only qId') cacheQueue rId qRec insertRef = do - sq <- mkQ True rId qRec -- loaded queue + sq <- liftIO $ mkQ True rId qRec -- loaded queue -- This lock prevents the scenario when the queue is added to cache, -- while another thread is proccessing the same queue in withAllMsgQueues -- without adding it to cache, possibly trying to open the same files twice. -- Alse see comment in idleDeleteExpiredMsgs. - withQueueLock sq "getQueue_" $ atomically $ - -- checking the cache again for concurrent reads, - -- use previously loaded queue if exists. - TM.lookup rId queues >>= \case - Just sq' -> pure sq' - Nothing -> do - insertRef rId - TM.insert rId sq queues - pure sq + ExceptT $ withQueueLock sq "getQueue_" $ runExceptT $ do + -- the queue could have been deleted after it was loaded and before the lock was taken + (_, qRec') <- loadQueueBy " WHERE recipient_id = ?" rId `catchE` \case + AUTH -> liftIO (removeQueueLock sq) >> throwE AUTH + e -> throwE e + atomically $ + -- checking the cache again for concurrent reads, + -- use previously loaded queue if exists. + TM.lookup rId queues >>= \case + Just sq' -> pure sq' + Nothing -> do + writeTVar (queueRec sq) $ Just qRec' + insertRef rId + TM.insert rId sq queues + pure sq getQueues_ :: forall p. BatchParty p => PostgresQueueStore q -> (Bool -> RecipientId -> QueueRec -> IO q) -> SParty p -> [QueueId] -> IO [Either ErrorType q] getQueues_ st mkQ party qIds @@ -253,7 +259,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where SRecipient -> do qs <- readTVarIO queues let qs' = map (\qId -> get qs qId qId) qIds - E.uninterruptibleMask_ $ loadQueues qs' " WHERE recipient_id IN ?" cacheRcvQueue + E.uninterruptibleMask_ $ loadQueues qs' " WHERE recipient_id IN ?" cacheRcvQueue >>= uncacheDeleted (lefts qs') SNotifier -> do ns <- readTVarIO notifiers qs <- readTVarIO queues @@ -294,6 +300,21 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where Just sq' -> pure sq' Nothing -> sq <$ TM.insert rId sq queues pure $ Just (rId, sq') + -- the queues could have been deleted after they were loaded and before they were cached + uncacheDeleted :: [RecipientId] -> [Either ErrorType q] -> IO [Either ErrorType q] + uncacheDeleted qIds' rs = case filter (`S.member` S.fromList qIds') [recipientId sq | Right sq <- rs] of + [] -> pure rs + rIds -> + runExceptT (withDB' "getQueues_" st $ \db -> DB.query db "SELECT recipient_id FROM msg_queues WHERE recipient_id IN ? AND deleted_at IS NULL" (Only (In rIds))) >>= \case + Right live -> mapM (uncache $ S.fromList rIds `S.difference` S.fromList (map fromOnly live)) rs + Left _ -> pure rs + uncache deleted = \case + Right sq | S.member (recipientId sq) deleted -> do + withQueueLock sq "getQueues_" $ do + atomically $ writeTVar (queueRec sq) Nothing >> TM.delete (recipientId sq) queues + removeQueueLock sq + pure $ Left AUTH + r -> pure r loadQueuesNoCache cond mkQueue' = do qs_ <- dbLoadQueues qIds cond mkQueue' pure $ map (result qs_) qIds @@ -322,17 +343,16 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where rId = recipientId sq addLink q update = do assertUpdated $ withDB' "addQueueLinkData" st update - atomically $ writeTVar (queueRec sq) $ Just q {queueData = Just (lnkId, d)} + atomically $ writeTVar (queueRec sq) $ Just q {queueData = Just $ withoutLinkData lnkId} withLog "addQueueLinkData" st $ \s -> logCreateLink s rId lnkId d qry = "UPDATE msg_queues SET fixed_data = ?, user_data = ?, link_id = ? WHERE recipient_id = ? AND deleted_at IS NULL" deleteQueueLinkData :: PostgresQueueStore q -> q -> IO (Either ErrorType ()) deleteQueueLinkData st sq = withQueueRec sq "deleteQueueLinkData" $ \q -> case queueData q of - Just (lnkId, _) -> do + Just _ -> 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 @@ -458,9 +478,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 + atomically $ do + TM.delete rId $ queues st + TM.delete (senderId q) $ senders st forM_ (notifier q) $ \NtfCreds {notifierId} -> do atomically $ TM.delete notifierId $ notifiers st atomically $ TM.delete notifierId $ notifierLocks st @@ -728,7 +748,7 @@ queueDataColumns = \case rowToQueueRec :: QueueRecRow -> (RecipientId, QueueRec) rowToQueueRec (rId, recipientKeys, rcvDhSecret, senderId, senderKey, queueMode, notifierId_, notifierKey_, rcvNtfDhSecret_, ntfServiceId, status, updatedAt, linkId_, rcvServiceId) = let notifier = mkNotifier (notifierId_, notifierKey_, rcvNtfDhSecret_) ntfServiceId - queueData = (,(EncDataBytes "", EncDataBytes "")) <$> linkId_ + queueData = withoutLinkData <$> linkId_ in (rId, QueueRec {recipientKeys, rcvDhSecret, senderId, senderKey, queueMode, queueData, notifier, status, updatedAt, rcvServiceId}) rowToQueueRecWithData :: QueueRecRow :. (Maybe EncDataBytes, Maybe EncDataBytes) -> (RecipientId, QueueRec) @@ -738,6 +758,10 @@ rowToQueueRecWithData ((rId, recipientKeys, rcvDhSecret, senderId, senderKey, qu queueData = (,(encData immutableData_, encData userData_)) <$> linkId_ in (rId, QueueRec {recipientKeys, rcvDhSecret, senderId, senderKey, queueMode, queueData, notifier, status, updatedAt, rcvServiceId}) +-- link data is read from the database when requested +withoutLinkData :: LinkId -> (LinkId, QueueLinkData) +withoutLinkData = (,(EncDataBytes "", EncDataBytes "")) + mkNotifier :: (Maybe NotifierId, Maybe NtfPublicAuthKey, Maybe RcvNtfDhSecret) -> Maybe ServiceId -> Maybe NtfCreds mkNotifier (Just notifierId, Just notifierKey, Just rcvNtfDhSecret) ntfServiceId = Just NtfCreds {notifierId, notifierKey, rcvNtfDhSecret, ntfServiceId} diff --git a/src/Simplex/Messaging/Server/QueueStore/Types.hs b/src/Simplex/Messaging/Server/QueueStore/Types.hs index 415a5f33c..d105357bf 100644 --- a/src/Simplex/Messaging/Server/QueueStore/Types.hs +++ b/src/Simplex/Messaging/Server/QueueStore/Types.hs @@ -27,6 +27,8 @@ class StoreQueueClass q where recipientId :: q -> RecipientId queueRec :: q -> TVar (Maybe QueueRec) withQueueLock :: q -> Text -> IO a -> IO a + -- must only be called for deleted queues + removeQueueLock :: q -> IO () class StoreQueueClass q => QueueStoreClass q s where type QueueStoreCfg s diff --git a/tests/CoreTests/MsgStoreTests.hs b/tests/CoreTests/MsgStoreTests.hs index a05082ffd..d56084cbe 100644 --- a/tests/CoreTests/MsgStoreTests.hs +++ b/tests/CoreTests/MsgStoreTests.hs @@ -18,6 +18,7 @@ module CoreTests.MsgStoreTests where import AgentTests.FunctionalAPITests (runRight, runRight_) import Control.Concurrent (threadDelay) +import Control.Concurrent.Async (concurrently) import Control.Concurrent.STM import Control.Exception (bracket) import Control.Monad @@ -26,10 +27,11 @@ import Control.Monad.Trans.Except import Crypto.Random (ChaChaDRG) import Data.ByteString.Char8 (ByteString) import qualified Data.ByteString.Char8 as B +import Data.Either (isRight) import Data.Int (Int64) import Data.List (isPrefixOf, isSuffixOf) import qualified Data.Map.Strict as M -import Data.Maybe (fromJust) +import Data.Maybe (fromJust, isNothing) import Data.Time.Clock (addUTCTime) import Data.Time.Clock.System (SystemTime (..), getSystemTime) import SMPClient (testStoreLogFile, testStoreMsgsDir, testStoreMsgsDir2, testStoreMsgsFile, testStoreMsgsFile2) @@ -48,6 +50,7 @@ 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 qualified Simplex.Messaging.TMap as TM import System.Directory (copyFile, createDirectoryIfMissing, listDirectory, removeFile, renameFile) import System.FilePath (()) import System.IO (IOMode (..), withFile) @@ -71,20 +74,23 @@ msgStoreTests = do someMsgStoreTests journalMsgStoreTests it "should export and import journal store" testExportImportStore - it "should remove deleted queues from queue store maps" $ testDeleteQueueMaps stmQueueMapSizes + it "should remove deleted queues from queue store maps" $ testDeleteQueueMaps stmQueueMapSizes (Just stmLinksSize) #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 + it "should remove deleted queues from queue cache maps" $ testDeleteQueueMaps postgresQueueMapSizes Nothing + it "should not cache queue deleted while loading" testDeletedQueueNotCached + it "should not keep link data in queue records" testQueueRecNoLinkData around (withMsgStore testPostgresStoreConfig) $ describe "Postgres-only message store" $ do someMsgStoreTests it "should correctly update message counts and canWrite flag" testUpdateMessageCounts it "tryDelPeekMsg (ACK not from NSE) should reset message counts when queue is empty" testResetMessageCounts it "should expire messages across commit batches" testExpireMessagesInBatches + it "should not keep link data in queue records" testQueueRecNoLinkData #endif describe "Journal message store: queue state backup expiration" $ do it "should remove old queue state backups" testRemoveQueueStateBackups @@ -106,6 +112,7 @@ msgStoreTests = do it "should get queue and store/read messages" testGetQueue it "should write/ack messages" testWriteAckMessages it "should not fail on EOF when changing read journal" testChangeReadJournal + it "should resolve sender ID equal to link ID of another queue" testLinkIdSenderIdCollision -- TODO constrain to STM stores? withMsgStore :: MsgStoreClass s => MsgStoreConfig s -> (s -> IO ()) -> IO () @@ -336,22 +343,27 @@ 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) +-- sizes of queues, senders and notifiers maps +type QueueMapSizes = (Int, Int, Int) stmQueueMapSizes :: JournalMsgStore 'QSMemory -> IO QueueMapSizes -stmQueueMapSizes ms = queueMapSizes queues senders links notifiers +stmQueueMapSizes ms = queueMapSizes queues senders notifiers where - STMQueueStore {queues, senders, links, notifiers} = stmQueueStore ms + STMQueueStore {queues, senders, 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 +stmLinksSize :: JournalMsgStore 'QSMemory -> IO Int +stmLinksSize ms = mapSize links where - size :: TMap k v -> IO Int - size = fmap M.size . readTVarIO + STMQueueStore {links} = stmQueueStore ms -testDeleteQueueMaps :: forall s. MsgStoreClass s => (s -> IO QueueMapSizes) -> s -> IO () -testDeleteQueueMaps mapSizes ms = do +queueMapSizes :: TMap RecipientId q -> TMap SenderId RecipientId -> TMap NotifierId RecipientId -> IO QueueMapSizes +queueMapSizes qs ss ns = (,,) <$> mapSize qs <*> mapSize ss <*> mapSize ns + +mapSize :: TMap k v -> IO Int +mapSize = fmap M.size . readTVarIO + +testDeleteQueueMaps :: forall s. MsgStoreClass s => (s -> IO QueueMapSizes) -> Maybe (s -> IO Int) -> s -> IO () +testDeleteQueueMaps mapSizes linksSize_ ms = do g <- C.newRandom ntfCreds <- testNtfCreds g let qd = (EncDataBytes "fixed data", EncDataBytes "user data") @@ -366,7 +378,7 @@ testDeleteQueueMaps mapSizes ms = do 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) + sizesShouldBe (0, 0, 0) 0 runRight_ $ do q1 <- ExceptT $ addQueue ms rId1 qr1 {notifier = Just ntfCreds} q2 <- ExceptT $ addQueue ms rId2 qr2 @@ -376,23 +388,90 @@ testDeleteQueueMaps mapSizes ms = do 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) + liftIO $ sizesShouldBe (4, 4, 1) 3 ExceptT $ deleteQueueLinkData (queueStore ms) q3 - liftIO $ mapSizes ms `shouldReturn` (4, 4, 2, 1) + liftIO $ sizesShouldBe (4, 4, 1) 2 forM_ ([q1, q2, q3, q4] :: [StoreQueue s]) $ void . ExceptT . deleteQueue ms - mapSizes ms `shouldReturn` (0, 0, 0, 0) + sizesShouldBe (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) + sizesShouldBe (0, 0, 0) 0 where expectAuth = either (`shouldBe` AUTH) (\_ -> expectationFailure "deleted queue is still found") + sizesShouldBe sizes linksSize = do + mapSizes ms `shouldReturn` sizes + forM_ linksSize_ $ \f -> f ms `shouldReturn` linksSize + +testLinkIdSenderIdCollision :: MsgStoreClass s => s -> IO () +testLinkIdSenderIdCollision ms = do + g <- C.newRandom + (rIdV, qrV) <- testNewQueueRec g QMContact + (rIdA, qrA) <- testNewQueueRec g QMContact + let sIdV = senderId qrV + runRight_ $ do + void $ ExceptT $ addQueue ms rIdV qrV + qA <- ExceptT $ addQueue ms rIdA qrA + ExceptT $ addQueueLinkData (queueStore ms) qA sIdV (EncDataBytes "fixed data", EncDataBytes "user data") + qA' <- ExceptT $ getQueue ms SSenderLink sIdV + liftIO $ recipientId qA' `shouldBe` rIdA + qV <- ExceptT $ getQueue ms SSender sIdV + liftIO $ recipientId qV `shouldBe` rIdV + +testQueueRecNoLinkData :: MsgStoreClass s => s -> IO () +testQueueRecNoLinkData ms = do + g <- C.newRandom + let qd = (EncDataBytes "fixed data", EncDataBytes "user data") + qd' = (EncDataBytes "fixed data", EncDataBytes "updated user data") + noData = (EncDataBytes "", EncDataBytes "") + newLinkId = atomically $ EntityId <$> C.randomBytes 24 g + lnkId1 <- newLinkId + lnkId2 <- newLinkId + (rId1, qr1) <- testNewQueueRecData g QMContact (Just (lnkId1, qd)) + (rId2, qr2) <- testNewQueueRec g QMContact + runRight_ $ do + q1 <- ExceptT $ addQueue ms rId1 qr1 + q2 <- ExceptT $ addQueue ms rId2 qr2 + ExceptT $ addQueueLinkData (queueStore ms) q2 lnkId2 qd + liftIO $ queueLinkData q1 `shouldReturn` Just (lnkId1, noData) + liftIO $ queueLinkData q2 `shouldReturn` Just (lnkId2, noData) + ExceptT (getQueueLinkData (queueStore ms) q1 lnkId1) >>= liftIO . (`shouldBe` qd) + ExceptT (getQueueLinkData (queueStore ms) q2 lnkId2) >>= liftIO . (`shouldBe` qd) + ExceptT $ addQueueLinkData (queueStore ms) q2 lnkId2 qd' + liftIO $ queueLinkData q2 `shouldReturn` Just (lnkId2, noData) + ExceptT (getQueueLinkData (queueStore ms) q2 lnkId2) >>= liftIO . (`shouldBe` qd') + where + queueLinkData q = (queueData =<<) <$> readTVarIO (queueRec q) #if defined(dbServerPostgres) postgresQueueMapSizes :: JournalMsgStore 'QSPostgres -> IO QueueMapSizes -postgresQueueMapSizes ms = queueMapSizes queues senders links notifiers +postgresQueueMapSizes ms = queueMapSizes queues senders notifiers where - PostgresQueueStore {queues, senders, links, notifiers} = postgresQueueStore ms + PostgresQueueStore {queues, senders, notifiers} = postgresQueueStore ms + +testDeletedQueueNotCached :: JournalMsgStore 'QSPostgres -> IO () +testDeletedQueueNotCached ms = do + g <- C.newRandom + -- the queue is cached without sender reference, as after subscription + loadWhileDeleting g getSndQueue $ \_ sId -> TM.delete sId senders + -- the queue is only in the database, as after server restart + loadWhileDeleting g getSndQueue evictQueue + loadWhileDeleting g getRcvQueues evictQueue + where + PostgresQueueStore {queues, senders} = postgresQueueStore ms + getSndQueue _ sId = getQueue ms SSender sId + getRcvQueues rId _ = head <$> getQueues ms SRecipient [rId] + evictQueue rId sId = TM.delete rId queues >> TM.delete sId senders + loadWhileDeleting g load evict = replicateM_ 100 $ do + (rId, qr) <- testNewQueueRec g QMMessaging + q <- either (fail . show) pure =<< addQueue ms rId qr + atomically $ evict rId (senderId qr) + (q_, deleted) <- concurrently (load rId (senderId qr)) (deleteQueue ms q) + deleted `shouldSatisfy` isRight + forM_ q_ $ \q' -> readTVarIO (queueRec q') >>= (`shouldSatisfy` isNothing) + TM.memberIO rId queues `shouldReturn` False + TM.lookupIO (senderId qr) senders `shouldReturn` Nothing + queueLockCount <$> loadedQueueCounts ms `shouldReturn` 0 testUpdateMessageCounts :: PostgresMsgStore -> IO () testUpdateMessageCounts ms = do