Files
iistwin/server/routes/mcp-rag.routes.ts

368 lines
17 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 { storage } from "../storage";
import { authenticateToken, requirePermission, type AuthenticatedRequest } from "../middleware/auth.middleware";
import { tenantIsolation } from "../middleware/tenant.middleware";
import { type ReminderRecipient } from "@shared/schema";
import { handleMcpRequest, handleMcpSse, handleMcpMessages } from "../mcp";
import { setupSwagger } from "../swagger";
import express, { type Request, type Response, type NextFunction } from 'express';
// In-process dedup: avoid spamming admins with repeated legacy-key notifications.
// Key: "orgId:keyPrefix" — reset only on process restart.
const _legacyKeyNotifiedSet = new Set<string>();
async function notifyAdminsLegacyKey(organizationId: number, keyPrefix: string, label: string): Promise<void> {
const dedupeKey = `${organizationId}:${keyPrefix}`;
if (_legacyKeyNotifiedSet.has(dedupeKey)) return;
_legacyKeyNotifiedSet.add(dedupeKey);
try {
const admins = await storage.getUsersByRole('admin', organizationId);
for (const admin of admins) {
await storage.createUserNotification({
userId: admin.id,
organizationId,
type: 'system',
title: 'Устаревший API-ключ',
message: `API-ключ "${label}" (${keyPrefix}…) использует устаревший алгоритм (SHA-256). Удалите его и создайте новый в Настройки → API-ключи.`,
isRead: false,
});
}
} catch {
_legacyKeyNotifiedSet.delete(dedupeKey);
}
}
async function requireMcpApiKey(req: Request, res: Response, next: NextFunction): Promise<void> {
const rawKey =
(req.headers['x-api-key'] as string | undefined) ||
(req.headers['authorization'] as string | undefined)?.replace(/^Bearer\s+/i, '') ||
(req.query?.['apiKey'] as string | undefined) ||
(req.query?.['key'] as string | undefined);
if (!rawKey) {
res.status(401).json({ error: 'Missing API key. Provide Authorization: Bearer <key>, X-Api-Key header, or ?key=<key> query parameter.' });
return;
}
const trimmed = rawKey.trim();
const apiKey = await storage.getApiKeyByHash(trimmed);
if (apiKey && apiKey.isActive) {
storage.touchApiKey(apiKey.id).catch(() => {});
next();
return;
}
const legacyKey = await storage.getApiKeyByLegacyHash(trimmed);
if (legacyKey && legacyKey.isActive) {
notifyAdminsLegacyKey(legacyKey.organizationId, legacyKey.keyPrefix, legacyKey.label).catch(() => {});
res.status(401).json({
error: 'API key is outdated (SHA-256). Please delete it and generate a new one in Settings → API Keys.',
});
return;
}
res.status(401).json({ error: 'Invalid or revoked API key.' });
}
export function registerMcpRagRoutes(app: import("express").Express): void {
/**
* @swagger
* /api/mcp-keys:
* get:
* tags: [MCP]
* summary: Список API-ключей для MCP (admin)
* description: API-ключи для подключения AI-агентов (Claude Desktop, Cursor, Cline). Требуется роль admin.
* responses:
* 200:
* description: Список ключей (без keyHash)
* post:
* tags: [MCP]
* summary: Создать API-ключ (admin)
* description: Ключ показывается **только один раз** в поле `key`. Формат — `wf_<random>`.
* /api/mcp-keys/{id}:
* delete:
* tags: [MCP]
* summary: Удалить API-ключ (admin)
*/
app.get('/api/mcp-keys', authenticateToken, tenantIsolation, requirePermission('settings.manage'), async (req: AuthenticatedRequest, res) => {
try {
const keys = await storage.listApiKeys(req.organizationId!);
const safeKeys = keys.map(({ keyHash: _h, ...rest }) => rest);
res.json({ success: true, keys: safeKeys });
} catch (error) {
console.error('List MCP keys error:', error);
res.status(500).json({ error: 'Ошибка получения ключей' });
}
});
app.post('/api/mcp-keys', authenticateToken, tenantIsolation, requirePermission('settings.manage'), async (req: AuthenticatedRequest, res) => {
try {
const label = (req.body?.label as string | undefined)?.trim() || 'Default';
const { key, record } = await storage.createApiKey(req.organizationId!, req.user!.id, label);
const { keyHash: _h, ...safeRecord } = record;
res.json({ success: true, key, record: safeRecord });
} catch (error) {
console.error('Create MCP key error:', error);
res.status(500).json({ error: 'Ошибка создания ключа' });
}
});
app.delete('/api/mcp-keys/:id', authenticateToken, tenantIsolation, requirePermission('settings.manage'), async (req: AuthenticatedRequest, res) => {
try {
const id = parseInt(req.params.id);
if (isNaN(id)) return res.status(400).json({ error: 'Неверный ID' });
await storage.deleteApiKey(id, req.organizationId!);
res.json({ success: true });
} catch (error) {
console.error('Delete MCP key error:', error);
res.status(500).json({ error: 'Ошибка удаления ключа' });
}
});
// RAG / Semantic Search Endpoints
app.get('/api/mcp/rag-stats', authenticateToken, tenantIsolation, requirePermission('settings.manage'), async (req: AuthenticatedRequest, res) => {
try {
const { getEmbeddingCounts } = await import('../services/embedding.service');
const rawCounts = await getEmbeddingCounts(req.organizationId!);
const counts = {
form: rawCounts['form'] ?? 0,
task: rawCounts['task'] ?? 0,
task_message: rawCounts['task_message'] ?? 0,
};
res.json({ success: true, counts });
} catch (error) {
console.error('RAG stats error:', error);
res.status(500).json({ success: false, error: 'Ошибка получения статистики индекса' });
}
});
interface ReindexProgress {
indexed: number;
total: number;
attempted: number;
failed: number;
status: 'running' | 'done' | 'error';
startedAt: number;
result?: { forms: number; tasks: number; messages: number; attempted: number; failed: number };
error?: string;
}
const reindexProgressStore = new Map<number, ReindexProgress>();
app.get('/api/mcp/reindex-progress', authenticateToken, tenantIsolation, requirePermission('settings.manage'), (req: AuthenticatedRequest, res) => {
const organizationId = req.organizationId!;
const progress = reindexProgressStore.get(organizationId);
if (!progress) {
return res.json({ success: true, status: 'idle' });
}
res.json({ success: true, ...progress });
});
app.post('/api/mcp/reindex', authenticateToken, tenantIsolation, requirePermission('settings.manage'), async (req: AuthenticatedRequest, res) => {
try {
const organizationId = req.organizationId!;
const existing = reindexProgressStore.get(organizationId);
if (existing && existing.status === 'running') {
return res.json({ success: true, message: 'Переиндексация уже выполняется', alreadyRunning: true });
}
const rawTypes: string[] = req.body?.entity_types ?? ['form', 'task', 'task_message'];
const validTypes = ['form', 'task', 'task_message'];
const entityTypes = rawTypes.filter(t => validTypes.includes(t)) as import('../services/embedding.service').EntityType[];
const { reindexOrganization } = await import('../services/embedding.service');
reindexProgressStore.set(organizationId, { indexed: 0, total: 0, attempted: 0, failed: 0, status: 'running', startedAt: Date.now() });
reindexOrganization(storage, organizationId, entityTypes, (indexed, total) => {
const prev = reindexProgressStore.get(organizationId);
reindexProgressStore.set(organizationId, {
indexed,
total,
attempted: prev?.attempted ?? 0,
failed: prev?.failed ?? 0,
status: 'running',
startedAt: prev?.startedAt ?? Date.now(),
});
}).then((result) => {
const finalTotal = result.attempted;
reindexProgressStore.set(organizationId, {
indexed: result.forms + result.tasks + result.messages,
total: finalTotal,
attempted: result.attempted,
failed: result.failed,
status: 'done',
startedAt: reindexProgressStore.get(organizationId)?.startedAt ?? Date.now(),
result,
});
setTimeout(() => reindexProgressStore.delete(organizationId), 5 * 60 * 1000);
}).catch((err) => {
console.error('[RAG] Reindex error:', err);
const prev = reindexProgressStore.get(organizationId);
reindexProgressStore.set(organizationId, {
indexed: prev?.indexed ?? 0,
total: prev?.total ?? 0,
attempted: prev?.attempted ?? 0,
failed: prev?.failed ?? 0,
status: 'error',
startedAt: prev?.startedAt ?? Date.now(),
error: String(err),
});
setTimeout(() => reindexProgressStore.delete(organizationId), 5 * 60 * 1000);
});
res.json({ success: true, message: 'Переиндексация запущена в фоне' });
} catch (error) {
console.error('Reindex error:', error);
res.status(500).json({ success: false, error: 'Ошибка запуска переиндексации' });
}
});
app.patch('/api/forms/:id/ai-summary', authenticateToken, tenantIsolation, requirePermission('settings.manage'), async (req: AuthenticatedRequest, res) => {
try {
const formId = parseInt(req.params.id);
if (isNaN(formId)) return res.status(400).json({ success: false, error: 'Неверный ID формы' });
const form = await storage.getForm(formId, req.organizationId!);
if (!form) return res.status(404).json({ success: false, error: 'Форма не найдена' });
const { aiSummary, aiSummaryIsAuto } = req.body;
if (aiSummaryIsAuto !== undefined && typeof aiSummaryIsAuto !== 'boolean') {
return res.status(400).json({ success: false, error: 'aiSummaryIsAuto должен быть булевым значением' });
}
const updated = await storage.updateForm(formId, req.organizationId!, {
aiSummary,
...(aiSummaryIsAuto !== undefined ? { aiSummaryIsAuto } : {}),
});
storage.enqueueEmbedding(req.organizationId!, 'form', formId, 'upsert').catch(() => {});
res.json({ success: true, form: updated });
} catch (error) {
console.error('Update AI summary error:', error);
res.status(500).json({ success: false, error: 'Ошибка обновления резюме' });
}
});
app.post('/api/forms/:id/regenerate-summary', authenticateToken, tenantIsolation, requirePermission('settings.manage'), async (req: AuthenticatedRequest, res) => {
try {
const formId = parseInt(req.params.id);
if (isNaN(formId)) return res.status(400).json({ success: false, error: 'Неверный ID формы' });
const form = await storage.getForm(formId, req.organizationId!);
if (!form) return res.status(404).json({ success: false, error: 'Форма не найдена' });
const { generateFormSummary } = await import('../services/embedding.service');
const [fields, statuses] = await Promise.all([
storage.getFormFields(formId, req.organizationId!),
storage.getFormStatuses(formId, req.organizationId!),
]);
const summary = await generateFormSummary(form, fields, statuses, req.organizationId!);
res.json({ success: true, aiSummary: summary ?? null });
} catch (error) {
console.error('Regenerate summary error:', error);
res.status(500).json({ success: false, error: 'Ошибка генерации резюме' });
}
});
// Task Reminders Endpoints
app.get('/api/tasks/:id/reminders', authenticateToken, tenantIsolation, async (req: AuthenticatedRequest, res) => {
const taskId = parseInt(req.params.id);
if (isNaN(taskId)) return res.status(400).json({ success: false, error: 'Неверный ID задачи' });
try {
const reminders = await storage.getTaskReminders(taskId, req.organizationId!);
return res.json({ success: true, reminders });
} catch (err: any) {
return res.status(500).json({ success: false, error: err.message });
}
});
app.post('/api/tasks/:id/reminders', authenticateToken, tenantIsolation, async (req: AuthenticatedRequest, res) => {
const taskId = parseInt(req.params.id);
if (isNaN(taskId)) return res.status(400).json({ success: false, error: 'Неверный ID задачи' });
const task = await storage.getTask(taskId, req.organizationId!);
if (!task) return res.status(404).json({ success: false, error: 'Задача не найдена' });
const { remindAt, note, recipients } = req.body;
if (!remindAt) return res.status(400).json({ success: false, error: 'Укажите дату и время напоминания' });
const remindAtDate = new Date(remindAt);
if (isNaN(remindAtDate.getTime())) {
return res.status(400).json({ success: false, error: 'Неверный формат даты напоминания' });
}
if (!Array.isArray(recipients) || recipients.length === 0) {
return res.status(400).json({ success: false, error: 'Укажите хотя бы одного получателя' });
}
const ALLOWED_SYSTEM_ROLES = new Set(['admin', 'user']);
const validatedRecipients: ReminderRecipient[] = [];
const hasUserRecipients = recipients.some((rec: any) => rec.type === 'user');
const orgUsersMap = hasUserRecipients
? new Map((await storage.getUsersByOrganization(req.organizationId!)).map(u => [u.id, u]))
: new Map();
for (const rec of recipients) {
if (rec.type === 'user') {
if (typeof rec.userId !== 'number') return res.status(400).json({ success: false, error: 'Неверный формат получателя' });
const orgUser = orgUsersMap.get(rec.userId);
if (!orgUser) {
return res.status(400).json({ success: false, error: `Пользователь ${rec.userId} не принадлежит организации` });
}
validatedRecipients.push({ type: 'user', userId: rec.userId });
} else if (rec.type === 'role') {
if (typeof rec.role !== 'string' || !ALLOWED_SYSTEM_ROLES.has(rec.role)) {
return res.status(400).json({ success: false, error: `Недопустимая системная роль: ${rec.role}` });
}
validatedRecipients.push({ type: 'role', role: rec.role });
} else if (rec.type === 'orgRole') {
if (typeof rec.roleId !== 'number') return res.status(400).json({ success: false, error: 'Неверный формат роли организации' });
await storage.getUsersByOrgRoleId(rec.roleId, req.organizationId!);
validatedRecipients.push({ type: 'orgRole', roleId: rec.roleId });
} else {
return res.status(400).json({ success: false, error: 'Неверный тип получателя' });
}
}
try {
const reminder = await storage.createTaskReminder({
taskId,
organizationId: req.organizationId!,
createdByUserId: req.user!.id,
remindAt: remindAtDate,
note: note || null,
recipients: validatedRecipients,
});
return res.json({ success: true, reminder });
} catch (err: any) {
return res.status(500).json({ success: false, error: err.message });
}
});
app.delete('/api/reminders/:id', authenticateToken, tenantIsolation, async (req: AuthenticatedRequest, res) => {
const id = parseInt(req.params.id);
if (isNaN(id)) return res.status(400).json({ success: false, error: 'Неверный ID' });
try {
await storage.deleteTaskReminder(id, req.organizationId!);
return res.json({ success: true });
} catch (err: any) {
return res.status(500).json({ success: false, error: err.message });
}
});
// MCP Server Endpoints — require valid org API key at route level
// Handle notifications before handleMcpRequest to avoid transport issues with consumed body streams.
// NOTE: we rely on the global express.json() middleware (already applied in server/index.ts).
app.all('/mcp', (req, res, next) => {
if (
req.method === 'POST' &&
req.body &&
!Array.isArray(req.body) &&
typeof (req.body as any).method === 'string' &&
(req.body as any).method.startsWith('notifications/')
) {
return res.status(202).send();
}
next();
}, handleMcpRequest);
app.get('/mcp/sse', requireMcpApiKey, handleMcpSse);
app.post('/mcp/messages', requireMcpApiKey, express.json(), handleMcpMessages);
// Swagger UI
setupSwagger(app);
}