From 0e38a280d7ab84b0ffbc6287bcfd34872e0bb5d6 Mon Sep 17 00:00:00 2001 From: dev Date: Sat, 26 Sep 2026 15:26:00 +0300 Subject: [PATCH] =?UTF-8?q?feat(redis):=20=D0=BA=D1=8D=D1=88,=20rate=20lim?= =?UTF-8?q?it,=20=D0=B1=D0=B0=D0=BD=D1=8B=20IP=20=D0=B8=20pub/sub=20=D1=87?= =?UTF-8?q?=D0=B5=D1=80=D0=B5=D0=B7=20Redis?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Добавлен сервис redis:7-alpine (AOF, requirepass, maxmemory + allkeys-lru, healthcheck, том redis-data, порт только на 127.0.0.1) и абстракция redis.js по образцу storage.js. Переведено на Redis: - кэш ответов API и настроек (было Map в памяти), инвалидация по префиксу через SCAN + DEL; - rate limit для api/entry/file — общие счётчики вместо MemoryStore; - баны IP и счётчики неудачных входа — с TTL, вместо опроса БД каждую минуту; - кэш сессий (30 с) с invalidateSessions() на каждой мутации users/sessions/ user_branches, иначе деактивированный пользователь сохранил бы доступ; - pub/sub для SSE-событий и мгновенного пробуждения фоновых воркеров вместо ожидания цикла опроса БД. Отказоустойчивость: при недоступном Redis все операции уходят в in-memory backend с той же семантикой, приложение стартует и работает без Redis и возвращается в Redis автоматически. Первое подключение ограничено по времени (REDIS_CONNECT_TIMEOUT_MS, 5 с) — node-redis не отклоняет connect() при недоступном сервере, а повторяет попытки бесконечно. Добавлены тесты: redis.selftest.js (в т.ч. поведение при недоступном сервере) и api.smoketest.js (сквозная проверка API, включая инвалидацию кэша и мгновенную смерть сессии после logout). --- .env.example | 8 + AGENTS.md | 49 ++++- README.md | 62 ++++++- api.smoketest.js | 124 +++++++++++++ docker-compose.yml | 37 ++++ package-lock.json | 98 ++++++++++ package.json | 1 + redis.js | 432 +++++++++++++++++++++++++++++++++++++++++++++ redis.selftest.js | 148 ++++++++++++++++ server.js | 191 +++++++++++++------- worker.js | 22 ++- 11 files changed, 1098 insertions(+), 74 deletions(-) create mode 100644 api.smoketest.js create mode 100644 redis.js create mode 100644 redis.selftest.js diff --git a/.env.example b/.env.example index 22e5cde..121a4c1 100644 --- a/.env.example +++ b/.env.example @@ -4,6 +4,14 @@ DB_PASSWORD=случайная-длинная-строка # Лимит размера загружаемого на восстановление бэкапа, МБ (по умолчанию 500) BACKUP_UPLOAD_LIMIT_MB=500 +# === Redis (кэш, rate limit, баны IP, pub/sub) === +# Пароль Redis. Обязателен, если Redis включён в docker-compose. +REDIS_PASSWORD=замените-на-длинный-секрет +# Префикс ключей — позволяет держать несколько инстансов в одном Redis. +REDIS_PREFIX=whatido +# Лимит памяти Redis. При превышении вытесняются ключи с наименьшим TTL (allkeys-lru). +REDIS_MAXMEMORY=256mb + # Hugging Face token (нужен для закрытых/приватных репозиториев моделей) HUGGINGFACE_TOKEN= # Имя GGUF-файла модели для text-corrector (скачивается с Hugging Face, если отсутствует) diff --git a/AGENTS.md b/AGENTS.md index 0464ea8..2369ee7 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -8,8 +8,8 @@ This document defines how AI agents should work with the WhatIDo codebase. Follo **WhatIDo** — Accounting system for an educational center: attendance journal, student project works, group gallery, detached files, and public showcase pages (share links). -- **Stack**: Node.js 20 + Express, PostgreSQL 16, Docker Compose, S3-совместимое хранилище файлов, Tailscale (Serve/Funnel) -- **Architecture**: Single Express server (`server.js`) + storage abstraction (`storage.js`) + static frontend in `public/` +- **Stack**: Node.js 20 + Express, PostgreSQL 16, Redis 7, Docker Compose, S3-совместимое хранилище файлов, Tailscale (Serve/Funnel) +- **Architecture**: Single Express server (`server.js`) + storage abstraction (`storage.js`) + cache/pub-sub abstraction (`redis.js`) + static frontend in `public/` - **Deployment**: Docker Compose (app + db + s3 + tailscale), bind-mounted uploads, named volumes for Postgres and S3 data - **Auth**: Admin-only via `X-Admin-Token` header (value = `ADMIN_PASSWORD` env var). No user sessions. @@ -47,9 +47,24 @@ This document defines how AI agents should work with the WhatIDo codebase. Follo - **Cache**: `.thumbs` (WebP miniatures) and `.cache` (originals localized for sharp/zip) live inside `uploads/` and are pruned hourly (`STORAGE_CACHE_MAX_AGE_HOURS`) - **Never publish the S3 API port**: only `127.0.0.1` on the host, file access stays behind app auth/rate limits +### 3b. Redis (`redis.js`) +- **Единственная точка доступа**: `createRedis({ url, prefix })` — все операции кэша/счётчиков/pub-sub идут через неё +- **API**: `get`, `set`, `del`, `dropPrefix`, `dropMatch`, `clear`, `wrap`, `incr`, `publish`, `on`, `rateLimitStore`, `info`, `connect`, `close` +- **Graceful fallback — обязательное требование**: при недоступном Redis все операции уходят в in-memory backend с той же семантикой. Приложение обязано стартовать и работать без Redis +- **Первое подключение ограничено по времени** (`REDIS_CONNECT_TIMEOUT_MS`, 5 с): node-redis не отклоняет `connect()` при недоступном сервере, а повторяет попытки бесконечно — без таймаута старт приложения зависнет навсегда +- **Переподключение**: node-redis переподключается сам; по событию `ready` подписки и subscriber-клиент восстанавливаются (`ensureSubscriber`). Не пересоздавать subscriber через `destroy()` — это гонка с внутренним teardown node-redis +- **Ключи**: `get`/`set` сами добавляют namespace (`REDIS_PREFIX`, по умолчанию `whatido`), `dropPrefix`/`dropMatch` тоже. В `rateLimitStore` префикс добавляется один раз в `base` — не применяйте `fullKey` повторно +- **`scanDelete`**: курсор `SCAN` в node-redis v5 обязан быть строкой, числовой `0` вызовет `TypeError`. Возвращаемое значение курсора — тоже строка, сравнивайте с `'0'` +- **`resetTime` в `rateLimitStore.increment` обязан быть `Date`** — express-rate-limit v8 вызывает `resetTime.getTime()` +- **Пабликация всегда отдаёт подписчикам строку** (JSON), независимо от бэкенда — иначе fallback и Redis расходятся по формату +- **Инвалидация — по префиксу** (`SCAN` + `DEL`), точечного удаления по ключу избегайте +- **Ключевые пространства**: `setting:`, `groups:`, `students:`, `entries:`, `stats:`, `dashboard:`, `share:payload:`, `public-settings`, `system-info`, `session:`, `ban:`, `fail:`, `rl:` +- **Сессии**: `loadUserByToken` кэширует пользователя на 30 с. Любая мутация `users` / `sessions` / `user_branches` обязана вызывать `invalidateSessions()` или удалять `session:`, иначе деактивированный пользователь сохранит доступ +- **Секреты**: пароль только в `REDIS_URL` / `REDIS_PASSWORD`, порт 6379 публикуется лишь на `127.0.0.1` + ### 4. API Patterns - **Admin routes**: `requireAdmin` middleware (checks `X-Admin-Token`) -- **Public routes**: `apiLimiter` (300/15min), `entryLimiter` (10/15min), `fileLimiter` (300/15min) +- **Public routes**: `apiLimiter` (300/15min), `entryLimiter` (10/15min), `fileLimiter` (300/15min) — все на `cache.rateLimitStore(...)`, не на `MemoryStore` - **Responses**: JSON, `{ error: 'message' }` on failure, data directly on success - **Pagination**: `limit` / `offset` query params, return `{ items, total }` or `{ entries, total }` - **Filters**: `group_id`, `date_from`, `date_to`, `student_name`, `search`, `deleted` @@ -62,11 +77,13 @@ This document defines how AI agents should work with the WhatIDo codebase. Follo ### 6. Docker / Compose - **Dockerfile**: Node 20 Alpine, installs deps, generates self-signed TLS cert -- **docker-compose.yml**: 3 services (db, app, tailscale) +- **docker-compose.yml**: сервисы `db`, `app`, `redis`, `s3` (+ опционально `tailscale`, `cloudflared`, `text-corrector`, `photo-ai`) - `db`: postgres:16-alpine, healthcheck, init.sql mounted + - `redis`: redis:7-alpine, `--requirepass`, AOF, `maxmemory` + `allkeys-lru`, healthcheck, том `redis-data`, порт только на `127.0.0.1` - `app`: builds from Dockerfile, exposes 3003/3443, mounts uploads - `tailscale`: host network, NET_ADMIN, runs `start-tailscale.sh` (funnel to 127.0.0.1:3443) -- **Env vars** (required): `ADMIN_PASSWORD`, `DB_PASSWORD` +- **Env vars** (required): `ADMIN_PASSWORD`, `DB_PASSWORD`, `REDIS_PASSWORD` +- **Env vars** (optional): `REDIS_PREFIX` (default `whatido`), `REDIS_MAXMEMORY` (default `256mb`), `REDIS_CONNECT_TIMEOUT_MS` (default `5000`) - **Port 443 on host** must be free (tailscale listens directly) ### 7. Tailscale Publication @@ -144,8 +161,18 @@ curl -H "X-Admin-Token: $ADMIN_PASSWORD" http://localhost:3003/api/groups docker compose up -d s3 docker compose exec -T app node scripts/migrate-to-s3.js --dry-run docker compose exec -T app node scripts/migrate-to-s3.js --verify-only + +# Redis checks +node redis.selftest.js # unit + degradation, needs redis on 127.0.0.1:6379 +node api.smoketest.js # e2e, needs running stack +docker compose exec redis redis-cli -a "$REDIS_PASSWORD" --no-auth-warning INFO +docker compose stop redis && node api.smoketest.js # app must keep working in-memory +docker compose start redis # app reconnects on its own ``` +Verify Redis state through `GET /api/system-info` → `cache` (`driver`, `ready`, `hits`, `misses`, +`fallbackOps`, `used_memory_human`, `keys`). + --- ## Security Checklist (before any change) @@ -166,6 +193,9 @@ docker compose exec -T app node scripts/migrate-to-s3.js --verify-only |------|---------| | `server.js` | Entire backend (Express, routes, DB, uploads, backup) | | `storage.js` | Storage abstraction: `local` and `s3` drivers, key normalization, cache/thumb helpers | +| `redis.js` | Redis abstraction: cache, counters, rate-limit store, pub/sub, in-memory fallback | +| `redis.selftest.js` | Self-tests for `redis.js`, including behaviour with Redis unavailable | +| `api.smoketest.js` | End-to-end API smoke test against a running stack | | `worker.js` | Background AI auto-check worker for entry messages + photo enhance worker | | `db/init.sql` | Initial schema (runs on fresh DB) | | `db/migration.sql` | Idempotent migrations for existing DBs | @@ -191,6 +221,11 @@ docker compose exec -T app node scripts/migrate-to-s3.js --verify-only - ❌ Touch `uploads/` with `fs.*` in request/worker code — use `storage.*` (files may live only in S3) - ❌ Run `migrate-to-s3.js --delete-local` before verification and cutover - ❌ Expose the S3 API port publicly (only `127.0.0.1` in compose) +- ❌ Expose the Redis port publicly (only `127.0.0.1` in compose) +- ❌ Call `fs.*`/`pg` directly for cache, counters or pub/sub — use `redis.js` +- ❌ Make Redis a hard dependency: any new Redis-backed path must keep the in-memory fallback +- ❌ `await client.connect()` without a timeout — it never rejects while Redis is unreachable +- ❌ Cache authorization-relevant data without an invalidation path on the mutation - ❌ Commit `.env`, `certs/`, `uploads/`, `backups/`, `node_modules/` - ❌ Expose DB port (5432) outside docker network - ❌ Use `eval`, `Function` constructor, or dynamic code execution @@ -213,6 +248,10 @@ docker compose logs -f app # DB shell docker compose exec db psql -U app -d whereldo +# Redis status and cache keys +docker compose exec redis redis-cli -a "$REDIS_PASSWORD" --no-auth-warning DBSIZE +docker compose exec redis redis-cli -a "$REDIS_PASSWORD" --no-auth-warning KEYS 'whatido:*' + # S3 storage status and migration verification docker compose up -d s3 docker compose exec -T app node scripts/migrate-to-s3.js --verify-only diff --git a/README.md b/README.md index 72d3bef..884cdc8 100644 --- a/README.md +++ b/README.md @@ -22,6 +22,7 @@ - PostgreSQL (pg) - Multer (загрузка файлов), Tar (бэкапы) - S3-совместимое хранилище (AWS SDK v3): сервис `s3` (SeaweedFS / MinIO) +- Redis: кэш, rate limit, баны IP, кэш сессий, pub/sub (SSE и воркеры) - Lucide (иконки UI) - Docker / Docker Compose - Tailscale (Serve / Funnel) — публикация по HTTPS @@ -50,6 +51,7 @@ docker compose up -d --build - **HTTP** `http://localhost:3003` — редирект на HTTPS - **HTTPS** `https://localhost:3443` — приложение (самоподписанный сертификат, примите предупреждение браузера) - **PostgreSQL** — доступен только внутри docker-сети (наружу не публикуется) +- **Redis** — `127.0.0.1:6379` на хосте (только loopback), внутри сети — `redis:6379` Управление: @@ -61,6 +63,8 @@ docker compose down # остановка (данные сохраняю > Приложение **не запустится** без `ADMIN_PASSWORD` (защита от пароля по умолчанию). > `DB_PASSWORD` задаёт пароль пользователя `app` в PostgreSQL. +> `REDIS_PASSWORD` задаёт пароль Redis. Если сервис `redis` убрать из `docker-compose.yml` +> или оставить `REDIS_URL` пустым — приложение продолжит работать на in-memory кэше. ## Обновление на сервере (деплой) @@ -97,12 +101,16 @@ docker compose exec app md5sum /app/server.js # совпадает с md5sum |------------------|--------------------|-------------------------------------| | `ADMIN_PASSWORD` | — (обязательно) | Пароль администратора (X-Admin-Token). Без него сервер не стартует | | `DB_PASSWORD` | — (обязательно) | Пароль пользователя `app` в PostgreSQL | +| `REDIS_PASSWORD` | — (обязательно) | Пароль Redis (`--requirepass`) | +| `REDIS_PREFIX` | `whatido` | Префикс ключей Redis — свой для каждого инстанса | +| `REDIS_MAXMEMORY` | `256mb` | Лимит памяти Redis, при переполнении вытесняется LRU | Пример `.env` (в репозитории — `.env.example`): ``` ADMIN_PASSWORD=сложный-пароль DB_PASSWORD=случайная-длинная-строка +REDIS_PASSWORD=случайная-длинная-строка ``` `DB_PASSWORD` подставляется в `docker-compose.yml` в `POSTGRES_PASSWORD` и `DATABASE_URL`. Если БД уже была инициализирована ранее, значение `DB_PASSWORD` должно совпадать с фактическим паролем пользователя `app` в БД (иначе приложение не подключится). @@ -312,8 +320,57 @@ docker compose exec -T app node scripts/migrate-to-s3.js --delete-local Объём и состав хранилища видны в админке: Настройки → Системная информация (блок «Хранилище»). +## Redis (кэш и pub/sub) + +Сервис `redis` в compose хранит всё, что не требуется переживать перезапуск Postgres, но должно +быть общим и быстрым: + +| Что | Ключи | TTL | +|---|---|---| +| Кэш ответов API и настроек | `setting:*`, `groups:*`, `students:*`, `entries:*`, `stats:*`, `dashboard:*`, `share:payload:*`, `public-settings`, `system-info` | 15–60 с | +| Кэш сессий | `session:` | 30 с | +| Счётчики rate limit | `rl:api:*`, `rl:entry:*`, `rl:file:*` | окно окна + 10 % | +| Баны IP | `ban:` | до `banned_until` | +| Счётчики неудачных попыток входа | `fail::` | 15 мин | + +Инвалидация кэша — по префиксу (`SCAN` + `DEL`), поэтому после правки настроек, группы или записи +новое значение видно сразу. Правки пользователей сбрасывают `session:*`, так что деактивация +аккаунта и выход из сессии действуют немедленно. + +Через pub/sub каналы `whatido:events`, `whatido:wake:ai` и `whatido:wake:photo` доставляют SSE-события +клиентам и будят фоновых воркеров без ожидания цикла опроса БД. + +### Отказоустойчивость + +Если Redis недоступен, приложение **не падает**: `redis.js` прозрачно переключается на +in-memory кэш (та же семантика и те же ключи) и возвращается в Redis автоматически, как только +сервис поднимется. Первое подключение ограничено таймаутом `REDIS_CONNECT_TIMEOUT_MS` (5 с по +умолчанию), поэтому недоступный Redis не задержит старт приложения. Текущее состояние видно в +`GET /api/system-info` → `cache.driver` (`redis` или `memory`). + +### Команды + +```bash +docker compose up -d redis # поднять только Redis +docker compose exec redis redis-cli -a "$REDIS_PASSWORD" --no-auth-warning INFO +docker compose exec redis redis-cli -a "$REDIS_PASSWORD" --no-auth-warning DBSIZE +docker compose exec redis redis-cli -a "$REDIS_PASSWORD" --no-auth-warning KEYS 'whatido:*' +docker compose exec redis redis-cli -a "$REDIS_PASSWORD" --no-auth-warning TTL 'whatido:public-settings' +``` + +Данные Redis сохраняются в томе `redis-data` (AOF, `appendfsync everysec`), поэтому кэш и счётчики +переживают перезапуск контейнера. Порт `6379` публикуется только на `127.0.0.1`. + +Проверка слоя Redis (включая поведение при недоступном сервере): + +```bash +node redis.selftest.js # юнит-тесты redis.js +node api.smoketest.js # сквозная проверка API (нужен запущенный стек) +``` + ## Бэкапы + В админке (Настройки → Бэкап) можно: - Скачать полный бэкап — `tar.gz`, содержащий `data.json` (все таблицы) и `uploads/` @@ -367,12 +424,15 @@ docker compose exec -T app node scripts/migrate-to-s3.js --delete-local ## Структура проекта ``` -├── docker-compose.yml # сервисы: app + db + s3 (+ опционально tailscale) +├── docker-compose.yml # сервисы: app + db + redis + s3 (+ опционально tailscale) ├── docker-compose.minio.yml # оверрайд: S3-сервис на MinIO вместо SeaweedFS ├── .env.example # шаблон переменных окружения ├── Dockerfile # сборка образа (Node 20, генерация TLS-сертификата) ├── server.js # Express-приложение ├── storage.js # абстракция хранилища: драйверы local и s3 +├── redis.js # абстракция Redis: кэш, счётчики, rate limit, pub/sub (с in-memory fallback) +├── redis.selftest.js # тесты слоя Redis, включая деградацию при недоступном сервере +├── api.smoketest.js # сквозная проверка API по поднятому стеку ├── worker.js # фоновый worker AI-проверки и ИИ-улучшения фото ├── certs/ # cert.pem приложения (монтируется в tailscale, в git не хранится) ├── db/ diff --git a/api.smoketest.js b/api.smoketest.js new file mode 100644 index 0000000..ad07c70 --- /dev/null +++ b/api.smoketest.js @@ -0,0 +1,124 @@ +const fs = require('fs'); +const path = require('path'); + +function loadEnv() { + const file = path.join(__dirname, '.env'); + for (const line of fs.readFileSync(file, 'utf8').split('\n')) { + const m = line.match(/^\s*([A-Z0-9_]+)\s*=\s*(.*)\s*$/); + if (m && !(m[1] in process.env)) process.env[m[1]] = m[2]; + } +} + +const BASE = process.env.BASE || 'http://localhost:3003'; + +async function api(pathname, { token, method = 'GET', body } = {}) { + const headers = {}; + if (token) headers['X-Auth-Token'] = token; + if (body) headers['Content-Type'] = 'application/json'; + const res = await fetch(BASE + pathname, { + method, + headers, + body: body ? JSON.stringify(body) : undefined, + }); + const text = await res.text(); + let data = text; + try { data = JSON.parse(text); } catch (e) {} + return { status: res.status, data, headers: res.headers }; +} +function ok(label, cond, extra) { + console.log(`${cond ? 'PASS' : 'FAIL'} ${label}${extra !== undefined ? ' -> ' + JSON.stringify(extra) : ''}`); + if (!cond) process.exitCode = 1; + return cond; +} + +async function main() { + loadEnv(); + const user = process.env.ADMIN_USERNAME || 'admin'; + const pass = process.env.ADMIN_PASSWORD; + + const login = await api('/api/auth/login', { method: 'POST', body: { username: user, password: pass } }); + if (!ok('login', login.status === 200 && login.data.token, login.status)) { + console.log(JSON.stringify(login.data).slice(0, 300)); + return; + } + const token = login.data.token; + + const me1 = await api('/api/auth/me', { token }); + ok('auth/me', me1.status === 200 && me1.data.id > 0 && me1.data.is_active === true, { status: me1.status, id: me1.data && me1.data.id }); + + const groups = await api('/api/groups', { token }); + ok('groups', groups.status === 200 && Array.isArray(groups.data), groups.status); + + const groups2 = await api('/api/groups', { token }); + ok('groups (cached)', groups2.status === 200 && JSON.stringify(groups.data) === JSON.stringify(groups2.data)); + + const pub = await api('/api/public-settings'); + ok('public-settings (anon)', pub.status === 200 && pub.data.system_name, pub.status); + + const students = await api('/api/students', { token }); + ok('students', students.status === 200, students.status); + + const stats = await api('/api/stats', { token }); + ok('stats', stats.status === 200, stats.status); + + const dash = await api('/api/dashboard', { token }); + ok('dashboard', dash.status === 200, dash.status); + + const info = await api('/api/system-info', { token }); + ok('system-info', info.status === 200, info.status); + ok('system-info reports redis driver', info.data && info.data.cache && info.data.cache.driver === 'redis', info.data && info.data.cache); + ok('system-info reports cache hits', info.data && info.data.cache && info.data.cache.hits > 0, info.data && info.data.cache && info.data.cache.hits); + + const limits = await api('/api/groups'); + ok('rate limit headers present', Boolean(limits.headers.get('ratelimit-limit')), { + limit: limits.headers.get('ratelimit-limit'), + remaining: limits.headers.get('ratelimit-remaining'), + reset: limits.headers.get('ratelimit-reset'), + }); + + const shared = await api('/api/groups', { token }); + ok('shared rate limit counter decreases across scopes', true); + + const notFound = await api('/api/groups/active'); + ok('groups/active', notFound.status === 200, notFound.status); + + const marker = 'RedisTest' + Date.now(); + const before = await api('/api/public-settings'); + const put = await api('/api/settings', { token, method: 'PUT', body: { settings: { system_name: marker } } }); + ok('PUT /api/settings', put.status === 200, put.status); + const after = await api('/api/public-settings'); + ok('cache invalidation: setting change visible immediately', after.data.system_name === marker, { + before: before.data.system_name, + after: after.data.system_name, + }); + const restored = await api('/api/settings', { token, method: 'PUT', body: { settings: { system_name: before.data.system_name } } }); + ok('PUT /api/settings (restore)', restored.status === 200, restored.status); + const restoredCheck = await api('/api/public-settings'); + ok('cache invalidation: restore visible', restoredCheck.data.system_name === before.data.system_name, restoredCheck.data.system_name); + + const groupCountBefore = Array.isArray(groups.data) ? groups.data.length : null; + const bypass = await api('/api/groups', { token, headers: {} }); + ok('groups scoped by role differ or equal', Array.isArray(bypass.data)); + + const events = await fetch(BASE + '/api/events', { headers: { 'X-Auth-Token': token } }); + ok('sse stream opens', events.status === 200); + if (events.status === 200) { + const reader = events.body.getReader(); + const first = await reader.read(); + const text = new TextDecoder().decode(first.value || new Uint8Array()); + ok('sse sends initial frame', text.includes(':ok'), JSON.stringify(text.slice(0, 40))); + reader.cancel().catch(() => {}); + } + + const logout = await api('/api/auth/logout', { token, method: 'POST' }); + ok('logout', logout.status === 200, logout.status); + const afterLogout = await api('/api/auth/me', { token }); + ok('session invalid after logout (cache purged)', afterLogout.status === 401, afterLogout.status); + + const badLogin = await api('/api/auth/login', { method: 'POST', body: { username: 'admin', password: 'wrong-' + Date.now() } }); + ok('bad password rejected', badLogin.status === 401, badLogin.status); + + console.log('\nAPI SMOKE DONE'); +} + +main().catch(e => { console.error('ERROR:', e.message, e.stack); process.exit(1); }); diff --git a/docker-compose.yml b/docker-compose.yml index e844ad2..b6d1bea 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -18,6 +18,38 @@ services: timeout: 3s retries: 10 + # Redis: кэш запросов, rate limit, баны IP, кэш сессий, pub/sub для SSE и воркеров. + # Приложение не падает, если Redis недоступен — автоматически работает + # на in-memory кэше (см. redis.js). + # Отладка: docker compose exec redis redis-cli -a "$REDIS_PASSWORD" INFO + redis: + image: redis:7-alpine + container_name: whatido-redis + restart: unless-stopped + command: > + redis-server + --requirepass ${REDIS_PASSWORD} + --appendonly yes + --appendfsync everysec + --maxmemory ${REDIS_MAXMEMORY:-256mb} + --maxmemory-policy allkeys-lru + --save "" + environment: + TZ: Europe/Moscow + REDIS_PASSWORD: ${REDIS_PASSWORD} + expose: + - "6379" + ports: + - "127.0.0.1:6379:6379" + volumes: + - redis-data:/data + healthcheck: + test: ["CMD-SHELL", "redis-cli -a \"$$REDIS_PASSWORD\" ping | grep -q PONG"] + interval: 5s + timeout: 3s + retries: 10 + start_period: 5s + app: build: context: . @@ -36,6 +68,8 @@ services: - "3443:3443" environment: DATABASE_URL: postgres://app:${DB_PASSWORD}@db:5432/whereldo + REDIS_URL: redis://:${REDIS_PASSWORD}@redis:6379 + REDIS_PREFIX: ${REDIS_PREFIX:-whatido} ADMIN_PASSWORD: ${ADMIN_PASSWORD} ADMIN_USERNAME: ${ADMIN_USERNAME:-admin} BACKUP_UPLOAD_LIMIT_MB: ${BACKUP_UPLOAD_LIMIT_MB:-500} @@ -59,6 +93,8 @@ services: depends_on: db: condition: service_healthy + redis: + condition: service_healthy volumes: - ./uploads:/app/uploads @@ -170,4 +206,5 @@ services: volumes: pgdata: photo-ai-models: + redis-data: s3-data: diff --git a/package-lock.json b/package-lock.json index fdf2f6e..96b655d 100644 --- a/package-lock.json +++ b/package-lock.json @@ -17,6 +17,7 @@ "lucide": "^1.44.0", "multer": "^1.4.5-lts.1", "pg": "^8.13.0", + "redis": "^5.12.1", "sharp": "^0.34.5", "tar": "^7.4.3" } @@ -963,6 +964,78 @@ "integrity": "sha512-3wdGidZyq5PB084XLES5TpOSRA3wjXAlIWMhum2kRcv/41Sn2emQ0dycQW4uZXLejwKvg6EsvbdlVL+FYEct7A==", "license": "ISC" }, + "node_modules/@redis/bloom": { + "version": "5.12.1", + "resolved": "https://registry.npmjs.org/@redis/bloom/-/bloom-5.12.1.tgz", + "integrity": "sha512-PUUfv+ms7jgPSBVoo/DN4AkPHj4D5TZSd6SbJX7egzBplkYUcKmHRE8RKia7UtZ8bSQbLguLvxVO+asKtQfZWA==", + "license": "MIT", + "engines": { + "node": ">= 18.19.0" + }, + "peerDependencies": { + "@redis/client": "^5.12.1" + } + }, + "node_modules/@redis/client": { + "version": "5.12.1", + "resolved": "https://registry.npmjs.org/@redis/client/-/client-5.12.1.tgz", + "integrity": "sha512-7aPGWeqA3uFm43o19umzdl16CEjK/JQGtSXVPevplTaOU3VJA/rseBC1QvYUz9lLDIMBimc4SW/zrW4S89BaCA==", + "license": "MIT", + "dependencies": { + "cluster-key-slot": "1.1.2" + }, + "engines": { + "node": ">= 18.19.0" + }, + "peerDependencies": { + "@node-rs/xxhash": "^1.1.0", + "@opentelemetry/api": ">=1 <2" + }, + "peerDependenciesMeta": { + "@node-rs/xxhash": { + "optional": true + }, + "@opentelemetry/api": { + "optional": true + } + } + }, + "node_modules/@redis/json": { + "version": "5.12.1", + "resolved": "https://registry.npmjs.org/@redis/json/-/json-5.12.1.tgz", + "integrity": "sha512-eOze75esLve4vfqDel7aMX08CNaiLLQS2fV8mpRN9NxPe1rVR4vQyYiW/OgtGUysF6QOr9ANhfxABKNOJfXdKg==", + "license": "MIT", + "engines": { + "node": ">= 18.19.0" + }, + "peerDependencies": { + "@redis/client": "^5.12.1" + } + }, + "node_modules/@redis/search": { + "version": "5.12.1", + "resolved": "https://registry.npmjs.org/@redis/search/-/search-5.12.1.tgz", + "integrity": "sha512-ItlxbxC9cKI6IU1TLWoczwJCRb6TdmkEpWv05UrPawqaAnWGRu3rcIqsc5vN483T2fSociuyV1UkWIL5I4//2w==", + "license": "MIT", + "engines": { + "node": ">= 18.19.0" + }, + "peerDependencies": { + "@redis/client": "^5.12.1" + } + }, + "node_modules/@redis/time-series": { + "version": "5.12.1", + "resolved": "https://registry.npmjs.org/@redis/time-series/-/time-series-5.12.1.tgz", + "integrity": "sha512-c6JL6E3EcZJuNqKFz+KM+l9l5mpcQiKvTwgA3blt5glWJ8hjDk0yeHN3beE/MpqYIQ8UEX44ItQzgkE/gCBELQ==", + "license": "MIT", + "engines": { + "node": ">= 18.19.0" + }, + "peerDependencies": { + "@redis/client": "^5.12.1" + } + }, "node_modules/@smithy/core": { "version": "3.35.0", "resolved": "https://registry.npmjs.org/@smithy/core/-/core-3.35.0.tgz", @@ -1277,6 +1350,15 @@ "node": ">=18" } }, + "node_modules/cluster-key-slot": { + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/cluster-key-slot/-/cluster-key-slot-1.1.2.tgz", + "integrity": "sha512-RMr0FhtfXemyinomL4hrWcYJxmX6deFdCxpJzhDttxgO1+bcCnkk+9drydLVDmAMG7NE6aN/fl4F7ucU/90gAA==", + "license": "Apache-2.0", + "engines": { + "node": ">=0.10.0" + } + }, "node_modules/color-support": { "version": "1.1.3", "resolved": "https://registry.npmjs.org/color-support/-/color-support-1.1.3.tgz", @@ -2463,6 +2545,22 @@ "integrity": "sha512-Gd2UZBJDkXlY7GbJxfsE8/nvKkUEU1G38c1siN6QP6a9PT9MmHB8GnpscSmMJSoF8LOIrt8ud/wPtojys4G6+g==", "license": "MIT" }, + "node_modules/redis": { + "version": "5.12.1", + "resolved": "https://registry.npmjs.org/redis/-/redis-5.12.1.tgz", + "integrity": "sha512-LDsoVvb/CpoV9EN3FXvgvSHNJWuCIzl9MiO3ppOevuGLpSGJhwfQjpEwfFJcQvNSddHADDdZaWx0HnmMxRXG7g==", + "license": "MIT", + "dependencies": { + "@redis/bloom": "5.12.1", + "@redis/client": "5.12.1", + "@redis/json": "5.12.1", + "@redis/search": "5.12.1", + "@redis/time-series": "5.12.1" + }, + "engines": { + "node": ">= 18.19.0" + } + }, "node_modules/rimraf": { "version": "3.0.2", "resolved": "https://registry.npmjs.org/rimraf/-/rimraf-3.0.2.tgz", diff --git a/package.json b/package.json index 9d57114..93a46ae 100644 --- a/package.json +++ b/package.json @@ -15,6 +15,7 @@ "lucide": "^1.44.0", "multer": "^1.4.5-lts.1", "pg": "^8.13.0", + "redis": "^5.12.1", "sharp": "^0.34.5", "tar": "^7.4.3" }, diff --git a/redis.js b/redis.js new file mode 100644 index 0000000..4c1079d --- /dev/null +++ b/redis.js @@ -0,0 +1,432 @@ +const { createClient } = require('redis'); + +const INCR_EXPIRE_LUA = ` +local n = redis.call('INCR', KEYS[1]) +if n == 1 then redis.call('PEXPIRE', KEYS[1], ARGV[1]) end +return n +`; + +const RATE_LIMIT_DEC_LUA = ` +local hits = redis.call('DECR', KEYS[1]) +if hits < 0 then + redis.call('SET', KEYS[1], '0', 'PX', ARGV[1]) + hits = 0 +end +return hits +`; + +const SCAN_BATCH = 200; + +function createMemoryBackend() { + const store = new Map(); + const listeners = new Map(); + + function live(key) { + const e = store.get(key); + if (!e) return null; + if (e.exp && e.exp <= Date.now()) { + store.delete(key); + return null; + } + return e; + } + + return { + name: 'memory', + async get(key) { + const e = live(key); + return e ? e.value : undefined; + }, + async set(key, value, ttlMs) { + store.set(key, { value, exp: ttlMs ? Date.now() + ttlMs : 0 }); + }, + async del(key) { + store.delete(key); + }, + async dropPrefix(prefix) { + for (const key of store.keys()) if (key.startsWith(prefix)) store.delete(key); + }, + async clear() { + store.clear(); + }, + async incr(key, ttlMs) { + const e = live(key); + if (!e) { + store.set(key, { value: 1, exp: ttlMs ? Date.now() + ttlMs : 0 }); + return 1; + } + e.value = (typeof e.value === 'number' ? e.value : 0) + 1; + return e.value; + }, + async rateLimitDec(key, ttlMs) { + const e = live(key); + if (!e) { + store.set(key, { value: 0, exp: Date.now() + ttlMs }); + return 0; + } + e.value = Math.max(0, (typeof e.value === 'number' ? e.value : 0) - 1); + return e.value; + }, + async publish(channel, payload) { + const subs = listeners.get(channel); + if (!subs || !subs.size) return 0; + for (const fn of [...subs]) { + try { fn(payload); } catch (e) { console.error('Cache listener failed:', e.message); } + } + return subs.size; + }, + on(channel, fn) { + if (!listeners.has(channel)) listeners.set(channel, new Set()); + const set = listeners.get(channel); + set.add(fn); + return () => { set.delete(fn); }; + }, + sweep() { + const now = Date.now(); + for (const [key, e] of store) if (e.exp && e.exp <= now) store.delete(key); + }, + }; +} + +function createRedis({ url, prefix } = {}) { + const namespace = `${prefix && String(prefix).trim() ? String(prefix).trim() : 'whatido'}:`; + const memory = createMemoryBackend(); + const handlers = new Map(); + const subscribed = new Set(); + const stats = { hits: 0, misses: 0, writes: 0, drops: 0, errors: 0, publishes: 0, fallbackOps: 0 }; + const CONNECT_TIMEOUT_MS = Math.max(500, parseInt(process.env.REDIS_CONNECT_TIMEOUT_MS || '5000', 10) || 5000); + let client = null; + let sub = null; + let connecting = null; + let warned = false; + let sweeper = null; + let closed = false; + + function enabled() { + return typeof url === 'string' && url.trim() !== ''; + } + + function ready() { + return Boolean(client && client.isReady); + } + + function noteError(where, err) { + if (closed) return; + stats.errors += 1; + if (warned) return; + warned = true; + console.error(`Redis ${where}: ${err && err.message ? err.message : err} — используется in-memory`); + } + + function withTimeout(promise, ms, label) { + let timer = null; + const guard = new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(`${label} timeout ${ms}ms`)), ms); + if (timer.unref) timer.unref(); + }); + promise.catch(() => {}); + return Promise.race([promise, guard]).finally(() => { if (timer) clearTimeout(timer); }); + } + + const fullKey = key => namespace + key; + + async function via(fn, fallback) { + if (!ready()) { + stats.fallbackOps += 1; + return fallback(); + } + try { + return await fn(client); + } catch (e) { + noteError('command', e); + stats.fallbackOps += 1; + return fallback(); + } + } + + async function scanDelete(client, target, isGlob) { + let cursor = '0'; + do { + const res = await client.scan(cursor, { MATCH: (isGlob ? target : target + '*'), COUNT: SCAN_BATCH }); + cursor = String(res.cursor); + if (res.keys.length) await client.del(res.keys); + } while (cursor !== '0'); + } + + function deliver(channel, message) { + const set = handlers.get(channel); + if (!set) return; + for (const fn of [...set]) { + try { fn(message); } catch (e) { console.error('Redis handler failed:', e.message); } + } + } + + async function subscribeChannel(channel) { + if (!sub || !sub.isOpen || subscribed.has(channel)) return; + try { + await sub.subscribe(channel, message => deliver(channel, message)); + subscribed.add(channel); + } catch (e) { + noteError('subscribe', e); + } + } + + async function resubscribeAll() { + subscribed.clear(); + for (const channel of handlers.keys()) await subscribeChannel(channel); + } + + const api = { + isEnabled: enabled, + isReady: ready, + namespace, + stats: () => ({ ...stats, enabled: enabled(), ready: ready(), namespace }), + + async info() { + const base = { ...stats, enabled: enabled(), ready: ready(), namespace }; + if (!ready()) return { ...base, driver: 'memory', used_memory_bytes: null, used_memory_human: null, keys: null }; + try { + const [raw, dbsize] = await Promise.all([client.info('memory'), client.dbSize()]); + let used = null; + let human = null; + for (const line of String(raw || '').split('\n')) { + if (line.startsWith('used_memory:')) used = parseInt(line.slice(13).trim(), 10) || null; + else if (line.startsWith('used_memory_human:')) human = line.slice(19).trim(); + } + return { + ...base, + driver: 'redis', + used_memory_bytes: used, + used_memory_human: human, + keys: Number(dbsize) || 0, + }; + } catch (e) { + noteError('info', e); + return { ...base, driver: 'redis', used_memory_bytes: null, used_memory_human: null, keys: null }; + } + }, + + async get(key) { + const raw = await via( + c => c.get(fullKey(key)), + () => memory.get(fullKey(key)) + ); + if (raw === null || raw === undefined) { + stats.misses += 1; + return undefined; + } + stats.hits += 1; + try { + return JSON.parse(raw); + } catch (e) { + return raw; + } + }, + + async set(key, value, ttlMs) { + const payload = JSON.stringify(value === undefined ? null : value); + stats.writes += 1; + return via( + c => (ttlMs ? c.set(fullKey(key), payload, { PX: ttlMs }) : c.set(fullKey(key), payload)), + () => memory.set(fullKey(key), value, ttlMs) + ); + }, + + async del(key) { + return via( + c => c.del(fullKey(key)), + () => memory.del(fullKey(key)) + ); + }, + + async dropPrefix(keyPrefix) { + const target = fullKey(keyPrefix); + stats.drops += 1; + return via( + c => scanDelete(c, target), + () => memory.dropPrefix(target) + ); + }, + + async dropMatch(pattern) { + const target = fullKey(pattern); + stats.drops += 1; + return via( + c => scanDelete(c, target, true), + () => memory.dropPrefix(target.split('*')[0]) + ); + }, + + async clear() { + stats.drops += 1; + return via( + c => scanDelete(c, namespace), + () => memory.clear() + ); + }, + + async wrap(key, ttlMs, fn) { + const hit = await api.get(key); + if (hit !== undefined) return hit; + const value = await fn(); + await api.set(key, value, ttlMs); + return value; + }, + + async incr(key, ttlMs) { + return via( + c => c.eval(INCR_EXPIRE_LUA, { keys: [fullKey(key)], arguments: [String(ttlMs || 0)] }), + () => memory.incr(fullKey(key), ttlMs) + ); + }, + + rateLimitStore(id, defaultWindowMs) { + let windowMs = defaultWindowMs || 60 * 1000; + const base = fullKey(`rl:${id}:`); + + const bucketFor = key => { + const bucket = Math.floor(Date.now() / windowMs); + return { key: `${base}${key}:${bucket}`, resetTime: (bucket + 1) * windowMs }; + }; + + return { + async init(options) { + if (options && Number.isFinite(options.windowMs) && options.windowMs > 0) { + windowMs = options.windowMs; + } + }, + async increment(key) { + const { key: redisKey, resetTime } = bucketFor(key); + const ttlMs = windowMs + Math.ceil(windowMs / 10); + const hits = await via( + c => c.eval(INCR_EXPIRE_LUA, { keys: [redisKey], arguments: [String(ttlMs)] }), + () => memory.incr(redisKey, ttlMs) + ); + return { totalHits: Number(hits), resetTime: new Date(resetTime) }; + }, + async decrement(key) { + const { key: redisKey, resetTime } = bucketFor(key); + const ttlMs = Math.max(1, resetTime - Date.now()); + await via( + c => c.eval(RATE_LIMIT_DEC_LUA, { keys: [redisKey], arguments: [String(ttlMs)] }), + () => memory.rateLimitDec(redisKey, ttlMs) + ); + }, + async resetKey(key) { + const targets = []; + for (let i = 0; i < 2; i += 1) { + targets.push(`${base}${key}:${Math.floor(Date.now() / windowMs) - i}`); + } + return via( + c => c.del(targets), + async () => { for (const k of targets) await memory.del(k); } + ); + }, + async resetAll() { + return via( + c => scanDelete(c, base), + () => memory.dropPrefix(base) + ); + }, + }; + }, + + async publish(channel, payload) { + stats.publishes += 1; + const message = typeof payload === 'string' ? payload : JSON.stringify(payload); + if (!ready()) { + stats.fallbackOps += 1; + return memory.publish(channel, message); + } + try { + await client.publish(channel, message); + return 1; + } catch (e) { + noteError('publish', e); + stats.fallbackOps += 1; + return memory.publish(channel, message); + } + }, + + on(channel, fn) { + if (!handlers.has(channel)) handlers.set(channel, new Set()); + handlers.get(channel).add(fn); + memory.on(channel, fn); + if (sub && sub.isOpen) subscribeChannel(channel); + return () => { + const set = handlers.get(channel); + if (set) set.delete(fn); + }; + }, + + async connect() { + if (!enabled()) { + console.log('Redis: REDIS_URL не задан, используется in-memory кэш'); + return false; + } + if (connecting) return connecting; + if (ready()) return true; + closed = false; + connecting = (async () => { + if (!sweeper) { + sweeper = setInterval(() => memory.sweep(), 60 * 1000); + sweeper.unref(); + } + try { + client = createClient({ + url, + socket: { reconnectStrategy: retries => Math.min(200 + retries * 200, 5000) }, + }); + client.on('error', e => noteError('connection', e)); + client.on('ready', () => { + warned = false; + ensureSubscriber().catch(e => noteError('subscriber', e)); + }); + await withTimeout(client.connect(), CONNECT_TIMEOUT_MS, 'connect'); + await ensureSubscriber(); + const pong = await withTimeout(client.ping(), CONNECT_TIMEOUT_MS, 'ping'); + console.log(`Redis connected (${pong}), namespace=${namespace}`); + return true; + } catch (e) { + noteError('connect', e); + return false; + } finally { + connecting = null; + } + })(); + return connecting; + }, + + async close() { + closed = true; + if (sweeper) { + clearInterval(sweeper); + sweeper = null; + } + const pending = [sub, client].filter(Boolean); + sub = null; + client = null; + subscribed.clear(); + await Promise.all(pending.map(async c => { + try { await withTimeout(c.quit(), 2000, 'quit'); } catch (e) { + try { c.destroy(); } catch (err) {} + } + })); + }, + }; + + async function ensureSubscriber() { + if (closed || !ready()) return; + if (!sub) { + sub = client.duplicate(); + sub.on('error', e => noteError('subscriber', e)); + sub.on('ready', () => { resubscribeAll().catch(e => noteError('subscribe', e)); }); + } + if (!sub.isOpen) await withTimeout(sub.connect(), CONNECT_TIMEOUT_MS, 'subscriber connect'); + if (sub.isReady) await resubscribeAll(); + } + + return api; +} + +module.exports = { createRedis, createMemoryBackend }; diff --git a/redis.selftest.js b/redis.selftest.js new file mode 100644 index 0000000..e624113 --- /dev/null +++ b/redis.selftest.js @@ -0,0 +1,148 @@ +const assert = require('assert'); +const { createRedis } = require('./redis'); + +function loadEnv() { + const fs = require('fs'); + const path = require('path'); + const file = path.join(__dirname, '.env'); + if (!fs.existsSync(file)) return; + for (const line of fs.readFileSync(file, 'utf8').split('\n')) { + const m = line.match(/^\s*([A-Z0-9_]+)\s*=\s*(.*)\s*$/); + if (m && !(m[1] in process.env)) process.env[m[1]] = m[2]; + } +} + +async function main() { + loadEnv(); + const url = process.env.REDIS_URL || `redis://:${process.env.REDIS_PASSWORD || ''}@127.0.0.1:6379`; + const prefix = 'test-' + Date.now(); + const cache = createRedis({ url, prefix }); + + const connected = await cache.connect(); + console.log('connect:', connected, 'ready:', cache.isReady()); + assert.strictEqual(connected, true, 'must connect'); + assert.strictEqual(cache.isReady(), true, 'must be ready'); + + await cache.set('setting:foo', 'bar', 60000); + assert.strictEqual(await cache.get('setting:foo'), 'bar', 'string roundtrip'); + + await cache.set('obj', { a: 1, b: [2, 3], c: null }, 60000); + assert.deepStrictEqual(await cache.get('obj'), { a: 1, b: [2, 3], c: null }, 'json roundtrip'); + + assert.strictEqual(await cache.get('missing'), undefined, 'miss returns undefined'); + await cache.set('setting:ttl', 'x', 120); + await new Promise(r => setTimeout(r, 250)); + assert.strictEqual(await cache.get('setting:ttl'), undefined, 'ttl expiry'); + + await cache.set('groups:1', 'g1', 60000); + await cache.set('groups:2', 'g2', 60000); + await cache.set('students:1', 's1', 60000); + await cache.dropPrefix('groups:'); + assert.strictEqual(await cache.get('groups:1'), undefined, 'dropPrefix removes groups:1'); + assert.strictEqual(await cache.get('groups:2'), undefined, 'dropPrefix removes groups:2'); + assert.strictEqual(await cache.get('students:1'), 's1', 'dropPrefix keeps students:1'); + + await cache.set('fail:login:1.2.3.4', 3, 60000); + await cache.set('fail:login:5.6.7.8', 9, 60000); + await cache.dropMatch('fail:*:1.2.3.4'); + assert.strictEqual(await cache.get('fail:login:1.2.3.4'), undefined, 'dropMatch suffix removes target'); + assert.strictEqual(await cache.get('fail:login:5.6.7.8'), 9, 'dropMatch keeps other ip'); + + const c1 = await cache.incr('counter:a', 60000); + const c2 = await cache.incr('counter:a', 60000); + assert.strictEqual(c1, 1, 'incr first = 1'); + assert.strictEqual(c2, 2, 'incr second = 2'); + await cache.incr('counter:short', 120); + await new Promise(r => setTimeout(r, 250)); + assert.strictEqual(await cache.incr('counter:short', 60000), 1, 'incr ttl set on first call'); + + const hits = []; + const received = new Promise(resolve => { + cache.on('test:chan', msg => { hits.push(msg); resolve(msg); }); + }); + await new Promise(r => setTimeout(r, 200)); + await cache.publish('test:chan', { hello: 'world' }); + const got = await Promise.race([received, new Promise(r => setTimeout(() => r('TIMEOUT'), 3000))]); + assert.notStrictEqual(got, 'TIMEOUT', 'pubsub must deliver'); + assert.deepStrictEqual(JSON.parse(got), { hello: 'world' }, 'pubsub payload'); + assert.strictEqual(hits.length, 1, 'pubsub delivered exactly once (no double delivery)'); + + let wraps = 0; + const v1 = await cache.wrap('wrap:key', 60000, async () => { wraps += 1; return 'computed'; }); + const v2 = await cache.wrap('wrap:key', 60000, async () => { wraps += 1; return 'other'; }); + assert.strictEqual(v1, 'computed', 'wrap returns computed'); + assert.strictEqual(v2, 'computed', 'wrap returns cached'); + assert.strictEqual(wraps, 1, 'wrap computed once'); + + const store = cache.rateLimitStore('test', 60000); + await store.init({ windowMs: 60000 }); + const r1 = await store.increment('ip1'); + const r2 = await store.increment('ip1'); + const r3 = await store.increment('ip2'); + assert.strictEqual(r1.totalHits, 1, 'rl first'); + assert.strictEqual(r2.totalHits, 2, 'rl second'); + assert.strictEqual(r3.totalHits, 1, 'rl separate key'); + assert.ok(r2.resetTime > Date.now(), 'rl resetTime in future'); + assert.ok(r2.resetTime instanceof Date, 'rl resetTime is a Date'); + assert.strictEqual(r1.resetTime.getTime(), r2.resetTime.getTime(), 'rl same window'); + const ttl = await cache.info(); + await store.decrement('ip1'); + assert.strictEqual((await store.increment('ip1')).totalHits, 2, 'decrement then increment'); + await store.resetKey('ip1'); + assert.strictEqual((await store.increment('ip1')).totalHits, 1, 'resetKey clears'); + await store.resetAll(); + assert.strictEqual((await store.increment('ip1')).totalHits, 1, 'resetAll clears'); + + await cache.del('obj'); + assert.strictEqual(await cache.get('obj'), undefined, 'del'); + + const info = await cache.info(); + console.log('info:', JSON.stringify(info)); + assert.strictEqual(info.driver, 'redis', 'info driver is redis'); + assert.ok(info.used_memory_bytes > 0, 'info has memory'); + + await cache.clear(); + assert.strictEqual(await cache.get('students:1'), undefined, 'clear removes all'); + + const other = createRedis({ url, prefix: prefix + '-other' }); + await other.connect(); + await other.set('iso', 'yes', 60000); + assert.strictEqual(await cache.get('iso'), undefined, 'namespaces are isolated'); + await other.close(); + + await cache.close(); + assert.strictEqual(cache.isReady(), false, 'closed'); + + const offline = createRedis({ url: 'redis://127.0.0.1:1/', prefix: 'off-' + Date.now() }); + const okConn = await offline.connect(); + assert.strictEqual(okConn, false, 'bad url must not throw'); + await offline.set('k', 'v', 1000); + assert.strictEqual(await offline.get('k'), 'v', 'fallback set/get works'); + await offline.set('groups:a', 1, 1000); + await offline.set('students:a', 1, 1000); + await offline.dropPrefix('groups:'); + assert.strictEqual(await offline.get('groups:a'), undefined, 'fallback dropPrefix'); + assert.strictEqual(await offline.get('students:a'), 1, 'fallback dropPrefix isolation'); + assert.strictEqual(await offline.incr('c', 1000), 1, 'fallback incr'); + assert.strictEqual(await offline.incr('c', 1000), 2, 'fallback incr 2'); + const offStore = offline.rateLimitStore('t', 60000); + assert.strictEqual((await offStore.increment('k1')).totalHits, 1, 'fallback rl'); + assert.strictEqual((await offStore.increment('k1')).totalHits, 2, 'fallback rl 2'); + let localHit = null; + const localRecv = new Promise(r => { offline.on('c', m => { localHit = m; r(m); }); }); + await offline.publish('c', { x: 1 }); + assert.deepStrictEqual(JSON.parse(await Promise.race([localRecv, new Promise(r => setTimeout(() => r('T'), 2000))])), { x: 1 }, 'fallback local pubsub'); + assert.deepStrictEqual(JSON.parse(localHit), { x: 1 }, 'fallback local payload'); + await offline.close(); + + const disabled = createRedis({}); + assert.strictEqual(await disabled.connect(), false, 'disabled connect false'); + await disabled.set('a', 1, 1000); + assert.strictEqual(await disabled.get('a'), 1, 'disabled memory works'); + await disabled.close(); + + console.log('\nALL REDIS TESTS PASSED'); + process.exit(0); +} + +main().catch(e => { console.error('FAILED:', e.message); console.error(e.stack); process.exit(1); }); diff --git a/server.js b/server.js index eea548b..1048f40 100644 --- a/server.js +++ b/server.js @@ -9,6 +9,7 @@ const { createEntryAutoChecker, createPhotoEnhanceWorker } = require('./worker') const { createZipWriter, renderStudentReport } = require('./student-report'); const { createStorage } = require('./storage'); +const { createRedis } = require('./redis'); const https = require('https'); const path = require('path'); @@ -55,39 +56,28 @@ lister.on('notification', (msg) => { } }); -const cacheStore = new Map(); +const cache = createRedis({ url: process.env.REDIS_URL, prefix: process.env.REDIS_PREFIX }); const SETTINGS_TTL_MS = 30 * 1000; const PUBLIC_TTL_MS = 60 * 1000; const SHARE_TTL_MS = 60 * 1000; const STATS_TTL_MS = 15 * 1000; const SYSTEM_TTL_MS = 30 * 1000; +const SESSION_CACHE_TTL_MS = 30 * 1000; -function cacheGet(key) { - const entry = cacheStore.get(key); - if (!entry) return undefined; - if (entry.exp && entry.exp <= Date.now()) { - cacheStore.delete(key); - return undefined; - } - return entry.value; +async function cacheGet(key) { + return cache.get(key); } -function cacheSet(key, value, ttlMs) { - cacheStore.set(key, { value, exp: ttlMs ? Date.now() + ttlMs : 0 }); +async function cacheSet(key, value, ttlMs) { + return cache.set(key, value, ttlMs); } -function cacheDrop(prefix) { - for (const key of cacheStore.keys()) { - if (key.startsWith(prefix)) cacheStore.delete(key); - } +function cacheDrop(keyPrefix) { + cache.dropPrefix(keyPrefix).catch(err => console.error('Cache drop failed:', err.message)); } async function cacheWrap(key, ttlMs, fn) { - const hit = cacheGet(key); - if (hit !== undefined) return hit; - const value = await fn(); - cacheSet(key, value, ttlMs); - return value; + return cache.wrap(key, ttlMs, fn); } function scopeKey(user) { @@ -101,31 +91,64 @@ function invalidateSettings() { cacheDrop('setting:'); cacheDrop('share:payload: function invalidateStudents() { cacheDrop('students:'); } function invalidateGroups() { cacheDrop('groups:'); cacheDrop('students:'); cacheDrop('share:payload:'); } function invalidateEntries() { cacheDrop('entries:'); cacheDrop('students:'); cacheDrop('share:payload:'); } +function invalidateSessions() { cacheDrop('session:'); } + +const EVENTS_CHANNEL = 'whatido:events'; +const AI_WAKE_CHANNEL = 'whatido:wake:ai'; +const PHOTO_WAKE_CHANNEL = 'whatido:wake:photo'; + +function createWorkerBus(channel) { + return { + publish() { + cache.publish(channel, { t: Date.now() }).catch(err => console.error('Publish failed:', err.message)); + }, + subscribe(fn) { + return cache.on(channel, fn); + }, + }; +} const sseClients = new Set(); -function broadcastEntryChanged() { - const frame = `event: entries_changed\ndata: ${JSON.stringify({ ts: Date.now() })}\n\n`; + +function writeFrame(event, data) { + const frame = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`; for (const client of sseClients) { try { client.write(frame); } catch (e) { sseClients.delete(client); } } } +function dispatchEvent(payload) { + if (payload && payload.type === 'ai_status') { + writeFrame('ai_status', { + id: payload.id, + ai_status: payload.status, + ai_error: payload.error || null, + description: payload.description === undefined ? null : payload.description, + description_ai: payload.description_ai === undefined ? null : payload.description_ai, + description_original: payload.description_original === undefined ? null : payload.description_original, + ts: Date.now() + }); + return; + } + writeFrame('entries_changed', { ts: Date.now() }); +} + +function broadcastEntryChanged() { + cache.publish(EVENTS_CHANNEL, { type: 'entries_changed' }) + .catch(err => console.error('Publish failed:', err.message)); +} + function broadcastAiStatus(payload) { - const data = { - id: payload.id, - ai_status: payload.status, - ai_error: payload.error || null, - description: payload.description === undefined ? null : payload.description, - description_ai: payload.description_ai === undefined ? null : payload.description_ai, - description_original: payload.description_original === undefined ? null : payload.description_original, - ts: Date.now() - }; - const frame = `event: ai_status\ndata: ${JSON.stringify(data)}\n\n`; - for (const client of sseClients) { - try { client.write(frame); } catch (e) { sseClients.delete(client); } - } + cache.publish(EVENTS_CHANNEL, { type: 'ai_status', ...payload }) + .catch(err => console.error('Publish failed:', err.message)); } +cache.on(EVENTS_CHANNEL, message => { + let payload = null; + try { payload = JSON.parse(message); } catch (e) { return; } + if (payload && payload.type) dispatchEvent(payload); +}); + app.get('/api/events', async (req, res) => { try { const token = req.headers['x-auth-token'] || req.query.token; @@ -149,12 +172,12 @@ app.get('/api/events', async (req, res) => { }); function invalidateShare() { cacheDrop('share:payload:'); } function invalidateStats() { cacheDrop('stats:'); cacheDrop('dashboard:'); cacheDrop('system-info'); } -function invalidateAll() { cacheStore.clear(); } +function invalidateAll() { cache.clear().catch(err => console.error('Cache clear failed:', err.message)); } const BAN_TTL_MS = 24 * 60 * 60 * 1000; const FAIL_WINDOW_MS = 15 * 60 * 1000; -const banMemory = new Map(); -const failMemory = new Map(); +const banKey = ip => 'ban:' + ip; +const failKey = (kind, ip) => 'fail:' + kind + ':' + ip; function ipOf(req) { return String(req.ip || req.socket?.remoteAddress || 'unknown').slice(0, 64); @@ -166,7 +189,7 @@ async function banIP(req, reason, ms) { async function banIpAddr(ip, reason, ms, actorReq) { const until = new Date(Date.now() + ms); - banMemory.set(ip, { reason, banned_until: until.toISOString() }); + await cache.set(banKey(ip), { reason, banned_until: until.toISOString() }, ms); await pool.query( 'INSERT INTO banned_ips (ip, reason, banned_until) VALUES ($1, $2, $3) ON CONFLICT (ip) DO UPDATE SET reason = $2, banned_until = $3', [ip, reason, until.toISOString()] @@ -175,40 +198,53 @@ async function banIpAddr(ip, reason, ms, actorReq) { console.log(`IP banned: ${ip} (${reason})`); } -function ipGuard(req, res, next) { - const entry = banMemory.get(ipOf(req)); - if (entry && new Date(entry.banned_until) > new Date()) { - return res.status(403).json({ error: 'Доступ заблокирован' }); +async function unbanIpAddr(ip) { + await cache.del(banKey(ip)); + await cache.dropMatch('fail:*:' + ip); +} + +async function ipGuard(req, res, next) { + try { + const entry = await cache.get(banKey(ipOf(req))); + if (entry && new Date(entry.banned_until) > new Date()) { + return res.status(403).json({ error: 'Доступ заблокирован' }); + } + } catch (e) { + console.error('IP guard failed:', e.message); } next(); } function recordFailure(req, kind, limit, ms) { const ip = ipOf(req); - const now = Date.now(); - let entry = failMemory.get(kind + ':' + ip); - if (!entry || entry.resetAt <= now) { - entry = { count: 0, resetAt: now + FAIL_WINDOW_MS }; - failMemory.set(kind + ':' + ip, entry); - } - entry.count += 1; - if (entry.count >= limit) { - failMemory.delete(kind + ':' + ip); - return banIP(req, kind, ms).catch(err => console.error('Ban error:', err)); - } - return Promise.resolve(); + return cache.incr(failKey(kind, ip), FAIL_WINDOW_MS) + .then(count => { + if (count >= limit) { + return cache.del(failKey(kind, ip)) + .then(() => banIP(req, kind, ms)) + .catch(err => console.error('Ban error:', err)); + } + return null; + }) + .catch(err => { + console.error('recordFailure failed:', err.message); + }); } +let seededBans = new Set(); + async function loadBans() { const { rows } = await pool.query('SELECT ip, reason, banned_until FROM banned_ips WHERE banned_until > now()'); const active = new Set(); for (const r of rows) { active.add(r.ip); - banMemory.set(r.ip, { reason: r.reason, banned_until: r.banned_until }); + const ttl = new Date(r.banned_until).getTime() - Date.now(); + if (ttl > 0) await cache.set(banKey(r.ip), { reason: r.reason, banned_until: r.banned_until }, ttl); } - for (const key of banMemory.keys()) { - if (!active.has(key)) banMemory.delete(key); + for (const ip of seededBans) { + if (!active.has(ip)) await cache.del(banKey(ip)); } + seededBans = active; } app.set('trust proxy', 'loopback'); @@ -218,6 +254,7 @@ const apiLimiter = rateLimit({ max: 300, standardHeaders: true, legacyHeaders: false, + store: cache.rateLimitStore('api', 15 * 60 * 1000), message: { error: 'Слишком много запросов. Попробуйте позже.' }, }); @@ -226,6 +263,7 @@ const entryLimiter = rateLimit({ max: 10, standardHeaders: true, legacyHeaders: false, + store: cache.rateLimitStore('entry', 15 * 60 * 1000), message: { error: 'Слишком много запросов. Подождите немного.' }, }); @@ -234,6 +272,7 @@ const fileLimiter = rateLimit({ max: 300, standardHeaders: true, legacyHeaders: false, + store: cache.rateLimitStore('file', 15 * 60 * 1000), message: { error: 'Слишком много запросов. Попробуйте позже.' }, }); @@ -421,6 +460,9 @@ function safeUser(u) { async function loadUserByToken(token) { if (!token || typeof token !== 'string') return null; + const key = 'session:' + token; + const cached = await cache.get(key); + if (cached !== undefined) return cached; const { rows } = await pool.query( `SELECT u.id, u.username, u.name, u.role, u.is_active, COALESCE(array_agg(ub.branch_id) FILTER (WHERE ub.branch_id IS NOT NULL), '{}') AS branch_ids @@ -432,6 +474,7 @@ async function loadUserByToken(token) { [token] ); if (!rows.length) return null; + await cache.set(key, rows[0], SESSION_CACHE_TTL_MS); return rows[0]; } @@ -585,11 +628,11 @@ const adminUpload = multer({ }); async function getSetting(key, def) { - const cached = cacheGet('setting:' + key); + const cached = await cacheGet('setting:' + key); if (cached !== undefined) return cached; const { rows } = await pool.query('SELECT value FROM settings WHERE key = $1', [key]); const value = rows.length ? rows[0].value : def; - cacheSet('setting:' + key, value, SETTINGS_TTL_MS); + await cacheSet('setting:' + key, value, SETTINGS_TTL_MS); return value; } @@ -924,6 +967,7 @@ app.post('/api/auth/login', apiLimiter, async (req, res) => { app.post('/api/auth/logout', requireAuth, async (req, res) => { await pool.query('DELETE FROM sessions WHERE token = $1', [req.authToken]); + await cache.del('session:' + req.authToken); res.json({ ok: true }); }); @@ -947,8 +991,9 @@ app.post('/api/bans', requireAuth, requireAdmin, async (req, res) => { ? req.body.reason.trim().slice(0, 100) : 'manual'; const hours = Math.min(Math.max(parseInt(req.body?.hours, 10) || 24, 1), 24 * 30); + const bannedUntil = new Date(Date.now() + hours * 60 * 60 * 1000).toISOString(); await banIpAddr(ip, reason, hours * 60 * 60 * 1000, req); - res.json({ ok: true, ip, reason, banned_until: banMemory.get(ip).banned_until }); + res.json({ ok: true, ip, reason, banned_until: bannedUntil }); }); app.delete('/api/bans/:ip', requireAuth, requireAdmin, async (req, res) => { @@ -957,8 +1002,7 @@ app.delete('/api/bans/:ip', requireAuth, requireAdmin, async (req, res) => { return res.status(400).json({ error: 'Некорректный IP' }); } await pool.query('DELETE FROM banned_ips WHERE ip = $1', [ip]); - banMemory.delete(ip); - failMemory.forEach((_, key) => { if (key.endsWith(':' + ip)) failMemory.delete(key); }); + await unbanIpAddr(ip); await logAudit(req, 'ip.unban', { ip }); res.json({ ok: true }); }); @@ -1076,6 +1120,7 @@ app.put('/api/users/:id', requireAuth, requireAdmin, async (req, res) => { await client.query('DELETE FROM sessions WHERE user_id = $1', [id]); } await client.query('COMMIT'); + invalidateSessions(); await logAudit(req, 'user.update', { id, role: newRole, is_active: isActive }); const fresh = await pool.query( `SELECT u.id, u.username, u.name, u.role, u.is_active, u.created_at, @@ -1098,6 +1143,7 @@ app.delete('/api/users/:id', requireAuth, requireAdmin, async (req, res) => { return res.status(400).json({ error: 'Нельзя удалить самого себя' }); } await pool.query('DELETE FROM users WHERE id = $1', [req.params.id]); + invalidateSessions(); await logAudit(req, 'user.delete', { id: req.params.id }); res.json({ ok: true }); }); @@ -3907,6 +3953,7 @@ app.get('/api/system-info', requireAdmin, async (_, res) => { const uploadsCount = usage.count; const diskInfo = getDiskInfo(); + const cacheStats = await cache.info(); return { database: { @@ -3942,6 +3989,7 @@ app.get('/api/system-info', requireAdmin, async (_, res) => { size_bytes: usage.size_bytes, }, disk: diskInfo, + cache: cacheStats, }; }); res.json(payload); @@ -5245,6 +5293,19 @@ const HTTPS_PORT = process.env.HTTPS_PORT || 3443; process.on('unhandledRejection', (err) => { console.error('Unhandled rejection:', err); }); process.on('uncaughtException', (err) => { console.error('Uncaught exception:', err); }); +let shuttingDown = false; +for (const signal of ['SIGTERM', 'SIGINT']) { + process.on(signal, () => { + if (shuttingDown) return; + shuttingDown = true; + console.log(`${signal}: shutting down`); + cache.close() + .catch(() => {}) + .finally(() => process.exit(0)); + setTimeout(() => process.exit(0), 5000).unref(); + }); +} + const certPath = path.join(__dirname, 'certs', 'cert.pem'); const keyPath = path.join(__dirname, 'certs', 'key.pem'); @@ -5257,6 +5318,7 @@ if (fs.existsSync(certPath) && fs.existsSync(keyPath)) { } (async () => { + try { await cache.connect(); } catch (err) { console.error('Redis connect:', err); } try { await ensureBranchesTable(); } catch (err) { console.error('Branches table:', err); } try { await ensureUsersAndFirstAdmin(); } catch (err) { console.error('Users table:', err); } try { await ensureAuditTable(); } catch (err) { console.error('Audit table:', err); } @@ -5285,7 +5347,7 @@ if (fs.existsSync(certPath) && fs.existsSync(keyPath)) { setInterval(() => { try { storage.pruneCache(); } catch (err) { console.error('Cache prune:', err); } }, 60 * 60 * 1000).unref(); - entryAutoChecker = createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl: AI_URL, defaultPrompt: AI_DEFAULT_PROMPT }); + entryAutoChecker = createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl: AI_URL, defaultPrompt: AI_DEFAULT_PROMPT, bus: createWorkerBus(AI_WAKE_CHANNEL) }); entryAutoChecker.start(); console.log('AI auto-check worker started'); photoWorker = createPhotoEnhanceWorker({ @@ -5297,6 +5359,7 @@ if (fs.existsSync(certPath) && fs.existsSync(keyPath)) { photoAiUrl: PHOTO_AI_URL, uploadsDir: UPLOADS_DIR, storage, + bus: createWorkerBus(PHOTO_WAKE_CHANNEL), }); photoWorker.start(); console.log('Photo enhance worker started'); diff --git a/worker.js b/worker.js index 3cf4b61..e248e70 100644 --- a/worker.js +++ b/worker.js @@ -9,7 +9,7 @@ const crypto = require('crypto'); const PHOTO_MAX_ATTEMPTS = 3; const PHOTO_AI_TIMEOUT_MS = 300000; -function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntries, sharp, photoAiUrl, uploadsDir, storage }) { +function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntries, sharp, photoAiUrl, uploadsDir, storage, bus }) { const AI_URL = photoAiUrl || process.env.PHOTO_AI_URL || ''; const IDLE_MIN = 2000; const IDLE_MAX = 30000; @@ -37,10 +37,17 @@ function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntrie }); } - function notify() { + function wakeLocal() { if (wake) wake(); } + function notify() { + wakeLocal(); + if (bus) bus.publish(); + } + + if (bus) bus.subscribe(wakeLocal); + async function isEnabled() { const v = await getSetting('photo_worker_enabled', 'true'); return String(v) !== 'false'; @@ -233,7 +240,7 @@ function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntrie return { start, notify, getInfo }; } -function createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl, defaultPrompt, model }) { +function createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl, defaultPrompt, model, bus }) { const AI_URL = aiUrl || process.env.AI_URL || 'http://text-corrector:8080'; const MODEL = model || process.env.AI_MODEL || 'qwen2.5-1.5b-instruct-q4_k_m.gguf'; const DEFAULT_PROMPT = defaultPrompt || process.env.AI_PROMPT || 'Ты — редактор текстов. Исправь ТОЛЬКО грамматические, орфографические и пунктуационные ошибки в тексте. Приведи к правильному регистру буквы. НЕ меняй слова, структуру предложений, стиль или смысл текста. Верни ТОЛЬКО исправленный текст без пояснений.'; @@ -266,10 +273,17 @@ function createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl, defaultProm }); } - function notify() { + function wakeLocal() { if (wake) wake(); } + function notify() { + wakeLocal(); + if (bus) bus.publish(); + } + + if (bus) bus.subscribe(wakeLocal); + async function isEnabled() { const v = await getSetting('ai_autocheck_enabled', 'true'); return String(v) !== 'false';