diff --git a/simplexmq.cabal b/simplexmq.cabal index 20f6d016b..8068fafd1 100644 --- a/simplexmq.cabal +++ b/simplexmq.cabal @@ -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 diff --git a/src/Simplex/Messaging/Server.hs b/src/Simplex/Messaging/Server.hs index 0e12bc6be..5aae3c8bd 100644 --- a/src/Simplex/Messaging/Server.hs +++ b/src/Simplex/Messaging/Server.hs @@ -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 diff --git a/src/Simplex/Messaging/Server/DataLog.hs b/src/Simplex/Messaging/Server/DataLog.hs new file mode 100644 index 000000000..7972de444 --- /dev/null +++ b/src/Simplex/Messaging/Server/DataLog.hs @@ -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 diff --git a/src/Simplex/Messaging/Server/Env/STM.hs b/src/Simplex/Messaging/Server/Env/STM.hs index 01e738ecc..4533ed885 100644 --- a/src/Simplex/Messaging/Server/Env/STM.hs +++ b/src/Simplex/Messaging/Server/Env/STM.hs @@ -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 diff --git a/src/Simplex/Messaging/Server/Main.hs b/src/Simplex/Messaging/Server/Main.hs index 784d0504a..90dfe750f 100644 --- a/src/Simplex/Messaging/Server/Main.hs +++ b/src/Simplex/Messaging/Server/Main.hs @@ -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 diff --git a/tests/SMPClient.hs b/tests/SMPClient.hs index 736016b3b..83b309348 100644 --- a/tests/SMPClient.hs +++ b/tests/SMPClient.hs @@ -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' = diff --git a/tests/ServerTests.hs b/tests/ServerTests.hs index 4325fe8ae..2282d4992 100644 --- a/tests/ServerTests.hs +++ b/tests/ServerTests.hs @@ -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="