From a5a21a2e9fb4ad567e3be6e83e7492dc7f78d35f Mon Sep 17 00:00:00 2001 From: "Evgeny @ SimpleX Chat" <259188159+evgeny-simplex@users.noreply.github.com> Date: Wed, 3 Jun 2026 17:34:12 +0000 Subject: [PATCH] render channel on any message changes etc --- src/Simplex/Chat.hs | 4 +- src/Simplex/Chat/Controller.hs | 24 +++- src/Simplex/Chat/Library/Commands.hs | 20 ++- src/Simplex/Chat/Library/Subscriber.hs | 14 ++- src/Simplex/Chat/Web.hs | 163 ++++++++++++++++++++++--- 5 files changed, 193 insertions(+), 32 deletions(-) diff --git a/src/Simplex/Chat.hs b/src/Simplex/Chat.hs index bc22ab2c3b..ee14c42a66 100644 --- a/src/Simplex/Chat.hs +++ b/src/Simplex/Chat.hs @@ -183,7 +183,7 @@ newChatController deliveryJobWorkers <- TM.emptyIO relayRequestWorkers <- TM.emptyIO relayGroupLinkChecksAsync <- newTVarIO Nothing - webPreviewAsync <- newTVarIO Nothing + webPreviewState <- forM webPreviewConfig $ \_ -> newWebPreviewState chatRelayTests <- TM.emptyIO expireCIThreads <- TM.emptyIO expireCIFlags <- TM.emptyIO @@ -228,7 +228,7 @@ newChatController deliveryJobWorkers, relayRequestWorkers, relayGroupLinkChecksAsync, - webPreviewAsync, + webPreviewState, chatRelayTests, expireCIThreads, expireCIFlags, diff --git a/src/Simplex/Chat/Controller.hs b/src/Simplex/Chat/Controller.hs index be8ebed788..8e845fab26 100644 --- a/src/Simplex/Chat/Controller.hs +++ b/src/Simplex/Chat/Controller.hs @@ -39,6 +39,7 @@ import Data.Char (ord) import Data.Int (Int64) import Data.List.NonEmpty (NonEmpty) import Data.Map.Strict (Map) +import Data.Set (Set) import qualified Data.Map.Strict as M import Data.Maybe (fromMaybe) import Data.String @@ -177,6 +178,27 @@ data WebPreviewConfig = WebPreviewConfig webUpdateInterval :: Int -- seconds } +data WebPreviewState = WebPreviewState + { publishableGroupIds :: TVar (Map Int64 FilePath), + priorityRender :: TQueue Int64, + filesToRemove :: TQueue FilePath, + corsNeeded :: TVar Bool, + routinePending :: TVar (Set Int64), + wakeSignal :: TMVar (), + webPreviewWorkerAsync :: TVar (Maybe (Async ())) + } + +newWebPreviewState :: IO WebPreviewState +newWebPreviewState = do + publishableGroupIds <- newTVarIO mempty + priorityRender <- newTQueueIO + filesToRemove <- newTQueueIO + corsNeeded <- newTVarIO False + routinePending <- newTVarIO mempty + wakeSignal <- newEmptyTMVarIO + webPreviewWorkerAsync <- newTVarIO Nothing + pure WebPreviewState {publishableGroupIds, priorityRender, filesToRemove, corsNeeded, routinePending, wakeSignal, webPreviewWorkerAsync} + data RandomAgentServers = RandomAgentServers { smpServers :: NonEmpty (ServerCfg 'PSMP), xftpServers :: NonEmpty (ServerCfg 'PXFTP) @@ -264,7 +286,7 @@ data ChatController = ChatController deliveryJobWorkers :: TMap DeliveryWorkerKey Worker, relayRequestWorkers :: TMap Int Worker, -- single global worker with key 1 is used to fit into existing worker management framework relayGroupLinkChecksAsync :: TVar (Maybe (Async ())), - webPreviewAsync :: TVar (Maybe (Async ())), + webPreviewState :: Maybe WebPreviewState, chatRelayTests :: TMap ConnId RelayTest, expireCIThreads :: TMap UserId (Maybe (Async ())), expireCIFlags :: TMap UserId Bool, diff --git a/src/Simplex/Chat/Library/Commands.hs b/src/Simplex/Chat/Library/Commands.hs index cba1e1d838..b5ac4c74a3 100644 --- a/src/Simplex/Chat/Library/Commands.hs +++ b/src/Simplex/Chat/Library/Commands.hs @@ -55,7 +55,7 @@ import Data.Type.Equality import qualified Data.UUID as UUID import qualified Data.UUID.V4 as V4 import Simplex.Chat.Library.Subscriber -import Simplex.Chat.Web (renderWebPreviews) +import Simplex.Chat.Web (webPreviewWorker) import Simplex.Chat.Call import Simplex.Chat.Controller import Simplex.Chat.Delivery (DeliveryJobScope (..), DeliveryJobSpec (..), DeliveryWorkerScope (..)) @@ -238,16 +238,14 @@ startChatController mainApp enableSndFiles = do ChatConfig {webPreviewConfig = cfg_} <- asks config case (relayUsers, cfg_) of (_ : _, Just cfg) -> do - wpAsync <- asks webPreviewAsync - readTVarIO wpAsync >>= \case - Nothing -> do - cc <- ask - a <- Just <$> async (liftIO $ forever $ do - forM_ relayUsers $ \relayUser -> - renderWebPreviews cfg cc relayUser - threadDelay (webUpdateInterval cfg * 1000000)) - atomically $ writeTVar wpAsync a - _ -> pure () + wps_ <- asks webPreviewState + forM_ wps_ $ \WebPreviewState {webPreviewWorkerAsync} -> + readTVarIO webPreviewWorkerAsync >>= \case + Nothing -> do + cc <- ask + a <- Just <$> async (liftIO $ webPreviewWorker cfg cc relayUsers) + atomically $ writeTVar webPreviewWorkerAsync a + _ -> pure () _ -> pure () startExpireCIs user = whenM shouldExpireChats $ do startExpireCIThread user diff --git a/src/Simplex/Chat/Library/Subscriber.hs b/src/Simplex/Chat/Library/Subscriber.hs index 93e6cdab26..14a36c6e90 100644 --- a/src/Simplex/Chat/Library/Subscriber.hs +++ b/src/Simplex/Chat/Library/Subscriber.hs @@ -47,6 +47,7 @@ import Simplex.Chat.Call import Simplex.Chat.Controller import Simplex.Chat.Delivery import Simplex.Chat.Library.Internal +import Simplex.Chat.Web (channelChanged, channelRemoved) import Simplex.Chat.Messages import Simplex.Chat.Messages.Batch (batchDeliveryTasks1, batchProfiles, batchProfilesWithBody, encodeBinaryBatch, encodeFwdElement, maxBatchElementSize) import Simplex.Chat.Messages.CIContent @@ -1017,6 +1018,7 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage = processEvent :: forall e. MsgEncodingI e => GroupInfo -> GroupMember -> VerifiedMsg e -> CM (Maybe NewMessageDeliveryTask) processEvent gInfo' m' verifiedMsg = do (m'', conn', msg@RcvMessage {msgId, chatMsgEvent = ACME _ event}) <- saveGroupRcvMsg user groupId m' conn msgMeta verifiedMsg + cc <- ask let ctx js = DeliveryTaskContext js False checkSendAsGroup :: Maybe Bool -> CM (Maybe DeliveryTaskContext) -> CM (Maybe DeliveryTaskContext) checkSendAsGroup asGroup_ a @@ -1073,7 +1075,17 @@ processAgentMessageConn vr user@User {userId} corrId agentConnId agentMessage = XInfoProbeOk probe -> Nothing <$ xInfoProbeOk (COMGroupMember m'') probe BFileChunk sharedMsgId chunk -> Nothing <$ bFileChunkGroup gInfo' sharedMsgId chunk msgMeta _ -> Nothing <$ messageError ("unsupported message: " <> tshow event) - forM deliveryTaskContext_ $ \taskContext -> + forM deliveryTaskContext_ $ \taskContext -> do + let contentChanged :: CM () + contentChanged = atomically $ channelChanged cc groupId False + case event of + XMsgNew {} -> contentChanged + XMsgUpdate {} -> contentChanged + XMsgDel {} -> contentChanged + XMsgReact {} -> contentChanged + XGrpInfo {} -> atomically $ channelChanged cc groupId True + XGrpDel {} -> atomically $ channelRemoved cc groupId + _ -> pure () pure $ NewMessageDeliveryTask {messageId = msgId, taskContext} checkSendRcpt :: [AParsedMsg] -> CM Bool checkSendRcpt aMsgs = do diff --git a/src/Simplex/Chat/Web.hs b/src/Simplex/Chat/Web.hs index 6d2ceeb215..94d4cc0afc 100644 --- a/src/Simplex/Chat/Web.hs +++ b/src/Simplex/Chat/Web.hs @@ -14,26 +14,33 @@ module Simplex.Chat.Web WebMemberProfile (..), WebFileInfo (..), CorsOrigin (..), - renderWebPreviews, - writeCorsConfig, + webPreviewWorker, + channelChanged, + channelRemoved, ) where -import Control.Monad (forM_) +import Control.Concurrent.STM (check) +import Control.Exception (SomeException, catch) +import Control.Logger.Simple +import Control.Monad (forM_, unless, void, when) +import Control.Monad.Except (runExceptT) import Data.Either (rights) +import Data.Int (Int64) import qualified Data.Aeson as J import qualified Data.Aeson.TH as JQ import qualified Data.ByteString.Char8 as B import qualified Data.ByteString.Lazy as LB import Data.List (nubBy) import Data.Map.Strict (Map) +import qualified Data.Map.Strict as M import qualified Data.Set as S import Data.Maybe (isJust, mapMaybe) import Data.Text (Text) import qualified Data.Text as T import qualified Data.Text.IO as TIO import Data.Time.Clock (UTCTime, getCurrentTime) -import Simplex.Chat.Controller (ChatConfig (..), ChatController (..), WebPreviewConfig (..)) +import Simplex.Chat.Controller (ChatConfig (..), ChatController (..), WebPreviewConfig (..), WebPreviewState (..)) import Simplex.Chat.Markdown (FormattedText (..), MarkdownList, parseMaybeMarkdownList) import Simplex.Chat.Messages ( CChatItem (..), @@ -51,6 +58,7 @@ import Simplex.Chat.Messages.CIContent (ciMsgContent) import Simplex.Chat.Protocol (MsgContent, MsgRef (..), QuotedMsg (..), isReport) import Simplex.Chat.Store.Groups (getGroupOwners, getRelayServedGroups) import Simplex.Chat.Store.Messages (getGroupWebPreviewItems) +import Simplex.Chat.Store.Shared (getGroupInfo) import Simplex.Chat.Types ( B64UrlByteString, GroupInfo (..), @@ -62,13 +70,14 @@ import Simplex.Chat.Types MemberName, PublicGroupAccess (..), PublicGroupProfile (..), - User, + User (..), ) import Simplex.Messaging.Agent.Store.Common (withTransaction) import Simplex.Messaging.Encoding.String (strEncode) import Simplex.Messaging.Parsers (defaultJSON) import System.Directory (createDirectoryIfMissing, listDirectory, removeFile) import System.FilePath (takeExtension, ()) +import UnliftIO.STM data WebFileInfo = WebFileInfo { fileName :: String, @@ -115,20 +124,113 @@ $(JQ.deriveJSON defaultJSON ''WebMessage) $(JQ.deriveJSON defaultJSON ''WebChannelPreview) -renderWebPreviews :: WebPreviewConfig -> ChatController -> User -> IO () -renderWebPreviews WebPreviewConfig {webJsonDir, webCorsFile} cc user = do - createDirectoryIfMissing True webJsonDir - groups <- withTransaction (chatStore cc) $ \db -> getRelayServedGroups db vr' user - let publishable = filter hasPublicGroup groups - activeFiles = S.fromList $ mapMaybe publicGroupFileName publishable - corsEntries <- mapMaybe id <$> mapM (renderGroupPreview webJsonDir cc user) publishable - removeStaleFiles webJsonDir activeFiles - forM_ webCorsFile $ writeCorsConfig corsEntries +webPreviewWorker :: WebPreviewConfig -> ChatController -> [User] -> IO () +webPreviewWorker WebPreviewConfig {webJsonDir, webCorsFile, webUpdateInterval} cc users = + forM_ (webPreviewState cc) $ \wps -> do + createDirectoryIfMissing True webJsonDir + initPublishableGroups wps + seedRoutinePending wps + workerLoop wps where vr' = chatVRange (config cc) - hasPublicGroup GroupInfo {groupProfile = GroupProfile {publicGroup}} = isJust publicGroup - publicGroupFileName GroupInfo {groupProfile = GroupProfile {publicGroup}} = - (\PublicGroupProfile {publicGroupId} -> publicGroupIdFileName publicGroupId <> ".json") <$> publicGroup + + workerLoop wps@WebPreviewState {priorityRender, filesToRemove, corsNeeded, routinePending, wakeSignal} = do + drainRemovals + drainPriority + handleCors + processOneRoutine + sleepOrWake + workerLoop wps + where + drainRemovals = atomically (tryReadTQueue filesToRemove) >>= \case + Nothing -> pure () + Just f -> do + removeFile (webJsonDir f) `catch` \(_ :: SomeException) -> pure () + drainRemovals + + drainPriority = atomically (tryReadTQueue priorityRender) >>= \case + Nothing -> pure () + Just gId -> do + renderOneGroup wps gId + drainPriority + + handleCors = do + needed <- atomically $ swapTVar corsNeeded False + when needed regenerateCors + + processOneRoutine = do + mGId <- atomically $ do + pending <- readTVar routinePending + case S.minView pending of + Nothing -> pure Nothing + Just (gId, rest) -> writeTVar routinePending rest >> pure (Just gId) + forM_ mGId $ renderOneGroup wps + + sleepOrWake = do + pending <- readTVarIO routinePending + if S.null pending + then do + cleanStaleFiles + seedRoutinePending wps + interruptibleSleep + else do + hasPriority <- atomically $ not <$> isEmptyTQueue priorityRender + unless hasPriority interruptibleSleep + + interruptibleSleep = do + delay <- registerDelay (webUpdateInterval * 1000000) + atomically $ + (readTVar delay >>= check) + `orElse` takeTMVar wakeSignal + + initPublishableGroups WebPreviewState {publishableGroupIds} = do + groups <- withTransaction (chatStore cc) $ \db -> + concat <$> mapM (getRelayServedGroups db vr') users + let gIds = M.fromList [(groupId, f) | g@GroupInfo {groupId} <- groups, Just f <- [publicGroupFileName g]] + atomically $ writeTVar publishableGroupIds gIds + + seedRoutinePending WebPreviewState {publishableGroupIds, routinePending} = + atomically $ M.keysSet <$> readTVar publishableGroupIds >>= writeTVar routinePending + + renderOneGroup WebPreviewState {publishableGroupIds} gId = do + publishable <- atomically $ M.member gId <$> readTVar publishableGroupIds + when publishable $ do + r <- withTransaction (chatStore cc) $ \db -> + findUser $ \u -> fmap (\g -> (u, g)) <$> runExceptT (getGroupInfo db vr' u gId) + case r of + Just (u, gInfo) | hasPublicGroup gInfo -> + void $ renderGroupPreview webJsonDir cc u gInfo + _ -> do + fName <- atomically $ do + ids <- readTVar publishableGroupIds + modifyTVar' publishableGroupIds (M.delete gId) + pure $ M.lookup gId ids + forM_ fName $ \f -> + removeFile (webJsonDir f) `catch` \(_ :: SomeException) -> pure () + logInfo $ "web preview: group " <> T.pack (show gId) <> " no longer publishable" + + findUser f = go users + where + go [] = pure Nothing + go (u : us) = f u >>= \case + Right a -> pure (Just a) + Left _ -> go us + + regenerateCors = do + groups <- withTransaction (chatStore cc) $ \db -> + concat <$> mapM (getRelayServedGroups db vr') users + let entries = mapMaybe groupCorsEntry groups + forM_ webCorsFile $ writeCorsConfig entries + + groupCorsEntry GroupInfo {groupProfile = GroupProfile {publicGroup}} = + publicGroup >>= \PublicGroupProfile {publicGroupId, publicGroupAccess} -> + corsEntry publicGroupId <$> publicGroupAccess + + cleanStaleFiles = do + groups <- withTransaction (chatStore cc) $ \db -> + concat <$> mapM (getRelayServedGroups db vr') users + let activeFiles = S.fromList $ mapMaybe publicGroupFileName [g | g <- groups, hasPublicGroup g] + removeStaleFiles webJsonDir activeFiles renderGroupPreview :: FilePath -> ChatController -> User -> GroupInfo -> IO (Maybe (Text, CorsOrigin)) renderGroupPreview webJsonDir cc user gInfo@GroupInfo {groupProfile = gp@GroupProfile {shortDescr = sd, description = wd, publicGroup}} = @@ -157,6 +259,26 @@ renderGroupPreview webJsonDir cc user gInfo@GroupInfo {groupProfile = gp@GroupPr where vr' = chatVRange (config cc) +channelChanged :: ChatController -> Int64 -> Bool -> STM () +channelChanged cc gId updateCors = + forM_ (webPreviewState cc) $ \WebPreviewState {publishableGroupIds, priorityRender, corsNeeded, routinePending, wakeSignal} -> do + ids <- readTVar publishableGroupIds + when (M.member gId ids) $ do + writeTQueue priorityRender gId + modifyTVar' routinePending (S.delete gId) + when updateCors $ writeTVar corsNeeded True + void $ tryPutTMVar wakeSignal () + +channelRemoved :: ChatController -> Int64 -> STM () +channelRemoved cc gId = + forM_ (webPreviewState cc) $ \WebPreviewState {publishableGroupIds, filesToRemove, corsNeeded, routinePending, wakeSignal} -> do + ids <- readTVar publishableGroupIds + forM_ (M.lookup gId ids) $ writeTQueue filesToRemove + modifyTVar' publishableGroupIds (M.delete gId) + modifyTVar' routinePending (S.delete gId) + writeTVar corsNeeded True + void $ tryPutTMVar wakeSignal () + toRenderedItem :: CChatItem 'CTGroup -> Maybe (WebMessage, Maybe WebMemberProfile) toRenderedItem (CChatItem _ ChatItem {chatDir, meta = CIMeta {itemTs, itemTimed, itemForwarded, itemEdited}, content, mentions, formattedText, quotedItem, reactions, file}) | isJust itemTimed = Nothing @@ -258,3 +380,10 @@ toFormattedText t = case parseMaybeMarkdownList t of publicGroupIdFileName :: B64UrlByteString -> String publicGroupIdFileName = B.unpack . strEncode + +hasPublicGroup :: GroupInfo -> Bool +hasPublicGroup GroupInfo {groupProfile = GroupProfile {publicGroup}} = isJust publicGroup + +publicGroupFileName :: GroupInfo -> Maybe FilePath +publicGroupFileName GroupInfo {groupProfile = GroupProfile {publicGroup}} = + (\PublicGroupProfile {publicGroupId} -> publicGroupIdFileName publicGroupId <> ".json") <$> publicGroup