From 2b1a02d7d49493bdbea09b75defa9a67ab4f2051 Mon Sep 17 00:00:00 2001 From: spaced4ndy <8711996+spaced4ndy@users.noreply.github.com> Date: Fri, 17 Mar 2023 13:50:49 +0400 Subject: [PATCH] xftp: reconnect XFTP client on replica retry (#689) * xftp: re-create XFTP client on replica retry * closeXFTPSessionClient * refactor --------- Co-authored-by: Evgeny Poberezkin <2769109+epoberezkin@users.noreply.github.com> --- src/Simplex/FileTransfer/Agent.hs | 3 ++- src/Simplex/Messaging/Agent/Client.hs | 30 +++++++++++++++-------- src/Simplex/Messaging/Agent/Env/SQLite.hs | 3 +++ 3 files changed, 25 insertions(+), 11 deletions(-) diff --git a/src/Simplex/FileTransfer/Agent.hs b/src/Simplex/FileTransfer/Agent.hs index 6e4a84c7e..aebf5c6a4 100644 --- a/src/Simplex/FileTransfer/Agent.hs +++ b/src/Simplex/FileTransfer/Agent.hs @@ -95,7 +95,7 @@ runXFTPWorker c srv doWork = do case nextChunk of Nothing -> noWorkToDo Just RcvFileChunk {rcvFileId, rcvFileEntityId, fileTmpPath, replicas = []} -> workerInternalError c rcvFileId rcvFileEntityId (Just fileTmpPath) "chunk has no replicas" - Just fc@RcvFileChunk {rcvFileId, rcvFileEntityId, fileTmpPath, replicas = replica@RcvFileChunkReplica {rcvChunkReplicaId, delay} : _} -> do + Just fc@RcvFileChunk {userId, rcvFileId, rcvFileEntityId, fileTmpPath, replicas = replica@RcvFileChunkReplica {rcvChunkReplicaId, delay} : _} -> do ri <- asks $ reconnectInterval . config let ri' = maybe ri (\d -> ri {initialInterval = d, increaseAfter = 0}) delay withRetryInterval ri' $ \delay' loop -> @@ -110,6 +110,7 @@ runXFTPWorker c srv doWork = do else done e where retryLoop = do + closeXFTPServerClient c userId replica withStore' c $ \db -> updateRcvChunkReplicaDelay db rcvChunkReplicaId replicaDelay atomically $ endAgentOperation c AORcvNetwork atomically $ throwWhenInactive c diff --git a/src/Simplex/Messaging/Agent/Client.hs b/src/Simplex/Messaging/Agent/Client.hs index 89559329f..8b8c1d9fa 100644 --- a/src/Simplex/Messaging/Agent/Client.hs +++ b/src/Simplex/Messaging/Agent/Client.hs @@ -25,6 +25,7 @@ module Simplex.Messaging.Agent.Client withConnLock, closeAgentClient, closeProtocolServerClients, + closeXFTPServerClient, runSMPServerTest, newRcvQueue, subscribeQueues, @@ -554,6 +555,7 @@ closeAgentClient c = liftIO $ do atomically $ writeTVar (active c) False closeProtocolServerClients c smpClients closeProtocolServerClients c ntfClients + closeProtocolServerClients c xftpClients cancelActions . actions $ reconnections c cancelActions . actions $ asyncClients c cancelActions $ smpQueueMsgDeliveries c @@ -586,14 +588,22 @@ throwWhenNoDelivery c SndQueue {server, sndId} = closeProtocolServerClients :: ProtocolServerClient err msg => AgentClient -> (AgentClient -> TMap (TransportSession msg) (ClientVar msg)) -> IO () closeProtocolServerClients c clientsSel = - atomically (swapTVar cs M.empty) >>= mapM_ (forkIO . closeClient) - where - cs = clientsSel c - closeClient cVar = do - NetworkConfig {tcpConnectTimeout} <- readTVarIO $ useNetworkConfig c - tcpConnectTimeout `timeout` atomically (readTMVar cVar) >>= \case - Just (Right client) -> closeProtocolServerClient client `catchAll_` pure () - _ -> pure () + atomically (clientsSel c `swapTVar` M.empty) >>= mapM_ (forkIO . closeClient_ c) + +closeClient :: ProtocolServerClient err msg => AgentClient -> (AgentClient -> TMap (TransportSession msg) (ClientVar msg)) -> TransportSession msg -> IO () +closeClient c clientSel tSess = + atomically (TM.lookupDelete tSess $ clientSel c) >>= mapM_ (closeClient_ c) + +closeClient_ :: ProtocolServerClient err msg => AgentClient -> ClientVar msg -> IO () +closeClient_ c cVar = do + NetworkConfig {tcpConnectTimeout} <- readTVarIO $ useNetworkConfig c + tcpConnectTimeout `timeout` atomically (readTMVar cVar) >>= \case + Just (Right client) -> closeProtocolServerClient client `catchAll_` pure () + _ -> pure () + +closeXFTPServerClient :: AgentMonad' m => AgentClient -> UserId -> RcvFileChunkReplica -> m () +closeXFTPServerClient c userId RcvFileChunkReplica {server, replicaId = ChunkReplicaId fId} = + mkTransportSession c userId server fId >>= liftIO . closeClient c xftpClients cancelActions :: (Foldable f, Monoid (f (Async ()))) => TVar (f (Async ())) -> IO () cancelActions as = atomically (swapTVar as mempty) >>= mapM_ (forkIO . uninterruptibleCancel) @@ -712,7 +722,7 @@ runSMPServerTest c userId (ProtoServerWithAuth srv auth) = do testErr :: SMPTestStep -> SMPClientError -> SMPTestFailure testErr step = SMPTestFailure step . protocolClientError SMP addr -mkTransportSession :: AgentMonad m => AgentClient -> UserId -> ProtoServer msg -> EntityId -> m (TransportSession msg) +mkTransportSession :: AgentMonad' m => AgentClient -> UserId -> ProtoServer msg -> EntityId -> m (TransportSession msg) mkTransportSession c userId srv entityId = mkTSession userId srv entityId <$> getSessionMode c mkTSession :: UserId -> ProtoServer msg -> EntityId -> TransportSessionMode -> TransportSession msg @@ -724,7 +734,7 @@ mkSMPTransportSession c q = mkSMPTSession q <$> getSessionMode c mkSMPTSession :: SMPQueueRec q => q -> TransportSessionMode -> SMPTransportSession mkSMPTSession q = mkTSession (qUserId q) (qServer q) (qConnId q) -getSessionMode :: AgentMonad m => AgentClient -> m TransportSessionMode +getSessionMode :: AgentMonad' m => AgentClient -> m TransportSessionMode getSessionMode = fmap sessionMode . readTVarIO . useNetworkConfig newRcvQueue :: AgentMonad m => AgentClient -> UserId -> ConnId -> SMPServerWithAuth -> VersionRange -> m (RcvQueue, SMPQueueUri) diff --git a/src/Simplex/Messaging/Agent/Env/SQLite.hs b/src/Simplex/Messaging/Agent/Env/SQLite.hs index e1263b90d..daf43b96b 100644 --- a/src/Simplex/Messaging/Agent/Env/SQLite.hs +++ b/src/Simplex/Messaging/Agent/Env/SQLite.hs @@ -12,6 +12,7 @@ module Simplex.Messaging.Agent.Env.SQLite ( AgentMonad, + AgentMonad', AgentConfig (..), AgentDatabase (..), databaseFile, @@ -64,6 +65,8 @@ import UnliftIO.STM -- | Agent monad with MonadReader Env and MonadError AgentErrorType type AgentMonad m = (MonadUnliftIO m, MonadReader Env m, MonadError AgentErrorType m) +type AgentMonad' m = (MonadUnliftIO m, MonadReader Env m) + data InitialAgentServers = InitialAgentServers { smp :: Map UserId (NonEmpty SMPServerWithAuth), ntf :: [NtfServer],