Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 65 additions & 0 deletions plans/2026-09-26-index-active-sub-conns.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
## Root cause: every subscription batch rebuilds the set of subscribed connections

Resubscribing a session allocates and discards a fresh `Set ConnId` covering *all* of that
session's active subscriptions, once per batch. On clients with many queues per session this
dominates both the CPU and the allocation of a reconnection, and it happens again on every
network change.

`subscribeSessQueues_` (`Agent/Client.hs`) needs to know which connections were already
subscribed, so that only newly subscribed ones are reported `UP`. It obtained that by folding
the session's whole subscription map:

```haskell
Just . S.fromList . map qConnId . M.elems <$> atomically (SS.getActiveSubs tSess $ currentSubs c)
```

`activeSubs` is keyed by `RecipientId` and holds every subscribed queue of the session, so this
walks all of them and builds a set of the same size — to answer a membership question about at
most `subsBatchSize` connections.

### Cost

`subscribeQueues` chunks the queues into batches of `subsBatchSize` (1350 by default) and calls
`subscribeSessQueues_` per batch. `activeSubs` grows as the resubscription proceeds, so batch *k*
folds roughly *k* × 1350 entries: resubscribing *N* queues in one session costs on the order of
*N²* / 2700 entry traversals, and allocates a `Set` of up to *N* elements per batch.

For a session holding 50k queues that is ~0.9M traversals per resubscription and ~37 sets of up
to 50k elements — tens of MB of short-lived allocation, repeated for each session. Under the
non-moving collector this allocation is also promotion pressure, which is what makes a
reconnection storm visible as resident memory rather than only as CPU.

Sessions are keyed `(userId, server, Nothing)` in `TSMSession` mode, so a client with many
connections concentrates its queues into few sessions and hits the worst case rather than
avoiding it.

## Fix

Maintain the connection index incrementally instead of deriving it per batch.

`SessSubs` gains `activeConns :: TMap ConnId Int` alongside `activeSubs`, and `getActiveConns`
reads it directly. `subscribeSessQueues_` then tests membership against that map.

A connection can hold more than one subscribed queue, so the index counts them rather than
storing a set: `incActiveConn` on a queue becoming active, `decActiveConn` when it stops, and the
entry is removed when its count reaches zero. That keeps "connection has at least one active
subscription" exact under partial subscription and partial failure.

It is maintained at every site that mutates `activeSubs`, and nowhere else:

| site | change |
|---|---|
| `addActiveSub'` | increments, but only when the queue was not already active |
| `batchAddActiveSubs` | increments for the queues actually added (`M.difference` against the previous map) |
| `deleteSub` | decrements for the queue if it was active |
| `batchDeleteSubs` | decrements for the removed queues that were active (`M.restrictKeys`) |
| `setSubsPending_` | clears the index with `activeSubs`, since all subscriptions become pending |

The guards matter: re-adding an already-active queue, or deleting one that was only pending,
must not move the count, otherwise the index drifts from `activeSubs` and connections are
reported `UP` twice or not at all.

## Behaviour

Unchanged. The same connections are reported `UP`, and the session-closing condition that
previously tested `S.null cs` now tests `M.null cs` over the same information.
6 changes: 3 additions & 3 deletions src/Simplex/Messaging/Agent/Client.hs
Original file line number Diff line number Diff line change
Expand Up @@ -1716,7 +1716,7 @@ subscribeSessQueues_ c withEvents qs = sendClientBatch_ "SUB" False subscribe_ c
rs <- sendBatch (\smp' _ -> subscribeSMPQueues smp') smp NRMBackground qs'
cs_ <-
if withEvents
then Just . S.fromList . map qConnId . M.elems <$> atomically (SS.getActiveSubs tSess $ currentSubs c)
then Just <$> atomically (SS.getActiveConns tSess $ currentSubs c)
else pure Nothing
active <- E.uninterruptibleMask_ $ do
(active, (serviceQs, notices)) <- atomically $ do
Expand All @@ -1734,13 +1734,13 @@ subscribeSessQueues_ c withEvents qs = sendClientBatch_ "SUB" False subscribe_ c
pure active
forM_ cs_ $ \cs -> do
let (errs, okConns) = partitionEithers $ map (\(RcvQueueSub {connId}, r) -> bimap (connId,) (const connId) r) $ L.toList rs
conns = filter (`S.notMember` cs) okConns
conns = filter (`M.notMember` cs) okConns
unless (null conns) $ notifySub c $ UP srv conns
forM_ (L.nonEmpty errs) $ \errs' -> do
let noFinalErrs = all (temporaryClientError . snd) errs'
addr = B.unpack $ strEncode srv
notifySub c $ ERRS $ L.map (second $ protocolClientError SMP addr) errs'
when (null okConns && S.null cs && noFinalErrs && active) $ liftIO $ do
when (null okConns && M.null cs && noFinalErrs && active) $ liftIO $ do
-- We only close the client session that was used to subscribe.
v_ <- atomically $ ifM (activeClientSession c tSess sessId) (TM.lookupDelete tSess $ smpClients c) (pure Nothing)
mapM_ (closeClient_ c) v_
Expand Down
46 changes: 41 additions & 5 deletions src/Simplex/Messaging/Agent/TSessionSubs.hs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ module Simplex.Messaging.Agent.TSessionSubs
getPendingSubs,
getPendingQueueSubs,
getActiveSubs,
getActiveConns,
setSubsPending,
updateClientNotices,
foldSessionSubs,
Expand All @@ -43,7 +44,7 @@ import Data.Map.Strict (Map)
import qualified Data.Map.Strict as M
import Data.Maybe (fromMaybe, isJust)
import qualified Data.Set as S
import Simplex.Messaging.Agent.Protocol (SMPQueue (..))
import Simplex.Messaging.Agent.Protocol (ConnId, SMPQueue (..))
import Simplex.Messaging.Agent.Store (RcvQueue, RcvQueueSub (..), ServiceAssoc, SomeRcvQueue, StoredRcvQueue (rcvServiceAssoc), rcvQueueSub)
import Simplex.Messaging.Client (SMPTransportSession, TransportSessionMode (..))
import Simplex.Messaging.Protocol (IdsHash, RecipientId, ServiceSub (..), queueIdHash)
Expand All @@ -59,6 +60,10 @@ data TSessionSubs = TSessionSubs
data SessSubs = SessSubs
{ subsSessId :: TVar (Maybe SessionId),
activeSubs :: TMap RecipientId RcvQueueSub,
-- Connections with active subscriptions, with the number of their active subscriptions.
-- Kept in sync with activeSubs, so that "is this connection already subscribed" does not
-- require folding activeSubs, which holds every subscribed queue of the session.
activeConns :: TMap ConnId Int,
pendingSubs :: TMap RecipientId RcvQueueSub,
activeServiceSub :: TVar (Maybe ServiceSub),
pendingServiceSub :: TVar (Maybe ServiceSub)
Expand All @@ -80,10 +85,21 @@ getSessSubs :: SMPTransportSession -> TSessionSubs -> STM SessSubs
getSessSubs tSess ss = lookupSubs tSess ss >>= maybe new pure
where
new = do
s <- SessSubs <$> newTVar Nothing <*> newTVar M.empty <*> newTVar M.empty <*> newTVar Nothing <*> newTVar Nothing
s <- SessSubs <$> newTVar Nothing <*> newTVar M.empty <*> newTVar M.empty <*> newTVar M.empty <*> newTVar Nothing <*> newTVar Nothing
TM.insert tSess s $ sessionSubs ss
pure s

incActiveConn :: SessSubs -> ConnId -> STM ()
incActiveConn s cId = TM.alter inc cId $ activeConns s
where
inc = Just . maybe 1 (+ 1)

decActiveConn :: SessSubs -> ConnId -> STM ()
decActiveConn s cId = TM.alter dec cId $ activeConns s
where
dec (Just n) | n > 1 = Just (n - 1)
dec _ = Nothing

hasActiveSub :: SMPTransportSession -> RecipientId -> TSessionSubs -> STM Bool
hasActiveSub = hasQueue_ activeSubs
{-# INLINE hasActiveSub #-}
Expand Down Expand Up @@ -135,7 +151,10 @@ addActiveSub' tSess sessId serviceId_ rq serviceAssoc ss = do
TM.delete rId $ pendingSubs s
case serviceId_ of
Just serviceId | serviceAssoc -> updateActiveService s serviceId 1 (queueIdHash rId)
_ -> TM.insert rId rq $ activeSubs s
_ -> do
prev <- TM.lookup rId $ activeSubs s
TM.insert rId rq $ activeSubs s
unless (isJust prev) $ incActiveConn s (connId rq)
else TM.insert rId rq $ pendingSubs s

batchAddActiveSubs :: SMPTransportSession -> SessionId -> Maybe ServiceId -> ([RcvQueueSub], [RcvQueueSub]) -> TSessionSubs -> STM ()
Expand All @@ -146,7 +165,9 @@ batchAddActiveSubs tSess sessId serviceId_ (rqs, serviceRQs) ss = do
serviceQs = queuesMap serviceRQs
if Just sessId == sessId'
then do
prev <- readTVar $ activeSubs s
TM.union qs $ activeSubs s
mapM_ (incActiveConn s . connId) $ M.elems $ qs `M.difference` prev
modifyTVar' (pendingSubs s) (`M.difference` qs)
unless (null serviceRQs) $ forM_ serviceId_ $ \serviceId -> do
modifyTVar' (pendingSubs s) (`M.difference` serviceQs)
Expand Down Expand Up @@ -178,13 +199,24 @@ batchDeletePendingSubs tSess rIds = lookupSubs tSess >=> mapM_ (delete . pending
delete = (`modifyTVar'` (`M.withoutKeys` rIds))

deleteSub :: SMPTransportSession -> RecipientId -> TSessionSubs -> STM ()
deleteSub tSess rId = lookupSubs tSess >=> mapM_ (\s -> TM.delete rId (activeSubs s) >> TM.delete rId (pendingSubs s))
deleteSub tSess rId = lookupSubs tSess >=> mapM_ delete_
where
delete_ s = do
prev <- TM.lookup rId $ activeSubs s
TM.delete rId (activeSubs s)
TM.delete rId (pendingSubs s)
mapM_ (decActiveConn s . connId) prev

batchDeleteSubs :: SomeRcvQueue q => SMPTransportSession -> [q] -> TSessionSubs -> STM ()
batchDeleteSubs tSess rqs = lookupSubs tSess >=> mapM_ (\s -> delete (activeSubs s) >> delete (pendingSubs s))
batchDeleteSubs tSess rqs = lookupSubs tSess >=> mapM_ delete_
where
rIds = S.fromList $ map queueId rqs
delete = (`modifyTVar'` (`M.withoutKeys` rIds))
delete_ s = do
prev <- readTVar $ activeSubs s
delete (activeSubs s)
delete (pendingSubs s)
mapM_ (decActiveConn s . connId) $ M.elems $ M.restrictKeys prev rIds

deleteServiceSub :: SMPTransportSession -> TSessionSubs -> STM ()
deleteServiceSub tSess = lookupSubs tSess >=> mapM_ (\s -> writeTVar (activeServiceSub s) Nothing >> writeTVar (pendingServiceSub s) Nothing)
Expand All @@ -204,6 +236,9 @@ getPendingQueueSubs :: SMPTransportSession -> TSessionSubs -> STM (Map Recipient
getPendingQueueSubs = getSubs_ pendingSubs
{-# INLINE getPendingQueueSubs #-}

getActiveConns :: SMPTransportSession -> TSessionSubs -> STM (Map ConnId Int)
getActiveConns tSess = lookupSubs tSess >=> maybe (pure M.empty) (readTVar . activeConns)

getActiveSubs :: SMPTransportSession -> TSessionSubs -> STM (Map RecipientId RcvQueueSub)
getActiveSubs = getSubs_ activeSubs
{-# INLINE getActiveSubs #-}
Expand Down Expand Up @@ -238,6 +273,7 @@ setSubsPending_ s sessId_ = do
subs <- readTVar as
unless (null subs) $ do
writeTVar as M.empty
TM.clear $ activeConns s
modifyTVar' (pendingSubs s) $ M.union subs
(subs,) <$> setServiceSubPending_ s

Expand Down
Loading