mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-10-05 20:57:16 +00:00
xftp: fix race when sending the file (#1879)
* xftp: fix race when sending the file * simpler test * fix tests * test --------- Co-authored-by: Evgeny @ SimpleX Chat <259188159+evgeny-simplex@users.noreply.github.com>
This commit is contained in:
co-authored by
Evgeny @ SimpleX Chat
parent
c7163afa3a
commit
35f1b145a7
@@ -17,11 +17,16 @@ module Simplex.FileTransfer.Agent
|
||||
toFSFilePath,
|
||||
-- Receiving files
|
||||
xftpReceiveFile',
|
||||
xftpPrepareReceiveFile',
|
||||
xftpStartReceiveFile',
|
||||
xftpDeleteRcvFile',
|
||||
xftpDeleteRcvFiles',
|
||||
-- Sending files
|
||||
xftpSendFile',
|
||||
xftpSendDescription',
|
||||
xftpPrepareSendFile',
|
||||
xftpPrepareSendDescription',
|
||||
xftpStartSendFile',
|
||||
deleteSndFileInternal,
|
||||
deleteSndFilesInternal,
|
||||
deleteSndFileRemote,
|
||||
@@ -127,7 +132,12 @@ closeXFTPAgent a = do
|
||||
stopWorkers workers = atomically (swapTVar workers M.empty) >>= mapM_ (liftIO . cancelWorker)
|
||||
|
||||
xftpReceiveFile' :: AgentClient -> UserId -> ValidFileDescription 'FRecipient -> Maybe CryptoFileArgs -> Bool -> AM RcvFileId
|
||||
xftpReceiveFile' c userId (ValidFileDescription fd@FileDescription {chunks, redirect}) cfArgs approvedRelays = do
|
||||
xftpReceiveFile' c userId vfd cfArgs approvedRelays = do
|
||||
fId <- xftpPrepareReceiveFile' c userId vfd cfArgs approvedRelays
|
||||
fId <$ xftpStartReceiveFile' c fId
|
||||
|
||||
xftpPrepareReceiveFile' :: AgentClient -> UserId -> ValidFileDescription 'FRecipient -> Maybe CryptoFileArgs -> Bool -> AM RcvFileId
|
||||
xftpPrepareReceiveFile' c userId (ValidFileDescription fd@FileDescription {redirect}) cfArgs approvedRelays = do
|
||||
g <- asks random
|
||||
prefixPath <- lift $ getPrefixPath "rcv.xftp"
|
||||
createDirectory prefixPath
|
||||
@@ -137,7 +147,7 @@ xftpReceiveFile' c userId (ValidFileDescription fd@FileDescription {chunks, redi
|
||||
lift $ createDirectory =<< toFSFilePath relTmpPath
|
||||
lift $ createEmptyFile =<< toFSFilePath relSavePath
|
||||
let saveFile = CryptoFile relSavePath cfArgs
|
||||
fId <- case redirect of
|
||||
case redirect of
|
||||
Nothing -> withStore c $ \db -> createRcvFile db g userId fd relPrefixPath relTmpPath saveFile approvedRelays
|
||||
Just _ -> do
|
||||
-- prepare description paths
|
||||
@@ -149,8 +159,11 @@ xftpReceiveFile' c userId (ValidFileDescription fd@FileDescription {chunks, redi
|
||||
let saveFileRedirect = CryptoFile relSavePathRedirect $ Just cfArgsRedirect
|
||||
-- create download tasks
|
||||
withStore c $ \db -> createRcvFileRedirect db g userId fd relPrefixPath relTmpPathRedirect saveFileRedirect relTmpPath saveFile approvedRelays
|
||||
forM_ chunks (downloadChunk c)
|
||||
pure fId
|
||||
|
||||
xftpStartReceiveFile' :: AgentClient -> RcvFileId -> AM ()
|
||||
xftpStartReceiveFile' c rcvFileEntityId = do
|
||||
srvs <- withStore c (`startPreparedRcvFile` rcvFileEntityId)
|
||||
lift $ forM_ srvs $ void . getXFTPRcvWorker True c . Just
|
||||
|
||||
downloadChunk :: AgentClient -> FileChunk -> AM ()
|
||||
downloadChunk c FileChunk {replicas = (FileChunkReplica {server} : _)} = do
|
||||
@@ -353,6 +366,11 @@ notify c entId cmd = atomically $ writeTBQueue (subQ c) ("", entId, AEvt (sAEnti
|
||||
|
||||
xftpSendFile' :: AgentClient -> UserId -> CryptoFile -> Int -> Maybe Word32 -> AM SndFileId
|
||||
xftpSendFile' c userId file numRecipients storageHours = do
|
||||
fId <- xftpPrepareSendFile' c userId file numRecipients storageHours
|
||||
fId <$ xftpStartSendFile' c fId
|
||||
|
||||
xftpPrepareSendFile' :: AgentClient -> UserId -> CryptoFile -> Int -> Maybe Word32 -> AM SndFileId
|
||||
xftpPrepareSendFile' c userId file numRecipients storageHours = do
|
||||
g <- asks random
|
||||
prefixPath <- lift $ getPrefixPath "snd.xftp"
|
||||
createDirectory prefixPath
|
||||
@@ -360,12 +378,15 @@ xftpSendFile' c userId file numRecipients storageHours = do
|
||||
key <- atomically $ C.randomSbKey g
|
||||
nonce <- atomically $ C.randomCbNonce g
|
||||
-- saving absolute filePath will not allow to restore file encryption after app update, but it's a short window
|
||||
fId <- withStore c $ \db -> createSndFile db g userId file numRecipients relPrefixPath key nonce Nothing storageHours
|
||||
lift . void $ getXFTPSndWorker True c Nothing
|
||||
pure fId
|
||||
withStore c $ \db -> createSndFile db g userId file numRecipients relPrefixPath key nonce Nothing storageHours
|
||||
|
||||
xftpSendDescription' :: AgentClient -> UserId -> ValidFileDescription 'FRecipient -> Int -> AM SndFileId
|
||||
xftpSendDescription' c userId (ValidFileDescription fdDirect@FileDescription {size, digest}) numRecipients = do
|
||||
xftpSendDescription' c userId vfd numRecipients = do
|
||||
fId <- xftpPrepareSendDescription' c userId vfd numRecipients
|
||||
fId <$ xftpStartSendFile' c fId
|
||||
|
||||
xftpPrepareSendDescription' :: AgentClient -> UserId -> ValidFileDescription 'FRecipient -> Int -> AM SndFileId
|
||||
xftpPrepareSendDescription' c userId (ValidFileDescription fdDirect@FileDescription {size, digest}) numRecipients = do
|
||||
g <- asks random
|
||||
prefixPath <- lift $ getPrefixPath "snd.xftp"
|
||||
createDirectory prefixPath
|
||||
@@ -376,9 +397,12 @@ xftpSendDescription' c userId (ValidFileDescription fdDirect@FileDescription {si
|
||||
liftError (FILE . FILE_IO . show) $ CF.writeFile file (LB.fromStrict $ strEncode fdDirect)
|
||||
key <- atomically $ C.randomSbKey g
|
||||
nonce <- atomically $ C.randomCbNonce g
|
||||
fId <- withStore c $ \db -> createSndFile db g userId file numRecipients relPrefixPath key nonce (Just RedirectFileInfo {size, digest}) Nothing
|
||||
withStore c $ \db -> createSndFile db g userId file numRecipients relPrefixPath key nonce (Just RedirectFileInfo {size, digest}) Nothing
|
||||
|
||||
xftpStartSendFile' :: AgentClient -> SndFileId -> AM ()
|
||||
xftpStartSendFile' c sndFileEntityId = do
|
||||
withStore c (`startPreparedSndFile` sndFileEntityId)
|
||||
lift . void $ getXFTPSndWorker True c Nothing
|
||||
pure fId
|
||||
|
||||
resumeXFTPSndWork :: AgentClient -> Maybe XFTPServer -> AM' ()
|
||||
resumeXFTPSndWork = void .: getXFTPSndWorker False
|
||||
|
||||
@@ -93,7 +93,8 @@ data RcvFile = RcvFile
|
||||
deriving (Show)
|
||||
|
||||
data RcvFileStatus
|
||||
= RFSReceiving
|
||||
= RFSPrepared
|
||||
| RFSReceiving
|
||||
| RFSReceived
|
||||
| RFSDecrypting
|
||||
| RFSComplete
|
||||
@@ -106,6 +107,7 @@ instance ToField RcvFileStatus where toField = toField . textEncode
|
||||
|
||||
instance TextEncoding RcvFileStatus where
|
||||
textDecode = \case
|
||||
"prepared" -> Just RFSPrepared
|
||||
"receiving" -> Just RFSReceiving
|
||||
"received" -> Just RFSReceived
|
||||
"decrypting" -> Just RFSDecrypting
|
||||
@@ -113,6 +115,7 @@ instance TextEncoding RcvFileStatus where
|
||||
"error" -> Just RFSError
|
||||
_ -> Nothing
|
||||
textEncode = \case
|
||||
RFSPrepared -> "prepared"
|
||||
RFSReceiving -> "receiving"
|
||||
RFSReceived -> "received"
|
||||
RFSDecrypting -> "decrypting"
|
||||
@@ -177,7 +180,8 @@ sndFileEncPath :: FilePath -> FilePath
|
||||
sndFileEncPath prefixPath = prefixPath </> "xftp.encrypted"
|
||||
|
||||
data SndFileStatus
|
||||
= SFSNew -- db record created
|
||||
= SFSPrepared
|
||||
| SFSNew -- db record created
|
||||
| SFSEncrypting -- encryption started
|
||||
| SFSEncrypted -- encryption complete
|
||||
| SFSUploading -- all chunk replicas are created on servers
|
||||
@@ -191,6 +195,7 @@ instance ToField SndFileStatus where toField = toField . textEncode
|
||||
|
||||
instance TextEncoding SndFileStatus where
|
||||
textDecode = \case
|
||||
"prepared" -> Just SFSPrepared
|
||||
"new" -> Just SFSNew
|
||||
"encrypting" -> Just SFSEncrypting
|
||||
"encrypted" -> Just SFSEncrypted
|
||||
@@ -199,6 +204,7 @@ instance TextEncoding SndFileStatus where
|
||||
"error" -> Just SFSError
|
||||
_ -> Nothing
|
||||
textEncode = \case
|
||||
SFSPrepared -> "prepared"
|
||||
SFSNew -> "new"
|
||||
SFSEncrypting -> "encrypting"
|
||||
SFSEncrypted -> "encrypted"
|
||||
|
||||
@@ -127,10 +127,15 @@ module Simplex.Messaging.Agent
|
||||
xftpStartWorkers,
|
||||
xftpStartSndWorkers,
|
||||
xftpReceiveFile,
|
||||
xftpPrepareReceiveFile,
|
||||
xftpStartReceiveFile,
|
||||
xftpDeleteRcvFile,
|
||||
xftpDeleteRcvFiles,
|
||||
xftpSendFile,
|
||||
xftpSendDescription,
|
||||
xftpPrepareSendFile,
|
||||
xftpPrepareSendDescription,
|
||||
xftpStartSendFile,
|
||||
xftpDeleteSndFileInternal,
|
||||
xftpDeleteSndFilesInternal,
|
||||
xftpDeleteSndFileRemote,
|
||||
@@ -190,7 +195,7 @@ import Data.Time.Clock
|
||||
import Data.Time.Clock.System (systemToUTCTime)
|
||||
import Data.Traversable (mapAccumL)
|
||||
import Data.Word (Word16, Word32)
|
||||
import Simplex.FileTransfer.Agent (closeXFTPAgent, deleteSndFileInternal, deleteSndFileRemote, deleteSndFilesInternal, deleteSndFilesRemote, startXFTPSndWorkers, startXFTPWorkers, toFSFilePath, xftpDeleteRcvFile', xftpDeleteRcvFiles', xftpReceiveFile', xftpSendDescription', xftpSendFile')
|
||||
import Simplex.FileTransfer.Agent (closeXFTPAgent, deleteSndFileInternal, deleteSndFileRemote, deleteSndFilesInternal, deleteSndFilesRemote, startXFTPSndWorkers, startXFTPWorkers, toFSFilePath, xftpDeleteRcvFile', xftpDeleteRcvFiles', xftpPrepareReceiveFile', xftpPrepareSendDescription', xftpPrepareSendFile', xftpReceiveFile', xftpSendDescription', xftpSendFile', xftpStartReceiveFile', xftpStartSendFile')
|
||||
import Simplex.FileTransfer.Description (ValidFileDescription)
|
||||
import Simplex.FileTransfer.Protocol (FileParty (..))
|
||||
import Simplex.FileTransfer.Types (RcvFileId, SndFileId)
|
||||
@@ -769,6 +774,14 @@ xftpReceiveFile :: AgentClient -> UserId -> ValidFileDescription 'FRecipient ->
|
||||
xftpReceiveFile c = withAgentEnv c .:: xftpReceiveFile' c
|
||||
{-# INLINE xftpReceiveFile #-}
|
||||
|
||||
xftpPrepareReceiveFile :: AgentClient -> UserId -> ValidFileDescription 'FRecipient -> Maybe CryptoFileArgs -> Bool -> AE RcvFileId
|
||||
xftpPrepareReceiveFile c = withAgentEnv c .:: xftpPrepareReceiveFile' c
|
||||
{-# INLINE xftpPrepareReceiveFile #-}
|
||||
|
||||
xftpStartReceiveFile :: AgentClient -> RcvFileId -> AE ()
|
||||
xftpStartReceiveFile c = withAgentEnv c . xftpStartReceiveFile' c
|
||||
{-# INLINE xftpStartReceiveFile #-}
|
||||
|
||||
-- | Delete XFTP rcv file (deletes work files from file system and db records)
|
||||
xftpDeleteRcvFile :: AgentClient -> RcvFileId -> IO ()
|
||||
xftpDeleteRcvFile c = withAgentEnv' c . xftpDeleteRcvFile' c
|
||||
@@ -789,6 +802,18 @@ xftpSendDescription :: AgentClient -> UserId -> ValidFileDescription 'FRecipient
|
||||
xftpSendDescription c = withAgentEnv c .:. xftpSendDescription' c
|
||||
{-# INLINE xftpSendDescription #-}
|
||||
|
||||
xftpPrepareSendFile :: AgentClient -> UserId -> CryptoFile -> Int -> Maybe Word32 -> AE SndFileId
|
||||
xftpPrepareSendFile c = withAgentEnv c .:: xftpPrepareSendFile' c
|
||||
{-# INLINE xftpPrepareSendFile #-}
|
||||
|
||||
xftpPrepareSendDescription :: AgentClient -> UserId -> ValidFileDescription 'FRecipient -> Int -> AE SndFileId
|
||||
xftpPrepareSendDescription c = withAgentEnv c .:. xftpPrepareSendDescription' c
|
||||
{-# INLINE xftpPrepareSendDescription #-}
|
||||
|
||||
xftpStartSendFile :: AgentClient -> SndFileId -> AE ()
|
||||
xftpStartSendFile c = withAgentEnv c . xftpStartSendFile' c
|
||||
{-# INLINE xftpStartSendFile #-}
|
||||
|
||||
-- | Delete XFTP snd file internally (deletes work files from file system and db records)
|
||||
xftpDeleteSndFileInternal :: AgentClient -> SndFileId -> IO ()
|
||||
xftpDeleteSndFileInternal c = withAgentEnv' c . deleteSndFileInternal c
|
||||
|
||||
@@ -215,6 +215,7 @@ module Simplex.Messaging.Agent.Store.AgentStore
|
||||
-- Rcv files
|
||||
createRcvFile,
|
||||
createRcvFileRedirect,
|
||||
startPreparedRcvFile,
|
||||
lockRcvFileForUpdate,
|
||||
getRcvFile,
|
||||
getRcvFileByEntityId,
|
||||
@@ -236,6 +237,7 @@ module Simplex.Messaging.Agent.Store.AgentStore
|
||||
getRcvFilesExpired,
|
||||
-- Snd files
|
||||
createSndFile,
|
||||
startPreparedSndFile,
|
||||
lockSndFileForUpdate,
|
||||
getSndFile,
|
||||
getSndFileByEntityId,
|
||||
@@ -3163,6 +3165,29 @@ createRcvFileRedirect db gVar userId redirectFd@FileDescription {chunks = redire
|
||||
chunks = []
|
||||
}
|
||||
|
||||
startPreparedRcvFile :: DB.Connection -> RcvFileId -> IO (Either StoreError [XFTPServer])
|
||||
startPreparedRcvFile db rcvFileEntityId = runExceptT $ do
|
||||
rcvFileId <- ExceptT $ getRcvFileIdByEntityId_ db rcvFileEntityId
|
||||
liftIO $ do
|
||||
updatedAt <- getCurrentTime
|
||||
DB.execute
|
||||
db
|
||||
"UPDATE rcv_files SET status = ?, updated_at = ? WHERE (rcv_file_id = ? OR redirect_id = ?) AND status = ?"
|
||||
(RFSReceiving, updatedAt, rcvFileId, rcvFileId, RFSPrepared)
|
||||
map toXFTPServer
|
||||
<$> DB.query
|
||||
db
|
||||
[sql|
|
||||
SELECT DISTINCT
|
||||
s.xftp_host, s.xftp_port, s.xftp_key_hash
|
||||
FROM rcv_file_chunk_replicas r
|
||||
JOIN xftp_servers s ON s.xftp_server_id = r.xftp_server_id
|
||||
JOIN rcv_file_chunks c ON c.rcv_file_chunk_id = r.rcv_file_chunk_id
|
||||
JOIN rcv_files f ON f.rcv_file_id = c.rcv_file_id
|
||||
WHERE (f.rcv_file_id = ? OR f.redirect_id = ?) AND r.replica_number = 1
|
||||
|]
|
||||
(rcvFileId, rcvFileId)
|
||||
|
||||
insertRcvFile :: DB.Connection -> TVar ChaChaDRG -> UserId -> FileDescription 'FRecipient -> FilePath -> FilePath -> CryptoFile -> Maybe DBRcvFileId -> Maybe RcvFileId -> Bool -> IO (Either StoreError (RcvFileId, DBRcvFileId))
|
||||
insertRcvFile db gVar userId FileDescription {size, digest, key, nonce, chunkSize, redirect} prefixPath tmpPath (CryptoFile savePath cfArgs) redirectId_ redirectEntityId_ approvedRelays = runExceptT $ do
|
||||
let (redirectDigest_, redirectSize_) = case redirect of
|
||||
@@ -3173,7 +3198,7 @@ insertRcvFile db gVar userId FileDescription {size, digest, key, nonce, chunkSiz
|
||||
DB.execute
|
||||
db
|
||||
"INSERT INTO rcv_files (rcv_file_entity_id, user_id, size, digest, key, nonce, chunk_size, prefix_path, tmp_path, save_path, save_file_key, save_file_nonce, status, redirect_id, redirect_entity_id, redirect_digest, redirect_size, approved_relays) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)"
|
||||
((Binary rcvFileEntityId, userId, size, digest, key, nonce, chunkSize, prefixPath, tmpPath) :. (savePath, fileKey <$> cfArgs, fileNonce <$> cfArgs, RFSReceiving, redirectId_, Binary <$> redirectEntityId_, redirectDigest_, redirectSize_, BI approvedRelays))
|
||||
((Binary rcvFileEntityId, userId, size, digest, key, nonce, chunkSize, prefixPath, tmpPath) :. (savePath, fileKey <$> cfArgs, fileNonce <$> cfArgs, RFSPrepared, redirectId_, Binary <$> redirectEntityId_, redirectDigest_, redirectSize_, BI approvedRelays))
|
||||
rcvFileId <- liftIO $ insertedRowId db
|
||||
pure (rcvFileEntityId, rcvFileId)
|
||||
|
||||
@@ -3478,13 +3503,23 @@ createSndFile db gVar userId (CryptoFile path cfArgs) numRecipients prefixPath k
|
||||
DB.execute
|
||||
db
|
||||
"INSERT INTO snd_files (snd_file_entity_id, user_id, path, src_file_key, src_file_nonce, num_recipients, prefix_path, key, nonce, status, redirect_size, redirect_digest, storage_time) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)"
|
||||
((Binary sndFileEntityId, userId, path, fileKey <$> cfArgs, fileNonce <$> cfArgs, numRecipients) :. (prefixPath, key, nonce, SFSNew, redirectSize_, redirectDigest_, storageHours))
|
||||
((Binary sndFileEntityId, userId, path, fileKey <$> cfArgs, fileNonce <$> cfArgs, numRecipients) :. (prefixPath, key, nonce, SFSPrepared, redirectSize_, redirectDigest_, storageHours))
|
||||
where
|
||||
(redirectSize_, redirectDigest_) =
|
||||
case redirect_ of
|
||||
Nothing -> (Nothing, Nothing)
|
||||
Just RedirectFileInfo {size, digest} -> (Just size, Just digest)
|
||||
|
||||
startPreparedSndFile :: DB.Connection -> SndFileId -> IO (Either StoreError ())
|
||||
startPreparedSndFile db sndFileEntityId = runExceptT $ do
|
||||
sndFileId <- ExceptT $ getSndFileIdByEntityId_ db sndFileEntityId
|
||||
liftIO $ do
|
||||
updatedAt <- getCurrentTime
|
||||
DB.execute
|
||||
db
|
||||
"UPDATE snd_files SET status = ?, updated_at = ? WHERE snd_file_id = ? AND status = ?"
|
||||
(SFSNew, updatedAt, sndFileId, SFSPrepared)
|
||||
|
||||
getSndFileByEntityId :: DB.Connection -> SndFileId -> IO (Either StoreError SndFile)
|
||||
getSndFileByEntityId db sndFileEntityId = runExceptT $ do
|
||||
sndFileId <- ExceptT $ getSndFileIdByEntityId_ db sndFileEntityId
|
||||
|
||||
@@ -4628,8 +4628,7 @@ testServerMultipleIdentities =
|
||||
testE2ERatchetParams12
|
||||
|
||||
testWaitForUserNetwork :: IO ()
|
||||
testWaitForUserNetwork = do
|
||||
a <- getSMPAgentClient' 1 aCfg initAgentServers testDB
|
||||
testWaitForUserNetwork = withAgent 1 aCfg initAgentServers testDB $ \a -> do
|
||||
noNetworkDelay a
|
||||
setUserNetworkInfo a $ UserNetworkInfo UNNone False
|
||||
networkDelay a 100000
|
||||
@@ -4646,8 +4645,7 @@ testWaitForUserNetwork = do
|
||||
aCfg = agentCfg {userNetworkInterval = 100000, userOfflineDelay = 0}
|
||||
|
||||
testDoNotResetOnlineToOffline :: IO ()
|
||||
testDoNotResetOnlineToOffline = do
|
||||
a <- getSMPAgentClient' 1 aCfg initAgentServers testDB
|
||||
testDoNotResetOnlineToOffline = withAgent 1 aCfg initAgentServers testDB $ \a -> do
|
||||
noNetworkDelay a
|
||||
setUserNetworkInfo a $ UserNetworkInfo UNWifi False
|
||||
networkDelay a 100000
|
||||
@@ -4668,8 +4666,7 @@ testDoNotResetOnlineToOffline = do
|
||||
aCfg = agentCfg {userNetworkInterval = 100000, userOfflineDelay = 0.1}
|
||||
|
||||
testResumeMultipleThreads :: IO ()
|
||||
testResumeMultipleThreads = do
|
||||
a <- getSMPAgentClient' 1 aCfg initAgentServers testDB
|
||||
testResumeMultipleThreads = withAgent 1 aCfg initAgentServers testDB $ \a -> do
|
||||
noNetworkDelay a
|
||||
setUserNetworkInfo a $ UserNetworkInfo UNNone False
|
||||
vs <-
|
||||
|
||||
@@ -743,9 +743,11 @@ testGetNextRcvChunkToDownload st = do
|
||||
withTransaction st $ \db -> do
|
||||
Right Nothing <- getNextRcvChunkToDownload db xftpServer1 86400
|
||||
|
||||
Right _ <- createRcvFile db g 1 rcvFileDescr1 "filepath" "filepath" (CryptoFile "filepath" Nothing) True
|
||||
Right fId1 <- createRcvFile db g 1 rcvFileDescr1 "filepath" "filepath" (CryptoFile "filepath" Nothing) True
|
||||
Right _ <- startPreparedRcvFile db fId1
|
||||
DB.execute_ db "UPDATE rcv_file_chunk_replicas SET replica_key = cast('bad' as blob) WHERE rcv_file_chunk_replica_id = 1"
|
||||
Right fId2 <- createRcvFile db g 1 rcvFileDescr1 "filepath" "filepath" (CryptoFile "filepath" Nothing) True
|
||||
Right _ <- startPreparedRcvFile db fId2
|
||||
|
||||
Left e <- getNextRcvChunkToDownload db xftpServer1 86400
|
||||
show e `shouldContain` "ConversionFailed"
|
||||
@@ -783,6 +785,7 @@ testGetNextSndFileToPrepare st = do
|
||||
-- Right _ <- createSndFile db g 1 (CryptoFile "filepath" Nothing) 1 "filepath" testFileSbKey testFileCbNonce Nothing
|
||||
-- DB.execute_ db "UPDATE snd_files SET status = 'new', num_recipients = 'bad' WHERE snd_file_id = 1"
|
||||
Right fId2 <- createSndFile db g 1 (CryptoFile "filepath" Nothing) 1 "filepath" testFileSbKey testFileCbNonce Nothing Nothing
|
||||
Right _ <- startPreparedSndFile db fId2
|
||||
DB.execute_ db "UPDATE snd_files SET status = 'new' WHERE snd_file_id = 2"
|
||||
|
||||
-- Left e <- getNextSndFileToPrepare db 86400
|
||||
|
||||
@@ -71,9 +71,7 @@ initServers =
|
||||
}
|
||||
|
||||
testChooseDifferentOperator :: IO ()
|
||||
testChooseDifferentOperator = do
|
||||
c <- getSMPAgentClient' 1 agentCfg initServers testDB
|
||||
runRight_ $ do
|
||||
testChooseDifferentOperator = withAgent 1 agentCfg initServers testDB $ \c -> runRight_ $ do
|
||||
-- chooses the only operator with storage role
|
||||
srv1 <- withAgentEnv c $ getNextServer c 1 storageSrvs []
|
||||
liftIO $ srv1 == testOp1Srv1 || srv1 == testOp1Srv2 `shouldBe` True
|
||||
|
||||
+53
-1
@@ -34,7 +34,7 @@ import Simplex.FileTransfer.Server.Env (AFStoreType, XFTPServerConfig (..), defa
|
||||
import Simplex.FileTransfer.Server.Store (STMFileStore)
|
||||
import Simplex.FileTransfer.Transport (XFTPErrorType (AUTH))
|
||||
import Simplex.FileTransfer.Types (RcvFileId, SndFileId)
|
||||
import Simplex.Messaging.Agent (AgentClient, testProtocolServer, xftpDeleteRcvFile, xftpDeleteSndFileInternal, xftpDeleteSndFileRemote, xftpReceiveFile, xftpSendDescription, xftpStartWorkers)
|
||||
import Simplex.Messaging.Agent (AgentClient, testProtocolServer, xftpDeleteRcvFile, xftpDeleteSndFileInternal, xftpDeleteSndFileRemote, xftpPrepareReceiveFile, xftpPrepareSendFile, xftpReceiveFile, xftpSendDescription, xftpStartReceiveFile, xftpStartSendFile, xftpStartWorkers)
|
||||
import qualified Simplex.Messaging.Agent as XA
|
||||
import Simplex.Messaging.Agent.Client (ProtocolTestFailure (..), ProtocolTestStep (..))
|
||||
import Simplex.Messaging.Agent.Env.SQLite (AgentConfig, xftpCfg)
|
||||
@@ -84,6 +84,8 @@ xftpAgentTests =
|
||||
it "should send and receive with encrypted local files" testXFTPAgentSendReceiveEncrypted
|
||||
it "should send and receive large file with a redirect" testXFTPAgentSendReceiveRedirect
|
||||
it "should send and receive small file without a redirect" testXFTPAgentSendReceiveNoRedirect
|
||||
it "should receive prepared file only after it is started" testXFTPAgentPrepareReceive
|
||||
it "should send prepared file only after it is started" testXFTPAgentPrepareSend
|
||||
it "should extend storage time with an entitlement proof and report the granted expiry" $ \_ -> testXFTPAgentEntitlement
|
||||
describe "sending and receiving with version negotiation" $ beforeWith (const (pure ())) testXFTPAgentSendReceiveMatrix
|
||||
it "should resume receiving file after restart" $ \_ -> testXFTPAgentReceiveRestore
|
||||
@@ -278,6 +280,56 @@ testXFTPAgentSendReceiveNoRedirect = withXFTPServer $ do
|
||||
inBytes <- B.readFile filePathIn
|
||||
B.readFile out `shouldReturn` inBytes
|
||||
|
||||
testXFTPAgentPrepareReceive :: HasCallStack => AFStoreType -> IO ()
|
||||
testXFTPAgentPrepareReceive = withXFTPServer $ do
|
||||
filePath <- createRandomFile
|
||||
(_, _, rfd1, rfd2) <- withAgent 1 agentCfg initAgentServers testDB $ \sndr -> runRight $ testSend sndr filePath
|
||||
rfId2 <- withAgent 2 agentCfg initAgentServers testDB2 $ \rcp -> runRight $ do
|
||||
xftpStartWorkers rcp (Just recipientFiles)
|
||||
rfId1 <- xftpReceiveFile rcp 1 rfd1 Nothing True
|
||||
rfId2 <- xftpPrepareReceiveFile rcp 1 rfd2 Nothing True
|
||||
rfProgress rcp $ mb 18
|
||||
("", rfId1', RFDONE _) <- rfGet rcp
|
||||
liftIO $ do
|
||||
rfId1' `shouldBe` rfId1
|
||||
timeout 300000 (rfGet rcp) `shouldReturn` Nothing
|
||||
pure rfId2
|
||||
withAgent 3 agentCfg initAgentServers testDB2 $ \rcp' -> runRight_ $ do
|
||||
xftpStartWorkers rcp' (Just recipientFiles)
|
||||
liftIO $ timeout 300000 (rfGet rcp') `shouldReturn` Nothing
|
||||
xftpStartReceiveFile rcp' rfId2
|
||||
rfProgress rcp' $ mb 18
|
||||
("", rfId2', RFDONE path) <- rfGet rcp'
|
||||
liftIO $ do
|
||||
rfId2' `shouldBe` rfId2
|
||||
file <- B.readFile filePath
|
||||
B.readFile path `shouldReturn` file
|
||||
|
||||
testXFTPAgentPrepareSend :: HasCallStack => AFStoreType -> IO ()
|
||||
testXFTPAgentPrepareSend = withXFTPServer $ do
|
||||
filePath1 <- createRandomFile' "testfile1"
|
||||
filePath2 <- createRandomFile' "testfile2"
|
||||
sfId2 <- withAgent 1 agentCfg initAgentServers testDB $ \sndr -> runRight $ do
|
||||
xftpStartWorkers sndr (Just senderFiles)
|
||||
sfId1 <- xftpSendFile sndr 1 (CF.plain filePath1) 1
|
||||
sfId2 <- xftpPrepareSendFile sndr 1 (CF.plain filePath2) 1 Nothing
|
||||
sfProgress sndr $ mb 18
|
||||
("", sfId1', SFDONE _ _) <- sfGet sndr
|
||||
liftIO $ do
|
||||
sfId1' `shouldBe` sfId1
|
||||
timeout 300000 (sfGet sndr) `shouldReturn` Nothing
|
||||
pure sfId2
|
||||
rfd <- withAgent 2 agentCfg initAgentServers testDB $ \sndr' -> runRight $ do
|
||||
xftpStartWorkers sndr' (Just senderFiles)
|
||||
liftIO $ timeout 300000 (sfGet sndr') `shouldReturn` Nothing
|
||||
xftpStartSendFile sndr' sfId2
|
||||
sfProgress sndr' $ mb 18
|
||||
("", sfId2', SFDONE _ [rfd]) <- sfGet sndr'
|
||||
liftIO $ sfId2' `shouldBe` sfId2
|
||||
pure rfd
|
||||
withAgent 3 agentCfg initAgentServers testDB2 $ \rcp ->
|
||||
runRight_ . void $ testReceive rcp rfd filePath2
|
||||
|
||||
testXFTPAgentSendReceiveMatrix :: Spec
|
||||
testXFTPAgentSendReceiveMatrix = do
|
||||
describe "old server" $ do
|
||||
|
||||
Reference in New Issue
Block a user