From 191b5128cc07620cc95a943b0e5062987e18d88d Mon Sep 17 00:00:00 2001 From: "Evgeny @ SimpleX Chat" <259188159+evgeny-simplex@users.noreply.github.com> Date: Mon, 17 Aug 2026 22:14:16 +0000 Subject: [PATCH] refactor --- src/Simplex/Messaging/Agent.hs | 31 ++++++++++++++-------------- src/Simplex/Messaging/Agent/Store.hs | 4 ++-- 2 files changed, 17 insertions(+), 18 deletions(-) diff --git a/src/Simplex/Messaging/Agent.hs b/src/Simplex/Messaging/Agent.hs index c64adff8c..63bf61055 100644 --- a/src/Simplex/Messaging/Agent.hs +++ b/src/Simplex/Messaging/Agent.hs @@ -2359,8 +2359,11 @@ enqueueMessagesB c reqs = do enqueueSavedMessageB c $ mapMaybe snd $ rights $ toList reqs' pure $ fst <$$> reqs' -isActiveSndQ :: SndQueue -> Bool -isActiveSndQ SndQueue {status} = status == Secured || status == Active +-- which tail send queues a fanned-out message is delivered to: active queues, plus a securing replacement (v8), minus a terminating one +isActiveSndQ :: ConnData -> SndQueue -> Bool +isActiveSndQ ConnData {connAgentVersion} sq@SndQueue {status, sndSwchStatus} = + sndSwchStatus /= Just SSSendingQEND + && (status == Secured || status == Active || (connAgentVersion >= rpcAddressSMPAgentVersion && securingSndQueue sq)) {-# INLINE isActiveSndQ #-} enqueueMessage :: AgentClient -> ConnData -> SndQueue -> MsgFlags -> AMessage -> AM (AgentMsgId, PQEncryption) @@ -2376,13 +2379,9 @@ enqueueMessageB c reqs = do mapAccumLM (\ids r -> storeSentMsg db cfg ids r `E.catchAny` \e -> (ids,) <$> handleInternal e) IM.empty reqs forME reqMids $ \((csqs_, _, _, _), InternalId msgId, pqSecr) -> forM csqs_ $ \(cData, sq :| sqs) -> do submitPendingMsg c sq - let sqs' = filter (sndDeliverTo cData) sqs + let sqs' = filter (isActiveSndQ cData) sqs pure ((msgId, pqSecr), if null sqs' then Nothing else Just (sqs', msgId)) where - sndDeliverTo :: ConnData -> SndQueue -> Bool - sndDeliverTo ConnData {connAgentVersion} sq@SndQueue {sndSwchStatus} = - sndSwchStatus /= Just SSSendingQEND - && (isActiveSndQ sq || (connAgentVersion >= rpcAddressSMPAgentVersion && securingSndQueue sq)) storeSentMsg :: DB.Connection -> AgentConfig -> @@ -2789,9 +2788,9 @@ abortConnectionSwitch' :: AgentClient -> ConnId -> AM ConnectionStats abortConnectionSwitch' c connId = withConnLock c connId "abortConnectionSwitch" $ withStore c (`getConn` connId) >>= \case - SomeConn _ (DuplexConnection cData@ConnData {connAgentVersion} rqs sqs) -> case switchingRQ rqs of + SomeConn _ (DuplexConnection cData rqs sqs) -> case switchingRQ rqs of Just rq - | canAbortRcvSwitch connAgentVersion rq -> do + | canAbortRcvSwitch cData rq -> do when (ratchetSyncSendProhibited cData) $ throwE $ CMD PROHIBITED "abortConnectionSwitch: send prohibited" -- multiple queues to which the connections switches were possible when repeating switch was allowed let (delRqs, keepRqs) = L.partition ((Just (dbQId rq) ==) . dbReplaceQId) rqs @@ -2820,7 +2819,7 @@ synchronizeRatchet' c connId pqSupport' force = withConnLock c connId "synchroni AgentConfig {e2eEncryptVRange} <- asks config g <- asks random (pks, e2eParams) <- liftIO $ CR.generateRcvE2EParams g (maxVersion e2eEncryptVRange) pqSupport' - enqueueRatchetKeyMsgs c sqs e2eParams + enqueueRatchetKeyMsgs c cData' sqs e2eParams withStore' c $ \db -> do setConnRatchetSync db connId RSStarted setRatchetX3dhKeys db connId pks @@ -3035,9 +3034,9 @@ connectionStats c = \case subStatus = Nothing } rcvQueueInfo :: ConnData -> RcvQueue -> AM RcvQueueInfo - rcvQueueInfo ConnData {connAgentVersion} rq@RcvQueue {server, status, rcvSwchStatus} = do + rcvQueueInfo cData rq@RcvQueue {server, status, rcvSwchStatus} = do subStatus <- atomically checkQueueSubStatus - pure $ RcvQueueInfo {rcvServer = server, status, rcvSwitchStatus = rcvSwchStatus, canAbortSwitch = canAbortRcvSwitch connAgentVersion rq, subStatus} + pure $ RcvQueueInfo {rcvServer = server, status, rcvSwitchStatus = rcvSwchStatus, canAbortSwitch = canAbortRcvSwitch cData rq, subStatus} where checkQueueSubStatus :: STM SubscriptionStatus checkQueueSubStatus = @@ -4129,7 +4128,7 @@ processSMPTransmissions c@AgentClient {subQ} (tSess@(userId, srv, _), THandlePar sendReplyKey = do g <- asks random (pks, e2eParams) <- liftIO $ CR.generateRcvE2EParams g e2eVersion pqSupport - enqueueRatchetKeyMsgs c sqs e2eParams + enqueueRatchetKeyMsgs c cData' sqs e2eParams pure pks notifyRatchetSyncError = do let cData'' = cData' {ratchetSyncState = RSRequired} :: ConnData @@ -4271,10 +4270,10 @@ storeConfirmation c cData@ConnData {connId, pqSupport, connAgentVersion = v} sq liftIO $ createSndMsg db connId msgData liftIO $ createSndMsgDelivery db sq internalId -enqueueRatchetKeyMsgs :: AgentClient -> NonEmpty SndQueue -> CR.RcvE2ERatchetParams 'C.X448 -> AM () -enqueueRatchetKeyMsgs c (sq :| sqs) e2eEncryption = do +enqueueRatchetKeyMsgs :: AgentClient -> ConnData -> NonEmpty SndQueue -> CR.RcvE2ERatchetParams 'C.X448 -> AM () +enqueueRatchetKeyMsgs c cData (sq :| sqs) e2eEncryption = do msgId <- enqueueRatchetKey c sq e2eEncryption - mapM_ (lift . enqueueSavedMessage c msgId) $ filter isActiveSndQ sqs + mapM_ (lift . enqueueSavedMessage c msgId) $ filter (isActiveSndQ cData) sqs enqueueRatchetKey :: AgentClient -> SndQueue -> CR.RcvE2ERatchetParams 'C.X448 -> AM AgentMsgId enqueueRatchetKey c sq@SndQueue {connId} e2eEncryption = do diff --git a/src/Simplex/Messaging/Agent/Store.hs b/src/Simplex/Messaging/Agent/Store.hs index 1054d1d17..abfae2b59 100644 --- a/src/Simplex/Messaging/Agent/Store.hs +++ b/src/Simplex/Messaging/Agent/Store.hs @@ -207,8 +207,8 @@ rcvSMPQueueAddress :: RcvQueue -> SMPQueueAddress rcvSMPQueueAddress RcvQueue {server, sndId, e2ePrivKey, queueMode} = SMPQueueAddress server sndId (C.publicKey e2ePrivKey) queueMode -canAbortRcvSwitch :: VersionSMPA -> RcvQueue -> Bool -canAbortRcvSwitch connAgentVersion = maybe False canAbort . rcvSwchStatus +canAbortRcvSwitch :: ConnData -> RcvQueue -> Bool +canAbortRcvSwitch ConnData {connAgentVersion} = maybe False canAbort . rcvSwchStatus where canAbort = \case RSSwitchStarted -> True