RAG: ночная пакетная индексация эмбеддингов, очередь embedding_queue, MCP get_organization_context, снижение OLLAMA_KEEP_ALIVE
This commit is contained in:
@@ -97,6 +97,9 @@ services:
|
||||
- with-ollama
|
||||
ports:
|
||||
- "11434:11434"
|
||||
environment:
|
||||
- OLLAMA_KEEP_ALIVE=${OLLAMA_KEEP_ALIVE:-5m}
|
||||
- OLLAMA_NUM_PARALLEL=1
|
||||
volumes:
|
||||
- ollama_data:/root/.ollama
|
||||
restart: unless-stopped
|
||||
|
||||
29
migrations/0053_embedding_queue.sql
Normal file
29
migrations/0053_embedding_queue.sql
Normal file
@@ -0,0 +1,29 @@
|
||||
-- Queue for deferred RAG embedding generation (nightly batch processing)
|
||||
CREATE TABLE IF NOT EXISTS embedding_queue (
|
||||
id SERIAL PRIMARY KEY,
|
||||
organization_id INTEGER NOT NULL REFERENCES organizations(id) ON DELETE CASCADE,
|
||||
entity_type VARCHAR(20) NOT NULL CHECK (entity_type IN ('form', 'task', 'task_message')),
|
||||
entity_id INTEGER NOT NULL,
|
||||
operation VARCHAR(20) NOT NULL DEFAULT 'upsert' CHECK (operation IN ('upsert', 'delete')),
|
||||
created_at TIMESTAMP NOT NULL DEFAULT NOW(),
|
||||
processed_at TIMESTAMP
|
||||
);
|
||||
|
||||
-- Only one pending operation per entity (deduplication)
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS embedding_queue_pending_entity_idx
|
||||
ON embedding_queue(organization_id, entity_type, entity_id)
|
||||
WHERE processed_at IS NULL;
|
||||
|
||||
CREATE INDEX IF NOT EXISTS embedding_queue_processed_at_idx
|
||||
ON embedding_queue(processed_at) WHERE processed_at IS NULL;
|
||||
|
||||
CREATE INDEX IF NOT EXISTS embedding_queue_created_at_idx
|
||||
ON embedding_queue(created_at);
|
||||
|
||||
-- Enable tenant isolation for embedding queue
|
||||
ALTER TABLE embedding_queue ENABLE ROW LEVEL SECURITY;
|
||||
DROP POLICY IF EXISTS tenant_iso ON embedding_queue;
|
||||
CREATE POLICY tenant_iso ON embedding_queue USING (
|
||||
current_setting('app.is_superadmin', true) = 'true'
|
||||
OR organization_id = NULLIF(current_setting('app.current_org_id', true), '')::int
|
||||
);
|
||||
@@ -1117,6 +1117,60 @@ async function runStartupDataPatches() {
|
||||
}, 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)`);
|
||||
|
||||
// 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.
|
||||
|
||||
@@ -2225,6 +2225,104 @@ To block task creation from task.before_create, set: ctx.result = { allow: false
|
||||
}
|
||||
);
|
||||
|
||||
// get_organization_context
|
||||
server.registerTool(
|
||||
"get_organization_context",
|
||||
{
|
||||
title: "Get Organization Context",
|
||||
description:
|
||||
"Returns a structured snapshot of the organization for AI assistants: list of forms with fields and statuses, " +
|
||||
"task counts by status, recent tasks, and recent chat messages. " +
|
||||
"Use this tool first when the user asks about the CRM structure, workflows, or overall organization state.",
|
||||
inputSchema: {
|
||||
include_recent_tasks: z.boolean().optional().describe("Include recent tasks for each form (default true)"),
|
||||
include_recent_messages: z.boolean().optional().describe("Include recent chat messages (default true)"),
|
||||
recent_tasks_per_form: z.number().int().min(0).max(50).optional().describe("Number of recent tasks per form (default 5)"),
|
||||
recent_messages_limit: z.number().int().min(0).max(100).optional().describe("Total number of recent messages (default 20)"),
|
||||
},
|
||||
},
|
||||
async ({ include_recent_tasks, include_recent_messages, recent_tasks_per_form, recent_messages_limit }) => {
|
||||
try {
|
||||
const forms = await storage.getFormsByOrganization(organizationId);
|
||||
const formContexts = await Promise.all(
|
||||
forms.map(async (form) => {
|
||||
const [fields, statuses, counts] = await Promise.all([
|
||||
storage.getFormFields(form.id, organizationId),
|
||||
storage.getFormStatuses(form.id, organizationId),
|
||||
storage.getTasksCountByForm(form.id),
|
||||
]);
|
||||
const statusCounts = await storage.getTaskCountsByStatus(form.id, organizationId).catch(() => []);
|
||||
const recentTasks = include_recent_tasks !== false
|
||||
? (await storage.getTasksByForm(form.id, organizationId)).slice(0, recent_tasks_per_form ?? 5)
|
||||
: [];
|
||||
return {
|
||||
id: form.id,
|
||||
name: form.name,
|
||||
description: form.description,
|
||||
ai_summary: form.aiSummary,
|
||||
fields: fields.map(f => ({ id: f.id, code: f.code, name: f.name, type: f.type, required: f.isRequired })),
|
||||
statuses: statuses.map(s => ({ id: s.id, name: s.name, color: s.color, is_initial: s.isInitial, is_final: s.isFinal })),
|
||||
task_counts: {
|
||||
total: counts.total,
|
||||
active: counts.active,
|
||||
completed: counts.completed,
|
||||
by_status: statusCounts.map(s => ({ status_id: s.statusId, name: s.name, count: s.count })),
|
||||
},
|
||||
recent_tasks: recentTasks.map(t => ({
|
||||
id: t.id,
|
||||
title: t.title,
|
||||
status_id: t.currentStatusId,
|
||||
assigned_to: t.assignedTo,
|
||||
created_at: t.createdAt,
|
||||
updated_at: t.updatedAt,
|
||||
})),
|
||||
};
|
||||
})
|
||||
);
|
||||
|
||||
let recentMessages: any[] = [];
|
||||
if (include_recent_messages !== false) {
|
||||
const allTasks = (await storage.getTasksByOrganization(organizationId, { limit: 200, minimal: false })) as Task[];
|
||||
const taskTitleMap = new Map(allTasks.map(t => [t.id, t.title]));
|
||||
const messageChunks = await Promise.all(
|
||||
allTasks.slice(0, 50).map(t => storage.getTaskMessages(t.id, organizationId).catch(() => []))
|
||||
);
|
||||
recentMessages = messageChunks
|
||||
.flat()
|
||||
.filter((m: any) => m.messageType === 'comment')
|
||||
.sort((a: any, b: any) => new Date(b.createdAt).getTime() - new Date(a.createdAt).getTime())
|
||||
.slice(0, recent_messages_limit ?? 20)
|
||||
.map((m: any) => ({
|
||||
id: m.id,
|
||||
task_id: m.taskId,
|
||||
task_title: taskTitleMap.get(m.taskId) ?? null,
|
||||
author: m.author ? formatUserName(m.author) : (m.authorId ? `user:${m.authorId}` : 'bot'),
|
||||
message: m.message,
|
||||
created_at: m.createdAt,
|
||||
}));
|
||||
}
|
||||
|
||||
return {
|
||||
content: [{
|
||||
type: "text" as const,
|
||||
text: JSON.stringify({
|
||||
organization_id: organizationId,
|
||||
forms: formContexts,
|
||||
recent_messages: recentMessages,
|
||||
note: "Use semantic_search for detailed natural-language lookups across forms, tasks, and messages.",
|
||||
}, null, 2),
|
||||
}],
|
||||
};
|
||||
} catch (err: unknown) {
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
return {
|
||||
content: [{ type: "text" as const, text: `get_organization_context error: ${msg}` }],
|
||||
isError: true,
|
||||
};
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
// list_field_templates
|
||||
server.registerTool(
|
||||
"list_field_templates",
|
||||
|
||||
@@ -13,7 +13,7 @@ import { webhookService } from "../services/webhook.service";
|
||||
import { generateBotServiceToken } from "../utils/jwt";
|
||||
import { sendWebhook } from "../utils/webhook";
|
||||
import { eventBus } from "./shared";
|
||||
import { buildMessageText, upsertEmbedding } from "../services/embedding.service";
|
||||
|
||||
import { handleAiBotMention, checkBotAccess } from "../services/ai-bot.service";
|
||||
import { db } from "../db";
|
||||
import { fileUploads } from "@shared/schema";
|
||||
@@ -300,22 +300,8 @@ export function registerChatMessageRoutes(router: Router): void {
|
||||
}
|
||||
|
||||
if (messageData.messageType === 'comment') {
|
||||
(async () => {
|
||||
try {
|
||||
const users = await storage.getUsersByOrganization(req.organizationId!);
|
||||
const usersMap = new Map(users.map(u => [u.id, u]));
|
||||
const authorUser = req.user?.id ? usersMap.get(req.user.id) : null;
|
||||
const authorName = authorUser
|
||||
? `${authorUser.firstName ?? ''} ${authorUser.middleName ?? ''} ${authorUser.lastName ?? ''}`.trim() || authorUser.email
|
||||
: 'Пользователь';
|
||||
const content = buildMessageText(createdMessage, authorName, task.title);
|
||||
await upsertEmbedding(req.organizationId!, 'task_message', createdMessage.id, content, {
|
||||
taskId, formId: task.formId, authorId: req.user?.id,
|
||||
});
|
||||
} catch (err) {
|
||||
console.error('[RAG] indexMessage error:', err);
|
||||
}
|
||||
})();
|
||||
storage.enqueueEmbedding(req.organizationId!, 'task_message', createdMessage.id, 'upsert')
|
||||
.catch(err => console.error('[RAG] enqueue message embedding error:', err));
|
||||
}
|
||||
|
||||
if (req.user?.id) {
|
||||
|
||||
@@ -230,13 +230,7 @@ export function registerMcpRagRoutes(app: import("express").Express): void {
|
||||
...(aiSummaryIsAuto !== undefined ? { aiSummaryIsAuto } : {}),
|
||||
});
|
||||
|
||||
const { buildFormText, upsertEmbedding } = await import('../services/embedding.service');
|
||||
const [fields, statuses] = await Promise.all([
|
||||
storage.getFormFields(formId, req.organizationId!),
|
||||
storage.getFormStatuses(formId, req.organizationId!),
|
||||
]);
|
||||
const content = buildFormText(updated, fields, statuses);
|
||||
upsertEmbedding(req.organizationId!, 'form', formId, content, { name: updated.name }).catch(() => {});
|
||||
storage.enqueueEmbedding(req.organizationId!, 'form', formId, 'upsert').catch(() => {});
|
||||
|
||||
res.json({ success: true, form: updated });
|
||||
} catch (error) {
|
||||
|
||||
@@ -1,22 +1,10 @@
|
||||
import type { Request } from 'express';
|
||||
import { storage } from "../storage";
|
||||
import {
|
||||
buildFormText, generateFormSummary, upsertEmbedding, buildTaskText,
|
||||
} from "../services/embedding.service";
|
||||
import { generateFormSummary } from "../services/embedding.service";
|
||||
|
||||
async function indexFormAsync(formId: number, organizationId: number): Promise<void> {
|
||||
try {
|
||||
const [form, fields, statuses] = await Promise.all([
|
||||
storage.getForm(formId, organizationId),
|
||||
storage.getFormFields(formId, organizationId),
|
||||
storage.getFormStatuses(formId, organizationId),
|
||||
]);
|
||||
if (!form) return;
|
||||
const content = buildFormText(form, fields, statuses);
|
||||
await upsertEmbedding(organizationId, 'form', formId, content, {
|
||||
name: form.name,
|
||||
description: form.description,
|
||||
});
|
||||
await storage.enqueueEmbedding(organizationId, 'form', formId, 'upsert');
|
||||
} catch (err) {
|
||||
console.error('[RAG] indexFormAsync error:', err);
|
||||
}
|
||||
@@ -24,22 +12,20 @@ async function indexFormAsync(formId: number, organizationId: number): Promise<v
|
||||
|
||||
async function indexFormWithSummaryAsync(formId: number, organizationId: number): Promise<void> {
|
||||
try {
|
||||
// Summary generation uses external chat provider (Kimi), so it can stay synchronous.
|
||||
const [form, fields, statuses] = await Promise.all([
|
||||
storage.getForm(formId, organizationId),
|
||||
storage.getFormFields(formId, organizationId),
|
||||
storage.getFormStatuses(formId, organizationId),
|
||||
]);
|
||||
if (!form) return;
|
||||
const summary = await generateFormSummary(form, fields, statuses);
|
||||
if (summary) {
|
||||
await storage.updateForm(formId, organizationId, { aiSummary: summary, aiSummaryIsAuto: true });
|
||||
if (form) {
|
||||
const summary = await generateFormSummary(form, fields, statuses);
|
||||
if (summary) {
|
||||
await storage.updateForm(formId, organizationId, { aiSummary: summary, aiSummaryIsAuto: true });
|
||||
}
|
||||
}
|
||||
const indexedForm = summary ? { ...form, aiSummary: summary } : form;
|
||||
const content = buildFormText(indexedForm, fields, statuses);
|
||||
await upsertEmbedding(organizationId, 'form', formId, content, {
|
||||
name: form.name,
|
||||
description: form.description,
|
||||
});
|
||||
// Actual embedding is deferred to the nightly batch.
|
||||
await storage.enqueueEmbedding(organizationId, 'form', formId, 'upsert');
|
||||
} catch (err) {
|
||||
console.error('[RAG] indexFormWithSummaryAsync error:', err);
|
||||
}
|
||||
@@ -47,40 +33,7 @@ async function indexFormWithSummaryAsync(formId: number, organizationId: number)
|
||||
|
||||
async function indexTaskAsync(taskId: number, organizationId: number): Promise<void> {
|
||||
try {
|
||||
const task = await storage.getTask(taskId, organizationId);
|
||||
if (!task) return;
|
||||
const [form, statuses, users, fieldValues, formFields] = await Promise.all([
|
||||
storage.getForm(task.formId, organizationId),
|
||||
storage.getFormStatuses(task.formId, organizationId),
|
||||
storage.getUsersByOrganization(organizationId),
|
||||
storage.getTaskFieldValues(taskId, organizationId),
|
||||
storage.getFormFields(task.formId, organizationId),
|
||||
]);
|
||||
if (!form) return;
|
||||
const statusMap = new Map(statuses.map(s => [s.id, s.name]));
|
||||
const usersMap = new Map(users.map(u => [u.id, u]));
|
||||
const assigneeUser = task.assignedTo ? usersMap.get(task.assignedTo) : null;
|
||||
const assigneeName = assigneeUser
|
||||
? `${assigneeUser.firstName ?? ''} ${assigneeUser.middleName ?? ''} ${assigneeUser.lastName ?? ''}`.trim() || assigneeUser.email
|
||||
: null;
|
||||
const fieldMap = new Map(formFields.map(f => [f.id, f]));
|
||||
const enrichedFieldValues = fieldValues.map(fv => ({
|
||||
name: fieldMap.get(fv.fieldId)?.name ?? '',
|
||||
value: fv.value,
|
||||
type: fieldMap.get(fv.fieldId)?.type ?? 'text',
|
||||
}));
|
||||
const content = buildTaskText(
|
||||
task,
|
||||
form.name,
|
||||
statusMap.get(task.currentStatusId) ?? '',
|
||||
assigneeName,
|
||||
enrichedFieldValues
|
||||
);
|
||||
await upsertEmbedding(organizationId, 'task', taskId, content, {
|
||||
formId: task.formId,
|
||||
formName: form.name,
|
||||
title: task.title,
|
||||
});
|
||||
await storage.enqueueEmbedding(organizationId, 'task', taskId, 'upsert');
|
||||
} catch (err) {
|
||||
console.error('[RAG] indexTaskAsync error:', err);
|
||||
}
|
||||
|
||||
@@ -471,6 +471,166 @@ export async function deleteRagEmbeddingsByForm(
|
||||
`);
|
||||
}
|
||||
|
||||
interface QueueItem {
|
||||
id: number;
|
||||
organizationId: number;
|
||||
entityType: EntityType;
|
||||
entityId: number;
|
||||
operation: 'upsert' | 'delete';
|
||||
}
|
||||
|
||||
export async function processEmbeddingQueue(
|
||||
stor: IStorage,
|
||||
items: QueueItem[],
|
||||
onProgress?: (indexed: number, total: number) => void
|
||||
): Promise<{ attempted: number; succeeded: number; failed: number }> {
|
||||
if (items.length === 0) return { attempted: 0, succeeded: 0, failed: 0 };
|
||||
|
||||
// Group by organization to reuse config and users cache
|
||||
const byOrg = new Map<number, QueueItem[]>();
|
||||
for (const item of items) {
|
||||
const list = byOrg.get(item.organizationId) || [];
|
||||
list.push(item);
|
||||
byOrg.set(item.organizationId, list);
|
||||
}
|
||||
|
||||
let attempted = 0;
|
||||
let succeeded = 0;
|
||||
let failed = 0;
|
||||
|
||||
for (const [organizationId, orgItems] of byOrg.entries()) {
|
||||
const config = await resolveOrgConfig(organizationId);
|
||||
const hasAccess = config.provider === 'ollama' || !!config.apiKey;
|
||||
if (!hasAccess) {
|
||||
console.warn(`[RAG] Skipping queue for org ${organizationId}: no API key configured`);
|
||||
failed += orgItems.length;
|
||||
continue;
|
||||
}
|
||||
|
||||
const users = await stor.getUsersByOrganization(organizationId);
|
||||
const usersMap = new Map(users.map(u => [u.id, u]));
|
||||
const formCache = new Map<number, { form: Awaited<ReturnType<typeof stor.getForm>>; statuses: Awaited<ReturnType<typeof stor.getFormStatuses>>; fields: Awaited<ReturnType<typeof stor.getFormFields>> }>();
|
||||
|
||||
async function getFormMeta(formId: number) {
|
||||
if (formCache.has(formId)) return formCache.get(formId)!;
|
||||
const [form, statuses, formFields] = await Promise.all([
|
||||
stor.getForm(formId, organizationId),
|
||||
stor.getFormStatuses(formId, organizationId),
|
||||
stor.getFormFields(formId, organizationId),
|
||||
]);
|
||||
const meta = { form, statuses, fields: formFields };
|
||||
formCache.set(formId, meta);
|
||||
return meta;
|
||||
}
|
||||
|
||||
const embeddingItems: EmbeddingItem[] = [];
|
||||
|
||||
for (const item of orgItems) {
|
||||
attempted++;
|
||||
try {
|
||||
if (item.operation === 'delete') {
|
||||
await deleteRagEmbedding(organizationId, item.entityType, item.entityId);
|
||||
succeeded++;
|
||||
if (onProgress) onProgress(succeeded + failed, attempted);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (item.entityType === 'form') {
|
||||
const form = await stor.getForm(item.entityId, organizationId);
|
||||
if (!form) {
|
||||
failed++;
|
||||
continue;
|
||||
}
|
||||
const [fields, statuses] = await Promise.all([
|
||||
stor.getFormFields(item.entityId, organizationId),
|
||||
stor.getFormStatuses(item.entityId, organizationId),
|
||||
]);
|
||||
const content = buildFormText(form, fields, statuses);
|
||||
embeddingItems.push({
|
||||
organizationId,
|
||||
entityType: 'form',
|
||||
entityId: item.entityId,
|
||||
content,
|
||||
metadata: { name: form.name },
|
||||
});
|
||||
} else if (item.entityType === 'task') {
|
||||
const task = await stor.getTask(item.entityId, organizationId);
|
||||
if (!task) {
|
||||
failed++;
|
||||
continue;
|
||||
}
|
||||
const { form, statuses, fields: formFields } = await getFormMeta(task.formId);
|
||||
if (!form) {
|
||||
failed++;
|
||||
continue;
|
||||
}
|
||||
const statusMap = new Map(statuses.map(s => [s.id, s.name]));
|
||||
const assigneeUser = task.assignedTo ? usersMap.get(task.assignedTo) : null;
|
||||
const assigneeName = assigneeUser
|
||||
? `${assigneeUser.firstName ?? ''} ${assigneeUser.middleName ?? ''} ${assigneeUser.lastName ?? ''}`.trim() || assigneeUser.email
|
||||
: null;
|
||||
const fieldMap = new Map(formFields.map(f => [f.id, f]));
|
||||
const fieldValues = await stor.getTaskFieldValues(task.id, organizationId);
|
||||
const enrichedFieldValues = fieldValues.map(fv => ({
|
||||
name: fieldMap.get(fv.fieldId)?.name ?? '',
|
||||
value: fv.value,
|
||||
type: fieldMap.get(fv.fieldId)?.type ?? 'text',
|
||||
}));
|
||||
const content = buildTaskText(
|
||||
task,
|
||||
form.name,
|
||||
statusMap.get(task.currentStatusId) ?? '',
|
||||
assigneeName,
|
||||
enrichedFieldValues
|
||||
);
|
||||
embeddingItems.push({
|
||||
organizationId,
|
||||
entityType: 'task',
|
||||
entityId: task.id,
|
||||
content,
|
||||
metadata: { formId: task.formId, formName: form.name, title: task.title },
|
||||
});
|
||||
} else if (item.entityType === 'task_message') {
|
||||
const message = await stor.getTaskMessage(item.entityId, organizationId);
|
||||
if (!message) {
|
||||
failed++;
|
||||
continue;
|
||||
}
|
||||
const task = await stor.getTask(message.taskId, organizationId);
|
||||
if (!task) {
|
||||
failed++;
|
||||
continue;
|
||||
}
|
||||
const authorUser = message.authorId ? usersMap.get(message.authorId) : null;
|
||||
const authorName = authorUser
|
||||
? `${authorUser.firstName ?? ''} ${authorUser.middleName ?? ''} ${authorUser.lastName ?? ''}`.trim() || authorUser.email
|
||||
: 'Бот';
|
||||
const content = buildMessageText(message, authorName, task.title);
|
||||
embeddingItems.push({
|
||||
organizationId,
|
||||
entityType: 'task_message',
|
||||
entityId: message.id,
|
||||
content,
|
||||
metadata: { taskId: task.id, formId: task.formId, authorId: message.authorId },
|
||||
});
|
||||
}
|
||||
} catch (err) {
|
||||
console.error(`[RAG] Failed to prepare queue item ${item.id}:`, err);
|
||||
failed++;
|
||||
if (onProgress) onProgress(succeeded + failed, attempted);
|
||||
}
|
||||
}
|
||||
|
||||
const batchSucceeded = await processBatchedItems(embeddingItems, config, (batchIndexed, total) => {
|
||||
if (onProgress) onProgress(succeeded + failed + batchIndexed, attempted);
|
||||
});
|
||||
succeeded += batchSucceeded;
|
||||
failed += embeddingItems.length - batchSucceeded;
|
||||
}
|
||||
|
||||
return { attempted, succeeded, failed };
|
||||
}
|
||||
|
||||
const TASK_CHUNK_SIZE = 100;
|
||||
|
||||
export async function reindexOrganization(
|
||||
|
||||
@@ -6,7 +6,7 @@ import { organizationApiKeys, type OrganizationApiKey } from "@shared/schema";
|
||||
import { automations, type Automation, type InsertAutomation } from "@shared/schema";
|
||||
import { taskRelations, type TaskRelation, type InsertTaskRelation } from "@shared/schema";
|
||||
import { taskReminders, type TaskReminder, type InsertTaskReminder } from "@shared/schema";
|
||||
import { ragSettings, type RagSetting } from "@shared/schema";
|
||||
import { ragSettings, type RagSetting, embeddingQueue } from "@shared/schema";
|
||||
import { llmProviders, type LlmProvider, type InsertLlmProvider } from "@shared/schema";
|
||||
import { systemConfig } from "@shared/schema";
|
||||
import { taskAuditLog, type TaskAuditLog, type InsertTaskAuditLog } from "@shared/schema";
|
||||
@@ -458,6 +458,12 @@ export interface IStorage {
|
||||
getRagSettings(organizationId: number): Promise<RagSetting | undefined>;
|
||||
upsertRagSettings(organizationId: number, updates: Partial<Omit<RagSetting, 'organizationId' | 'updatedAt'>>): Promise<RagSetting>;
|
||||
|
||||
// Embedding Queue
|
||||
enqueueEmbedding(organizationId: number, entityType: 'form' | 'task' | 'task_message', entityId: number, operation?: 'upsert' | 'delete'): Promise<void>;
|
||||
getPendingEmbeddingQueue(limit?: number): Promise<Array<{ id: number; organizationId: number; entityType: 'form' | 'task' | 'task_message'; entityId: number; operation: 'upsert' | 'delete' }>>;
|
||||
markEmbeddingQueueProcessed(ids: number[]): Promise<void>;
|
||||
clearEmbeddingQueue(): Promise<void>;
|
||||
|
||||
// LLM Providers
|
||||
getLlmProviders(organizationId: number): Promise<LlmProvider[]>;
|
||||
getLlmProvider(id: number, organizationId: number): Promise<LlmProvider | undefined>;
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { users, forms } from "@shared/schema";
|
||||
import { ragSettings, type RagSetting } from "@shared/schema";
|
||||
import { embeddingQueue } from "@shared/schema";
|
||||
import { llmProviders, type LlmProvider, type InsertLlmProvider } from "@shared/schema";
|
||||
import { taskRoles, type TaskRole, taskUserAccess, type TaskUserAccess, roleMembers, roles, taskAssignees, type TaskAssignee, formAccessRules, type FormAccessRule, type InsertFormAccessRule } from "@shared/schema";
|
||||
import { db } from "../db";
|
||||
@@ -439,4 +440,56 @@ export class SystemStorage extends RolesStorage {
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
// =====================
|
||||
// Embedding Queue
|
||||
// =====================
|
||||
|
||||
async enqueueEmbedding(
|
||||
organizationId: number,
|
||||
entityType: 'form' | 'task' | 'task_message',
|
||||
entityId: number,
|
||||
operation: 'upsert' | 'delete' = 'upsert'
|
||||
): Promise<void> {
|
||||
await db
|
||||
.insert(embeddingQueue)
|
||||
.values({
|
||||
organizationId,
|
||||
entityType,
|
||||
entityId,
|
||||
operation,
|
||||
createdAt: new Date(),
|
||||
})
|
||||
.onConflictDoNothing({ target: [embeddingQueue.organizationId, embeddingQueue.entityType, embeddingQueue.entityId] });
|
||||
}
|
||||
|
||||
async getPendingEmbeddingQueue(
|
||||
limit: number = 1000
|
||||
): Promise<Array<{ id: number; organizationId: number; entityType: 'form' | 'task' | 'task_message'; entityId: number; operation: 'upsert' | 'delete' }>> {
|
||||
const rows = await db
|
||||
.select({
|
||||
id: embeddingQueue.id,
|
||||
organizationId: embeddingQueue.organizationId,
|
||||
entityType: embeddingQueue.entityType,
|
||||
entityId: embeddingQueue.entityId,
|
||||
operation: embeddingQueue.operation,
|
||||
})
|
||||
.from(embeddingQueue)
|
||||
.where(isNull(embeddingQueue.processedAt))
|
||||
.orderBy(asc(embeddingQueue.createdAt))
|
||||
.limit(limit);
|
||||
return rows as Array<{ id: number; organizationId: number; entityType: 'form' | 'task' | 'task_message'; entityId: number; operation: 'upsert' | 'delete' }>;
|
||||
}
|
||||
|
||||
async markEmbeddingQueueProcessed(ids: number[]): Promise<void> {
|
||||
if (ids.length === 0) return;
|
||||
await db
|
||||
.update(embeddingQueue)
|
||||
.set({ processedAt: new Date() })
|
||||
.where(inArray(embeddingQueue.id, ids));
|
||||
}
|
||||
|
||||
async clearEmbeddingQueue(): Promise<void> {
|
||||
await db.delete(embeddingQueue);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2995,6 +2995,27 @@ export const ragEmbeddings = pgTable("rag_embeddings", {
|
||||
export type RagEmbedding = typeof ragEmbeddings.$inferSelect;
|
||||
export type InsertRagEmbedding = typeof ragEmbeddings.$inferInsert;
|
||||
|
||||
// =====================
|
||||
// Embedding Queue (отложенная генерация эмбеддингов для RAG)
|
||||
// =====================
|
||||
export const embeddingQueue = pgTable("embedding_queue", {
|
||||
id: serial("id").primaryKey(),
|
||||
organizationId: integer("organization_id").notNull().references(() => organizations.id, { onDelete: "cascade" }),
|
||||
entityType: varchar("entity_type", { length: 20 }).notNull(), // 'form' | 'task' | 'task_message'
|
||||
entityId: integer("entity_id").notNull(),
|
||||
operation: varchar("operation", { length: 20 }).notNull().default("upsert"), // 'upsert' | 'delete'
|
||||
createdAt: timestamp("created_at").defaultNow(),
|
||||
processedAt: timestamp("processed_at"),
|
||||
}, (table) => ({
|
||||
// Only one pending operation per entity
|
||||
pendingEntityIdx: uniqueIndex("embedding_queue_pending_entity_idx").on(table.organizationId, table.entityType, table.entityId).where(sql`${table.processedAt} IS NULL`),
|
||||
pendingIdx: index("embedding_queue_processed_at_idx").on(table.processedAt).where(sql`${table.processedAt} IS NULL`),
|
||||
createdAtIdx: index("embedding_queue_created_at_idx").on(table.createdAt),
|
||||
}));
|
||||
|
||||
export type EmbeddingQueue = typeof embeddingQueue.$inferSelect;
|
||||
export type InsertEmbeddingQueue = typeof embeddingQueue.$inferInsert;
|
||||
|
||||
// =====================
|
||||
// RAG Settings (настройки LLM/эмбеддингов на уровне организации)
|
||||
// =====================
|
||||
|
||||
Reference in New Issue
Block a user