fix DebugDelivery

This commit is contained in:
Alexander Bondarenko
2024-04-26 15:41:53 +03:00
parent ce5cb3137c
commit 19a3ab6230
3 changed files with 84 additions and 112 deletions
+62 -101
View File
@@ -231,7 +231,7 @@ newChatController
inputQ <- newTBQueueIO tbqSize
outputQ <- newTBQueueIO tbqSize
connNetworkStatuses <- atomically TM.empty
agentConnStatuses <- atomically TM.empty
agentDeliveryStatuses <- atomically TM.empty
subscriptionMode <- newTVarIO SMSubscribe
chatLock <- newEmptyTMVarIO
entityLocks <- atomically TM.empty
@@ -268,7 +268,7 @@ newChatController
inputQ,
outputQ,
connNetworkStatuses,
agentConnStatuses,
agentDeliveryStatuses,
subscriptionMode,
chatLock,
entityLocks,
@@ -2153,55 +2153,10 @@ processChatCommand' vr = \case
chatMigrations <- map upMigration <$> withStore' (Migrations.getCurrent . DB.conn)
agentMigrations <- withAgent getAgentMigrations
pure $ CRVersionInfo {versionInfo, chatMigrations, agentMigrations}
DebugAcks -> do
acs <- mapM readTVarIO =<< readTVarIO =<< asks agentConnStatuses
fmap (CRDebugAcks . M.fromList) . forM (M.toList acs) $ \(acId@(AgentConnId acId'), agentConnStatus) -> do
user_ <- withStore' (`getUserByAConnId` acId)
rq' <- withAgent $ \ac -> liftIO . (`runReaderT` agentEnv ac) . runExceptT $ do
AS.RcvQueue {server, rcvId, status, smpClientVersion} <- AC.withStore ac (`ADB.getPrimaryRcvQueue` acId')
pure (server, (rcvId, status, smpClientVersion))
(inActive, inPending) <- case (user_, rq') of
(Just User {userId}, Right (srv, (rcvId', _, _))) -> do
let tSess = (userId, srv, Just acId')
withAgent $ \AgentClient {activeSubs, pendingSubs} -> do
active <- atomically (RQ.getSessQueues tSess activeSubs)
pending <- atomically (RQ.getSessQueues tSess pendingSubs)
pure
( any (\AS.RcvQueue {rcvId} -> rcvId == rcvId') active,
any (\AS.RcvQueue {rcvId} -> rcvId == rcvId') pending
)
_ -> pure (False, False)
entity_ <- forM user_ $ \user -> withStore (\db -> getConnectionEntity db vr user acId)
conn_ <- forM entity_ $ \case
RcvDirectMsgConnection {entityConnection} -> pure entityConnection
RcvGroupMsgConnection {entityConnection} -> pure entityConnection
SndFileConnection {entityConnection} -> pure entityConnection
RcvFileConnection {entityConnection} -> pure entityConnection
UserContactConnection {entityConnection} -> pure entityConnection
let AgentConnStatus {lastCmd, lastCmdTag, lastMsg, ackSent, okRcvd} = agentConnStatus
pure
( decodeLatin1 $ strEncode acId,
DebugAck
{ lastCmd = Just (decodeLatin1 $ strEncode lastCmdTag, lastCmd),
lastMsg,
lastAck = ackSent,
lastOK = okRcvd,
inActive,
inPending,
server = either (const Nothing) (\(SMP.ProtocolServer {host, port}, _) -> Just (decodeLatin1 . strEncode $ L.head host, port)) rq',
hasSMPClient = True, -- TODO
hasSubWorker = False, -- TODO
hasDeliveryWorker = False, -- TODO
connStatus_ = (\Connection {connStatus} -> connStatus) <$> conn_,
-- agent connstatus?
connAuthErrors = (\c@Connection {authErrCounter} -> (authErrCounter, connDisabled c)) <$> conn_,
createdAt = (\Connection {createdAt} -> createdAt) <$> conn_
}
)
DebugDelivery -> do
ads <- mapM readTVarIO =<< readTVarIO =<< asks agentDeliveryStatuses
let collect (acId, ds) = if agentDeliveryOk ds then Nothing else Just (decodeLatin1 $ strEncode acId, ds)
pure $ CRDebugDelivery . M.fromList . mapMaybe collect $ M.toList ads
DebugLocks -> lift $ do
chatLockName <- atomically . tryReadTMVar =<< asks chatLock
chatEntityLocks <- getLocks =<< asks entityLocks
@@ -2528,8 +2483,6 @@ processChatCommand' vr = \case
toView $ CRNewChatItem user (AChatItem SCTDirect SMDSnd (DirectChat ct) ci)
forM_ (timed_ >>= timedDeleteAt') $
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 =
tryChatError (withStore (`getUser` userId)) >>= \case
@@ -3583,17 +3536,12 @@ expireChatItems user@User {userId} ttl sync = do
forM_ membersToDelete $ \m -> withStore' $ \db -> deleteGroupMember db user m
processAgentMessage :: ACorrId -> ConnId -> ACommand 'Agent 'AEConn -> CM ()
processAgentMessage _ connId (DEL_RCVQ srv qId err_) = do
let acId = AgentConnId connId
asks agentConnStatuses >>= atomically . TM.delete acId
toView $ CRAgentRcvQueueDeleted acId srv (AgentQueueId qId) err_
processAgentMessage _ connId DEL_CONN = do
let acId = AgentConnId connId
asks agentConnStatuses >>= atomically . TM.delete acId
toView $ CRAgentConnDeleted acId
processAgentMessage _ connId (DEL_RCVQ srv qId err_) = toView $ CRAgentRcvQueueDeleted (AgentConnId connId) srv (AgentQueueId qId) err_
processAgentMessage _ connId DEL_CONN = toView $ CRAgentConnDeleted (AgentConnId connId)
processAgentMessage corrId connId msg = do
let acId = AgentConnId connId
lift $ trackAgentConn acId msg
acTag = aCommandTag msg
when (acTag == MSG_ || acTag == RCVD_) . lift $ trackNewDelivery acId (acTag == MSG_)
lockEntity <- critical (withStore (`getChatLockEntity` acId))
withEntityLock "processAgentMessage" lockEntity $ do
vr <- chatVersionRange
@@ -3602,24 +3550,15 @@ processAgentMessage corrId connId msg = do
Just user -> processAgentMessageConn vr user corrId connId msg `catchChatError` (toView . CRChatError (Just user))
_ -> throwChatError $ CENoConnectionUser acId
trackAgentConn :: AgentConnId -> ACommand 'Agent 'AEConn -> CM' ()
trackAgentConn acId msg = do
-- TODO: clean up deliveries
trackNewDelivery :: AgentConnId -> Bool -> CM' ()
trackNewDelivery acId isMSG = do
now <- liftIO getCurrentTime
asks agentConnStatuses >>= atomically . TM.alterF (updateConn now) acId
asks agentDeliveryStatuses >>= atomically . TM.alterF (updateConn now) acId
where
updateConn now = \case
Nothing -> Just <$> newTVar (status now Nothing Nothing Nothing)
Just v -> Just v <$ modifyTVar' v (\AgentConnStatus {lastMsg, ackSent, okRcvd} -> status now lastMsg ackSent okRcvd)
status now lastMsg ackSent okRcvd = AgentConnStatus
{ lastCmd = now,
lastCmdTag,
lastMsg = if isMSG then Just now else lastMsg,
ackSent,
okRcvd = if isOK then Just now else okRcvd
}
lastCmdTag = aCommandTag msg
isMSG = lastCmdTag == MSG_
isOK = lastCmdTag == OK_
updateConn lastCmd = \case
Nothing -> Just <$> newTVar AgentDeliveryStatus {lastCmd, isMSG, connId = Nothing, eventTag = Nothing,ackSent = Nothing, pendingAcks = M.empty}
Just v -> Just v <$ modifyTVar' v (\AgentDeliveryStatus {pendingAcks} -> AgentDeliveryStatus {lastCmd, isMSG, connId = Nothing, 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.
-- SEDBBusyError will only be thrown on IO exceptions or SQLError during DB queries,
@@ -3834,15 +3773,20 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
_ -> toView $ CRSubscriptionEnd user entity
MSGNTF smpMsgInfo -> toView $ CRNtfMessage user entity $ ntfMsgInfo smpMsgInfo
_ -> case entity of
RcvDirectMsgConnection conn contact_ ->
RcvDirectMsgConnection conn contact_ -> do
storeDeliveryConn conn
processDirectMessage agentMessage entity conn contact_
RcvGroupMsgConnection conn gInfo m ->
RcvGroupMsgConnection conn gInfo m -> do
storeDeliveryConn conn
processGroupMessage agentMessage entity conn gInfo m
RcvFileConnection conn ft ->
RcvFileConnection conn ft -> do
storeDeliveryConn conn
processRcvFileConn agentMessage entity conn ft
SndFileConnection conn ft ->
SndFileConnection conn ft -> do
storeDeliveryConn conn
processSndFileConn agentMessage entity conn ft
UserContactConnection conn uc ->
UserContactConnection conn uc -> do
storeDeliveryConn conn
processUserContactRequest agentMessage entity conn uc
where
updateConnStatus :: ConnectionEntity -> CM ConnectionEntity
@@ -3852,7 +3796,8 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
withStore' $ \db -> updateConnectionStatus db conn connStatus
pure $ updateEntityConnStatus acEntity connStatus
Nothing -> pure acEntity
storeDeliveryConn :: Connection -> CM ()
storeDeliveryConn Connection {connId} = lift $ agentDeliveryStatus (AgentConnId agentConnId) $ \ad -> ad {connId = Just connId}
agentMsgConnStatus :: ACommand 'Agent e -> Maybe ConnStatus
agentMsgConnStatus = \case
CONF {} -> Just ConnRequested
@@ -3931,6 +3876,7 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
(ct', conn') <- updateContactPQRcv user ct conn pqEncryption
checkIntegrityCreateItem (CDDirectRcv ct') msgMeta `catchChatError` \_ -> pure ()
(conn'', msg@RcvMessage {chatMsgEvent = ACME _ event}) <- saveDirectRcvMSG conn' msgMeta msgBody
lift $ agentDeliveryStatus (AgentConnId agentConnId) $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag event)}
let ct'' = ct' {activeConn = Just conn''} :: Contact
assertDirectAllowed user MDRcv ct'' $ toCMEventTag event
case event of
@@ -4351,6 +4297,7 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
processEvent :: MsgEncodingI e => ChatMessage e -> CM ()
processEvent chatMsg = do
(m', conn', msg@RcvMessage {chatMsgEvent = ACME _ event}) <- saveGroupRcvMsg user groupId m conn msgMeta msgBody chatMsg
lift $ agentDeliveryStatus (AgentConnId agentConnId) $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag event)}
case event of
XMsgNew mc -> memberCanSend m' $ newGroupContentMessage gInfo m' mc msg brokerTs False
XMsgFileDescr sharedMsgId fileDescr -> memberCanSend m' $ groupMessageFileDescription gInfo m' sharedMsgId fileDescr
@@ -4689,16 +4636,20 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
withCompletedCommand :: forall e. AEntityI e => Connection -> ACommand 'Agent e -> (CommandData -> CM ()) -> CM ()
withCompletedCommand Connection {connId} agentMsg action = do
let agentMsgTag = APCT (sAEntity @e) $ aCommandTag agentMsg
cmdData_ <- withStore' $ \db -> getCommandDataByCorrId db user corrId
case cmdData_ of
Just cmdData@CommandData {cmdId, cmdConnId = Just cmdConnId', cmdFunction}
| connId == cmdConnId' && (agentMsgTag == commandExpectedResponse cmdFunction || agentMsgTag == APCT SAEConn ERR_) -> do
withStore' $ \db -> deleteCommand db user cmdId
action cmdData
| otherwise -> err cmdId $ "not matching connection id or unexpected response, corrId = " <> show corrId
Just CommandData {cmdId, cmdConnId = Nothing} -> err cmdId $ "no command connection id, corrId = " <> show corrId
Nothing -> throwChatError . CEAgentCommandError $ "command not found, corrId = " <> show corrId
if agentMsgTag == APCT SAEConn OK_ && corrId /= "" then markDelivery else do
cmdData_ <- withStore' $ \db -> getCommandDataByCorrId db user corrId
case cmdData_ of
Just cmdData@CommandData {cmdId, cmdConnId = Just cmdConnId', cmdFunction}
| connId == cmdConnId' && (agentMsgTag == commandExpectedResponse cmdFunction || agentMsgTag == APCT SAEConn ERR_) -> do
withStore' $ \db -> deleteCommand db user cmdId
action cmdData
| otherwise -> err cmdId $ "not matching connection id or unexpected response, corrId = " <> show corrId
Just CommandData {cmdId, cmdConnId = Nothing} -> err cmdId $ "no command connection id, corrId = " <> show corrId
Nothing -> throwChatError . CEAgentCommandError $ "command not found, corrId = " <> show corrId
where
markDelivery = lift $ agentDeliveryStatus (AgentConnId agentConnId) $ \ad@AgentDeliveryStatus {pendingAcks} -> ad {pendingAcks = M.adjust (const True) ackKey pendingAcks}
where
ackKey = decodeLatin1 $ strEncode corrId
err cmdId msg = do
withStore' $ \db -> updateCommandStatus db user cmdId CSError
throwChatError . CEAgentCommandError $ msg
@@ -4724,12 +4675,12 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
where
ackMsg :: MsgMeta -> Maybe MsgReceiptInfo -> CM ()
ackMsg MsgMeta {recipient = (msgId, _)} rcpt = do
withAgent $ \a -> ackMessageAsync a "" cId msgId rcpt
acs <- asks agentConnStatuses
liftIO $ do
let acId = AgentConnId cId
now <- getCurrentTime
atomically $ TM.lookup acId acs >>= mapM_ (\v -> modifyTVar' v $ \cs -> cs {ackSent = Just now})
ackId <- drgRandomBytes 24
withAgent $ \a -> ackMessageAsync a ackId cId msgId rcpt
now <- liftIO getCurrentTime
let ackKey = decodeLatin1 $ strEncode ackId
lift . agentDeliveryStatus (AgentConnId agentConnId) $ \ad@AgentDeliveryStatus {pendingAcks} ->
ad {ackSent = Just (now, ackKey), pendingAcks = M.insert ackKey False pendingAcks}
sentMsgDeliveryEvent :: Connection -> AgentMsgId -> CM ()
sentMsgDeliveryEvent Connection {connId} msgId =
@@ -6028,6 +5979,7 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage =
processForwardedMsg author chatMsg = do
let body = LB.toStrict $ J.encode msg
rcvMsg@RcvMessage {chatMsgEvent = ACME _ event} <- saveGroupFwdRcvMsg user groupId m author body chatMsg
lift $ agentDeliveryStatus (AgentConnId agentConnId) $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag event)}
case event of
XMsgNew mc -> memberCanSend author $ newGroupContentMessage gInfo author mc rcvMsg msgTs True
XMsgFileDescr sharedMsgId fileDescr -> memberCanSend author $ groupMessageFileDescription gInfo author sharedMsgId fileDescr
@@ -6652,9 +6604,10 @@ sendPendingGroupMessages user GroupMember {groupMemberId, localDisplayName} conn
-- TODO [batch send] refactor direct message processing same as groups (e.g. checkIntegrity before processing)
saveDirectRcvMSG :: Connection -> MsgMeta -> MsgBody -> CM (Connection, RcvMessage)
saveDirectRcvMSG conn@Connection {connId} agentMsgMeta msgBody =
saveDirectRcvMSG conn@Connection {connId, agentConnId} agentMsgMeta msgBody =
case parseChatMessages msgBody of
[Right (ACMsg _ ChatMessage {chatVRange, msgId = sharedMsgId_, chatMsgEvent})] -> do
lift $ agentDeliveryStatus agentConnId $ \ad -> ad {eventTag = Just $! tshow (toCMEventTag chatMsgEvent)}
conn' <- updatePeerChatVRange conn chatVRange
let agentMsgId = fst $ recipient agentMsgMeta
newMsg = NewRcvMessage {chatMsgEvent, msgBody}
@@ -7357,7 +7310,7 @@ chatCommandP =
"/_download " *> (APIDownloadStandaloneFile <$> A.decimal <* A.space <*> strP_ <*> cryptoFileP),
("/quit" <|> "/q" <|> "/exit") $> QuitChat,
("/version" <|> "/v") $> ShowVersion,
"/debug acks" $> DebugAcks,
"/debug delivery" $> DebugDelivery,
"/debug locks" $> DebugLocks,
"/debug event " *> (DebugEvent <$> jsonP),
"/get stats" $> GetAgentStats,
@@ -7519,6 +7472,9 @@ timeItToView s action = do
toView' $ CRTimedAction s diff
pure a
drgRandomBytes :: Int -> CM ByteString
drgRandomBytes n = asks random >>= atomically . C.randomBytes n
mkValidName :: String -> String
mkValidName = reverse . dropWhile isSpace . fst3 . foldl' addChar ("", '\NUL', 0 :: Int)
where
@@ -7563,3 +7519,8 @@ xftpSndFileRedirect user ftId vfd = do
dummyFileDescr :: FileDescr
dummyFileDescr = FileDescr {fileDescrText = "", fileDescrPartNo = 0, fileDescrComplete = False}
agentDeliveryStatus :: AgentConnId -> (AgentDeliveryStatus -> AgentDeliveryStatus) -> CM' ()
agentDeliveryStatus acId f = do
ads <- asks agentDeliveryStatuses
atomically $ TM.lookup acId ads >>= mapM_ (\v -> modifyTVar' v f)
+21 -10
View File
@@ -207,7 +207,7 @@ data ChatController = ChatController
inputQ :: TBQueue String,
outputQ :: TBQueue (Maybe CorrId, Maybe RemoteHostId, ChatResponse),
connNetworkStatuses :: TMap AgentConnId NetworkStatus,
agentConnStatuses :: TMap AgentConnId (TVar AgentConnStatus),
agentDeliveryStatuses :: TMap AgentConnId (TVar AgentDeliveryStatus),
subscriptionMode :: TVar SubscriptionMode,
chatLock :: Lock,
entityLocks :: TMap ChatLockEntity Lock,
@@ -234,14 +234,20 @@ data ChatController = ChatController
contactMergeEnabled :: TVar Bool
}
data AgentConnStatus = AgentConnStatus
data AgentDeliveryStatus = AgentDeliveryStatus
{ lastCmd :: UTCTime,
lastCmdTag :: ACommandTag 'Agent 'AEConn,
lastMsg :: Maybe UTCTime, -- no message yet / got an MSG to ack
ackSent :: Maybe UTCTime,
okRcvd :: Maybe UTCTime -- ACK delivered, resulting in OK or MSG
isMSG :: Bool, -- False for RCVD
connId :: Maybe Int64, -- chat connection ID
eventTag :: Maybe Text, -- tshow of ACMEventTag (for JSON instances)
ackSent :: Maybe (UTCTime, Text), -- strEncode of random CorrId
pendingAcks :: Map Text Bool
}
deriving (Show)
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
deriving (Show)
@@ -498,7 +504,8 @@ data ChatCommand
| APIStandaloneFileInfo FileDescriptionURI
| QuitChat
| ShowVersion
| DebugAcks
| DebugDelivery
-- | DebugConnection Int64
| DebugLocks
| DebugEvent ChatResponse
| GetAgentStats
@@ -746,7 +753,7 @@ data ChatResponse
| CRContactPQEnabled {user :: User, contact :: Contact, pqEnabled :: PQEncryption}
| CRSQLResult {rows :: [Text]}
| CRSlowSQLQueries {chatQueries :: [SlowSQLQuery], agentQueries :: [SlowSQLQuery]}
| CRDebugAcks {debugAcks :: Map DebugAckKey DebugAck}
| CRDebugDelivery {debugDelivery :: Map Text AgentDeliveryStatus}
| CRDebugLocks {chatLockName :: Maybe String, chatEntityLocks :: Map String String, agentLocks :: AgentLocks}
| CRAgentStats {agentStats :: [[String]]}
| CRAgentWorkersDetails {agentWorkersDetails :: AgentWorkersDetails}
@@ -1495,7 +1502,11 @@ $(JQ.deriveJSON (sumTypeJSON $ dropPrefix "RCSR") ''RemoteCtrlStopReason)
$(JQ.deriveJSON (sumTypeJSON $ dropPrefix "RHSR") ''RemoteHostStopReason)
$(JQ.deriveJSON defaultJSON ''DebugAck)
$(JQ.deriveJSON defaultJSON ''AgentDeliveryStatus)
-- $(JQ.deriveJSON defaultJSON ''DebugDelivery)
-- $(JQ.deriveJSON defaultJSON ''DebugConnection)
$(JQ.deriveJSON (sumTypeJSON $ dropPrefix "CR") ''ChatResponse)
+1 -1
View File
@@ -351,7 +351,7 @@ responseToView hu@(currentRH, user_) ChatConfig {logLevel, showReactions, showRe
<> (" :: avg: " <> sShow timeAvg <> " ms")
<> (" :: " <> plain (T.unwords $ T.lines query))
in ("Chat queries" : map viewQuery chatQueries) <> [""] <> ("Agent queries" : map viewQuery agentQueries)
CRDebugAcks {debugAcks} -> [plain $ LB.unpack (J.encode debugAcks)]
CRDebugDelivery ads -> [plain $ LB.unpack (J.encode ads)]
CRDebugLocks {chatLockName, chatEntityLocks, agentLocks} ->
[ maybe "no chat lock" (("chat lock: " <>) . plain) chatLockName,
plain $ "chat entity locks: " <> LB.unpack (J.encode chatEntityLocks),