mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-07 16:17:25 +00:00
smp-server: remove queue lock when queue is not added
This commit is contained in:
@@ -176,7 +176,7 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where
|
||||
withQueueLock sq "addQueue_" $ E.uninterruptibleMask_ $ runExceptT $ do
|
||||
void $ withDB "addQueue_" st $ \db ->
|
||||
E.try (DB.execute db insertQueueQuery $ queueRecToRow (rId, qr))
|
||||
>>= bimapM handleDuplicate pure
|
||||
>>= bimapM (\e -> unless (isRecipientIdViolation e) (removeQueueLock sq) >> handleDuplicate e) pure
|
||||
when useCache $ do
|
||||
atomically $ TM.insert rId sq queues
|
||||
atomically $ TM.insert (senderId qr) rId senders
|
||||
@@ -185,6 +185,8 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where
|
||||
pure sq
|
||||
where
|
||||
PostgresQueueStore {queues, senders, notifiers, useCache} = st
|
||||
-- the lock is kept when another queue has the same recipient ID
|
||||
isRecipientIdViolation e = constraintViolation e == Just (UniqueViolation "msg_queues_pkey")
|
||||
-- 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]
|
||||
@@ -209,20 +211,19 @@ instance StoreQueueClass q => QueueStoreClass q (PostgresQueueStore q) where
|
||||
loadRcvQueue = do
|
||||
(rId, qRec) <- loadQueue " WHERE recipient_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
|
||||
loadSndQueue = loadSndQueue_ " WHERE sender_id = ?" cacheSender
|
||||
-- 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
|
||||
loadLinkQueue = loadSndQueue_ " WHERE link_id = ?" $ \_ -> 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_ condition insertRef = do
|
||||
(rId, qRec) <- loadQueue condition
|
||||
-- checking recipient map first, ref is only cached for the queue in the map
|
||||
atomically (TM.lookup rId queues >>= mapM (\sq -> sq <$ insertRef rId))
|
||||
>>= maybe (cacheQueue rId qRec insertRef) pure
|
||||
loadQueueNoCache cond = mask $ loadQueue cond >>= liftIO . uncurry (mkQ True)
|
||||
mask = E.uninterruptibleMask_ . runExceptT
|
||||
cacheSender rId = TM.insert qId rId senders
|
||||
|
||||
Reference in New Issue
Block a user