diff --git a/backend/package.json b/backend/package.json index 27772897..4d0dbe5e 100644 --- a/backend/package.json +++ b/backend/package.json @@ -6,7 +6,7 @@ "scripts": { "predev": "node scripts/free-ports.js 8080", "dev": "nodemon --watch server --watch .env --ignore logs --ignore uploads server", - "build": "node --check server/index.js && node --check server/server.js && node --check server/bootstrap.js && node --check server/config.js && node --check server/db/sequelize.js && node --check server/db/connect.js && node --check server/db/models/inbox.js && node --check server/db/models/chat.js && node --check server/db/models/chatRoomCounter.js && node --check server/db/models/messageReceipt.js && node --check server/db/models/chatDraft.js && node --check server/db/models/messageRequest.js && node --check server/db/models/chatTopic.js && node --check server/db/models/e2eeDeviceKey.js && node --check server/db/models/resumableUpload.js && node --check server/db/models/chatAiConfig.js && node --check server/db/models/adminSocialAuthConfig.js && node --check server/db/models/profile.js && node --check server/db/models/group.js && node --check server/db/models/channel.js && node --check server/db/models/adminStorageConfig.js && node --check server/db/models/adminCallConfig.js && node --check server/db/models/callHistory.js && node --check server/db/models/profilePhoto.js && node --check server/db/models/nativePushDevice.js && node --check server/helpers/storageConfig.js && node --check server/helpers/callConfig.js && node --check server/helpers/callHistory.js && node --check server/helpers/callState.js && node --check server/helpers/avatarDefaults.js && node --check server/helpers/ensureAvatarDefaults.js && node --check server/helpers/ensureProfile.js && node --check server/helpers/privacy.js && node --check server/helpers/profilePhotos.js && node --check server/helpers/livekit.js && node --check server/helpers/nativePushConfig.js && node --check server/helpers/nativePush.js && node --check server/helpers/storage.js && node --check server/helpers/videoPipeline.js && node --check server/helpers/accountArchive.js && node --check server/helpers/accountExport.js && node --check server/helpers/logger.js && node --check server/helpers/chatAbuse.js && node --check server/helpers/chatMentions.js && node --check server/helpers/messageRequests.js && node --check server/helpers/chatReliability.js && node --check server/helpers/chatAiConfig.js && node --check server/helpers/socialAuthConfig.js && node --check server/helpers/e2eeKeyDirectory.js && node --check server/helpers/chatMaintenance.js && node --check server/middleware/chatSendIdempotency.js && node --check server/controllers/avatar.js && node --check server/controllers/profile.js && node --check server/controllers/chatUpload.js && node --check server/controllers/chatSecurity.js && node --check server/controllers/chatDeletion.js && node --check server/controllers/chatV2.js && node --check server/controllers/chatAiAdmin.js && node --check server/controllers/chatSuggestions.js && node --check server/controllers/chatResumableUpload.js && node --check server/controllers/socialAuth.js && node --check server/controllers/socialAuthAdmin.js && node --check server/controllers/storageAdmin.js && node --check server/controllers/callingAdmin.js && node --check server/controllers/callingPushAdmin.js && node --check server/controllers/callingConfig.js && node --check server/controllers/nativePush.js && node --check server/controllers/accountStorage.js && node --check server/routes/profile.js && node --check server/routes/inbox.js && node --check server/routes/chat.js && node --check server/routes/chatV2.js && node --check server/routes/chatAiAdmin.js && node --check server/routes/socialAuthAdmin.js && node --check server/routes/storageAdmin.js && node --check server/routes/callingAdmin.js && node --check server/routes/callingPushAdmin.js && node --check server/routes/callingConfig.js && node --check server/routes/setting.js && node --check server/helpers/socketAdapter.js && node --check server/socket/auth.js && node --check server/socket/index.js && node --check server/socket/events/chatV2.js && node --check server/socket/events/calling.js && node --check server/socket/events/backgroundCalling.js && node --check server/socket/events/groupModeration.js && node --check server/socket/events/room.js && node --check server/routes/cron.js && node --check api/index.js", + "build": "node --check server/index.js && node --check server/server.js && node --check server/bootstrap.js && node --check server/config.js && node --check server/db/sequelize.js && node --check server/db/connect.js && node --check server/db/models/inbox.js && node --check server/db/models/chat.js && node --check server/db/models/chatRoomCounter.js && node --check server/db/models/messageReceipt.js && node --check server/db/models/chatDraft.js && node --check server/db/models/messageRequest.js && node --check server/db/models/chatTopic.js && node --check server/db/models/e2eeDeviceKey.js && node --check server/db/models/resumableUpload.js && node --check server/db/models/chatAiConfig.js && node --check server/db/models/adminSocialAuthConfig.js && node --check server/db/models/profile.js && node --check server/db/models/group.js && node --check server/db/models/channel.js && node --check server/db/models/adminStorageConfig.js && node --check server/db/models/adminCallConfig.js && node --check server/db/models/callHistory.js && node --check server/db/models/profilePhoto.js && node --check server/db/models/nativePushDevice.js && node --check server/helpers/storageConfig.js && node --check server/helpers/callConfig.js && node --check server/helpers/callHistory.js && node --check server/helpers/callState.js && node --check server/helpers/avatarDefaults.js && node --check server/helpers/ensureAvatarDefaults.js && node --check server/helpers/ensureProfile.js && node --check server/helpers/privacy.js && node --check server/helpers/profilePhotos.js && node --check server/helpers/livekit.js && node --check server/helpers/nativePushConfig.js && node --check server/helpers/nativePush.js && node --check server/helpers/storage.js && node --check server/helpers/videoPipeline.js && node --check server/helpers/accountArchive.js && node --check server/helpers/accountExport.js && node --check server/helpers/logger.js && node --check server/helpers/chatAbuse.js && node --check server/helpers/chatMentions.js && node --check server/helpers/messageRequests.js && node --check server/helpers/chatReliability.js && node --check server/helpers/chatAiConfig.js && node --check server/helpers/socialAuthConfig.js && node --check server/helpers/e2eeKeyDirectory.js && node --check server/helpers/chatMaintenance.js && node --check server/middleware/chatSendIdempotency.js && node --check server/controllers/avatar.js && node --check server/controllers/profile.js && node --check server/controllers/chatUpload.js && node --check server/controllers/chatSecurity.js && node --check server/controllers/chatDeletion.js && node --check server/controllers/chatV2.js && node --check server/controllers/chatAiAdmin.js && node --check server/controllers/chatSuggestions.js && node --check server/controllers/chatResumableUpload.js && node --check server/controllers/socialAuth.js && node --check server/controllers/socialAuthAdmin.js && node --check server/controllers/storageAdmin.js && node --check server/controllers/callingAdmin.js && node --check server/controllers/callingPushAdmin.js && node --check server/controllers/callingConfig.js && node --check server/controllers/nativePush.js && node --check server/controllers/accountStorage.js && node --check server/routes/profile.js && node --check server/routes/inbox.js && node --check server/routes/chat.js && node --check server/routes/chatV2.js && node --check server/routes/chatAiAdmin.js && node --check server/routes/socialAuthAdmin.js && node --check server/routes/storageAdmin.js && node --check server/routes/callingAdmin.js && node --check server/routes/callingPushAdmin.js && node --check server/routes/callingConfig.js && node --check server/routes/setting.js && node --check server/helpers/socketAdapter.js && node --check server/socket/auth.js && node --check server/socket/index.js && node --check server/socket/events/chatV2.js && node --check server/socket/events/calling.js && node --check server/socket/events/callRealtime.js && node --check server/socket/events/backgroundCalling.js && node --check server/socket/events/groupModeration.js && node --check server/socket/events/room.js && node --check server/routes/cron.js && node --check api/index.js", "start": "node server" }, "dependencies": { diff --git a/backend/server/bootstrap.js b/backend/server/bootstrap.js index b55bbc42..3d7fc6e8 100644 --- a/backend/server/bootstrap.js +++ b/backend/server/bootstrap.js @@ -3,6 +3,7 @@ const { startScheduledMessageWorker, } = require('./helpers/scheduledMessages'); const { configureSocketAdapter } = require('./helpers/socketAdapter'); +const { getCallConfig } = require('./helpers/callConfig'); const ensureAvatarDefaults = require('./helpers/ensureAvatarDefaults'); const ensureChatIndexes = require('./helpers/chatIndexes'); const logger = require('./helpers/logger'); @@ -14,13 +15,17 @@ const bootstrap = ({ startScheduledWorker = true } = {}) => { if (!bootstrapPromise) { bootstrapPromise = (async () => { await connectDb(); - await ensureChatIndexes(); - await ensureAvatarDefaults(); - await configureSocketAdapter(global.io); + await Promise.all([ + ensureChatIndexes(), + ensureAvatarDefaults(), + configureSocketAdapter(global.io), + getCallConfig(), + ]); logger.info('RUNTIME_READY', { vercel: process.env.VERCEL === '1', redis: Boolean(process.env.REDIS_URL), + realtimeWarm: true, }); return true; diff --git a/backend/server/helpers/callState.js b/backend/server/helpers/callState.js index 4dbcd821..48811828 100644 --- a/backend/server/helpers/callState.js +++ b/backend/server/helpers/callState.js @@ -1,5 +1,8 @@ -const { createClient } = require('redis'); const logger = require('./logger'); +const { + getSocketRedisCommandClient, + isRedisConfigured, +} = require('./socketAdapter'); const DEFAULT_TTL_SEC = 24 * 60 * 60; const MIN_TTL_SEC = 5 * 60; @@ -12,10 +15,8 @@ const clampTtl = (value) => { }; const stateTtlSec = clampTtl(process.env.CALL_STATE_TTL_SEC); -const redisUrl = String(process.env.REDIS_URL || '').trim(); const production = String(process.env.NODE_ENV || '').toLowerCase() === 'production'; -let redisPromise = null; let warnedFallback = false; const memory = new Map(); @@ -60,7 +61,7 @@ const memoryDelete = (key) => { }; const getRedis = async () => { - if (!redisUrl) { + if (!isRedisConfigured()) { if (production) { const error = new Error( 'REDIS_URL is required for durable production call state' @@ -77,24 +78,12 @@ const getRedis = async () => { return null; } - if (redisPromise) return redisPromise; + const client = await getSocketRedisCommandClient(); + if (client) return client; - redisPromise = (async () => { - const client = createClient({ url: redisUrl }); - client.on('error', (error) => { - logger.error('CALL_STATE_REDIS_ERROR', { message: error.message }); - }); - await client.connect(); - logger.info('CALL_STATE_REDIS_READY', { - ttlSec: stateTtlSec, - }); - return client; - })().catch((error) => { - redisPromise = null; - throw error; - }); - - return redisPromise; + const error = new Error('Shared Redis runtime is not ready for call state'); + error.code = 'CALL_STATE_REDIS_NOT_READY'; + throw error; }; const encodeState = (state, ttlSec = stateTtlSec) => ({ diff --git a/backend/server/helpers/chatAbuse.js b/backend/server/helpers/chatAbuse.js index e0f581a4..34d14d54 100644 --- a/backend/server/helpers/chatAbuse.js +++ b/backend/server/helpers/chatAbuse.js @@ -1,9 +1,11 @@ const crypto = require('crypto'); -const { createClient } = require('redis'); const logger = require('./logger'); +const { + getSocketRedisCommandClient, + isRedisConfigured, +} = require('./socketAdapter'); const localWindows = new Map(); -let redisPromise = null; const MAX_MESSAGES = Math.max(5, Number(process.env.CHAT_RATE_LIMIT_MESSAGES || 30)); const WINDOW_SECONDS = Math.max(5, Number(process.env.CHAT_RATE_LIMIT_WINDOW_SEC || 10)); @@ -11,23 +13,14 @@ const MAX_DUPLICATES = Math.max(2, Number(process.env.CHAT_DUPLICATE_LIMIT || 6) const DUPLICATE_WINDOW_SECONDS = Math.max(10, Number(process.env.CHAT_DUPLICATE_WINDOW_SEC || 30)); const getRedis = async () => { - const url = String(process.env.REDIS_URL || '').trim(); - if (!url) return null; - if (!redisPromise) { - redisPromise = (async () => { - const client = createClient({ url }); - client.on('error', (error) => { - logger.warn('CHAT_ABUSE_REDIS_ERROR', { message: error.message }); - }); - await client.connect(); - return client; - })().catch((error) => { - redisPromise = null; - logger.warn('CHAT_ABUSE_REDIS_CONNECT_FAILED', { message: error.message }); - return null; - }); - } - return redisPromise; + if (!isRedisConfigured()) return null; + const client = await getSocketRedisCommandClient(); + if (client) return client; + + logger.warn('CHAT_ABUSE_REDIS_NOT_READY', { + message: 'Shared Redis runtime is not ready; using local rate limits for this request', + }); + return null; }; const createRateError = (message, code) => { diff --git a/backend/server/helpers/chatReliability.js b/backend/server/helpers/chatReliability.js index 85537d33..f04b28f8 100644 --- a/backend/server/helpers/chatReliability.js +++ b/backend/server/helpers/chatReliability.js @@ -91,14 +91,32 @@ const wrapReliableChatInsert = (socket) => { return; } + // The warm Redis abuse guard must pass before any optional lookup starts. + // Sequence allocation and mention resolution can then run concurrently. await assertChatSendAllowed({ userId: socket.userId, text: args.text || '' }); - const sequence = await nextSequence(args.roomId); - const mentions = await resolveMentions({ - text: args.text || '', - roomId: args.roomId, - roomType: args.roomType, - senderId: socket.userId, - }); + const [sequence, mentions] = await Promise.all([ + nextSequence(args.roomId), + resolveMentions({ + text: args.text || '', + roomId: args.roomId, + roomType: args.roomType, + senderId: socket.userId, + }), + ]); + + const topicId = args.topicId || null; + const e2eeEnvelope = args.e2eeEnvelope && typeof args.e2eeEnvelope === 'object' + ? args.e2eeEnvelope + : null; + const transcript = String(args.transcript || '').slice(0, 8000); + + // Persist reliability metadata in the original create instead of doing a + // second database update after the message has already been written. + args.sequence = sequence; + args.mentionUserIds = mentions.mentionedUserIds; + args.topicId = topicId; + args.e2eeEnvelope = e2eeEnvelope; + args.transcript = transcript; await original(args); @@ -119,19 +137,21 @@ const wrapReliableChatInsert = (socket) => { const inbox = await InboxModel.findOne({ where: { roomId: args.roomId } }); const ownerIds = asArray(toPlain(inbox)?.ownersId || args.ownersId); - const topicId = args.topicId || null; - const e2eeEnvelope = args.e2eeEnvelope && typeof args.e2eeEnvelope === 'object' - ? args.e2eeEnvelope - : null; - const transcript = String(args.transcript || '').slice(0, 8000); - await created.update({ - clientMessageId: args.clientMessageId, - sequence, - mentionUserIds: mentions.mentionedUserIds, - topicId, - e2eeEnvelope, - transcript, + // Ack as soon as the durable chat and room ownership are confirmed. + // Message-request and mention side effects should not hold the sender UI. + emitMeta({ + socket, + chat: toPlain(created), + meta: { + clientMessageId: args.clientMessageId, + sequence, + mentionUserIds: mentions.mentionedUserIds, + topicId, + e2eeEnvelope, + transcript, + ownerIds, + }, }); if (args.roomType === 'private' && ownerIds.length === 2) { @@ -163,20 +183,6 @@ const wrapReliableChatInsert = (socket) => { }); }); } - - emitMeta({ - socket, - chat: toPlain(created), - meta: { - clientMessageId: args.clientMessageId, - sequence, - mentionUserIds: mentions.mentionedUserIds, - topicId, - e2eeEnvelope, - transcript, - ownerIds, - }, - }); } catch (error0) { logger.warn('CHAT_RELIABILITY_REJECTED', { userId: socket.userId, diff --git a/backend/server/helpers/socketAdapter.js b/backend/server/helpers/socketAdapter.js index eb7b7149..67948f66 100644 --- a/backend/server/helpers/socketAdapter.js +++ b/backend/server/helpers/socketAdapter.js @@ -44,6 +44,15 @@ const configureSocketAdapter = async (io = global.io) => { return adapterPromise; }; +const getSocketRedisCommandClient = async () => { + if (!isRedisConfigured()) return null; + if (redisClients?.pubClient?.isReady) return redisClients.pubClient; + if (!adapterPromise) return null; + + const clients = await adapterPromise; + return clients?.pubClient?.isReady ? clients.pubClient : null; +}; + const closeSocketAdapter = async () => { const clients = redisClients; redisClients = null; @@ -58,6 +67,7 @@ const closeSocketAdapter = async () => { module.exports = { configureSocketAdapter, + getSocketRedisCommandClient, closeSocketAdapter, isRedisConfigured, }; diff --git a/backend/server/socket/events/callRealtime.js b/backend/server/socket/events/callRealtime.js new file mode 100644 index 00000000..b4a88591 --- /dev/null +++ b/backend/server/socket/events/callRealtime.js @@ -0,0 +1,107 @@ +const { getActiveCallByRoom } = require('../../helpers/callState'); +const logger = require('../../helpers/logger'); + +const INFLIGHT_DEDUPE_MS = 2000; +const RECENT_START_TTL_MS = 2 * 60 * 1000; + +const emitExistingStarted = (socket, current) => { + socket.emit('call/started', { + callId: current.callId, + roomId: current.roomId, + roomType: current.roomType, + mediaType: current.mediaType, + }); +}; + +const attachCallRealtimeEvents = (socket) => { + if (!socket || socket.__syncchatCallRealtime) return; + socket.__syncchatCallRealtime = true; + + const listeners = socket.listeners('call/start'); + if (listeners.length) { + const original = listeners[0]; + const recentStarts = new Map(); + socket.removeAllListeners('call/start'); + + socket.on('call/start', async (rawArgs = {}) => { + const args = rawArgs && typeof rawArgs === 'object' ? rawArgs : {}; + const fromUserId = socket.userId || args.fromUserId; + const roomId = String(args.roomId || ''); + const mediaType = args.mediaType === 'video' ? 'video' : 'audio'; + if (!roomId || !fromUserId) return; + + const key = `${roomId}:${fromUserId}:${mediaType}`; + const now = Date.now(); + const previous = Number(recentStarts.get(key) || 0); + + if (previous) { + const current = await getActiveCallByRoom(roomId).catch(() => null); + if ( + current?.initiatorId === fromUserId && + current?.mediaType === mediaType + ) { + emitExistingStarted(socket, current); + return; + } + + // Only suppress a duplicate while the first call/start handler is + // still reserving state. If the first attempt failed, allow a retry. + if (now - previous < INFLIGHT_DEDUPE_MS) return; + } + + recentStarts.set(key, now); + for (const [entryKey, timestamp] of recentStarts.entries()) { + if (now - timestamp >= RECENT_START_TTL_MS) recentStarts.delete(entryKey); + } + + await original({ + ...args, + fromUserId, + mediaType, + }); + }); + + listeners.slice(1).forEach((listener) => socket.on('call/start', listener)); + } + + // When an outgoing call is signalled before getUserMedia finishes, a very + // fast recipient can join first. Once the initiator's media is ready and it + // joins, replay already-joined recipients so WebRTC offer creation cannot be + // missed because of that race. + socket.on('call/join', async ({ callId, roomId, userId, mediaType }) => { + const joiningUserId = socket.userId || userId; + if ((!callId && !roomId) || !joiningUserId) return; + + try { + // Let the primary calling handler process the same join event first. + await new Promise((resolve) => setTimeout(resolve, 0)); + const state = roomId + ? await getActiveCallByRoom(roomId) + : null; + if (!state || state.initiatorId !== joiningUserId) return; + + const joinedRecipients = (Array.isArray(state.joinedUserIds) + ? state.joinedUserIds + : [] + ).filter((id) => id && id !== joiningUserId); + + joinedRecipients.forEach((joinedUserId) => { + socket.emit('call/user-joined', { + callId: state.callId, + roomId: state.roomId, + userId: joinedUserId, + mediaType: mediaType === 'video' ? 'video' : state.mediaType, + catchUp: true, + }); + }); + } catch (error0) { + logger.warn('CALL_REALTIME_CATCHUP_ERROR', { + roomId: roomId || null, + userId: joiningUserId, + message: error0.message, + }); + } + }); +}; + +module.exports = attachCallRealtimeEvents; diff --git a/backend/server/socket/events/room.js b/backend/server/socket/events/room.js index aa952957..e4ef2dc6 100644 --- a/backend/server/socket/events/room.js +++ b/backend/server/socket/events/room.js @@ -1,5 +1,6 @@ const { io } = global; const attachCallingEvents = require('./calling'); +const attachCallRealtimeEvents = require('./callRealtime'); const attachBackgroundCallingEvents = require('./backgroundCalling'); const attachGroupModerationEvents = require('./groupModeration'); @@ -11,6 +12,7 @@ module.exports = (socket) => { }); attachCallingEvents(socket); + attachCallRealtimeEvents(socket); attachBackgroundCallingEvents(socket); attachGroupModerationEvents(socket); }; diff --git a/frontend/client/components/calling/globalCallLayer.jsx b/frontend/client/components/calling/globalCallLayer.jsx index 321f08fb..511572fd 100644 --- a/frontend/client/components/calling/globalCallLayer.jsx +++ b/frontend/client/components/calling/globalCallLayer.jsx @@ -115,6 +115,65 @@ function GlobalCallLayer() { }; }, [dispatch, master?._id, callPanel]); + // Start server-side call reservation/ringing as soon as the outgoing panel + // opens. Camera/microphone acquisition can then happen in parallel instead + // of blocking the recipient notification path. + useEffect(() => { + if ( + !master?._id || + callPanel?.mode !== 'outgoing' || + !callPanel?.roomId + ) { + return; + } + + socket.emit('call/start', { + roomId: callPanel.roomId, + roomType: callPanel.roomType === 'group' ? 'group' : 'private', + fromUserId: master._id, + mediaType: callPanel.mediaType === 'video' ? 'video' : 'audio', + fromName: callPanel.fromName || master.fullname || '', + fromUsername: callPanel.fromUsername || master.username || '', + recipientsId: Array.isArray(callPanel.recipientsId) + ? callPanel.recipientsId + : [], + }); + }, [ + master?._id, + callPanel?.mode, + callPanel?.roomId, + callPanel?.roomType, + callPanel?.mediaType, + ]); + + useEffect(() => { + const onStarted = (payload = {}) => { + if ( + callPanel?.mode !== 'outgoing' || + !payload.callId || + payload.roomId !== callPanel.roomId || + callPanel.callId === payload.callId + ) { + return; + } + + dispatch( + setModal({ + target: 'callPanel', + data: { + ...callPanel, + callId: payload.callId, + roomType: payload.roomType || callPanel.roomType, + mediaType: payload.mediaType || callPanel.mediaType, + }, + }) + ); + }; + + socket.on('call/started', onStarted); + return () => socket.off('call/started', onStarted); + }, [dispatch, callPanel]); + useEffect(() => { if (!('serviceWorker' in navigator)) return undefined; @@ -125,9 +184,7 @@ function GlobalCallLayer() { }; navigator.serviceWorker.addEventListener('message', onMessage); - return () => { - navigator.serviceWorker.removeEventListener('message', onMessage); - }; + return () => navigator.serviceWorker.removeEventListener('message', onMessage); }, [dispatch, master?._id, callPanel]); useEffect(() => { diff --git a/frontend/client/containers/chat/room.jsx b/frontend/client/containers/chat/room.jsx index c3fd9fcf..aa8083e9 100644 --- a/frontend/client/containers/chat/room.jsx +++ b/frontend/client/containers/chat/room.jsx @@ -49,6 +49,27 @@ function Room() { } }; + const mergePendingRows = (serverRows, currentRows) => { + const rows = Array.isArray(serverRows) ? serverRows : []; + const pending = (currentRows || []).filter( + (item) => item?.pending || item?.sendFailed + ); + if (!pending.length) return rows; + + const serverIds = new Set(rows.map((item) => item?._id).filter(Boolean)); + const serverClientIds = new Set( + rows.map((item) => item?.clientMessageId).filter(Boolean) + ); + return [ + ...rows, + ...pending.filter( + (item) => + !serverIds.has(item?._id) && + !serverClientIds.has(item?.clientMessageId) + ), + ]; + }; + const handleGetChats = async (signal) => { try { const { data } = await axios.get(`/chats/${chatRoom.data.roomId}`, { @@ -57,11 +78,11 @@ function Room() { }); if (data.payload.length > 0) { - setChats(data.payload); + setChats((prev) => mergePendingRows(data.payload, prev)); const callback = (mutationlist, observer) => { const monitor = document.querySelector('#monitor'); - monitor.scrollTop = monitor.scrollHeight; + if (monitor) monitor.scrollTop = monitor.scrollHeight; setLoaded(true); @@ -71,14 +92,16 @@ function Room() { const observer = new MutationObserver(callback); const elem = document.querySelector('#monitor-content'); - observer.observe(elem, { childList: true }); + if (elem) observer.observe(elem, { childList: true }); + else setLoaded(true); return; } + setChats((prev) => mergePendingRows([], prev)); setLoaded(true); } catch (error0) { - console.error(error0.response.data.message); + console.error(error0?.response?.data?.message || error0.message); } }; @@ -149,7 +172,9 @@ function Room() { params: { skip: 0, limit: control.limit }, signal: abortCtrl.signal, }); - setChats(Array.isArray(data?.payload) ? data.payload : []); + setChats((prev) => + mergePendingRows(Array.isArray(data?.payload) ? data.payload : [], prev) + ); } catch (error0) { console.error(error0?.response?.data?.message || error0.message); } finally { @@ -166,6 +191,106 @@ function Room() { }; }, [chatRoom?.data?.roomId, control.limit]); + useEffect(() => { + const roomId = chatRoom?.data?.roomId; + if (!roomId) return undefined; + + const appendWithLimit = (list, payload) => { + const current = Array.isArray(list) ? list : []; + if (current.length >= control.limit) { + return [...current.slice(1), payload]; + } + return [...current, payload]; + }; + + const onOptimistic = (event) => { + const payload = event?.detail; + if (!payload?.clientMessageId || payload.roomId !== roomId) return; + + setChats((prev) => { + const list = Array.isArray(prev) ? prev : []; + if ( + list.some( + (item) => item?.clientMessageId === payload.clientMessageId + ) + ) { + return list; + } + return appendWithLimit(list, payload); + }); + }; + + const onConfirmed = (event) => { + const payload = event?.detail; + if (!payload?._id || payload.roomId !== roomId) return; + + setChats((prev) => { + const list = Array.isArray(prev) ? prev : []; + const index = list.findIndex( + (item) => + item?._id === payload._id || + (payload.clientMessageId && + item?.clientMessageId === payload.clientMessageId) + ); + + if (index < 0) { + return appendWithLimit(list, { + ...payload, + pending: false, + sendFailed: false, + }); + } + + const current = list[index]; + const preserveLocalEncryptedText = + current?.pending && + payload?.e2eeEnvelope && + payload?.text === 'Encrypted message'; + const next = [...list]; + next[index] = { + ...current, + ...payload, + text: preserveLocalEncryptedText ? current.text : payload.text, + pending: false, + sendFailed: false, + }; + return next; + }); + }; + + const onFailed = (event) => { + const payload = event?.detail || {}; + if (payload.roomId && payload.roomId !== roomId) return; + if (!payload.clientMessageId) return; + + setChats((prev) => + (prev || []).map((item) => + item?.clientMessageId === payload.clientMessageId + ? { + ...item, + pending: false, + sendFailed: true, + sendError: + payload.message || payload.code || payload.reason || 'Send failed', + } + : item + ) + ); + }; + + window.addEventListener('syncchat:optimistic-message', onOptimistic); + window.addEventListener('syncchat:message-confirmed', onConfirmed); + window.addEventListener('syncchat:optimistic-message-failed', onFailed); + window.addEventListener('syncchat:outbox-failed', onFailed); + + return () => { + window.removeEventListener('syncchat:optimistic-message', onOptimistic); + window.removeEventListener('syncchat:message-confirmed', onConfirmed); + window.removeEventListener('syncchat:optimistic-message-failed', onFailed); + window.removeEventListener('syncchat:outbox-failed', onFailed); + }; + }, [chatRoom?.data?.roomId, control.limit]); + return (