mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-10 02:47:31 +00:00
smp-server: share subscription key between maps
This commit is contained in:
@@ -277,17 +277,17 @@ smpServer started cfg@ServerConfig {transports, transportConfig = tCfg, startOpt
|
||||
Nothing -> updateSubDisconnected
|
||||
where
|
||||
updateSubConnected c = case clntSub of
|
||||
CSClient qId prevServiceId serviceId_ -> do
|
||||
CSClient qId k prevServiceId serviceId_ -> do
|
||||
modifyTVar' subClients $ IS.insert clntId -- add ID to server's subscribed cients
|
||||
as'' <- if prevServiceId == serviceId_ then pure [] else endServiceSub prevServiceId qId END
|
||||
case serviceId_ of
|
||||
Just serviceId -> do
|
||||
modifyTVar' totalServiceSubs $ addServiceSubs (1, queueIdHash qId) -- server count and IDs hash for all services
|
||||
as <- endQueueSub qId END
|
||||
as' <- cancelServiceSubs serviceId =<< upsertSubscribedClient serviceId c serviceSubscribers
|
||||
as' <- cancelServiceSubs serviceId =<< upsertSubscribedClient (subKey serviceId) c serviceSubscribers
|
||||
pure $ as ++ as' ++ as''
|
||||
Nothing -> do
|
||||
as <- prevSub qId END (CSAEndSub qId) =<< upsertSubscribedClient qId c queueSubscribers
|
||||
as <- prevSub qId END (CSAEndSub qId) =<< upsertSubscribedClient k c queueSubscribers
|
||||
pure $ as ++ as''
|
||||
CSDeleted qId serviceId -> do
|
||||
removeWhenNoSubs c
|
||||
@@ -297,10 +297,10 @@ smpServer started cfg@ServerConfig {transports, transportConfig = tCfg, startOpt
|
||||
CSService serviceId changedSubs -> do
|
||||
modifyTVar' subClients $ IS.insert clntId -- add ID to server's subscribed cients
|
||||
modifyTVar' totalServiceSubs $ addServiceSubs changedSubs -- server count and IDs hash for all services
|
||||
cancelServiceSubs serviceId =<< upsertSubscribedClient serviceId c serviceSubscribers
|
||||
cancelServiceSubs serviceId =<< upsertSubscribedClient (subKey serviceId) c serviceSubscribers
|
||||
updateSubDisconnected = case clntSub of
|
||||
-- do not insert client if it is already disconnected, but send END/DELD to any other client subscribed to this queue or service
|
||||
CSClient qId prevServiceId serviceId -> do
|
||||
CSClient qId _ prevServiceId serviceId -> do
|
||||
as <- endQueueSub qId END
|
||||
as' <- endServiceSub serviceId qId END
|
||||
as'' <- if prevServiceId == serviceId then pure [] else endServiceSub prevServiceId qId END
|
||||
@@ -1720,12 +1720,13 @@ client
|
||||
|
||||
subscribeNewQueue :: RecipientId -> QueueRec -> M s ()
|
||||
subscribeNewQueue rId QueueRec {rcvServiceId} = do
|
||||
let k = subKey rId
|
||||
case rcvServiceId of
|
||||
Just _ -> atomically $ modifyTVar' (serviceSubsCount clnt) $ addServiceSubs (1, queueIdHash rId)
|
||||
Nothing -> do
|
||||
sub <- atomically $ newSubscription NoSub
|
||||
atomically $ TM.insert (subKey rId) sub $ subscriptions clnt
|
||||
atomically $ writeTQueue (subQ subscribers) (CSClient rId rcvServiceId rcvServiceId, clientId)
|
||||
atomically $ TM.insert k sub $ subscriptions clnt
|
||||
atomically $ writeTQueue (subQ subscribers) (CSClient rId k rcvServiceId rcvServiceId, clientId)
|
||||
|
||||
-- clients that use GET are not added to server subscribers
|
||||
getMessage :: StoreQueue s -> QueueRec -> M s (Transmission BrokerMsg)
|
||||
@@ -1804,7 +1805,8 @@ client
|
||||
sharedSubscribeQueue q queueServiceId srvSubscribers clientSubs clientServiceSubs mkSub servicesSel = do
|
||||
stats <- asks serverStats
|
||||
let incSrvStat sel = incStat $ sel $ servicesSel stats
|
||||
writeSub = writeTQueue (subQ srvSubscribers) (CSClient entId queueServiceId clntServiceId, clientId)
|
||||
k = subKey entId
|
||||
writeSub = writeTQueue (subQ srvSubscribers) (CSClient entId k queueServiceId clntServiceId, clientId)
|
||||
liftIO $ case clntServiceId of
|
||||
Just serviceId
|
||||
| queueServiceId == Just serviceId -> do
|
||||
@@ -1842,12 +1844,12 @@ client
|
||||
unless hasSub $ atomically writeSub
|
||||
pure r
|
||||
where
|
||||
getSubscription = TM.lookup (subKey entId) $ clientSubs clnt
|
||||
getSubscription = TM.lookup k $ clientSubs clnt
|
||||
newSub = \case
|
||||
Just sub -> pure (True, Just sub)
|
||||
Nothing -> do
|
||||
sub <- mkSub
|
||||
TM.insert (subKey entId) sub $ clientSubs clnt
|
||||
TM.insert k sub $ clientSubs clnt
|
||||
pure (False, Just sub)
|
||||
|
||||
subscribeServiceMessages :: ServiceId -> (Int64, IdsHash) -> M s BrokerMsg
|
||||
|
||||
@@ -410,16 +410,14 @@ getSubscribedClient entId (SubscribedClients cs) = TM.lookupIO (subKey entId) cs
|
||||
{-# INLINE getSubscribedClient #-}
|
||||
|
||||
-- insert subscribed and current client, return previously subscribed client if it is different
|
||||
upsertSubscribedClient :: EntityId -> Client s -> SubscribedClients s -> STM (Maybe (Client s))
|
||||
upsertSubscribedClient entId c (SubscribedClients cs) =
|
||||
upsertSubscribedClient :: SubKey -> Client s -> SubscribedClients s -> STM (Maybe (Client s))
|
||||
upsertSubscribedClient k c (SubscribedClients cs) =
|
||||
TM.lookup k cs >>= \case
|
||||
Nothing -> Nothing <$ TM.insertM k (newTVar (Just c)) cs
|
||||
Just cv ->
|
||||
readTVar cv >>= \case
|
||||
Just c' | sameClientId c c' -> pure Nothing
|
||||
c_ -> c_ <$ writeTVar cv (Just c)
|
||||
where
|
||||
k = subKey entId
|
||||
|
||||
lookupSubscribedClient :: EntityId -> SubscribedClients s -> STM (Maybe (Client s))
|
||||
lookupSubscribedClient entId (SubscribedClients cs) = TM.lookup (subKey entId) cs $>>= readTVar
|
||||
@@ -451,7 +449,7 @@ sameClient c cv = maybe False (sameClientId c) <$> readTVar cv
|
||||
{-# INLINE sameClient #-}
|
||||
|
||||
data ClientSub
|
||||
= CSClient QueueId (Maybe ServiceId) (Maybe ServiceId) -- includes previous and new associated service IDs
|
||||
= CSClient QueueId SubKey (Maybe ServiceId) (Maybe ServiceId) -- includes the key in client subscriptions, previous and new associated service IDs
|
||||
| CSDeleted QueueId (Maybe ServiceId) -- includes previously associated service IDs
|
||||
| CSService ServiceId (Int64, IdsHash) -- only send END to idividual client subs on message delivery, not of SSUB/NSSUB
|
||||
|
||||
|
||||
Reference in New Issue
Block a user