fix: rate limiting по пользователям из JWT + cleanup WebPush подписок и cooldown sync-push
This commit is contained in:
@@ -1002,6 +1002,29 @@ async function runStartupDataPatches() {
|
||||
log(`Startup: page recovery error for '${seedName}': ${err instanceof Error ? err.message : String(err)}`);
|
||||
}
|
||||
}
|
||||
|
||||
// Cleanup stale web push subscriptions: keep the 3 most recent per user.
|
||||
// This prevents old/reinstalled PWA registrations from accumulating and
|
||||
// receiving unnecessary silent sync pushes.
|
||||
try {
|
||||
const cleanupResult = await db.execute(sql`
|
||||
DELETE FROM web_push_subscriptions
|
||||
WHERE id IN (
|
||||
SELECT id FROM (
|
||||
SELECT id,
|
||||
ROW_NUMBER() OVER (PARTITION BY user_id, organization_id ORDER BY id DESC) AS rn
|
||||
FROM web_push_subscriptions
|
||||
) ranked
|
||||
WHERE rn > 3
|
||||
)
|
||||
`);
|
||||
const deleted = (cleanupResult as any)?.rowCount ?? 0;
|
||||
if (deleted > 0) {
|
||||
log(`Startup: cleaned up ${deleted} stale web push subscription(s)`);
|
||||
}
|
||||
} catch (err) {
|
||||
log(`Startup: web push subscription cleanup error: ${err instanceof Error ? err.message : String(err)}`);
|
||||
}
|
||||
}
|
||||
|
||||
(async () => {
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import rateLimit, { ipKeyGenerator } from 'express-rate-limit';
|
||||
import type { Request } from 'express';
|
||||
import { storage } from '../storage';
|
||||
import { logAudit } from '../utils/audit';
|
||||
import { db } from '../db';
|
||||
@@ -6,6 +7,54 @@ import { formatUserName } from '../utils/formatUserName';
|
||||
import { tasks, forms, roles } from '@shared/schema';
|
||||
import { eq, and, inArray } from 'drizzle-orm';
|
||||
|
||||
// ── Rate-limit helpers ───────────────────────────────────────────────────────
|
||||
// In offices with a single public IP and many users, IP-only rate limiting
|
||||
// puts every user into the same bucket. We extract the caller identity from
|
||||
// the JWT (Authorization header or access_token cookie) so each user/bot has
|
||||
// their own bucket while still keeping the IP as a fallback namespace.
|
||||
|
||||
function extractJwtIdentifier(token: string): string | null {
|
||||
try {
|
||||
const parts = token.split('.');
|
||||
if (parts.length !== 3) return null;
|
||||
const payload = JSON.parse(Buffer.from(parts[1], 'base64url').toString('utf8'));
|
||||
if (typeof payload.userId === 'number') return `u:${payload.userId}`;
|
||||
if (typeof payload.botId === 'number') return `b:${payload.botId}`;
|
||||
if (typeof payload.superAdminId === 'number') return `sa:${payload.superAdminId}`;
|
||||
return null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function getBearerToken(req: Request): string | null {
|
||||
const auth = req.headers['authorization'];
|
||||
if (auth && typeof auth === 'string' && auth.startsWith('Bearer ')) {
|
||||
return auth.slice(7);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function getCookieToken(req: Request): string | null {
|
||||
const cookieHeader = req.headers['cookie'];
|
||||
if (!cookieHeader || typeof cookieHeader !== 'string') return null;
|
||||
const match = cookieHeader.match(/\baccess_token=([^;]+)/);
|
||||
return match ? decodeURIComponent(match[1]) : null;
|
||||
}
|
||||
|
||||
function buildRateLimitKey(req: Request): string {
|
||||
const ip = ipKeyGenerator(req.ip ?? 'unknown');
|
||||
const token = getBearerToken(req) || getCookieToken(req);
|
||||
if (token) {
|
||||
const id = extractJwtIdentifier(token);
|
||||
if (id) return `${ip}:${id}`;
|
||||
// API keys and other opaque tokens: include the first 32 chars so
|
||||
// different keys do not share a bucket.
|
||||
return `${ip}:${token.slice(0, 32)}`;
|
||||
}
|
||||
return ip;
|
||||
}
|
||||
|
||||
// Event Bus для SSE соединений
|
||||
export interface SSEConnection {
|
||||
id: string;
|
||||
@@ -246,24 +295,16 @@ export const refreshLimiter = rateLimit({
|
||||
|
||||
// Общий API лимит — 1500 запросов за 15 минут
|
||||
// Auth, SSE и health check пропускаются — у них свои лимитеры
|
||||
// Для аутентифицированных запросов ключ включает хэш токена, чтобы
|
||||
// пользователи за одним NAT не делили между собой один лимит.
|
||||
// Для аутентифицированных запросов ключ включает идентификатор пользователя
|
||||
// из JWT (Authorization или cookie), чтобы пользователи за одним NAT не делили
|
||||
// между собой один лимит.
|
||||
export const apiLimiter = rateLimit({
|
||||
windowMs: 15 * 60 * 1000,
|
||||
max: 1500,
|
||||
message: { error: 'Превышен лимит запросов к API. Попробуйте через несколько минут.' },
|
||||
standardHeaders: true,
|
||||
legacyHeaders: false,
|
||||
keyGenerator: (req) => {
|
||||
const auth = req.headers['authorization'];
|
||||
const ip = ipKeyGenerator(req.ip ?? 'unknown');
|
||||
if (auth && typeof auth === 'string' && auth.startsWith('Bearer ')) {
|
||||
// Используем первые 16 символов токена как proxy для user identity.
|
||||
// Это позволяет разным пользователям за одним NAT иметь разные бакеты.
|
||||
return ip + ':' + auth.slice(7, 23);
|
||||
}
|
||||
return ip;
|
||||
},
|
||||
keyGenerator: buildRateLimitKey,
|
||||
skip: (req) => {
|
||||
if (req.method === 'HEAD') return true;
|
||||
// req.path внутри app.use('/api', ...) — без префикса /api
|
||||
@@ -286,14 +327,16 @@ export const apiLimiter = rateLimit({
|
||||
},
|
||||
});
|
||||
|
||||
// SSE соединения — не более 600 подключений за час на IP
|
||||
// (несколько компонентов на странице × перезагрузки × HMR в dev)
|
||||
// SSE соединения — не более 600 подключений за час на пользователя/IP.
|
||||
// Как и для API, разделяем бакеты по JWT-идентификатору, чтобы десяток
|
||||
// человек за одним офисным NAT не исчерпал общий лимит.
|
||||
export const sseLimiter = rateLimit({
|
||||
windowMs: 60 * 60 * 1000,
|
||||
max: 600,
|
||||
message: { error: 'Превышен лимит SSE-подключений. Попробуйте через час.' },
|
||||
standardHeaders: true,
|
||||
legacyHeaders: false,
|
||||
keyGenerator: buildRateLimitKey,
|
||||
handler: (req, res) => {
|
||||
logAudit({
|
||||
action: 'rate_limit.sse',
|
||||
|
||||
@@ -480,6 +480,23 @@ export class ContentStorage extends DataTablesStorage {
|
||||
async registerWebPushSubscription(userId: number, organizationId: number, endpoint: string, subscription: string): Promise<void> {
|
||||
await db.delete(webPushSubscriptions).where(eq(webPushSubscriptions.endpoint, endpoint));
|
||||
await db.insert(webPushSubscriptions).values({ userId, organizationId, endpoint, subscription });
|
||||
|
||||
// Keep at most 3 subscriptions per user to prevent stale devices from
|
||||
// accumulating and receiving unnecessary sync pushes.
|
||||
const MAX_SUBSCRIPTIONS_PER_USER = 3;
|
||||
const userSubs = await db
|
||||
.select({ id: webPushSubscriptions.id })
|
||||
.from(webPushSubscriptions)
|
||||
.where(and(
|
||||
eq(webPushSubscriptions.userId, userId),
|
||||
eq(webPushSubscriptions.organizationId, organizationId)
|
||||
))
|
||||
.orderBy(webPushSubscriptions.id);
|
||||
|
||||
if (userSubs.length > MAX_SUBSCRIPTIONS_PER_USER) {
|
||||
const toDelete = userSubs.slice(0, userSubs.length - MAX_SUBSCRIPTIONS_PER_USER).map(s => s.id);
|
||||
await db.delete(webPushSubscriptions).where(inArray(webPushSubscriptions.id, toDelete));
|
||||
}
|
||||
}
|
||||
|
||||
async unregisterWebPushSubscription(endpoint: string): Promise<void> {
|
||||
|
||||
@@ -2,7 +2,7 @@ import { storage } from '../storage';
|
||||
import { webPushService } from '../services/web-push.service';
|
||||
import { eventBus } from '../routes/shared';
|
||||
|
||||
const SYNC_PUSH_COOLDOWN_MS = 60_000;
|
||||
const SYNC_PUSH_COOLDOWN_MS = 5 * 60 * 1000; // 5 minutes
|
||||
const lastPushByOrganization = new Map<number, number>();
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user