Soundboard: naming a clip used window.prompt, which Electron does not implement — it returned nothing, the flow aborted in silence, and adding a sound worked in the browser while doing nothing at all in the desktop app. Replaced with a two-step field inside the popover, identical in both. Spotify, three separate defects behind the two symptoms reported: Out of sync — a 20s poll stacked on the activity store's 5s debounce left everyone else on the previous track for up to 25s. The next poll is now scheduled just past the current track's end instead of on a fixed interval, and a track change bypasses the debounce (it happens once every few minutes; the debounce exists for chatty producers). Vanishing — a paused track, and the silent gap Spotify reports between two songs, both cleared the activity outright. Pausing is now carried as state rather than absence, and an empty answer is tolerated for 25s before the block comes down. Progress bar — timestamps are computed with the server's clock and were drawn against the viewer's, so any drift displaced the bar; and it kept advancing after a pause until the next poll. The ready payload now carries server time so each client can correct its own offset, and the bar freezes when paused. Tray, native notifications and system audio in screen share were all found already implemented and wired end to end; recorded in the roadmap rather than built again.
1881 lines
69 KiB
TypeScript
1881 lines
69 KiB
TypeScript
import type { FastifyInstance } from 'fastify';
|
||
import type { WebSocket } from 'ws';
|
||
import { verifyJwt } from '../utils/auth.js';
|
||
import { openVoiceSession, closeVoiceSession } from '../utils/voiceSessions.js';
|
||
import { getDb, schema } from '../db/index.js';
|
||
import { eq, and, or, inArray, isNull, desc, sql } from 'drizzle-orm';
|
||
import { handleClientEvent } from './events.js';
|
||
import { computePermissions, PermissionBits, permissionsToString } from '../utils/permissions.js';
|
||
import type {
|
||
User,
|
||
Space,
|
||
SpaceWithChannelsAndMembers,
|
||
MemberWithUser,
|
||
Channel,
|
||
ChannelCategory,
|
||
DmChannel,
|
||
ServerEvent,
|
||
SpaceFolder,
|
||
SpaceLayoutItem,
|
||
ReadState,
|
||
ActiveCallInfo,
|
||
Activity,
|
||
} from '@backspace/shared';
|
||
import { sanitizeUser } from '../utils/sanitize.js';
|
||
import { collectProfileBroadcastTargetIds } from '../utils/userDeletion.js';
|
||
|
||
// ─── Heartbeat State ──────────────────────────────────────────────────────────
|
||
const wsIsAlive: WeakMap<WebSocket, boolean> = new WeakMap();
|
||
let heartbeatInterval: ReturnType<typeof setInterval> | null = null;
|
||
|
||
// 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<TId, TResult>(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;
|
||
}
|
||
|
||
export interface AuthenticatedSocket {
|
||
ws: WebSocket;
|
||
userId: string;
|
||
username: string;
|
||
}
|
||
|
||
// ─── VoiceRoom Abstraction ─────────────────────────────────────────────────
|
||
|
||
export interface SpaceRoomMeta {
|
||
type: 'space';
|
||
spaceId: string;
|
||
}
|
||
|
||
export interface DmRoomMeta {
|
||
type: 'dm';
|
||
callerId: string;
|
||
state: 'ringing' | 'active';
|
||
}
|
||
|
||
/** In-memory registry for federated calls on REMOTE instances. */
|
||
export interface FederatedCallEntry {
|
||
dmChannelId: string | null; // null for Path B (no local DM), late-bound when DM created mid-call
|
||
federatedId: string; // primary key — cross-instance stable
|
||
callerId: string; // local stub userId of the caller
|
||
callerHomeUserId: string;
|
||
federatedCallHost: string; // peer origin of the host instance
|
||
livekitUrl: string;
|
||
tokens: Map<string, string>; // homeUserId → LiveKit token
|
||
ringedUserIds: string[]; // local userIds that received dm_call_incoming
|
||
state: 'ringing' | 'active';
|
||
startedAt: number;
|
||
}
|
||
|
||
export interface VoiceRoom {
|
||
roomId: string;
|
||
roomType: 'space' | 'dm';
|
||
participants: Set<string>;
|
||
metadata: SpaceRoomMeta | DmRoomMeta;
|
||
startedAt: number;
|
||
}
|
||
|
||
// ─── ConnectionManager ─────────────────────────────────────────────────────
|
||
|
||
class ConnectionManager {
|
||
// userId → Set of WebSocket connections (multiple tabs)
|
||
private connections: Map<string, Set<WebSocket>> = new Map();
|
||
// userId → Set of space IDs the user belongs to
|
||
private userSpaces: Map<string, Set<string>> = new Map();
|
||
// ws → userId (reverse lookup)
|
||
private wsToUser: Map<WebSocket, string> = new Map();
|
||
// Unified voice room tracking (replaces voiceStates + activeCalls)
|
||
private voiceRooms: Map<string, VoiceRoom> = new Map();
|
||
// O(1) reverse index: userId → roomId
|
||
private userToRoom: Map<string, string> = new Map();
|
||
// userId → { isMuted, isDeafened, isCameraOn, isScreenSharing } — voice user status
|
||
private voiceUserStates: Map<string, { isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean }> = new Map();
|
||
// userId → Timeout
|
||
private pendingOfflineTimeouts: Map<string, NodeJS.Timeout> = new Map();
|
||
// roomId → Timeout for ringing DM rooms (60s auto-cleanup)
|
||
private ringingTimeouts: Map<string, NodeJS.Timeout> = new Map();
|
||
// Callback registered by events.ts to fan dm_call_end out to peers on ring timeout.
|
||
// Null during startup — ring timeouts that fire before registration simply no-op (there are no peers to notify before boot completes).
|
||
private ringTimeoutFanoutHook: ((dmChannelId: string, callerId: string) => Promise<void>) | null = null;
|
||
/** Federated calls where this instance is NOT the host. Keyed by federatedId. */
|
||
private federatedCalls: Map<string, FederatedCallEntry> = new Map();
|
||
private federatedCallTimeouts: Map<string, NodeJS.Timeout> = new Map();
|
||
// Space-muted/deafened users (moderator action)
|
||
private spaceMutedUsers: Set<string> = new Set(); // Stores spaceId:userId
|
||
private spaceDeafenedUsers: Set<string> = new Set(); // Stores spaceId:userId
|
||
// Permission-muted users (SPEAK permission revoked while in voice)
|
||
private permissionMutedUsers: Set<string> = new Set(); // Stores spaceId:userId
|
||
// The specific WebSocket that initiated voice_join / DM call for this user.
|
||
// When THIS socket closes, voice state is cleaned up immediately.
|
||
private voiceWs: Map<string, WebSocket> = new Map();
|
||
// Per-user WebSocket rate limiters (shared across all tabs/connections)
|
||
private userRateLimiters: Map<string, WsRateLimiter> = new Map();
|
||
|
||
// ─── Rich Presence ──────────────────────────────────────────────────────
|
||
// userId → Activity[] (ephemeral, same lifecycle as voiceUserStates)
|
||
private userActivities: Map<string, Activity[]> = new Map();
|
||
// userId → boolean (cached from DB at auth time, updated via REST)
|
||
private userShowActivity: Map<string, boolean> = new Map();
|
||
// userId → status string (cached at auth, updated on presence_update)
|
||
private userStatuses: Map<string, string> = new Map();
|
||
// userId → timestamp of last activity_update (rate limiting)
|
||
private lastActivityUpdate: Map<string, number> = new Map();
|
||
|
||
addConnection(userId: string, ws: WebSocket): void {
|
||
if (!this.connections.has(userId)) {
|
||
this.connections.set(userId, new Set());
|
||
}
|
||
this.connections.get(userId)!.add(ws);
|
||
this.wsToUser.set(ws, userId);
|
||
|
||
// If they were pending offline, cancel it!
|
||
this.cancelDisconnect(userId);
|
||
}
|
||
|
||
removeConnection(ws: WebSocket): string | undefined {
|
||
const userId = this.wsToUser.get(ws);
|
||
if (!userId) return undefined;
|
||
|
||
this.wsToUser.delete(ws);
|
||
const userConnections = this.connections.get(userId);
|
||
if (userConnections) {
|
||
userConnections.delete(ws);
|
||
|
||
// ── Immediate voice cleanup if this was the voice-active socket ──
|
||
if (this.voiceWs.get(userId) === ws) {
|
||
this.voiceWs.delete(userId);
|
||
|
||
// Leave voice room (space or DM)
|
||
const left = this.leaveCurrentRoom(userId);
|
||
this.clearVoiceUserStatus(userId);
|
||
if (left) {
|
||
if (left.room.roomType === 'space') {
|
||
const meta = left.room.metadata as SpaceRoomMeta;
|
||
this.sendToSpace(meta.spaceId, {
|
||
type: 'voice_state_update',
|
||
channelId: left.roomId,
|
||
userId,
|
||
action: 'leave',
|
||
});
|
||
} else {
|
||
this.sendToDmMembers(left.roomId, {
|
||
type: 'voice_state_update',
|
||
channelId: left.roomId,
|
||
userId,
|
||
action: 'leave',
|
||
});
|
||
// Auto-end empty active DM calls
|
||
const updatedRoom = this.voiceRooms.get(left.roomId);
|
||
if (updatedRoom && updatedRoom.participants.size === 0
|
||
&& (updatedRoom.metadata as DmRoomMeta).state === 'active') {
|
||
this.destroyRoom(left.roomId);
|
||
this.sendToDmMembers(left.roomId, {
|
||
type: 'dm_call_ended',
|
||
dmChannelId: left.roomId,
|
||
});
|
||
}
|
||
}
|
||
}
|
||
|
||
// Clean up ringing DM rooms where this user is the caller
|
||
for (const [roomId, room] of this.voiceRooms) {
|
||
if (room.roomType === 'dm') {
|
||
const meta = room.metadata as DmRoomMeta;
|
||
if (meta.state === 'ringing' && meta.callerId === userId) {
|
||
this.destroyRoom(roomId);
|
||
this.sendToDmMembers(roomId, {
|
||
type: 'dm_call_ended',
|
||
dmChannelId: roomId,
|
||
});
|
||
}
|
||
}
|
||
}
|
||
|
||
// Notify the user's remaining tabs so their UI updates
|
||
if (userConnections.size > 0 && left) {
|
||
this.sendToUser(userId, {
|
||
type: 'voice_disconnected',
|
||
userId,
|
||
channelId: left.roomId,
|
||
reason: 'session_closed',
|
||
});
|
||
}
|
||
}
|
||
|
||
if (userConnections.size === 0) {
|
||
this.connections.delete(userId);
|
||
// Schedule disconnect cleanup (presence/offline, NOT voice — already handled above)
|
||
this.scheduleDisconnect(userId);
|
||
}
|
||
}
|
||
return userId;
|
||
}
|
||
|
||
private scheduleDisconnect(userId: string) {
|
||
if (this.pendingOfflineTimeouts.has(userId)) return;
|
||
|
||
const timeout = setTimeout(() => {
|
||
this.finalizeDisconnect(userId);
|
||
this.pendingOfflineTimeouts.delete(userId);
|
||
}, 5000); // 5 second grace period
|
||
|
||
this.pendingOfflineTimeouts.set(userId, timeout);
|
||
}
|
||
|
||
private cancelDisconnect(userId: string) {
|
||
const timeout = this.pendingOfflineTimeouts.get(userId);
|
||
if (timeout) {
|
||
clearTimeout(timeout);
|
||
this.pendingOfflineTimeouts.delete(userId);
|
||
console.log(`[ConnectionManager] Rescued session for user ${userId}`);
|
||
}
|
||
}
|
||
|
||
private finalizeDisconnect(userId: string) {
|
||
// Double check they are still offline
|
||
if (this.isUserOnline(userId)) return;
|
||
|
||
console.log(`[ConnectionManager] Finalizing disconnect for user ${userId}`);
|
||
const db = getDb();
|
||
db.update(schema.users).set({ status: 'offline' }).where(eq(schema.users.id, userId)).run();
|
||
|
||
// Leave voice room if in one (handles both space and DM rooms)
|
||
const left = this.leaveCurrentRoom(userId);
|
||
this.clearVoiceUserStatus(userId);
|
||
this.voiceWs.delete(userId);
|
||
if (left) {
|
||
if (left.room.roomType === 'space') {
|
||
const meta = left.room.metadata as SpaceRoomMeta;
|
||
this.sendToSpace(meta.spaceId, {
|
||
type: 'voice_state_update',
|
||
channelId: left.roomId,
|
||
userId: userId,
|
||
action: 'leave',
|
||
});
|
||
} else {
|
||
// DM room — broadcast leave and auto-end if empty
|
||
this.sendToDmMembers(left.roomId, {
|
||
type: 'voice_state_update',
|
||
channelId: left.roomId,
|
||
userId: userId,
|
||
action: 'leave',
|
||
});
|
||
const updatedRoom = this.voiceRooms.get(left.roomId);
|
||
if (updatedRoom && updatedRoom.participants.size === 0 && (updatedRoom.metadata as DmRoomMeta).state === 'active') {
|
||
this.destroyRoom(left.roomId);
|
||
this.sendToDmMembers(left.roomId, {
|
||
type: 'dm_call_ended',
|
||
dmChannelId: left.roomId,
|
||
});
|
||
}
|
||
}
|
||
}
|
||
|
||
// Destroy any ringing DM rooms where this user is the caller
|
||
for (const [roomId, room] of this.voiceRooms) {
|
||
if (room.roomType === 'dm') {
|
||
const meta = room.metadata as DmRoomMeta;
|
||
if (meta.state === 'ringing' && meta.callerId === userId) {
|
||
this.destroyRoom(roomId);
|
||
this.sendToDmMembers(roomId, {
|
||
type: 'dm_call_ended',
|
||
dmChannelId: roomId,
|
||
});
|
||
}
|
||
}
|
||
}
|
||
|
||
// Clear activity state
|
||
this.clearUserActivities(userId);
|
||
this.userShowActivity.delete(userId);
|
||
this.userStatuses.delete(userId);
|
||
this.lastActivityUpdate.delete(userId);
|
||
|
||
// Broadcast offline to friends + DM co-members + space co-members.
|
||
// Mirrors collectProfileBroadcastTargetIds (the recipient set used by
|
||
// user_updated). Two locally-friended users with no shared space now see
|
||
// each other's offline transitions live, instead of being space-only.
|
||
const offlinePayload = {
|
||
type: 'presence_update' as const,
|
||
userId,
|
||
status: 'offline' as const,
|
||
activities: [] as Activity[],
|
||
};
|
||
const offlineTargets = collectProfileBroadcastTargetIds(userId);
|
||
for (const uid of offlineTargets) this.sendToUser(uid, offlinePayload);
|
||
|
||
// S2S: project offline to all active peers (mirrors profile_update fanout).
|
||
// Imported lazily to avoid circular import (federationPresence → db → ws/handler).
|
||
void import('../utils/federationPresence.js').then(({ queuePresenceRelay }) => {
|
||
try { queuePresenceRelay(userId, 'offline', []); } catch (e) { console.warn('[ws] queuePresenceRelay(offline) failed', e); }
|
||
});
|
||
|
||
// Clean up userSpaces (re-populated on next connect via setUserSpaces)
|
||
this.userSpaces.delete(userId);
|
||
|
||
// Clean up per-user rate limiter
|
||
this.userRateLimiters.delete(userId);
|
||
}
|
||
|
||
getUserConnections(userId: string): Set<WebSocket> {
|
||
return this.connections.get(userId) ?? new Set();
|
||
}
|
||
|
||
isUserOnline(userId: string): boolean {
|
||
const conns = this.connections.get(userId);
|
||
return conns !== undefined && conns.size > 0;
|
||
}
|
||
|
||
getUserRateLimiter(userId: string): WsRateLimiter {
|
||
let limiter = this.userRateLimiters.get(userId);
|
||
if (!limiter) {
|
||
limiter = new WsRateLimiter();
|
||
this.userRateLimiters.set(userId, limiter);
|
||
}
|
||
return limiter;
|
||
}
|
||
|
||
// ─── Activity accessors ─────────────────────────────────────────────────
|
||
|
||
setUserActivities(userId: string, activities: Activity[]): void {
|
||
if (activities.length === 0) {
|
||
this.userActivities.delete(userId);
|
||
} else {
|
||
this.userActivities.set(userId, activities);
|
||
}
|
||
}
|
||
|
||
getUserActivities(userId: string): Activity[] {
|
||
return this.userActivities.get(userId) ?? [];
|
||
}
|
||
|
||
clearUserActivities(userId: string): void {
|
||
this.userActivities.delete(userId);
|
||
}
|
||
|
||
setUserShowActivity(userId: string, show: boolean): void {
|
||
this.userShowActivity.set(userId, show);
|
||
}
|
||
|
||
getUserShowActivity(userId: string): boolean {
|
||
return this.userShowActivity.get(userId) ?? true;
|
||
}
|
||
|
||
setUserStatus(userId: string, status: string): void {
|
||
this.userStatuses.set(userId, status);
|
||
}
|
||
|
||
getUserStatus(userId: string): string {
|
||
return this.userStatuses.get(userId) ?? 'offline';
|
||
}
|
||
|
||
checkActivityRateLimit(userId: string): boolean {
|
||
const now = Date.now();
|
||
const last = this.lastActivityUpdate.get(userId) ?? 0;
|
||
if (now - last < 3000) return false;
|
||
this.lastActivityUpdate.set(userId, now);
|
||
return true;
|
||
}
|
||
|
||
setUserSpaces(userId: string, spaceIds: string[]): void {
|
||
this.userSpaces.set(userId, new Set(spaceIds));
|
||
}
|
||
|
||
addUserSpace(userId: string, spaceId: string): void {
|
||
if (!this.userSpaces.has(userId)) {
|
||
this.userSpaces.set(userId, new Set());
|
||
}
|
||
this.userSpaces.get(userId)!.add(spaceId);
|
||
|
||
// A user joining a space mid-session must be bootstrapped with that space's
|
||
// current voice presence. The `ready` payload only carries voice state at
|
||
// connect time (see buildReadyPayload), so without this push, members already
|
||
// sitting in a voice channel stay invisible in the new member's channel
|
||
// sidebar until a full page reload. We deliver a scoped snapshot over the same
|
||
// ordered WebSocket as the `voice_state_update` deltas, so there is no
|
||
// snapshot-vs-stream race (a join/leave that happens after this snapshot is
|
||
// emitted strictly afterwards on the same socket). `addUserSpace` is the single
|
||
// chokepoint every join path funnels through (invite, public join, join-request
|
||
// approval) and is NOT used on reconnect (that path uses setUserSpaces), so this
|
||
// fires exactly once per genuine join. Space creation hits this too but produces
|
||
// an empty snapshot and is skipped below.
|
||
const snapshot = this.buildSpaceVoiceState(spaceId, userId);
|
||
if (Object.keys(snapshot.voiceStates).length === 0
|
||
&& Object.keys(snapshot.spaceVoiceStates).length === 0) {
|
||
return;
|
||
}
|
||
this.sendToUser(userId, {
|
||
type: 'space_voice_state',
|
||
spaceId,
|
||
voiceStates: snapshot.voiceStates,
|
||
voiceRoomStarts: snapshot.voiceRoomStarts,
|
||
voiceUserStates: snapshot.voiceUserStates,
|
||
spaceVoiceStates: snapshot.spaceVoiceStates,
|
||
});
|
||
}
|
||
|
||
getUserSpaces(userId: string): Set<string> {
|
||
return this.userSpaces.get(userId) ?? new Set();
|
||
}
|
||
|
||
/**
|
||
* Build the current voice-presence snapshot for a single space, from the
|
||
* perspective of `userId`:
|
||
* - which voice channels the user can VIEW have participants, and who they are,
|
||
* - each participant's per-user status (mute/deafen/camera/screenshare),
|
||
* - space-level mute/deafen (persisted) + permission-mute (ephemeral)
|
||
* restrictions, keyed `spaceId:userId`.
|
||
*
|
||
* Voice presence is VIEW_CHANNEL-filtered per `computePermissions` exactly as
|
||
* `buildReadyPayload` does — a user must never learn who is sitting in a voice
|
||
* channel they cannot see.
|
||
*
|
||
* Single source of truth shared by `buildReadyPayload` (connect-time bootstrap,
|
||
* looped across all of a user's spaces) and `addUserSpace` (mid-session join
|
||
* push). Keep these two consumers in sync by changing only this method.
|
||
*/
|
||
buildSpaceVoiceState(spaceId: string, userId: string): {
|
||
voiceStates: Record<string, string[]>;
|
||
voiceRoomStarts: Record<string, number>;
|
||
voiceUserStates: Record<string, { isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean }>;
|
||
spaceVoiceStates: Record<string, { spaceMuted: boolean; spaceDeafened: boolean; permissionMuted: boolean }>;
|
||
} {
|
||
const db = getDb();
|
||
const voiceStates: Record<string, string[]> = {};
|
||
const voiceRoomStarts: Record<string, number> = {};
|
||
const voiceUserStates: Record<string, { isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean }> = {};
|
||
const spaceVoiceStates: Record<string, { spaceMuted: boolean; spaceDeafened: boolean; permissionMuted: boolean }> = {};
|
||
|
||
// Who is currently in each of this space's voice channels the user can VIEW.
|
||
const voiceChannels = db.select({ id: schema.channels.id })
|
||
.from(schema.channels)
|
||
.where(and(eq(schema.channels.spaceId, spaceId), eq(schema.channels.type, 'voice')))
|
||
.all();
|
||
for (const ch of voiceChannels) {
|
||
const chPerms = computePermissions(userId, spaceId, ch.id);
|
||
const hasView = (chPerms & PermissionBits.VIEW_CHANNEL) !== 0n || (chPerms & PermissionBits.ADMINISTRATOR) !== 0n;
|
||
if (!hasView) continue;
|
||
const participants = this.getRoomParticipants(ch.id);
|
||
if (participants.size > 0) {
|
||
const ids = Array.from(participants);
|
||
voiceStates[ch.id] = ids;
|
||
const startedAt = this.getRoomStartedAt(ch.id);
|
||
if (startedAt !== null) voiceRoomStarts[ch.id] = startedAt;
|
||
for (const uid of ids) {
|
||
const status = this.getVoiceUserStatus(uid);
|
||
if (status) voiceUserStates[uid] = status;
|
||
}
|
||
}
|
||
}
|
||
|
||
// Space mute/deafen — persisted, authoritative (survives reconnect). These are
|
||
// space-level flags (they do not reveal which channel a user is in), so they
|
||
// are not channel-filtered, mirroring buildReadyPayload.
|
||
const restrictions = db.select()
|
||
.from(schema.voiceRestrictions)
|
||
.where(eq(schema.voiceRestrictions.spaceId, spaceId))
|
||
.all();
|
||
for (const r of restrictions) {
|
||
const key = `${r.spaceId}:${r.userId}`;
|
||
const existing = spaceVoiceStates[key] ?? { spaceMuted: false, spaceDeafened: false, permissionMuted: false };
|
||
if (r.restrictionType === 'mute') existing.spaceMuted = true;
|
||
if (r.restrictionType === 'deafen') existing.spaceDeafened = true;
|
||
spaceVoiceStates[key] = existing;
|
||
}
|
||
// Permission-mute — ephemeral, derived from in-memory state for every
|
||
// participant currently in this space's voice rooms (mirrors buildReadyPayload).
|
||
for (const [, room] of this.voiceRooms) {
|
||
if (room.roomType !== 'space') continue;
|
||
const meta = room.metadata as SpaceRoomMeta;
|
||
if (meta.spaceId !== spaceId) continue;
|
||
for (const participantId of room.participants) {
|
||
if (this.isPermissionMuted(spaceId, participantId)) {
|
||
const key = `${spaceId}:${participantId}`;
|
||
const existing = spaceVoiceStates[key] ?? { spaceMuted: false, spaceDeafened: false, permissionMuted: false };
|
||
existing.permissionMuted = true;
|
||
spaceVoiceStates[key] = existing;
|
||
}
|
||
}
|
||
}
|
||
|
||
return { voiceStates, voiceRoomStarts, voiceUserStates, spaceVoiceStates };
|
||
}
|
||
|
||
// ─── Unified VoiceRoom API ─────────────────────────────────────────────────
|
||
|
||
/** Create a room. Returns false if room already exists. */
|
||
createRoom(roomId: string, roomType: 'space' | 'dm', metadata: SpaceRoomMeta | DmRoomMeta): boolean {
|
||
if (this.voiceRooms.has(roomId)) return false;
|
||
this.voiceRooms.set(roomId, {
|
||
roomId,
|
||
roomType,
|
||
participants: new Set(),
|
||
metadata,
|
||
startedAt: Date.now(),
|
||
});
|
||
return true;
|
||
}
|
||
|
||
/** Register a fan-out callback invoked when a ringing DM room hits its 60s timeout. */
|
||
setRingTimeoutFanoutHook(fn: (dmChannelId: string, callerId: string) => Promise<void>): void {
|
||
this.ringTimeoutFanoutHook = fn;
|
||
}
|
||
|
||
/** Create a DM room in ringing state with 60s auto-cleanup. */
|
||
createDmRoom(dmChannelId: string, callerId: string): boolean {
|
||
const created = this.createRoom(dmChannelId, 'dm', {
|
||
type: 'dm',
|
||
callerId,
|
||
state: 'ringing',
|
||
});
|
||
if (!created) return false;
|
||
|
||
// 60s ringing timeout — auto-destroy if still ringing
|
||
const timeout = setTimeout(() => {
|
||
this.ringingTimeouts.delete(dmChannelId);
|
||
const room = this.voiceRooms.get(dmChannelId);
|
||
if (room && room.roomType === 'dm' && (room.metadata as DmRoomMeta).state === 'ringing') {
|
||
const ringedCallerId = (room.metadata as DmRoomMeta).callerId;
|
||
this.destroyRoom(dmChannelId);
|
||
this.sendToDmMembers(dmChannelId, {
|
||
type: 'dm_call_ended',
|
||
dmChannelId,
|
||
});
|
||
// Fan dm_call_end out to remote peers so stranded Path-A/B ringees exit the ring.
|
||
// Without this, an accept-relay failure → Alice's 60s auto-clean leaves Bob's FederatedCallEntry lingering with no terminal event.
|
||
if (this.ringTimeoutFanoutHook) {
|
||
this.ringTimeoutFanoutHook(dmChannelId, ringedCallerId).catch(err =>
|
||
console.error('[ws] ring-timeout fan-out error:', err),
|
||
);
|
||
}
|
||
}
|
||
}, 60_000);
|
||
this.ringingTimeouts.set(dmChannelId, timeout);
|
||
|
||
return true;
|
||
}
|
||
|
||
/** Transition a DM room from ringing → active. Returns false if not found or not ringing. */
|
||
activateDmRoom(dmChannelId: string): boolean {
|
||
const room = this.voiceRooms.get(dmChannelId);
|
||
if (!room || room.roomType !== 'dm') return false;
|
||
const meta = room.metadata as DmRoomMeta;
|
||
if (meta.state !== 'ringing') return false;
|
||
meta.state = 'active';
|
||
|
||
// Clear ringing timeout
|
||
const timeout = this.ringingTimeouts.get(dmChannelId);
|
||
if (timeout) {
|
||
clearTimeout(timeout);
|
||
this.ringingTimeouts.delete(dmChannelId);
|
||
}
|
||
return true;
|
||
}
|
||
|
||
/** Register a federated call received via S2S. Adds 60s ringing timeout. */
|
||
createFederatedCall(entry: FederatedCallEntry): void {
|
||
this.clearFederatedCall(entry.federatedId);
|
||
this.federatedCalls.set(entry.federatedId, entry);
|
||
|
||
const timeout = setTimeout(() => {
|
||
this.federatedCallTimeouts.delete(entry.federatedId);
|
||
const call = this.federatedCalls.get(entry.federatedId);
|
||
if (call && call.state === 'ringing') {
|
||
this.federatedCalls.delete(entry.federatedId);
|
||
const endEvent = {
|
||
type: 'dm_call_ended',
|
||
dmChannelId: call.dmChannelId,
|
||
federatedCallId: call.federatedId,
|
||
};
|
||
for (const uid of call.ringedUserIds) {
|
||
this.sendToUser(uid, endEvent as ServerEvent);
|
||
}
|
||
}
|
||
}, 60_000);
|
||
this.federatedCallTimeouts.set(entry.federatedId, timeout);
|
||
}
|
||
|
||
/** Get a federated call entry by federatedId (primary lookup). */
|
||
getFederatedCall(federatedId: string): FederatedCallEntry | undefined {
|
||
return this.federatedCalls.get(federatedId);
|
||
}
|
||
|
||
/** Get a federated call entry by local dmChannelId (convenience reverse lookup). */
|
||
getFederatedCallByDmChannel(dmChannelId: string): FederatedCallEntry | undefined {
|
||
for (const entry of this.federatedCalls.values()) {
|
||
if (entry.dmChannelId === dmChannelId) return entry;
|
||
}
|
||
return undefined;
|
||
}
|
||
|
||
/** Transition a federated call from ringing → active. */
|
||
activateFederatedCall(federatedId: string): boolean {
|
||
const call = this.federatedCalls.get(federatedId);
|
||
if (!call || call.state !== 'ringing') return false;
|
||
call.state = 'active';
|
||
const timeout = this.federatedCallTimeouts.get(federatedId);
|
||
if (timeout) {
|
||
clearTimeout(timeout);
|
||
this.federatedCallTimeouts.delete(federatedId);
|
||
}
|
||
return true;
|
||
}
|
||
|
||
/** Remove a federated call entry and clear its timeout. */
|
||
clearFederatedCall(federatedId: string): void {
|
||
this.federatedCalls.delete(federatedId);
|
||
const timeout = this.federatedCallTimeouts.get(federatedId);
|
||
if (timeout) {
|
||
clearTimeout(timeout);
|
||
this.federatedCallTimeouts.delete(federatedId);
|
||
}
|
||
}
|
||
|
||
/**
|
||
* Evict all FederatedCallEntry objects whose federatedCallHost matches the given peer origin.
|
||
* Emits dm_call_undeliverable { phase: 'host_unreachable', terminal: true } to each entry's
|
||
* ringedUserIds, then clears the entry (and its 60s ring timer if still armed).
|
||
*
|
||
* Idempotent: re-invocation with an already-evicted host returns 0.
|
||
* Called from onPeerDeactivated (signal 1) and the 30s sentinel (signal 2 / backstop).
|
||
*/
|
||
evictFederatedCallsForHost(
|
||
peerOrigin: string,
|
||
ctx: {
|
||
reason: 'peer_transient_failure' | 'peer_rejected';
|
||
peerLabel?: string;
|
||
},
|
||
): number {
|
||
const matches: FederatedCallEntry[] = [];
|
||
for (const entry of this.federatedCalls.values()) {
|
||
if (entry.federatedCallHost === peerOrigin) matches.push(entry);
|
||
}
|
||
if (matches.length === 0) return 0;
|
||
|
||
let evicted = 0;
|
||
for (const entry of matches) {
|
||
// Re-check — concurrent teardown may have removed it between collect and broadcast.
|
||
if (!this.federatedCalls.has(entry.federatedId)) continue;
|
||
|
||
const event: ServerEvent = {
|
||
type: 'dm_call_undeliverable',
|
||
dmChannelId: entry.dmChannelId,
|
||
federatedCallId: entry.federatedId,
|
||
terminal: true,
|
||
phase: 'host_unreachable',
|
||
failures: [{
|
||
reason: ctx.reason,
|
||
peerOrigin,
|
||
peerLabel: ctx.peerLabel,
|
||
}],
|
||
};
|
||
|
||
for (const uid of entry.ringedUserIds) {
|
||
this.sendToUser(uid, event);
|
||
}
|
||
|
||
this.clearFederatedCall(entry.federatedId);
|
||
evicted += 1;
|
||
}
|
||
|
||
return evicted;
|
||
}
|
||
|
||
/** Late-bind a dmChannelId onto a Path B FederatedCallEntry. */
|
||
lateBindFederatedCall(federatedId: string, dmChannelId: string): void {
|
||
const call = this.federatedCalls.get(federatedId);
|
||
if (call && call.dmChannelId === null) {
|
||
call.dmChannelId = dmChannelId;
|
||
}
|
||
}
|
||
|
||
/** Expose federated calls for ready payload assembly. */
|
||
getAllFederatedCalls(): Map<string, FederatedCallEntry> {
|
||
return this.federatedCalls;
|
||
}
|
||
|
||
/** Add a user to a room. Enforces one-room-per-user invariant. Returns the room or null if room doesn't exist. */
|
||
joinRoom(roomId: string, userId: string): VoiceRoom | null {
|
||
const room = this.voiceRooms.get(roomId);
|
||
if (!room) return null;
|
||
|
||
// Enforce one-room-per-user invariant: silently remove from old room
|
||
const currentRoomId = this.userToRoom.get(userId);
|
||
if (currentRoomId && currentRoomId !== roomId) {
|
||
const oldRoom = this.voiceRooms.get(currentRoomId);
|
||
if (oldRoom) {
|
||
oldRoom.participants.delete(userId);
|
||
if (oldRoom.participants.size === 0 && oldRoom.roomType === 'space') {
|
||
this.voiceRooms.delete(currentRoomId);
|
||
}
|
||
}
|
||
}
|
||
|
||
room.participants.add(userId);
|
||
this.userToRoom.set(userId, roomId);
|
||
|
||
// Recorded here rather than at the seven call sites that lead into voice:
|
||
// every path — join, move, DM call, reconnect — funnels through this
|
||
// method, so hooking it cannot miss one.
|
||
openVoiceSession({
|
||
spaceId: room.roomType === 'space' ? (room.metadata as SpaceRoomMeta).spaceId : null,
|
||
channelId: roomId,
|
||
userId,
|
||
});
|
||
|
||
return room;
|
||
}
|
||
|
||
/** Remove a user from a specific room. Returns the room or null if not found. */
|
||
leaveRoom(roomId: string, userId: string): VoiceRoom | null {
|
||
const room = this.voiceRooms.get(roomId);
|
||
if (!room || !room.participants.has(userId)) return null;
|
||
|
||
room.participants.delete(userId);
|
||
this.userToRoom.delete(userId);
|
||
|
||
if (room.roomType === 'space') {
|
||
const meta = room.metadata as SpaceRoomMeta;
|
||
this.clearSpaceVoiceState(meta.spaceId, userId);
|
||
}
|
||
|
||
// Auto-cleanup empty space rooms (they're lazy-created)
|
||
if (room.participants.size === 0 && room.roomType === 'space') {
|
||
this.voiceRooms.delete(roomId);
|
||
}
|
||
|
||
return room;
|
||
}
|
||
|
||
/** Leave whatever room the user is in. Returns { roomId, room } or null. */
|
||
leaveCurrentRoom(userId: string): { roomId: string; room: VoiceRoom } | null {
|
||
const roomId = this.userToRoom.get(userId);
|
||
if (!roomId) return null;
|
||
|
||
const room = this.leaveRoom(roomId, userId);
|
||
if (!room) return null;
|
||
|
||
closeVoiceSession(userId);
|
||
|
||
return { roomId, room };
|
||
}
|
||
|
||
/** Destroy a room entirely. Returns displaced userIds. */
|
||
destroyRoom(roomId: string): string[] {
|
||
const room = this.voiceRooms.get(roomId);
|
||
if (!room) return [];
|
||
|
||
const displaced: string[] = [];
|
||
for (const userId of room.participants) {
|
||
this.userToRoom.delete(userId);
|
||
// Destroying a room bypasses leaveCurrentRoom, so these sessions would
|
||
// otherwise stay open until the next restart swept them away.
|
||
closeVoiceSession(userId);
|
||
displaced.push(userId);
|
||
}
|
||
|
||
this.voiceRooms.delete(roomId);
|
||
|
||
// Clear ringing timeout if any
|
||
const timeout = this.ringingTimeouts.get(roomId);
|
||
if (timeout) {
|
||
clearTimeout(timeout);
|
||
this.ringingTimeouts.delete(roomId);
|
||
}
|
||
|
||
return displaced;
|
||
}
|
||
|
||
/** Get a room by ID. */
|
||
getRoom(roomId: string): VoiceRoom | undefined {
|
||
return this.voiceRooms.get(roomId);
|
||
}
|
||
|
||
/** Get participants in a room. */
|
||
getRoomParticipants(roomId: string): Set<string> {
|
||
return this.voiceRooms.get(roomId)?.participants ?? new Set();
|
||
}
|
||
|
||
/** Get the room a user is currently in. Returns { roomId, room } or null. */
|
||
getUserRoom(userId: string): { roomId: string; room: VoiceRoom } | null {
|
||
const roomId = this.userToRoom.get(userId);
|
||
if (!roomId) return null;
|
||
const room = this.voiceRooms.get(roomId);
|
||
if (!room) return null;
|
||
return { roomId, room };
|
||
}
|
||
|
||
/** Read-only access to all rooms. */
|
||
getAllRooms(): Map<string, VoiceRoom> {
|
||
return this.voiceRooms;
|
||
}
|
||
|
||
// ─── Voice User Status (unchanged) ────────────────────────────────────────
|
||
|
||
setVoiceUserStatus(userId: string, isMuted: boolean, isDeafened: boolean, isCameraOn: boolean, isScreenSharing: boolean): void {
|
||
this.voiceUserStates.set(userId, { isMuted, isDeafened, isCameraOn, isScreenSharing });
|
||
}
|
||
|
||
getVoiceUserStatus(userId: string): { isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean } | undefined {
|
||
return this.voiceUserStates.get(userId);
|
||
}
|
||
|
||
clearVoiceUserStatus(userId: string): void {
|
||
this.voiceUserStates.delete(userId);
|
||
}
|
||
|
||
// ─── Voice WebSocket Binding ───────────────────────────────────────────────
|
||
|
||
/** Store which ws owns the voice session for this user. */
|
||
setVoiceWs(userId: string, ws: WebSocket): void {
|
||
this.voiceWs.set(userId, ws);
|
||
}
|
||
|
||
/** Get the voice-owning ws for this user. */
|
||
getVoiceWs(userId: string): WebSocket | undefined {
|
||
return this.voiceWs.get(userId);
|
||
}
|
||
|
||
/** Clear the voice ws binding for this user. */
|
||
clearVoiceWs(userId: string): void {
|
||
this.voiceWs.delete(userId);
|
||
}
|
||
|
||
setSpaceMuted(spaceId: string, userId: string, muted: boolean): void {
|
||
const key = `${spaceId}:${userId}`;
|
||
if (muted) this.spaceMutedUsers.add(key);
|
||
else this.spaceMutedUsers.delete(key);
|
||
}
|
||
|
||
isSpaceMuted(spaceId: string, userId: string): boolean {
|
||
return this.spaceMutedUsers.has(`${spaceId}:${userId}`);
|
||
}
|
||
|
||
setSpaceDeafened(spaceId: string, userId: string, deafened: boolean): void {
|
||
const key = `${spaceId}:${userId}`;
|
||
if (deafened) this.spaceDeafenedUsers.add(key);
|
||
else this.spaceDeafenedUsers.delete(key);
|
||
}
|
||
|
||
isSpaceDeafened(spaceId: string, userId: string): boolean {
|
||
return this.spaceDeafenedUsers.has(`${spaceId}:${userId}`);
|
||
}
|
||
|
||
clearSpaceVoiceState(spaceId: string, userId: string): void {
|
||
this.spaceMutedUsers.delete(`${spaceId}:${userId}`);
|
||
this.spaceDeafenedUsers.delete(`${spaceId}:${userId}`);
|
||
this.permissionMutedUsers.delete(`${spaceId}:${userId}`);
|
||
}
|
||
|
||
setPermissionMuted(spaceId: string, userId: string, muted: boolean): void {
|
||
const key = `${spaceId}:${userId}`;
|
||
if (muted) this.permissionMutedUsers.add(key);
|
||
else this.permissionMutedUsers.delete(key);
|
||
}
|
||
|
||
isPermissionMuted(spaceId: string, userId: string): boolean {
|
||
return this.permissionMutedUsers.has(`${spaceId}:${userId}`);
|
||
}
|
||
|
||
getAllVoiceUserStates(): Map<string, { isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean }> {
|
||
return this.voiceUserStates;
|
||
}
|
||
|
||
// ─── Broadcasting ─────────────────────────────────────────────────────────
|
||
|
||
/** Send to a specific user (all their connections). */
|
||
sendToUser(userId: string, event: ServerEvent): void {
|
||
const connections = this.getUserConnections(userId);
|
||
const message = JSON.stringify(event);
|
||
for (const ws of connections) {
|
||
if (ws.readyState === 1) { // WebSocket.OPEN
|
||
ws.send(message);
|
||
}
|
||
}
|
||
}
|
||
|
||
/** Send to all members of a space. */
|
||
sendToSpace(spaceId: string, event: ServerEvent, excludeUserId?: string): void {
|
||
const message = JSON.stringify(event);
|
||
for (const [userId, spaceIds] of this.userSpaces) {
|
||
if (spaceIds.has(spaceId) && userId !== excludeUserId) {
|
||
const connections = this.getUserConnections(userId);
|
||
for (const ws of connections) {
|
||
if (ws.readyState === 1) {
|
||
ws.send(message);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
/** Send to space members who have VIEW_CHANNEL on the given channel. */
|
||
sendToChannel(spaceId: string, channelId: string, event: ServerEvent, excludeUserId?: string): void {
|
||
const message = JSON.stringify(event);
|
||
for (const [userId, spaceIds] of this.userSpaces) {
|
||
if (spaceIds.has(spaceId) && userId !== excludeUserId) {
|
||
const perms = computePermissions(userId, spaceId, channelId);
|
||
if ((perms & PermissionBits.VIEW_CHANNEL) !== 0n) {
|
||
const connections = this.getUserConnections(userId);
|
||
for (const ws of connections) {
|
||
if (ws.readyState === 1) {
|
||
ws.send(message);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
/** Expose userSpaces iterator for pre-delete viewer collection. */
|
||
getUserSpaceEntries(): IterableIterator<[string, Set<string>]> {
|
||
return this.userSpaces.entries();
|
||
}
|
||
|
||
/** Send to all DM channel members (queries dm_members table). */
|
||
sendToDmMembers(dmChannelId: string, event: ServerEvent, excludeUserId?: string): void {
|
||
const db = getDb();
|
||
const dmMembers = db.select()
|
||
.from(schema.dmMembers)
|
||
.where(eq(schema.dmMembers.dmChannelId, dmChannelId))
|
||
.all();
|
||
|
||
for (const member of dmMembers) {
|
||
if (member.userId !== excludeUserId) {
|
||
this.sendToUser(member.userId, event);
|
||
}
|
||
}
|
||
}
|
||
|
||
/** Send event to users who were ringed for a federated call.
|
||
* ALWAYS uses ringedUserIds, never sendToDmMembers — sendToDmMembers would
|
||
* also reach the caller's replicated stub, causing cross-instance event contamination
|
||
* (the caller's multi-instance WS gets dm_call_accepted with the wrong dmChannelId). */
|
||
sendToFederatedCallUsers(federatedId: string, event: ServerEvent, excludeUserId?: string): void {
|
||
const call = this.federatedCalls.get(federatedId);
|
||
if (!call) return;
|
||
for (const uid of call.ringedUserIds) {
|
||
if (uid !== excludeUserId) {
|
||
this.sendToUser(uid, event);
|
||
}
|
||
}
|
||
}
|
||
|
||
/** Send to a room — routes to sendToSpace (space rooms) or sendToDmMembers (DM rooms). */
|
||
sendToRoom(roomId: string, event: ServerEvent, excludeUserId?: string): void {
|
||
const room = this.voiceRooms.get(roomId);
|
||
if (!room) return;
|
||
|
||
if (room.roomType === 'space') {
|
||
const meta = room.metadata as SpaceRoomMeta;
|
||
this.sendToSpace(meta.spaceId, event, excludeUserId);
|
||
} else {
|
||
this.sendToDmMembers(roomId, event, excludeUserId);
|
||
}
|
||
}
|
||
|
||
/**
|
||
* When the current occupancy of a room began. Null when nobody is in it —
|
||
* empty space rooms are destroyed, which is what makes the call timer reset
|
||
* once the last person leaves.
|
||
*/
|
||
getRoomStartedAt(roomId: string): number | null {
|
||
return this.voiceRooms.get(roomId)?.startedAt ?? null;
|
||
}
|
||
|
||
/**
|
||
* Send only to the people actually inside a room.
|
||
*
|
||
* Distinct from `sendToRoom`, which fans a space room out to the whole
|
||
* space — right for presence updates the sidebar shows, wrong for anything
|
||
* audible: a soundboard clip must reach the call, not everyone online.
|
||
*/
|
||
sendToRoomParticipants(roomId: string, event: ServerEvent): void {
|
||
const room = this.voiceRooms.get(roomId);
|
||
if (!room) return;
|
||
for (const userId of room.participants) {
|
||
this.sendToUser(userId, event);
|
||
}
|
||
}
|
||
|
||
/** Send to all connections of all online users. */
|
||
sendToAll(event: ServerEvent, excludeUserId?: string): void {
|
||
const message = JSON.stringify(event);
|
||
for (const [userId, connections] of this.connections) {
|
||
if (userId !== excludeUserId) {
|
||
for (const ws of connections) {
|
||
if (ws.readyState === 1) {
|
||
ws.send(message);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
/** Send to a specific WebSocket instance (not all of a user's connections). */
|
||
sendToWs(ws: WebSocket, event: ServerEvent): void {
|
||
if (ws.readyState === 1) { // WebSocket.OPEN
|
||
ws.send(JSON.stringify(event));
|
||
}
|
||
}
|
||
|
||
/** Force-disconnect all WebSocket connections for a user (e.g. account deletion). */
|
||
forceDisconnectUser(userId: string): void {
|
||
// Cancel any pending offline timeout
|
||
const timeout = this.pendingOfflineTimeouts.get(userId);
|
||
if (timeout) {
|
||
clearTimeout(timeout);
|
||
this.pendingOfflineTimeouts.delete(userId);
|
||
}
|
||
|
||
// Leave voice room if in one
|
||
const left = this.leaveCurrentRoom(userId);
|
||
this.clearVoiceUserStatus(userId);
|
||
this.voiceWs.delete(userId);
|
||
if (left) {
|
||
if (left.room.roomType === 'space') {
|
||
const meta = left.room.metadata as SpaceRoomMeta;
|
||
this.sendToSpace(meta.spaceId, {
|
||
type: 'voice_state_update',
|
||
channelId: left.roomId,
|
||
userId,
|
||
action: 'leave',
|
||
});
|
||
} else {
|
||
this.sendToDmMembers(left.roomId, {
|
||
type: 'voice_state_update',
|
||
channelId: left.roomId,
|
||
userId,
|
||
action: 'leave',
|
||
});
|
||
}
|
||
}
|
||
|
||
// Destroy any ringing DM rooms where this user is the caller
|
||
for (const [roomId, room] of this.voiceRooms) {
|
||
if (room.roomType === 'dm') {
|
||
const meta = room.metadata as DmRoomMeta;
|
||
if (meta.state === 'ringing' && meta.callerId === userId) {
|
||
this.destroyRoom(roomId);
|
||
this.sendToDmMembers(roomId, {
|
||
type: 'dm_call_ended',
|
||
dmChannelId: roomId,
|
||
});
|
||
}
|
||
}
|
||
}
|
||
|
||
// Clear activity state
|
||
this.clearUserActivities(userId);
|
||
this.userShowActivity.delete(userId);
|
||
this.userStatuses.delete(userId);
|
||
this.lastActivityUpdate.delete(userId);
|
||
|
||
// Close all WebSocket connections
|
||
const connections = this.connections.get(userId);
|
||
if (connections) {
|
||
for (const ws of connections) {
|
||
this.wsToUser.delete(ws);
|
||
try { ws.close(4001, 'Account deleted'); } catch { /* ignore */ }
|
||
}
|
||
this.connections.delete(userId);
|
||
}
|
||
|
||
// Clean up user spaces
|
||
this.userSpaces.delete(userId);
|
||
}
|
||
|
||
getAllOnlineUserIds(): string[] {
|
||
return Array.from(this.connections.keys());
|
||
}
|
||
|
||
getAllConnections(): Map<string, Set<WebSocket>> {
|
||
return this.connections;
|
||
}
|
||
|
||
/** Send an event to all connected admin users. */
|
||
sendToAdmins(event: ServerEvent): void {
|
||
const db = getDb();
|
||
for (const userId of this.connections.keys()) {
|
||
const user = db.select({ isAdmin: schema.users.isAdmin })
|
||
.from(schema.users).where(eq(schema.users.id, userId)).get();
|
||
if (user?.isAdmin === 1) {
|
||
this.sendToUser(userId, event);
|
||
}
|
||
}
|
||
}
|
||
|
||
/** Push a fresh ready payload to a specific user, forcing full store re-sync. */
|
||
pushReadyPayload(userId: string): void {
|
||
const connections = this.getUserConnections(userId);
|
||
if (connections.size === 0) return;
|
||
|
||
const readyData = buildReadyPayload(userId);
|
||
const message = JSON.stringify({ type: 'ready', serverTime: Date.now(), ...readyData });
|
||
for (const ws of connections) {
|
||
if (ws.readyState === 1) {
|
||
ws.send(message);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
export const connectionManager = new ConnectionManager();
|
||
|
||
// ─── WebSocket Rate Limiter (Token Bucket) ─────────────────────────────────
|
||
|
||
class WsRateLimiter {
|
||
private tokens: number;
|
||
private readonly maxTokens: number;
|
||
private readonly refillRate: number; // tokens per second
|
||
private lastRefill: number;
|
||
|
||
constructor(maxTokens = 30, refillRate = 2) {
|
||
this.maxTokens = maxTokens;
|
||
this.tokens = maxTokens;
|
||
this.refillRate = refillRate;
|
||
this.lastRefill = Date.now();
|
||
}
|
||
|
||
consume(): boolean {
|
||
const now = Date.now();
|
||
const elapsed = (now - this.lastRefill) / 1000;
|
||
this.tokens = Math.min(this.maxTokens, this.tokens + elapsed * this.refillRate);
|
||
this.lastRefill = now;
|
||
|
||
if (this.tokens >= 1) {
|
||
this.tokens -= 1;
|
||
return true;
|
||
}
|
||
return false;
|
||
}
|
||
}
|
||
|
||
function buildReadyPayload(userId: string): {
|
||
user: User;
|
||
spaces: SpaceWithChannelsAndMembers[];
|
||
dmChannels: DmChannel[];
|
||
folders: SpaceFolder[];
|
||
spaceLayout: SpaceLayoutItem[] | null;
|
||
layoutUpdatedAt: number | null;
|
||
voiceStates: Record<string, string[]>;
|
||
voiceUserStates: Record<string, { isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean }>;
|
||
spaceVoiceStates: Record<string, { spaceMuted: boolean; spaceDeafened: boolean; permissionMuted: boolean }>;
|
||
readStates: ReadState[];
|
||
activeCalls: ActiveCallInfo[];
|
||
userActivities: Record<string, Activity[]>;
|
||
rejectedPeerOrigins: string[];
|
||
awaitingApprovalPeerOrigins: string[];
|
||
activePeerOrigins: string[];
|
||
pendingApprovalCount: number;
|
||
} {
|
||
const db = getDb();
|
||
|
||
// Get user
|
||
const userRow = db.select().from(schema.users).where(eq(schema.users.id, userId)).get();
|
||
if (!userRow) {
|
||
throw new Error('User not found');
|
||
}
|
||
const user = sanitizeUser(userRow, true);
|
||
const isFederated = !!userRow.homeInstance;
|
||
|
||
// Cache showActivity and status for Rich Presence
|
||
connectionManager.setUserShowActivity(userId, userRow.showActivity !== 0);
|
||
connectionManager.setUserStatus(userId, (userRow.status ?? 'offline') as string);
|
||
|
||
// Get user's space memberships
|
||
const memberships = db.select()
|
||
.from(schema.spaceMembers)
|
||
.where(eq(schema.spaceMembers.userId, userId))
|
||
.all();
|
||
|
||
const spaceIds = memberships.map(m => m.spaceId);
|
||
|
||
const visibleChannelIdSet = new Set<string>();
|
||
const spaces: SpaceWithChannelsAndMembers[] = [];
|
||
|
||
if (spaceIds.length > 0) {
|
||
const spaceRows = db.select()
|
||
.from(schema.spaces)
|
||
.where(inArray(schema.spaces.id, spaceIds))
|
||
.all();
|
||
|
||
// Batch: all channels for all spaces (1 query instead of N)
|
||
const allChannels = batchInArray(
|
||
spaceIds,
|
||
ids => db.select().from(schema.channels).where(inArray(schema.channels.spaceId, ids)).all(),
|
||
);
|
||
const channelsBySpace = new Map<string, (typeof allChannels)>();
|
||
for (const ch of allChannels) {
|
||
let arr = channelsBySpace.get(ch.spaceId);
|
||
if (!arr) { arr = []; channelsBySpace.set(ch.spaceId, arr); }
|
||
arr.push(ch);
|
||
}
|
||
|
||
// Batch: determine which channels are private (VIEW_CHANNEL denied on @everyone)
|
||
// @everyone role ID equals the space ID, so we query for overrides targeting role = spaceId
|
||
const allEveroneOverrides = batchInArray(
|
||
spaceIds,
|
||
ids => db.select().from(schema.channelOverrides).where(
|
||
and(
|
||
eq(schema.channelOverrides.targetType, 'role'),
|
||
inArray(schema.channelOverrides.targetId, ids),
|
||
)
|
||
).all(),
|
||
);
|
||
const privateChannelIds = new Set<string>();
|
||
for (const o of allEveroneOverrides) {
|
||
const denyBits = BigInt(o.deny || '0');
|
||
if ((denyBits & PermissionBits.VIEW_CHANNEL) !== 0n) {
|
||
privateChannelIds.add(o.channelId);
|
||
}
|
||
}
|
||
|
||
// Batch: all categories for all spaces (1 query instead of N)
|
||
const allCategories = batchInArray(
|
||
spaceIds,
|
||
ids => db.select().from(schema.channelCategories).where(inArray(schema.channelCategories.spaceId, ids)).all(),
|
||
);
|
||
const categoriesBySpace = new Map<string, ChannelCategory[]>();
|
||
for (const cat of allCategories) {
|
||
let arr = categoriesBySpace.get(cat.spaceId);
|
||
if (!arr) { arr = []; categoriesBySpace.set(cat.spaceId, arr); }
|
||
arr.push({
|
||
id: cat.id,
|
||
spaceId: cat.spaceId,
|
||
name: cat.name,
|
||
position: cat.position ?? 0,
|
||
createdAt: cat.createdAt,
|
||
});
|
||
}
|
||
|
||
// Batch: last message ID per channel (1 query instead of N×C)
|
||
const allChannelIds = allChannels.map(ch => ch.id);
|
||
const lastMsgMap = new Map<string, string>();
|
||
if (allChannelIds.length > 0) {
|
||
const lastMsgRows = batchInArray(
|
||
allChannelIds,
|
||
ids => db.select({
|
||
channelId: schema.messages.channelId,
|
||
lastId: sql<string>`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 spaceRow of spaceRows) {
|
||
const channels = channelsBySpace.get(spaceRow.id) ?? [];
|
||
|
||
const roles = db.select()
|
||
.from(schema.roles)
|
||
.where(eq(schema.roles.spaceId, spaceRow.id))
|
||
.orderBy(schema.roles.position)
|
||
.all();
|
||
|
||
const memberRows = db.select()
|
||
.from(schema.spaceMembers)
|
||
.where(eq(schema.spaceMembers.spaceId, spaceRow.id))
|
||
.all();
|
||
|
||
const memberUserIds = memberRows.map(m => m.userId);
|
||
const users = memberUserIds.length > 0
|
||
? 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]));
|
||
|
||
const memberRoleRows = db.select()
|
||
.from(schema.memberRoles)
|
||
.where(eq(schema.memberRoles.spaceId, spaceRow.id))
|
||
.all();
|
||
|
||
const members: MemberWithUser[] = memberRows
|
||
.map(m => {
|
||
const u = userMap.get(m.userId);
|
||
if (!u) return null;
|
||
|
||
const assignedRoleIds = memberRoleRows
|
||
.filter(mr => mr.userId === m.userId)
|
||
.map(mr => mr.roleId);
|
||
|
||
const memberRoles = roles
|
||
.filter(r => assignedRoleIds.includes(r.id))
|
||
.map(r => ({
|
||
id: r.id,
|
||
spaceId: r.spaceId,
|
||
name: r.name,
|
||
color: r.color ?? '#b9bbbe',
|
||
position: r.position ?? 0,
|
||
createdAt: r.createdAt,
|
||
}));
|
||
|
||
return {
|
||
spaceId: m.spaceId,
|
||
userId: m.userId,
|
||
nickname: m.nickname,
|
||
joinedAt: m.joinedAt,
|
||
user: sanitizeUser(u),
|
||
roles: memberRoles,
|
||
};
|
||
})
|
||
.filter((m): m is MemberWithUser => m !== null);
|
||
|
||
// Compute space-level permissions for this user
|
||
const spacePerms = computePermissions(userId, spaceRow.id);
|
||
|
||
// Filter channels by VIEW_CHANNEL and attach per-channel permissions
|
||
const visibleChannels: Channel[] = [];
|
||
for (const ch of channels) {
|
||
const chPerms = computePermissions(userId, spaceRow.id, ch.id);
|
||
const hasView = (chPerms & PermissionBits.VIEW_CHANNEL) !== 0n || (chPerms & PermissionBits.ADMINISTRATOR) !== 0n;
|
||
if (hasView) {
|
||
visibleChannelIdSet.add(ch.id);
|
||
visibleChannels.push({
|
||
id: ch.id,
|
||
spaceId: ch.spaceId,
|
||
name: ch.name,
|
||
type: ch.type as Channel['type'],
|
||
topic: ch.topic,
|
||
position: ch.position ?? 0,
|
||
categoryId: ch.categoryId ?? null,
|
||
isPrivate: privateChannelIds.has(ch.id),
|
||
createdAt: ch.createdAt,
|
||
lastMessageId: lastMsgMap.get(ch.id) ?? null,
|
||
myPermissions: permissionsToString(chPerms),
|
||
});
|
||
}
|
||
}
|
||
|
||
spaces.push({
|
||
id: spaceRow.id,
|
||
name: spaceRow.name,
|
||
icon: spaceRow.icon,
|
||
banner: spaceRow.banner ?? null,
|
||
avatarColor: (spaceRow.avatarColor as Space['avatarColor']) ?? null,
|
||
ownerId: spaceRow.ownerId,
|
||
inviteCode: spaceRow.inviteCode,
|
||
visibility: (spaceRow.visibility ?? 'private') as SpaceWithChannelsAndMembers['visibility'],
|
||
description: spaceRow.description ?? null,
|
||
createdAt: spaceRow.createdAt,
|
||
channels: visibleChannels,
|
||
categories: categoriesBySpace.get(spaceRow.id) ?? [],
|
||
members,
|
||
roles: roles.map(r => ({
|
||
id: r.id,
|
||
spaceId: r.spaceId,
|
||
name: r.name,
|
||
color: r.color ?? '#b9bbbe',
|
||
position: r.position ?? 0,
|
||
permissions: r.permissions ?? undefined,
|
||
isEveryone: r.id === spaceRow.id,
|
||
createdAt: r.createdAt,
|
||
})),
|
||
myPermissions: permissionsToString(spacePerms),
|
||
});
|
||
}
|
||
}
|
||
|
||
// Store user's space IDs for broadcasting
|
||
connectionManager.setUserSpaces(userId, spaceIds);
|
||
|
||
// Get DM channels
|
||
const dmMemberships = db.select()
|
||
.from(schema.dmMembers)
|
||
.where(and(
|
||
eq(schema.dmMembers.userId, userId),
|
||
eq(schema.dmMembers.closed, 0),
|
||
))
|
||
.all();
|
||
|
||
const dmChannelIds = dmMemberships.map(dm => dm.dmChannelId);
|
||
const dmChannels: DmChannel[] = [];
|
||
|
||
if (dmChannelIds.length > 0) {
|
||
// Batch: all DM channels (1 query, exclude soft-deleted)
|
||
const allDmChannelRows = batchInArray(
|
||
dmChannelIds,
|
||
ids => db.select().from(schema.dmChannels).where(and(inArray(schema.dmChannels.id, ids), isNull(schema.dmChannels.deletedAt))).all(),
|
||
);
|
||
const dmChannelMap = new Map(allDmChannelRows.map(c => [c.id, c]));
|
||
|
||
// 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(),
|
||
);
|
||
|
||
// 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]));
|
||
|
||
// Batch: last message per DM channel.
|
||
// Two-step approach (same as GET /api/dm): get MAX(created_at) per channel,
|
||
// then fetch the actual message rows matching those timestamps.
|
||
const dmMaxTimestamps = batchInArray(
|
||
dmChannelIds,
|
||
ids => db.select({
|
||
dmChannelId: schema.dmMessages.dmChannelId,
|
||
maxCreatedAt: sql<number>`MAX(${schema.dmMessages.createdAt})`.as('max_created_at'),
|
||
}).from(schema.dmMessages).where(inArray(schema.dmMessages.dmChannelId, ids)).groupBy(schema.dmMessages.dmChannelId).all(),
|
||
);
|
||
const dmLastMsgMap = new Map<string, typeof schema.dmMessages.$inferSelect>();
|
||
if (dmMaxTimestamps.length > 0) {
|
||
const conditions = dmMaxTimestamps.map(t =>
|
||
and(eq(schema.dmMessages.dmChannelId, t.dmChannelId), eq(schema.dmMessages.createdAt, t.maxCreatedAt!))
|
||
);
|
||
const dmLastMessages = db.select().from(schema.dmMessages).where(or(...conditions)).all();
|
||
for (const m of dmLastMessages) {
|
||
if (!dmLastMsgMap.has(m.dmChannelId)) {
|
||
dmLastMsgMap.set(m.dmChannelId, m);
|
||
}
|
||
}
|
||
}
|
||
const dmLastMsgIds = [...dmLastMsgMap.values()].map(m => m.id);
|
||
|
||
// Batch: attachments for last messages (1 query)
|
||
const dmLastMsgAttachments = dmLastMsgIds.length > 0
|
||
? batchInArray(dmLastMsgIds, ids =>
|
||
db.select({
|
||
dmMessageId: schema.attachments.dmMessageId,
|
||
type: schema.attachments.mimetype,
|
||
filename: schema.attachments.originalName,
|
||
}).from(schema.attachments).where(inArray(schema.attachments.dmMessageId, ids)).all()
|
||
)
|
||
: [];
|
||
const dmLastMsgAttachmentMap = new Map<string, Array<{ type: string; filename: string }>>();
|
||
for (const a of dmLastMsgAttachments) {
|
||
if (!a.dmMessageId) continue;
|
||
const arr = dmLastMsgAttachmentMap.get(a.dmMessageId) ?? [];
|
||
arr.push({ type: a.type, filename: a.filename });
|
||
dmLastMsgAttachmentMap.set(a.dmMessageId, arr);
|
||
}
|
||
|
||
// Assemble DM channels with zero additional queries
|
||
for (const dm of dmMemberships) {
|
||
const dmChannel = dmChannelMap.get(dm.dmChannelId);
|
||
if (!dmChannel) continue;
|
||
|
||
const memberRows = allDmMemberRows.filter(m => m.dmChannelId === dm.dmChannelId);
|
||
const members = memberRows
|
||
.map(m => dmUserMap.get(m.userId))
|
||
.filter((u): u is NonNullable<typeof u> => u != null)
|
||
.map(u => sanitizeUser(u));
|
||
|
||
const last = dmLastMsgMap.get(dm.dmChannelId) ?? null;
|
||
|
||
dmChannels.push({
|
||
id: dmChannel.id,
|
||
federatedId: dmChannel.federatedId ?? null,
|
||
ownerId: dmChannel.ownerId ?? null,
|
||
ownerHomeUserId: dmChannel.ownerHomeUserId ?? null,
|
||
ownerHomeInstance: dmChannel.ownerHomeInstance ?? null,
|
||
createdAt: dmChannel.createdAt,
|
||
name: dmChannel.name ?? null,
|
||
icon: dmChannel.icon ?? null,
|
||
metadataUpdatedAt: dmChannel.metadataUpdatedAt ?? 0,
|
||
members,
|
||
lastMessage: last ? {
|
||
id: last.id,
|
||
dmChannelId: last.dmChannelId,
|
||
userId: last.userId,
|
||
content: last.content,
|
||
createdAt: last.createdAt,
|
||
type: last.type === 'system' ? 'system' : 'user',
|
||
attachments: dmLastMsgAttachmentMap.get(last.id) ?? [],
|
||
} : null,
|
||
});
|
||
}
|
||
|
||
}
|
||
|
||
// Include DM channel IDs in the visible set for read state filtering
|
||
for (const dm of dmChannels) {
|
||
visibleChannelIdSet.add(dm.id);
|
||
}
|
||
|
||
// Seed read states for federated users' DM channels that have no existing read state.
|
||
// This handles the bootstrap: DMs existed before cross-instance access was enabled,
|
||
// so the remote instance has no read state history. Mark as read (latest message).
|
||
// Going forward, the S2S read_state_update relay keeps things in sync.
|
||
if (isFederated && dmChannels.length > 0) {
|
||
const dmIds = dmChannels.map(dm => dm.id);
|
||
const existingDmReadStates = batchInArray(
|
||
dmIds,
|
||
ids => db.select({ channelId: schema.readStates.channelId })
|
||
.from(schema.readStates)
|
||
.where(and(eq(schema.readStates.userId, userId), inArray(schema.readStates.channelId, ids)))
|
||
.all(),
|
||
);
|
||
const hasReadState = new Set(existingDmReadStates.map(rs => rs.channelId));
|
||
const now = Date.now();
|
||
for (const dm of dmChannels) {
|
||
if (!hasReadState.has(dm.id) && dm.lastMessage) {
|
||
db.insert(schema.readStates).values({
|
||
userId,
|
||
channelId: dm.id,
|
||
lastReadMessageId: dm.lastMessage.id,
|
||
updatedAt: now,
|
||
}).run();
|
||
}
|
||
}
|
||
}
|
||
|
||
// Get Space Folders
|
||
const folderRows = db.select()
|
||
.from(schema.spaceFolders)
|
||
.where(eq(schema.spaceFolders.userId, userId))
|
||
.orderBy(schema.spaceFolders.position)
|
||
.all();
|
||
|
||
const folders: SpaceFolder[] = [];
|
||
for (const folder of folderRows) {
|
||
const folderSpaceIds = db.select()
|
||
.from(schema.spaceFolderMembers)
|
||
.where(eq(schema.spaceFolderMembers.folderId, folder.id))
|
||
.orderBy(schema.spaceFolderMembers.position)
|
||
.all()
|
||
.map(m => m.spaceId);
|
||
|
||
folders.push({
|
||
id: folder.id,
|
||
userId: folder.userId,
|
||
name: folder.name,
|
||
color: folder.color,
|
||
position: folder.position ?? 0,
|
||
spaceIds: folderSpaceIds,
|
||
});
|
||
}
|
||
|
||
// Get user space layout
|
||
const layoutRow = db.select().from(schema.userSpaceLayout)
|
||
.where(eq(schema.userSpaceLayout.userId, userId)).get();
|
||
const spaceLayout: SpaceLayoutItem[] | null = layoutRow ? JSON.parse(layoutRow.layout) : null;
|
||
const layoutUpdatedAt: number | null = layoutRow?.updatedAt ?? null;
|
||
|
||
// Build voice states — who is currently in voice channels, plus space mute/
|
||
// deafen and permission-mute, across all the user's spaces. Delegates to the
|
||
// shared per-space helper (also used for the mid-session join push in
|
||
// ConnectionManager.addUserSpace) so the two code paths can never diverge.
|
||
// The helper applies the same VIEW_CHANNEL filtering used when building the
|
||
// `spaces` array above.
|
||
const voiceStates: Record<string, string[]> = {};
|
||
const spaceVoiceStates: Record<string, { spaceMuted: boolean; spaceDeafened: boolean; permissionMuted: boolean }> = {};
|
||
for (const space of spaces) {
|
||
const snap = connectionManager.buildSpaceVoiceState(space.id, userId);
|
||
Object.assign(voiceStates, snap.voiceStates);
|
||
Object.assign(spaceVoiceStates, snap.spaceVoiceStates);
|
||
}
|
||
|
||
// Build active calls from user's DM memberships
|
||
const activeCalls: ActiveCallInfo[] = [];
|
||
for (const dm of dmMemberships) {
|
||
const room = connectionManager.getRoom(dm.dmChannelId);
|
||
if (room && room.roomType === 'dm') {
|
||
const dmMeta = room.metadata as DmRoomMeta;
|
||
activeCalls.push({
|
||
dmChannelId: dm.dmChannelId,
|
||
callerId: dmMeta.callerId,
|
||
participants: Array.from(room.participants),
|
||
startedAt: room.startedAt,
|
||
state: dmMeta.state,
|
||
});
|
||
// Inject DM call participants into voiceStates so frontend's generic handler works
|
||
if (room.participants.size > 0) {
|
||
voiceStates[dm.dmChannelId] = Array.from(room.participants);
|
||
}
|
||
}
|
||
}
|
||
|
||
// Resolve this user's homeUserId for token lookup
|
||
const readyUser = db.select({ homeUserId: schema.users.homeUserId })
|
||
.from(schema.users)
|
||
.where(eq(schema.users.id, userId))
|
||
.get();
|
||
const myHomeUserId = readyUser?.homeUserId || userId;
|
||
|
||
// Also include federated calls (this instance is NOT the host)
|
||
for (const [_fedId, fedCall] of connectionManager.getAllFederatedCalls()) {
|
||
const isParticipant = fedCall.ringedUserIds.includes(userId);
|
||
const isDmMember = fedCall.dmChannelId && dmMemberships.some(dm => dm.dmChannelId === fedCall.dmChannelId);
|
||
if (isParticipant || isDmMember) {
|
||
activeCalls.push({
|
||
dmChannelId: fedCall.dmChannelId,
|
||
federatedCallId: fedCall.federatedId,
|
||
callerId: fedCall.callerId,
|
||
participants: [],
|
||
startedAt: fedCall.startedAt,
|
||
state: fedCall.state,
|
||
federatedCallHost: fedCall.federatedCallHost,
|
||
livekitUrl: fedCall.livekitUrl,
|
||
livekitToken: fedCall.tokens.get(myHomeUserId),
|
||
});
|
||
}
|
||
}
|
||
|
||
// Build voice user states — includes both space and DM participants now
|
||
const voiceUserStates: Record<string, { isMuted: boolean; isDeafened: boolean; isCameraOn: boolean; isScreenSharing: boolean }> = {};
|
||
for (const chId of Object.keys(voiceStates)) {
|
||
const usersInChannel = voiceStates[chId];
|
||
if (usersInChannel) {
|
||
for (const uid of usersInChannel) {
|
||
const status = connectionManager.getVoiceUserStatus(uid);
|
||
if (status) {
|
||
voiceUserStates[uid] = status;
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// Fetch read states for unread tracking
|
||
const readStateRows = db.select()
|
||
.from(schema.readStates)
|
||
.where(eq(schema.readStates.userId, userId))
|
||
.all();
|
||
|
||
const readStates: ReadState[] = readStateRows
|
||
.filter(rs => !isFederated || visibleChannelIdSet.has(rs.channelId))
|
||
.map(rs => ({
|
||
channelId: rs.channelId,
|
||
lastReadMessageId: rs.lastReadMessageId,
|
||
}));
|
||
|
||
// Build user activities snapshot for all visible users
|
||
// Auto-inject customStatus as a 'custom' activity for users with no ephemeral activities
|
||
const userActivities: Record<string, Activity[]> = {};
|
||
const seenUserIds = new Set<string>();
|
||
|
||
function collectUserActivities(uid: string, customStatus: string | null) {
|
||
if (seenUserIds.has(uid)) return;
|
||
seenUserIds.add(uid);
|
||
let acts = connectionManager.getUserActivities(uid);
|
||
if (acts.length === 0 && customStatus) {
|
||
acts = [{ type: 'custom', name: customStatus }];
|
||
}
|
||
if (acts.length > 0) {
|
||
userActivities[uid] = acts;
|
||
}
|
||
}
|
||
|
||
for (const space of spaces) {
|
||
for (const member of space.members) {
|
||
collectUserActivities(member.userId, member.user?.customStatus ?? null);
|
||
}
|
||
}
|
||
for (const dm of dmChannels) {
|
||
for (const member of dm.members) {
|
||
collectUserActivities(member.id, member.customStatus ?? null);
|
||
}
|
||
}
|
||
|
||
// Rejected peer origins for unreachable member indicators
|
||
const rejectedPeers = db
|
||
.select({ origin: schema.federationPeers.origin })
|
||
.from(schema.federationPeers)
|
||
.where(eq(schema.federationPeers.status, 'rejected'))
|
||
.all();
|
||
const rejectedPeerOrigins = rejectedPeers.map(p => p.origin);
|
||
|
||
// Awaiting-approval peer origins for softer unreachable indicators
|
||
const awaitingApprovalPeers = db
|
||
.select({ origin: schema.federationPeers.origin })
|
||
.from(schema.federationPeers)
|
||
.where(eq(schema.federationPeers.status, 'awaiting_approval'))
|
||
.all();
|
||
const awaitingApprovalPeerOrigins = awaitingApprovalPeers.map(p => p.origin);
|
||
|
||
// Active peer origins — client uses this allowlist to gate DM events from remote instances
|
||
const activePeers = db
|
||
.select({ origin: schema.federationPeers.origin })
|
||
.from(schema.federationPeers)
|
||
.where(eq(schema.federationPeers.status, 'active'))
|
||
.all();
|
||
const activePeerOrigins = activePeers.map(p => p.origin);
|
||
|
||
// Pending approval count for admin notification
|
||
let pendingApprovalCount = 0;
|
||
if (userRow?.isAdmin === 1) {
|
||
const countResult = db
|
||
.select({ count: sql<number>`count(*)` })
|
||
.from(schema.peerApprovalRequests)
|
||
.get();
|
||
pendingApprovalCount = countResult?.count ?? 0;
|
||
}
|
||
|
||
return { user, spaces, dmChannels, folders, spaceLayout, layoutUpdatedAt, voiceStates, voiceUserStates, spaceVoiceStates, readStates, activeCalls, userActivities, rejectedPeerOrigins, awaitingApprovalPeerOrigins, activePeerOrigins, pendingApprovalCount };
|
||
}
|
||
|
||
export async function registerWebSocket(app: FastifyInstance): Promise<void> {
|
||
app.get('/ws', { websocket: true }, (socket, request) => {
|
||
const ws = socket as unknown as WebSocket;
|
||
let authenticated = false;
|
||
let userId: string | undefined;
|
||
let username: string | undefined;
|
||
let isFederated = false;
|
||
|
||
// Set auth timeout - must authenticate within 10 seconds
|
||
const authTimeout = setTimeout(() => {
|
||
if (!authenticated) {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'Authentication timeout' }));
|
||
ws.close();
|
||
}
|
||
}, 10000);
|
||
|
||
ws.on('message', (data: Buffer | string) => {
|
||
let parsed: Record<string, unknown>;
|
||
try {
|
||
const raw = typeof data === 'string' ? data : data.toString('utf-8');
|
||
parsed = JSON.parse(raw) as Record<string, unknown>;
|
||
} catch {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'Invalid JSON' }));
|
||
return;
|
||
}
|
||
|
||
// Any received message proves liveness
|
||
wsIsAlive.set(ws, true);
|
||
|
||
if (!authenticated) {
|
||
// First message must be auth
|
||
if (parsed.type !== 'auth' || typeof parsed.token !== 'string') {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'First message must be auth' }));
|
||
ws.close();
|
||
return;
|
||
}
|
||
|
||
try {
|
||
const payload = verifyJwt(parsed.token);
|
||
userId = payload.userId;
|
||
username = payload.username;
|
||
|
||
// Reject deleted users and revoked tokens
|
||
const db = getDb();
|
||
const userRow = db.select().from(schema.users).where(eq(schema.users.id, userId)).get();
|
||
if (!userRow || userRow.isDeleted) {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'This account has been deleted' }));
|
||
ws.close();
|
||
return;
|
||
}
|
||
// Token revocation: reject tokens issued before last password change
|
||
if (userRow.passwordChangedAt && payload.iat) {
|
||
if (payload.iat < Math.floor(userRow.passwordChangedAt / 1000)) {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'Token has been revoked' }));
|
||
ws.close();
|
||
return;
|
||
}
|
||
}
|
||
|
||
authenticated = true;
|
||
isFederated = !!userRow.homeInstance;
|
||
clearTimeout(authTimeout);
|
||
|
||
// Update user status to online
|
||
db.update(schema.users).set({ status: 'online' }).where(eq(schema.users.id, userId)).run();
|
||
|
||
// Add connection
|
||
connectionManager.addConnection(userId, ws);
|
||
|
||
// Mark alive for heartbeat detection; browsers auto-respond to ping frames (RFC 6455)
|
||
wsIsAlive.set(ws, true);
|
||
ws.on('pong', () => { wsIsAlive.set(ws, true); });
|
||
|
||
// Build and send ready payload
|
||
const readyData = buildReadyPayload(userId);
|
||
ws.send(JSON.stringify({
|
||
type: 'ready',
|
||
// Lets each client measure its own offset from this server, so activity
|
||
// timestamps computed here render correctly on a machine whose clock drifts.
|
||
serverTime: Date.now(),
|
||
...readyData,
|
||
}));
|
||
|
||
// Broadcast online to friends + DM co-members + space co-members.
|
||
const onlinePayload = { type: 'presence_update' as const, userId, status: 'online' as const };
|
||
const onlineTargets = collectProfileBroadcastTargetIds(userId);
|
||
for (const uid of onlineTargets) connectionManager.sendToUser(uid, onlinePayload);
|
||
|
||
// S2S: project online to all active peers (mirrors profile_update fanout).
|
||
const _uid = userId;
|
||
void import('../utils/federationPresence.js').then(({ queuePresenceRelay }) => {
|
||
try { queuePresenceRelay(_uid, 'online', []); } catch (e) { console.warn('[ws] queuePresenceRelay(online) failed', e); }
|
||
});
|
||
} catch {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'Invalid token' }));
|
||
ws.close();
|
||
}
|
||
return;
|
||
}
|
||
|
||
// Fast-path heartbeat — never reaches business logic
|
||
if (parsed.type === 'ping') {
|
||
ws.send(JSON.stringify({ type: 'pong' }));
|
||
return;
|
||
}
|
||
|
||
// Rate limit all post-auth, non-ping messages (per-user, shared across tabs)
|
||
if (!connectionManager.getUserRateLimiter(userId!).consume()) {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'Rate limited' }));
|
||
return;
|
||
}
|
||
|
||
// Handle authenticated events
|
||
if (userId && username) {
|
||
try {
|
||
handleClientEvent(parsed, userId, username, ws, isFederated);
|
||
} catch (err) {
|
||
app.log.error({ err, eventType: parsed.type, userId }, 'Unhandled error in WS event handler');
|
||
try {
|
||
ws.send(JSON.stringify({ type: 'error', message: 'Internal server error' }));
|
||
} catch { /* ws may already be closed */ }
|
||
}
|
||
}
|
||
});
|
||
|
||
ws.on('close', () => {
|
||
clearTimeout(authTimeout);
|
||
if (userId) {
|
||
connectionManager.removeConnection(ws);
|
||
}
|
||
});
|
||
|
||
ws.on('error', () => {
|
||
clearTimeout(authTimeout);
|
||
});
|
||
});
|
||
|
||
// ─── Heartbeat Sweep ──────────────────────────────────────────────────────
|
||
// Detect dead connections (e.g. PC shut off without TCP FIN).
|
||
// Sends protocol-level ping frames; browsers auto-respond with pong (RFC 6455).
|
||
// Worst-case detection: 30s + 30s + 5s grace = ~65s.
|
||
const HEARTBEAT_INTERVAL_MS = 30_000;
|
||
|
||
heartbeatInterval = setInterval(() => {
|
||
for (const [, userConnections] of connectionManager.getAllConnections()) {
|
||
for (const ws of userConnections) {
|
||
if (wsIsAlive.get(ws) === false) {
|
||
ws.terminate(); // Emits 'close' → removeConnection → scheduleDisconnect → finalizeDisconnect
|
||
continue;
|
||
}
|
||
wsIsAlive.set(ws, false);
|
||
if (ws.readyState === 1) ws.ping();
|
||
}
|
||
}
|
||
}, HEARTBEAT_INTERVAL_MS);
|
||
|
||
app.addHook('onClose', async () => {
|
||
if (heartbeatInterval) {
|
||
clearInterval(heartbeatInterval);
|
||
heartbeatInterval = null;
|
||
}
|
||
});
|
||
}
|