render channel on any message changes etc

This commit is contained in:
Evgeny @ SimpleX Chat
2026-06-03 17:34:12 +00:00
parent 751d5fd3f0
commit a5a21a2e9f
5 changed files with 193 additions and 32 deletions
+2 -2
View File
@@ -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,
+23 -1
View File
@@ -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,
+9 -11
View File
@@ -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
+13 -1
View File
@@ -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
+146 -17
View File
@@ -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