Шаги 1.4 и 1.5 плана production-готовности: - SseConnectionIndex (byOrg/byUser), один heartbeat-таймер, протокол SSE не тронут - getTaskTree/getTaskParentChain: WITH RECURSIVE CTE (было 1+2N запросов) - sync initial/delta: Promise.all по формам - embedding queue/reindex: батч-предзагрузки вместо поштучных запросов - 8 новых тестов (96/96)
98 lines
3.3 KiB
TypeScript
98 lines
3.3 KiB
TypeScript
// Индекс SSE-соединений по организации и пользователю.
|
||
// Позволяет publishEvent доставлять событие только соединениям нужной
|
||
// организации/пользователя, не перебирая все открытые соединения.
|
||
|
||
export interface IndexedSSEConnection {
|
||
id: string;
|
||
userId: number;
|
||
organizationId: number;
|
||
}
|
||
|
||
export class SseConnectionIndex<T extends IndexedSSEConnection> {
|
||
private byId = new Map<string, T>();
|
||
private byOrg = new Map<number, Set<T>>();
|
||
private byUser = new Map<number, Set<T>>();
|
||
|
||
add(connection: T): void {
|
||
this.byId.set(connection.id, connection);
|
||
|
||
let orgSet = this.byOrg.get(connection.organizationId);
|
||
if (!orgSet) {
|
||
orgSet = new Set();
|
||
this.byOrg.set(connection.organizationId, orgSet);
|
||
}
|
||
orgSet.add(connection);
|
||
|
||
let userSet = this.byUser.get(connection.userId);
|
||
if (!userSet) {
|
||
userSet = new Set();
|
||
this.byUser.set(connection.userId, userSet);
|
||
}
|
||
userSet.add(connection);
|
||
}
|
||
|
||
remove(connectionId: string): T | undefined {
|
||
const connection = this.byId.get(connectionId);
|
||
if (!connection) return undefined;
|
||
|
||
this.byId.delete(connectionId);
|
||
|
||
const orgSet = this.byOrg.get(connection.organizationId);
|
||
if (orgSet) {
|
||
orgSet.delete(connection);
|
||
if (orgSet.size === 0) this.byOrg.delete(connection.organizationId);
|
||
}
|
||
|
||
const userSet = this.byUser.get(connection.userId);
|
||
if (userSet) {
|
||
userSet.delete(connection);
|
||
if (userSet.size === 0) this.byUser.delete(connection.userId);
|
||
}
|
||
|
||
return connection;
|
||
}
|
||
|
||
get size(): number {
|
||
return this.byId.size;
|
||
}
|
||
|
||
/** Все соединения (для глобального heartbeat) */
|
||
forEach(callback: (connection: T) => void): void {
|
||
this.byId.forEach(callback);
|
||
}
|
||
|
||
/**
|
||
* Соединения-получатели события. Семантика соответствует прежнему
|
||
* полному перебору с фильтрами tenant isolation:
|
||
* - userId + organizationId — соединения этого пользователя в этой организации;
|
||
* - только userId — все соединения пользователя (в любой организации);
|
||
* - только organizationId — все соединения организации;
|
||
* - ни того ни другого — broadcast на все соединения.
|
||
*/
|
||
getForEventTarget(organizationId?: number, userId?: number): Iterable<T> {
|
||
if (userId !== undefined) {
|
||
const userSet = this.byUser.get(userId);
|
||
if (!userSet) return [];
|
||
if (organizationId === undefined) return userSet;
|
||
const filtered: T[] = [];
|
||
for (const connection of userSet) {
|
||
if (connection.organizationId === organizationId) filtered.push(connection);
|
||
}
|
||
return filtered;
|
||
}
|
||
if (organizationId !== undefined) {
|
||
return this.byOrg.get(organizationId) ?? [];
|
||
}
|
||
return this.byId.values();
|
||
}
|
||
|
||
isUserConnected(userId: number, organizationId: number): boolean {
|
||
const userSet = this.byUser.get(userId);
|
||
if (!userSet) return false;
|
||
for (const connection of userSet) {
|
||
if (connection.organizationId === organizationId) return true;
|
||
}
|
||
return false;
|
||
}
|
||
}
|