109 lines
3.8 KiB
TypeScript
109 lines
3.8 KiB
TypeScript
import { Router } from "express";
|
||
import crypto from 'crypto';
|
||
import { storage } from "../storage";
|
||
import { authenticateToken, type AuthenticatedRequest } from "../middleware/auth.middleware";
|
||
import { sseLimiter, eventBus, type SSEConnection } from "./shared";
|
||
import { withTenant } from "../db";
|
||
|
||
interface SseTokenEntry {
|
||
userId: number;
|
||
organizationId: number;
|
||
expiresAt: number;
|
||
}
|
||
|
||
const sseTokens = new Map<string, SseTokenEntry>();
|
||
const SSE_TOKEN_TTL_MS = 30_000;
|
||
|
||
setInterval(() => {
|
||
const now = Date.now();
|
||
for (const [t, entry] of Array.from(sseTokens.entries())) {
|
||
if (entry.expiresAt < now) sseTokens.delete(t);
|
||
}
|
||
}, 60_000).unref();
|
||
|
||
export function registerSseRoutes(router: Router): void {
|
||
// POST /api/events/token — exchange long-lived JWT for a one-time SSE token
|
||
router.post('/api/events/token',
|
||
authenticateToken,
|
||
(req: AuthenticatedRequest, res) => {
|
||
const token = crypto.randomUUID();
|
||
sseTokens.set(token, {
|
||
userId: req.user!.id,
|
||
organizationId: req.organizationId!,
|
||
expiresAt: Date.now() + SSE_TOKEN_TTL_MS,
|
||
});
|
||
res.json({ token });
|
||
}
|
||
);
|
||
|
||
// GET /api/events — SSE stream, authenticated via one-time token
|
||
router.get('/api/events',
|
||
sseLimiter,
|
||
async (req, res) => {
|
||
try {
|
||
const rawToken = req.query.token as string;
|
||
if (!rawToken) {
|
||
return res.status(401).json({ error: 'Токен не предоставлен' });
|
||
}
|
||
|
||
const now = Date.now();
|
||
const entry = sseTokens.get(rawToken);
|
||
sseTokens.delete(rawToken);
|
||
|
||
if (!entry || entry.expiresAt < now) {
|
||
return res.status(401).json({ error: 'Недействительный или просроченный SSE-токен' });
|
||
}
|
||
|
||
const user = await withTenant(entry.organizationId, () =>
|
||
storage.getUser(entry.userId)
|
||
);
|
||
if (!user || !user.isActive) {
|
||
return res.status(401).json({ error: 'Пользователь не найден или неактивен' });
|
||
}
|
||
|
||
res.writeHead(200, {
|
||
'Content-Type': 'text/event-stream',
|
||
'Cache-Control': 'no-cache',
|
||
'Connection': 'keep-alive',
|
||
'Access-Control-Allow-Origin': '*',
|
||
'Access-Control-Allow-Headers': 'Cache-Control'
|
||
});
|
||
|
||
const connectionId = `${Date.now()}-${Math.random()}`;
|
||
const connection: SSEConnection = {
|
||
id: connectionId,
|
||
userId: user.id,
|
||
organizationId: user.organizationId,
|
||
res
|
||
};
|
||
|
||
eventBus.addConnection(connection);
|
||
|
||
// Восстанавливаем события, пропущенные с момента lastEventId, до отправки connected.
|
||
// Это гарантирует, что клиент не потеряет сообщения во время reconnect.
|
||
// Поддерживаем как нативный заголовок EventSource, так и query-param для ручного reconnect.
|
||
const lastEventId =
|
||
(req.headers['last-event-id'] as string | undefined) ||
|
||
(req.query.lastEventId as string | undefined);
|
||
if (lastEventId) {
|
||
eventBus.replayBufferedEvents(connection, lastEventId);
|
||
}
|
||
|
||
res.write(`data: {"type":"connected","userId":${user.id},"organizationId":${user.organizationId}}\n\n`);
|
||
|
||
req.on('close', () => {
|
||
eventBus.removeConnection(connectionId);
|
||
});
|
||
|
||
req.on('error', (error) => {
|
||
console.error(`SSE connection error ${connectionId}:`, error);
|
||
eventBus.removeConnection(connectionId);
|
||
});
|
||
} catch (error) {
|
||
console.error('SSE endpoint error:', error);
|
||
res.status(500).json({ error: 'Ошибка подключения к событиям' });
|
||
}
|
||
}
|
||
);
|
||
}
|