diff --git a/packages/server/src/routes/dm.ts b/packages/server/src/routes/dm.ts index 35fcbc2e..a99f7785 100644 --- a/packages/server/src/routes/dm.ts +++ b/packages/server/src/routes/dm.ts @@ -374,24 +374,26 @@ export async function dmRoutes(app: FastifyInstance): Promise { } } - // Create new DM channel + // Create new DM channel with both members atomically const dmChannelId = generateSnowflake(); const now = Date.now(); - db.insert(schema.dmChannels).values({ - id: dmChannelId, - createdAt: now, - }).run(); + db.transaction((tx) => { + tx.insert(schema.dmChannels).values({ + id: dmChannelId, + createdAt: now, + }).run(); - db.insert(schema.dmMembers).values({ - dmChannelId, - userId: request.userId, - }).run(); + tx.insert(schema.dmMembers).values({ + dmChannelId, + userId: request.userId, + }).run(); - db.insert(schema.dmMembers).values({ - dmChannelId, - userId, - }).run(); + tx.insert(schema.dmMembers).values({ + dmChannelId, + userId, + }).run(); + }); const currentUserRow = db.select().from(schema.users).where(eq(schema.users.id, request.userId)).get(); const members = [currentUserRow, targetUser] @@ -791,24 +793,26 @@ export async function dmRoutes(app: FastifyInstance): Promise { const messageId = generateSnowflake(); const now = Date.now(); - db.insert(schema.dmMessages).values({ - id: messageId, - dmChannelId: id, - userId: request.userId, - replyToId: replyToId || null, - content: content?.trim() || null, - createdAt: now, - }).run(); + // Insert message and link attachments atomically + db.transaction((tx) => { + tx.insert(schema.dmMessages).values({ + id: messageId, + dmChannelId: id, + userId: request.userId, + replyToId: replyToId || null, + content: content?.trim() || null, + createdAt: now, + }).run(); - // Link attachments to this DM message - if (attachmentIds && attachmentIds.length > 0) { - for (const attId of attachmentIds) { - db.update(schema.attachments) - .set({ dmMessageId: messageId }) - .where(eq(schema.attachments.id, attId)) - .run(); + if (attachmentIds && attachmentIds.length > 0) { + for (const attId of attachmentIds) { + tx.update(schema.attachments) + .set({ dmMessageId: messageId }) + .where(eq(schema.attachments.id, attId)) + .run(); + } } - } + }); const message = getDmMessageWithUser(messageId); if (!message) { @@ -886,20 +890,20 @@ export async function dmRoutes(app: FastifyInstance): Promise { return reply.code(403).send({ error: 'You can only delete your own messages', statusCode: 403 }); } - // Delete attachments linked to this DM message - db.delete(schema.attachments) - .where(eq(schema.attachments.dmMessageId, id)) - .run(); + // Delete attachments, reactions, and message atomically + db.transaction((tx) => { + tx.delete(schema.attachments) + .where(eq(schema.attachments.dmMessageId, id)) + .run(); - // Delete reactions - db.delete(schema.dmReactions) - .where(eq(schema.dmReactions.dmMessageId, id)) - .run(); + tx.delete(schema.dmReactions) + .where(eq(schema.dmReactions.dmMessageId, id)) + .run(); - // Delete message - db.delete(schema.dmMessages) - .where(eq(schema.dmMessages.id, id)) - .run(); + tx.delete(schema.dmMessages) + .where(eq(schema.dmMessages.id, id)) + .run(); + }); // Broadcast to all DM members const dmMembers = db.select() diff --git a/packages/server/src/routes/messages.ts b/packages/server/src/routes/messages.ts index e09220cf..82477448 100644 --- a/packages/server/src/routes/messages.ts +++ b/packages/server/src/routes/messages.ts @@ -274,24 +274,26 @@ export async function messageRoutes(app: FastifyInstance): Promise { const messageId = generateSnowflake(); const now = Date.now(); - db.insert(schema.messages).values({ - id: messageId, - channelId: id, - userId: request.userId, - replyToId: replyToId || null, - content: content?.trim() || null, - createdAt: now, - }).run(); + // Insert message and link attachments atomically + db.transaction((tx) => { + tx.insert(schema.messages).values({ + id: messageId, + channelId: id, + userId: request.userId, + replyToId: replyToId || null, + content: content?.trim() || null, + createdAt: now, + }).run(); - // Link attachments to message - if (attachmentIds && attachmentIds.length > 0) { - for (const attId of attachmentIds) { - db.update(schema.attachments) - .set({ messageId }) - .where(eq(schema.attachments.id, attId)) - .run(); + if (attachmentIds && attachmentIds.length > 0) { + for (const attId of attachmentIds) { + tx.update(schema.attachments) + .set({ messageId }) + .where(eq(schema.attachments.id, attId)) + .run(); + } } - } + }); const user = db.select().from(schema.users).where(eq(schema.users.id, request.userId)).get(); if (!user) { @@ -415,9 +417,11 @@ export async function messageRoutes(app: FastifyInstance): Promise { return reply.code(403).send({ error: 'You cannot delete this message', statusCode: 403 }); } - // Delete attachments then message - db.delete(schema.attachments).where(eq(schema.attachments.messageId, id)).run(); - db.delete(schema.messages).where(eq(schema.messages.id, id)).run(); + // Delete attachments and message atomically + db.transaction((tx) => { + tx.delete(schema.attachments).where(eq(schema.attachments.messageId, id)).run(); + tx.delete(schema.messages).where(eq(schema.messages.id, id)).run(); + }); // Broadcast deletion connectionManager.sendToServer(serverId, { diff --git a/packages/server/src/routes/servers.ts b/packages/server/src/routes/servers.ts index 54b43aa8..fb747c68 100644 --- a/packages/server/src/routes/servers.ts +++ b/packages/server/src/routes/servers.ts @@ -80,33 +80,33 @@ export async function serverRoutes(app: FastifyInstance): Promise { const now = Date.now(); const inviteCode = generateInviteCode(); - // Create the server - db.insert(schema.servers).values({ - id: serverId, - name: trimmedName, - icon: icon ?? null, - ownerId: request.userId, - inviteCode, - createdAt: now, - }).run(); + // Create server, owner membership, and default channel atomically + db.transaction((tx) => { + tx.insert(schema.servers).values({ + id: serverId, + name: trimmedName, + icon: icon ?? null, + ownerId: request.userId, + inviteCode, + createdAt: now, + }).run(); - // Add owner as member with 'owner' role - db.insert(schema.serverMembers).values({ - serverId, - userId: request.userId, - role: 'owner', - joinedAt: now, - }).run(); + tx.insert(schema.serverMembers).values({ + serverId, + userId: request.userId, + role: 'owner', + joinedAt: now, + }).run(); - // Create default #general text channel - db.insert(schema.channels).values({ - id: channelId, - serverId, - name: 'general', - type: 'text', - position: 0, - createdAt: now, - }).run(); + tx.insert(schema.channels).values({ + id: channelId, + serverId, + name: 'general', + type: 'text', + position: 0, + createdAt: now, + }).run(); + }); const server = db.select().from(schema.servers).where(eq(schema.servers.id, serverId)).get(); if (!server) { @@ -302,10 +302,12 @@ export async function serverRoutes(app: FastifyInstance): Promise { return reply.code(403).send({ error: 'Only the server owner can delete the server', statusCode: 403 }); } - // Delete all channels (messages cascade), members, then server - db.delete(schema.channels).where(eq(schema.channels.serverId, id)).run(); - db.delete(schema.serverMembers).where(eq(schema.serverMembers.serverId, id)).run(); - db.delete(schema.servers).where(eq(schema.servers.id, id)).run(); + // Delete all channels (messages cascade), members, then server atomically + db.transaction((tx) => { + tx.delete(schema.channels).where(eq(schema.channels.serverId, id)).run(); + tx.delete(schema.serverMembers).where(eq(schema.serverMembers.serverId, id)).run(); + tx.delete(schema.servers).where(eq(schema.servers.id, id)).run(); + }); return reply.code(200).send({ success: true }); }); diff --git a/packages/server/src/routes/social.ts b/packages/server/src/routes/social.ts index 93ff49d9..1def8fd2 100644 --- a/packages/server/src/routes/social.ts +++ b/packages/server/src/routes/social.ts @@ -207,15 +207,22 @@ export async function socialRoutes(app: FastifyInstance): Promise { } if (status === 'accepted') { - // Add to friends table + // Insert friend and update request status atomically const now = Date.now(); - db.insert(schema.friends).values({ - userId: friendRequest.fromId, - friendId: friendRequest.toId, - createdAt: now, - }).run(); + db.transaction((tx) => { + tx.insert(schema.friends).values({ + userId: friendRequest.fromId, + friendId: friendRequest.toId, + createdAt: now, + }).run(); - // Get the accepting user's data for the WS event + tx.update(schema.friendRequests) + .set({ status }) + .where(eq(schema.friendRequests.id, id)) + .run(); + }); + + // WS broadcast AFTER transaction commits const acceptingUser = db.select().from(schema.users).where(eq(schema.users.id, request.userId)).get(); if (acceptingUser) { const friend: Friend = { @@ -228,14 +235,14 @@ export async function socialRoutes(app: FastifyInstance): Promise { requestId: id, }); } + } else { + // For declined, just update the status (single write, no transaction needed) + db.update(schema.friendRequests) + .set({ status }) + .where(eq(schema.friendRequests.id, id)) + .run(); } - // Update request status - db.update(schema.friendRequests) - .set({ status }) - .where(eq(schema.friendRequests.id, id)) - .run(); - return reply.code(200).send({ success: true }); }); diff --git a/packages/server/src/ws/handler.ts b/packages/server/src/ws/handler.ts index 1f0a26c0..7d5ab083 100644 --- a/packages/server/src/ws/handler.ts +++ b/packages/server/src/ws/handler.ts @@ -2,7 +2,7 @@ import type { FastifyInstance } from 'fastify'; import type { WebSocket } from 'ws'; import { verifyJwt } from '../utils/auth.js'; import { getDb, schema } from '../db/index.js'; -import { eq, inArray, desc } from 'drizzle-orm'; +import { eq, inArray, desc, sql } from 'drizzle-orm'; import { handleClientEvent } from './events.js'; import type { User, @@ -16,6 +16,19 @@ import type { ActiveCallInfo, } from '@opencord/shared'; +// SQLite's SQLITE_MAX_VARIABLE_NUMBER default is 999. +// Chunk inArray() calls to stay safely under this limit. +const BATCH_CHUNK_SIZE = 500; + +function batchInArray(ids: TId[], queryFn: (chunk: TId[]) => TResult[]): TResult[] { + if (ids.length <= BATCH_CHUNK_SIZE) return queryFn(ids); + const results: TResult[] = []; + for (let i = 0; i < ids.length; i += BATCH_CHUNK_SIZE) { + results.push(...queryFn(ids.slice(i, i + BATCH_CHUNK_SIZE))); + } + return results; +} + function sanitizeUser(row: typeof schema.users.$inferSelect): User { return { id: row.id, @@ -495,11 +508,36 @@ function buildReadyPayload(userId: string): { .where(inArray(schema.servers.id, serverIds)) .all(); + // Batch: all channels for all servers (1 query instead of N) + const allChannels = batchInArray( + serverIds, + ids => db.select().from(schema.channels).where(inArray(schema.channels.serverId, ids)).all(), + ); + const channelsByServer = new Map(); + for (const ch of allChannels) { + let arr = channelsByServer.get(ch.serverId); + if (!arr) { arr = []; channelsByServer.set(ch.serverId, arr); } + arr.push(ch); + } + + // Batch: last message ID per channel (1 query instead of N×C) + const allChannelIds = allChannels.map(ch => ch.id); + const lastMsgMap = new Map(); + if (allChannelIds.length > 0) { + const lastMsgRows = batchInArray( + allChannelIds, + ids => db.select({ + channelId: schema.messages.channelId, + lastId: sql`max(${schema.messages.id})`, + }).from(schema.messages).where(inArray(schema.messages.channelId, ids)).groupBy(schema.messages.channelId).all(), + ); + for (const row of lastMsgRows) { + if (row.lastId) lastMsgMap.set(row.channelId, row.lastId); + } + } + for (const serverRow of serverRows) { - const channels = db.select() - .from(schema.channels) - .where(eq(schema.channels.serverId, serverRow.id)) - .all(); + const channels = channelsByServer.get(serverRow.id) ?? []; const roles = db.select() .from(schema.roles) @@ -514,7 +552,7 @@ function buildReadyPayload(userId: string): { const memberUserIds = memberRows.map(m => m.userId); const users = memberUserIds.length > 0 - ? db.select().from(schema.users).where(inArray(schema.users.id, memberUserIds)).all() + ? batchInArray(memberUserIds, ids => db.select().from(schema.users).where(inArray(schema.users.id, ids)).all()) : []; const userMap = new Map(users.map(u => [u.id, u])); @@ -562,24 +600,16 @@ function buildReadyPayload(userId: string): { ownerId: serverRow.ownerId, inviteCode: serverRow.inviteCode, createdAt: serverRow.createdAt, - channels: channels.map(ch => { - const lastMsg = db.select({ id: schema.messages.id }) - .from(schema.messages) - .where(eq(schema.messages.channelId, ch.id)) - .orderBy(desc(schema.messages.createdAt)) - .limit(1) - .get(); - return { - id: ch.id, - serverId: ch.serverId, - name: ch.name, - type: ch.type as Channel['type'], - topic: ch.topic, - position: ch.position ?? 0, - createdAt: ch.createdAt, - lastMessageId: lastMsg?.id ?? null, - }; - }), + channels: channels.map(ch => ({ + id: ch.id, + serverId: ch.serverId, + name: ch.name, + type: ch.type as Channel['type'], + topic: ch.topic, + position: ch.position ?? 0, + createdAt: ch.createdAt, + lastMessageId: lastMsgMap.get(ch.id) ?? null, + })), members, roles: roles.map(r => ({ id: r.id, @@ -602,47 +632,70 @@ function buildReadyPayload(userId: string): { .where(eq(schema.dmMembers.userId, userId)) .all(); + const dmChannelIds = dmMemberships.map(dm => dm.dmChannelId); const dmChannels: DmChannel[] = []; - for (const dm of dmMemberships) { - const dmChannel = db.select() - .from(schema.dmChannels) - .where(eq(schema.dmChannels.id, dm.dmChannelId)) - .get(); + if (dmChannelIds.length > 0) { + // Batch: all DM channels (1 query) + const allDmChannelRows = batchInArray( + dmChannelIds, + ids => db.select().from(schema.dmChannels).where(inArray(schema.dmChannels.id, ids)).all(), + ); + const dmChannelMap = new Map(allDmChannelRows.map(c => [c.id, c])); - if (!dmChannel) continue; + // Batch: all DM members across all channels (1 query) + const allDmMemberRows = batchInArray( + dmChannelIds, + ids => db.select().from(schema.dmMembers).where(inArray(schema.dmMembers.dmChannelId, ids)).all(), + ); - const dmMemberRows = db.select() - .from(schema.dmMembers) - .where(eq(schema.dmMembers.dmChannelId, dm.dmChannelId)) - .all(); - - const dmMemberUserIds = dmMemberRows.map(m => m.userId); - const dmUsers = dmMemberUserIds.length > 0 - ? db.select().from(schema.users).where(inArray(schema.users.id, dmMemberUserIds)).all() + // Batch: all unique users from DM members (1 query) + const allDmUserIds = [...new Set(allDmMemberRows.map(m => m.userId))]; + const allDmUsers = allDmUserIds.length > 0 + ? batchInArray(allDmUserIds, ids => db.select().from(schema.users).where(inArray(schema.users.id, ids)).all()) : []; + const dmUserMap = new Map(allDmUsers.map(u => [u.id, u])); - // Get last message - const lastMessage = db.select() - .from(schema.dmMessages) - .where(eq(schema.dmMessages.dmChannelId, dm.dmChannelId)) - .orderBy(schema.dmMessages.createdAt) - .all(); + // Batch: last message per DM channel (1 query — fixes the full-table-scan bug) + const dmLastMsgIdRows = batchInArray( + dmChannelIds, + ids => db.select({ + dmChannelId: schema.dmMessages.dmChannelId, + lastId: sql`max(${schema.dmMessages.id})`, + }).from(schema.dmMessages).where(inArray(schema.dmMessages.dmChannelId, ids)).groupBy(schema.dmMessages.dmChannelId).all(), + ); + const dmLastMsgIds = dmLastMsgIdRows.map(r => r.lastId).filter((id): id is string => id != null); + const dmLastMessages = dmLastMsgIds.length > 0 + ? batchInArray(dmLastMsgIds, ids => db.select().from(schema.dmMessages).where(inArray(schema.dmMessages.id, ids)).all()) + : []; + const dmLastMsgMap = new Map(dmLastMessages.map(m => [m.dmChannelId, m])); - const last = lastMessage.length > 0 ? lastMessage[lastMessage.length - 1] : null; + // Assemble DM channels with zero additional queries + for (const dm of dmMemberships) { + const dmChannel = dmChannelMap.get(dm.dmChannelId); + if (!dmChannel) continue; - dmChannels.push({ - id: dmChannel.id, - createdAt: dmChannel.createdAt, - members: dmUsers.map(sanitizeUser), - lastMessage: last ? { - id: last.id, - dmChannelId: last.dmChannelId, - userId: last.userId, - content: last.content, - createdAt: last.createdAt, - } : null, - }); + const memberRows = allDmMemberRows.filter(m => m.dmChannelId === dm.dmChannelId); + const members = memberRows + .map(m => dmUserMap.get(m.userId)) + .filter((u): u is NonNullable => u != null) + .map(sanitizeUser); + + const last = dmLastMsgMap.get(dm.dmChannelId) ?? null; + + dmChannels.push({ + id: dmChannel.id, + createdAt: dmChannel.createdAt, + members, + lastMessage: last ? { + id: last.id, + dmChannelId: last.dmChannelId, + userId: last.userId, + content: last.content, + createdAt: last.createdAt, + } : null, + }); + } } // Get Server Folders diff --git a/packages/web/src/hooks/useLiveKit.ts b/packages/web/src/hooks/useLiveKit.ts index 24bff4be..0776904f 100644 --- a/packages/web/src/hooks/useLiveKit.ts +++ b/packages/web/src/hooks/useLiveKit.ts @@ -25,6 +25,7 @@ import { stopScreenShare, handleScreenShareUnpublished, } from '../utils/screenShare'; +import { getMediaStreamTrack } from '../utils/livekitInternals'; let _activeRoom: Room | null = null; @@ -522,7 +523,7 @@ export function useLiveKit() { const opts = buildScreenShareOptions(screenShareConfig); const screenPub = room.localParticipant.getTrackPublications().find(p => p.source === Track.Source.ScreenShare); if (screenPub?.videoTrack) { - const mediaTrack = (screenPub.videoTrack as any).mediaStreamTrack as MediaStreamTrack; + const mediaTrack = getMediaStreamTrack(screenPub.videoTrack); if (mediaTrack) { await mediaTrack.applyConstraints({ width: { ideal: opts.capture.width }, height: { ideal: opts.capture.height }, frameRate: { ideal: opts.capture.frameRate } }); mediaTrack.contentHint = opts.contentHint; diff --git a/packages/web/src/hooks/useTrackStats.ts b/packages/web/src/hooks/useTrackStats.ts index 1427ceed..55ff42cf 100644 --- a/packages/web/src/hooks/useTrackStats.ts +++ b/packages/web/src/hooks/useTrackStats.ts @@ -1,6 +1,7 @@ import { useState, useEffect, useRef } from 'react'; import { Track } from 'livekit-client'; import { getActiveRoom } from './useLiveKit'; +import { discoverPeerConnections } from '../utils/livekitInternals'; // ── Types ── @@ -96,40 +97,6 @@ function reportKind(report: any): 'audio' | 'video' | null { return null; } -/** - * Discover all unique RTCPeerConnections from the LiveKit Room engine. - * Different livekit-client versions expose the PC at different internal paths. - */ -function discoverPeerConnections(room: any): RTCPeerConnection[] { - const engine = room?.engine; - if (!engine) return []; - - const pcs: RTCPeerConnection[] = []; - const seen = new WeakSet(); - - const tryAdd = (val: any) => { - if (val && typeof val.getStats === 'function' && !seen.has(val)) { - seen.add(val); - pcs.push(val); - } - }; - - // Current livekit-client (1.x+): engine.pcManager.{publisher,subscriber}.pc - tryAdd(engine.pcManager?.publisher?.pc); - tryAdd(engine.pcManager?.subscriber?.pc); - // Private backing field fallback - tryAdd(engine.pcManager?.publisher?._pc); - tryAdd(engine.pcManager?.subscriber?._pc); - // Older livekit-client paths - tryAdd(engine.publisher?.pc); - tryAdd(engine.subscriber?.pc); - // Unified-plan single PC - tryAdd(engine.pc); - tryAdd(room.pc); - - return pcs; -} - function inferSimulcastLayer(width: number | null, height: number | null): string | null { if (height !== null && height > 0) { if (height >= 1000) return 'High'; diff --git a/packages/web/src/hooks/useWebSocket.ts b/packages/web/src/hooks/useWebSocket.ts index 10bb2b23..0b178ea0 100644 --- a/packages/web/src/hooks/useWebSocket.ts +++ b/packages/web/src/hooks/useWebSocket.ts @@ -357,7 +357,7 @@ function connect(): void { if (ws.readyState === WebSocket.OPEN) { ws.send(JSON.stringify({ type: 'ping' })); } - }, 30_000); + }, 15_000); }; ws.onmessage = (e) => { diff --git a/packages/web/src/stores/authStore.ts b/packages/web/src/stores/authStore.ts index 8207d4cd..c83f5db0 100644 --- a/packages/web/src/stores/authStore.ts +++ b/packages/web/src/stores/authStore.ts @@ -1,6 +1,10 @@ import { create } from 'zustand'; import type { User } from '@opencord/shared'; import { api } from '../api/client'; +import { useChatStore } from './chatStore'; +import { useServerStore } from './serverStore'; +import { useSocialStore } from './socialStore'; +import { useVoiceStore } from './voiceStore'; interface AuthState { token: string | null; @@ -48,6 +52,11 @@ export const useAuthStore = create((set, get) => ({ logout: () => { localStorage.removeItem('opencord_token'); + // Clear all user-scoped state to prevent data leaking between sessions + useChatStore.getState().clearAllMessages(); + useServerStore.getState().populateFromReady([], [], []); + useSocialStore.getState().reset(); + useVoiceStore.getState().clearAllVoiceUsers(); set({ token: null, user: null }); }, diff --git a/packages/web/src/stores/chatStore.ts b/packages/web/src/stores/chatStore.ts index 09c424fd..52467ca8 100644 --- a/packages/web/src/stores/chatStore.ts +++ b/packages/web/src/stores/chatStore.ts @@ -5,6 +5,10 @@ import { wsSend } from '../hooks/useWebSocket'; import { isDmChannel, useServerStore } from './serverStore'; import { useAuthStore } from './authStore'; +const MAX_MESSAGES_PER_CHANNEL = 200; +const MAX_CACHED_CHANNELS = 20; +const EVICT_TO_CHANNELS = 15; + interface TypingUser { userId: string; username: string; @@ -27,6 +31,7 @@ interface ChatState { readStates: Map; unreadChannels: Set; realtimeMessageEvents: RealtimeMessageEvent[]; + channelAccessTimes: Map; setCurrentChannel: (channelId: string | null) => void; setReplyTo: (message: MessageWithUser | null) => void; loadMessages: (channelId: string, force?: boolean) => Promise; @@ -64,14 +69,59 @@ export const useChatStore = create((set, get) => ({ readStates: new Map(), unreadChannels: new Set(), realtimeMessageEvents: [], + channelAccessTimes: new Map(), - setCurrentChannel: (channelId) => set({ currentChannelId: channelId }), + setCurrentChannel: (channelId) => { + set((state) => { + const newAccessTimes = new Map(state.channelAccessTimes); + if (channelId) { + newAccessTimes.set(channelId, Date.now()); + } + + // Evict stale channels if we have too many cached + let newMessages = state.messages; + let newHasMore = state.hasMore; + if (state.messages.size > MAX_CACHED_CHANNELS) { + const entries = [...newAccessTimes.entries()] + .filter(([id]) => id !== channelId) + .sort((a, b) => a[1] - b[1]); + const toEvict = state.messages.size - EVICT_TO_CHANNELS; + const evictIds = new Set(entries.slice(0, toEvict).map(([id]) => id)); + if (evictIds.size > 0) { + newMessages = new Map(state.messages); + newHasMore = new Map(state.hasMore); + for (const id of evictIds) { + newMessages.delete(id); + newHasMore.delete(id); + newAccessTimes.delete(id); + } + } + } + + return { + currentChannelId: channelId, + channelAccessTimes: newAccessTimes, + messages: newMessages, + hasMore: newHasMore, + }; + }); + }, setReplyTo: (message) => set({ replyTo: message }), - clearAllMessages: () => set({ messages: new Map(), hasMore: new Map() }), + clearAllMessages: () => set({ + messages: new Map(), + hasMore: new Map(), + typingUsers: new Map(), + readStates: new Map(), + unreadChannels: new Set(), + realtimeMessageEvents: [], + channelAccessTimes: new Map(), + currentChannelId: null, + replyTo: null, + }), loadMessages: async (channelId: string, force?: boolean) => { - if (!force && get().messages.has(channelId)) return; + if (!force && get().hasMore.has(channelId)) return; set({ isLoading: true, loadError: null }); try { const isDm = isDmChannel(channelId); @@ -84,7 +134,9 @@ export const useChatStore = create((set, get) => ({ newMessages.set(channelId, messages as MessageWithUser[]); const newHasMore = new Map(state.hasMore); newHasMore.set(channelId, messages.length >= 50); - return { messages: newMessages, hasMore: newHasMore, isLoading: false, loadError: null }; + const newAccessTimes = new Map(state.channelAccessTimes); + newAccessTimes.set(channelId, Date.now()); + return { messages: newMessages, hasMore: newHasMore, channelAccessTimes: newAccessTimes, isLoading: false, loadError: null }; }); } catch (err) { set({ isLoading: false, loadError: (err as Error).message || 'Failed to load messages' }); @@ -232,7 +284,12 @@ export const useChatStore = create((set, get) => ({ if (!m.id.startsWith('temp_') || m.userId !== message.userId) return true; return m.content !== message.content; }); - newMessages.set(channelId, [...filtered, message]); + let updated = [...filtered, message]; + // Cap per-channel messages to prevent memory growth + if (updated.length > MAX_MESSAGES_PER_CHANNEL) { + updated = updated.slice(updated.length - MAX_MESSAGES_PER_CHANNEL); + } + newMessages.set(channelId, updated); return { messages: newMessages }; }); }, @@ -248,7 +305,12 @@ export const useChatStore = create((set, get) => ({ if (!m.id.startsWith('temp_') || m.userId !== message.userId) return true; return m.content !== message.content; }); - newMessages.set(channelId, [...filtered, message]); + let updated = [...filtered, message]; + // Cap per-channel messages to prevent memory growth + if (updated.length > MAX_MESSAGES_PER_CHANNEL) { + updated = updated.slice(updated.length - MAX_MESSAGES_PER_CHANNEL); + } + newMessages.set(channelId, updated); // Append to realtimeMessageEvents (capped at 50) const newEvents = [...state.realtimeMessageEvents, { channelId, message }]; if (newEvents.length > 50) newEvents.splice(0, newEvents.length - 50); diff --git a/packages/web/src/stores/socialStore.ts b/packages/web/src/stores/socialStore.ts index 0403548c..241c430c 100644 --- a/packages/web/src/stores/socialStore.ts +++ b/packages/web/src/stores/socialStore.ts @@ -18,6 +18,7 @@ interface SocialState { addFriendFromAccepted: (friend: Friend, requestId: string) => void; updateFriendPresence: (userId: string, status: string) => void; removeFriendLocally: (userId: string) => void; + reset: () => void; } export const useSocialStore = create((set, get) => ({ @@ -139,4 +140,6 @@ export const useSocialStore = create((set, get) => ({ ), })); }, + + reset: () => set({ friends: [], requests: [], isLoading: false, error: null }), })); diff --git a/packages/web/src/utils/livekitInternals.ts b/packages/web/src/utils/livekitInternals.ts new file mode 100644 index 00000000..3bf462cb --- /dev/null +++ b/packages/web/src/utils/livekitInternals.ts @@ -0,0 +1,62 @@ +import type { Room } from 'livekit-client'; + +/** + * Discover all unique RTCPeerConnections from the LiveKit Room engine. + * Different livekit-client versions expose the PC at different internal paths. + */ +export function discoverPeerConnections(room: Room): RTCPeerConnection[] { + const engine = (room as any)?.engine; + if (!engine) return []; + + const pcs: RTCPeerConnection[] = []; + const seen = new WeakSet(); + + const tryAdd = (val: any) => { + if (val && typeof val.getStats === 'function' && !seen.has(val)) { + seen.add(val); + pcs.push(val); + } + }; + + // Current livekit-client (1.x+): engine.pcManager.{publisher,subscriber}.pc + tryAdd(engine.pcManager?.publisher?.pc); + tryAdd(engine.pcManager?.subscriber?.pc); + // Private backing field fallback + tryAdd(engine.pcManager?.publisher?._pc); + tryAdd(engine.pcManager?.subscriber?._pc); + // Older livekit-client paths + tryAdd(engine.publisher?.pc); + tryAdd(engine.subscriber?.pc); + // Unified-plan single PC + tryAdd(engine.pc); + tryAdd((room as any).pc); + + return pcs; +} + +/** + * Get the publisher RTCPeerConnection from a LiveKit Room. + * Used by overdrive to inject RTP sender parameters. + */ +export function getPublisherPC(room: Room): RTCPeerConnection | null { + const engine = (room as any)?.engine; + if (!engine) return null; + + return ( + engine.pcManager?.publisher?.pc ?? + engine.pcManager?.publisher?._pc ?? + engine.publisher?.pc ?? + engine.pc ?? + null + ); +} + +/** + * Safely extract the underlying MediaStreamTrack from a LiveKit track object. + * Handles both public `.mediaStreamTrack` and private `._mediaStreamTrack`. + */ +export function getMediaStreamTrack(track: unknown): MediaStreamTrack | null { + if (!track) return null; + const t = track as any; + return t.mediaStreamTrack ?? t._mediaStreamTrack ?? null; +} diff --git a/packages/web/src/utils/screenShare.ts b/packages/web/src/utils/screenShare.ts index b31e6e8c..0aac3914 100644 --- a/packages/web/src/utils/screenShare.ts +++ b/packages/web/src/utils/screenShare.ts @@ -3,6 +3,7 @@ import { useVoiceStore } from '../stores/voiceStore'; import type { ScreenShareConfig } from '../stores/voiceStore'; import { AudioManager } from '../audio/AudioManager'; import { wsSend } from '../hooks/useWebSocket'; +import { getPublisherPC, getMediaStreamTrack } from './livekitInternals'; // --------------------------------------------------------------------------- // Types @@ -87,12 +88,12 @@ export async function applyOverdrive( const pub = room.localParticipant.getTrackPublications().find(p => p.source === source); if (!pub?.track) return; - const engine = (room as any).engine; - const pc = engine?.pcManager?.publisher?.pc || engine?.publisher?.pc || engine?.pc; + const pc = getPublisherPC(room); if (!pc) return; - const senders = (pc as RTCPeerConnection).getSenders(); - const sender = senders.find(s => s.track?.id === (pub.track as any).mediaStreamTrack?.id); + const pubMediaTrack = getMediaStreamTrack(pub.track); + const senders = pc.getSenders(); + const sender = senders.find(s => s.track?.id === pubMediaTrack?.id); if (!sender) return; const params = sender.getParameters();