mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-08-23 07:49:58 +00:00
agent: prevent batch deletions from failing on one connection sql error (#628)
* agent: prevent error reading one connection failing batched subscription * prevent batch deletions from failing on one connection sql error * rename
This commit is contained in:
@@ -12,6 +12,7 @@
|
||||
{-# LANGUAGE RankNTypes #-}
|
||||
{-# LANGUAGE ScopedTypeVariables #-}
|
||||
{-# LANGUAGE TupleSections #-}
|
||||
{-# LANGUAGE TypeApplications #-}
|
||||
|
||||
-- |
|
||||
-- Module : Simplex.Messaging.Agent
|
||||
@@ -472,7 +473,7 @@ deleteConnectionsAsync_ :: forall m. AgentMonad m => m () -> AgentClient -> [Con
|
||||
deleteConnectionsAsync_ onSuccess c connIds = case connIds of
|
||||
[] -> onSuccess
|
||||
_ -> do
|
||||
(_, rqs, connIds') <- prepareDeleteConnections_ getConn c connIds
|
||||
(_, rqs, connIds') <- prepareDeleteConnections_ getConns c connIds
|
||||
withStore' c $ forM_ connIds' . setConnDeleted
|
||||
void . forkIO $
|
||||
withLock (deleteLock c) "deleteConnectionsAsync" $
|
||||
@@ -628,7 +629,7 @@ type QCmdResult = (QueueStatus, Either AgentErrorType ())
|
||||
subscribeConnections' :: forall m. AgentMonad m => AgentClient -> [ConnId] -> m (Map ConnId (Either AgentErrorType ()))
|
||||
subscribeConnections' _ [] = pure M.empty
|
||||
subscribeConnections' c connIds = do
|
||||
conns :: Map ConnId (Either StoreError SomeConn) <- M.fromList . zip connIds <$> withStore' c (forM connIds . getConn)
|
||||
conns :: Map ConnId (Either StoreError SomeConn) <- M.fromList . zip connIds <$> withStore' c (`getConns` connIds)
|
||||
let (errs, cs) = M.mapEither id conns
|
||||
errs' = M.map (Left . storeError) errs
|
||||
(subRs, rcvQs) = M.mapEither rcvQueueOrResult cs
|
||||
@@ -1193,20 +1194,20 @@ disableConn c connId = do
|
||||
|
||||
-- Unlike deleteConnectionsAsync, this function does not mark connections as deleted in case of deletion failure.
|
||||
deleteConnections' :: forall m. AgentMonad m => AgentClient -> [ConnId] -> m (Map ConnId (Either AgentErrorType ()))
|
||||
deleteConnections' = deleteConnections_ getConn False
|
||||
deleteConnections' = deleteConnections_ getConns False
|
||||
|
||||
deleteDeletedConns :: forall m. AgentMonad m => AgentClient -> [ConnId] -> m (Map ConnId (Either AgentErrorType ()))
|
||||
deleteDeletedConns = deleteConnections_ getDeletedConn True
|
||||
deleteDeletedConns = deleteConnections_ getDeletedConns True
|
||||
|
||||
prepareDeleteConnections_ ::
|
||||
forall m.
|
||||
AgentMonad m =>
|
||||
(DB.Connection -> ConnId -> IO (Either StoreError SomeConn)) ->
|
||||
(DB.Connection -> [ConnId] -> IO [Either StoreError SomeConn]) ->
|
||||
AgentClient ->
|
||||
[ConnId] ->
|
||||
m (Map ConnId (Either AgentErrorType ()), [RcvQueue], [ConnId])
|
||||
prepareDeleteConnections_ getConnection c connIds = do
|
||||
conns :: Map ConnId (Either StoreError SomeConn) <- M.fromList . zip connIds <$> withStore' c (forM connIds . getConnection)
|
||||
prepareDeleteConnections_ getConnections c connIds = do
|
||||
conns :: Map ConnId (Either StoreError SomeConn) <- M.fromList . zip connIds <$> withStore' c (`getConnections` connIds)
|
||||
let (errs, cs) = M.mapEither id conns
|
||||
errs' = M.map (Left . storeError) errs
|
||||
(delRs, rcvQs) = M.mapEither rcvQueues cs
|
||||
@@ -1259,14 +1260,14 @@ deleteConnQueues c ntf rqs = do
|
||||
deleteConnections_ ::
|
||||
forall m.
|
||||
AgentMonad m =>
|
||||
(DB.Connection -> ConnId -> IO (Either StoreError SomeConn)) ->
|
||||
(DB.Connection -> [ConnId] -> IO [Either StoreError SomeConn]) ->
|
||||
Bool ->
|
||||
AgentClient ->
|
||||
[ConnId] ->
|
||||
m (Map ConnId (Either AgentErrorType ()))
|
||||
deleteConnections_ _ _ _ [] = pure M.empty
|
||||
deleteConnections_ getConnection ntf c connIds = do
|
||||
(rs, rqs, _) <- prepareDeleteConnections_ getConnection c connIds
|
||||
deleteConnections_ getConnections ntf c connIds = do
|
||||
(rs, rqs, _) <- prepareDeleteConnections_ getConnections c connIds
|
||||
rcvRs <- deleteConnQueues c ntf rqs
|
||||
let rs' = M.union rs rcvRs
|
||||
notifyResultError rs'
|
||||
@@ -1576,7 +1577,7 @@ cleanupManager c = do
|
||||
forever $ do
|
||||
void . runExceptT $
|
||||
withLock (deleteLock c) "cleanupManager" $ do
|
||||
void $ withStore' c getDeletedConns >>= deleteDeletedConns c
|
||||
void $ withStore' c getDeletedConnIds >>= deleteDeletedConns c
|
||||
withStore' c deleteUsersWithoutConns >>= mapM_ notifyUserDeleted
|
||||
threadDelay int
|
||||
where
|
||||
|
||||
@@ -43,9 +43,11 @@ module Simplex.Messaging.Agent.Store.SQLite
|
||||
createSndConn,
|
||||
getConn,
|
||||
getDeletedConn,
|
||||
getConns,
|
||||
getDeletedConns,
|
||||
getConnData,
|
||||
setConnDeleted,
|
||||
getDeletedConns,
|
||||
getDeletedConnIds,
|
||||
getRcvConn,
|
||||
deleteConn,
|
||||
upgradeRcvConnToDuplex,
|
||||
@@ -1404,6 +1406,18 @@ getAnyConn deleted' dbConn connId =
|
||||
(Nothing, Nothing, _) -> Right $ SomeConn SCNew (NewConnection cData)
|
||||
_ -> Left SEConnNotFound
|
||||
|
||||
getConns :: DB.Connection -> [ConnId] -> IO [Either StoreError SomeConn]
|
||||
getConns = getAnyConns_ False
|
||||
|
||||
getDeletedConns :: DB.Connection -> [ConnId] -> IO [Either StoreError SomeConn]
|
||||
getDeletedConns = getAnyConns_ True
|
||||
|
||||
getAnyConns_ :: Bool -> DB.Connection -> [ConnId] -> IO [Either StoreError SomeConn]
|
||||
getAnyConns_ deleted' db connIds = forM connIds $ E.handle handleDBError . getAnyConn deleted' db
|
||||
where
|
||||
handleDBError :: E.SomeException -> IO (Either StoreError SomeConn)
|
||||
handleDBError = pure . Left . SEInternal . bshow
|
||||
|
||||
getConnData :: DB.Connection -> ConnId -> IO (Maybe (ConnData, ConnectionMode))
|
||||
getConnData dbConn connId' =
|
||||
maybeFirstRow cData $ DB.query dbConn "SELECT user_id, conn_id, conn_mode, smp_agent_version, enable_ntfs, duplex_handshake, deleted FROM connections WHERE conn_id = ?;" (Only connId')
|
||||
@@ -1413,8 +1427,8 @@ getConnData dbConn connId' =
|
||||
setConnDeleted :: DB.Connection -> ConnId -> IO ()
|
||||
setConnDeleted db connId = DB.execute db "UPDATE connections SET deleted = ? WHERE conn_id = ?" (True, connId)
|
||||
|
||||
getDeletedConns :: DB.Connection -> IO [ConnId]
|
||||
getDeletedConns db = map fromOnly <$> DB.query db "SELECT conn_id FROM connections WHERE deleted = ?" (Only True)
|
||||
getDeletedConnIds :: DB.Connection -> IO [ConnId]
|
||||
getDeletedConnIds db = map fromOnly <$> DB.query db "SELECT conn_id FROM connections WHERE deleted = ?" (Only True)
|
||||
|
||||
-- | returns all connection queues, the first queue is the primary one
|
||||
getRcvQueuesByConnId_ :: DB.Connection -> ConnId -> IO (Maybe (NonEmpty RcvQueue))
|
||||
|
||||
Reference in New Issue
Block a user