Compare commits

...

19 Commits

Author SHA1 Message Date
Alexander Bondarenko d49244b654 add /debug subs 2024-04-30 22:00:57 +03:00
Alexander Bondarenko d0f21b943c debug conn by chat connection id from contact info 2024-04-30 20:10:25 +03:00
IC Rainbow ef5a144b7b Merge remote-tracking branch 'origin/master' into ab/debug-subs 2024-04-29 21:14:10 +03:00
Alexander Bondarenko 8d0d8208ce fix undefined error 2024-04-29 20:38:54 +03:00
Evgeny Poberezkin 83ab8d6ef2 more tracking 2024-04-28 18:27:43 +01:00
Evgeny Poberezkin 7d4d0dd2c0 Merge branch 'master' into ab/debug-subs 2024-04-28 09:09:30 +01:00
Alexander Bondarenko 982c4afd7b WIP: don't crash on missing corrId from acks 2024-04-26 22:52:43 +03:00
Alexander Bondarenko 24e718c7d5 one more /_stop 2024-04-26 22:05:51 +03:00
Alexander Bondarenko d3e6d638af fix late event in test 2024-04-26 21:17:47 +03:00
Alexander Bondarenko f6f23119cf don't omit nulls in debug objects 2024-04-26 18:57:11 +03:00
Alexander Bondarenko 33c36c8101 add Connection wholesale 2024-04-26 18:43:49 +03:00
Alexander Bondarenko 0622a52e00 fix withCompletedCommand 2024-04-26 18:25:03 +03:00
Evgeny Poberezkin 4634788ac0 Merge branch 'master' into ab/debug-subs 2024-04-26 15:47:07 +01:00
Alexander Bondarenko 40e717e7f4 read conn networkStatus 2024-04-26 17:39:19 +03:00
Alexander Bondarenko 97bdb9484b format 2024-04-26 17:25:33 +03:00
Alexander Bondarenko 3afc025a0d add connection debug 2024-04-26 17:10:38 +03:00
Alexander Bondarenko 19a3ab6230 fix DebugDelivery 2024-04-26 15:41:53 +03:00
Alexander Bondarenko ce5cb3137c collect some fields 2024-04-25 21:56:25 +03:00
Alexander Bondarenko 05d0554fef chat: delivery troubleshooting helper 2024-04-25 19:29:59 +03:00
5 changed files with 211 additions and 25 deletions
+157 -25
View File
@@ -46,6 +46,7 @@ import qualified Data.List.NonEmpty as L
import Data.Map.Strict (Map) import Data.Map.Strict (Map)
import qualified Data.Map.Strict as M import qualified Data.Map.Strict as M
import Data.Maybe (catMaybes, fromMaybe, isJust, isNothing, listToMaybe, mapMaybe, maybeToList) import Data.Maybe (catMaybes, fromMaybe, isJust, isNothing, listToMaybe, mapMaybe, maybeToList)
import qualified Data.Set as S
import Data.Text (Text) import Data.Text (Text)
import qualified Data.Text as T import qualified Data.Text as T
import Data.Text.Encoding (decodeLatin1, encodeUtf8) import Data.Text.Encoding (decodeLatin1, encodeUtf8)
@@ -91,14 +92,18 @@ import qualified Simplex.FileTransfer.Description as FD
import Simplex.FileTransfer.Protocol (FileParty (..), FilePartyI) import Simplex.FileTransfer.Protocol (FileParty (..), FilePartyI)
import Simplex.Messaging.Agent as Agent import Simplex.Messaging.Agent as Agent
import Simplex.Messaging.Agent.Client (AgentStatsKey (..), SubInfo (..), agentClientStore, getAgentWorkersDetails, getAgentWorkersSummary, temporaryAgentError, withLockMap) import Simplex.Messaging.Agent.Client (AgentStatsKey (..), SubInfo (..), agentClientStore, getAgentWorkersDetails, getAgentWorkersSummary, temporaryAgentError, withLockMap)
import qualified Simplex.Messaging.Agent.Client as AC
import Simplex.Messaging.Agent.Env.SQLite (AgentConfig (..), InitialAgentServers (..), createAgentStore, defaultAgentConfig) import Simplex.Messaging.Agent.Env.SQLite (AgentConfig (..), InitialAgentServers (..), createAgentStore, defaultAgentConfig)
import Simplex.Messaging.Agent.Lock (withLock) import Simplex.Messaging.Agent.Lock (withLock)
import Simplex.Messaging.Agent.Protocol import Simplex.Messaging.Agent.Protocol
import qualified Simplex.Messaging.Agent.Protocol as AP (AgentErrorType (..)) import qualified Simplex.Messaging.Agent.Protocol as AP (AgentErrorType (..))
import qualified Simplex.Messaging.Agent.Store as AS
import Simplex.Messaging.Agent.Store.SQLite (MigrationConfirmation (..), MigrationError, SQLiteStore (dbNew), execSQL, upMigration, withConnection) import Simplex.Messaging.Agent.Store.SQLite (MigrationConfirmation (..), MigrationError, SQLiteStore (dbNew), execSQL, upMigration, withConnection)
import qualified Simplex.Messaging.Agent.Store.SQLite as ADB
import Simplex.Messaging.Agent.Store.SQLite.DB (SlowQueryStats (..)) import Simplex.Messaging.Agent.Store.SQLite.DB (SlowQueryStats (..))
import qualified Simplex.Messaging.Agent.Store.SQLite.DB as DB import qualified Simplex.Messaging.Agent.Store.SQLite.DB as DB
import qualified Simplex.Messaging.Agent.Store.SQLite.Migrations as Migrations import qualified Simplex.Messaging.Agent.Store.SQLite.Migrations as Migrations
import qualified Simplex.Messaging.Agent.TRcvQueues as RQ
import Simplex.Messaging.Client (defaultNetworkConfig) import Simplex.Messaging.Client (defaultNetworkConfig)
import qualified Simplex.Messaging.Crypto as C import qualified Simplex.Messaging.Crypto as C
import Simplex.Messaging.Crypto.File (CryptoFile (..), CryptoFileArgs (..)) import Simplex.Messaging.Crypto.File (CryptoFile (..), CryptoFileArgs (..))
@@ -227,6 +232,8 @@ newChatController
inputQ <- newTBQueueIO tbqSize inputQ <- newTBQueueIO tbqSize
outputQ <- newTBQueueIO tbqSize outputQ <- newTBQueueIO tbqSize
connNetworkStatuses <- atomically TM.empty connNetworkStatuses <- atomically TM.empty
agentDeliveryStatuses <- atomically TM.empty
agentSubscriptions <- newTVarIO S.empty
subscriptionMode <- newTVarIO SMSubscribe subscriptionMode <- newTVarIO SMSubscribe
chatLock <- newEmptyTMVarIO chatLock <- newEmptyTMVarIO
entityLocks <- atomically TM.empty entityLocks <- atomically TM.empty
@@ -263,6 +270,8 @@ newChatController
inputQ, inputQ,
outputQ, outputQ,
connNetworkStatuses, connNetworkStatuses,
agentDeliveryStatuses,
agentSubscriptions,
subscriptionMode, subscriptionMode,
chatLock, chatLock,
entityLocks, entityLocks,
@@ -402,6 +411,7 @@ startChatController mainApp = do
subscribeUsers :: Bool -> [User] -> CM' () subscribeUsers :: Bool -> [User] -> CM' ()
subscribeUsers onlyNeeded users = do subscribeUsers onlyNeeded users = do
let (us, us') = partition activeUser users let (us, us') = partition activeUser users
asks agentSubscriptions >>= atomically . (`writeTVar` S.empty)
vr <- chatVersionRange' vr <- chatVersionRange'
subscribe vr us subscribe vr us
subscribe vr us' subscribe vr us'
@@ -2147,6 +2157,67 @@ processChatCommand' vr = \case
chatMigrations <- map upMigration <$> withStore' (Migrations.getCurrent . DB.conn) chatMigrations <- map upMigration <$> withStore' (Migrations.getCurrent . DB.conn)
agentMigrations <- withAgent getAgentMigrations agentMigrations <- withAgent getAgentMigrations
pure $ CRVersionInfo {versionInfo, chatMigrations, agentMigrations} pure $ CRVersionInfo {versionInfo, chatMigrations, agentMigrations}
DebugDelivery showAll -> do
ads <- mapM readTVarIO =<< readTVarIO =<< asks agentDeliveryStatuses
let collect (acId, ds) = if not showAll && agentDeliveryOk ds then Nothing else Just (decodeLatin1 $ strEncode acId, ds)
pure $ CRDebugDelivery . M.fromList . mapMaybe collect $ M.toList ads
DebugConnection (Left connId) -> withUser $ \user -> do
Connection {agentConnId} <- withStore $ \db -> getConnectionById db vr user connId
processChatCommand $ DebugConnection (Right agentConnId)
DebugConnection (Right acId@(AgentConnId acId')) -> do
user@User {userId} <- withStore' (`getUserByAConnId` acId) >>= maybe (throwError . ChatErrorStore $ SEUserNotFoundByAConnId acId) pure
AS.RcvQueue {server, rcvId = rcvId', status = rqStatus} <- withAgent $ \ac ->
-- dive into agent internals
liftIOEither . (`runReaderT` agentEnv ac) . runExceptT $ AC.withStore ac $ \adb ->
ADB.getPrimaryRcvQueue adb acId'
let tSess = (userId, server, Just acId')
SMP.ProtocolServer {host, port} = server
AgentClient {activeSubs, pendingSubs, smpClients, smpSubWorkers} <- withAgent pure
inActive <- any (\AS.RcvQueue {rcvId} -> rcvId == rcvId') <$> atomically (RQ.getSessQueues tSess activeSubs)
inPending <- any (\AS.RcvQueue {rcvId} -> rcvId == rcvId') <$> atomically (RQ.getSessQueues tSess pendingSubs)
smpClient <- atomically (TM.lookup (tSess $> Nothing) smpClients) >>= mapM (\AC.SessionVar {sessionVar} -> atomically (tryReadTMVar sessionVar))
smpClientIsolated <- atomically (TM.lookup tSess smpClients) >>= mapM (\AC.SessionVar {sessionVar} -> atomically (tryReadTMVar sessionVar))
let smpClientStatus = case smpClient <|> smpClientIsolated of
Nothing -> "missing"
Just Nothing -> "connecting"
Just (Just (Left err)) -> tshow err
Just (Just (Right _client)) -> "connected"
subWorker <- atomically (TM.lookup (tSess $> Nothing) smpSubWorkers) >>= mapM (\AC.SessionVar {sessionVar} -> atomically (tryReadTMVar sessionVar))
subWorkerIsolated <- atomically (TM.lookup tSess smpSubWorkers) >>= mapM (\AC.SessionVar {sessionVar} -> atomically (tryReadTMVar sessionVar))
let subWorkerStatus = case subWorker <|> subWorkerIsolated of
Nothing -> "idle"
Just Nothing -> "waiting"
Just (Just _async) -> "working"
connection <-
withStore (\db -> getConnectionEntity db vr user acId) >>= \case
RcvDirectMsgConnection {entityConnection} -> pure entityConnection
RcvGroupMsgConnection {entityConnection} -> pure entityConnection
SndFileConnection {entityConnection} -> pure entityConnection
RcvFileConnection {entityConnection} -> pure entityConnection
UserContactConnection {entityConnection} -> pure entityConnection
deliveryStatus <- mapM readTVarIO . M.lookup acId =<< readTVarIO =<< asks agentDeliveryStatuses
networkStatus <- M.lookup acId <$> chatReadVar connNetworkStatuses
pure $
CRDebugConnection
DebugConnectionStatus
{ deliveryStatus,
networkStatus,
inActive,
inPending,
server = (decodeLatin1 $ strEncode host, port),
smpClientStatus,
subWorkerStatus,
queueStatus = tshow rqStatus,
connection
}
DebugSubs -> do
ads <- mapM readTVarIO =<< readTVarIO =<< asks agentDeliveryStatuses
conns <- readTVarIO =<< asks agentSubscriptions
pure . CRDebugSubs . S.toList $ foldl' (\cs (acId, ds) -> if agentDeliveryOk ds then S.delete acId cs else cs) conns (M.toList ads)
DebugSubsDetails -> do
CRDebugSubs broken <- processChatCommand' vr DebugSubs
crs <- mapM (processChatCommand' vr . DebugConnection . Right) broken
pure $ CRDebugSubsDetails [dcs | CRDebugConnection dcs <- crs]
DebugLocks -> lift $ do DebugLocks -> lift $ do
chatLockName <- atomically . tryReadTMVar =<< asks chatLock chatLockName <- atomically . tryReadTMVar =<< asks chatLock
chatEntityLocks <- getLocks =<< asks entityLocks chatEntityLocks <- getLocks =<< asks entityLocks
@@ -2473,8 +2544,6 @@ processChatCommand' vr = \case
toView $ CRNewChatItem user (AChatItem SCTDirect SMDSnd (DirectChat ct) ci) toView $ CRNewChatItem user (AChatItem SCTDirect SMDSnd (DirectChat ct) ci)
forM_ (timed_ >>= timedDeleteAt') $ forM_ (timed_ >>= timedDeleteAt') $
startProximateTimedItemThread user (ChatRef CTDirect contactId, chatItemId' ci) startProximateTimedItemThread user (ChatRef CTDirect contactId, chatItemId' ci)
drgRandomBytes :: Int -> CM ByteString
drgRandomBytes n = asks random >>= atomically . C.randomBytes n
privateGetUser :: UserId -> CM User privateGetUser :: UserId -> CM User
privateGetUser userId = privateGetUser userId =
tryChatError (withStore (`getUser` userId)) >>= \case tryChatError (withStore (`getUser` userId)) >>= \case
@@ -3244,6 +3313,7 @@ subscribeUserConnections vr onlyNeeded agentBatchSubscribe user = do
let conns = concat [ctConns, ucConns, mConns, sftConns, rftConns, pcConns] let conns = concat [ctConns, ucConns, mConns, sftConns, rftConns, pcConns]
pure (conns, cts, ucs, gs, ms, sfts, rfts, pcs) pure (conns, cts, ucs, gs, ms, sfts, rfts, pcs)
-- subscribe using batched commands -- subscribe using batched commands
asks agentSubscriptions >>= \v -> atomically . modifyTVar' v $ \cs -> foldl' (\cs' acId -> S.insert (AgentConnId acId) cs') cs conns
rs <- withAgent $ \a -> agentBatchSubscribe a conns rs <- withAgent $ \a -> agentBatchSubscribe a conns
-- send connection events to view -- send connection events to view
contactSubsToView rs cts ce contactSubsToView rs cts ce
@@ -3533,13 +3603,29 @@ processAgentMessage _ connId (DEL_RCVQ srv qId err_) =
processAgentMessage _ connId DEL_CONN = processAgentMessage _ connId DEL_CONN =
toView $ CRAgentConnDeleted (AgentConnId connId) toView $ CRAgentConnDeleted (AgentConnId connId)
processAgentMessage corrId connId msg = do processAgentMessage corrId connId msg = do
lockEntity <- critical (withStore (`getChatLockEntity` AgentConnId connId)) let acId = AgentConnId connId
acTag = aCommandTag msg
when (acTag == MSG_ || acTag == RCVD_) . lift $ trackNewDelivery acId msg
lockEntity <- critical (withStore (`getChatLockEntity` acId))
withEntityLock "processAgentMessage" lockEntity $ do withEntityLock "processAgentMessage" lockEntity $ do
vr <- chatVersionRange vr <- chatVersionRange
-- getUserByAConnId never throws logical errors, only SEDBBusyError can be thrown here -- getUserByAConnId never throws logical errors, only SEDBBusyError can be thrown here
critical (withStore' (`getUserByAConnId` AgentConnId connId)) >>= \case critical (withStore' (`getUserByAConnId` acId)) >>= \case
Just user -> processAgentMessageConn vr user corrId connId msg `catchChatError` (toView . CRChatError (Just user)) Just user -> processAgentMessageConn vr user corrId connId msg `catchChatError` (toView . CRChatError (Just user))
_ -> throwChatError $ CENoConnectionUser (AgentConnId connId) _ -> throwChatError $ CENoConnectionUser acId
-- TODO: clean up deliveries
trackNewDelivery :: AgentConnId -> ACommand 'Agent 'AEConn -> CM' ()
trackNewDelivery acId msg = do
now <- liftIO getCurrentTime
asks agentDeliveryStatuses >>= atomically . TM.alterF (updateConn now) acId
where
(isMSG, msgBodyPfx) = case msg of
MSG _ _ msgBody -> (True, T.take 1000 $ safeDecodeUtf8 msgBody)
_ -> (False, "")
updateConn lastCmd = \case
Nothing -> Just <$> newTVar AgentDeliveryStatus {lastCmd, tracking = "create", isMSG, connId = Nothing, msgBodyPfx, eventTag = Nothing, ackSent = Nothing, pendingAcks = M.empty}
Just v -> Just v <$ modifyTVar' v (\AgentDeliveryStatus {pendingAcks} -> AgentDeliveryStatus {lastCmd, tracking = "create", isMSG, connId = Nothing, msgBodyPfx, eventTag = Nothing, ackSent = Nothing, pendingAcks = M.filter not pendingAcks})
-- CRITICAL error will be shown to the user as alert with restart button in Android/desktop apps. -- CRITICAL error will be shown to the user as alert with restart button in Android/desktop apps.
-- SEDBBusyError will only be thrown on IO exceptions or SQLError during DB queries, -- SEDBBusyError will only be thrown on IO exceptions or SQLError during DB queries,
@@ -3764,15 +3850,20 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
_ -> toView $ CRSubscriptionEnd user entity _ -> toView $ CRSubscriptionEnd user entity
MSGNTF smpMsgInfo -> toView $ CRNtfMessage user entity $ ntfMsgInfo smpMsgInfo MSGNTF smpMsgInfo -> toView $ CRNtfMessage user entity $ ntfMsgInfo smpMsgInfo
_ -> case entity of _ -> case entity of
RcvDirectMsgConnection conn contact_ -> RcvDirectMsgConnection conn contact_ -> do
storeDeliveryConn conn
processDirectMessage agentMessage entity conn contact_ processDirectMessage agentMessage entity conn contact_
RcvGroupMsgConnection conn gInfo m -> RcvGroupMsgConnection conn gInfo m -> do
storeDeliveryConn conn
processGroupMessage agentMessage entity conn gInfo m processGroupMessage agentMessage entity conn gInfo m
RcvFileConnection conn ft -> RcvFileConnection conn ft -> do
storeDeliveryConn conn
processRcvFileConn agentMessage entity conn ft processRcvFileConn agentMessage entity conn ft
SndFileConnection conn ft -> SndFileConnection conn ft -> do
storeDeliveryConn conn
processSndFileConn agentMessage entity conn ft processSndFileConn agentMessage entity conn ft
UserContactConnection conn uc -> UserContactConnection conn uc -> do
storeDeliveryConn conn
processUserContactRequest agentMessage entity conn uc processUserContactRequest agentMessage entity conn uc
where where
updateConnStatus :: ConnectionEntity -> CM ConnectionEntity updateConnStatus :: ConnectionEntity -> CM ConnectionEntity
@@ -3782,7 +3873,8 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
withStore' $ \db -> updateConnectionStatus db conn connStatus withStore' $ \db -> updateConnectionStatus db conn connStatus
pure $ updateEntityConnStatus acEntity connStatus pure $ updateEntityConnStatus acEntity connStatus
Nothing -> pure acEntity Nothing -> pure acEntity
storeDeliveryConn :: Connection -> CM ()
storeDeliveryConn Connection {connId} = lift $ agentDeliveryStatus (AgentConnId agentConnId) "storeDeliveryConn" $ \ad -> ad {connId = Just connId}
agentMsgConnStatus :: ACommand 'Agent e -> Maybe ConnStatus agentMsgConnStatus :: ACommand 'Agent e -> Maybe ConnStatus
agentMsgConnStatus = \case agentMsgConnStatus = \case
CONF {} -> Just ConnRequested CONF {} -> Just ConnRequested
@@ -3861,6 +3953,7 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
(ct', conn') <- updateContactPQRcv user ct conn pqEncryption (ct', conn') <- updateContactPQRcv user ct conn pqEncryption
checkIntegrityCreateItem (CDDirectRcv ct') msgMeta `catchChatError` \_ -> pure () checkIntegrityCreateItem (CDDirectRcv ct') msgMeta `catchChatError` \_ -> pure ()
(conn'', msg@RcvMessage {chatMsgEvent = ACME _ event}) <- saveDirectRcvMSG conn' msgMeta msgBody (conn'', msg@RcvMessage {chatMsgEvent = ACME _ event}) <- saveDirectRcvMSG conn' msgMeta msgBody
lift $ agentDeliveryStatus (AgentConnId agentConnId) "eventTag" $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag event)}
let ct'' = ct' {activeConn = Just conn''} :: Contact let ct'' = ct' {activeConn = Just conn''} :: Contact
assertDirectAllowed user MDRcv ct'' $ toCMEventTag event assertDirectAllowed user MDRcv ct'' $ toCMEventTag event
case event of case event of
@@ -4264,11 +4357,17 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
_ -> messageWarning "sendXGrpMemCon: member category GCPreMember or GCPostMember is expected" _ -> messageWarning "sendXGrpMemCon: member category GCPreMember or GCPostMember is expected"
MSG msgMeta _msgFlags msgBody -> do MSG msgMeta _msgFlags msgBody -> do
withAckMessage agentConnId msgMeta True $ do withAckMessage agentConnId msgMeta True $ do
lift $ agentDeliveryStatus (AgentConnId agentConnId) "group MSG" id
checkIntegrityCreateItem (CDGroupRcv gInfo m) msgMeta `catchChatError` \_ -> pure () checkIntegrityCreateItem (CDGroupRcv gInfo m) msgMeta `catchChatError` \_ -> pure ()
lift $ agentDeliveryStatus (AgentConnId agentConnId) "after checkIntegrityCreateItem" id
forM_ aChatMsgs $ \case forM_ aChatMsgs $ \case
Right (ACMsg _ chatMsg) -> Right (ACMsg _ chatMsg) -> do
lift $ agentDeliveryStatus (AgentConnId agentConnId) "has event" id
processEvent chatMsg `catchChatError` \e -> toView $ CRChatError (Just user) e processEvent chatMsg `catchChatError` \e -> toView $ CRChatError (Just user) e
Left e -> toView $ CRChatError (Just user) (ChatError . CEException $ "error parsing chat message: " <> e) Left e -> do
lift $ agentDeliveryStatus (AgentConnId agentConnId) "has error" id
toView $ CRChatError (Just user) (ChatError . CEException $ "error parsing chat message: " <> e)
lift $ agentDeliveryStatus (AgentConnId agentConnId) "after forM_" id
forwardMsg_ `catchChatError` \_ -> pure () forwardMsg_ `catchChatError` \_ -> pure ()
checkSendRcpt $ rights aChatMsgs checkSendRcpt $ rights aChatMsgs
where where
@@ -4276,7 +4375,9 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
brokerTs = metaBrokerTs msgMeta brokerTs = metaBrokerTs msgMeta
processEvent :: MsgEncodingI e => ChatMessage e -> CM () processEvent :: MsgEncodingI e => ChatMessage e -> CM ()
processEvent chatMsg = do processEvent chatMsg = do
lift $ agentDeliveryStatus (AgentConnId agentConnId) "processEvent" id
(m', conn', msg@RcvMessage {chatMsgEvent = ACME _ event}) <- saveGroupRcvMsg user groupId m conn msgMeta msgBody chatMsg (m', conn', msg@RcvMessage {chatMsgEvent = ACME _ event}) <- saveGroupRcvMsg user groupId m conn msgMeta msgBody chatMsg
lift $ agentDeliveryStatus (AgentConnId agentConnId) "eventTag in processEvent" $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag event)}
case event of case event of
XMsgNew mc -> memberCanSend m' $ newGroupContentMessage gInfo m' mc msg brokerTs False XMsgNew mc -> memberCanSend m' $ newGroupContentMessage gInfo m' mc msg brokerTs False
XMsgFileDescr sharedMsgId fileDescr -> memberCanSend m' $ groupMessageFileDescription gInfo m' sharedMsgId fileDescr XMsgFileDescr sharedMsgId fileDescr -> memberCanSend m' $ groupMessageFileDescription gInfo m' sharedMsgId fileDescr
@@ -4620,16 +4721,22 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
withCompletedCommand :: forall e. AEntityI e => Connection -> ACommand 'Agent e -> (CommandData -> CM ()) -> CM () withCompletedCommand :: forall e. AEntityI e => Connection -> ACommand 'Agent e -> (CommandData -> CM ()) -> CM ()
withCompletedCommand Connection {connId} agentMsg action = do withCompletedCommand Connection {connId} agentMsg action = do
let agentMsgTag = APCT (sAEntity @e) $ aCommandTag agentMsg let agentMsgTag = APCT (sAEntity @e) $ aCommandTag agentMsg
cmdData_ <- withStore' $ \db -> getCommandDataByCorrId db user corrId pending_ <- mapM readTVarIO =<< atomically . TM.lookup acId =<< asks agentDeliveryStatuses
case cmdData_ of if agentMsgTag == APCT SAEConn OK_ && corrId /= "" && maybe False (M.member ackKey . pendingAcks) pending_
Just cmdData@CommandData {cmdId, cmdConnId = Just cmdConnId', cmdFunction} then lift $ agentDeliveryStatus acId "withCompletedCommand" $ \ad@AgentDeliveryStatus {pendingAcks} -> ad {pendingAcks = M.adjust (const True) ackKey pendingAcks}
| connId == cmdConnId' && (agentMsgTag == commandExpectedResponse cmdFunction || agentMsgTag == APCT SAEConn ERR_) -> do else do
withStore' $ \db -> deleteCommand db user cmdId cmdData_ <- withStore' $ \db -> getCommandDataByCorrId db user corrId
action cmdData case cmdData_ of
| otherwise -> err cmdId $ "not matching connection id or unexpected response, corrId = " <> show corrId Just cmdData@CommandData {cmdId, cmdConnId = Just cmdConnId', cmdFunction}
Just CommandData {cmdId, cmdConnId = Nothing} -> err cmdId $ "no command connection id, corrId = " <> show corrId | connId == cmdConnId' && (agentMsgTag == commandExpectedResponse cmdFunction || agentMsgTag == APCT SAEConn ERR_) -> do
Nothing -> throwChatError . CEAgentCommandError $ "command not found, corrId = " <> show corrId withStore' $ \db -> deleteCommand db user cmdId
action cmdData
| otherwise -> logWarn $ "not matching connection id or unexpected response, corrId = " <> tshow corrId
Just CommandData {cmdId, cmdConnId = Nothing} -> logWarn $ "no command connection id, corrId = " <> tshow corrId
Nothing -> logWarn $ "command not found, corrId = " <> tshow corrId
where where
acId = AgentConnId agentConnId
ackKey = decodeLatin1 $ strEncode corrId
err cmdId msg = do err cmdId msg = do
withStore' $ \db -> updateCommandStatus db user cmdId CSError withStore' $ \db -> updateCommandStatus db user cmdId CSError
throwChatError . CEAgentCommandError $ msg throwChatError . CEAgentCommandError $ msg
@@ -4650,11 +4757,20 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
Right withRcpt -> ackMsg msgMeta $ if withRcpt then Just "" else Nothing Right withRcpt -> ackMsg msgMeta $ if withRcpt then Just "" else Nothing
-- If showCritical is True, then these errors don't result in ACK and show user visible alert -- If showCritical is True, then these errors don't result in ACK and show user visible alert
-- This prevents losing the message that failed to be processed. -- This prevents losing the message that failed to be processed.
Left (ChatErrorStore SEDBBusyError {message}) | showCritical -> throwError $ ChatErrorAgent (CRITICAL True message) Nothing Left (ChatErrorStore SEDBBusyError {message}) | showCritical -> do
lift $ agentDeliveryStatus (AgentConnId agentConnId) "SEDBBusyError" id
throwError $ ChatErrorAgent (CRITICAL True message) Nothing
Left e -> ackMsg msgMeta Nothing >> throwError e Left e -> ackMsg msgMeta Nothing >> throwError e
where where
ackMsg :: MsgMeta -> Maybe MsgReceiptInfo -> CM () ackMsg :: MsgMeta -> Maybe MsgReceiptInfo -> CM ()
ackMsg MsgMeta {recipient = (msgId, _)} rcpt = withAgent $ \a -> ackMessageAsync a "" cId msgId rcpt ackMsg MsgMeta {recipient = (msgId, _)} rcpt = do
ackCorrId <- drgRandomBytes 24
lift $ agentDeliveryStatus (AgentConnId agentConnId) "ackMsg" id
withAgent $ \a -> ackMessageAsync a ackCorrId cId msgId rcpt
now <- liftIO getCurrentTime
let ackKey = decodeLatin1 $ strEncode ackCorrId
lift . agentDeliveryStatus (AgentConnId agentConnId) "after ackMessageAsync" $ \ad@AgentDeliveryStatus {pendingAcks} ->
ad {ackSent = Just (now, ackKey), pendingAcks = M.insert ackKey False pendingAcks}
sentMsgDeliveryEvent :: Connection -> AgentMsgId -> CM () sentMsgDeliveryEvent :: Connection -> AgentMsgId -> CM ()
sentMsgDeliveryEvent Connection {connId} msgId = sentMsgDeliveryEvent Connection {connId} msgId =
@@ -5953,6 +6069,7 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
processForwardedMsg author chatMsg = do processForwardedMsg author chatMsg = do
let body = LB.toStrict $ J.encode msg let body = LB.toStrict $ J.encode msg
rcvMsg@RcvMessage {chatMsgEvent = ACME _ event} <- saveGroupFwdRcvMsg user groupId m author body chatMsg rcvMsg@RcvMessage {chatMsgEvent = ACME _ event} <- saveGroupFwdRcvMsg user groupId m author body chatMsg
lift $ agentDeliveryStatus (AgentConnId agentConnId) "eventTag in processForwardedMsg" $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag event)}
case event of case event of
XMsgNew mc -> memberCanSend author $ newGroupContentMessage gInfo author mc rcvMsg msgTs True XMsgNew mc -> memberCanSend author $ newGroupContentMessage gInfo author mc rcvMsg msgTs True
XMsgFileDescr sharedMsgId fileDescr -> memberCanSend author $ groupMessageFileDescription gInfo author sharedMsgId fileDescr XMsgFileDescr sharedMsgId fileDescr -> memberCanSend author $ groupMessageFileDescription gInfo author sharedMsgId fileDescr
@@ -6577,9 +6694,10 @@ sendPendingGroupMessages user GroupMember {groupMemberId, localDisplayName} conn
-- TODO [batch send] refactor direct message processing same as groups (e.g. checkIntegrity before processing) -- TODO [batch send] refactor direct message processing same as groups (e.g. checkIntegrity before processing)
saveDirectRcvMSG :: Connection -> MsgMeta -> MsgBody -> CM (Connection, RcvMessage) saveDirectRcvMSG :: Connection -> MsgMeta -> MsgBody -> CM (Connection, RcvMessage)
saveDirectRcvMSG conn@Connection {connId} agentMsgMeta msgBody = saveDirectRcvMSG conn@Connection {connId, agentConnId} agentMsgMeta msgBody =
case parseChatMessages msgBody of case parseChatMessages msgBody of
[Right (ACMsg _ ChatMessage {chatVRange, msgId = sharedMsgId_, chatMsgEvent})] -> do [Right (ACMsg _ ChatMessage {chatVRange, msgId = sharedMsgId_, chatMsgEvent})] -> do
lift $ agentDeliveryStatus agentConnId "eventTag in saveDirectRcvMSG" $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag chatMsgEvent)}
conn' <- updatePeerChatVRange conn chatVRange conn' <- updatePeerChatVRange conn chatVRange
let agentMsgId = fst $ recipient agentMsgMeta let agentMsgId = fst $ recipient agentMsgMeta
newMsg = NewRcvMessage {chatMsgEvent, msgBody} newMsg = NewRcvMessage {chatMsgEvent, msgBody}
@@ -6744,6 +6862,7 @@ deleteAgentConnectionAsync user acId = deleteAgentConnectionAsync' user acId Fal
deleteAgentConnectionAsync' :: User -> ConnId -> Bool -> CM () deleteAgentConnectionAsync' :: User -> ConnId -> Bool -> CM ()
deleteAgentConnectionAsync' user acId waitDelivery = do deleteAgentConnectionAsync' user acId waitDelivery = do
withAgent (\a -> deleteConnectionAsync a waitDelivery acId) `catchChatError` (toView . CRChatError (Just user)) withAgent (\a -> deleteConnectionAsync a waitDelivery acId) `catchChatError` (toView . CRChatError (Just user))
asks agentSubscriptions >>= \v -> atomically $ modifyTVar v $ S.delete (AgentConnId acId)
deleteAgentConnectionsAsync :: User -> [ConnId] -> CM () deleteAgentConnectionsAsync :: User -> [ConnId] -> CM ()
deleteAgentConnectionsAsync user acIds = deleteAgentConnectionsAsync' user acIds False deleteAgentConnectionsAsync user acIds = deleteAgentConnectionsAsync' user acIds False
@@ -6752,6 +6871,7 @@ deleteAgentConnectionsAsync' :: User -> [ConnId] -> Bool -> CM ()
deleteAgentConnectionsAsync' _ [] _ = pure () deleteAgentConnectionsAsync' _ [] _ = pure ()
deleteAgentConnectionsAsync' user acIds waitDelivery = do deleteAgentConnectionsAsync' user acIds waitDelivery = do
withAgent (\a -> deleteConnectionsAsync a waitDelivery acIds) `catchChatError` (toView . CRChatError (Just user)) withAgent (\a -> deleteConnectionsAsync a waitDelivery acIds) `catchChatError` (toView . CRChatError (Just user))
asks agentSubscriptions >>= \v -> atomically $ modifyTVar v $ \as -> foldl' (\as' acId -> S.delete (AgentConnId acId) as') as acIds
agentXFTPDeleteRcvFile :: RcvFileId -> FileTransferId -> CM () agentXFTPDeleteRcvFile :: RcvFileId -> FileTransferId -> CM ()
agentXFTPDeleteRcvFile aFileId fileId = do agentXFTPDeleteRcvFile aFileId fileId = do
@@ -7282,6 +7402,10 @@ chatCommandP =
"/_download " *> (APIDownloadStandaloneFile <$> A.decimal <* A.space <*> strP_ <*> cryptoFileP), "/_download " *> (APIDownloadStandaloneFile <$> A.decimal <* A.space <*> strP_ <*> cryptoFileP),
("/quit" <|> "/q" <|> "/exit") $> QuitChat, ("/quit" <|> "/q" <|> "/exit") $> QuitChat,
("/version" <|> "/v") $> ShowVersion, ("/version" <|> "/v") $> ShowVersion,
"/debug delivery" *> (DebugDelivery <$> (" all" $> True <|> pure False)),
"/debug conn " *> (DebugConnection <$> ((Left <$> A.decimal) <|> (Right . AgentConnId <$> base64P))),
"/debug subs" $> DebugSubs,
"/debug subs details" $> DebugSubsDetails,
"/debug locks" $> DebugLocks, "/debug locks" $> DebugLocks,
"/debug event " *> (DebugEvent <$> jsonP), "/debug event " *> (DebugEvent <$> jsonP),
"/get stats" $> GetAgentStats, "/get stats" $> GetAgentStats,
@@ -7443,6 +7567,9 @@ timeItToView s action = do
toView' $ CRTimedAction s diff toView' $ CRTimedAction s diff
pure a pure a
drgRandomBytes :: Int -> CM ByteString
drgRandomBytes n = asks random >>= atomically . C.randomBytes n
mkValidName :: String -> String mkValidName :: String -> String
mkValidName = reverse . dropWhile isSpace . fst3 . foldl' addChar ("", '\NUL', 0 :: Int) mkValidName = reverse . dropWhile isSpace . fst3 . foldl' addChar ("", '\NUL', 0 :: Int)
where where
@@ -7487,3 +7614,8 @@ xftpSndFileRedirect user ftId vfd = do
dummyFileDescr :: FileDescr dummyFileDescr :: FileDescr
dummyFileDescr = FileDescr {fileDescrText = "", fileDescrPartNo = 0, fileDescrComplete = False} dummyFileDescr = FileDescr {fileDescrText = "", fileDescrPartNo = 0, fileDescrComplete = False}
agentDeliveryStatus :: AgentConnId -> Text -> (AgentDeliveryStatus -> AgentDeliveryStatus) -> CM' ()
agentDeliveryStatus acId tracking f = do
ads <- asks agentDeliveryStatuses
atomically $ TM.lookup acId ads >>= mapM_ (`modifyTVar'` \ds -> (f ds) {tracking})
+47
View File
@@ -39,6 +39,7 @@ import Data.Int (Int64)
import Data.List.NonEmpty (NonEmpty) import Data.List.NonEmpty (NonEmpty)
import Data.Map.Strict (Map) import Data.Map.Strict (Map)
import qualified Data.Map.Strict as M import qualified Data.Map.Strict as M
import Data.Set (Set)
import Data.String import Data.String
import Data.Text (Text) import Data.Text (Text)
import Data.Text.Encoding (decodeLatin1) import Data.Text.Encoding (decodeLatin1)
@@ -207,6 +208,8 @@ data ChatController = ChatController
inputQ :: TBQueue String, inputQ :: TBQueue String,
outputQ :: TBQueue (Maybe CorrId, Maybe RemoteHostId, ChatResponse), outputQ :: TBQueue (Maybe CorrId, Maybe RemoteHostId, ChatResponse),
connNetworkStatuses :: TMap AgentConnId NetworkStatus, connNetworkStatuses :: TMap AgentConnId NetworkStatus,
agentDeliveryStatuses :: TMap AgentConnId (TVar AgentDeliveryStatus),
agentSubscriptions :: TVar (Set AgentConnId),
subscriptionMode :: TVar SubscriptionMode, subscriptionMode :: TVar SubscriptionMode,
chatLock :: Lock, chatLock :: Lock,
entityLocks :: TMap ChatLockEntity Lock, entityLocks :: TMap ChatLockEntity Lock,
@@ -233,6 +236,23 @@ data ChatController = ChatController
contactMergeEnabled :: TVar Bool contactMergeEnabled :: TVar Bool
} }
data AgentDeliveryStatus = AgentDeliveryStatus
{ lastCmd :: UTCTime,
tracking :: Text,
isMSG :: Bool, -- False for RCVD
connId :: Maybe Int64, -- chat connection ID
msgBodyPfx :: Text,
eventTag :: Maybe Text, -- tshow of ACMEventTag (for JSON instances)
ackSent :: Maybe (UTCTime, Text), -- strEncode of random CorrId
pendingAcks :: Map Text Bool
}
deriving (Show) -- for ChatResponse
agentDeliveryOk :: AgentDeliveryStatus -> Bool
agentDeliveryOk AgentDeliveryStatus {ackSent, pendingAcks} = case ackSent of
Nothing -> False
Just (_, corrId) -> M.lookup corrId pendingAcks == Just True && and pendingAcks
data HelpSection = HSMain | HSFiles | HSGroups | HSContacts | HSMyAddress | HSIncognito | HSMarkdown | HSMessages | HSRemote | HSSettings | HSDatabase data HelpSection = HSMain | HSFiles | HSGroups | HSContacts | HSMyAddress | HSIncognito | HSMarkdown | HSMessages | HSRemote | HSSettings | HSDatabase
deriving (Show) deriving (Show)
@@ -488,6 +508,10 @@ data ChatCommand
| APIStandaloneFileInfo FileDescriptionURI | APIStandaloneFileInfo FileDescriptionURI
| QuitChat | QuitChat
| ShowVersion | ShowVersion
| DebugDelivery Bool
| DebugConnection (Either Int64 AgentConnId)
| DebugSubs -- Check that all `subscribeUserConnections` are in place and healthy
| DebugSubsDetails -- Run DebugConnection on each unhealthy sub
| DebugLocks | DebugLocks
| DebugEvent ChatResponse | DebugEvent ChatResponse
| GetAgentStats | GetAgentStats
@@ -735,6 +759,10 @@ data ChatResponse
| CRContactPQEnabled {user :: User, contact :: Contact, pqEnabled :: PQEncryption} | CRContactPQEnabled {user :: User, contact :: Contact, pqEnabled :: PQEncryption}
| CRSQLResult {rows :: [Text]} | CRSQLResult {rows :: [Text]}
| CRSlowSQLQueries {chatQueries :: [SlowSQLQuery], agentQueries :: [SlowSQLQuery]} | CRSlowSQLQueries {chatQueries :: [SlowSQLQuery], agentQueries :: [SlowSQLQuery]}
| CRDebugDelivery {debugDelivery :: Map Text AgentDeliveryStatus}
| CRDebugConnection {debugConnection :: DebugConnectionStatus}
| CRDebugSubs {brokenSubs :: [AgentConnId]}
| CRDebugSubsDetails {debugConnections :: [DebugConnectionStatus]}
| CRDebugLocks {chatLockName :: Maybe String, chatEntityLocks :: Map String String, agentLocks :: AgentLocks} | CRDebugLocks {chatLockName :: Maybe String, chatEntityLocks :: Map String String, agentLocks :: AgentLocks}
| CRAgentStats {agentStats :: [[String]]} | CRAgentStats {agentStats :: [[String]]}
| CRAgentWorkersDetails {agentWorkersDetails :: AgentWorkersDetails} | CRAgentWorkersDetails {agentWorkersDetails :: AgentWorkersDetails}
@@ -755,6 +783,21 @@ data ChatResponse
| CRCustomChatResponse {user_ :: Maybe User, response :: Text} | CRCustomChatResponse {user_ :: Maybe User, response :: Text}
deriving (Show) deriving (Show)
data DebugConnectionStatus = DebugConnectionStatus
{ deliveryStatus :: Maybe AgentDeliveryStatus,
networkStatus :: Maybe NetworkStatus,
-- from agent's TRecvQ via rcvId
inActive :: Bool, -- should the delivery work right now?
inPending :: Bool, -- is there a temporary error?
-- from receive queue
queueStatus :: Text,
server :: (Text, String), -- what's the server for this connection? -- XXX: reveals private servers and association
smpClientStatus :: Text, -- is there an active client for it?
subWorkerStatus :: Text, -- a session was recently restarted and tries to resubscribe
connection :: Connection
}
deriving (Show)
-- some of these can only be used as command responses -- some of these can only be used as command responses
allowRemoteEvent :: ChatResponse -> Bool allowRemoteEvent :: ChatResponse -> Bool
allowRemoteEvent = \case allowRemoteEvent = \case
@@ -1456,6 +1499,10 @@ $(JQ.deriveJSON (sumTypeJSON $ dropPrefix "RCSR") ''RemoteCtrlStopReason)
$(JQ.deriveJSON (sumTypeJSON $ dropPrefix "RHSR") ''RemoteHostStopReason) $(JQ.deriveJSON (sumTypeJSON $ dropPrefix "RHSR") ''RemoteHostStopReason)
$(JQ.deriveJSON J.defaultOptions ''AgentDeliveryStatus)
$(JQ.deriveJSON J.defaultOptions ''DebugConnectionStatus)
$(JQ.deriveJSON (sumTypeJSON $ dropPrefix "CR") ''ChatResponse) $(JQ.deriveJSON (sumTypeJSON $ dropPrefix "CR") ''ChatResponse)
$(JQ.deriveFromJSON defaultJSON ''ArchiveConfig) $(JQ.deriveFromJSON defaultJSON ''ArchiveConfig)
+1
View File
@@ -59,6 +59,7 @@ data StoreError
= SEDuplicateName = SEDuplicateName
| SEUserNotFound {userId :: UserId} | SEUserNotFound {userId :: UserId}
| SEUserNotFoundByName {contactName :: ContactName} | SEUserNotFoundByName {contactName :: ContactName}
| SEUserNotFoundByAConnId {agentConnId :: AgentConnId}
| SEUserNotFoundByContactId {contactId :: ContactId} | SEUserNotFoundByContactId {contactId :: ContactId}
| SEUserNotFoundByGroupId {groupId :: GroupId} | SEUserNotFoundByGroupId {groupId :: GroupId}
| SEUserNotFoundByFileId {fileId :: FileTransferId} | SEUserNotFoundByFileId {fileId :: FileTransferId}
+4
View File
@@ -351,6 +351,10 @@ responseToView hu@(currentRH, user_) ChatConfig {logLevel, showReactions, showRe
<> (" :: avg: " <> sShow timeAvg <> " ms") <> (" :: avg: " <> sShow timeAvg <> " ms")
<> (" :: " <> plain (T.unwords $ T.lines query)) <> (" :: " <> plain (T.unwords $ T.lines query))
in ("Chat queries" : map viewQuery chatQueries) <> [""] <> ("Agent queries" : map viewQuery agentQueries) in ("Chat queries" : map viewQuery chatQueries) <> [""] <> ("Agent queries" : map viewQuery agentQueries)
CRDebugDelivery ads -> [plain $ LB.unpack (J.encode ads)]
CRDebugConnection cs -> [plain $ LB.unpack (J.encode cs)]
CRDebugSubs cs -> [plain $ LB.unpack (J.encode cs)]
CRDebugSubsDetails cs -> [plain $ LB.unpack (J.encode cs)]
CRDebugLocks {chatLockName, chatEntityLocks, agentLocks} -> CRDebugLocks {chatLockName, chatEntityLocks, agentLocks} ->
[ maybe "no chat lock" (("chat lock: " <>) . plain) chatLockName, [ maybe "no chat lock" (("chat lock: " <>) . plain) chatLockName,
plain $ "chat entity locks: " <> LB.unpack (J.encode chatEntityLocks), plain $ "chat entity locks: " <> LB.unpack (J.encode chatEntityLocks),
+2
View File
@@ -1069,6 +1069,7 @@ testMaintenanceMode tmp = do
bob <# "alice> hi again" bob <# "alice> hi again"
bob #> "@alice hello" bob #> "@alice hello"
alice <# "bob> hello" alice <# "bob> hello"
threadDelay 100000
-- export / delete / import -- export / delete / import
alice ##> "/_stop" alice ##> "/_stop"
alice <## "chat stopped" alice <## "chat stopped"
@@ -1149,6 +1150,7 @@ testDatabaseEncryption tmp = do
alice <## "error: chat not stopped" alice <## "error: chat not stopped"
alice ##> "/db decrypt mykey" alice ##> "/db decrypt mykey"
alice <## "error: chat not stopped" alice <## "error: chat not stopped"
threadDelay 100000
alice ##> "/_stop" alice ##> "/_stop"
alice <## "chat stopped" alice <## "chat stopped"
alice ##> "/db decrypt mykey" alice ##> "/db decrypt mykey"