// 레드팀 r2-0: 로컬 RAG 오케스트레이터를 가짜 포트로 검증한다(실제 fetch·SQLite 없음). // - 임베딩 모델이 없거나 전부 실패하면 색인 실패를 알리고 진행 표시를 닫는다. // - 색인 중에 질의가 끝나도 상태가 idle로 덮이지 않는다. // - 추출·청킹·순위·Ollama 어댑터는 순수 모듈로 따로 검증한다. import fs from 'fs' import os from 'os' import path from 'path' import { deflateRawSync, deflateSync } from 'zlib' import { afterEach, describe, expect, it } from 'vitest' import { IPC_CHANNELS } from '@d3ro/core/ipc-channels' import { D3ROError, ErrorCode } from '@d3ro/core/errors' import type { RAGDocument } from '@d3ro/core/types' import { RAGService, defaultRAGServiceDeps, resetRAGServiceForTests, type RAGIndexFailure, type RAGServiceDeps, } from '../../../src/main/services/RAGService' import type { ChunkStore, StoredChunk } from '../../../src/main/services/rag/chunk-store' import type { EmbeddingPort } from '../../../src/main/services/rag/embedding-port' import { OllamaEmbeddingAdapter, isSameOllamaModel } from '../../../src/main/services/rag/embedding-port' import { CHUNK_SIZE, chunkText, extractBinaryDocumentText, fileTypeForExtension, } from '../../../src/main/services/rag/document-text' import { extractPdfText } from '../../../src/main/services/rag/pdf-text' import { buildAnswerSystemPrompt, cosineSimilarity, rankChunks } from '../../../src/main/services/rag/retrieval' class MemoryChunkStore implements ChunkStore { docs = new Map() chunks: StoredChunk[] = [] listDocuments(): RAGDocument[] { return [...this.docs.values()] } getDocument(id: string): RAGDocument | null { return this.docs.get(id) ?? null } hasDocument(id: string): boolean { return this.docs.has(id) } insertDocument(doc: RAGDocument): void { this.docs.set(doc.id, { ...doc }) } updateDocument(id: string, patch: Partial>): void { const doc = this.docs.get(id) if (doc) Object.assign(doc, patch) } removeDocument(id: string): boolean { this.chunks = this.chunks.filter((c) => c.documentId !== id) return this.docs.delete(id) } replaceChunks(documentId: string, chunks: readonly string[]): void { this.chunks = this.chunks.filter((c) => c.documentId !== documentId) chunks.forEach((content, chunkIndex) => this.chunks.push({ id: crypto.randomUUID(), documentId, content, embedding: '', chunkIndex }) ) } listChunks(documentId: string): StoredChunk[] { return this.chunks.filter((c) => c.documentId === documentId).sort((a, b) => a.chunkIndex - b.chunkIndex) } setEmbedding(chunkId: string, embedding: readonly number[]): void { const chunk = this.chunks.find((c) => c.id === chunkId) if (chunk) chunk.embedding = JSON.stringify(embedding) } listEmbeddedChunks(): StoredChunk[] { return this.chunks.filter((c) => c.embedding.length > 0) } counts(): { documentCount: number; totalChunks: number } { return { documentCount: this.docs.size, totalChunks: this.chunks.length } } } class FakeEmbedder implements EmbeddingPort { readonly model = 'fake-embed' modelInstalled = true failWhen: (text: string) => boolean = () => false embedCalls = 0 gate: Promise | null = null async ensureModel(): Promise { if (!this.modelInstalled) { throw new D3ROError(ErrorCode.RAGEmbeddingFailed, 'Embedding model "fake-embed" is not installed', { reason: 'model_missing', }) } } async embed(text: string): Promise { this.embedCalls++ if (this.gate) await this.gate if (this.failWhen(text)) throw new D3ROError(ErrorCode.RAGEmbeddingFailed, 'Embed API error: 500') return [text.length, 1] } } interface Harness { service: RAGService store: MemoryChunkStore embedder: FakeEmbedder events: Array<{ channel: string; data: unknown }> failures: RAGIndexFailure[] pushed: string[] } function harness(overrides: Partial = {}): Harness { const store = new MemoryChunkStore() const embedder = new FakeEmbedder() const events: Array<{ channel: string; data: unknown }> = [] const pushed: string[] = [] const service = new RAGService({ store, embedder, answerer: () => ({ generate: async () => ({ text: ' answer ' }) }), sync: () => ({ pushOne: (_entity, id) => pushed.push(`up:${id}`), pushDelete: (_entity, id) => pushed.push(`del:${id}`), }), notify: (channel, data) => events.push({ channel, data }), readFile: (filePath) => fs.promises.readFile(filePath), fileExists: (filePath) => fs.existsSync(filePath), assertLicensed: async () => undefined, yieldMs: 0, ...overrides, }) const failures: RAGIndexFailure[] = [] service.on('index-failed', (f) => failures.push(f)) return { service, store, embedder, events, failures, pushed } } function writeTemp(name: string, body: string): string { const p = path.join(os.tmpdir(), `d3ro-rag-r2-${crypto.randomUUID()}-${name}`) fs.writeFileSync(p, body, 'utf-8') return p } async function settle(): Promise { for (let i = 0; i < 20; i++) await new Promise((r) => setTimeout(r, 0)) } function channels(h: Harness): string[] { return h.events.map((e) => e.channel) } const LONG_TEXT = 'Knowledge base paragraph about quarterly planning and milestones. '.repeat(30) afterEach(() => { resetRAGServiceForTests() }) describe('RAGService 색인 실패 알림', () => { it('임베딩 모델이 없으면 청크마다 기다리지 않고 실패를 알리고 진행 표시를 닫는다', async () => { const h = harness() h.embedder.modelInstalled = false const doc = await h.service.addDocument(writeTemp('a.txt', LONG_TEXT)) await settle() expect(h.embedder.embedCalls).toBe(0) expect(h.store.getDocument(doc.id)?.indexed).toBe(false) expect(h.failures).toEqual([ expect.objectContaining({ documentId: doc.id, code: ErrorCode.RAGEmbeddingFailed }), ]) expect(channels(h)).toEqual([IPC_CHANNELS.RAG.INDEX_FAILED, IPC_CHANNELS.RAG.INDEX_COMPLETE]) // 원문은 남아 재색인·동기화가 가능하다 expect(h.store.listChunks(doc.id).length).toBeGreaterThan(0) expect(h.pushed).toEqual([`up:${doc.id}`]) expect(h.service.state).toBe('idle') }) it('모든 청크 임베딩이 실패하면 indexed=false 로 두고 실패 + 실행 종료를 알린다', async () => { const h = harness() h.embedder.failWhen = () => true const doc = await h.service.addDocument(writeTemp('b.txt', LONG_TEXT)) await settle() expect(h.store.getDocument(doc.id)?.indexed).toBe(false) expect(h.failures).toHaveLength(1) const tail = channels(h).slice(-2) expect(tail).toEqual([IPC_CHANNELS.RAG.INDEX_FAILED, IPC_CHANNELS.RAG.INDEX_COMPLETE]) expect(channels(h)).toContain(IPC_CHANNELS.RAG.INDEX_PROGRESS) }) it('일부만 성공해도 색인됨으로 표시하고 실패는 알리지 않는다', async () => { const h = harness() let n = 0 h.embedder.failWhen = () => n++ % 2 === 0 const doc = await h.service.addDocument(writeTemp('c.txt', LONG_TEXT)) await settle() const stored = h.store.getDocument(doc.id) expect(stored?.indexed).toBe(true) expect(stored?.chunkCount).toBe(h.store.listChunks(doc.id).length) expect(h.failures).toEqual([]) expect(channels(h).at(-1)).toBe(IPC_CHANNELS.RAG.INDEX_COMPLETE) expect(channels(h)).not.toContain(IPC_CHANNELS.RAG.INDEX_FAILED) }) it('동기화로 받은 문서도 모델이 없으면 실패를 알린다', async () => { const h = harness() h.embedder.modelInstalled = false const applied = h.service.applyRemoteDocument({ id: crypto.randomUUID(), fileName: 'phone.txt', fileType: 'txt', chunks: ['first', 'second'], addedAt: 1, }) await settle() expect(applied).toBe(true) expect(h.failures).toHaveLength(1) }) it('재색인이 실패하면 호출자에게 RAGEmbeddingFailed 로 돌려준다', async () => { const h = harness() const id = crypto.randomUUID() h.service.applyRemoteDocument({ id, fileName: 'x.txt', fileType: 'txt', chunks: ['alpha', 'beta'], addedAt: 1 }) await settle() expect(h.store.getDocument(id)?.indexed).toBe(true) h.embedder.modelInstalled = false await expect(h.service.reindex(id)).rejects.toMatchObject({ code: ErrorCode.RAGEmbeddingFailed }) expect(h.store.getDocument(id)?.indexed).toBe(false) }) }) describe('RAGService 상태', () => { it('색인 중에 질의가 끝나도 상태는 indexing 으로 남는다', async () => { const h = harness() const seeded = crypto.randomUUID() h.service.applyRemoteDocument({ id: seeded, fileName: 's.txt', fileType: 'txt', chunks: ['seed chunk'], addedAt: 1 }) await settle() let release: () => void = () => undefined h.embedder.gate = new Promise((r) => { release = r }) h.service.applyRemoteDocument({ id: crypto.randomUUID(), fileName: 'slow.txt', fileType: 'txt', chunks: ['slow'], addedAt: 2 }) await settle() expect(h.service.state).toBe('indexing') const query = h.service.query('seed?') await settle() expect(h.service.state).toBe('indexing') release() h.embedder.gate = null const result = await query expect(result.answer).toBe('answer') await settle() expect(h.service.state).toBe('idle') }) it('색인된 청크가 없으면 질의는 RAGQueryFailed 다', async () => { const h = harness() await expect(h.service.query('anything')).rejects.toMatchObject({ code: ErrorCode.RAGQueryFailed }) expect(h.service.state).toBe('idle') }) it('resetRAGServiceForTests 로 포트를 주입할 수 있다', () => { const store = new MemoryChunkStore() resetRAGServiceForTests({ store }) expect(defaultRAGServiceDeps().yieldMs).toBeGreaterThanOrEqual(0) }) }) describe('문서 텍스트 (순수 함수)', () => { it('확장자 → 종류', () => { expect(fileTypeForExtension('.pdf')).toBe('pdf') expect(fileTypeForExtension('.bin')).toBeNull() }) it('청킹은 CHUNK_SIZE 창으로 겹쳐 자르고 짧은 조각은 버린다', () => { const chunks = chunkText('a'.repeat(CHUNK_SIZE * 2)) expect(chunks.length).toBe(3) expect(chunks.every((c) => c.length <= CHUNK_SIZE)).toBe(true) expect(chunkText('short')).toEqual([]) }) it('압축된 PDF 스트림에서 Tj/TJ 텍스트를 뽑는다', () => { const content = 'BT (Hello RAG world) Tj [(second) 120 (part)] TJ ET' const stream = deflateSync(Buffer.from(content, 'binary')) const pdf = Buffer.concat([ Buffer.from('%PDF-1.4\n1 0 obj << /Filter /FlateDecode >>\nstream\n', 'binary'), stream, Buffer.from('\nendstream\nendobj\n', 'binary'), ]) expect(extractPdfText(pdf)).toBe('Hello RAG world second part') expect(extractBinaryDocumentText(pdf, 'pdf')).toBe('Hello RAG world second part') expect(() => extractBinaryDocumentText(Buffer.from('%PDF-1.4 nothing'), 'pdf')).toThrow() }) it('DOCX 는 ZIP 안의 document.xml 에서 문단을 뽑는다', () => { const xml = Buffer.from('Docx body text') const name = Buffer.from('word/document.xml') const data = deflateRawSync(xml) const local = Buffer.alloc(30) local.writeUInt32LE(0x04034b50, 0) local.writeUInt16LE(8, 8) local.writeUInt32LE(data.length, 18) local.writeUInt32LE(xml.length, 22) local.writeUInt16LE(name.length, 26) const central = Buffer.alloc(46) central.writeUInt32LE(0x02014b50, 0) central.writeUInt16LE(8, 10) central.writeUInt32LE(data.length, 20) central.writeUInt32LE(xml.length, 24) central.writeUInt16LE(name.length, 28) central.writeUInt32LE(0, 42) const centralOffset = local.length + name.length + data.length const eocd = Buffer.alloc(22) eocd.writeUInt32LE(0x06054b50, 0) eocd.writeUInt16LE(1, 8) eocd.writeUInt16LE(1, 10) eocd.writeUInt32LE(central.length + name.length, 12) eocd.writeUInt32LE(centralOffset, 16) const zip = Buffer.concat([local, name, data, central, name, eocd]) expect(extractBinaryDocumentText(zip, 'docx')).toBe('Docx body text') }) }) describe('검색 (순수 함수)', () => { it('코사인 유사도 순으로 topK 를 고른다', () => { expect(cosineSimilarity([1, 0], [1, 0])).toBeCloseTo(1) expect(cosineSimilarity([1, 0], [1, 0, 0])).toBe(0) const ranked = rankChunks( [1, 0], [ { documentId: 'd1', content: 'far', embedding: JSON.stringify([0, 1]) }, { documentId: 'd2', content: 'near', embedding: JSON.stringify([1, 0.1]) }, ], new Map([['d2', 'near.txt']]), 1 ) expect(ranked).toEqual([expect.objectContaining({ content: 'near', fileName: 'near.txt' })]) expect(buildAnswerSystemPrompt(ranked)).toContain('[1] (near.txt)\nnear') }) }) describe('OllamaEmbeddingAdapter', () => { function jsonResponse(body: unknown, status = 200): Response { return new Response(JSON.stringify(body), { status, headers: { 'Content-Type': 'application/json' } }) } it('설치된 모델이면 통과하고, 없으면 model_missing 으로 실패한다', async () => { const calls: string[] = [] const installed = new OllamaEmbeddingAdapter({ serverUrl: () => 'http://127.0.0.1:11434', fetchFn: async (url) => { calls.push(url) return jsonResponse({ models: [{ name: 'gemma4:e4b' }, { name: 'nomic-embed-text:latest' }] }) }, }) await expect(installed.ensureModel()).resolves.toBeUndefined() expect(calls).toEqual(['http://127.0.0.1:11434/api/tags']) const missing = new OllamaEmbeddingAdapter({ serverUrl: () => 'http://127.0.0.1:11434', fetchFn: async () => jsonResponse({ models: [{ name: 'gemma4:e4b' }] }), }) await expect(missing.ensureModel()).rejects.toMatchObject({ code: ErrorCode.RAGEmbeddingFailed, details: { reason: 'model_missing' }, }) }) it('서버에 닿지 않으면 server_unreachable 로 실패한다', async () => { const adapter = new OllamaEmbeddingAdapter({ serverUrl: () => 'http://127.0.0.1:1', fetchFn: async () => { throw new TypeError('fetch failed') }, }) await expect(adapter.ensureModel()).rejects.toMatchObject({ details: { reason: 'server_unreachable' } }) }) it('embed 는 첫 벡터를 돌려주고 HTTP 오류는 RAGEmbeddingFailed 다', async () => { const ok = new OllamaEmbeddingAdapter({ serverUrl: () => 'http://x', fetchFn: async () => jsonResponse({ embeddings: [[0.1, 0.2]] }), }) expect(await ok.embed('hi')).toEqual([0.1, 0.2]) const notFound = new OllamaEmbeddingAdapter({ serverUrl: () => 'http://x', fetchFn: async () => jsonResponse({ error: 'model not found' }, 404), }) await expect(notFound.embed('hi')).rejects.toMatchObject({ code: ErrorCode.RAGEmbeddingFailed }) }) it('모델 이름의 :latest 태그를 같은 모델로 본다', () => { expect(isSameOllamaModel('nomic-embed-text:latest', 'nomic-embed-text')).toBe(true) expect(isSameOllamaModel('nomic-embed-text:v1.5', 'nomic-embed-text')).toBe(false) }) })