diff --git a/docker-compose.yml b/docker-compose.yml index 5fe6dd3..7d19ae7 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -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 diff --git a/migrations/0053_embedding_queue.sql b/migrations/0053_embedding_queue.sql new file mode 100644 index 0000000..e8b56b3 --- /dev/null +++ b/migrations/0053_embedding_queue.sql @@ -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 +); diff --git a/server/index.ts b/server/index.ts index dc2f430..5190dd0 100644 --- a/server/index.ts +++ b/server/index.ts @@ -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. diff --git a/server/mcp.ts b/server/mcp.ts index 3860143..2e8c241 100644 --- a/server/mcp.ts +++ b/server/mcp.ts @@ -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", diff --git a/server/routes/chat.messages.routes.ts b/server/routes/chat.messages.routes.ts index e8d619b..2ecf902 100644 --- a/server/routes/chat.messages.routes.ts +++ b/server/routes/chat.messages.routes.ts @@ -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) { diff --git a/server/routes/mcp-rag.routes.ts b/server/routes/mcp-rag.routes.ts index 7b414ad..3dbd1a4 100644 --- a/server/routes/mcp-rag.routes.ts +++ b/server/routes/mcp-rag.routes.ts @@ -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) { diff --git a/server/routes/task-helpers.ts b/server/routes/task-helpers.ts index 9a6d702..0880d38 100644 --- a/server/routes/task-helpers.ts +++ b/server/routes/task-helpers.ts @@ -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 { 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 { 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 { 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); } diff --git a/server/services/embedding.service.ts b/server/services/embedding.service.ts index 0828c7c..40a4977 100644 --- a/server/services/embedding.service.ts +++ b/server/services/embedding.service.ts @@ -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(); + 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>; statuses: Awaited>; fields: Awaited> }>(); + + 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( diff --git a/server/storage.ts b/server/storage.ts index a7d8e72..2cdbd39 100644 --- a/server/storage.ts +++ b/server/storage.ts @@ -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; upsertRagSettings(organizationId: number, updates: Partial>): Promise; + // Embedding Queue + enqueueEmbedding(organizationId: number, entityType: 'form' | 'task' | 'task_message', entityId: number, operation?: 'upsert' | 'delete'): Promise; + getPendingEmbeddingQueue(limit?: number): Promise>; + markEmbeddingQueueProcessed(ids: number[]): Promise; + clearEmbeddingQueue(): Promise; + // LLM Providers getLlmProviders(organizationId: number): Promise; getLlmProvider(id: number, organizationId: number): Promise; diff --git a/server/storage/system.storage.ts b/server/storage/system.storage.ts index 85361b9..9e097c8 100644 --- a/server/storage/system.storage.ts +++ b/server/storage/system.storage.ts @@ -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 { + 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> { + 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 { + if (ids.length === 0) return; + await db + .update(embeddingQueue) + .set({ processedAt: new Date() }) + .where(inArray(embeddingQueue.id, ids)); + } + + async clearEmbeddingQueue(): Promise { + await db.delete(embeddingQueue); + } } diff --git a/shared/schema.ts b/shared/schema.ts index 2bd4e03..f110845 100644 --- a/shared/schema.ts +++ b/shared/schema.ts @@ -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/эмбеддингов на уровне организации) // =====================