import { db } from "../db"; import { sql } from "drizzle-orm"; import { storage } from "../storage"; import { decrypt } from "../crypto"; import type { Form, FormField, FormStatus, Task, TaskMessage } from "@shared/schema"; import type { IStorage } from "../storage"; import { generateEmbedding, generateEmbeddingBatch, generateSummary, type LlmProviderConfig, } from "./llm-provider"; const GLOBAL_OPENAI_API_KEY = process.env.OPENAI_API_KEY; const DEFAULT_EMBEDDING_MODEL = "text-embedding-3-small"; const BATCH_SIZE = 20; export type EntityType = "form" | "task" | "task_message"; export interface RagResult { entityType: EntityType; entityId: number; score: number; content: string; metadata: Record; } interface EmbeddingItem { organizationId: number; entityType: EntityType; entityId: number; content: string; metadata: Record; } export interface OrgEmbeddingConfig { provider: string; apiKey: string | null; baseUrl: string | null; embeddingModel: string; customHeaders: Record; } async function resolveOrgConfig(organizationId: number): Promise { const settings = await storage.getRagSettings(organizationId); if (settings) { // Если задан embeddingProviderId — загружаем провайдера из реестра if (settings.embeddingProviderId) { try { const provider = await storage.getLlmProvider(settings.embeddingProviderId, organizationId); if (provider && provider.isActive) { const apiKey = provider.apiKey ? decrypt(provider.apiKey) : null; const resolvedKey = apiKey || (provider.providerType === 'openai' ? GLOBAL_OPENAI_API_KEY || null : null); const embeddingModel = settings.embeddingModel || (provider.enabledModels && provider.enabledModels.length > 0 ? provider.enabledModels[0] : null) || DEFAULT_EMBEDDING_MODEL; return { provider: provider.providerType, apiKey: resolvedKey, baseUrl: provider.baseUrl ?? null, embeddingModel, chatModel: settings.chatModel ?? null, maxChunkChars: settings.maxChunkChars ?? null, summarizationEnabled: settings.summarizationEnabled ?? true, customHeaders: (provider.customHeaders ?? {}) as Record, }; } } catch (err) { console.warn("[RAG] resolveOrgConfig: failed to load embeddingProvider", settings.embeddingProviderId, err); } } const apiKey = settings.apiKey ? decrypt(settings.apiKey) : null; return { provider: settings.provider, apiKey: apiKey || GLOBAL_OPENAI_API_KEY || null, baseUrl: settings.baseUrl, embeddingModel: settings.embeddingModel || DEFAULT_EMBEDDING_MODEL, chatModel: settings.chatModel ?? null, maxChunkChars: settings.maxChunkChars ?? null, summarizationEnabled: settings.summarizationEnabled ?? true, customHeaders: {}, }; } return { provider: "openai", apiKey: GLOBAL_OPENAI_API_KEY || null, baseUrl: null, embeddingModel: DEFAULT_EMBEDDING_MODEL, chatModel: null, maxChunkChars: null, summarizationEnabled: true, customHeaders: {}, }; } async function doEmbedding(text: string, config: LlmProviderConfig): Promise { return generateEmbedding(text, config); } async function doEmbeddingBatch(texts: string[], config: LlmProviderConfig): Promise<(number[] | null)[]> { return generateEmbeddingBatch(texts, config); } export async function resolveChatConfig(organizationId: number): Promise { const settings = await storage.getRagSettings(organizationId); if (settings) { // chatProviderId — отдельный провайдер для суммаризации; // если не задан — наследуем embeddingProviderId (UI: "Тот же провайдер") const effectiveChatProviderId = settings.chatProviderId ?? settings.embeddingProviderId; if (effectiveChatProviderId) { try { const provider = await storage.getLlmProvider(effectiveChatProviderId, organizationId); if (provider && provider.isActive) { const apiKey = provider.apiKey ? decrypt(provider.apiKey) : null; const resolvedKey = apiKey || (provider.providerType === "openai" ? GLOBAL_OPENAI_API_KEY || null : null); const chatModel = settings.chatModel ?? (provider.enabledModels && provider.enabledModels.length > 0 ? provider.enabledModels[0] : null) ?? (provider.providerType === "ollama" ? "llama3" : "gpt-4o-mini"); return { provider: provider.providerType, apiKey: resolvedKey, baseUrl: provider.baseUrl ?? null, embeddingModel: DEFAULT_EMBEDDING_MODEL, chatModel, maxChunkChars: settings.maxChunkChars ?? null, summarizationEnabled: settings.summarizationEnabled ?? true, customHeaders: (provider.customHeaders ?? {}) as Record, }; } } catch (err) { console.warn("[RAG] resolveChatConfig: failed to load chatProvider", effectiveChatProviderId, err); } } // Fallback: legacy rag_settings const apiKey = settings.apiKey ? decrypt(settings.apiKey) : null; return { provider: settings.provider, apiKey: apiKey || GLOBAL_OPENAI_API_KEY || null, baseUrl: settings.baseUrl, embeddingModel: settings.embeddingModel || DEFAULT_EMBEDDING_MODEL, chatModel: settings.chatModel ?? null, maxChunkChars: settings.maxChunkChars ?? null, summarizationEnabled: settings.summarizationEnabled ?? true, customHeaders: {}, }; } return { provider: "openai", apiKey: GLOBAL_OPENAI_API_KEY || null, baseUrl: null, embeddingModel: DEFAULT_EMBEDDING_MODEL, chatModel: null, maxChunkChars: null, summarizationEnabled: true, customHeaders: {}, }; } export async function embedText(text: string, organizationId?: number): Promise { if (organizationId !== undefined) { const config = await resolveOrgConfig(organizationId); return doEmbedding(text, config); } const config: LlmProviderConfig = { provider: "openai", apiKey: GLOBAL_OPENAI_API_KEY || null, baseUrl: null, embeddingModel: DEFAULT_EMBEDDING_MODEL, chatModel: null, maxChunkChars: null, summarizationEnabled: true, customHeaders: {}, }; return doEmbedding(text, config); } export async function generateFormSummary(form: Form, fields: FormField[], statuses: FormStatus[], organizationId?: number): Promise { // Use chat config (chatProviderId) for summarization, not embedding config const config: LlmProviderConfig = organizationId ? await resolveChatConfig(organizationId) : { provider: "openai", apiKey: GLOBAL_OPENAI_API_KEY || null, baseUrl: null, embeddingModel: DEFAULT_EMBEDDING_MODEL, chatModel: null, maxChunkChars: null, summarizationEnabled: true, customHeaders: {}, }; // Respect the summarization toggle if (!config.summarizationEnabled) return null; // For non-Ollama providers we need an API key if (config.provider !== "ollama" && !config.apiKey) return null; const fieldList = fields.map(f => `${f.name} (${f.type})`).join(", "); const statusList = statuses.map(s => `${s.name}${s.isInitial ? " [начальный]" : ""}${s.isFinal ? " [финальный]" : ""}`).join(", "); const prompt = `Напиши краткое резюме формы (2-4 предложения) для корпоративной системы управления задачами. Название: ${form.name} Описание: ${form.description || "не указано"} Поля: ${fieldList || "нет"} Статусы: ${statusList || "нет"} Объясни: для чего эта форма, что означают её статусы, кто обычно работает с ней.`; try { return await generateSummary(prompt, config); } catch (err) { console.error("[RAG] generateFormSummary error:", err); return null; } } export function buildFormText(form: Form, fields: FormField[], statuses: FormStatus[]): string { const parts: string[] = []; parts.push(`Форма: ${form.name}`); if (form.description) parts.push(`Описание: ${form.description}`); if (form.aiSummary) parts.push(`Резюме: ${form.aiSummary}`); if (statuses.length > 0) { const statusStr = statuses .map(s => `${s.name}${s.isInitial ? " (начальный)" : ""}${s.isFinal ? " (финальный)" : ""}`) .join(", "); parts.push(`Статусы: ${statusStr}`); } if (fields.length > 0) { const fieldStr = fields.map(f => `${f.name} [${f.type}]${f.isRequired ? "*" : ""}`).join(", "); parts.push(`Поля: ${fieldStr}`); } return parts.join("\n"); } export function buildTaskText( task: Task, formName: string, statusName: string, assigneeName: string | null, fieldValues: Array<{ name: string; value: unknown; type: string }> ): string { const parts: string[] = []; parts.push(`Задача: ${task.title}`); parts.push(`Форма: ${formName}`); parts.push(`Статус: ${statusName}`); if (assigneeName) parts.push(`Ответственный: ${assigneeName}`); if (task.description) parts.push(`Описание: ${task.description}`); const relevantFields = fieldValues.filter(fv => fv.value !== null && fv.value !== undefined && fv.value !== "" && !["file", "table"].includes(fv.type) ); if (relevantFields.length > 0) { const lines = relevantFields.map(fv => { const val = typeof fv.value === "object" ? JSON.stringify(fv.value) : String(fv.value); return `${fv.name}: ${val}`; }); parts.push(`Поля:\n${lines.join("\n")}`); } return parts.join("\n"); } export function buildMessageText(message: TaskMessage, authorName: string, taskTitle: string): string { const parts: string[] = []; parts.push(`Сообщение в задаче: ${taskTitle}`); parts.push(`Автор: ${authorName}`); parts.push(`Текст: ${message.message}`); return parts.join("\n"); } async function upsertEmbeddingWithVector(item: EmbeddingItem, embedding: number[]): Promise { const vectorStr = `[${embedding.join(",")}]`; await db.execute(sql` INSERT INTO rag_embeddings (organization_id, entity_type, entity_id, content, embedding, metadata, updated_at) VALUES (${item.organizationId}, ${item.entityType}, ${item.entityId}, ${item.content}, ${vectorStr}::vector, ${JSON.stringify(item.metadata)}::jsonb, NOW()) ON CONFLICT (organization_id, entity_type, entity_id) DO UPDATE SET content = EXCLUDED.content, embedding = EXCLUDED.embedding, metadata = EXCLUDED.metadata, updated_at = NOW() `); } async function processBatchedItems( items: EmbeddingItem[], config: LlmProviderConfig, onProgress?: (indexed: number, total: number) => void ): Promise { const hasAccess = config.provider === "ollama" || !!config.apiKey; if (!hasAccess) return 0; const total = items.length; let successCount = 0; let indexed = 0; for (let i = 0; i < items.length; i += BATCH_SIZE) { const batch = items.slice(i, i + BATCH_SIZE); const texts = batch.map(item => item.content); let embeddings = await doEmbeddingBatch(texts, config); const failedIndices: number[] = []; for (let j = 0; j < embeddings.length; j++) { if (embeddings[j] === null) failedIndices.push(j); } if (failedIndices.length > 0) { console.log(`[RAG] Retrying ${failedIndices.length} failed items in batch starting at ${i}`); for (const idx of failedIndices) { const retried = await doEmbedding(texts[idx], config); embeddings[idx] = retried; } } await Promise.all( batch.map(async (item, j) => { if (embeddings[j]) { try { await upsertEmbeddingWithVector(item, embeddings[j]!); successCount++; } catch (err) { console.error(`[RAG] Failed to upsert ${item.entityType} ${item.entityId}:`, err); } } }) ); indexed = Math.min(i + BATCH_SIZE, total); if (onProgress) onProgress(indexed, total); } return successCount; } export async function upsertEmbedding( organizationId: number, entityType: EntityType, entityId: number, content: string, metadata: Record ): Promise { const config = await resolveOrgConfig(organizationId); const hasAccess = config.provider === "ollama" || !!config.apiKey; if (!hasAccess) return; const embedding = await doEmbedding(content, config); if (!embedding) return; const vectorStr = `[${embedding.join(",")}]`; await db.execute(sql` INSERT INTO rag_embeddings (organization_id, entity_type, entity_id, content, embedding, metadata, updated_at) VALUES (${organizationId}, ${entityType}, ${entityId}, ${content}, ${vectorStr}::vector, ${JSON.stringify(metadata)}::jsonb, NOW()) ON CONFLICT (organization_id, entity_type, entity_id) DO UPDATE SET content = EXCLUDED.content, embedding = EXCLUDED.embedding, metadata = EXCLUDED.metadata, updated_at = NOW() `); } export async function semanticSearch( organizationId: number, query: string, entityTypes?: EntityType[], limit: number = 10 ): Promise { const config = await resolveOrgConfig(organizationId); const hasAccess = config.provider === "ollama" || !!config.apiKey; if (!hasAccess) return []; const embedding = await doEmbedding(query, config); if (!embedding) return []; const vectorStr = `[${embedding.join(",")}]`; let rows; if (entityTypes && entityTypes.length > 0) { rows = await db.execute(sql` SELECT entity_type, entity_id, content, metadata, 1 - (embedding <=> ${vectorStr}::vector) AS score FROM rag_embeddings WHERE organization_id = ${organizationId} AND entity_type = ANY(ARRAY[${sql.join(entityTypes.map(t => sql`${t}`), sql`, `)}]) ORDER BY embedding <=> ${vectorStr}::vector LIMIT ${limit} `); } else { rows = await db.execute(sql` SELECT entity_type, entity_id, content, metadata, 1 - (embedding <=> ${vectorStr}::vector) AS score FROM rag_embeddings WHERE organization_id = ${organizationId} ORDER BY embedding <=> ${vectorStr}::vector LIMIT ${limit} `); } return (rows.rows as Array<{ entity_type: string; entity_id: number; content: string; metadata: unknown; score: number; }>).map(r => ({ entityType: r.entity_type as EntityType, entityId: Number(r.entity_id), score: parseFloat(String(r.score)), content: r.content, metadata: (r.metadata as Record) ?? {}, })); } export async function getEmbeddingCounts(organizationId: number): Promise> { const rows = await db.execute(sql` SELECT entity_type, COUNT(*) as count FROM rag_embeddings WHERE organization_id = ${organizationId} GROUP BY entity_type `); const result: Record = {}; for (const row of rows.rows as Array<{ entity_type: string; count: string }>) { result[row.entity_type] = parseInt(row.count); } return result; } export async function deleteRagEmbedding( organizationId: number, entityType: EntityType, entityId: number ): Promise { await db.execute(sql` DELETE FROM rag_embeddings WHERE organization_id = ${organizationId} AND entity_type = ${entityType} AND entity_id = ${entityId} `); } export async function deleteRagEmbeddingsByTask( organizationId: number, taskId: number ): Promise { await db.execute(sql` DELETE FROM rag_embeddings WHERE organization_id = ${organizationId} AND ( (entity_type = 'task' AND entity_id = ${taskId}) OR (entity_type = 'task_message' AND entity_id IN ( SELECT id FROM task_messages WHERE task_id = ${taskId} )) ) `); } export async function deleteRagEmbeddingsByForm( organizationId: number, formId: number ): Promise { await db.execute(sql` DELETE FROM rag_embeddings WHERE organization_id = ${organizationId} AND ( (entity_type = 'form' AND entity_id = ${formId}) OR (entity_type = 'task' AND entity_id IN ( SELECT id FROM tasks WHERE form_id = ${formId} )) OR (entity_type = 'task_message' AND entity_id IN ( SELECT id FROM task_messages WHERE task_id IN ( SELECT id FROM tasks WHERE form_id = ${formId} ) )) ) `); } 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( stor: IStorage, organizationId: number, entityTypes: EntityType[] = ["form", "task", "task_message"], onProgress?: (indexed: number, total: number) => void ): Promise<{ forms: number; tasks: number; messages: number; attempted: number; failed: number }> { const config = await resolveOrgConfig(organizationId); const hasAccess = config.provider === "ollama" || !!config.apiKey; if (!hasAccess) { console.warn(`[RAG] Skipping reindex for org ${organizationId}: no API key configured`); return { forms: 0, tasks: 0, messages: 0, attempted: 0, failed: 0 }; } let formsIndexed = 0; let tasksIndexed = 0; let messagesIndexed = 0; let totalAttempted = 0; let totalFailed = 0; // ── 1. Pre-count: точный COUNT(*) без выгрузки строк (список теперь капнут на 200) ── let taskCount = 0; if (entityTypes.includes("task") || entityTypes.includes("task_message")) { taskCount = await stor.getTasksCountByOrganization(organizationId); } const formCount = entityTypes.includes("form") ? (await stor.getFormsByOrganization(organizationId)).length : 0; const approxTotal = formCount + taskCount; if (onProgress) onProgress(0, approxTotal); let cumulativeIndexed = 0; // ── 2. Index forms ───────────────────────────────────────────────────────── if (entityTypes.includes("form")) { const forms = await stor.getFormsByOrganization(organizationId); const formItems: EmbeddingItem[] = []; for (const baseForm of forms) { const [fields, statuses] = await Promise.all([ stor.getFormFields(baseForm.id, organizationId), stor.getFormStatuses(baseForm.id, organizationId), ]); const summary = await generateFormSummary(baseForm, fields, statuses, organizationId); if (summary) { await stor.updateForm(baseForm.id, organizationId, { aiSummary: summary, aiSummaryIsAuto: true }); } const indexedForm = summary ? { ...baseForm, aiSummary: summary } : baseForm; const content = buildFormText(indexedForm, fields, statuses); formItems.push({ organizationId, entityType: "form", entityId: baseForm.id, content, metadata: { name: baseForm.name }, }); } totalAttempted += formItems.length; const succeeded = await processBatchedItems(formItems, config, (batchIndexed) => { if (onProgress) onProgress(cumulativeIndexed + batchIndexed, approxTotal); }); formsIndexed = succeeded; totalFailed += formItems.length - succeeded; cumulativeIndexed += formItems.length; if (onProgress) onProgress(cumulativeIndexed, approxTotal); } // ── 3. Index tasks and messages in chunks ────────────────────────────────── if (entityTypes.includes("task") || entityTypes.includes("task_message")) { 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; } let chunkOffset = 0; while (true) { const chunk = await stor.getTasksByOrganization(organizationId, { sortBy: "createdAt", order: "asc", minimal: false, limit: TASK_CHUNK_SIZE, offset: chunkOffset, }) as Task[]; if (chunk.length === 0) break; const taskItems: EmbeddingItem[] = []; const messageItems: EmbeddingItem[] = []; for (const task of chunk) { if (entityTypes.includes("task")) { const { form, statuses, fields: formFields } = await getFormMeta(task.formId); if (form) { const statusMap = new Map(statuses.map(s => [s.id, s.name])); const fieldMap = new Map(formFields.map(f => [f.id, f])); const assigneeUser = task.assignedTo ? usersMap.get(task.assignedTo) : null; const assigneeName = assigneeUser ? `${assigneeUser.firstName ?? ""} ${assigneeUser.middleName ?? ""} ${assigneeUser.lastName ?? ""}`.trim() || assigneeUser.email : null; 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); taskItems.push({ organizationId, entityType: "task", entityId: task.id, content, metadata: { formId: task.formId, formName: form.name, title: task.title }, }); } } if (entityTypes.includes("task_message")) { const messages = await stor.getTaskMessages(task.id, organizationId); for (const msg of messages) { if (msg.messageType !== "comment") continue; const authorUser = msg.authorId ? usersMap.get(msg.authorId) : null; const authorName = authorUser ? `${authorUser.firstName ?? ""} ${authorUser.middleName ?? ""} ${authorUser.lastName ?? ""}`.trim() || authorUser.email : "Бот"; const content = buildMessageText(msg, authorName, task.title); messageItems.push({ organizationId, entityType: "task_message", entityId: msg.id, content, metadata: { taskId: task.id, formId: task.formId, authorId: msg.authorId }, }); } } } totalAttempted += taskItems.length; const chunkTasksSucceeded = await processBatchedItems(taskItems, config, (batchIndexed) => { if (onProgress) onProgress(cumulativeIndexed + batchIndexed, approxTotal); }); tasksIndexed += chunkTasksSucceeded; totalFailed += taskItems.length - chunkTasksSucceeded; cumulativeIndexed += taskItems.length; totalAttempted += messageItems.length; const chunkMsgsSucceeded = await processBatchedItems(messageItems, config, (batchIndexed) => { if (onProgress) onProgress(cumulativeIndexed + batchIndexed, approxTotal + messageItems.length); }); messagesIndexed += chunkMsgsSucceeded; totalFailed += messageItems.length - chunkMsgsSucceeded; cumulativeIndexed += messageItems.length; if (onProgress) onProgress(cumulativeIndexed, Math.max(approxTotal, cumulativeIndexed)); if (chunk.length < TASK_CHUNK_SIZE) break; chunkOffset += TASK_CHUNK_SIZE; } } if (onProgress) onProgress(cumulativeIndexed, cumulativeIndexed); console.log(`[RAG] Reindex done for org ${organizationId}: ${formsIndexed} forms, ${tasksIndexed} tasks, ${messagesIndexed} messages (attempted=${totalAttempted}, failed=${totalFailed})`); return { forms: formsIndexed, tasks: tasksIndexed, messages: messagesIndexed, attempted: totalAttempted, failed: totalFailed }; }