Sfoglia il codice sorgente

update no rag sobre os atendimentos

leonardo 2 mesi fa
parent
commit
5f98f45e46

+ 12 - 0
.env.example

@@ -88,3 +88,15 @@ IFBOT_TOKEN=
 # base pública dos arquivos de mídia (/data/img, /data/aud); padrão: IFBOT_BASE_URL sem o sufixo /api
 # IFBOT_MEDIA_BASE_URL=https://star.ifbot.com.br
 ATENDIMENTOS_SYNC_INTERVAL_MINUTES=30
+
+# RAG de atendimentos (ingestão de conversas no Qdrant, coleção separada da de documentos)
+# ATENDIMENTOS_RAG_COLLECTION=atendimentos_rag
+# ATENDIMENTOS_RAG_INTERVAL_MINUTES=60
+# ATENDIMENTOS_RAG_BATCH=10
+# ATENDIMENTOS_RAG_CHUNK_SIZE=1800
+# ATENDIMENTOS_RAG_CHUNK_OVERLAP=200
+# ATENDIMENTOS_RAG_TOP_K=20
+# ATENDIMENTOS_RAG_MIN_SCORE=0.5
+# habilita a busca vetorial no modo "atendimentos" do chat — deixar false até calibrar
+# ATENDIMENTOS_RAG_MIN_SCORE com perguntas reais após o backfill inicial
+# ATENDIMENTOS_RAG_SEARCH_ENABLED=false

+ 16 - 0
db/migrations/20260715180000_create_atendimento_rag_index_table.cjs

@@ -0,0 +1,16 @@
+exports.up = function (knex) {
+  return knex.schema.createTable("atendimento_rag_index", (t) => {
+    t.integer("AtendimentoId").unsigned().primary();
+    t.bigInteger("UltimaMensagemId").unsigned().nullable(); 
+    t.integer("ChunksCount").unsigned().nullable(); 
+    t.string("EmbeddingModel", 100).nullable(); 
+    t.timestamp("CreatedAt").notNullable().defaultTo(knex.fn.now());
+    t.timestamp("UpdatedAt").notNullable().defaultTo(knex.fn.now());
+
+    t.foreign("AtendimentoId").references("atendimentos.Id").onDelete("CASCADE");
+  });
+};
+
+exports.down = function (knex) {
+  return knex.schema.dropTable("atendimento_rag_index");
+};

+ 77 - 0
scripts/ingestarAtendimentosRag.js

@@ -0,0 +1,77 @@
+import { config } from "../src/config/index.js";
+import { ingestarPendentes, contarPendentesRag } from "../src/services/atendimentoRagService.js";
+
+const APLICAR = process.argv.includes("--apply");
+
+const loteIdx = process.argv.indexOf("--lote");
+const LOTE = loteIdx > -1 ? Number(process.argv[loteIdx + 1]) : config.atendimentosRag.batchSize;
+
+const setorIdx = process.argv.indexOf("--setor");
+const SETOR = setorIdx > -1 ? process.argv[setorIdx + 1] : null;
+
+const limiteIdx = process.argv.indexOf("--limite");
+const LIMITE = limiteIdx > -1 ? Number(process.argv[limiteIdx + 1]) : null;
+
+async function main() {
+  const totalPendentes = await contarPendentesRag({ setor: SETOR });
+  console.log(
+    `[ingestar-atendimentos-rag] ${totalPendentes} atendimentos pendentes de indexação` +
+    `${SETOR ? ` no setor ${SETOR}` : ""} (coleção: ${config.atendimentosRag.collection}).`
+  );
+
+  if (!APLICAR) {
+    console.log("[ingestar-atendimentos-rag] dry-run (nenhuma alteração feita). Use --apply para executar.");
+    process.exit(0);
+  }
+  if (totalPendentes === 0) {
+    console.log("[ingestar-atendimentos-rag] nada para ingerir.");
+    process.exit(0);
+  }
+
+  const alvo = LIMITE ? Math.min(LIMITE, totalPendentes) : totalPendentes;
+  console.log(`[ingestar-atendimentos-rag] ingerindo até ${alvo} atendimentos, lotes de ${LOTE}${LIMITE ? ` (--limite ${LIMITE})` : ""}...`);
+
+  let processados = 0;
+  let totalIngeridos = 0;
+  let totalSemConteudo = 0;
+  const erros = [];
+  const inicio = Date.now();
+
+  while (processados < alvo) {
+    const limiteLote = Math.min(LOTE, alvo - processados);
+    const t0 = Date.now();
+    const r = await ingestarPendentes({ limite: limiteLote, setor: SETOR });
+    processados += r.pendentes;
+    totalIngeridos += r.ingeridos;
+    totalSemConteudo += r.semConteudo;
+    erros.push(...r.erros);
+
+    const restantes = await contarPendentesRag({ setor: SETOR });
+    const segundosLote = (Date.now() - t0) / 1000;
+    const ritmo = r.pendentes ? segundosLote / r.pendentes : 0;
+    const etaMin = ritmo && restantes ? ((ritmo * restantes) / 60).toFixed(0) : "?";
+
+    console.log(
+      `[ingestar-atendimentos-rag] +${r.ingeridos} ingeridos, ${r.semConteudo} sem conteúdo, ${r.erros.length} erros ` +
+      `(${segundosLote.toFixed(0)}s, ~${ritmo.toFixed(1)}s/atendimento) — restam ${restantes} (ETA ~${etaMin}min)`
+    );
+    if (r.erros.length) r.erros.slice(0, 3).forEach((e) => console.log(`  erro: ${e.atendimentoId} ${e.erro}`));
+
+    if (r.pendentes === 0) {
+      console.log("[ingestar-atendimentos-rag] lote sem nenhum progresso, abortando para evitar loop infinito.");
+      break;
+    }
+  }
+
+  console.log(
+    `[ingestar-atendimentos-rag] concluído em ${((Date.now() - inicio) / 1000 / 60).toFixed(1)}min: ` +
+    `${totalIngeridos} ingeridos, ${totalSemConteudo} sem conteúdo, ${erros.length} erros.`
+  );
+
+  process.exit(erros.length ? 1 : 0);
+}
+
+main().catch((e) => {
+  console.error("[ingestar-atendimentos-rag] falha:", e.message);
+  process.exit(1);
+});

+ 178 - 0
scripts/reavaliarComExpediente.js

@@ -0,0 +1,178 @@
+import fs from "node:fs";
+import path from "node:path";
+import { fileURLToPath } from "node:url";
+import { AtendimentoAvaliacao } from "../src/models/AtendimentoAvaliacao.model.js";
+import { AtendimentoMensagem } from "../src/models/AtendimentoMensagem.model.js";
+import { avaliarPendentes, contarPendentes } from "../src/services/atendimentoAvaliacaoService.js";
+import { sincronizarTodosProtocolosAbertos } from "../src/services/atendimentosService.js";
+import { EXPEDIENTE_POR_SETOR, estaForaExpediente } from "../src/config/expediente.js";
+
+const __dirname = path.dirname(fileURLToPath(import.meta.url));
+const APLICAR = process.argv.includes("--apply");
+const FORCAR_TODOS = process.argv.includes("--forcar-todos");
+const SEM_SYNC = process.argv.includes("--sem-sync");
+
+const loteIdx = process.argv.indexOf("--lote");
+const LOTE = loteIdx > -1 ? Number(process.argv[loteIdx + 1]) : 20;
+
+const setorIdx = process.argv.indexOf("--setor");
+const SETOR = setorIdx > -1 ? process.argv[setorIdx + 1] : null;
+
+const limiteIdx = process.argv.indexOf("--limite");
+const LIMITE = limiteIdx > -1 ? Number(process.argv[limiteIdx + 1]) : null;
+
+const expectIdx = process.argv.indexOf("--expect-count");
+const EXPECT_COUNT = expectIdx > -1 ? Number(process.argv[expectIdx + 1]) : null;
+
+const SETORES_COM_GRADE = Object.keys(EXPEDIENTE_POR_SETOR);
+
+async function buscarCandidatos() {
+  const query = AtendimentoAvaliacao.query()
+    .where("Avaliavel", true)
+    .join("atendimentos as a", "a.Id", "atendimento_avaliacoes.AtendimentoId")
+    .select("atendimento_avaliacoes.AtendimentoId", "atendimento_avaliacoes.ScoreAtendente", "a.Setor")
+    .orderBy("atendimento_avaliacoes.AtendimentoId", "asc");
+
+  if (SETOR) query.where("a.Setor", SETOR);
+  if (!FORCAR_TODOS) query.whereIn("a.Setor", SETORES_COM_GRADE);
+
+  return query;
+}
+
+async function temMensagemForaExpediente(atendimentoId, setor) {
+  const mensagensCliente = await AtendimentoMensagem.query()
+    .where("AtendimentoId", atendimentoId)
+    .where("Resposta", 0)
+    .select("Timestamp");
+
+  return mensagensCliente.some((m) => estaForaExpediente(setor, m.Timestamp));
+}
+
+async function filtrarAlvos(candidatos) {
+  if (FORCAR_TODOS) return { alvos: candidatos, semMensagemForaExpediente: 0 };
+
+  const alvos = [];
+  let semMensagemForaExpediente = 0;
+
+  for (const c of candidatos) {
+    if (await temMensagemForaExpediente(c.AtendimentoId, c.Setor)) {
+      alvos.push(c);
+    } else {
+      semMensagemForaExpediente += 1;
+    }
+  }
+
+  return { alvos, semMensagemForaExpediente };
+}
+
+async function main() {
+  const totalQuery = AtendimentoAvaliacao.query()
+    .where("Avaliavel", true)
+    .join("atendimentos as a", "a.Id", "atendimento_avaliacoes.AtendimentoId");
+  if (SETOR) totalQuery.where("a.Setor", SETOR);
+  const totalAvaliavel = await totalQuery.resultSize();
+
+  const candidatos = await buscarCandidatos();
+  const descartadosSetorSemGrade = totalAvaliavel - candidatos.length;
+
+  console.log(`[reavaliar-expediente] ${totalAvaliavel} avaliações Avaliavel=true${SETOR ? ` no setor ${SETOR}` : ""}.`);
+  if (!FORCAR_TODOS) {
+    console.log(`[reavaliar-expediente] setores com grade de horário: ${SETORES_COM_GRADE.join(", ")}.`);
+    console.log(`[reavaliar-expediente] ${descartadosSetorSemGrade} descartados: setor sem grade de horário.`);
+  }
+  console.log(`[reavaliar-expediente] ${candidatos.length} candidatos após filtro de setor.`);
+
+  const { alvos: todosAlvos, semMensagemForaExpediente } = await filtrarAlvos(candidatos);
+  if (!FORCAR_TODOS) {
+    console.log(`[reavaliar-expediente] ${semMensagemForaExpediente} descartados: nenhuma mensagem do cliente fora do expediente.`);
+  }
+
+  const alvos = LIMITE ? todosAlvos.slice(0, LIMITE) : todosAlvos;
+  console.log(`[reavaliar-expediente] ${alvos.length} avaliações serão resetadas para pendente${LIMITE ? ` (--limite ${LIMITE})` : ""}.`);
+
+  if (EXPECT_COUNT !== null && alvos.length !== EXPECT_COUNT) {
+    console.error(`[reavaliar-expediente] abortando: esperava ${EXPECT_COUNT}, achei ${alvos.length}.`);
+    process.exit(1);
+  }
+
+  if (!APLICAR) {
+    console.log("[reavaliar-expediente] dry-run (nenhuma alteração feita). Use --apply para executar.");
+    process.exit(0);
+  }
+  if (alvos.length === 0) {
+    console.log("[reavaliar-expediente] nada para reavaliar.");
+    process.exit(0);
+  }
+
+  const backupDir = path.resolve(__dirname, "logs");
+  fs.mkdirSync(backupDir, { recursive: true });
+  const backupPath = path.join(backupDir, `reavaliar-expediente-${new Date().toISOString().replace(/[:.]/g, "-")}.json`);
+  fs.writeFileSync(backupPath, JSON.stringify(alvos, null, 2));
+  console.log(`[reavaliar-expediente] backup dos scores antigos salvo em ${backupPath}`);
+
+  if (!SEM_SYNC) {
+    console.log("[reavaliar-expediente] passo 1/3: sincronizando protocolos abertos (sincronizarTodosProtocolosAbertos)...");
+    try {
+      const sync = await sincronizarTodosProtocolosAbertos();
+      console.log(
+        `[reavaliar-expediente] sync concluído: ${sync.importados}/${sync.processados} importados, ` +
+        `${sync.mensagens} mensagens, ${sync.erros.length} erros, ${sync.paginas} página(s).`
+      );
+    } catch (err) {
+      console.error(`[reavaliar-expediente] sync falhou (${err.message}), seguindo mesmo assim.`);
+    }
+  } else {
+    console.log("[reavaliar-expediente] passo 1/3: pulado (--sem-sync).");
+  }
+
+  const ids = alvos.map((a) => a.AtendimentoId);
+  const pendentesAntes = await contarPendentes();
+  await AtendimentoAvaliacao.query().whereIn("AtendimentoId", ids).patch({ UltimaMensagemId: null });
+  const pendentesDepois = await contarPendentes();
+  console.log(`[reavaliar-expediente] passo 2/3: ${ids.length} avaliações resetadas — fila de pendentes ${pendentesAntes} -> ${pendentesDepois}.`);
+
+  console.log(`[reavaliar-expediente] passo 3/3: drenando fila de pendentes (mesmo fluxo de avaliarTudo.js, lotes de ${LOTE})...`);
+  let restantes = pendentesDepois;
+  let totalAvaliados = 0;
+  let totalSemDialogo = 0;
+  const erros = [];
+  const inicio = Date.now();
+
+  while (restantes > 0) {
+    const t0 = Date.now();
+    const r = await avaliarPendentes({ limite: LOTE });
+    restantes = await contarPendentes();
+    totalAvaliados += r.avaliados;
+    totalSemDialogo += r.semDialogo;
+    erros.push(...r.erros);
+
+    console.log(
+      `[reavaliar-expediente] +${r.avaliados} avaliados, ${r.semDialogo} sem diálogo, ${r.erros.length} erros ` +
+      `(${((Date.now() - t0) / 1000).toFixed(0)}s) — restam ${restantes}`
+    );
+    if (r.erros.length) r.erros.slice(0, 3).forEach((e) => console.log(`  erro: ${e.atendimentoId} ${e.erro}`));
+
+    if (r.avaliados === 0 && r.semDialogo === 0 && r.erros.length === 0) {
+      console.error("[reavaliar-expediente] lote sem nenhum progresso, abortando para evitar loop infinito.");
+      break;
+    }
+  }
+
+  console.log(
+    `[reavaliar-expediente] fila drenada em ${(((Date.now() - inicio)) / 1000 / 60).toFixed(1)}min: ` +
+    `${totalAvaliados} avaliados, ${totalSemDialogo} sem diálogo, ${erros.length} erros.`
+  );
+
+  const idsSet = new Set(ids);
+  const scoreAntigoPorId = new Map(alvos.map((a) => [a.AtendimentoId, a.ScoreAtendente]));
+  const atuais = await AtendimentoAvaliacao.query().whereIn("AtendimentoId", ids).select("AtendimentoId", "ScoreAtendente");
+  const mudouScore = atuais.filter((a) => idsSet.has(a.AtendimentoId) && a.ScoreAtendente !== scoreAntigoPorId.get(a.AtendimentoId)).length;
+  console.log(`[reavaliar-expediente] dos ${ids.length} alvos do backfill, ${mudouScore} tiveram o score alterado.`);
+
+  process.exit(erros.length ? 1 : 0);
+}
+
+main().catch((e) => {
+  console.error("[reavaliar-expediente] falha:", e.message);
+  process.exit(1);
+});

+ 10 - 0
src/config/index.js

@@ -85,6 +85,16 @@ export const config = {
     intervalMinutes: Number(process.env.ATENDIMENTOS_AVALIACAO_INTERVAL_MINUTES ?? 60),
     batchSize: Number(process.env.ATENDIMENTOS_AVALIACAO_BATCH ?? 10),
     maxChars: Number(process.env.ATENDIMENTOS_AVALIACAO_MAX_CHARS ?? 12000)
+  },
+  atendimentosRag: {
+    collection: process.env.ATENDIMENTOS_RAG_COLLECTION ?? "atendimentos_rag",
+    intervalMinutes: Number(process.env.ATENDIMENTOS_RAG_INTERVAL_MINUTES ?? 60),
+    batchSize: Number(process.env.ATENDIMENTOS_RAG_BATCH ?? 10),
+    chunkSize: Number(process.env.ATENDIMENTOS_RAG_CHUNK_SIZE ?? 1800),
+    chunkOverlap: Number(process.env.ATENDIMENTOS_RAG_CHUNK_OVERLAP ?? 200),
+    topK: Number(process.env.ATENDIMENTOS_RAG_TOP_K ?? 20),
+    minScore: Number(process.env.ATENDIMENTOS_RAG_MIN_SCORE ?? 0.5),
+    searchEnabled: process.env.ATENDIMENTOS_RAG_SEARCH_ENABLED === "true"
   }
 };
 

+ 2 - 0
src/factories/Server.factory.js

@@ -11,6 +11,7 @@ import { Roteamento } from "../routes/index.js";
 import { scheduleCleanup } from "../jobs/cleanupTokens.js";
 import { scheduleSyncAtendimentos } from "../jobs/syncAtendimentos.js";
 import { scheduleAvaliarAtendimentos } from "../jobs/avaliarAtendimentos.js";
+import { scheduleIngestAtendimentosRag } from "../jobs/ingestAtendimentosRag.js";
 
 export class ServerFactory {
   static Iniciar() {
@@ -57,6 +58,7 @@ export class ServerFactory {
       scheduleCleanup();
       scheduleSyncAtendimentos();
       scheduleAvaliarAtendimentos();
+      scheduleIngestAtendimentosRag();
     });
 
     this.app = app;

+ 26 - 0
src/jobs/ingestAtendimentosRag.js

@@ -0,0 +1,26 @@
+import { config } from "../config/index.js";
+import { ingestarPendentes } from "../services/atendimentoRagService.js";
+
+export async function runIngestAtendimentosRag() {
+  const r = await ingestarPendentes();
+  console.log(
+    `[ingest-atendimentos-rag] ${r.ingeridos} ingeridos, ${r.semConteudo} sem conteúdo, ` +
+    `${r.erros.length} erros (${r.pendentes} pendentes no lote)`
+  );
+  return r;
+}
+
+export function scheduleIngestAtendimentosRag(intervalMs = config.atendimentosRag.intervalMinutes * 60_000) {
+  if (!intervalMs || intervalMs <= 0) {
+    console.warn("[ingest-atendimentos-rag] desabilitado — ATENDIMENTOS_RAG_INTERVAL_MINUTES=0");
+    return null;
+  }
+
+  console.log(`[ingest-atendimentos-rag] agendado a cada ${Math.round(intervalMs / 60_000)} min`);
+  return setInterval(() => {
+    runIngestAtendimentosRag().catch((e) => {
+      if (e?.message === "rag_ingestao_em_andamento") return;
+      console.error("[ingest-atendimentos-rag] falha:", e.message);
+    });
+  }, intervalMs).unref();
+}

+ 6 - 0
src/models/Atendimento.model.js

@@ -3,6 +3,7 @@ import "../config/db.config.js";
 import { AtendimentoCliente } from "./AtendimentoCliente.model.js";
 import { AtendimentoMensagem } from "./AtendimentoMensagem.model.js";
 import { AtendimentoAvaliacao } from "./AtendimentoAvaliacao.model.js";
+import { AtendimentoRagIndex } from "./AtendimentoRagIndex.model.js";
 
 export class Atendimento extends Model {
   static get tableName() {
@@ -29,6 +30,11 @@ export class Atendimento extends Model {
         relation: Model.HasOneRelation,
         modelClass: AtendimentoAvaliacao,
         join: { from: "atendimentos.Id", to: "atendimento_avaliacoes.AtendimentoId" }
+      },
+      ragIndex: {
+        relation: Model.HasOneRelation,
+        modelClass: AtendimentoRagIndex,
+        join: { from: "atendimentos.Id", to: "atendimento_rag_index.AtendimentoId" }
       }
     };
   }

+ 14 - 0
src/models/AtendimentoRagIndex.model.js

@@ -0,0 +1,14 @@
+import { Model } from "objection";
+import "../config/db.config.js";
+
+export class AtendimentoRagIndex extends Model {
+  static get tableName() {
+    return "atendimento_rag_index";
+  }
+
+  static get idColumn() {
+    return "AtendimentoId";
+  }
+}
+
+export default AtendimentoRagIndex;

+ 2 - 28
src/services/atendimentoAvaliacaoService.js

@@ -1,7 +1,6 @@
 import { config } from "../config/index.js";
 import { chatCompletion } from "./ollamaClient.js";
-import { formatMensagem, formatDateTime } from "../utils/atendimentoFormat.js";
-import { estaForaExpediente } from "../config/expediente.js";
+import { buildConversaTexto, calcularUltimaMensagemId } from "../utils/atendimentoFormat.js";
 import { Atendimento } from "../models/Atendimento.model.js";
 import { AtendimentoAvaliacao } from "../models/AtendimentoAvaliacao.model.js";
 import { AtendimentoMensagem } from "../models/AtendimentoMensagem.model.js";
@@ -105,31 +104,6 @@ const AVALIACAO_JSON_SCHEMA = {
 };
 
 
-function buildConversaTexto({ atendimento, cliente, mensagens }) {
-  const ordenadas = [...mensagens].sort((a, b) => {
-    const ta = a.Timestamp ? new Date(a.Timestamp).getTime() : 0;
-    const tb = b.Timestamp ? new Date(b.Timestamp).getTime() : 0;
-    return ta - tb || (Number(a.Id) || 0) - (Number(b.Id) || 0);
-  });
-
-  const cabecalho = [
-    `Atendimento ${atendimento.Codigo}` +
-      (atendimento.Setor ? ` — Setor: ${atendimento.Setor}` : "") +
-      (atendimento.Status !== null && atendimento.Status !== undefined ? ` — Status: ${atendimento.Status}` : ""),
-    `Cliente: ${cliente?.Nome ?? "desconhecido"}`,
-    atendimento.Abertura ? `Aberto em: ${formatDateTime(atendimento.Abertura)}` : null
-  ].filter(Boolean);
-
-  const linhas = ordenadas
-    .map((m) =>
-      formatMensagem(m, cliente?.Nome, {
-        foraExpediente: m.Resposta === 0 && estaForaExpediente(atendimento.Setor, m.Timestamp)
-      })
-    )
-    .filter(Boolean);
-  return [...cabecalho, "", ...linhas].join("\n").trim();
-}
-
 const TRUNCATION_MARKER = "\n[... trecho intermediário da conversa omitido ...]\n";
 
 
@@ -232,7 +206,7 @@ export async function avaliarAtendimento(atendimentoId) {
   if (!atendimento) throw new NotFoundError("atendimento_not_found");
 
   const mensagens = atendimento.mensagens ?? [];
-  const ultimaMensagemId = mensagens.reduce((max, m) => (Number(m.Id) > max ? Number(m.Id) : max), 0) || null;
+  const ultimaMensagemId = calcularUltimaMensagemId(mensagens);
   const { cliente: msgsCliente, atendente: msgsAtendente } = contarInteracoes(mensagens);
 
   // sem diálogo real não há o que avaliar; grava como não-avaliável para sair da fila

+ 140 - 0
src/services/atendimentoRagService.js

@@ -0,0 +1,140 @@
+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 { NotFoundError } from "../shared/errors/index.js";
+
+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;
+
+// pendente = atendimento com mensagens e sem índice RAG, ou com mensagens mais novas que a última indexada
+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: [] };
+
+    // sequencial de propósito: o Ollama local não se beneficia de concorrência aqui
+    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
+}) {
+  return searchDocs({
+    query,
+    topK,
+    minScore,
+    collectionName: config.atendimentosRag.collection,
+    embedRole: "query"
+  });
+}

+ 51 - 9
src/services/atendimentosQueryService.js

@@ -1,8 +1,10 @@
+import { config } from "../config/index.js";
 import { chatCompletion } from "./ollamaClient.js";
 import { Atendimento } from "../models/Atendimento.model.js";
 import { AtendimentoCliente } from "../models/AtendimentoCliente.model.js";
 import { AtendimentoMensagem } from "../models/AtendimentoMensagem.model.js";
 import { formatMensagem, formatDateTime, parseAtendenteBody } from "../utils/atendimentoFormat.js";
+import { buscarAtendimentosSemelhantes } from "./atendimentoRagService.js";
 
 const BUSCA_LIMIT_ATENDIMENTOS = 5;
 
@@ -17,10 +19,11 @@ function extrairCodigosProtocolo(texto) {
 
 async function buscarPorCodigos(codigos) {
   if (!codigos.length) return [];
-  return Atendimento.query()
+  const atendimentos = await Atendimento.query()
     .whereIn("Codigo", codigos)
     .withGraphFetched("[cliente, mensagens]")
     .modifyGraph("mensagens", (q) => q.orderBy("Timestamp", "asc").orderBy("Id", "asc"));
+  return atendimentos.map((atendimento) => ({ atendimento, score: null }));
 }
 
 // nome do cliente vive em atendimento_clientes.Nome — assim como o código, não
@@ -35,12 +38,13 @@ async function buscarPorNomeCliente(nome, limit = BUSCA_LIMIT_ATENDIMENTOS) {
   const clienteIds = clientes.map((c) => c.Id);
   if (!clienteIds.length) return [];
 
-  return Atendimento.query()
+  const atendimentos = await Atendimento.query()
     .whereIn("ClienteId", clienteIds)
     .withGraphFetched("[cliente, mensagens]")
     .modifyGraph("mensagens", (q) => q.orderBy("Timestamp", "asc").orderBy("Id", "asc"))
     .orderByRaw("COALESCE(Abertura, IngestedAt) DESC")
     .limit(limit);
+  return atendimentos.map((atendimento) => ({ atendimento, score: null }));
 }
 
 function classificacaoSystemPrompt() {
@@ -170,7 +174,43 @@ async function buscarAtendimentosPorTexto({ termos, setor, status, limit = BUSCA
   const atendimentos = await query;
   const ordem = new Map(ids.map((id, idx) => [id, idx]));
   atendimentos.sort((a, b) => (ordem.get(a.Id) ?? 0) - (ordem.get(b.Id) ?? 0));
-  return atendimentos;
+  return atendimentos.map((atendimento) => ({ atendimento, score: null }));
+}
+
+// busca semântica via Qdrant (coleção separada de atendimentos, config.atendimentosRag) —
+// alternativa ao FULLTEXT acima, atrás de config.atendimentosRag.searchEnabled. Os hits do
+// Qdrant vêm em nível de CHUNK; agrupamos por atendimentoId mantendo o melhor score entre
+// os chunks do mesmo atendimento antes de recarregar os registros completos do MySQL.
+async function buscarAtendimentosPorVetor({ termos, setor, status, limit = BUSCA_LIMIT_ATENDIMENTOS }) {
+  const termo = String(termos ?? "").trim();
+  if (!termo) return [];
+
+  const hits = await buscarAtendimentosSemelhantes({ query: termo });
+
+  const melhorScorePorId = new Map();
+  for (const h of hits) {
+    const id = h.metadata?.atendimentoId;
+    if (!id) continue;
+    if (!melhorScorePorId.has(id) || h.score > melhorScorePorId.get(id)) {
+      melhorScorePorId.set(id, h.score);
+    }
+  }
+  const ids = [...melhorScorePorId.keys()];
+  if (!ids.length) return [];
+
+  let query = Atendimento.query()
+    .findByIds(ids)
+    .withGraphFetched("[cliente, mensagens]")
+    .modifyGraph("mensagens", (q) => q.orderBy("Timestamp", "asc").orderBy("Id", "asc"));
+
+  if (setor) query = query.where("Setor", setor);
+  if (status !== null && status !== undefined) query = query.where("Status", status);
+
+  const atendimentos = await query;
+  return atendimentos
+    .map((atendimento) => ({ atendimento, score: melhorScorePorId.get(atendimento.Id) ?? null }))
+    .sort((a, b) => (b.score ?? 0) - (a.score ?? 0))
+    .slice(0, limit);
 }
 
 // calcula, de forma determinística (não depende do LLM "contar certo"), quais
@@ -259,20 +299,22 @@ export async function resolverContextoAtendimentos(message, history = []) {
     codigos = extrairCodigosProtocolo(history.map((m) => m.Content).join("\n"));
   }
 
-  let atendimentos;
+  let resultados;
   if (codigos.length) {
-    atendimentos = await buscarPorCodigos(codigos);
+    resultados = await buscarPorCodigos(codigos);
   } else if (filtros.clienteNome) {
-    atendimentos = await buscarPorNomeCliente(filtros.clienteNome);
+    resultados = await buscarPorNomeCliente(filtros.clienteNome);
+  } else if (config.atendimentosRag.searchEnabled) {
+    resultados = await buscarAtendimentosPorVetor(filtros);
   } else {
-    atendimentos = await buscarAtendimentosPorTexto(filtros);
+    resultados = await buscarAtendimentosPorTexto(filtros);
   }
-  return atendimentos.map((a) => ({
+  return resultados.map(({ atendimento: a, score }) => ({
     id: `atendimento:${a.Codigo}`,
     source: `atendimento:${a.Codigo}`,
     text: formatarTrechoAtendimento(a),
     chunkIndex: 0,
     metadata: { type: "atendimento", codigo: a.Codigo, setor: a.Setor, status: a.Status, clienteNome: a.cliente?.Nome ?? null },
-    score: null
+    score
   }));
 }

+ 1 - 2
src/services/collectionService.js

@@ -1,8 +1,7 @@
 import { config } from "../config/index.js";
 import { qdrant } from "./qdrantClient.js";
 
-export async function ensureCollection({ vectorSize }) {
-  const collectionName = config.qdrant.collection;
+export async function ensureCollection({ vectorSize, collectionName = config.qdrant.collection }) {
   const existing = await qdrant.getCollections();
   const has = (existing?.collections ?? []).some((c) => c.name === collectionName);
 

+ 2 - 5
src/services/documentsService.js

@@ -1,9 +1,7 @@
 import { config } from "../config/index.js";
 import { qdrant } from "./qdrantClient.js";
 
-export async function listDocuments({ limit, offset }) {
-  const collectionName = config.qdrant.collection;
-
+export async function listDocuments({ limit, offset, collectionName = config.qdrant.collection }) {
   let points;
   try {
     points = await qdrant.scroll(collectionName, {
@@ -32,8 +30,7 @@ export async function listDocuments({ limit, offset }) {
   return { items, nextOffset };
 }
 
-export async function deleteDocumentsBySource(source) {
-  const collectionName = config.qdrant.collection;
+export async function deleteDocumentsBySource(source, { collectionName = config.qdrant.collection } = {}) {
   await qdrant.delete(collectionName, {
     wait: true,
     filter: {

+ 9 - 7
src/services/ingestService.js

@@ -46,8 +46,11 @@ function stableUuid(seed) {
   ].join("-");
 }
 
-export async function ingestDocuments(documents) {
-  const collectionName = config.qdrant.collection;
+export async function ingestDocuments(documents, {
+  collectionName = config.qdrant.collection,
+  chunkSize = config.rag.chunkSize,
+  chunkOverlap = config.rag.chunkOverlap
+} = {}) {
   const allChunks = [];
   const ingestedAt = new Date().toISOString();
 
@@ -55,8 +58,8 @@ export async function ingestDocuments(documents) {
     const docHash = createHash("sha256").update(doc.text).digest("hex");
     const baseSeed = doc.id ?? docHash;
     const chunks = chunkText(doc.text, {
-      chunkSize: config.rag.chunkSize,
-      chunkOverlap: config.rag.chunkOverlap
+      chunkSize,
+      chunkOverlap
     });
     chunks.forEach((chunk, idx) => {
       allChunks.push({
@@ -81,12 +84,11 @@ export async function ingestDocuments(documents) {
     throw err;
   }
 
-  await ensureCollection({ vectorSize });
+  await ensureCollection({ vectorSize, collectionName });
 
-  
   const sources = [...new Set(documents.map((d) => d.source).filter(Boolean))];
   for (const source of sources) {
-    await deleteDocumentsBySource(source);
+    await deleteDocumentsBySource(source, { collectionName });
   }
 
   const points = allChunks.map((c, idx) => ({

+ 1 - 2
src/services/searchService.js

@@ -2,8 +2,7 @@ import { config } from "../config/index.js";
 import { embedTexts } from "./ollamaClient.js";
 import { qdrant } from "./qdrantClient.js";
 
-export async function searchDocs({ query, topK, embedRole = "query", minScore }) {
-  const collectionName = config.qdrant.collection;
+export async function searchDocs({ query, topK, embedRole = "query", minScore, collectionName = config.qdrant.collection }) {
   const [vector] = await embedTexts([query], { role: embedRole });
 
   let result;

+ 31 - 0
src/utils/atendimentoFormat.js

@@ -1,3 +1,5 @@
+import { estaForaExpediente } from "../config/expediente.js";
+
 // mensagens do ifbot: Resposta 0 = cliente, 1 = atendente; status de sistema usam
 // códigos variados (8, 9, 10, 21, 22, ...) sempre com Tipodemidia "ifstatus".
 // o campo Autor normalmente vem null; o nome do atendente é embutido no próprio
@@ -113,3 +115,32 @@ export function formatMensagem(m, clienteNome, { foraExpediente = false } = {})
   const timestamp = formatDateTime(m.Timestamp);
   return timestamp ? `[${timestamp}] ${label}: ${texto}` : `${label}: ${texto}`;
 }
+
+export function buildConversaTexto({ atendimento, cliente, mensagens }) {
+  const ordenadas = [...mensagens].sort((a, b) => {
+    const ta = a.Timestamp ? new Date(a.Timestamp).getTime() : 0;
+    const tb = b.Timestamp ? new Date(b.Timestamp).getTime() : 0;
+    return ta - tb || (Number(a.Id) || 0) - (Number(b.Id) || 0);
+  });
+
+  const cabecalho = [
+    `Atendimento ${atendimento.Codigo}` +
+      (atendimento.Setor ? ` — Setor: ${atendimento.Setor}` : "") +
+      (atendimento.Status !== null && atendimento.Status !== undefined ? ` — Status: ${atendimento.Status}` : ""),
+    `Cliente: ${cliente?.Nome ?? "desconhecido"}`,
+    atendimento.Abertura ? `Aberto em: ${formatDateTime(atendimento.Abertura)}` : null
+  ].filter(Boolean);
+
+  const linhas = ordenadas
+    .map((m) =>
+      formatMensagem(m, cliente?.Nome, {
+        foraExpediente: m.Resposta === 0 && estaForaExpediente(atendimento.Setor, m.Timestamp)
+      })
+    )
+    .filter(Boolean);
+  return [...cabecalho, "", ...linhas].join("\n").trim();
+}
+
+export function calcularUltimaMensagemId(mensagens) {
+  return mensagens.reduce((max, m) => (Number(m.Id) > max ? Number(m.Id) : max), 0) || null;
+}