822 lines
30 KiB
TypeScript
822 lines
30 KiB
TypeScript
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 };
|
||
}
|