import { config } from "../config/index.js"; import { Atendimento } from "../models/Atendimento.model.js"; import { AtendimentoMensagem } from "../models/AtendimentoMensagem.model.js"; import { AtendimentoRagIndex } from "../models/AtendimentoRagIndex.model.js"; import { buildConversaTexto, calcularUltimaMensagemId } from "../utils/atendimentoFormat.js"; import { ingestDocuments } from "./ingestService.js"; import { searchDocs } from "./searchService.js"; import { generateHydeDocument } from "./hydeService.js"; import { rerankHits } from "./rerankService.js"; import { NotFoundError } from "../shared/errors/index.js"; const ATENDIMENTO_HYDE_SYSTEM_PROMPT = [ "Você é um atendente experiente de suporte ao cliente. Dada a pergunta ou descrição", "abaixo, escreva um parágrafo curto (3 a 6 frases) no estilo de um trecho real de", "uma conversa de atendimento ao cliente (chat de suporte) que trate exatamente desse", "assunto — mensagens típicas de cliente e de atendente, termos usados nesse tipo de", "conversa (ex.: protocolo, boleto, fatura, cancelamento, setor, prazo), quando fizer", "sentido. Não inclua a pergunta original, saudações genéricas ou ressalvas de", "incerteza. Se não tiver certeza do conteúdo exato, escreva de forma plausível no", "mesmo estilo, pois o texto será usado apenas para busca por similaridade, nunca", "mostrado ao usuário. Responda em português brasileiro." ].join("\n"); const ATENDIMENTO_HYDE_MIN_QUERY_LENGTH = 15; function shouldSkipAtendimentoHyde(query) { return query.trim().length < ATENDIMENTO_HYDE_MIN_QUERY_LENGTH; } async function upsertRagIndex(atendimentoId, dados) { const now = new Date(); const row = { AtendimentoId: atendimentoId, ...dados, UpdatedAt: now }; const merge = Object.keys(dados).concat("UpdatedAt"); await AtendimentoRagIndex.query().insert(row).onConflict("AtendimentoId").merge(merge); return AtendimentoRagIndex.query().findById(atendimentoId); } export async function ingestarAtendimento(atendimentoId) { const atendimento = await Atendimento.query() .findById(atendimentoId) .withGraphFetched("[cliente, mensagens]") .modifyGraph("mensagens", (q) => q.orderBy("Timestamp", "asc").orderBy("Id", "asc")); if (!atendimento) throw new NotFoundError("atendimento_not_found"); const mensagens = atendimento.mensagens ?? []; const ultimaMensagemId = calcularUltimaMensagemId(mensagens); const conversaTexto = buildConversaTexto({ atendimento, cliente: atendimento.cliente, mensagens }); // conversa sem conteúdo indexável (ex.: só mensagens de sistema) — grava com // ChunksCount:0 pra sair da fila de pendentes; volta a aparecer se chegarem mensagens novas if (!conversaTexto) { await upsertRagIndex(atendimento.Id, { UltimaMensagemId: ultimaMensagemId, ChunksCount: 0, EmbeddingModel: null }); return { ingestado: false, chunks: 0 }; } const { upserted } = await ingestDocuments( [ { id: atendimento.Id, source: `atendimento:${atendimento.Codigo}`, title: `Atendimento ${atendimento.Codigo}`, metadata: { atendimentoId: atendimento.Id, codigo: atendimento.Codigo, setor: atendimento.Setor }, text: conversaTexto } ], { collectionName: config.atendimentosRag.collection, chunkSize: config.atendimentosRag.chunkSize, chunkOverlap: config.atendimentosRag.chunkOverlap } ); await upsertRagIndex(atendimento.Id, { UltimaMensagemId: ultimaMensagemId, ChunksCount: upserted, EmbeddingModel: config.ollama.embeddingsModel }); return { ingestado: true, chunks: upserted }; } let ragIngestaoEmAndamento = false; function pendentesQuery({ setor } = {}) { const subMax = AtendimentoMensagem.query() .select("AtendimentoId") .max("Id as MaxMsgId") .groupBy("AtendimentoId") .as("m"); const query = Atendimento.query() .select("atendimentos.Id") .join(subMax, "m.AtendimentoId", "atendimentos.Id") .leftJoin("atendimento_rag_index as ri", "ri.AtendimentoId", "atendimentos.Id") .where((q) => { q.whereNull("ri.AtendimentoId") .orWhereNull("ri.UltimaMensagemId") .orWhereRaw("ri.UltimaMensagemId < m.MaxMsgId"); }); if (setor) query.where("atendimentos.Setor", setor); return query; } function buscarPendentes(limite, { setor } = {}) { return pendentesQuery({ setor }).orderBy("atendimentos.IngestedAt", "desc").limit(limite); } export async function contarPendentesRag(params = {}) { return pendentesQuery(params).resultSize(); } export async function ingestarPendentes({ limite = config.atendimentosRag.batchSize, setor } = {}) { if (ragIngestaoEmAndamento) { const err = new Error("rag_ingestao_em_andamento"); err.statusCode = 409; throw err; } ragIngestaoEmAndamento = true; try { const pendentes = await buscarPendentes(limite, { setor }); const resultado = { pendentes: pendentes.length, ingeridos: 0, semConteudo: 0, erros: [] }; for (const { Id } of pendentes) { try { const r = await ingestarAtendimento(Id); if (r.ingestado) resultado.ingeridos += 1; else resultado.semConteudo += 1; } catch (err) { resultado.erros.push({ atendimentoId: Id, erro: err?.message ?? "erro_desconhecido" }); } } return resultado; } finally { ragIngestaoEmAndamento = false; } } export async function buscarAtendimentosSemelhantes({ query, topK = config.atendimentosRag.topK, minScore = config.atendimentosRag.minScore, hyde = config.atendimentosRag.hyde, collectionName = config.atendimentosRag.collection, embeddingsModel = config.ollama.embeddingsModel }) { let embeddingQuery = query; let embedRole = "query"; if (hyde && !shouldSkipAtendimentoHyde(query)) { const hydeText = await generateHydeDocument(query, ATENDIMENTO_HYDE_SYSTEM_PROMPT, { numCtx: config.llm.atendimentosNumCtx }); if (hydeText) { embeddingQuery = hydeText; embedRole = "passage"; } } const hits = await searchDocs({ query: embeddingQuery, topK, minScore, collectionName, embedRole, embeddingsModel }); if (!config.atendimentosRag.rerank || hits.length <= 1) return hits; return rerankHits(query, hits, { topN: config.atendimentosRag.rerankTopN, numCtx: config.llm.atendimentosNumCtx }); }