mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-09-01 18:08:36 +00:00
save, restore stats wip
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user