From cfe995325a8468d69ad29e8974240244d27396ea Mon Sep 17 00:00:00 2001 From: Evgeny Poberezkin <2769109+epoberezkin@users.noreply.github.com> Date: Sat, 4 Feb 2023 20:46:45 +0000 Subject: [PATCH] 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 --- src/Simplex/Messaging/Agent.hs | 23 +++++++++++---------- src/Simplex/Messaging/Agent/Store/SQLite.hs | 20 +++++++++++++++--- 2 files changed, 29 insertions(+), 14 deletions(-) diff --git a/src/Simplex/Messaging/Agent.hs b/src/Simplex/Messaging/Agent.hs index d45f15c43..e132cab41 100644 --- a/src/Simplex/Messaging/Agent.hs +++ b/src/Simplex/Messaging/Agent.hs @@ -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 diff --git a/src/Simplex/Messaging/Agent/Store/SQLite.hs b/src/Simplex/Messaging/Agent/Store/SQLite.hs index a8d7a7bab..810f68b6d 100644 --- a/src/Simplex/Messaging/Agent/Store/SQLite.hs +++ b/src/Simplex/Messaging/Agent/Store/SQLite.hs @@ -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))