Chuyển đến nội dung chính

Lesson 11: RAG Engine — Retrieval-Augmented Generation

Integrate RAG into Agent: retrieve, rerank, inject context. Collections management, hybrid search, web ingestion. Analytics & quality metrics.

🧠 AI & ML — Lesson 10 Lesson 11: RAG Engine — Retrieval-Augmented Generation

Building AI Agent Platform from Zero — Real battle with xClaw

Part 3: RAG Pipeline

xdev.asia

Introduction

RAG Engine is the orchestrator that connects the Document Processor, Embedding Provider, and Vector Store into a complete pipeline. This article implements the full RAG Engine with ingestion, retrieval, and reranking.


1. RAG Engine Class

// packages/core/src/rag/rag-engine.ts
export class RagEngine {
  private docProcessor: DocumentProcessor;
  private embedding: EmbeddingProvider;
  private vectorStore: VectorStore;
  private webCrawler: WebCrawler;

  constructor(
    docProcessor: DocumentProcessor,
    embedding: EmbeddingProvider,
    vectorStore: VectorStore,
  ) {
    this.docProcessor = docProcessor;
    this.embedding = embedding;
    this.vectorStore = vectorStore;
    this.webCrawler = new WebCrawler();
  }

  // Ingest plain text
  async ingestText(
    text: string,
    collectionId: string,
    metadata?: Record<string, unknown>,
  ): Promise<{ chunksCreated: number }> {
    const chunks = await this.docProcessor.process(text, 'text/plain', {
      source: 'direct-text',
      sourceType: 'text',
      chunkSize: 1000,
      chunkOverlap: 200,
    });

    const vectors = await this.embedding.embed(chunks.map(c => c.content));
    const documents: VectorDocument[] = chunks.map((chunk, i) => ({
      id: chunk.id,
      vector: vectors[i],
      content: chunk.content,
      metadata: { ...chunk.metadata, ...metadata, collectionId },
    }));

    await this.vectorStore.upsert(documents);
    return { chunksCreated: chunks.length };
  }

  // Ingest from URL
  async ingestUrl(
    url: string,
    collectionId: string,
  ): Promise<{ chunksCreated: number }> {
    const text = await this.webCrawler.crawl(url);
    const chunks = await this.docProcessor.process(text, 'text/html', {
      source: url,
      sourceType: 'url',
      chunkSize: 1000,
      chunkOverlap: 200,
    });

    const vectors = await this.embedding.embed(chunks.map(c => c.content));
    const documents: VectorDocument[] = chunks.map((chunk, i) => ({
      id: chunk.id,
      vector: vectors[i],
      content: chunk.content,
      metadata: { ...chunk.metadata, collectionId },
    }));

    await this.vectorStore.upsert(documents);
    return { chunksCreated: chunks.length };
  }
}

2. Retrieval & Reranking

// Retrieve relevant documents
async retrieve(
  query: string,
  collectionId: string,
  options: { topK?: number; minScore?: number } = {},
): Promise<SearchResult[]> {
  const [queryVector] = await this.embedding.embed([query]);

  const results = await this.vectorStore.search(queryVector, {
    collection: collectionId,
    topK: options.topK || 5,
    minScore: options.minScore || 0.7,
    filter: { 'metadata.collectionId': collectionId },
  });

  return results;
}

// Search with reranking — better accuracy
async searchWithReranking(
  query: string,
  collectionId: string,
  options: { topK?: number; rerankTopK?: number } = {},
): Promise<SearchResult[]> {
  // Step 1: Retrieve more candidates
  const candidates = await this.retrieve(query, collectionId, {
    topK: options.rerankTopK || 20,
    minScore: 0.5, // Lower threshold for candidates
  });

  if (candidates.length === 0) return [];

  // Step 2: Rerank using cross-encoder or LLM
  const reranked = await this.rerankResults(query, candidates);

  // Step 3: Return top K
  return reranked.slice(0, options.topK || 5);
}

private async rerankResults(
  query: string,
  candidates: SearchResult[],
): Promise<SearchResult[]> {
  // Simple reranking: score by keyword overlap + vector score
  return candidates
    .map(result => {
      const queryTerms = query.toLowerCase().split(/\s+/);
      const contentLower = result.content.toLowerCase();
      const keywordScore = queryTerms.filter(t => contentLower.includes(t)).length / queryTerms.length;
      const combinedScore = result.score * 0.7 + keywordScore * 0.3;

      return { ...result, score: combinedScore };
    })
    .sort((a, b) => b.score - a.score);
}

3. Collections Management

// Collection CRUD
async createCollection(name: string, tenantId: string): Promise<string> {
  const collectionId = crypto.randomUUID();

  await this.vectorStore.createCollection(collectionId);

  // Save metadata to MongoDB
  await this.db.collection('rag_collections').insertOne({
    _id: collectionId,
    name,
    tenantId,
    documentCount: 0,
    chunkCount: 0,
    createdAt: new Date(),
  });

  return collectionId;
}

async listCollections(tenantId: string) {
  return this.db.collection('rag_collections')
    .find({ tenantId })
    .sort({ createdAt: -1 })
    .toArray();
}

async deleteCollection(collectionId: string) {
  await this.vectorStore.delete(
    await this.getChunkIds(collectionId),
  );
  await this.db.collection('rag_collections').deleteOne({ _id: collectionId });
}

4. RAG in Agent Flow

// Khi Agent nhận được message, RAG context được inject:
const ragContext = await ragEngine.searchWithReranking(
  userMessage,
  activeCollectionId,
  { topK: 5 },
);

if (ragContext.length > 0) {
  messages.push({
    role: 'system',
    content: [
      '### Relevant Knowledge Base Context',
      '',
      ...ragContext.map((r, i) =>
        `**[${i + 1}]** (score: ${r.score.toFixed(2)})\n${r.content}`
      ),
      '',
      'Use the above context to answer the user\'s question. Cite sources using [1], [2], etc.',
    ].join('\n'),
  });
}

5. Summary

ComponentsRole
DocumentProcessorParse & chunk documents
EmbeddingProviderText → vectors
VectorStoreStore & search vectors
RAG EngineOrchestrate pipeline
RerankingImprove retrieval accuracy

Next article: Workflow Engine — Visual automation for AI tasks.