ntf server: optimize in-memory storage (#1516)

* ntf server: optimize in-memory storage

* test

* ntf server: fix store log parser for token status
This commit is contained in:
Evgeny
2025-04-21 17:12:16 +01:00
committed by GitHub
parent 1e29f7c811
commit afb338a41a
4 changed files with 79 additions and 51 deletions
@@ -8,6 +8,7 @@
{-# LANGUAGE OverloadedLists #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE StandaloneDeriving #-}
module Simplex.Messaging.Notifications.Server.Store where
@@ -16,15 +17,16 @@ import Control.Monad
import Data.ByteString.Char8 (ByteString)
import Data.Functor (($>))
import Data.List.NonEmpty (NonEmpty (..), (<|))
import Data.Map.Strict (Map)
import qualified Data.Map.Strict as M
import Data.Maybe (catMaybes)
import Data.Maybe (isNothing)
import Data.Set (Set)
import qualified Data.Set as S
import Data.Word (Word16)
import qualified Simplex.Messaging.Crypto as C
import Simplex.Messaging.Encoding.String
import Simplex.Messaging.Notifications.Protocol
import Simplex.Messaging.Protocol (NtfPrivateAuthKey, NtfPublicAuthKey, SMPServer)
import Simplex.Messaging.Protocol (NotifierId, NtfPrivateAuthKey, NtfPublicAuthKey, SMPServer)
import Simplex.Messaging.Server.QueueStore (RoundedSystemTime)
import Simplex.Messaging.TMap (TMap)
import qualified Simplex.Messaging.TMap as TM
@@ -35,8 +37,11 @@ data NtfStore = NtfStore
-- multiple registrations exist to protect from malicious registrations if token is compromised
tokenRegistrations :: TMap DeviceToken (TMap ByteString NtfTokenId),
subscriptions :: TMap NtfSubscriptionId NtfSubData,
tokenSubscriptions :: TMap NtfTokenId (TVar (Set NtfSubscriptionId)),
subscriptionLookup :: TMap SMPQueueNtf NtfSubscriptionId,
-- the first set is used to delete from `subscriptions` when token is deleted, the second - to cancel SMP subsriptions.
-- TODO [notifications] it can be simplified once NtfSubData is fully removed.
tokenSubscriptions :: TMap NtfTokenId (TMap SMPServer (TVar (Set NtfSubscriptionId), TVar (Set NotifierId))),
-- TODO [notifications] for subscriptions that "migrated" to server subscription, we may replace NtfSubData with NtfTokenId here (Either NtfSubData NtfTokenId).
subscriptionLookup :: TMap SMPServer (TMap NotifierId NtfSubData),
tokenLastNtfs :: TMap NtfTokenId (TVar (NonEmpty PNMessageData))
}
@@ -134,7 +139,7 @@ removeTokenRegistration st NtfTknData {ntfTknId = tId, token, tknVerifyKey} =
>>= mapM_ (\tId' -> when (tId == tId') $ TM.delete k regs)
k = C.toPubKey C.pubKeyBytes tknVerifyKey
deleteNtfToken :: NtfStore -> NtfTokenId -> STM [SMPQueueNtf]
deleteNtfToken :: NtfStore -> NtfTokenId -> STM (Map SMPServer (Set NotifierId))
deleteNtfToken st tknId = do
void $
TM.lookupDelete tknId (tokens st) $>>= \NtfTknData {token, tknVerifyKey} ->
@@ -147,25 +152,25 @@ deleteNtfToken st tknId = do
regs = tokenRegistrations st
regKey = C.toPubKey C.pubKeyBytes
deleteTokenSubs :: NtfStore -> NtfTokenId -> STM [SMPQueueNtf]
deleteTokenSubs st tknId = do
qs <-
TM.lookupDelete tknId (tokenSubscriptions st)
>>= mapM (readTVar >=> mapM deleteSub . S.toList)
pure $ maybe [] catMaybes qs
deleteTokenSubs :: NtfStore -> NtfTokenId -> STM (Map SMPServer (Set NotifierId))
deleteTokenSubs st tknId =
TM.lookupDelete tknId (tokenSubscriptions st)
>>= maybe (pure M.empty) (readTVar >=> deleteSrvSubs)
where
deleteSub subId = do
TM.lookupDelete subId (subscriptions st)
$>>= \NtfSubData {smpQueue} ->
TM.delete smpQueue (subscriptionLookup st) $> Just smpQueue
deleteSrvSubs :: Map SMPServer (TVar (Set NtfSubscriptionId), TVar (Set NotifierId)) -> STM (Map SMPServer (Set NotifierId))
deleteSrvSubs = M.traverseWithKey $ \smpServer (sVar, nVar) -> do
sIds <- readTVar sVar
modifyTVar' (subscriptions st) (`M.withoutKeys` sIds)
nIds <- readTVar nVar
TM.lookup smpServer (subscriptionLookup st) >>= mapM_ (`modifyTVar'` (`M.withoutKeys` nIds))
pure nIds
getNtfSubscriptionIO :: NtfStore -> NtfSubscriptionId -> IO (Maybe NtfSubData)
getNtfSubscriptionIO st subId = TM.lookupIO subId (subscriptions st)
findNtfSubscription :: NtfStore -> SMPQueueNtf -> STM (Maybe NtfSubData)
findNtfSubscription st smpQueue = do
TM.lookup smpQueue (subscriptionLookup st)
$>>= \subId -> TM.lookup subId (subscriptions st)
findNtfSubscription st SMPQueueNtf {smpServer, notifierId} =
TM.lookup smpServer (subscriptionLookup st) $>>= TM.lookup notifierId
findNtfSubscriptionToken :: NtfStore -> SMPQueueNtf -> STM (Maybe NtfTknData)
findNtfSubscriptionToken st smpQueue = do
@@ -183,30 +188,45 @@ mkNtfSubData ntfSubId (NewNtfSub tokenId smpQueue notifierKey) = do
subStatus <- newTVar NSNew
pure NtfSubData {ntfSubId, smpQueue, tokenId, subStatus, notifierKey}
addNtfSubscription :: NtfStore -> NtfSubscriptionId -> NtfSubData -> STM (Maybe ())
addNtfSubscription st subId sub@NtfSubData {smpQueue, tokenId} =
TM.lookup tokenId (tokenSubscriptions st) >>= maybe newTokenSub pure >>= insertSub
-- returns False if subscription existed before
addNtfSubscription :: NtfStore -> NtfSubscriptionId -> NtfSubData -> STM Bool
addNtfSubscription st subId sub@NtfSubData {smpQueue = SMPQueueNtf {smpServer, notifierId}, tokenId} =
TM.lookup tokenId (tokenSubscriptions st)
>>= maybe newTokenSubs pure
>>= \ts -> TM.lookup smpServer ts
>>= maybe (newTokenSrvSubs ts) pure
>>= insertSub
where
newTokenSub = do
ts <- newTVar S.empty
newTokenSubs = do
ts <- newTVar M.empty
TM.insert tokenId ts $ tokenSubscriptions st
pure ts
insertSub ts = do
modifyTVar' ts $ S.insert subId
newTokenSrvSubs ts = do
tss <- (,) <$> newTVar S.empty <*> newTVar S.empty
TM.insert smpServer tss ts
pure tss
insertSub :: (TVar (Set NtfSubscriptionId), TVar (Set NotifierId)) -> STM Bool
insertSub (sIds, nIds) = do
modifyTVar' sIds $ S.insert subId
modifyTVar' nIds $ S.insert notifierId
TM.insert subId sub $ subscriptions st
TM.insert smpQueue subId (subscriptionLookup st)
-- return Nothing if subscription existed before
pure $ Just ()
TM.lookup smpServer (subscriptionLookup st)
>>= maybe newSubs pure
>>= fmap isNothing . TM.lookupInsert notifierId sub
newSubs = do
ss <- newTVar M.empty
TM.insert smpServer ss $ subscriptionLookup st
pure ss
deleteNtfSubscription :: NtfStore -> NtfSubscriptionId -> STM ()
deleteNtfSubscription st subId = do
TM.lookupDelete subId (subscriptions st)
>>= mapM_
( \NtfSubData {smpQueue, tokenId} -> do
TM.delete smpQueue $ subscriptionLookup st
ts_ <- TM.lookup tokenId (tokenSubscriptions st)
forM_ ts_ $ \ts -> modifyTVar' ts $ S.delete subId
)
deleteNtfSubscription st subId = TM.lookupDelete subId (subscriptions st) >>= mapM_ deleteSubIndices
where
deleteSubIndices NtfSubData {smpQueue = SMPQueueNtf {smpServer, notifierId}, tokenId} = do
TM.lookup smpServer (subscriptionLookup st) >>= mapM_ (TM.delete notifierId)
tss_ <- TM.lookup tokenId (tokenSubscriptions st) $>>= TM.lookup smpServer
forM_ tss_ $ \(sIds, nIds) -> do
modifyTVar' sIds $ S.delete subId
modifyTVar' nIds $ S.delete notifierId
addTokenLastNtf :: NtfStore -> NtfTokenId -> PNMessageData -> IO (NonEmpty PNMessageData)
addTokenLastNtf st tknId newNtf =