From d65d790a20aba28beb21f6cf21a0ac0e19d2041f Mon Sep 17 00:00:00 2001 From: sh <37271604+shumvgolove@users.noreply.github.com> Date: Fri, 31 Jul 2026 18:19:43 +0400 Subject: [PATCH 1/2] ntf server: shard push workers per notification token (#1840) * ntf server: shard push workers per notification token * ntf server: inline push worker shard calculation --- src/Simplex/Messaging/Notifications/Server.hs | 13 ++++++++----- src/Simplex/Messaging/Notifications/Server/Env.hs | 2 +- 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/src/Simplex/Messaging/Notifications/Server.hs b/src/Simplex/Messaging/Notifications/Server.hs index f04ca4e35..edcebb96f 100644 --- a/src/Simplex/Messaging/Notifications/Server.hs +++ b/src/Simplex/Messaging/Notifications/Server.hs @@ -34,6 +34,7 @@ import Data.ByteString.Char8 (ByteString) import qualified Data.ByteString.Char8 as B import Data.Either (partitionEithers) import Data.Functor (($>)) +import Data.Hashable (hash) import Data.IORef import Data.Int (Int64) import qualified Data.IntSet as IS @@ -640,11 +641,13 @@ showServer' :: SMPServer -> Text showServer' = decodeLatin1 . strEncode . host pushNotification :: NtfPushServer -> Maybe T.Text -> OwnServer -> NtfTknRec -> PushNotification -> M () -pushNotification s srvHost_ isOwn tkn@NtfTknRec {token = token@(DeviceToken pp _)} ntf = +pushNotification s srvHost_ isOwn tkn@NtfTknRec {ntfTknId, token = token@(DeviceToken pp _)} ntf = ifM (pushProviderAllowed token) - (getOrCreatePushWorker s (srvHost_, pp) isOwn >>= atomically . (`writeTBQueue` (tkn, ntf))) + (getOrCreatePushWorker s (srvHost_, pp, hash (unEntityId ntfTknId) `mod` pushWorkersPerServer) isOwn >>= atomically . (`writeTBQueue` (tkn, ntf))) (logWarn "skipping disabled APNS test push provider") + where + pushWorkersPerServer = 8 pushProviderAllowed :: DeviceToken -> M Bool pushProviderAllowed (DeviceToken PPApnsTest _) = asks (allowTestPushProvider . config) @@ -657,8 +660,8 @@ guardPushProvider token action = action (pure $ NRErr $ CMD SMP.PROHIBITED) -getOrCreatePushWorker :: NtfPushServer -> (Maybe T.Text, PushProvider) -> OwnServer -> M (TBQueue (NtfTknRec, PushNotification)) -getOrCreatePushWorker s@NtfPushServer {pushWorkers, pushWorkerSeq, pushQSize} key@(srvHost_, _) isOwn = do +getOrCreatePushWorker :: NtfPushServer -> (Maybe T.Text, PushProvider, Int) -> OwnServer -> M (TBQueue (NtfTknRec, PushNotification)) +getOrCreatePushWorker s@NtfPushServer {pushWorkers, pushWorkerSeq, pushQSize} key@(srvHost_, _, _) isOwn = do ts <- liftIO getCurrentTime withGetSessVar' pushWorkerSeq key pushWorkers ts createWorker existingWorker where @@ -731,7 +734,7 @@ runPushWorker s srvHost_ isOwn q = forever $ do _ -> err e err e = logError ("Push provider error (" <> tshow pp <> ", " <> tshow ntfTknId <> "): " <> tshow e) $> Left e -pushWorkersQLength :: TMap (Maybe T.Text, PushProvider) PushWorkerVar -> IO Natural +pushWorkersQLength :: TMap (Maybe T.Text, PushProvider, Int) PushWorkerVar -> IO Natural pushWorkersQLength workers = do ws <- readTVarIO workers foldM addQLength 0 ws diff --git a/src/Simplex/Messaging/Notifications/Server/Env.hs b/src/Simplex/Messaging/Notifications/Server/Env.hs index 6f9416db4..d7f772a87 100644 --- a/src/Simplex/Messaging/Notifications/Server/Env.hs +++ b/src/Simplex/Messaging/Notifications/Server/Env.hs @@ -174,7 +174,7 @@ data SMPSubscriber = SMPSubscriber } data NtfPushServer = NtfPushServer - { pushWorkers :: TMap (Maybe T.Text, PushProvider) PushWorkerVar, + { pushWorkers :: TMap (Maybe T.Text, PushProvider, Int) PushWorkerVar, -- Int is the worker shard pushWorkerSeq :: TVar Int, pushQSize :: Natural, pushClients :: TMap PushProvider PushClientVar, From 27a37387be98d9c7ec0e62373e125539675d0095 Mon Sep 17 00:00:00 2001 From: Evgeny Poberezkin Date: Fri, 31 Jul 2026 15:20:41 +0100 Subject: [PATCH 2/2] 7.0.1.0 --- simplexmq.cabal | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/simplexmq.cabal b/simplexmq.cabal index 694ee40b9..ea0f1e2df 100644 --- a/simplexmq.cabal +++ b/simplexmq.cabal @@ -1,7 +1,7 @@ cabal-version: 3.0 name: simplexmq -version: 7.0.0.6 +version: 7.0.1.0 synopsis: SimpleXMQ message broker description: This package includes <./docs/Simplex-Messaging-Server.html server>, <./docs/Simplex-Messaging-Client.html client> and