From 27cfc4c62b223f96a310a15a3fd56e00ea12c91d Mon Sep 17 00:00:00 2001 From: Evgeny Poberezkin <2769109+epoberezkin@users.noreply.github.com> Date: Fri, 14 Apr 2023 10:52:18 +0200 Subject: [PATCH] move assertForeground (#724) --- src/Simplex/FileTransfer/Agent.hs | 27 ++++++++++++++++----------- src/Simplex/Messaging/Agent/Client.hs | 6 ------ 2 files changed, 16 insertions(+), 17 deletions(-) diff --git a/src/Simplex/FileTransfer/Agent.hs b/src/Simplex/FileTransfer/Agent.hs index e9a063866..c2c1c122a 100644 --- a/src/Simplex/FileTransfer/Agent.hs +++ b/src/Simplex/FileTransfer/Agent.hs @@ -171,7 +171,7 @@ runXFTPRcvWorker :: forall m. AgentMonad m => AgentClient -> XFTPServer -> TMVar runXFTPRcvWorker c srv doWork = do forever $ do void . atomically $ readTMVar doWork - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c runXFTPOperation where noWorkToDo = void . atomically $ tryTakeTMVar doWork @@ -195,7 +195,7 @@ runXFTPRcvWorker c srv doWork = do when notifyOnRetry $ notify c rcvFileEntityId $ RFERR e closeXFTPServerClient c userId server $ bshow rcvChunkId withStore' c $ \db -> updateRcvChunkReplicaDelay db rcvChunkReplicaId replicaDelay - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c loop retryDone e = rcvWorkerInternalError c rcvFileId rcvFileEntityId (Just fileTmpPath) (show e) downloadFileChunk :: RcvFileChunk -> RcvFileChunkReplica -> m () @@ -239,7 +239,7 @@ runXFTPRcvLocalWorker :: forall m. AgentMonad m => AgentClient -> TMVar () -> m runXFTPRcvLocalWorker c doWork = do forever $ do void . atomically $ readTMVar doWork - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c runXFTPOperation where runXFTPOperation :: m () @@ -360,7 +360,7 @@ runXFTPSndPrepareWorker :: forall m. AgentMonad m => AgentClient -> TMVar () -> runXFTPSndPrepareWorker c doWork = do forever $ do void . atomically $ readTMVar doWork - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c runXFTPOperation where runXFTPOperation :: m () @@ -414,7 +414,7 @@ runXFTPSndPrepareWorker c doWork = do any (\SndFileChunkReplica {replicaStatus} -> replicaStatus == SFRSCreated) replicas createChunk :: Int -> SndFileChunk -> m () createChunk numRecipients' ch = do - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c (replica, ProtoServerWithAuth srv _) <- agentOperationBracket c AOSndNetwork throwWhenInactive tryCreate withStore' c $ \db -> createSndFileReplica db ch replica addXFTPSndWorker c $ Just srv @@ -426,7 +426,7 @@ runXFTPSndPrepareWorker c doWork = do createWithNextSrv usedSrvs `catchError` \e -> retryOnError "XFTP prepare worker" (retryLoop loop) (throwError e) e where - retryLoop loop = atomically (checkAgentForeground c) >> loop + retryLoop loop = atomically (assertAgentForeground c) >> loop createWithNextSrv usedSrvs = do deleted <- withStore' c $ \db -> getSndFileDeleted db sndFileId when deleted $ throwError $ INTERNAL "file deleted, aborting chunk creation" @@ -444,7 +444,7 @@ runXFTPSndWorker :: forall m. AgentMonad m => AgentClient -> XFTPServer -> TMVar runXFTPSndWorker c srv doWork = do forever $ do void . atomically $ readTMVar doWork - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c agentOperationBracket c AOSndNetwork throwWhenInactive runXFTPOperation where noWorkToDo = void . atomically $ tryTakeTMVar doWork @@ -468,7 +468,7 @@ runXFTPSndWorker c srv doWork = do when notifyOnRetry $ notify c sndFileEntityId $ SFERR e closeXFTPServerClient c userId server $ bshow sndChunkId withStore' c $ \db -> updateSndChunkReplicaDelay db sndChunkReplicaId replicaDelay - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c loop retryDone e = sndWorkerInternalError c sndFileId sndFileEntityId (Just filePrefixPath) (show e) uploadFileChunk :: SndFileChunk -> SndFileChunkReplica -> m () @@ -476,7 +476,7 @@ runXFTPSndWorker c srv doWork = do replica'@SndFileChunkReplica {sndChunkReplicaId} <- addRecipients sndFileChunk replica fsFilePath <- toFSFilePath filePath let chunkSpec' = chunkSpec {filePath = fsFilePath} :: XFTPChunkSpec - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c agentXFTPUploadChunk c userId sndChunkId replica' chunkSpec' sf@SndFile {sndFileEntityId, prefixPath, chunks} <- withStore c $ \db -> do updateSndChunkReplicaStatus db sndChunkReplicaId SFRSUploaded @@ -602,7 +602,7 @@ runXFTPDelWorker :: forall m. AgentMonad m => AgentClient -> XFTPServer -> TMVar runXFTPDelWorker c srv doWork = do forever $ do void . atomically $ readTMVar doWork - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c runXFTPOperation where noWorkToDo = void . atomically $ tryTakeTMVar doWork @@ -626,7 +626,7 @@ runXFTPDelWorker c srv doWork = do when notifyOnRetry $ notify c "" $ SFERR e closeXFTPServerClient c userId server replId withStore' c $ \db -> updateDeletedSndChunkReplicaDelay db deletedSndChunkReplicaId replicaDelay - atomically $ checkAgentForeground c + atomically $ assertAgentForeground c loop retryDone e = delWorkerInternalError c deletedSndChunkReplicaId e deleteChunkReplica :: DeletedSndChunkReplica -> m () @@ -638,3 +638,8 @@ delWorkerInternalError :: AgentMonad m => AgentClient -> Int64 -> AgentErrorType delWorkerInternalError c deletedSndChunkReplicaId e = do withStore' c $ \db -> deleteDeletedSndChunkReplica db deletedSndChunkReplicaId notify c "" $ SFERR e + +assertAgentForeground :: AgentClient -> STM () +assertAgentForeground c = do + throwWhenInactive c + waitUntilForeground c diff --git a/src/Simplex/Messaging/Agent/Client.hs b/src/Simplex/Messaging/Agent/Client.hs index 20844d938..434bdc3b7 100644 --- a/src/Simplex/Messaging/Agent/Client.hs +++ b/src/Simplex/Messaging/Agent/Client.hs @@ -82,7 +82,6 @@ module Simplex.Messaging.Agent.Client beginAgentOperation, endAgentOperation, waitUntilForeground, - checkAgentForeground, suspendSendingAndDatabase, suspendOperation, notifySuspended, @@ -1217,11 +1216,6 @@ agentOperationBracket c op check action = waitUntilForeground :: AgentClient -> STM () waitUntilForeground c = unlessM ((ASForeground ==) <$> readTVar (agentState c)) retry -checkAgentForeground :: AgentClient -> STM () -checkAgentForeground c = do - throwWhenInactive c - waitUntilForeground c - withStore' :: AgentMonad m => AgentClient -> (DB.Connection -> IO a) -> m a withStore' c action = withStore c $ fmap Right . action