mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-03 16:58:17 +00:00
smp-server: remove deleted queues from store maps
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 ()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user