Files
iistwin/server/services/embedding.service.ts

822 lines
30 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 { 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<string, unknown>;
}
interface EmbeddingItem {
organizationId: number;
entityType: EntityType;
entityId: number;
content: string;
metadata: Record<string, unknown>;
}
export interface OrgEmbeddingConfig {
provider: string;
apiKey: string | null;
baseUrl: string | null;
embeddingModel: string;
customHeaders: Record<string, string>;
}
async function resolveOrgConfig(organizationId: number): Promise<LlmProviderConfig> {
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<string, string>,
};
}
} 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<number[] | null> {
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<LlmProviderConfig> {
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<string, string>,
};
}
} 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<number[] | null> {
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<string | null> {
// 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<void> {
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<number> {
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<string, unknown>
): Promise<void> {
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<RagResult[]> {
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<string, unknown>) ?? {},
}));
}
export async function getEmbeddingCounts(organizationId: number): Promise<Record<string, number>> {
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<string, number> = {};
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<void> {
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<void> {
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<void> {
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<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(
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: fetch minimal task list to know total for progress ─────────
let taskCount = 0;
if (entityTypes.includes("task") || entityTypes.includes("task_message")) {
const minimalTasks = await stor.getTasksByOrganization(organizationId, {
sortBy: "createdAt",
order: "asc",
minimal: true,
limit: 1_000_000,
});
taskCount = minimalTasks.length;
}
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<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;
}
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 };
}