From f892c661c1efa4dc2cb12da18126d5dd55446995 Mon Sep 17 00:00:00 2001 From: spaced4ndy <8711996+spaced4ndy@users.noreply.github.com> Date: Fri, 14 Jun 2024 14:34:23 +0400 Subject: [PATCH] save, restore stats wip --- src/Simplex/Messaging/Agent.hs | 39 +++++++- src/Simplex/Messaging/Agent/Env/SQLite.hs | 6 +- src/Simplex/Messaging/Agent/Stats.hs | 109 ++++++++++++++++++++-- 3 files changed, 144 insertions(+), 10 deletions(-) diff --git a/src/Simplex/Messaging/Agent.hs b/src/Simplex/Messaging/Agent.hs index c550ba04a..e4e33deae 100644 --- a/src/Simplex/Messaging/Agent.hs +++ b/src/Simplex/Messaging/Agent.hs @@ -126,6 +126,7 @@ import qualified Data.Aeson as J import Data.Bifunctor (bimap, first, second) import Data.ByteString.Char8 (ByteString) import qualified Data.ByteString.Char8 as B +import qualified Data.ByteString.Lazy as LB import Data.Composition ((.:), (.:.), (.::), (.::.)) import Data.Either (isRight, rights) import Data.Foldable (foldl', toList) @@ -156,6 +157,7 @@ import Simplex.Messaging.Agent.Lock (withLock, withLock') import Simplex.Messaging.Agent.NtfSubSupervisor import Simplex.Messaging.Agent.Protocol import Simplex.Messaging.Agent.RetryInterval +import Simplex.Messaging.Agent.Stats import Simplex.Messaging.Agent.Store import Simplex.Messaging.Agent.Store.SQLite import qualified Simplex.Messaging.Agent.Store.SQLite.DB as DB @@ -182,6 +184,7 @@ import Simplex.Messaging.Version import Simplex.RemoteControl.Client import Simplex.RemoteControl.Invitation import Simplex.RemoteControl.Types +import System.Directory import System.Mem.Weak (deRefWeak) import UnliftIO.Concurrent (forkFinally, forkIO, killThread, mkWeakThreadId, threadDelay) import qualified UnliftIO.Exception as E @@ -207,17 +210,51 @@ getSMPAgentClient_ clientId cfg initServers store backgroundMode = pure c runAgentThreads c | backgroundMode = run c "subscriber" $ subscriber c - | otherwise = + | otherwise = do + restoreAgentStats c raceAny_ [ run c "subscriber" $ subscriber c, run c "runNtfSupervisor" $ runNtfSupervisor c, run c "cleanupManager" $ cleanupManager c ] + `E.finally` saveAgentStats c run AgentClient {subQ, acThread} name a = a `E.catchAny` \e -> whenM (isJust <$> readTVarIO acThread) $ do logError $ "Agent thread " <> name <> " crashed: " <> tshow e atomically $ writeTBQueue subQ ("", "", AEvt SAEConn $ ERR $ CRITICAL True $ show e) +saveAgentStats :: AgentClient -> AM' () +saveAgentStats AgentClient {smpServersStats, xftpServersStats} = + asks (agentStatsBackupFile . config) >>= mapM_ (liftIO . saveStats) + where + saveStats f = do + sss <- readTVarIO smpServersStats + smpServersStatsData <- mapM (atomically . getAgentSMPServerStats) sss + xss <- readTVarIO xftpServersStats + xftpServersStatsData <- mapM (atomically . getAgentXFTPServerStats) xss + let stats = AgentPersistedServerStats {smpServersStatsData, xftpServersStatsData} + logInfo $ "saving agent stats to file " <> T.pack f + B.writeFile f $ LB.toStrict $ J.encode stats + logInfo "agent stats saved" + +restoreAgentStats :: AgentClient -> AM' () +restoreAgentStats AgentClient {smpServersStats, xftpServersStats} = + asks (agentStatsBackupFile . config) >>= mapM_ (liftIO . restoreStats) + where + restoreStats f = whenM (doesFileExist f) $ do + logInfo $ "restoring agent stats from file " <> T.pack f + liftIO (J.decode . LB.fromStrict <$> B.readFile f) >>= \case + Just AgentPersistedServerStats {smpServersStatsData, xftpServersStatsData} -> do + sss <- mapM (atomically . newAgentSMPServerStats') smpServersStatsData + atomically $ writeTVar smpServersStats sss + xss <- mapM (atomically . newAgentXFTPServerStats') xftpServersStatsData + atomically $ writeTVar xftpServersStats xss + renameFile f $ f <> ".bak" + logInfo "server stats restored" + Nothing -> do + logInfo $ "error restoring server stats" + renameFile f $ f <> ".bak" + disconnectAgentClient :: AgentClient -> IO () disconnectAgentClient c@AgentClient {agentEnv = Env {ntfSupervisor = ns, xftpAgent = xa}} = do closeAgentClient c diff --git a/src/Simplex/Messaging/Agent/Env/SQLite.hs b/src/Simplex/Messaging/Agent/Env/SQLite.hs index ee2bb16cc..54605482b 100644 --- a/src/Simplex/Messaging/Agent/Env/SQLite.hs +++ b/src/Simplex/Messaging/Agent/Env/SQLite.hs @@ -118,7 +118,8 @@ data AgentConfig = AgentConfig certificateFile :: FilePath, e2eEncryptVRange :: VersionRangeE2E, smpAgentVRange :: VersionRangeSMPA, - smpClientVRange :: VersionRangeSMPC + smpClientVRange :: VersionRangeSMPC, + agentStatsBackupFile :: Maybe FilePath } defaultReconnectInterval :: RetryInterval @@ -190,7 +191,8 @@ defaultAgentConfig = certificateFile = "/etc/opt/simplex-agent/agent.crt", e2eEncryptVRange = supportedE2EEncryptVRange, smpAgentVRange = supportedSMPAgentVRange, - smpClientVRange = supportedSMPClientVRange + smpClientVRange = supportedSMPClientVRange, + agentStatsBackupFile = Nothing } data Env = Env diff --git a/src/Simplex/Messaging/Agent/Stats.hs b/src/Simplex/Messaging/Agent/Stats.hs index f48334df0..aec117223 100644 --- a/src/Simplex/Messaging/Agent/Stats.hs +++ b/src/Simplex/Messaging/Agent/Stats.hs @@ -5,17 +5,13 @@ module Simplex.Messaging.Agent.Stats where import qualified Data.Aeson.TH as J +import Data.Map (Map) import Data.Time.Clock (UTCTime (..)) +import Simplex.Messaging.Agent.Protocol (UserId) import Simplex.Messaging.Parsers (defaultJSON) +import Simplex.Messaging.Protocol (SMPServer, XFTPServer) import UnliftIO.STM --- Add single type for gathering both smp and xftp stats across all users and servers? --- Idea is to provide agent with a single path to save/restore stats, instead of managing it in UI --- and providing agent with "initial stats". Agent would use this unifying type to write/read --- json representation of all stats and populating AgentClient maps: --- smpServersStats :: TMap (UserId, SMPServer) AgentSMPServerStats, --- xftpServersStats :: TMap (UserId, XFTPServer) AgentXFTPServerStats - data AgentSMPServerStats = AgentSMPServerStats { fromTime :: TVar UTCTime, msgSent :: TVar Int, -- total messages sent to server directly @@ -113,6 +109,54 @@ newAgentSMPServerStats ts = do subErr } +newAgentSMPServerStats' :: AgentSMPServerStatsData -> STM AgentSMPServerStats +newAgentSMPServerStats' s@AgentSMPServerStatsData {_fromTime} = do + fromTime <- newTVar _fromTime + msgSent <- newTVar $ _msgSent s + msgSentRetries <- newTVar $ _msgSentRetries s + msgSentSuccesses <- newTVar $ _msgSentSuccesses s + msgSentAuth <- newTVar $ _msgSentAuth s + msgSentQuota <- newTVar $ _msgSentQuota s + msgSentExpired <- newTVar $ _msgSentExpired s + msgSentErr <- newTVar $ _msgSentErr s + msgProx <- newTVar $ _msgProx s + msgProxRetries <- newTVar $ _msgProxRetries s + msgProxSuccesses <- newTVar $ _msgProxSuccesses s + msgProxExpired <- newTVar $ _msgProxExpired s + msgProxAuth <- newTVar $ _msgProxAuth s + msgProxQuota <- newTVar $ _msgProxQuota s + msgProxErr <- newTVar $ _msgProxErr s + msgRecv <- newTVar $ _msgRecv s + msgRecvDuplicate <- newTVar $ _msgRecvDuplicate s + msgRecvErr <- newTVar $ _msgRecvErr s + sub <- newTVar $ _sub s + subRetries <- newTVar $ _subRetries s + subErr <- newTVar $ _subErr s + pure + AgentSMPServerStats + { fromTime, + msgSent, + msgSentRetries, + msgSentSuccesses, + msgSentAuth, + msgSentQuota, + msgSentExpired, + msgSentErr, + msgProx, + msgProxRetries, + msgProxSuccesses, + msgProxExpired, + msgProxAuth, + msgProxQuota, + msgProxErr, + msgRecv, + msgRecvDuplicate, + msgRecvErr, + sub, + subRetries, + subErr + } + getAgentSMPServerStats :: AgentSMPServerStats -> STM AgentSMPServerStatsData getAgentSMPServerStats s@AgentSMPServerStats {fromTime} = do _fromTime <- readTVar fromTime @@ -251,6 +295,40 @@ newAgentXFTPServerStats ts = do replDeleteErr } +newAgentXFTPServerStats' :: AgentXFTPServerStatsData -> STM AgentXFTPServerStats +newAgentXFTPServerStats' s@AgentXFTPServerStatsData {_fromTime} = do + fromTime <- newTVar _fromTime + replUpload <- newTVar $ _replUpload s + replUploadRetries <- newTVar $ _replUploadRetries s + replUploadSuccesses <- newTVar $ _replUploadSuccesses s + replUploadErr <- newTVar $ _replUploadErr s + replDownload <- newTVar $ _replDownload s + replDownloadRetries <- newTVar $ _replDownloadRetries s + replDownloadSuccesses <- newTVar $ _replDownloadSuccesses s + replDownloadAuth <- newTVar $ _replDownloadAuth s + replDownloadErr <- newTVar $ _replDownloadErr s + replDelete <- newTVar $ _replDelete s + replDeleteRetries <- newTVar $ _replDeleteRetries s + replDeleteSuccesses <- newTVar $ _replDeleteSuccesses s + replDeleteErr <- newTVar $ _replDeleteErr s + pure + AgentXFTPServerStats + { fromTime, + replUpload, + replUploadRetries, + replUploadSuccesses, + replUploadErr, + replDownload, + replDownloadRetries, + replDownloadSuccesses, + replDownloadAuth, + replDownloadErr, + replDelete, + replDeleteRetries, + replDeleteSuccesses, + replDeleteErr + } + getAgentXFTPServerStats :: AgentXFTPServerStats -> STM AgentXFTPServerStatsData getAgentXFTPServerStats s@AgentXFTPServerStats {fromTime} = do _fromTime <- readTVar fromTime @@ -302,6 +380,23 @@ setAgentXFTPServerStats s@AgentXFTPServerStats {fromTime} d@AgentXFTPServerStats writeTVar (replDeleteSuccesses s) $! _replDeleteSuccesses d writeTVar (replDeleteErr s) $! _replDeleteErr d +-- Type for gathering both smp and xftp stats across all users and servers. +-- +-- Idea is to provide agent with a single path to save/restore stats, instead of managing it in UI +-- and providing agent with "initial stats". +-- +-- Agent would use this unifying type to write/read json representation of all stats +-- and populating AgentClient maps: +-- smpServersStats :: TMap (UserId, SMPServer) AgentSMPServerStats, +-- xftpServersStats :: TMap (UserId, XFTPServer) AgentXFTPServerStats +data AgentPersistedServerStats = AgentPersistedServerStats + { smpServersStatsData :: Map (UserId, SMPServer) AgentSMPServerStatsData, + xftpServersStatsData :: Map (UserId, XFTPServer) AgentXFTPServerStatsData + } + deriving (Show) + $(J.deriveJSON defaultJSON ''AgentSMPServerStatsData) $(J.deriveJSON defaultJSON ''AgentXFTPServerStatsData) + +$(J.deriveJSON defaultJSON ''AgentPersistedServerStats)