smp server: persistence for data blobs (#1245)

This commit is contained in:
Evgeny Poberezkin
2024-07-25 21:46:47 +01:00
committed by GitHub
parent eee8c0ba78
commit fd009fe0d9
7 changed files with 154 additions and 10 deletions
+1
View File
@@ -167,6 +167,7 @@ library
Simplex.Messaging.Server
Simplex.Messaging.Server.CLI
Simplex.Messaging.Server.Control
Simplex.Messaging.Server.DataLog
Simplex.Messaging.Server.DataStore
Simplex.Messaging.Server.Env.STM
Simplex.Messaging.Server.Expiration
+19 -6
View File
@@ -83,6 +83,7 @@ import Simplex.Messaging.Encoding
import Simplex.Messaging.Encoding.String
import Simplex.Messaging.Protocol
import Simplex.Messaging.Server.Control
import Simplex.Messaging.Server.DataLog
import Simplex.Messaging.Server.DataStore
import Simplex.Messaging.Server.Env.STM as Env
import Simplex.Messaging.Server.Expiration
@@ -156,7 +157,11 @@ smpServer started cfg@ServerConfig {transports, transportConfig = tCfg} = do
fromTLSCredentials (_, pk) = C.x509ToPrivate (pk, []) >>= C.privKey
saveServer :: Bool -> M ()
saveServer keepMsgs = withLog closeStoreLog >> saveServerMessages keepMsgs >> saveServerStats
saveServer keepMsgs = do
withLog closeStoreLog
withLog' dataLog closeStoreLog
saveServerMessages keepMsgs
saveServerStats
closeServer :: M ()
closeServer = asks (smpAgent . proxyAgent) >>= liftIO . closeSMPClientAgent
@@ -1436,12 +1441,18 @@ client thParams' clnt@Client {subscriptions, ntfSubscriptions, rcvQ, sndQ, sessi
storeDataBlob :: DataPublicAuthKey -> DataBlob -> M (Transmission BrokerMsg)
storeDataBlob dataKey dataBlob
| B.length (dataBody dataBlob) > e2eEncMessageLength = pure $ err LARGE_MSG
| otherwise = ok <$ (atomically . TM.insert entId d =<< asks dataStore)
| otherwise = do
atomically . TM.insert entId d =<< asks dataStore
withLog' dataLog (`logCreateBlob` d)
pure ok
where
d = DataRec {dataId = entId, dataKey, dataBlob}
deleteDataBlob :: M (Transmission BrokerMsg)
deleteDataBlob = ok <$ (atomically . TM.delete entId =<< asks dataStore)
deleteDataBlob = do
atomically . TM.delete entId =<< asks dataStore
withLog' dataLog (`logDeleteBlob` entId)
pure ok
getDataBlob :: M (Transmission BrokerMsg)
getDataBlob = case vRes of
@@ -1482,9 +1493,11 @@ incStat v = atomically $ modifyTVar' v (+ 1)
{-# INLINE incStat #-}
withLog :: (StoreLog 'WriteMode -> IO a) -> M ()
withLog action = do
env <- ask
liftIO . mapM_ action $ storeLog (env :: Env)
withLog = withLog' storeLog
{-# INLINE withLog #-}
withLog' :: (Env -> Maybe (StoreLog 'WriteMode)) -> (StoreLog 'WriteMode -> IO a) -> M ()
withLog' sel action = liftIO . mapM_ action =<< asks sel
timed :: T.Text -> RecipientId -> M a -> M a
timed name qId a = do
+62
View File
@@ -0,0 +1,62 @@
{-# LANGUAGE DataKinds #-}
{-# LANGUAGE LambdaCase #-}
{-# LANGUAGE OverloadedStrings #-}
module Simplex.Messaging.Server.DataLog where
import Control.Applicative ((<|>))
import Control.Monad (foldM)
import qualified Data.ByteString.Char8 as B
import qualified Data.ByteString.Lazy.Char8 as LB
import Data.Map.Strict (Map)
import qualified Data.Map.Strict as M
import Simplex.Messaging.Protocol (BlobId)
import Simplex.Messaging.Server.DataStore
import Simplex.Messaging.Server.StoreLog
import Simplex.Messaging.Encoding.String
import Simplex.Messaging.Transport.Buffer (trimCR)
import Simplex.Messaging.Util (ifM)
import System.Directory (doesFileExist)
import System.IO
data DataLogRecord = CreateBlob DataRec | DeleteBlob BlobId
instance StrEncoding DataLogRecord where
strEncode = \case
CreateBlob d -> strEncode (Str "CREATE", d)
DeleteBlob dId -> strEncode (Str "DELETE", dId)
strP =
"CREATE " *> (CreateBlob <$> strP)
<|> "DELETE " *> (DeleteBlob <$> strP)
logCreateBlob :: StoreLog 'WriteMode -> DataRec -> IO ()
logCreateBlob s = writeStoreLogRecord s . CreateBlob
logDeleteBlob :: StoreLog 'WriteMode -> BlobId -> IO ()
logDeleteBlob s = writeStoreLogRecord s . DeleteBlob
readWriteDataLog :: FilePath -> IO (Map BlobId DataRec, StoreLog 'WriteMode)
readWriteDataLog f = do
ds <- ifM (doesFileExist f) (readDataBlobs f) (pure M.empty)
s <- openWriteStoreLog f
writeDataBlobs s ds
pure (ds, s)
writeDataBlobs :: StoreLog 'WriteMode -> Map BlobId DataRec -> IO ()
writeDataBlobs = mapM_ . logCreateBlob
readDataBlobs :: FilePath -> IO (Map BlobId DataRec)
readDataBlobs f = foldM processLine M.empty . LB.lines =<< LB.readFile f
where
processLine :: Map BlobId DataRec -> LB.ByteString -> IO (Map BlobId DataRec)
processLine m s' = case strDecode $ trimCR s of
Right r -> pure $ procLogRecord r
Left e -> m <$ printError e
where
s = LB.toStrict s'
procLogRecord :: DataLogRecord -> Map BlobId DataRec
procLogRecord = \case
CreateBlob d -> M.insert (dataId d) d m
DeleteBlob dId -> M.delete dId m
printError :: String -> IO ()
printError e = B.putStrLn $ "Error parsing log: " <> B.pack e <> " - " <> s
+14 -2
View File
@@ -32,6 +32,7 @@ import Simplex.Messaging.Client.Agent (SMPClientAgent, SMPClientAgentConfig, new
import Simplex.Messaging.Crypto (KeyHash (..))
import qualified Simplex.Messaging.Crypto as C
import Simplex.Messaging.Protocol
import Simplex.Messaging.Server.DataLog
import Simplex.Messaging.Server.DataStore
import Simplex.Messaging.Server.Expiration
import Simplex.Messaging.Server.Information
@@ -56,6 +57,7 @@ data ServerConfig = ServerConfig
queueIdBytes :: Int,
msgIdBytes :: Int,
storeLogFile :: Maybe FilePath,
dataLogFile :: Maybe FilePath,
storeMsgsFile :: Maybe FilePath,
-- | set to False to prohibit creating new queues
allowNewQueues :: Bool,
@@ -126,6 +128,7 @@ data Env = Env
dataStore :: TMap BlobId DataRec,
random :: TVar ChaChaDRG,
storeLog :: Maybe (StoreLog 'WriteMode),
dataLog :: Maybe (StoreLog 'WriteMode),
tlsServerParams :: T.ServerParams,
serverStats :: ServerStats,
sockets :: SocketState,
@@ -216,7 +219,7 @@ newProhibitedSub = do
return Sub {subThread = ProhibitSub, delivered}
newEnv :: ServerConfig -> IO Env
newEnv config@ServerConfig {caCertificateFile, certificateFile, privateKeyFile, storeLogFile, smpAgentCfg, transportConfig, information, messageExpiration} = do
newEnv config@ServerConfig {caCertificateFile, certificateFile, privateKeyFile, storeLogFile, dataLogFile, smpAgentCfg, transportConfig, information, messageExpiration} = do
server <- atomically newServer
queueStore <- atomically newQueueStore
msgStore <- atomically newMsgStore
@@ -226,6 +229,10 @@ newEnv config@ServerConfig {caCertificateFile, certificateFile, privateKeyFile,
forM storeLogFile $ \f -> do
logInfo $ "restoring queues from file " <> T.pack f
restoreQueues queueStore f
dataLog <-
forM dataLogFile $ \f -> do
logInfo $ "restoring data blobs from file " <> T.pack f
restoreDataBlobs dataStore f
tlsServerParams <- loadTLSServerParams caCertificateFile certificateFile privateKeyFile (alpn transportConfig)
Fingerprint fp <- loadFingerprint caCertificateFile
let serverIdentity = KeyHash fp
@@ -234,7 +241,7 @@ newEnv config@ServerConfig {caCertificateFile, certificateFile, privateKeyFile,
clientSeq <- newTVarIO 0
clients <- newTVarIO mempty
proxyAgent <- atomically $ newSMPProxyAgent smpAgentCfg random
pure Env {config, serverInfo, server, serverIdentity, queueStore, msgStore, dataStore, random, storeLog, tlsServerParams, serverStats, sockets, clientSeq, clients, proxyAgent}
pure Env {config, serverInfo, server, serverIdentity, queueStore, msgStore, dataStore, random, storeLog, dataLog, tlsServerParams, serverStats, sockets, clientSeq, clients, proxyAgent}
where
restoreQueues :: QueueStore -> FilePath -> IO (StoreLog 'WriteMode)
restoreQueues QueueStore {queues, senders, notifiers} f = do
@@ -244,6 +251,11 @@ newEnv config@ServerConfig {caCertificateFile, certificateFile, privateKeyFile,
writeTVar senders $! M.foldr' addSender M.empty qs
writeTVar notifiers $! M.foldr' addNotifier M.empty qs
pure s
restoreDataBlobs :: TMap BlobId DataRec -> FilePath -> IO (StoreLog 'WriteMode)
restoreDataBlobs dataStore f = do
(ds, s) <- readWriteDataLog f
atomically $ writeTVar dataStore ds
pure s
addSender :: QueueRec -> Map SenderId RecipientId -> Map SenderId RecipientId
addSender q = M.insert (senderId q) (recipientId q)
addNotifier :: QueueRec -> Map NotifierId RecipientId -> Map NotifierId RecipientId
+2
View File
@@ -79,6 +79,7 @@ smpServerCLI_ generateSite serveStaticFiles cfgPath logPath =
defaultServerPort = "5223"
executableName = "smp-server"
storeLogFilePath = combine logPath "smp-server-store.log"
dataLogFilePath = combine logPath "smp-server-data.log"
httpsCertFile = combine cfgPath "web.cert"
httpsKeyFile = combine cfgPath "web.key"
defaultStaticPath = combine logPath "www"
@@ -262,6 +263,7 @@ smpServerCLI_ generateSite serveStaticFiles cfgPath logPath =
privateKeyFile = c serverKeyFile,
certificateFile = c serverCrtFile,
storeLogFile = enableStoreLog $> storeLogFilePath,
dataLogFile = enableStoreLog $> dataLogFilePath,
storeMsgsFile =
let messagesPath = combine logPath "smp-server-messages.log"
in case iniOnOff "STORE_LOG" "restore_messages" ini of
+5 -1
View File
@@ -57,6 +57,9 @@ testStoreLogFile = "tests/tmp/smp-server-store.log"
testStoreLogFile2 :: FilePath
testStoreLogFile2 = "tests/tmp/smp-server-store.log.2"
testDataLogFile :: FilePath
testDataLogFile = "tests/tmp/smp-server-data.log"
testStoreMsgsFile :: FilePath
testStoreMsgsFile = "tests/tmp/smp-server-messages.log"
@@ -104,6 +107,7 @@ cfg =
queueIdBytes = 24,
msgIdBytes = 24,
storeLogFile = Nothing,
dataLogFile = Nothing,
storeMsgsFile = Nothing,
allowNewQueues = True,
newQueueBasicAuth = Nothing,
@@ -158,7 +162,7 @@ withSmpServerStoreMsgLogOn :: HasCallStack => ATransport -> ServiceName -> (HasC
withSmpServerStoreMsgLogOn t = withSmpServerConfigOn t cfg {storeLogFile = Just testStoreLogFile, storeMsgsFile = Just testStoreMsgsFile, serverStatsBackupFile = Just testServerStatsBackupFile}
withSmpServerStoreLogOn :: HasCallStack => ATransport -> ServiceName -> (HasCallStack => ThreadId -> IO a) -> IO a
withSmpServerStoreLogOn t = withSmpServerConfigOn t cfg {storeLogFile = Just testStoreLogFile, serverStatsBackupFile = Just testServerStatsBackupFile}
withSmpServerStoreLogOn t = withSmpServerConfigOn t cfg {storeLogFile = Just testStoreLogFile, dataLogFile = Just testDataLogFile, serverStatsBackupFile = Just testServerStatsBackupFile}
withSmpServerConfigOn :: HasCallStack => ATransport -> ServerConfig -> ServiceName -> (HasCallStack => ThreadId -> IO a) -> IO a
withSmpServerConfigOn t cfg' port' =
+51 -1
View File
@@ -71,7 +71,9 @@ serverTests t@(ATransport t') = do
testMsgExpireOnSend t'
testMsgExpireOnInterval t'
testMsgNOTExpireOnInterval t'
describe "Data blobs" $ testDataBlobs t'
describe "Data blobs" $ do
testDataBlobs t'
testDataBlobsWithLog t
pattern Resp :: CorrId -> QueueId -> BrokerMsg -> SignedTransmission ErrorType BrokerMsg
pattern Resp corrId queueId command <- (_, _, (corrId, queueId, Right command))
@@ -991,6 +993,54 @@ testDataBlobs t =
Resp "11" _ (ERR AUTH) <- sendRecv s ("", "11", sBlobId, READ)
pure ()
testDataBlobsWithLog :: ATransport -> Spec
testDataBlobsWithLog at@(ATransport t) =
it "should store data blob to log and restore after server restart" $ do
g <- C.newRandom
(C.PublicKeyX25519 k, pk) <- atomically $ C.generateKeyPair @'C.X25519 g
(blobKey, blobPKey) <- atomically $ C.generateAuthKeyPair C.SEd25519 g
dataNonce <- atomically $ C.randomCbNonce g
let kBytes = BA.convert k :: ByteString
rBlobId = C.sha256Hash kBytes
sBlobId = kBytes
blob = DataBlob {dataNonce, dataBody = "some random encrypted data"} -- the previous test shows e2e blob encryption
blob2 = DataBlob {dataNonce, dataBody = "some other encrypted data"}
clientServer t $ \h -> do
Resp "1" _ OK <- signSendRecv h blobPKey ("1", rBlobId, WRT blobKey blob)
pure ()
clientServer t $ \h -> do
testGetBlob h "2" pk sBlobId blob
-- update blob
Resp "3" _ OK <- signSendRecv h blobPKey ("3", rBlobId, WRT blobKey blob2)
testGetBlob h "4" pk sBlobId blob2
clientServer t $ \h -> do
-- updated after restart
testGetBlob h "5" pk sBlobId blob2
-- delete blob
Resp "6" _ OK <- signSendRecv h blobPKey ("6", rBlobId, CLR)
Resp "7" _ (ERR AUTH) <- sendRecv h ("", "7", sBlobId, READ)
pure ()
clientServer t $ \h -> do
-- deleted after restart
Resp "8" _ (ERR AUTH) <- sendRecv h ("", "8", sBlobId, READ)
pure ()
where
clientServer :: Transport c => TProxy c -> (THandleSMP c 'TClient -> IO ()) -> IO ()
clientServer _ test' =
withSmpServerStoreLogOn at testPort $ \server -> do
testSMPClient test' `shouldReturn` ()
killThread server
testGetBlob h corrId pk sBlobId expectedBlob = do
Resp (CorrId corrId') _ (DATA encBlob) <- sendRecv h ("", corrId, sBlobId, READ)
corrId' `shouldBe` corrId
THandle {params = THandleParams {thAuth = Just THAuthClient {serverPeerPubKey}}} <- pure h
let ss = C.dh' serverPeerPubKey pk
respNonce = C.cbNonce corrId -- correlation ID sent in READ request
Right blobStr <- pure $ C.cbDecrypt ss respNonce encBlob
Right blob' <- pure $ smpDecode blobStr
blob' `shouldBe` expectedBlob
samplePubKey :: C.APublicVerifyKey
samplePubKey = C.APublicVerifyKey C.SEd25519 "MCowBQYDK2VwAyEAfAOflyvbJv1fszgzkQ6buiZJVgSpQWsucXq7U6zjMgY="