Files
iistwin/server/index.ts
Ильяс Султанов 254950f151 DI2 внедрение, таск 1: серверная интеграция Data-Insight2 (server/finance-di2)
- Копия DI2-сервера в server/finance-di2/ с правками: schema/db-client/cache/table-config/override-tables/google-auth/audit-agent/routes
- db-client — обёртка над существующим пулом server/finance/db-client (второй пул не создаётся, добавлен экспорт isConnectionError)
- ai-config.ts — чистое IO конфигов из ai-agent.ts без Telegram-поллинга; роуты /api/ai/toggle и /api/ai/status удалены, без compression/startAuditScheduler/autoStartIfEnabled
- Обёртка registerDi2Routes с auth-gate (authenticateToken + finance.manage), регистрация строго перед registerFinanceRoutes
- Статика /di2 из dist/public-di2 перед веткой vite/static (dev и prod)
- Фикс предсуществующего бага DI2: buildGlobalExclusionConditions без cats в /api/profitability
2026-09-06 19:27:43 +03:00

1310 lines
55 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import express, { type Request, Response, NextFunction } from "express";
import helmet from "helmet";
import cookieParser from "cookie-parser";
import { registerRoutes, publishNotificationSSE } from "./routes";
import { setupVite, serveStatic, log } from "./vite";
import { pushService } from "./services/push.service";
import { notificationService } from "./services/notification.service";
import { webPushService } from "./services/web-push.service";
import { startMedScheduleWorker } from "./medschedule/worker";
import { startAutomationScheduler } from "./workers/automation-scheduler";
import { startGpsWorker } from "./gps/worker";
import { storage } from "./storage";
import type { ReminderRecipient } from "@shared/schema";
import { db, withSuperAdmin } from "./db"; // lazy proxy — safe at module load; throws on first use if DATABASE_URL missing
import { sql, eq } from "drizzle-orm";
import fs from "fs";
import path from "path";
import { fileURLToPath } from "url";
import { isS3Enabled, ensureS3Bucket, streamFromS3 } from "./utils/s3";
import { getKnownEmbeddingDim, detectEmbeddingDimension, type LlmProviderConfig } from "./services/llm-provider";
import { decrypt as decryptSecret } from "./crypto";
import { fileUploads, errorLogs } from "@shared/schema";
import { authenticateToken, authenticateFileToken, type AuthenticatedRequest } from "./middleware/auth.middleware";
import { EXT_TO_MIME, getFileExt } from "./utils/upload";
import crypto from "crypto";
const app = express();
// Configure Express to trust proxy for accurate client IP identification
// This prevents ValidationError with express-rate-limit when X-Forwarded-For header is present
app.set('trust proxy', 1);
// Disable ETag to prevent stale API responses (critical for real-time notifications)
app.disable('etag');
// Ensure API responses are never cached by browsers/CDNs
app.use('/api', (req, res, next) => {
res.setHeader('Cache-Control', 'no-store, no-cache, must-revalidate, proxy-revalidate');
res.setHeader('Pragma', 'no-cache');
res.setHeader('Expires', '0');
next();
});
// Parse cookies early: file-serving routes below need the HttpOnly access_token
// cookie for authenticateFileToken / authenticateToken.
app.use(cookieParser());
// ── Presigned URLs for file preview (temporary unauthenticated access) ──────
interface PresignedToken {
fileKey: string;
organizationId: number;
expiresAt: number;
}
const presignedTokens = new Map<string, PresignedToken>();
const PRESIGNED_TTL_MS = 5 * 60 * 1000; // 5 minutes
function generatePresignedToken(fileKey: string, organizationId: number): string {
const token = crypto.randomBytes(32).toString('hex');
presignedTokens.set(token, { fileKey, organizationId, expiresAt: Date.now() + PRESIGNED_TTL_MS });
return token;
}
function validatePresignedToken(token: string, fileKey: string): number | null {
const entry = presignedTokens.get(token);
if (!entry) return null;
if (entry.fileKey !== fileKey) return null;
if (entry.expiresAt <= Date.now()) {
presignedTokens.delete(token);
return null;
}
return entry.organizationId;
}
function sweepPresignedTokens(): void {
const now = Date.now();
for (const [token, entry] of presignedTokens.entries()) {
if (entry.expiresAt <= now) {
presignedTokens.delete(token);
}
}
}
setInterval(sweepPresignedTokens, PRESIGNED_TTL_MS);
// ──────────────────────────────────────────────────────────────────────────────
// Helper: check file ownership — fail-closed (deny if untracked or DB error).
// Legacy files are backfilled from task_field_values / task_messages at startup
// so untracked after startup == foreign or unknown file.
async function canAccessFile(fileKey: string, organizationId: number): Promise<boolean> {
try {
const rows = await db.select({ orgId: fileUploads.organizationId })
.from(fileUploads)
.where(eq(fileUploads.fileKey, fileKey))
.limit(1);
if (rows.length === 0) {
// Not tracked even after startup backfill — deny (fail-closed)
return false;
}
return rows[0].orgId === organizationId;
} catch {
// DB error — fail closed to prevent cross-tenant leakage
return false;
}
}
// Middleware: try presigned token first, fall back to authenticateFileToken.
function tryPresignedOrAuth() {
return async (req: AuthenticatedRequest, res: Response, next: NextFunction) => {
const presigned = typeof req.query.presigned === 'string' ? req.query.presigned : null;
const fileKey = req.params.key || req.params.filename;
if (presigned && fileKey) {
const orgId = validatePresignedToken(presigned, fileKey);
if (orgId !== null) {
req.organizationId = orgId;
return next();
}
return res.status(403).json({ error: 'Недействительная или истёкшая presigned-ссылка' });
}
return authenticateFileToken(req, res, next);
};
}
// Endpoint to generate a presigned URL for file preview.
app.get('/api/files/:key/presigned', authenticateToken, async (req: AuthenticatedRequest, res: Response) => {
const key = req.params.key;
if (!key || key.includes('..') || key.includes('/')) {
return res.status(400).json({ error: 'Неверный ключ файла' });
}
const allowed = await canAccessFile(key, req.organizationId!);
if (!allowed) return res.status(403).json({ error: 'Доступ запрещён' });
const token = generatePresignedToken(key, req.organizationId!);
const presignedUrl = isS3Enabled
? `/api/files/${key}?presigned=${token}`
: `/uploads/${key}?presigned=${token}`;
res.json({ presignedUrl, expiresAt: new Date(Date.now() + PRESIGNED_TTL_MS).toISOString() });
});
// Helper to set Content-Type and Content-Disposition headers based on file extension.
function setFileResponseHeaders(
res: Response,
fileKey: string,
overrideContentType?: string,
overrideContentDisposition?: string,
originalName?: string | null,
) {
const ext = getFileExt(fileKey);
const contentType = overrideContentType || EXT_TO_MIME[ext]?.[0] || 'application/octet-stream';
const isPdf = ext === 'pdf';
const utf8Name = originalName ? `filename*=UTF-8''${encodeURIComponent(originalName)}` : '';
// PDFs are always shown inline (browser iframe preview).
// Other files use S3-provided disposition if available, defaulting to attachment.
// Prefer file_uploads.originalName (UTF-8) over raw fileKey to avoid encoding issues.
const contentDisposition = isPdf
? (utf8Name ? `inline; ${utf8Name}` : `inline; filename="${fileKey}"`)
: (overrideContentDisposition || (utf8Name ? `attachment; ${utf8Name}` : `attachment; filename="${fileKey}"`));
res.setHeader('Content-Type', contentType);
res.setHeader('Content-Disposition', contentDisposition);
res.setHeader('X-Content-Type-Options', 'nosniff');
}
async function getFileOriginalName(fileKey: string): Promise<string | undefined> {
try {
const rows = await db.select({ originalName: fileUploads.originalName })
.from(fileUploads)
.where(eq(fileUploads.fileKey, fileKey))
.limit(1);
return rows[0]?.originalName || undefined;
} catch {
return undefined;
}
}
if (isS3Enabled) {
// S3/MinIO mode: files served via /api/files/:key.
app.use('/api/files/:key', tryPresignedOrAuth(), async (req: AuthenticatedRequest, res: Response) => {
const key = req.params.key;
if (!key || key.includes('..') || key.includes('/')) {
return res.status(400).end();
}
const allowed = await canAccessFile(key, req.organizationId!);
if (!allowed) return res.status(403).json({ error: 'Доступ запрещён' });
const obj = await streamFromS3(key);
if (!obj) return res.status(404).end();
const originalName = await getFileOriginalName(key);
setFileResponseHeaders(res, key, obj.contentType, obj.contentDisposition || undefined, originalName);
if (obj.contentLength) res.setHeader('Content-Length', obj.contentLength);
(obj.body as any).pipe(res);
});
} else {
// Local mode: files served from uploads/.
const uploadsDir = path.resolve(process.cwd(), 'uploads');
if (!fs.existsSync(uploadsDir)) {
fs.mkdirSync(uploadsDir, { recursive: true });
}
app.use('/uploads/:filename', tryPresignedOrAuth(), async (req: AuthenticatedRequest, res: Response) => {
const filename = req.params.filename;
if (!filename || filename.includes('..') || filename.includes('/')) {
return res.status(400).end();
}
const allowed = await canAccessFile(filename, req.organizationId!);
if (!allowed) return res.status(403).json({ error: 'Доступ запрещён' });
const filePath = path.join(uploadsDir, filename);
if (!fs.existsSync(filePath)) return res.status(404).end();
const originalName = await getFileOriginalName(filename);
setFileResponseHeaders(res, filename, undefined, undefined, originalName);
res.sendFile(filePath);
});
}
// Security headers
// HSTS and CSP upgrade-insecure-requests are enabled only when the server is behind
// an HTTPS reverse proxy. This is detected per-request via:
// 1. APP_HTTPS_PROXY=true env var (static override, e.g. in docker-compose)
// 2. X-Forwarded-Proto: https request header (set by Nginx/Traefik automatically)
// In plain HTTP deployments neither condition is true, so these headers are omitted.
const httpsProxyEnv = process.env.APP_HTTPS_PROXY === 'true';
function isHttpsRequest(req: Request): boolean {
if (httpsProxyEnv) return true;
if (req.secure) return true;
const forwarded = req.headers['x-forwarded-proto'];
if (typeof forwarded === 'string') {
return forwarded.split(',')[0].trim() === 'https';
}
return false;
}
app.use(helmet({
contentSecurityPolicy: {
directives: {
defaultSrc: ["'self'"],
scriptSrc: ["'self'", "'unsafe-inline'", "'unsafe-eval'", "blob:", "https://api-maps.yandex.ru", "https://yastatic.net"],
styleSrc: ["'self'", "'unsafe-inline'", "https://fonts.googleapis.com"],
imgSrc: ["'self'", "data:", "blob:", "https:", "http:"],
connectSrc: ["'self'", "ws:", "wss:", "http:", "https:"],
fontSrc: ["'self'", "data:", "https://fonts.gstatic.com"],
workerSrc: ["'self'", "blob:"],
frameSrc: ["'self'", "https://docs.google.com"],
objectSrc: ["'none'"],
frameAncestors: ["'self'"],
},
},
frameguard: false,
hsts: false,
crossOriginEmbedderPolicy: false,
}));
// Conditionally add HSTS and upgrade-insecure-requests per request when HTTPS is detected
app.use((req: Request, res: Response, next: NextFunction) => {
if (isHttpsRequest(req)) {
res.setHeader('Strict-Transport-Security', 'max-age=31536000; includeSubDomains; preload');
const existingCsp = res.getHeader('Content-Security-Policy') as string | undefined;
if (existingCsp && !existingCsp.includes('upgrade-insecure-requests')) {
res.setHeader('Content-Security-Policy', existingCsp + '; upgrade-insecure-requests');
}
}
next();
});
app.use(express.json({ limit: '10mb' }));
app.use(express.urlencoded({ extended: false, limit: '10mb' }));
app.use((req, res, next) => {
const start = Date.now();
const path = req.path;
let capturedJsonResponse: Record<string, any> | undefined = undefined;
const originalResJson = res.json;
res.json = function (bodyJson, ...args) {
capturedJsonResponse = bodyJson;
return originalResJson.apply(res, [bodyJson, ...args]);
};
res.on("finish", () => {
const duration = Date.now() - start;
if (path.startsWith("/api")) {
let logLine = `${req.method} ${path} ${res.statusCode} in ${duration}ms`;
if (capturedJsonResponse) {
logLine += ` :: ${JSON.stringify(capturedJsonResponse)}`;
}
if (logLine.length > 80) {
logLine = logLine.slice(0, 79) + "…";
}
log(logLine);
// Журнал ошибок API (просмотр админом в Настройках → «Журнал ошибок»)
captureErrorLog(req, res.statusCode, path, capturedJsonResponse);
}
});
next();
});
// --- Захват ошибок API в таблицу error_logs ---
// Шум исключаем: 401 (просроченные токены/SSE), health-check и SSE-поток.
// Дедупликация: одинаковые (method, path, status, message) не чаще раза в 60 сек.
const ERROR_LOG_EXCLUDED_PATHS = new Set(['/api/health', '/api/events']);
const errorLogDedupe = new Map<string, number>();
function captureErrorLog(req: Request, statusCode: number, path: string, body?: Record<string, any>) {
try {
if (statusCode < 400 || statusCode === 401) return;
if (ERROR_LOG_EXCLUDED_PATHS.has(path)) return;
const rawMessage = (body && (body.error || body.message)) || '';
const message = String(rawMessage).slice(0, 2000) || `HTTP ${statusCode}`;
const key = `${req.method} ${path} ${statusCode} ${message.slice(0, 200)}`;
const now = Date.now();
const lastSeen = errorLogDedupe.get(key) ?? 0;
if (now - lastSeen < 60_000) return;
errorLogDedupe.set(key, now);
// Периодическая чистка map, чтобы не разрастался
if (errorLogDedupe.size > 5000) {
for (const [k, ts] of errorLogDedupe) {
if (now - ts > 60_000) errorLogDedupe.delete(k);
}
}
const authReq = req as AuthenticatedRequest;
db.insert(errorLogs).values({
organizationId: authReq.organizationId ?? null,
userId: authReq.user?.id ?? null,
method: req.method,
path: path.slice(0, 500),
statusCode,
message,
}).catch((e: unknown) => console.error('[error-logs] insert failed:', e));
} catch (e) {
console.error('[error-logs] capture failed:', e);
}
}
/**
* Page seeds — canonical JS code stored in server/page-seeds/ as fallback/recovery.
* Seeds are applied ONLY when the current DB code contains known critical bugs.
* Good MCP edits are preserved across restarts; only broken code gets recovered.
*
* Mapping: { custom_pages.id => seed filename (without .js) }
*/
const PAGE_SEEDS: Record<number, string> = {
2: "contract-dashboard",
};
/** Returns true if the page code has any known critical bugs that break rendering. */
function pageCodeIsBroken(code: string): boolean {
// Wrong task navigation route — causes 404
if (/ctx\.navigate\(['"`]\/tasks\//.test(code)) return true;
// SelectItem used without required sibling components — causes React error
if (
code.includes("SelectItem") &&
(!code.includes("SelectTrigger") || !code.includes("SelectContent") || !code.includes("SelectValue"))
) return true;
// Missing Widget return — page won't render at all
if (!/return\s+function\s+Widget/.test(code)) return true;
return false;
}
async function runStartupDataPatches() {
// Ensure system_audit_log table exists
await db.execute(sql`
CREATE TABLE IF NOT EXISTS system_audit_log (
id SERIAL PRIMARY KEY,
action VARCHAR(100) NOT NULL,
user_id INTEGER REFERENCES users(id) ON DELETE SET NULL,
organization_id INTEGER REFERENCES organizations(id) ON DELETE SET NULL,
task_id INTEGER,
details JSONB,
ip VARCHAR(45),
user_agent TEXT,
created_at TIMESTAMP DEFAULT NOW()
)
`);
await db.execute(sql`
CREATE INDEX IF NOT EXISTS system_audit_log_action_created_idx
ON system_audit_log (action, created_at DESC)
`);
await db.execute(sql`
CREATE INDEX IF NOT EXISTS system_audit_log_org_created_idx
ON system_audit_log (organization_id, created_at DESC)
`);
// Ensure task_audit_log table exists
await db.execute(sql`
CREATE TABLE IF NOT EXISTS task_audit_log (
id SERIAL PRIMARY KEY,
task_id INTEGER NOT NULL REFERENCES tasks(id) ON DELETE CASCADE,
organization_id INTEGER NOT NULL REFERENCES organizations(id) ON DELETE CASCADE,
action VARCHAR(50) NOT NULL,
field_id INTEGER,
field_name TEXT,
old_value JSONB,
new_value JSONB,
changed_by INTEGER REFERENCES users(id),
changed_by_name TEXT,
metadata JSONB,
created_at TIMESTAMP DEFAULT NOW()
)
`);
// Ensure task_reminders table exists
await db.execute(sql`
CREATE TABLE IF NOT EXISTS task_reminders (
id SERIAL PRIMARY KEY,
task_id INTEGER NOT NULL REFERENCES tasks(id) ON DELETE CASCADE,
organization_id INTEGER NOT NULL REFERENCES organizations(id) ON DELETE CASCADE,
created_by_user_id INTEGER NOT NULL REFERENCES users(id),
remind_at TIMESTAMP NOT NULL,
note TEXT,
recipients JSONB NOT NULL DEFAULT '[]',
is_sent BOOLEAN DEFAULT FALSE,
created_at TIMESTAMP DEFAULT NOW()
)
`);
// Ensure task.reminder event type exists in notification_event_types
await db.execute(sql`
INSERT INTO notification_event_types (code, name, description, category, is_active)
VALUES ('task.reminder', 'Напоминание', 'Напоминание о задаче', 'task', true)
ON CONFLICT (code) DO NOTHING
`);
// Billing: add billing_blocked column to organizations if missing
await db.execute(sql`
ALTER TABLE organizations ADD COLUMN IF NOT EXISTS billing_blocked BOOLEAN DEFAULT FALSE
`);
// Billing settings table (1:1 with organizations)
await db.execute(sql`
CREATE TABLE IF NOT EXISTS organization_billing (
organization_id INTEGER PRIMARY KEY REFERENCES organizations(id) ON DELETE CASCADE,
price_per_user NUMERIC(10,2) NOT NULL DEFAULT 0,
balance NUMERIC(10,2) NOT NULL DEFAULT 0,
currency VARCHAR(10) NOT NULL DEFAULT 'RUB',
next_billing_date DATE,
blocked_at TIMESTAMP,
updated_at TIMESTAMP DEFAULT NOW()
)
`);
// Billing transactions table
await db.execute(sql`
CREATE TABLE IF NOT EXISTS billing_transactions (
id SERIAL PRIMARY KEY,
organization_id INTEGER NOT NULL REFERENCES organizations(id) ON DELETE CASCADE,
amount NUMERIC(10,2) NOT NULL,
type VARCHAR(10) NOT NULL,
description TEXT,
created_at TIMESTAMP DEFAULT NOW(),
created_by INTEGER
)
`);
// Денормализованный organization_id в tasks — добавляем если отсутствует и заполняем из forms
await db.execute(sql`ALTER TABLE tasks ADD COLUMN IF NOT EXISTS organization_id INTEGER REFERENCES organizations(id) ON DELETE CASCADE`);
await db.execute(sql`UPDATE tasks SET organization_id = forms.organization_id FROM forms WHERE tasks.form_id = forms.id AND tasks.organization_id IS NULL`);
// Enforce NOT NULL after backfill (idempotent — PG ignores if already NOT NULL)
await db.execute(sql`ALTER TABLE tasks ALTER COLUMN organization_id SET NOT NULL`).catch((e: unknown) => {
log(`Startup: tasks.organization_id SET NOT NULL skipped — ${e instanceof Error ? e.message : String(e)}`);
});
// Удаляем устаревшую колонку hierarchy_number из tasks
await db.execute(sql`ALTER TABLE tasks DROP COLUMN IF EXISTS hierarchy_number`);
// Группа «Наша Фирма»: добавляем флаг is_default в conversations
await db.execute(sql`ALTER TABLE conversations ADD COLUMN IF NOT EXISTS is_default BOOLEAN DEFAULT FALSE`);
// Миграция: создать группу «Наша Фирма» для существующих организаций, где её ещё нет
await db.execute(sql`
INSERT INTO conversations (organization_id, type, name, created_by, is_default)
SELECT o.id, 'group', 'Наша Фирма', (
SELECT u.id FROM users u WHERE u.organization_id = o.id ORDER BY u.id LIMIT 1
), TRUE
FROM organizations o
WHERE NOT EXISTS (
SELECT 1 FROM conversations c
WHERE c.organization_id = o.id AND c.is_default = TRUE
)
AND EXISTS (
SELECT 1 FROM users u WHERE u.organization_id = o.id
)
`);
// Добавить всех активных пользователей в группу «Наша Фирма» если они ещё не участники
await db.execute(sql`
INSERT INTO conversation_members (conversation_id, user_id)
SELECT c.id, u.id
FROM conversations c
JOIN users u ON u.organization_id = c.organization_id
WHERE c.is_default = TRUE
AND u.is_active = TRUE
ON CONFLICT DO NOTHING
`);
// Аудит колонок tasks (задача #199): все остальные колонки используются активно:
// id, form_id, title, description, current_status_id, assigned_to, created_by,
// created_at, updated_at, completed_at, is_completed, due_date,
// parent_task_id, depth, position, organization_id
// Лишних колонок помимо hierarchy_number не обнаружено.
// RAG: enable pgvector extension and create rag_embeddings table
// All steps are individually guarded so a missing pgvector extension does
// not prevent the server from starting — RAG simply stays disabled.
let pgvectorAvailable = false;
try {
await db.execute(sql`CREATE EXTENSION IF NOT EXISTS vector`);
pgvectorAvailable = true;
} catch {
log('Startup: pgvector extension not available — RAG semantic search will be disabled');
}
// Always ensure the ai_summary column exists on forms (no pgvector dependency)
await db.execute(sql`ALTER TABLE forms ADD COLUMN IF NOT EXISTS ai_summary TEXT`).catch((err: unknown) => {
log(`Startup: could not add ai_summary column: ${err instanceof Error ? err.message : String(err)}`);
});
await db.execute(sql`ALTER TABLE forms ADD COLUMN IF NOT EXISTS ai_summary_is_auto BOOLEAN DEFAULT TRUE`).catch((err: unknown) => {
log(`Startup: could not add ai_summary_is_auto column: ${err instanceof Error ? err.message : String(err)}`);
});
if (pgvectorAvailable) {
try {
await db.execute(sql`
CREATE TABLE IF NOT EXISTS rag_embeddings (
id SERIAL PRIMARY KEY,
organization_id INTEGER NOT NULL REFERENCES organizations(id) ON DELETE CASCADE,
entity_type VARCHAR(20) NOT NULL,
entity_id INTEGER NOT NULL,
content TEXT NOT NULL,
embedding vector(1),
metadata JSONB DEFAULT '{}',
updated_at TIMESTAMP DEFAULT NOW()
)
`);
await db.execute(sql`
CREATE UNIQUE INDEX IF NOT EXISTS rag_embeddings_org_entity_type_id_key
ON rag_embeddings (organization_id, entity_type, entity_id)
`);
await db.execute(sql`
CREATE INDEX IF NOT EXISTS rag_embeddings_org_type_idx
ON rag_embeddings (organization_id, entity_type)
`);
// HNSW index for cosine similarity
await db.execute(sql`
CREATE INDEX IF NOT EXISTS rag_embeddings_vector_hnsw_idx
ON rag_embeddings USING hnsw (embedding vector_cosine_ops)
`).catch(() => {
log('Startup: HNSW vector index creation skipped (may already exist or require superuser)');
});
log('Startup: RAG embeddings table and indexes are ready');
} catch (err: unknown) {
log(`Startup: RAG table/index setup error: ${err instanceof Error ? err.message : String(err)}`);
}
}
// RAG Settings table — ensure it exists and has all current columns
await db.execute(sql`
CREATE TABLE IF NOT EXISTS rag_settings (
organization_id INTEGER PRIMARY KEY REFERENCES organizations(id) ON DELETE CASCADE,
provider VARCHAR(50) NOT NULL DEFAULT 'openai',
api_key TEXT,
base_url TEXT,
embedding_model VARCHAR(200) NOT NULL DEFAULT 'text-embedding-3-small',
updated_at TIMESTAMP DEFAULT NOW()
)
`).catch(() => {});
await db.execute(sql`ALTER TABLE rag_settings ADD COLUMN IF NOT EXISTS chat_model VARCHAR(200)`).catch(() => {});
await db.execute(sql`ALTER TABLE rag_settings ADD COLUMN IF NOT EXISTS max_chunk_chars INTEGER`).catch(() => {});
await db.execute(sql`ALTER TABLE rag_settings ADD COLUMN IF NOT EXISTS summarization_enabled BOOLEAN DEFAULT TRUE`).catch(() => {});
// RAG: auto-migrate embedding column dimension when model changes.
// Determines the correct vector size from the configured embedding model and
// aligns the DB column if they differ. Steps:
// 1. Read current column dimension via format_type() — unambiguous across providers
// 2. Find the effective model from rag_settings + llm_providers (ORDER BY org for determinism)
// 3. Lookup dim from known-models table; fallback to a live test embedding (5 s timeout)
// Live test decrypts API key so it works for OpenAI/openai_compatible too
// 4. Default to 1536 (OpenAI) when no config or provider is unreachable
// 5. On mismatch: drop HNSW idx → DELETE rows first → ALTER column → recreate HNSW idx
// 6. Emit one authoritative log line: "embedding dimension set to N (model)"
if (pgvectorAvailable) {
try {
// Step 1 — read current column type as string e.g. "vector(1536)" via format_type()
// This is unambiguous and avoids the atttypmod offset debate between pg types.
const dimResult = await db.execute(sql`
SELECT format_type(a.atttypid, a.atttypmod) AS col_type
FROM pg_attribute a
JOIN pg_class c ON c.oid = a.attrelid
WHERE c.relname = 'rag_embeddings'
AND a.attname = 'embedding'
AND a.attnum > 0
AND NOT a.attisdropped
`);
const colTypeStr = dimResult.rows[0]
? (dimResult.rows[0] as { col_type: string }).col_type
: null;
const dimMatch = colTypeStr?.match(/vector\((\d+)\)/);
const currentDim = dimMatch ? Number(dimMatch[1]) : null;
// Step 2 — find effective model (legacy embedding_model OR registry enabled_models[0])
// ORDER BY organization_id ASC for deterministic selection in multi-org setups.
let targetDim: number | null = null;
let targetModel = 'text-embedding-3-small (default)';
let liveTestConfig: LlmProviderConfig | null = null;
try {
const modelsResult = await db.execute(sql`
SELECT
COALESCE(
NULLIF(rs.embedding_model, ''),
(lp.enabled_models->>0)
) AS effective_model,
COALESCE(rs.provider, lp.provider_type, 'openai') AS provider,
COALESCE(rs.base_url, lp.base_url) AS base_url,
COALESCE(rs.api_key, lp.api_key) AS raw_api_key
FROM rag_settings rs
LEFT JOIN llm_providers lp ON lp.id = rs.embedding_provider_id
WHERE rs.embedding_model IS NOT NULL
OR lp.id IS NOT NULL
ORDER BY rs.organization_id ASC
LIMIT 10
`);
type ConfigRow = {
effective_model: string | null;
provider: string;
base_url: string | null;
raw_api_key: string | null;
};
const seenDims = new Set<number>();
for (const row of modelsResult.rows as ConfigRow[]) {
if (!row.effective_model) continue;
// Step 3a — lookup table (instant, no network)
const knownDim = getKnownEmbeddingDim(row.effective_model);
if (knownDim) {
seenDims.add(knownDim);
if (!targetDim) {
targetDim = knownDim;
targetModel = row.effective_model;
// don't break — keep scanning to detect multi-org ambiguity
}
}
// Step 3b — build config for live test (model not in lookup table)
if (!liveTestConfig && !knownDim) {
// Decrypt API key so OpenAI/openai_compatible providers can be probed too
let apiKey: string | null = null;
if (row.raw_api_key) {
try { apiKey = decryptSecret(row.raw_api_key); } catch { apiKey = row.raw_api_key; }
}
liveTestConfig = {
provider: row.provider,
apiKey,
baseUrl: row.base_url ?? null,
embeddingModel: row.effective_model,
chatModel: null,
maxChunkChars: null,
summarizationEnabled: false,
customHeaders: {},
};
if (!targetModel || targetModel === 'text-embedding-3-small (default)') {
targetModel = row.effective_model;
}
}
}
// Warn when orgs have different embedding models with different dimensions
// (only one global vector dimension is supported — first org wins)
if (seenDims.size > 1) {
log(`Startup RAG: WARNING — multiple orgs have different embedding model dimensions (${[...seenDims].join(', ')}). Using first org's model "${targetModel}" (dim=${targetDim}). All orgs should use the same embedding model.`);
}
} catch {
// rag_settings / llm_providers may not exist on very first startup — safe to ignore
}
// Step 3b — live test for models not in the lookup table (5 s timeout, non-blocking)
if (!targetDim && liveTestConfig) {
try {
const liveDim = await Promise.race<number | null>([
detectEmbeddingDimension(liveTestConfig),
new Promise<null>((resolve) => setTimeout(() => resolve(null), 5000)),
]);
if (liveDim) targetDim = liveDim;
} catch {
// Provider unreachable — handled below
}
}
// OpenAI default dimension — used only when no model is configured at all
const RAG_DEFAULT_DIM = 1536;
// Step 4 — resolve final target dimension
if (!targetDim) {
if (liveTestConfig) {
// Unknown model is configured but provider is temporarily unreachable.
// Do NOT migrate (we don't know the correct dimension) — keep current column as-is.
if (currentDim) {
log(`Startup RAG: provider for model "${targetModel}" unreachable — keeping current dimension vector(${currentDim}) to avoid destructive migration`);
targetDim = currentDim; // skip step 5
} else {
// Brand-new install, no column yet — use OpenAI default and warn
targetDim = RAG_DEFAULT_DIM;
targetModel = `text-embedding-3-small (default, model "${targetModel}" unreachable)`;
}
} else {
// No embedding model configured anywhere — safe to default to OpenAI
targetDim = RAG_DEFAULT_DIM;
targetModel = 'text-embedding-3-small (default)';
}
}
// Step 5 — migrate column when dimension has changed.
// Correct order: drop index → delete rows (ALTER TYPE fails on incompatible data) → ALTER → recreate index
if (currentDim && currentDim !== targetDim) {
log(`Startup RAG: dimension mismatch — column=vector(${currentDim}), model "${targetModel}"=vector(${targetDim}). Migrating…`);
await db.execute(sql`DROP INDEX IF EXISTS rag_embeddings_vector_hnsw_idx`).catch(() => {});
await db.execute(sql`DELETE FROM rag_embeddings`);
await db.execute(sql.raw(`ALTER TABLE rag_embeddings ALTER COLUMN embedding TYPE vector(${targetDim})`));
await db.execute(
sql.raw(`CREATE INDEX IF NOT EXISTS rag_embeddings_vector_hnsw_idx ON rag_embeddings USING hnsw (embedding vector_cosine_ops)`)
).catch(() => {
log('Startup RAG: HNSW index recreation skipped (may require superuser)');
});
}
// Step 6 — single authoritative log line emitted regardless of whether migration ran
log(`Startup RAG: embedding dimension set to ${targetDim} (${targetModel})`);
} catch (err: unknown) {
log(`Startup RAG: dimension check error: ${err instanceof Error ? err.message : String(err)}`);
}
}
await db.execute(sql`ALTER TABLE conversation_members ADD COLUMN IF NOT EXISTS reaction_unread_count INTEGER NOT NULL DEFAULT 0`).catch(() => {});
await db.execute(sql`ALTER TABLE conversation_members ADD COLUMN IF NOT EXISTS muted_at TIMESTAMP`).catch(() => {});
// HMAC API keys: add is_legacy column — existing rows marked TRUE (SHA-256), new rows default FALSE (HMAC)
await db.execute(sql`ALTER TABLE organization_api_keys ADD COLUMN IF NOT EXISTS is_legacy BOOLEAN NOT NULL DEFAULT TRUE`).catch(() => {});
// Ensure new rows inserted after this migration start with is_legacy=false
await db.execute(sql`ALTER TABLE organization_api_keys ALTER COLUMN is_legacy SET DEFAULT FALSE`).catch(() => {});
// Encrypt existing plaintext bot secrets (webhookSecret + mcpServers[].apiKey).
// decrypt() is backward-compatible: returns input unchanged if it doesn't look encrypted.
// We detect plaintext by checking for the ':' separator used in our AES-256-GCM format.
// This patch is idempotent — already-encrypted values are skipped.
try {
const { encrypt: _encrypt } = await import('./crypto');
const SESSION_SECRET_PRESENT = !!process.env.SESSION_SECRET;
if (SESSION_SECRET_PRESENT) {
const allBots = await db.execute(sql`SELECT id, webhook_secret, mcp_servers FROM bots`);
for (const row of allBots.rows as Array<{ id: number; webhook_secret: string | null; mcp_servers: unknown }>) {
const updates: Record<string, unknown> = {};
// Detect AES-GCM encrypted values by checking for the 4-part hex format (salt:iv:authTag:ciphertext)
// or the legacy 3-part format (iv:authTag:ciphertext). A plain secret might contain ':' too,
// so we validate structurally: each part must be a non-empty hex string.
const AES_GCM_RE = /^[0-9a-f]+:[0-9a-f]+:[0-9a-f]+:[0-9a-f]+$/i;
const AES_GCM_V1_RE = /^[0-9a-f]+:[0-9a-f]+:[0-9a-f]+$/i;
const isAlreadyEncrypted = (v: string) => AES_GCM_RE.test(v) || AES_GCM_V1_RE.test(v);
// Encrypt plaintext webhookSecret
if (row.webhook_secret && !isAlreadyEncrypted(row.webhook_secret)) {
updates.webhook_secret = _encrypt(row.webhook_secret);
}
// Encrypt plaintext mcpServers[].apiKey entries
if (Array.isArray(row.mcp_servers)) {
let changed = false;
const encryptedServers = row.mcp_servers.map((srv: { url?: string; apiKey?: string; name?: string }) => {
if (srv.apiKey && !isAlreadyEncrypted(srv.apiKey)) {
changed = true;
return { ...srv, apiKey: _encrypt(srv.apiKey) };
}
return srv;
});
if (changed) updates.mcp_servers = JSON.stringify(encryptedServers);
}
if (updates.webhook_secret !== undefined) {
const encSecret = updates.webhook_secret as string;
await db.execute(sql`UPDATE bots SET webhook_secret = ${encSecret} WHERE id = ${row.id}`).catch((e: unknown) => {
log(`Startup: bot ${row.id} webhookSecret encryption error: ${e instanceof Error ? e.message : String(e)}`);
});
}
if (updates.mcp_servers !== undefined) {
const encServersJson = updates.mcp_servers as string;
await db.execute(sql`UPDATE bots SET mcp_servers = ${encServersJson}::jsonb WHERE id = ${row.id}`).catch((e: unknown) => {
log(`Startup: bot ${row.id} mcpServers encryption error: ${e instanceof Error ? e.message : String(e)}`);
});
}
if (Object.keys(updates).length > 0) {
log(`Startup: encrypted plaintext secrets for bot id=${row.id}`);
}
}
} else {
log('Startup: SESSION_SECRET not set — skipping bot secret backfill encryption');
}
} catch (err: unknown) {
log(`Startup: bot secret backfill error: ${err instanceof Error ? err.message : String(err)}`);
}
// System Config table — generic key-value store for server-side configuration
await db.execute(sql`
CREATE TABLE IF NOT EXISTS system_config (
key VARCHAR(255) PRIMARY KEY,
value TEXT NOT NULL,
updated_at TIMESTAMP DEFAULT NOW()
)
`).catch(() => {});
// Backfill file_uploads from historical task field values and chat attachments.
// This runs once at startup so canAccessFile can use fail-closed logic for all files.
// Single-file field values: { url, name, size }
await db.execute(sql`
INSERT INTO file_uploads (organization_id, uploaded_by, file_key, original_name, size_bytes, task_id, field_id)
SELECT DISTINCT
t.organization_id,
t.created_by,
CASE
WHEN tfv.value->>'url' LIKE '/uploads/%' THEN substring(tfv.value->>'url' FROM 10)
WHEN tfv.value->>'url' LIKE '/api/files/%' THEN substring(tfv.value->>'url' FROM 12)
END AS file_key,
tfv.value->>'name',
(tfv.value->>'size')::integer,
tfv.task_id,
tfv.field_id
FROM task_field_values tfv
JOIN tasks t ON tfv.task_id = t.id
JOIN form_fields ff ON tfv.field_id = ff.id
WHERE ff.type = 'file'
AND jsonb_typeof(tfv.value) = 'object'
AND (tfv.value->>'url' LIKE '/uploads/%' OR tfv.value->>'url' LIKE '/api/files/%')
ON CONFLICT (file_key) DO NOTHING
`).catch(() => {});
// Multi-file field values: [{ url, name, size }, ...]
await db.execute(sql`
INSERT INTO file_uploads (organization_id, uploaded_by, file_key, original_name, size_bytes, task_id, field_id)
SELECT DISTINCT
t.organization_id,
t.created_by,
CASE
WHEN elem->>'url' LIKE '/uploads/%' THEN substring(elem->>'url' FROM 10)
WHEN elem->>'url' LIKE '/api/files/%' THEN substring(elem->>'url' FROM 12)
END AS file_key,
elem->>'name',
(elem->>'size')::integer,
tfv.task_id,
tfv.field_id
FROM task_field_values tfv
JOIN tasks t ON tfv.task_id = t.id
JOIN form_fields ff ON tfv.field_id = ff.id
CROSS JOIN LATERAL jsonb_array_elements(tfv.value) AS elem
WHERE ff.type = 'file'
AND jsonb_typeof(tfv.value) = 'array'
AND (elem->>'url' LIKE '/uploads/%' OR elem->>'url' LIKE '/api/files/%')
ON CONFLICT (file_key) DO NOTHING
`).catch(() => {});
// Chat attachments: task_messages.attachments = [{ url, name, size }, ...]
await db.execute(sql`
INSERT INTO file_uploads (organization_id, uploaded_by, file_key, original_name, size_bytes, task_id)
SELECT DISTINCT
f.organization_id,
COALESCE(tm.author_id, t.created_by),
CASE
WHEN att->>'url' LIKE '/uploads/%' THEN substring(att->>'url' FROM 10)
WHEN att->>'url' LIKE '/api/files/%' THEN substring(att->>'url' FROM 12)
END AS file_key,
att->>'name',
(att->>'size')::integer,
tm.task_id
FROM task_messages tm
JOIN tasks t ON tm.task_id = t.id
JOIN forms f ON t.form_id = f.id
CROSS JOIN LATERAL jsonb_array_elements(tm.attachments) AS att
WHERE tm.attachments IS NOT NULL
AND jsonb_typeof(tm.attachments) = 'array'
AND (att->>'url' LIKE '/uploads/%' OR att->>'url' LIKE '/api/files/%')
ON CONFLICT (file_key) DO NOTHING
`).catch(() => {});
// RLS policies are created by migration 0011_rls_tasks_notnull.sql.
// Startup only activates ENABLE + FORCE ROW LEVEL SECURITY when ENABLE_RLS=true,
// then verifies all tables are correctly protected before allowing traffic.
const ALL_RLS_TABLES = [
"tasks", "users", "forms", "task_relations", "user_notifications",
"message_reads", "bookmarks", "bookmark_folders", "bots", "bot_subscriptions",
"data_tables", "invitations", "user_custom_fields", "device_tokens",
"web_push_subscriptions", "organization_api_keys", "automations",
"task_reminders", "task_audit_log", "conversations", "field_templates",
"custom_tab_modules", "external_services", "notification_subscriptions",
"form_fields", "form_statuses", "form_tabs", "status_transitions",
"task_messages", "task_field_values", "field_history",
];
if (process.env.ENABLE_RLS === "true") {
const tableList = ALL_RLS_TABLES.map(t => `'${t}'`).join(",");
// Verify all expected tables exist before activating RLS to detect schema drift early.
const existsCheck = await db.execute(sql.raw(`
SELECT array_agg(tablename) AS found
FROM pg_tables
WHERE schemaname = 'public' AND tablename IN (${tableList})
`));
const foundTables = new Set<string>(
((existsCheck.rows[0] as Record<string, unknown>)?.found as string[] | null) ?? []
);
const missingTables = ALL_RLS_TABLES.filter(t => !foundTables.has(t));
if (missingTables.length > 0) {
throw new Error(
`Startup RLS: schema drift detected — tables not found: ${missingTables.join(", ")}. ` +
`Run all pending migrations before enabling RLS.`
);
}
for (const tbl of ALL_RLS_TABLES) {
await db.execute(sql.raw(`ALTER TABLE "${tbl}" ENABLE ROW LEVEL SECURITY`));
await db.execute(sql.raw(`ALTER TABLE "${tbl}" FORCE ROW LEVEL SECURITY`));
}
const rlsCheck = await db.execute(sql.raw(`
SELECT COUNT(*) AS cnt
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE n.nspname = 'public'
AND c.relname IN (${tableList})
AND c.relrowsecurity = TRUE
AND c.relforcerowsecurity = TRUE
`));
const rlsVerified = Number((rlsCheck.rows[0] as Record<string, unknown>)?.cnt ?? 0);
if (rlsVerified < ALL_RLS_TABLES.length) {
throw new Error(
`Startup RLS: ENABLE+FORCE verification failed — only ${rlsVerified}/${ALL_RLS_TABLES.length} tables active`
);
}
const policyCheck = await db.execute(sql.raw(`
SELECT COUNT(*) AS cnt
FROM pg_policies
WHERE schemaname = 'public' AND policyname = 'tenant_iso'
AND tablename IN (${tableList})
`));
const policyVerified = Number((policyCheck.rows[0] as Record<string, unknown>)?.cnt ?? 0);
if (policyVerified < ALL_RLS_TABLES.length) {
throw new Error(
`Startup RLS: policy verification failed — only ${policyVerified}/${ALL_RLS_TABLES.length} tenant_iso policies found. ` +
`Run migration 0011_rls_tasks_notnull.sql first.`
);
}
log(`Startup RLS: active — ENABLE+FORCE verified on ${rlsVerified} tables, ${policyVerified} policies`);
} else {
log("Startup RLS: ENABLE_RLS not set — policies exist but inactive (safe default)");
}
const seedDir = path.resolve(
path.dirname(
typeof __filename !== "undefined"
? __filename
: fileURLToPath(import.meta.url)
),
"page-seeds"
);
for (const [pageId, seedName] of Object.entries(PAGE_SEEDS)) {
const seedPath = path.join(seedDir, `${seedName}.js`);
try {
const rows = await db.execute(
sql`SELECT code FROM custom_pages WHERE id = ${parseInt(pageId, 10)}`
);
const row = rows.rows[0] as { code: string } | undefined;
if (!row) continue;
if (!pageCodeIsBroken(row.code)) {
// Code is fine — preserve it, skip recovery
continue;
}
if (!fs.existsSync(seedPath)) {
log(`Startup: page ${pageId} has broken code but seed file is missing: ${seedPath}`);
continue;
}
const seedCode = fs.readFileSync(seedPath, "utf8");
await db.execute(
sql`UPDATE custom_pages SET code = ${seedCode}, updated_at = NOW() WHERE id = ${parseInt(pageId, 10)}`
);
log(`Startup: recovered broken page '${seedName}' (id=${pageId}) from seed file`);
} catch (err) {
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 () => {
// Verify database connectivity. db is a lazy proxy: first access triggers pool initialisation,
// which throws a descriptive Error if DATABASE_URL is not set.
// Errors are caught here, logged, and re-thrown; Node.js terminates via unhandled rejection.
try {
await db.execute(sql`SELECT 1`);
} catch (dbConnErr) {
const msg = dbConnErr instanceof Error ? dbConnErr.message : String(dbConnErr);
console.error(`[STARTUP] FATAL: ${msg}`);
console.error('[STARTUP] Check DATABASE_URL and ensure the PostgreSQL server is reachable.');
throw dbConnErr; // propagate — Node.js exits via unhandled promise rejection
}
await runStartupDataPatches();
// Check VAPID key rotation after DB migrations are complete
await webPushService.checkVapidKeyRotation();
// Инициализируем S3 bucket при старте (если настроен MinIO/S3)
if (isS3Enabled) {
await ensureS3Bucket();
log('S3/MinIO хранилище инициализировано');
}
const server = await registerRoutes(app);
// Static files for VPN client apps (olcbox APK/IPA)
app.use('/apps', express.static('/app/apps'));
// Data-Insight2 sub-app (финансы v2) — сборка кладётся в dist/public-di2 (Task 2).
// Блок стоит ПЕРЕД setupVite/serveStatic, чтобы и в dev, и в prod /di2/* не уходил в main index.html.
const di2Dist = path.resolve(import.meta.dirname, "public-di2");
if (fs.existsSync(di2Dist)) {
app.use("/di2", express.static(di2Dist));
app.use("/di2", (_req, res) => res.sendFile(path.resolve(di2Dist, "index.html")));
}
// Seed application roles and permissions
try {
await storage.seedAppRolesAndPermissions();
log('App roles and permissions seeded');
} catch (e) {
console.error('Failed to seed app roles:', e);
}
app.use((err: any, _req: Request, res: Response, _next: NextFunction) => {
const status = err.status || err.statusCode || 500;
const message = err.message || "Internal Server Error";
if (res.headersSent) {
console.error('[express] Error after headers sent:', err);
return;
}
res.status(status).json({ message });
if (status >= 500) {
console.error('[express] Server error:', err);
}
});
// importantly only setup vite in development and after
// setting up all the other routes so the catch-all route
// doesn't interfere with the other routes
if (app.get("env") === "development") {
await setupVite(app, server);
} else {
serveStatic(app);
}
// Start push notification worker (Pushy)
pushService.startQueueProcessor(5000);
log('Push notification worker started (Pushy)');
// Start MedSchedule worker (medication reminders, daily check-ins, expiry alerts)
startMedScheduleWorker();
log('MedSchedule worker started');
// Start automation scheduler (schedule-trigger automations, daily at configured MSK time)
startAutomationScheduler();
log('Automation scheduler started');
// Start GPS worker (offline detection for tracked assets)
startGpsWorker();
log('GPS worker started');
// Clean up duplicate subscriptions then backfill defaults for existing users
(async () => {
try {
await notificationService.cleanupDuplicateSubscriptions();
const { db, withSuperAdmin } = await import('./db');
const { users } = await import('@shared/schema');
const allUsers = await withSuperAdmin(() =>
db.select({ id: users.id, organizationId: users.organizationId }).from(users)
);
let count = 0;
for (const u of allUsers) {
await notificationService.createDefaultSubscriptions(u.id, u.organizationId);
count++;
}
if (count > 0) log(`[Notifications] Backfilled default subscriptions for ${count} users`);
} catch (err) {
console.error('[Notifications] Backfill error:', err);
}
})();
// Start task reminder worker (every 60 seconds)
setInterval(async () => {
try {
const dueReminders = await withSuperAdmin(() => storage.getDueReminders());
for (const reminder of dueReminders) {
try {
const recipientUserIds = new Set<number>();
const recipients = reminder.recipients as ReminderRecipient[];
const [roleResults, task] = await withSuperAdmin(async () => {
const roleResultsInner: number[] = [];
for (const rec of recipients) {
if (rec.type === 'user') {
roleResultsInner.push(rec.userId);
} else if (rec.type === 'role') {
const us = await storage.getUsersByRole(rec.role, reminder.organizationId);
for (const u of us) roleResultsInner.push(u.id);
} else if (rec.type === 'orgRole') {
const us = await storage.getUsersByOrgRoleId(rec.roleId, reminder.organizationId);
for (const u of us) roleResultsInner.push(u.id);
}
}
const t = await storage.getTask(reminder.taskId, reminder.organizationId);
return [roleResultsInner, t] as const;
});
for (const uid of roleResults) recipientUserIds.add(uid);
const taskTitle = task?.title || `#${reminder.taskId}`;
const formId = task?.formId;
if (recipientUserIds.size > 0) {
const note = reminder.note;
const notifiedUserIds = await notificationService.processEvent({
type: 'task.reminder',
organizationId: reminder.organizationId,
triggeredBy: reminder.createdByUserId,
taskId: reminder.taskId,
formId,
targetUserIds: Array.from(recipientUserIds),
payload: {
taskTitle,
message: note ? `Напоминание: ${note}` : `Напоминание о задаче "${taskTitle}"`,
},
timestamp: new Date(),
});
for (const uid of notifiedUserIds) {
publishNotificationSSE(uid, reminder.organizationId, {
type: 'task.reminder',
taskId: reminder.taskId,
});
}
}
await withSuperAdmin(() => storage.markReminderSent(reminder.id));
} catch (reminderErr) {
log(`Reminder worker: failed to process reminder ${reminder.id}: ${reminderErr instanceof Error ? reminderErr.message : String(reminderErr)}`);
await withSuperAdmin(() => storage.markReminderSent(reminder.id)).catch(() => {});
}
}
} catch (err) {
log(`Reminder worker error: ${err instanceof Error ? err.message : String(err)}`);
}
}, 60_000);
log('Task reminder worker started (60s interval)');
// Billing cycle worker — runs once a day (every 24 hours)
async function runBillingCycle() {
try {
await withSuperAdmin(() => storage.processBillingCycle());
log('Billing cycle processed');
} catch (err) {
log(`Billing cycle error: ${err instanceof Error ? err.message : String(err)}`);
}
}
// Run once 30s after startup, then every 24h
setTimeout(async () => {
await runBillingCycle();
setInterval(runBillingCycle, 24 * 60 * 60 * 1000);
}, 30_000);
log('Billing worker started (24h interval)');
// Embedding queue worker — process deferred RAG embeddings nightly at 23:00 MSK
function getMsUntilMskHour(targetHour: number): number {
const now = new Date();
const formatter = new Intl.DateTimeFormat('en-US', {
timeZone: 'Europe/Moscow',
year: 'numeric',
month: 'numeric',
day: 'numeric',
hour: 'numeric',
minute: 'numeric',
second: 'numeric',
hour12: false,
});
const parts = formatter.formatToParts(now);
const getPart = (type: string) => parseInt(parts.find(p => p.type === type)?.value ?? '0', 10);
const year = getPart('year');
const month = getPart('month') - 1;
const day = getPart('day');
const hour = getPart('hour');
const minute = getPart('minute');
const second = getPart('second');
const mskNow = new Date(year, month, day, hour, minute, second).getTime();
let mskTarget = new Date(year, month, day, targetHour, 0, 0).getTime();
if (mskTarget <= mskNow) {
mskTarget += 24 * 60 * 60 * 1000;
}
return mskTarget - mskNow;
}
async function runEmbeddingQueueWorker() {
try {
const { processEmbeddingQueue } = await import('./services/embedding.service');
const pending = await storage.getPendingEmbeddingQueue(10_000);
if (pending.length === 0) {
log('[RAG Queue] No pending embeddings');
return;
}
log(`[RAG Queue] Processing ${pending.length} pending embedding items`);
const result = await processEmbeddingQueue(storage, pending);
const ids = pending.map(i => i.id);
await storage.markEmbeddingQueueProcessed(ids);
log(`[RAG Queue] Done: attempted=${result.attempted}, succeeded=${result.succeeded}, failed=${result.failed}`);
} catch (err) {
log(`[RAG Queue] Worker error: ${err instanceof Error ? err.message : String(err)}`);
}
}
const msUntil23Msk = getMsUntilMskHour(23);
setTimeout(() => {
runEmbeddingQueueWorker();
setInterval(runEmbeddingQueueWorker, 24 * 60 * 60 * 1000);
}, msUntil23Msk);
log(`Embedding queue worker scheduled at 23:00 MSK (in ${Math.round(msUntil23Msk / 1000 / 60)} minutes)`);
// VPN room validity worker — checks Yandex Telemost rooms every 6 hours
async function runVpnRoomValidityWorker() {
try {
const { checkVpnRoomsValidity } = await import('./vpn/vpn.service');
const { notifyVpnRoomExpired } = await import('./vpn/vpn-bot.service');
const expired = await checkVpnRoomsValidity();
if (expired.length === 0) {
log('[VPN] Room validity check done: no expired rooms');
return;
}
log(`[VPN] Room validity check done: ${expired.length} expired room(s)`);
for (const room of expired) {
try {
await notifyVpnRoomExpired(room.organizationId, room.createdBy, room.roomUrl);
} catch (err) {
console.error(`[VPN] Failed to notify user ${room.createdBy} about expired room:`, err);
}
}
} catch (err) {
log(`[VPN] Room validity worker error: ${err instanceof Error ? err.message : String(err)}`);
}
}
// Run once 2 minutes after startup, then every 6 hours
setTimeout(() => {
runVpnRoomValidityWorker();
setInterval(runVpnRoomValidityWorker, 6 * 60 * 60 * 1000);
}, 120_000);
log('VPN room validity worker started (6h interval)');
// ALWAYS serve the app on the port specified in the environment variable PORT
// Other ports are firewalled. Default to 5000 if not specified.
// this serves both the API and the client.
// It is the only port that is not firewalled.
const port = parseInt(process.env.PORT || '5000', 10);
// reusePort is not supported on Windows; skip it for native Windows dev
const isWindows = process.platform === 'win32';
server.listen({
port,
host: "0.0.0.0",
...(isWindows ? {} : { reusePort: true }),
}, () => {
log(`serving on port ${port}`);
});
})();