leonardo преди 2 месеца
родител
ревизия
535fc37ebf

+ 16 - 0
.env.example

@@ -30,6 +30,14 @@ OLLAMA_EMBEDDINGS_MODEL=nomic-embed-text
 OLLAMA_CHAT_MODEL=llama3
 OLLAMA_VISION_MODEL=llava
 
+# Prefixos opcionais prependados ao texto antes de gerar embedding (útil para modelos
+# assimétricos que recomendam prefixo diferente para query de busca e passagem indexada,
+# ex.: "search_query: "/"search_document: " no nomic-embed-text, "query: "/"passage: " em
+# modelos da família e5). Deixe vazio se o modelo não precisar. Validar empiricamente com
+# backend/scripts/evalRetrieval.mjs antes de habilitar em produção.
+OLLAMA_EMBED_QUERY_PREFIX=
+OLLAMA_EMBED_PASSAGE_PREFIX=
+
 # RAG (valores padrão se não definidos)
 RAG_TOP_K=8
 RAG_CHUNK_SIZE=900
@@ -47,3 +55,11 @@ LLM_REPEAT_PENALTY=1.1
 # RAG avançado
 RAG_QUERY_REWRITE=false
 RAG_MIN_SCORE=0.40
+RAG_RERANK=false
+RAG_RERANK_TOP_N=8
+
+# HyDE (Hypothetical Document Embeddings): gera um trecho hipotético de resposta via LLM
+# e usa ele (em vez da query crua) para o embedding de busca, o que tende a elevar o score
+# de similaridade porque vira matching passagem-passagem. Custa uma chamada LLM extra por
+# pergunta. Validar com backend/scripts/evalRetrieval.mjs antes de habilitar em produção.
+RAG_HYDE=false

+ 3 - 1
docs/setup.md

@@ -62,12 +62,14 @@ Crie o arquivo `backend/.env` com base nas variáveis abaixo:
 
 | Variável | Descrição | Exemplo |
 |---|---|---|
-| `RAG_QUERY_REWRITE` | Reescrever query antes da busca | `true` |
+| `RAG_QUERY_REWRITE` | Reescrever query antes da busca, considerando o histórico da conversa (resolve follow-ups tipo "e o segundo caso?") | `true` |
 | `RAG_HISTORY_TOKEN_BUDGET` | Tokens de histórico enviados ao LLM | `1500` |
 | `RAG_TOP_K` | Número de chunks recuperados (default `8`) | `8` |
 | `RAG_MIN_SCORE` | Score mínimo de similaridade (default `0.40`) | `0.40` |
 | `RAG_CHUNK_SIZE` | Tamanho de cada chunk (caracteres) | `900` |
 | `RAG_CHUNK_OVERLAP` | Sobreposição entre chunks | `150` |
+| `RAG_RERANK` | Reordenar os chunks recuperados por relevância via LLM antes de montar o contexto | `true` |
+| `RAG_RERANK_TOP_N` | Quantos chunks manter após o reranking (default = `RAG_TOP_K`) | `8` |
 
 ### LLM (parâmetros de geração)
 

+ 112 - 0
scripts/evalRetrieval.mjs

@@ -0,0 +1,112 @@
+).
+
+import { config } from "#config/index.js";
+import { searchDocs } from "#services/searchService.js";
+import { rewriteQuery, generateHydeDocument } from "#chat/chatChain.js";
+
+function parseArgs(argv) {
+  const flags = { hyde: false, rewrite: false, minScore: 0, topK: 5, label: "" };
+  for (const arg of argv) {
+    if (arg === "--hyde") flags.hyde = true;
+    else if (arg === "--no-hyde") flags.hyde = false;
+    else if (arg === "--rewrite") flags.rewrite = true;
+    else if (arg === "--no-rewrite") flags.rewrite = false;
+    else if (arg.startsWith("--min-score=")) flags.minScore = Number(arg.split("=")[1]);
+    else if (arg.startsWith("--top-k=")) flags.topK = Number(arg.split("=")[1]);
+    else if (arg.startsWith("--label=")) flags.label = arg.split("=").slice(1).join("=");
+  }
+  return flags;
+}
+
+
+const CASES = [
+  { query: "como ver os usuários conectados no equipamento", expectedSourceIncludes: "Huawei", type: "positive" },
+  { query: "comando para exibir usuários de acesso", expectedSourceIncludes: "Huawei", type: "positive" },
+  { query: "como configurar Eth-Trunk", expectedSourceIncludes: "Huawei", type: "positive" },
+  { query: "como configurar OSPF no equipamento", expectedSourceIncludes: "Huawei", type: "positive" },
+  { query: "comando display access-user", expectedSourceIncludes: "Huawei", type: "positive" },
+  { query: "como resetar o roteador TP-Link", expectedSourceIncludes: "TL-WR940N", type: "positive" },
+  { query: "como instalar o roteador TP-Link pela primeira vez", expectedSourceIncludes: "TL-WR940N", type: "positive" },
+  { query: "qual o endereço padrão de acesso ao roteador", expectedSourceIncludes: "TL-WR940N", type: "positive" },
+  { query: "qual o valor do vale refeição da empresa", expectedSourceIncludes: null, type: "negative" },
+  { query: "quantos dias de férias tenho direito por ano", expectedSourceIncludes: null, type: "negative" }
+];
+
+function sourceMatches(hit, expectedIncludes) {
+  if (!expectedIncludes) return false;
+  return String(hit?.source ?? "").toLowerCase().includes(expectedIncludes.toLowerCase());
+}
+
+async function resolveEmbeddingQuery(originalQuery, flags) {
+  const searchQuery = flags.rewrite ? await rewriteQuery(originalQuery, []) : originalQuery;
+  if (!flags.hyde) return { searchQuery, embeddingQuery: searchQuery, embedRole: "query" };
+  const hydeText = await generateHydeDocument(searchQuery);
+  return hydeText
+    ? { searchQuery, embeddingQuery: hydeText, embedRole: "passage" }
+    : { searchQuery, embeddingQuery: searchQuery, embedRole: "query" };
+}
+
+async function runCase(testCase, flags) {
+  const { searchQuery, embeddingQuery, embedRole } = await resolveEmbeddingQuery(testCase.query, flags);
+  const hits = await searchDocs({
+    query: embeddingQuery,
+    topK: flags.topK,
+    embedRole,
+    minScore: flags.minScore
+  });
+
+  const top1 = hits[0] ?? null;
+  const top3 = hits.slice(0, 3);
+  const top1Match = testCase.type === "positive" ? sourceMatches(top1, testCase.expectedSourceIncludes) : null;
+  const top3Match = testCase.type === "positive" ? top3.some((h) => sourceMatches(h, testCase.expectedSourceIncludes)) : null;
+
+  return {
+    query: testCase.query,
+    searchQuery: searchQuery !== testCase.query ? searchQuery : "",
+    type: testCase.type,
+    top1_score: top1 ? top1.score.toFixed(4) : "-",
+    top1_source: top1?.source ?? "-",
+    top1_match: testCase.type === "positive" ? (top1Match ? "OK" : "MISS") : "-",
+    top3_match: testCase.type === "positive" ? (top3Match ? "OK" : "MISS") : "-"
+  };
+}
+
+async function main() {
+  const flags = parseArgs(process.argv.slice(2));
+
+  console.log(`\n=== evalRetrieval ${flags.label ? `[${flags.label}] ` : ""}===`);
+  console.log(
+    `modelo=${config.ollama.embeddingsModel} minScore=${flags.minScore} topK=${flags.topK} ` +
+    `rewrite=${flags.rewrite} hyde=${flags.hyde} colecao=${config.qdrant.collection}\n`
+  );
+
+  const rows = [];
+  for (const testCase of CASES) {
+    
+    rows.push(await runCase(testCase, flags));
+  }
+
+  console.table(rows);
+
+  const positives = rows.filter((r) => r.type === "positive");
+  const negatives = rows.filter((r) => r.type === "negative");
+
+  const positiveScores = positives.map((r) => Number(r.top1_score)).filter((n) => !Number.isNaN(n));
+  const negativeScores = negatives.map((r) => Number(r.top1_score)).filter((n) => !Number.isNaN(n));
+
+  const avg = (arr) => (arr.length ? arr.reduce((a, b) => a + b, 0) / arr.length : NaN);
+  const hitRate = (arr, key) => (arr.length ? arr.filter((r) => r[key] === "OK").length / arr.length : NaN);
+
+  console.log("\n--- resumo ---");
+  console.log(`positivos: score médio top-1 = ${avg(positiveScores).toFixed(4)} | hit-rate top-1 = ${(hitRate(positives, "top1_match") * 100).toFixed(0)}% | hit-rate top-3 = ${(hitRate(positives, "top3_match") * 100).toFixed(0)}%`);
+  console.log(`positivos: score mínimo top-1 (piso de segurança) = ${positiveScores.length ? Math.min(...positiveScores).toFixed(4) : "-"}`);
+  console.log(`negativos: score máximo top-1 (piso de ruído) = ${negativeScores.length ? Math.max(...negativeScores).toFixed(4) : "-"}`);
+  console.log("");
+}
+
+main()
+  .then(() => process.exit(0))
+  .catch((err) => {
+    console.error("[evalRetrieval] erro:", err);
+    process.exit(1);
+  });

+ 152 - 21
src/chat/chatChain.js

@@ -51,21 +51,137 @@ function hitsToSources(hits) {
   }));
 }
 
-async function rewriteQuery(query) {
+const REWRITE_SYSTEM_PROMPT = [
+  "Você é um motor de busca interno. Reescreva a pergunta abaixo de forma mais específica e completa",
+  "para busca em documentos internos de empresa, sempre em linguagem natural (nunca em código, SQL,",
+  "fórmulas ou pseudocódigo). Responda APENAS com a query reescrita, em uma frase, sem",
+  "explicações, sem aspas ou formatação extra."
+].join(" ");
+
+const REWRITE_SYSTEM_PROMPT_WITH_HISTORY = [
+  "Você é um motor de busca interno. Sua tarefa é transformar a ÚLTIMA pergunta do usuário",
+  "em uma query de busca autocontida e específica para buscar documentos internos da empresa,",
+  "sempre em linguagem natural (nunca em código, SQL, fórmulas ou pseudocódigo).",
+  "Use o HISTÓRICO DA CONVERSA abaixo apenas para resolver pronomes, referências e continuidade",
+  "de assunto (ex.: 'e o segundo caso?', 'isso', 'e sobre X?').",
+  "Responda APENAS com a query reescrita, em uma frase, sem explicações, sem aspas."
+].join("\n");
+
+export async function rewriteQuery(query, history = []) {
+  try {
+    const hasHistory = history.length > 0;
+    const messages = hasHistory
+      ? [
+          { role: "system", content: REWRITE_SYSTEM_PROMPT_WITH_HISTORY },
+          {
+            role: "system",
+            content: `HISTÓRICO DA CONVERSA:\n${history.map((m) => `${m.Role}: ${m.Content}`).join("\n")}`
+          },
+          { role: "user", content: query }
+        ]
+      : [
+          { role: "system", content: REWRITE_SYSTEM_PROMPT },
+          { role: "user", content: query }
+        ];
+    const { content } = await chatCompletion({ messages, options: { temperature: 0.1 } });
+    return content?.trim() || query;
+  } catch (err) {
+    console.warn("[chatChain] rewriteQuery falhou, usando query original:", err.message);
+    return query;
+  }
+}
+
+const RERANK_SYSTEM_PROMPT = [
+  "Você reordena trechos de documentos internos por relevância à pergunta do usuário.",
+  'Responda APENAS com um objeto JSON no formato {"order": [...]}, contendo TODOS os índices',
+  '(0-based) de cada trecho, em ordem decrescente de relevância. Exemplo: {"order": [2,0,3,1]}.',
+  "Não inclua explicações, texto extra ou markdown."
+].join("\n");
+
+const RERANK_SNIPPET_LENGTH = 300;
+
+function buildRerankPrompt(query, hits) {
+  const listing = hits
+    .map((h, i) => `[${i}] ${h.text.slice(0, RERANK_SNIPPET_LENGTH)}`)
+    .join("\n\n");
+  return `Pergunta: ${query}\n\nTrechos:\n${listing}`;
+}
+
+function parseRerankOrder(content, len) {
+  try {
+    const cleaned = String(content ?? "")
+      .trim()
+      .replace(/^```json\s*/i, "")
+      .replace(/^```\s*/, "")
+      .replace(/```$/, "")
+      .trim();
+    const parsed = JSON.parse(cleaned);
+    const arr = Array.isArray(parsed) ? parsed : Array.isArray(parsed?.order) ? parsed.order : null;
+    if (!arr) return null;
+
+    const seen = new Set();
+    const valid = [];
+    for (const v of arr) {
+      const idx = Number(v);
+      if (Number.isInteger(idx) && idx >= 0 && idx < len && !seen.has(idx)) {
+        seen.add(idx);
+        valid.push(idx);
+      }
+    }
+    if (valid.length === 0) return null;
+    for (let i = 0; i < len; i++) if (!seen.has(i)) valid.push(i);
+    return valid;
+  } catch {
+    return null;
+  }
+}
+
+async function rerankHits(query, hits) {
+  if (hits.length <= 1) return hits;
+  try {
+    const { content } = await chatCompletion({
+      messages: [
+        { role: "system", content: RERANK_SYSTEM_PROMPT },
+        { role: "user", content: buildRerankPrompt(query, hits) }
+      ],
+      format: "json",
+      options: { temperature: 0.1 }
+    });
+    const order = parseRerankOrder(content, hits.length);
+    if (!order) console.warn("[chatChain] rerankHits: resposta não parseável, mantendo ordem original");
+    const reordered = order ? order.map((i) => hits[i]) : hits;
+    return reordered.slice(0, config.rag.rerankTopN);
+  } catch (err) {
+    console.warn("[chatChain] rerankHits falhou, mantendo ordem original:", err.message);
+    return hits.slice(0, config.rag.rerankTopN);
+  }
+}
+
+const HYDE_SYSTEM_PROMPT = [
+  "Você é um redator técnico que escreve trechos de manuais internos de redes e",
+  "equipamentos (comandos Huawei, configuração de roteadores TP-Link, processos e",
+  "políticas internas da empresa, etc.). Dada a pergunta abaixo, escreva um parágrafo",
+  "curto (3 a 6 frases) no estilo de um trecho real de documentação técnica que",
+  "responda a essa pergunta — nomes de comandos, telas, campos, passos, quando fizer",
+  "sentido. Não inclua a pergunta, saudações 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");
+
+export async function generateHydeDocument(query) {
   try {
     const { content } = await chatCompletion({
       messages: [
-        {
-          role: "system",
-          content:
-            "Você é um motor de busca interno. Reescreva a pergunta abaixo de forma mais específica e completa para busca em documentos internos de empresa. Responda APENAS com a query reescrita, sem explicações ou formatação extra."
-        },
+        { role: "system", content: HYDE_SYSTEM_PROMPT },
         { role: "user", content: query }
-      ]
+      ],
+      options: { temperature: 0.2 }
     });
-    return content?.trim() || query;
-  } catch {
-    return query;
+    return content?.trim() || null;
+  } catch (err) {
+    console.warn("[chatChain] generateHydeDocument falhou, usando query normal para embedding:", err.message);
+    return null;
   }
 }
 
@@ -73,7 +189,8 @@ async function loadHistory(conversationId, userId) {
   if (!conversationId) return [];
   try {
     return await getRecentMessages(conversationId, userId, config.rag.historyLimit);
-  } catch {
+  } catch (err) {
+    console.warn("[chatChain] loadHistory falhou, seguindo sem histórico:", err.message);
     return [];
   }
 }
@@ -90,12 +207,20 @@ function buildDefaultOptions(override) {
     : defaults;
 }
 
+async function resolveEmbeddingQuery(searchQuery) {
+  if (!config.rag.hyde) return { embeddingQuery: searchQuery, embedRole: "query" };
+  const hydeText = await generateHydeDocument(searchQuery);
+  return hydeText
+    ? { embeddingQuery: hydeText, embedRole: "passage" }
+    : { embeddingQuery: searchQuery, embedRole: "query" };
+}
+
 export async function answerWithContext({ message, conversationId, userId, options }) {
-  const searchQuery = config.rag.queryRewrite ? await rewriteQuery(message) : message;
-  const [hits, history] = await Promise.all([
-    searchDocs({ query: searchQuery, topK: config.rag.topK }),
-    loadHistory(conversationId, userId)
-  ]);
+  const history = await loadHistory(conversationId, userId);
+  const searchQuery = config.rag.queryRewrite ? await rewriteQuery(message, history) : message;
+  const { embeddingQuery, embedRole } = await resolveEmbeddingQuery(searchQuery);
+  let hits = await searchDocs({ query: embeddingQuery, topK: config.rag.topK, embedRole });
+  if (config.rag.rerank && hits.length > 1) hits = await rerankHits(searchQuery, hits);
   const context = buildContextBlock(hits);
   const messages = buildMessages(message, context, history);
   const completion = await chatCompletion({ messages, options: buildDefaultOptions(options) });
@@ -107,13 +232,19 @@ export async function answerWithContext({ message, conversationId, userId, optio
 }
 
 export async function answerWithContextStream({ message, conversationId, userId, options, onChunk, onSources, onStatus, signal }) {
-  const searchQuery = config.rag.queryRewrite ? await rewriteQuery(message) : message;
   onStatus?.("buscando");
-  const [hits, history] = await Promise.all([
-    searchDocs({ query: searchQuery, topK: config.rag.topK }),
-    loadHistory(conversationId, userId)
-  ]);
+  const history = await loadHistory(conversationId, userId);
+  const searchQuery = config.rag.queryRewrite ? await rewriteQuery(message, history) : message;
+  if (config.rag.hyde) onStatus?.("gerando_hipotese");
+  const { embeddingQuery, embedRole } = await resolveEmbeddingQuery(searchQuery);
+  let hits = await searchDocs({ query: embeddingQuery, topK: config.rag.topK, embedRole });
   onStatus?.("encontrou", hits.length);
+
+  if (config.rag.rerank && hits.length > 1) {
+    onStatus?.("reordenando");
+    hits = await rerankHits(searchQuery, hits);
+  }
+
   const sources = hitsToSources(hits);
   onSources?.(sources);
   const context = buildContextBlock(hits);

+ 6 - 1
src/config/index.js

@@ -31,6 +31,8 @@ export const config = {
   ollama: {
     url: process.env.OLLAMA_URL ?? "http://localhost:11434",
     embeddingsModel: process.env.OLLAMA_EMBEDDINGS_MODEL ?? "nomic-embed-text",
+    embedQueryPrefix: process.env.OLLAMA_EMBED_QUERY_PREFIX ?? "",
+    embedPassagePrefix: process.env.OLLAMA_EMBED_PASSAGE_PREFIX ?? "",
     chatModel: process.env.OLLAMA_CHAT_MODEL ?? "llama3.1",
     visionModel: process.env.OLLAMA_VISION_MODEL ?? "llava:latest"
   },
@@ -40,7 +42,10 @@ export const config = {
     chunkSize: Number(process.env.RAG_CHUNK_SIZE ?? 900),
     chunkOverlap: Number(process.env.RAG_CHUNK_OVERLAP ?? 150),
     queryRewrite: process.env.RAG_QUERY_REWRITE === "true",
-    historyLimit: Number(process.env.RAG_HISTORY_LIMIT ?? 6)
+    historyLimit: Number(process.env.RAG_HISTORY_LIMIT ?? 6),
+    rerank: process.env.RAG_RERANK === "true",
+    rerankTopN: Number(process.env.RAG_RERANK_TOP_N ?? process.env.RAG_TOP_K ?? 8),
+    hyde: process.env.RAG_HYDE === "true"
   },
   llm: {
     temperature: Number(process.env.LLM_TEMPERATURE ?? 0.2),

+ 3 - 4
src/controllers/Chat.Controller.js

@@ -6,10 +6,9 @@ async function tryPersistMessages(conversationId, userId, userContent, result) {
   try {
     const owns = await verifyConversationOwner(conversationId, userId);
     if (!owns) return;
-    await Promise.all([
-      addMessage(conversationId, { role: "user", content: userContent }),
-      addMessage(conversationId, { role: "assistant", content: result.answer, sources: result.sources })
-    ]);
+
+    await addMessage(conversationId, { role: "user", content: userContent });
+    await addMessage(conversationId, { role: "assistant", content: result.answer, sources: result.sources });
   } catch (err) {
     console.error("[chat] falha ao salvar mensagem:", err);
   }

+ 6 - 2
src/controllers/Documents.Controller.js

@@ -5,7 +5,11 @@ export const DocumentsController = {
   Listar: async function (req, res, next) {
     try {
       const limit = z.coerce.number().int().positive().max(200).catch(50).parse(req.query.limit);
-      const offset = z.coerce.number().int().nonnegative().optional().parse(req.query.offset);
+     
+      const offset = z
+        .union([z.coerce.number().int().nonnegative(), z.string().min(1).max(64)])
+        .optional()
+        .parse(req.query.offset);
       const result = await listDocuments({ limit, offset });
       res.json(result);
     } catch (err) {
@@ -15,7 +19,7 @@ export const DocumentsController = {
 
   RemoverPorFonte: async function (req, res, next) {
     try {
-      const source = z.string().min(1).max(255).parse(decodeURIComponent(req.params.source));
+      const source = z.string().min(1).max(255).parse(req.params.source);
       await deleteDocumentsBySource(source);
       res.json({ ok: true });
     } catch (err) {

+ 6 - 2
src/controllers/Usuario.Controller.js

@@ -100,8 +100,12 @@ export const UsuarioController = {
       const usuario = await Usuario.query().findById(id);
       if (!usuario) throw new NotFoundError("Usuário não encontrado");
 
-      const senhaValida = bcrypt.compareSync(senhaAtual, String(usuario.Senha ?? ""));
-      if (!senhaValida) throw new AppError("Senha atual incorreta", 401);
+      
+      if (callerIsOwner) {
+        if (!senhaAtual) throw new BadRequestError("Senha atual é obrigatória");
+        const senhaValida = bcrypt.compareSync(senhaAtual, String(usuario.Senha ?? ""));
+        if (!senhaValida) throw new AppError("Senha atual incorreta", 401);
+      }
 
       const novaHash = bcrypt.hashSync(senhaNova, 10);
       await usuario.$query().patch({ Senha: novaHash });

+ 2 - 1
src/middleware/schemas/Usuario.Schema.js

@@ -19,6 +19,7 @@ export const atualizarSchema = z.object({
 });
 
 export const senhaSchema = z.object({
-  senhaAtual: z.string().min(1),
+ 
+  senhaAtual: z.string().min(1).optional(),
   senhaNova:  z.string().min(6).max(128)
 });

+ 12 - 1
src/services/collectionService.js

@@ -6,7 +6,18 @@ export async function ensureCollection({ vectorSize }) {
   const existing = await qdrant.getCollections();
   const has = (existing?.collections ?? []).some((c) => c.name === collectionName);
 
-  if (has) return;
+  if (has) {
+    const info = await qdrant.getCollection(collectionName);
+    const existingSize = info?.config?.params?.vectors?.size;
+    if (existingSize && existingSize !== vectorSize) {
+      const err = new Error(
+        `vector_size_mismatch:collection=${collectionName}:existing=${existingSize}:incoming=${vectorSize}`
+      );
+      err.statusCode = 409;
+      throw err;
+    }
+    return;
+  }
 
   await qdrant.createCollection(collectionName, {
     vectors: {

+ 3 - 0
src/services/conversationsService.js

@@ -28,6 +28,7 @@ export async function getConversationMessages(conversationId, userId) {
   const msgs = await Message.query()
     .where({ ConversationId: conversationId })
     .orderBy("SentAt", "asc")
+    .orderBy("Id", "asc")
     .select("Id", "Role", "Content", "Sources", "SentAt");
 
   return msgs;
@@ -46,6 +47,7 @@ export async function getRecentMessages(conversationId, userId, limit = 12) {
   const msgs = await Message.query()
     .where({ ConversationId: conversationId })
     .orderBy("SentAt", "desc")
+    .orderBy("Id", "desc")
     .limit(limit)
     .select("Role", "Content");
   return msgs.reverse();
@@ -82,6 +84,7 @@ export async function getConversationForExport(conversationId, userId) {
   const msgs = await Message.query()
     .where({ ConversationId: conversationId })
     .orderBy("SentAt", "asc")
+    .orderBy("Id", "asc")
     .select("Role", "Content", "SentAt");
 
   return { title: conv.Title, createdAt: conv.CreatedAt, messages: msgs };

+ 37 - 7
src/services/ingestService.js

@@ -2,6 +2,7 @@ import { config } from "../config/index.js";
 import { chunkText } from "./textChunker.js";
 import { embedTexts, visionExtractFromImage } from "./ollamaClient.js";
 import { ensureCollection } from "./collectionService.js";
+import { deleteDocumentsBySource } from "./documentsService.js";
 import { qdrant } from "./qdrantClient.js";
 import { isRefusal } from "../utils/isRefusal.js";
 import { createHash } from "node:crypto";
@@ -72,7 +73,7 @@ export async function ingestDocuments(documents) {
 
   if (allChunks.length === 0) return { upserted: 0 };
 
-  const vectors = await embedTexts(allChunks.map((c) => c.text));
+  const vectors = await embedTexts(allChunks.map((c) => c.text), { role: "passage" });
   const vectorSize = vectors[0]?.length ?? 0;
   if (!vectorSize) {
     const err = new Error("embeddings_empty");
@@ -82,6 +83,12 @@ export async function ingestDocuments(documents) {
 
   await ensureCollection({ vectorSize });
 
+  
+  const sources = [...new Set(documents.map((d) => d.source).filter(Boolean))];
+  for (const source of sources) {
+    await deleteDocumentsBySource(source);
+  }
+
   const points = allChunks.map((c, idx) => ({
     id: c.id,
     vector: vectors[idx],
@@ -113,7 +120,9 @@ function isPrivateUrl(urlStr) {
       /^10\./.test(hostname) ||
       /^192\.168\./.test(hostname) ||
       /^172\.(1[6-9]|2\d|3[01])\./.test(hostname) ||
+      /^169\.254\./.test(hostname) ||
       hostname === "0.0.0.0" ||
+      hostname.includes(":") ||
       hostname.endsWith(".local")
     );
   } catch {
@@ -121,7 +130,7 @@ function isPrivateUrl(urlStr) {
   }
 }
 
-export async function fetchUrlText(urlStr) {
+function assertPublicHttpsUrl(urlStr) {
   if (!urlStr.startsWith("https://")) {
     const err = new Error("url_must_be_https");
     err.statusCode = 400;
@@ -132,12 +141,33 @@ export async function fetchUrlText(urlStr) {
     err.statusCode = 400;
     throw err;
   }
+}
 
-  const res = await fetch(urlStr, {
-    headers: { "User-Agent": "Mozilla/5.0 star-oraculo/1.0" },
-    redirect: "follow",
-    signal: AbortSignal.timeout(15_000)
-  });
+const MAX_REDIRECTS = 5;
+
+export async function fetchUrlText(urlStr) {
+  
+  let currentUrl = urlStr;
+  let res;
+  for (let hop = 0; ; hop += 1) {
+    assertPublicHttpsUrl(currentUrl);
+
+    res = await fetch(currentUrl, {
+      headers: { "User-Agent": "Mozilla/5.0 star-oraculo/1.0" },
+      redirect: "manual",
+      signal: AbortSignal.timeout(15_000)
+    });
+
+    if (res.status < 300 || res.status >= 400) break;
+
+    const location = res.headers.get("location");
+    if (!location || hop >= MAX_REDIRECTS) {
+      const err = new Error(`url_fetch_error:${res.status}`);
+      err.statusCode = 502;
+      throw err;
+    }
+    currentUrl = new URL(location, currentUrl).toString();
+  }
 
   if (!res.ok) {
     const err = new Error(`url_fetch_error:${res.status}`);

+ 34 - 17
src/services/ollamaClient.js

@@ -34,30 +34,38 @@ async function ollamaFetch(path, body) {
   return res.json();
 }
 
-async function embedSingle(text) {
+function embedPrefixFor(role) {
+  if (role === "query") return config.ollama.embedQueryPrefix;
+  if (role === "passage") return config.ollama.embedPassagePrefix;
+  return "";
+}
+
+async function embedSingle(text, { role } = {}) {
+  const prefix = embedPrefixFor(role);
   const data = await ollamaFetch("/api/embeddings", {
     model: config.ollama.embeddingsModel,
-    prompt: text
+    prompt: `${prefix}${text}`
   });
   return data.embedding;
 }
 
-export async function embedTexts(texts) {
+export async function embedTexts(texts, { role } = {}) {
   const BATCH = 10;
   const results = [];
   for (let i = 0; i < texts.length; i += BATCH) {
     const batch = texts.slice(i, i + BATCH);
-    const embeddings = await Promise.all(batch.map(embedSingle));
+    const embeddings = await Promise.all(batch.map((text) => embedSingle(text, { role })));
     results.push(...embeddings);
   }
   return results;
 }
 
-export async function chatCompletion({ messages, options }) {
+export async function chatCompletion({ messages, options, format }) {
   const data = await ollamaFetch("/api/chat", {
     model: config.ollama.chatModel,
     messages,
     stream: false,
+    ...(format && { format }),
     ...(options && Object.keys(options).length > 0 && { options })
   });
 
@@ -89,24 +97,33 @@ export async function chatCompletionStream({ messages, onChunk, signal, options
   const reader = res.body.getReader();
   const decoder = new TextDecoder();
   let fullContent = "";
+  let buffer = "";
+
+  const processLine = (line) => {
+    if (!line.trim()) return;
+    try {
+      const parsed = JSON.parse(line);
+      const delta = parsed?.message?.content ?? "";
+      if (delta) {
+        fullContent += delta;
+        onChunk(delta);
+      }
+    } catch {}
+  };
 
   while (true) {
     const { done, value } = await reader.read();
     if (done) break;
-    const text = decoder.decode(value, { stream: true });
-    for (const line of text.split("\n")) {
-      if (!line.trim()) continue;
-      try {
-        const parsed = JSON.parse(line);
-        const delta = parsed?.message?.content ?? "";
-        if (delta) {
-          fullContent += delta;
-          onChunk(delta);
-        }
-      } catch {}
-    }
+    
+    buffer += decoder.decode(value, { stream: true });
+    const lines = buffer.split("\n");
+    buffer = lines.pop() ?? "";
+    for (const line of lines) processLine(line);
   }
 
+  buffer += decoder.decode();
+  processLine(buffer);
+
   return { content: fullContent };
 }
 

+ 16 - 9
src/services/searchService.js

@@ -2,17 +2,24 @@ import { config } from "../config/index.js";
 import { embedTexts } from "./ollamaClient.js";
 import { qdrant } from "./qdrantClient.js";
 
-export async function searchDocs({ query, topK }) {
+export async function searchDocs({ query, topK, embedRole = "query", minScore }) {
   const collectionName = config.qdrant.collection;
-  const [vector] = await embedTexts([query]);
+  const [vector] = await embedTexts([query], { role: embedRole });
 
-  const result = await qdrant.search(collectionName, {
-    vector,
-    limit: topK ?? config.rag.topK,
-    score_threshold: config.rag.minScore,
-    with_payload: true,
-    with_vector: false
-  });
+  let result;
+  try {
+    result = await qdrant.search(collectionName, {
+      vector,
+      limit: topK ?? config.rag.topK,
+      score_threshold: minScore ?? config.rag.minScore,
+      with_payload: true,
+      with_vector: false
+    });
+  } catch (err) {
+    const status = err?.status ?? err?.statusCode ?? err?.response?.status ?? null;
+    if (Number(status) === 404) return [];
+    throw err;
+  }
 
   return (result ?? [])
     .map((r) => ({

+ 39 - 7
src/services/textChunker.js

@@ -16,6 +16,28 @@ function splitBySentence(text, chunkSize) {
   return chunks;
 }
 
+function splitByLength(text, chunkSize) {
+  const chunks = [];
+  for (let i = 0; i < text.length; i += chunkSize) {
+    chunks.push(text.slice(i, i + chunkSize));
+  }
+  return chunks;
+}
+
+
+function splitOversized(text, chunkSize) {
+  const bySentence = splitBySentence(text, chunkSize);
+  const result = [];
+  for (const piece of bySentence) {
+    if (piece.length <= chunkSize) {
+      result.push(piece);
+    } else {
+      result.push(...splitByLength(piece, chunkSize));
+    }
+  }
+  return result;
+}
+
 function getOverlapTail(text, overlapSize) {
   if (text.length <= overlapSize) return text;
   const tail = text.slice(-overlapSize * 2);
@@ -34,7 +56,13 @@ export function chunkText(text, { chunkSize = 900, chunkOverlap = 150 }) {
 
   for (const para of paragraphs) {
     if (!current) {
-      current = para;
+      if (para.length > chunkSize) {
+        const pieces = splitOversized(para, chunkSize);
+        chunks.push(...pieces.slice(0, -1).map((p) => p.trim()).filter(Boolean));
+        current = pieces[pieces.length - 1] ?? "";
+      } else {
+        current = para;
+      }
       continue;
     }
 
@@ -46,11 +74,9 @@ export function chunkText(text, { chunkSize = 900, chunkOverlap = 150 }) {
       chunks.push(current.trim());
 
       if (para.length > chunkSize) {
-        const sentences = splitBySentence(para, chunkSize);
-        for (let i = 0; i < sentences.length - 1; i++) {
-          if (sentences[i].trim()) chunks.push(sentences[i].trim());
-        }
-        current = sentences[sentences.length - 1] ?? "";
+        const pieces = splitOversized(para, chunkSize);
+        chunks.push(...pieces.slice(0, -1).map((p) => p.trim()).filter(Boolean));
+        current = pieces[pieces.length - 1] ?? "";
       } else {
         const overlap = getOverlapTail(current, chunkOverlap);
         current = overlap ? overlap + "\n\n" + para : para;
@@ -58,7 +84,13 @@ export function chunkText(text, { chunkSize = 900, chunkOverlap = 150 }) {
     }
   }
 
-  if (current.trim()) chunks.push(current.trim());
+  if (current.trim()) {
+    if (current.length > chunkSize) {
+      chunks.push(...splitOversized(current, chunkSize).map((p) => p.trim()).filter(Boolean));
+    } else {
+      chunks.push(current.trim());
+    }
+  }
 
   return chunks.filter(Boolean);
 }