diff --git a/apps/desktop/src/main/services/RAGService.ts b/apps/desktop/src/main/services/RAGService.ts index b075f4d..9910463 100644 --- a/apps/desktop/src/main/services/RAGService.ts +++ b/apps/desktop/src/main/services/RAGService.ts @@ -1,100 +1,158 @@ // src/main/services/RAGService.ts -// Phase 13.2: 로컬 RAG — 문서 임베딩 + 코사인 유사도 검색 + LLM 컨텍스트 주입 -// Ollama nomic-embed-text 모델, SQLite에 JSON 직렬화 벡터 저장. +// Phase 13.2: 로컬 RAG — 문서 임베딩 + 코사인 유사도 검색 + LLM 컨텍스트 주입. +// +// 이 서비스는 오케스트레이션만 한다. 세부 책임은 포트/순수 모듈로 나뉜다. +// - rag/document-text.ts : 원문 추출(PDF/DOCX)·청킹 (순수 함수) +// - rag/embedding-port.ts: 임베딩 포트 + Ollama 어댑터(모델 존재 확인 포함) +// - rag/chunk-store.ts : 문서·청크 저장소 포트 + SQLite 구현 +// - rag/retrieval.ts : 유사도 순위·답변 프롬프트 (순수 함수) import { EventEmitter } from 'events' import path from 'path' import fs from 'fs' -import { asc, eq } from 'drizzle-orm' import { getLogger } from './LoggerService' import { getLlmGateway } from './llm/LlmGateway' -import { extractDocxText } from './rag/docx-text' import { getOllamaServerUrl } from './LocalLLMService' -import { getDatabase } from '../db' -import { ragDocuments, ragChunks } from '../db/schema' import { getMainWindow } from '../windows/WindowManager' import { getCloudSyncService } from './CloudSyncService' import { IPC_CHANNELS } from '@d3ro/core/ipc-channels' import { D3ROError, ErrorCode } from '@d3ro/core/errors' -import type { - RAGDocument, - RAGQueryResult, - RAGState, - RAGStateInfo, - RAGIndexProgress, -} from '@d3ro/core/types' +import type { RAGDocument, RAGQueryResult, RAGState, RAGStateInfo, RAGIndexProgress } from '@d3ro/core/types' +import { + chunkText, + documentParseError, + extractBinaryDocumentText, + fileTypeForExtension, + isBinaryDocumentType, + isPlainTextType, + unsupportedFormatError, + type RagFileType, +} from './rag/document-text' +import { OllamaEmbeddingAdapter, type EmbeddingPort } from './rag/embedding-port' +import { SqliteChunkStore, type ChunkStore } from './rag/chunk-store' +import { buildAnswerSystemPrompt, rankChunks } from './rag/retrieval' const logger = getLogger('RAGService') -/** 임베딩 모델 */ -const EMBED_MODEL = 'nomic-embed-text' -/** 청크 크기 (자) */ -const CHUNK_SIZE = 500 -/** 청크 오버랩 (자) */ -const CHUNK_OVERLAP = 50 /** 검색 기본 topK */ const DEFAULT_TOP_K = 5 +/** 청크마다 이벤트 루프를 양보하는 시간(ms) — UI 블로킹 방지 */ +const INDEX_YIELD_MS = 10 -class RAGService extends EventEmitter { - private _state: RAGState = 'idle' +/** 색인 실패 알림 (서비스 이벤트 'index-failed' 와 IPC rag:indexFailed 의 내용) */ +export interface RAGIndexFailure { + documentId: string + fileName: string + code: ErrorCode + message: string +} + +/** 답변 생성기 — 사용자가 고른 LLM 백엔드(LlmGateway)의 최소 면 */ +export interface RagAnswerGenerator { + generate(text: string, options?: { systemPrompt?: string }): Promise<{ text: string }> +} + +/** 로컬 변경을 다른 기기로 알리는 면 (CloudSyncService의 최소 면) */ +export interface RagSyncSink { + pushOne(entity: 'knowledge_documents', id: string): void + pushDelete(entity: 'knowledge_documents', id: string): void +} + +export interface RAGServiceDeps { + store: ChunkStore + embedder: EmbeddingPort + /** 호출 때마다 현재 게이트웨이를 받는다(백엔드 전환을 따른다) */ + answerer: () => RagAnswerGenerator + sync: () => RagSyncSink + /** 렌더러로 이벤트 전송 */ + notify: (channel: string, data: unknown) => void + readFile: (filePath: string) => Promise + fileExists: (filePath: string) => boolean + /** Pro+ 라이선스 확인. 허용되지 않으면 D3ROError를 던진다. */ + assertLicensed: () => Promise + /** 청크 사이 양보 시간(ms) */ + yieldMs: number +} + +interface RAGServiceEvents { + 'index-failed': (failure: RAGIndexFailure) => void +} + +function errorText(err: unknown): string { + return err instanceof Error ? err.message : String(err) +} + +function sendToMainWindow(channel: string, data: unknown): void { + try { + const mainWindow = getMainWindow() + if (mainWindow && !mainWindow.isDestroyed()) mainWindow.webContents.send(channel, data) + } catch { + // 창이 없거나 닫히는 중 — 알림만 건너뛴다 + } +} + +async function assertLocalRagLicense(): Promise { + try { + const { getLicenseService } = await import('./LicenseService') + const { Feature } = await import('@d3ro/core/types') + const license = getLicenseService() + const access = license.canUse(Feature.LOCAL_RAG) + if (!access.allowed) { + license.promptUpgrade(Feature.LOCAL_RAG, 'tier_required') + throw new D3ROError(ErrorCode.FeatureNotAvailable, 'Pro+ required for Local RAG') + } + } catch (err) { + if (err instanceof D3ROError) throw err + } +} + +export function defaultRAGServiceDeps(): RAGServiceDeps { + return { + store: new SqliteChunkStore(), + embedder: new OllamaEmbeddingAdapter({ serverUrl: getOllamaServerUrl }), + answerer: () => getLlmGateway(), + sync: () => getCloudSyncService(), + notify: sendToMainWindow, + readFile: (filePath) => fs.promises.readFile(filePath), + fileExists: (filePath) => fs.existsSync(filePath), + assertLicensed: assertLocalRagLicense, + yieldMs: INDEX_YIELD_MS, + } +} + +export class RAGService extends EventEmitter { + private readonly deps: RAGServiceDeps + /** 진행 중인 색인·질의 수. 상태는 여기서 파생한다 — 겹쳐 실행돼도 서로의 상태를 덮지 않는다. */ + private activeIndexing = 0 + private activeQueries = 0 + + constructor(deps: RAGServiceDeps = defaultRAGServiceDeps()) { + super() + this.deps = deps + } get state(): RAGState { - return this._state + if (this.activeIndexing > 0) return 'indexing' + if (this.activeQueries > 0) return 'querying' + return 'idle' } getStateInfo(): RAGStateInfo { - const db = getDatabase() - const docs = db.select().from(ragDocuments).all() - const chunks = db.select().from(ragChunks).all() - return { - state: this._state, - documentCount: docs.length, - totalChunks: chunks.length, - } + return { state: this.state, ...this.deps.store.counts() } } getDocuments(): RAGDocument[] { - const db = getDatabase() - const rows = db.select().from(ragDocuments).all() - return rows.map((r) => ({ - id: r.id, - fileName: r.fileName, - filePath: r.filePath, - fileType: r.fileType as RAGDocument['fileType'], - chunkCount: r.chunkCount, - indexed: r.indexed, - indexedAt: r.indexedAt, - addedAt: r.addedAt, - })) + return this.deps.store.listDocuments() } /** * 문서 추가 + 인덱싱 (청킹 → 임베딩 → DB 저장) */ async addDocument(filePath: string): Promise { - // 라이센스 체크 - try { - const { getLicenseService } = await import('./LicenseService') - const { Feature } = await import('@d3ro/core/types') - const license = getLicenseService() - const access = license.canUse(Feature.LOCAL_RAG) - if (!access.allowed) { - license.promptUpgrade(Feature.LOCAL_RAG, 'tier_required') - throw new D3ROError(ErrorCode.FeatureNotAvailable, 'Pro+ required for Local RAG') - } - } catch (err) { - if (err instanceof D3ROError) throw err - } + await this.deps.assertLicensed() const ext = path.extname(filePath).toLowerCase() - const supportedTypes: Record = { - '.txt': 'txt', - '.md': 'md', - '.pdf': 'pdf', - '.docx': 'docx', - } - - const fileType = supportedTypes[ext] + const fileType = fileTypeForExtension(ext) if (!fileType) { throw new D3ROError(ErrorCode.RAGUnsupportedFormat, `Unsupported format: ${ext}. Supported: .txt, .md, .pdf, .docx`) } @@ -108,7 +166,7 @@ class RAGService extends EventEmitter { content = await this._extractText(filePath, fileType) } catch (err) { if (err instanceof D3ROError) throw err - throw new D3ROError(ErrorCode.RAGIndexingFailed, `Text extraction failed: ${err instanceof Error ? err.message : String(err)}`) + throw new D3ROError(ErrorCode.RAGIndexingFailed, `Text extraction failed: ${errorText(err)}`) } if (!content || content.trim().length < 20) { @@ -117,35 +175,12 @@ class RAGService extends EventEmitter { logger.info(`RAG text extracted: ${fileName} (${content.length} chars)`) - // 청킹 - const chunks = this._chunkText(content) + const chunks = chunkText(content) if (chunks.length === 0) { throw new D3ROError(ErrorCode.RAGIndexingFailed, 'Document produced no valid text chunks') } - // DB에 문서 레코드 삽입 - const db = getDatabase() - db.insert(ragDocuments).values({ - id: docId, - fileName, - filePath, - fileType, - chunkCount: chunks.length, - indexed: false, - indexedAt: null, - addedAt: Date.now(), - }).run() - - // 청크 원문을 먼저 저장한다 — 임베딩이 실패해도 원문은 남아 재색인·기기 간 동기화가 가능하다. - this._storeChunks(docId, chunks) - getCloudSyncService().pushOne('knowledge_documents', docId) - - // 비동기 인덱싱 (임베딩 생성) - this._embedStoredChunks(docId, fileName).catch((err) => { - logger.error(`Indexing failed for ${fileName}:`, err) - }) - - return { + const doc: RAGDocument = { id: docId, fileName, filePath, @@ -155,6 +190,18 @@ class RAGService extends EventEmitter { indexedAt: null, addedAt: Date.now(), } + this.deps.store.insertDocument(doc) + + // 청크 원문을 먼저 저장한다 — 임베딩이 실패해도 원문은 남아 재색인·기기 간 동기화가 가능하다. + this.deps.store.replaceChunks(docId, chunks) + this.deps.sync().pushOne('knowledge_documents', docId) + + // 비동기 인덱싱 (임베딩 생성). 실패는 rag:indexFailed 로 렌더러에 알린다. + this._embedStoredChunks(docId, fileName).catch((err) => { + logger.error(`Indexing failed for ${fileName}:`, err) + }) + + return doc } /** @@ -162,15 +209,13 @@ class RAGService extends EventEmitter { */ removeDocument(documentId: string): void { this.removeRemote(documentId) - getCloudSyncService().pushDelete('knowledge_documents', documentId) + this.deps.sync().pushDelete('knowledge_documents', documentId) logger.info(`RAG document removed: ${documentId}`) } /** 동기화: 다른 기기에서 지운 문서를 지운다(outbox에 넣지 않는다). */ removeRemote(documentId: string): boolean { - const db = getDatabase() - db.delete(ragChunks).where(eq(ragChunks.documentId, documentId)).run() - return db.delete(ragDocuments).where(eq(ragDocuments.id, documentId)).run().changes > 0 + return this.deps.store.removeDocument(documentId) } /** @@ -184,11 +229,11 @@ class RAGService extends EventEmitter { chunks: string[] addedAt: number }): boolean { - const db = getDatabase() - if (db.select({ id: ragDocuments.id }).from(ragDocuments).where(eq(ragDocuments.id, doc.id)).get()) return false + const { store } = this.deps + if (store.hasDocument(doc.id)) return false const chunks = doc.chunks.filter((c) => c.trim().length > 0) if (chunks.length === 0) return false - db.insert(ragDocuments).values({ + store.insertDocument({ id: doc.id, fileName: doc.fileName, filePath: '', @@ -197,8 +242,8 @@ class RAGService extends EventEmitter { indexed: false, indexedAt: null, addedAt: doc.addedAt, - }).run() - this._storeChunks(doc.id, chunks) + }) + store.replaceChunks(doc.id, chunks) this._embedStoredChunks(doc.id, doc.fileName).catch((err) => { logger.warn(`Synced document ${doc.fileName} is stored but not embedded yet:`, err) }) @@ -207,43 +252,31 @@ class RAGService extends EventEmitter { /** 저장된 청크 원문(chunkIndex 순) — 동기화 업로드용 */ getStoredChunks(documentId: string): string[] { - return getDatabase() - .select({ content: ragChunks.content }) - .from(ragChunks) - .where(eq(ragChunks.documentId, documentId)) - .orderBy(asc(ragChunks.chunkIndex)) - .all() - .map((r) => r.content) + return this.deps.store.listChunks(documentId).map((c) => c.content) } /** * 문서 재인덱싱 */ async reindex(documentId: string): Promise { - const db = getDatabase() - const rows = db.select().from(ragDocuments).where(eq(ragDocuments.id, documentId)).all() - if (rows.length === 0) { + const doc = this.deps.store.getDocument(documentId) + if (!doc) { throw new D3ROError(ErrorCode.RAGDocumentNotFound, 'Document not found') } - const doc = rows[0] // 원본 파일이 없으면(다른 기기에서 동기화된 문서 등) 저장된 원문 청크로 다시 임베딩한다. - if (!doc.filePath || !fs.existsSync(doc.filePath)) { + if (!doc.filePath || !this.deps.fileExists(doc.filePath)) { await this._embedStoredChunks(documentId, doc.fileName, { force: true }) return } // 텍스트 재추출 + 재인덱싱 - const content = await this._extractText(doc.filePath, doc.fileType as RAGDocument['fileType']) - const chunks = this._chunkText(content) + const content = await this._extractText(doc.filePath, doc.fileType) + const chunks = chunkText(content) - db.update(ragDocuments) - .set({ chunkCount: chunks.length, indexed: false }) - .where(eq(ragDocuments.id, documentId)) - .run() - - this._storeChunks(documentId, chunks) - getCloudSyncService().pushOne('knowledge_documents', documentId) + this.deps.store.updateDocument(documentId, { chunkCount: chunks.length, indexed: false }) + this.deps.store.replaceChunks(documentId, chunks) + this.deps.sync().pushOne('knowledge_documents', documentId) await this._embedStoredChunks(documentId, doc.fileName) } @@ -251,104 +284,84 @@ class RAGService extends EventEmitter { * 벡터 검색 + LLM 답변 생성 */ async query(queryText: string, topK: number = DEFAULT_TOP_K): Promise { - this._state = 'querying' - + this.activeQueries++ try { - const db = getDatabase() // 원문만 있고 아직 임베딩되지 않은 청크는 검색 대상이 아니다. - const allChunks = db.select().from(ragChunks).all().filter((c) => c.embedding.length > 0) + const allChunks = this.deps.store.listEmbeddedChunks() if (allChunks.length === 0) { throw new D3ROError(ErrorCode.RAGQueryFailed, 'No indexed chunks to query') } - // 쿼리 임베딩 - const queryEmbedding = await this._embed(queryText) - - const allDocs = db.select().from(ragDocuments).all() - const docMap = new Map(allDocs.map((d) => [d.id, d.fileName])) - - const scored = allChunks.map((chunk) => { - const embedding = JSON.parse(chunk.embedding) as number[] - const similarity = this._cosineSimilarity(queryEmbedding, embedding) - return { - documentId: chunk.documentId, - fileName: docMap.get(chunk.documentId) ?? 'unknown', - content: chunk.content, - similarity, - } - }) - - // 상위 topK - scored.sort((a, b) => b.similarity - a.similarity) - const topResults = scored.slice(0, topK) - - // LLM에 컨텍스트 주입 - const context = topResults - .map((r, i) => `[${i + 1}] (${r.fileName})\n${r.content}`) - .join('\n\n') - - const systemPrompt = `You are a helpful assistant. Answer the user's question based on the following documents. If the documents don't contain relevant information, say so. Respond in the same language as the question. - -Documents: -${context}` + const queryEmbedding = await this.deps.embedder.embed(queryText) + const fileNames = new Map(this.deps.store.listDocuments().map((d) => [d.id, d.fileName])) + const topResults = rankChunks(queryEmbedding, allChunks, fileNames, topK) + const systemPrompt = buildAnswerSystemPrompt(topResults) // 답변 LLM 은 사용자가 고른 백엔드를 따른다(로컬 선택 시 문서 조각을 클라우드로 보내지 않는다) - const result = await getLlmGateway().generate(queryText, { systemPrompt }) + const result = await this.deps.answerer().generate(queryText, { systemPrompt }) - return { - query: queryText, - results: topResults, - answer: result.text.trim(), - } + return { query: queryText, results: topResults, answer: result.text.trim() } } finally { - this._state = 'idle' + this.activeQueries-- } } // ── 내부 메서드 ── - /** 청크 원문을 (다시) 저장한다. 임베딩은 비워 두고 _embedStoredChunks 가 채운다. */ - private _storeChunks(docId: string, chunks: string[]): void { - const db = getDatabase() - db.transaction((tx) => { - tx.delete(ragChunks).where(eq(ragChunks.documentId, docId)).run() - chunks.forEach((content, index) => { - tx.insert(ragChunks).values({ - id: crypto.randomUUID(), - documentId: docId, - content, - embedding: '', - chunkIndex: index, - }).run() - }) - }) + private async _extractText(filePath: string, fileType: RagFileType): Promise { + if (isPlainTextType(fileType)) { + return (await this.deps.readFile(filePath)).toString('utf-8') + } + if (!isBinaryDocumentType(fileType)) throw unsupportedFormatError(fileType) + try { + return extractBinaryDocumentText(await this.deps.readFile(filePath), fileType) + } catch (err) { + logger.warn(`${fileType.toUpperCase()} parsing failed: ${errorText(err)}`) + throw documentParseError(fileType, err) + } } - /** 저장된 청크 중 임베딩이 없는 것(force면 전부)을 이 기기의 임베딩 모델로 채운다. */ - private async _embedStoredChunks( - docId: string, - fileName: string, - options: { force?: boolean } = {} - ): Promise { - this._state = 'indexing' - const db = getDatabase() - const chunks = db - .select() - .from(ragChunks) - .where(eq(ragChunks.documentId, docId)) - .orderBy(asc(ragChunks.chunkIndex)) - .all() + /** + * 저장된 청크 중 임베딩이 없는 것(force면 전부)을 이 기기의 임베딩 모델로 채운다. + * 실패하면 문서를 색인 안 됨으로 두고 'index-failed'·rag:indexFailed 를 알린 뒤, + * 렌더러의 진행 표시가 멈춰 있지 않도록 rag:indexComplete(실행 종료)도 보낸다. + */ + private async _embedStoredChunks(docId: string, fileName: string, options: { force?: boolean } = {}): Promise { + this.activeIndexing++ + try { + await this._indexChunks(docId, fileName, options) + } catch (err) { + this._reportIndexFailure(docId, fileName, err) + throw err + } finally { + this.activeIndexing-- + } + } + + private async _indexChunks(docId: string, fileName: string, options: { force?: boolean }): Promise { + const { store, embedder, notify } = this.deps + const chunks = store.listChunks(docId) logger.info(`RAG indexing started: ${fileName} (${chunks.length} chunks)`) + // 임베딩할 청크가 있으면 먼저 모델을 확인한다 — 모델이 없으면 청크마다 404를 기다리지 않고 바로 실패한다. + if (chunks.some((c) => options.force || c.embedding.length === 0)) { + try { + await embedder.ensureModel() + } catch (err) { + store.updateDocument(docId, { indexed: false, indexedAt: null }) + throw err + } + } + // 초기 진행률 즉시 전송 - this._sendToRenderer(IPC_CHANNELS.RAG.INDEX_PROGRESS, { + notify(IPC_CHANNELS.RAG.INDEX_PROGRESS, { documentId: docId, fileName, currentChunk: 0, totalChunks: chunks.length, percent: 0, - } as RAGIndexProgress) + } satisfies RAGIndexProgress) let successCount = 0 for (let i = 0; i < chunks.length; i++) { @@ -358,278 +371,72 @@ ${context}` continue } try { - const embedding = await this._embed(chunk.content) - db.update(ragChunks) - .set({ embedding: JSON.stringify(embedding) }) - .where(eq(ragChunks.id, chunk.id)) - .run() + store.setEmbedding(chunk.id, await embedder.embed(chunk.content)) successCount++ } catch (err) { logger.warn(`RAG embedding failed for chunk ${i}/${chunks.length} of ${fileName}:`, err) // 개별 청크 실패는 건너뛰고 계속 진행 — 원문은 남아 있어 재색인할 수 있다 } - const progress: RAGIndexProgress = { + notify(IPC_CHANNELS.RAG.INDEX_PROGRESS, { documentId: docId, fileName, currentChunk: i + 1, totalChunks: chunks.length, percent: Math.round(((i + 1) / chunks.length) * 100), - } - this._sendToRenderer(IPC_CHANNELS.RAG.INDEX_PROGRESS, progress) + } satisfies RAGIndexProgress) - // 이벤트 루프 양보 (UI 블로킹 방지) - await new Promise((r) => setTimeout(r, 10)) + await new Promise((r) => setTimeout(r, this.deps.yieldMs)) } - this._state = 'idle' - // 한 청크도 임베딩하지 못했으면(임베딩 서버 없음 등) 색인됐다고 표시하지 않는다. // indexed=true + 0 chunks 로 두면 문서가 검색 가능한 것처럼 보이지만 질의에 걸리지 않는다. if (successCount === 0) { - db.update(ragDocuments) - .set({ indexed: false, indexedAt: null }) - .where(eq(ragDocuments.id, docId)) - .run() - throw new D3ROError( - ErrorCode.RAGEmbeddingFailed, - `No chunk of ${fileName} could be embedded (${chunks.length} attempted)`, - ) + store.updateDocument(docId, { indexed: false, indexedAt: null }) + throw new D3ROError(ErrorCode.RAGEmbeddingFailed, `No chunk of ${fileName} could be embedded (${chunks.length} attempted)`) } // 인덱싱 완료 표시 (chunkCount는 원문 청크 수 — 서버·모바일과 같은 기준) - db.update(ragDocuments) - .set({ indexed: true, indexedAt: Date.now(), chunkCount: chunks.length }) - .where(eq(ragDocuments.id, docId)) - .run() + store.updateDocument(docId, { indexed: true, indexedAt: Date.now(), chunkCount: chunks.length }) - this._sendToRenderer(IPC_CHANNELS.RAG.INDEX_COMPLETE, { documentId: docId, fileName }) + notify(IPC_CHANNELS.RAG.INDEX_COMPLETE, { documentId: docId, fileName }) logger.info(`RAG indexing complete: ${fileName} (${successCount}/${chunks.length} chunks embedded)`) } - /** - * Ollama /api/embed 엔드포인트로 텍스트 임베딩 - */ - private async _embed(text: string): Promise { - const serverUrl = getOllamaServerUrl() - - try { - const response = await fetch(`${serverUrl}/api/embed`, { - method: 'POST', - headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ - model: EMBED_MODEL, - input: text, - }), - signal: AbortSignal.timeout(30000), - }) - - if (!response.ok) { - throw new D3ROError(ErrorCode.RAGEmbeddingFailed, `Embed API error: ${response.status}`) - } - - const data = (await response.json()) as { embeddings: number[][] } - if (!data.embeddings || data.embeddings.length === 0) { - throw new D3ROError(ErrorCode.RAGEmbeddingFailed, 'No embeddings returned') - } - - return data.embeddings[0] - } catch (err) { - if (err instanceof D3ROError) throw err - throw new D3ROError( - ErrorCode.RAGEmbeddingFailed, - `Embedding failed: ${err instanceof Error ? err.message : String(err)}`, - ) - } - } - - private _chunkText(text: string): string[] { - // 텍스트 크기 제한 (500KB 초과 시 잘라냄) - const MAX_TEXT_LENGTH = 500000 - const safeText = text.length > MAX_TEXT_LENGTH ? text.slice(0, MAX_TEXT_LENGTH) : text - - const chunks: string[] = [] - let start = 0 - const MAX_CHUNKS = 2000 - while (start < safeText.length && chunks.length < MAX_CHUNKS) { - const end = Math.min(start + CHUNK_SIZE, safeText.length) - const chunk = safeText.slice(start, end).trim() - if (chunk.length > 10) { - chunks.push(chunk) - } - if (end >= safeText.length) break - start = Math.max(start + 1, end - CHUNK_OVERLAP) - if (start >= safeText.length) break - } - return chunks - } - - private async _extractText(filePath: string, fileType: string): Promise { - const { promises: fsp } = await import('fs') - - if (fileType === 'txt' || fileType === 'md') { - return fsp.readFile(filePath, 'utf-8') - } - - if (fileType === 'pdf') { - try { - const buffer = await fsp.readFile(filePath) - const text = this._extractPdfText(buffer) - if (text.trim().length < 10) { - throw new Error('No readable text found in PDF') - } - return text.trim() - } catch (err) { - logger.error('PDF parsing failed:', err) - throw new D3ROError(ErrorCode.RAGUnsupportedFormat, 'PDF parsing failed. Ensure the PDF contains readable text.') - } - } - - if (fileType === 'docx') { - try { - const buffer = await fsp.readFile(filePath) - // DOCX는 ZIP 안의 (deflate 압축된) word/document.xml — 압축을 풀어 문단 텍스트를 뽑는다. - // 텍스트를 찾지 못하면 바이너리 잡음을 색인하지 않고 실패시킨다. - return extractDocxText(buffer).slice(0, 500000) - } catch (err) { - logger.warn(`DOCX parsing failed: ${err instanceof Error ? err.message : String(err)}`) - throw new D3ROError(ErrorCode.RAGUnsupportedFormat, 'DOCX parsing failed') - } - } - - throw new D3ROError(ErrorCode.RAGUnsupportedFormat, `Unsupported: ${fileType}`) - } - - /** - * PDF 바이너리에서 텍스트 추출 (순수 JS + zlib). - * FlateDecode 압축 해제 후 BT...ET 블록 내 Tj/TJ 파싱. - */ - private _extractPdfText(buffer: Buffer): string { - const zlib = require('zlib') as typeof import('zlib') - const raw = buffer.toString('binary') - const textParts: string[] = [] - - // 스트림 블록 추출 — 바이너리 오프셋 기반 - const streamMarker = 'stream\r\n' - const streamMarker2 = 'stream\n' - const endMarker = 'endstream' - - let pos = 0 - while (pos < raw.length) { - let streamStart = raw.indexOf(streamMarker, pos) - let offset = streamMarker.length - if (streamStart === -1) { - streamStart = raw.indexOf(streamMarker2, pos) - offset = streamMarker2.length - } - if (streamStart === -1) break - - const dataStart = streamStart + offset - const streamEnd = raw.indexOf(endMarker, dataStart) - if (streamEnd === -1) break - - const streamData = Buffer.from(raw.slice(dataStart, streamEnd), 'binary') - pos = streamEnd + endMarker.length - - // FlateDecode 해제 시도 - let decoded: string - try { - const inflated = zlib.inflateSync(streamData) - decoded = inflated.toString('binary') - } catch { - // 압축 안 된 스트림 - decoded = streamData.toString('binary') - } - - // BT...ET 블록 파싱 - this._extractTextFromStream(decoded, textParts) - } - - // PDF 이스케이프 디코딩 + 정리 - let text = textParts.join(' ') - text = text - .replace(/\\n/g, '\n') - .replace(/\\r/g, '\r') - .replace(/\\t/g, '\t') - .replace(/\\\(/g, '(') - .replace(/\\\)/g, ')') - .replace(/\\\\/g, '\\') - .replace(/[^\x20-\x7E\u00A0-\u00FF\u3000-\u9FFF\uAC00-\uD7AF\n]/g, ' ') - .replace(/\s+/g, ' ') - - return text.trim() - } - - private _extractTextFromStream(content: string, textParts: string[]): void { - const btRegex = /BT([\s\S]*?)ET/g - let btMatch: RegExpExecArray | null = null - while ((btMatch = btRegex.exec(content)) !== null) { - const block = btMatch[1] - - // Tj: (text) Tj - const tjRegex = /\(([^)]*)\)\s*Tj/g - let tjMatch: RegExpExecArray | null = null - while ((tjMatch = tjRegex.exec(block)) !== null) { - if (tjMatch[1].trim()) textParts.push(tjMatch[1]) - } - - // TJ: [(text) num (text)] TJ - const tjArrayRegex = /\[(.*?)\]\s*TJ/g - let tjArrayMatch: RegExpExecArray | null = null - while ((tjArrayMatch = tjArrayRegex.exec(block)) !== null) { - const items = tjArrayMatch[1] - const itemRegex = /\(([^)]*)\)/g - let itemMatch: RegExpExecArray | null = null - while ((itemMatch = itemRegex.exec(items)) !== null) { - if (itemMatch[1].trim()) textParts.push(itemMatch[1]) - } - } - - // ' 연산자: (text) ' - const quoteRegex = /\(([^)]*)\)\s*'/g - let quoteMatch: RegExpExecArray | null = null - while ((quoteMatch = quoteRegex.exec(block)) !== null) { - if (quoteMatch[1].trim()) textParts.push(quoteMatch[1]) - } - } - } - - private _cosineSimilarity(a: number[], b: number[]): number { - if (a.length !== b.length) return 0 - let dotProduct = 0 - let normA = 0 - let normB = 0 - for (let i = 0; i < a.length; i++) { - dotProduct += a[i] * b[i] - normA += a[i] * a[i] - normB += b[i] * b[i] - } - const denom = Math.sqrt(normA) * Math.sqrt(normB) - return denom === 0 ? 0 : dotProduct / denom - } - - private _sendToRenderer(channel: string, data: unknown): void { - try { - const mainWindow = getMainWindow() - if (mainWindow && !mainWindow.isDestroyed()) { - mainWindow.webContents.send(channel, data) - } - } catch { - // ignore + private _reportIndexFailure(documentId: string, fileName: string, err: unknown): void { + const failure: RAGIndexFailure = { + documentId, + fileName, + code: err instanceof D3ROError ? err.code : ErrorCode.RAGIndexingFailed, + message: errorText(err), } + this.deps.notify(IPC_CHANNELS.RAG.INDEX_FAILED, failure) + this.deps.notify(IPC_CHANNELS.RAG.INDEX_COMPLETE, { documentId, fileName }) + this.emit('index-failed', failure) } dispose(): void { this.removeAllListeners() } + + // ── 타입 안전 이벤트 ── + + override on(event: K, listener: RAGServiceEvents[K]): this { + return super.on(event, listener) + } + + override emit(event: K, ...args: Parameters): boolean { + return super.emit(event, ...args) + } } // ── 싱글톤 ── let instance: RAGService | null = null -export function resetRAGServiceForTests(): void { +/** 테스트용: 싱글톤을 비운다. deps를 주면 그 포트로 새 인스턴스를 만든다. */ +export function resetRAGServiceForTests(deps?: Partial): void { if (instance) instance.removeAllListeners() - instance = null + instance = deps ? new RAGService({ ...defaultRAGServiceDeps(), ...deps }) : null } export function getRAGService(): RAGService { diff --git a/apps/desktop/src/main/services/rag/chunk-store.ts b/apps/desktop/src/main/services/rag/chunk-store.ts new file mode 100644 index 0000000..25bab76 --- /dev/null +++ b/apps/desktop/src/main/services/rag/chunk-store.ts @@ -0,0 +1,115 @@ +// src/main/services/rag/chunk-store.ts +// 지식 문서·청크 저장소 포트와 SQLite(drizzle) 구현. RAGService는 getDatabase()를 직접 부르지 않는다. + +import { asc, eq } from 'drizzle-orm' +import { getDatabase } from '../../db' +import { ragChunks, ragDocuments } from '../../db/schema' +import type { RAGDocument } from '@d3ro/core/types' + +export type StoredDocument = RAGDocument + +export interface StoredChunk { + id: string + documentId: string + content: string + /** JSON 직렬화 float[]. 아직 임베딩하지 않았으면 빈 문자열. */ + embedding: string + chunkIndex: number +} + +export interface ChunkStore { + listDocuments(): StoredDocument[] + getDocument(id: string): StoredDocument | null + hasDocument(id: string): boolean + insertDocument(doc: StoredDocument): void + /** 색인 상태·청크 수 갱신 */ + updateDocument(id: string, patch: Partial>): void + /** 문서와 청크를 지운다. 문서 행이 있었으면 true. */ + removeDocument(id: string): boolean + /** 문서의 청크 원문을 (다시) 저장한다. 임베딩은 비워 둔다. */ + replaceChunks(documentId: string, chunks: readonly string[]): void + /** 문서의 청크(chunkIndex 순) */ + listChunks(documentId: string): StoredChunk[] + setEmbedding(chunkId: string, embedding: readonly number[]): void + /** 임베딩이 채워진 모든 청크(검색 대상) */ + listEmbeddedChunks(): StoredChunk[] + counts(): { documentCount: number; totalChunks: number } +} + +function toDocument(r: typeof ragDocuments.$inferSelect): StoredDocument { + return { + id: r.id, + fileName: r.fileName, + filePath: r.filePath, + fileType: r.fileType, + chunkCount: r.chunkCount, + indexed: r.indexed, + indexedAt: r.indexedAt, + addedAt: r.addedAt, + } +} + +export class SqliteChunkStore implements ChunkStore { + listDocuments(): StoredDocument[] { + return getDatabase().select().from(ragDocuments).all().map(toDocument) + } + + getDocument(id: string): StoredDocument | null { + const row = getDatabase().select().from(ragDocuments).where(eq(ragDocuments.id, id)).get() + return row ? toDocument(row) : null + } + + hasDocument(id: string): boolean { + return getDatabase().select({ id: ragDocuments.id }).from(ragDocuments).where(eq(ragDocuments.id, id)).get() !== undefined + } + + insertDocument(doc: StoredDocument): void { + getDatabase().insert(ragDocuments).values(doc).run() + } + + updateDocument(id: string, patch: Partial>): void { + getDatabase().update(ragDocuments).set(patch).where(eq(ragDocuments.id, id)).run() + } + + removeDocument(id: string): boolean { + const db = getDatabase() + db.delete(ragChunks).where(eq(ragChunks.documentId, id)).run() + return db.delete(ragDocuments).where(eq(ragDocuments.id, id)).run().changes > 0 + } + + replaceChunks(documentId: string, chunks: readonly string[]): void { + getDatabase().transaction((tx) => { + tx.delete(ragChunks).where(eq(ragChunks.documentId, documentId)).run() + chunks.forEach((content, index) => { + tx.insert(ragChunks) + .values({ id: crypto.randomUUID(), documentId, content, embedding: '', chunkIndex: index }) + .run() + }) + }) + } + + listChunks(documentId: string): StoredChunk[] { + return getDatabase() + .select() + .from(ragChunks) + .where(eq(ragChunks.documentId, documentId)) + .orderBy(asc(ragChunks.chunkIndex)) + .all() + } + + setEmbedding(chunkId: string, embedding: readonly number[]): void { + getDatabase().update(ragChunks).set({ embedding: JSON.stringify(embedding) }).where(eq(ragChunks.id, chunkId)).run() + } + + listEmbeddedChunks(): StoredChunk[] { + return getDatabase().select().from(ragChunks).all().filter((c) => c.embedding.length > 0) + } + + counts(): { documentCount: number; totalChunks: number } { + const db = getDatabase() + return { + documentCount: db.select({ id: ragDocuments.id }).from(ragDocuments).all().length, + totalChunks: db.select({ id: ragChunks.id }).from(ragChunks).all().length, + } + } +} diff --git a/apps/desktop/src/main/services/rag/document-text.ts b/apps/desktop/src/main/services/rag/document-text.ts new file mode 100644 index 0000000..38d1146 --- /dev/null +++ b/apps/desktop/src/main/services/rag/document-text.ts @@ -0,0 +1,86 @@ +// src/main/services/rag/document-text.ts +// 지식 문서 원문 → 텍스트 → 청크. 파일 I/O·DB·네트워크 없는 순수 함수만 둔다(RAGService에서 옮김, 동작 동일). + +import { D3ROError, ErrorCode } from '@d3ro/core/errors' +import type { RAGDocument } from '@d3ro/core/types' +import { extractDocxText } from './docx-text' +import { extractPdfText } from './pdf-text' + +export type RagFileType = RAGDocument['fileType'] + +/** 청크 크기 (자) */ +export const CHUNK_SIZE = 500 +/** 청크 오버랩 (자) */ +export const CHUNK_OVERLAP = 50 +/** 색인할 원문 상한 (500KB 초과 시 잘라낸다) */ +export const MAX_TEXT_LENGTH = 500_000 +/** 문서 하나의 청크 상한 */ +export const MAX_CHUNKS = 2000 +/** 이보다 짧은 청크는 버린다 */ +const MIN_CHUNK_LENGTH = 10 +/** PDF에서 이보다 짧은 텍스트만 나오면 읽을 수 없는 PDF로 본다 */ +const MIN_PDF_TEXT_LENGTH = 10 + +const EXTENSION_TYPES: Readonly> = { + '.txt': 'txt', + '.md': 'md', + '.pdf': 'pdf', + '.docx': 'docx', +} + +/** 확장자(점 포함, 소문자)로 문서 종류를 고른다. 지원하지 않으면 null. */ +export function fileTypeForExtension(ext: string): RagFileType | null { + return EXTENSION_TYPES[ext] ?? null +} + +/** 그대로 UTF-8로 읽는 문서인지 */ +export function isPlainTextType(fileType: string): fileType is 'txt' | 'md' { + return fileType === 'txt' || fileType === 'md' +} + +/** 압축을 풀어 파싱해야 하는 문서인지 */ +export function isBinaryDocumentType(fileType: string): fileType is 'pdf' | 'docx' { + return fileType === 'pdf' || fileType === 'docx' +} + +export function unsupportedFormatError(fileType: string): D3ROError { + return new D3ROError(ErrorCode.RAGUnsupportedFormat, `Unsupported: ${fileType}`) +} + +/** PDF/DOCX 파싱 실패를 사용자에게 보이는 오류로 바꾼다(원인은 details.cause). */ +export function documentParseError(fileType: 'pdf' | 'docx', cause: unknown): D3ROError { + const details = { cause: cause instanceof Error ? cause.message : String(cause) } + return fileType === 'pdf' + ? new D3ROError(ErrorCode.RAGUnsupportedFormat, 'PDF parsing failed. Ensure the PDF contains readable text.', details) + : new D3ROError(ErrorCode.RAGUnsupportedFormat, 'DOCX parsing failed', details) +} + +/** + * PDF/DOCX 버퍼에서 본문 텍스트를 뽑는다. 읽을 텍스트가 없으면 Error를 던진다 + * (호출자가 documentParseError로 감싼다). + */ +export function extractBinaryDocumentText(buffer: Buffer, fileType: 'pdf' | 'docx'): string { + if (fileType === 'pdf') { + const text = extractPdfText(buffer) + if (text.trim().length < MIN_PDF_TEXT_LENGTH) throw new Error('No readable text found in PDF') + return text.trim() + } + // DOCX는 ZIP 안의 (deflate 압축된) word/document.xml — 텍스트를 찾지 못하면 바이너리 잡음을 색인하지 않고 실패시킨다. + return extractDocxText(buffer).slice(0, MAX_TEXT_LENGTH) +} + +/** 텍스트를 CHUNK_SIZE 창으로 CHUNK_OVERLAP 만큼 겹쳐 자른다. */ +export function chunkText(text: string): string[] { + const safeText = text.length > MAX_TEXT_LENGTH ? text.slice(0, MAX_TEXT_LENGTH) : text + const chunks: string[] = [] + let start = 0 + while (start < safeText.length && chunks.length < MAX_CHUNKS) { + const end = Math.min(start + CHUNK_SIZE, safeText.length) + const chunk = safeText.slice(start, end).trim() + if (chunk.length > MIN_CHUNK_LENGTH) chunks.push(chunk) + if (end >= safeText.length) break + start = Math.max(start + 1, end - CHUNK_OVERLAP) + if (start >= safeText.length) break + } + return chunks +} diff --git a/apps/desktop/src/main/services/rag/embedding-port.ts b/apps/desktop/src/main/services/rag/embedding-port.ts new file mode 100644 index 0000000..78b41d6 --- /dev/null +++ b/apps/desktop/src/main/services/rag/embedding-port.ts @@ -0,0 +1,105 @@ +// src/main/services/rag/embedding-port.ts +// 로컬 RAG 임베딩 포트와 Ollama 어댑터. RAGService는 이 인터페이스만 본다(테스트는 가짜로 바꿔 끼운다). + +import { D3ROError, ErrorCode } from '@d3ro/core/errors' + +/** 데스크톱 로컬 임베딩 모델 (768차원). 서버(1536)와 공간이 달라 벡터는 기기 간에 옮기지 않는다. */ +export const EMBED_MODEL = 'nomic-embed-text' + +const EMBED_TIMEOUT_MS = 30_000 +const TAGS_TIMEOUT_MS = 5_000 + +export interface EmbeddingPort { + /** 사용하는 모델 이름 (오류 메시지·로그용) */ + readonly model: string + /** + * 임베딩을 만들 수 있는지 확인한다. 서버에 닿지 않거나 모델이 설치돼 있지 않으면 + * D3ROError(RAGEmbeddingFailed)를 던진다. 모델을 내려받지는 않는다. + */ + ensureModel(signal?: AbortSignal): Promise + /** 텍스트 하나의 임베딩 벡터. 실패하면 D3ROError(RAGEmbeddingFailed). */ + embed(text: string, signal?: AbortSignal): Promise +} + +type FetchFn = (input: string, init?: RequestInit) => Promise + +export interface OllamaEmbeddingOptions { + /** Ollama 서버 주소(호출 때마다 읽는다 — 설정 변경을 따른다) */ + serverUrl: () => string + model?: string + fetchFn?: FetchFn +} + +function errorText(err: unknown): string { + return err instanceof Error ? err.message : String(err) +} + +function withTimeout(ms: number, signal?: AbortSignal): AbortSignal { + const timeout = AbortSignal.timeout(ms) + return signal ? AbortSignal.any([signal, timeout]) : timeout +} + +/** /api/tags 의 모델 이름이 찾는 모델인지 ('nomic-embed-text' ↔ 'nomic-embed-text:latest'). */ +export function isSameOllamaModel(installed: string, wanted: string): boolean { + const normalize = (name: string): string => (name.includes(':') ? name : `${name}:latest`) + return normalize(installed) === normalize(wanted) +} + +export class OllamaEmbeddingAdapter implements EmbeddingPort { + readonly model: string + private readonly serverUrl: () => string + private readonly fetchFn: FetchFn + + constructor(options: OllamaEmbeddingOptions) { + this.serverUrl = options.serverUrl + this.model = options.model ?? EMBED_MODEL + this.fetchFn = options.fetchFn ?? ((input, init) => fetch(input, init)) + } + + async ensureModel(signal?: AbortSignal): Promise { + const serverUrl = this.serverUrl() + let names: string[] + try { + const response = await this.fetchFn(`${serverUrl}/api/tags`, { signal: withTimeout(TAGS_TIMEOUT_MS, signal) }) + if (!response.ok) throw new Error(`HTTP ${response.status}`) + const data = (await response.json()) as { models?: Array<{ name?: unknown; model?: unknown }> } + names = (data.models ?? []).flatMap((m) => + [m.name, m.model].filter((n): n is string => typeof n === 'string') + ) + } catch (err) { + throw new D3ROError(ErrorCode.RAGEmbeddingFailed, `Ollama is not reachable at ${serverUrl}: ${errorText(err)}`, { + reason: 'server_unreachable', + model: this.model, + }) + } + if (!names.some((n) => isSameOllamaModel(n, this.model))) { + throw new D3ROError( + ErrorCode.RAGEmbeddingFailed, + `Embedding model "${this.model}" is not installed in Ollama (run: ollama pull ${this.model})`, + { reason: 'model_missing', model: this.model } + ) + } + } + + async embed(text: string, signal?: AbortSignal): Promise { + try { + const response = await this.fetchFn(`${this.serverUrl()}/api/embed`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ model: this.model, input: text }), + signal: withTimeout(EMBED_TIMEOUT_MS, signal), + }) + if (!response.ok) { + throw new D3ROError(ErrorCode.RAGEmbeddingFailed, `Embed API error: ${response.status}`) + } + const data = (await response.json()) as { embeddings?: number[][] } + if (!data.embeddings || data.embeddings.length === 0) { + throw new D3ROError(ErrorCode.RAGEmbeddingFailed, 'No embeddings returned') + } + return data.embeddings[0] + } catch (err) { + if (err instanceof D3ROError) throw err + throw new D3ROError(ErrorCode.RAGEmbeddingFailed, `Embedding failed: ${errorText(err)}`) + } + } +} diff --git a/apps/desktop/src/main/services/rag/pdf-text.ts b/apps/desktop/src/main/services/rag/pdf-text.ts new file mode 100644 index 0000000..1fbafd7 --- /dev/null +++ b/apps/desktop/src/main/services/rag/pdf-text.ts @@ -0,0 +1,94 @@ +// src/main/services/rag/pdf-text.ts +// PDF 바이너리에서 텍스트를 뽑는다 — 순수 JS + zlib, 외부 의존 없음(RAGService에서 옮김, 동작 동일). +// FlateDecode 스트림을 풀고 BT...ET 블록 안의 Tj/TJ/' 연산자 문자열을 모은다. + +import { inflateSync } from 'zlib' + +const STREAM_MARKER_CRLF = 'stream\r\n' +const STREAM_MARKER_LF = 'stream\n' +const END_MARKER = 'endstream' + +/** 텍스트 블록(BT...ET) 하나에서 문자열 조각을 모은다. */ +export function extractTextOperators(content: string, textParts: string[]): void { + const btRegex = /BT([\s\S]*?)ET/g + let btMatch: RegExpExecArray | null = null + while ((btMatch = btRegex.exec(content)) !== null) { + const block = btMatch[1] + + // Tj: (text) Tj + const tjRegex = /\(([^)]*)\)\s*Tj/g + let tjMatch: RegExpExecArray | null = null + while ((tjMatch = tjRegex.exec(block)) !== null) { + if (tjMatch[1].trim()) textParts.push(tjMatch[1]) + } + + // TJ: [(text) num (text)] TJ + const tjArrayRegex = /\[(.*?)\]\s*TJ/g + let tjArrayMatch: RegExpExecArray | null = null + while ((tjArrayMatch = tjArrayRegex.exec(block)) !== null) { + const items = tjArrayMatch[1] + const itemRegex = /\(([^)]*)\)/g + let itemMatch: RegExpExecArray | null = null + while ((itemMatch = itemRegex.exec(items)) !== null) { + if (itemMatch[1].trim()) textParts.push(itemMatch[1]) + } + } + + // ' 연산자: (text) ' + const quoteRegex = /\(([^)]*)\)\s*'/g + let quoteMatch: RegExpExecArray | null = null + while ((quoteMatch = quoteRegex.exec(block)) !== null) { + if (quoteMatch[1].trim()) textParts.push(quoteMatch[1]) + } + } +} + +/** PDF 문자열 이스케이프를 풀고 표시할 수 없는 문자를 공백으로 정리한다. */ +export function normalizePdfText(parts: readonly string[]): string { + return parts + .join(' ') + .replace(/\\n/g, '\n') + .replace(/\\r/g, '\r') + .replace(/\\t/g, '\t') + .replace(/\\\(/g, '(') + .replace(/\\\)/g, ')') + .replace(/\\\\/g, '\\') + .replace(/[^\x20-\x7E\u00A0-\u00FF\u3000-\u9FFF\uAC00-\uD7AF\n]/g, ' ') + .replace(/\s+/g, ' ') + .trim() +} + +/** PDF 바이너리에서 텍스트 추출. 읽을 텍스트가 없으면 빈 문자열. */ +export function extractPdfText(buffer: Buffer): string { + const raw = buffer.toString('binary') + const textParts: string[] = [] + + let pos = 0 + while (pos < raw.length) { + let streamStart = raw.indexOf(STREAM_MARKER_CRLF, pos) + let offset = STREAM_MARKER_CRLF.length + if (streamStart === -1) { + streamStart = raw.indexOf(STREAM_MARKER_LF, pos) + offset = STREAM_MARKER_LF.length + } + if (streamStart === -1) break + + const dataStart = streamStart + offset + const streamEnd = raw.indexOf(END_MARKER, dataStart) + if (streamEnd === -1) break + + const streamData = Buffer.from(raw.slice(dataStart, streamEnd), 'binary') + pos = streamEnd + END_MARKER.length + + let decoded: string + try { + decoded = inflateSync(streamData).toString('binary') + } catch { + // 압축 안 된 스트림 + decoded = streamData.toString('binary') + } + extractTextOperators(decoded, textParts) + } + + return normalizePdfText(textParts) +} diff --git a/apps/desktop/src/main/services/rag/retrieval.ts b/apps/desktop/src/main/services/rag/retrieval.ts new file mode 100644 index 0000000..2caee38 --- /dev/null +++ b/apps/desktop/src/main/services/rag/retrieval.ts @@ -0,0 +1,53 @@ +// src/main/services/rag/retrieval.ts +// 벡터 검색과 답변 프롬프트 조립 — 순수 함수(RAGService에서 옮김, 동작 동일). + +import type { RAGQueryResult } from '@d3ro/core/types' + +export type RetrievedChunk = RAGQueryResult['results'][number] + +export interface EmbeddedChunk { + documentId: string + content: string + /** JSON 직렬화 float[] */ + embedding: string +} + +export function cosineSimilarity(a: readonly number[], b: readonly number[]): number { + if (a.length !== b.length) return 0 + let dotProduct = 0 + let normA = 0 + let normB = 0 + for (let i = 0; i < a.length; i++) { + dotProduct += a[i] * b[i] + normA += a[i] * a[i] + normB += b[i] * b[i] + } + const denom = Math.sqrt(normA) * Math.sqrt(normB) + return denom === 0 ? 0 : dotProduct / denom +} + +/** 질의 임베딩과 가장 가까운 청크 topK개 (유사도 내림차순). */ +export function rankChunks( + queryEmbedding: readonly number[], + chunks: readonly EmbeddedChunk[], + fileNames: ReadonlyMap, + topK: number +): RetrievedChunk[] { + const scored = chunks.map((chunk) => ({ + documentId: chunk.documentId, + fileName: fileNames.get(chunk.documentId) ?? 'unknown', + content: chunk.content, + similarity: cosineSimilarity(queryEmbedding, JSON.parse(chunk.embedding) as number[]), + })) + scored.sort((a, b) => b.similarity - a.similarity) + return scored.slice(0, topK) +} + +/** 검색된 문서 조각을 넣은 답변용 시스템 프롬프트. */ +export function buildAnswerSystemPrompt(results: readonly RetrievedChunk[]): string { + const context = results.map((r, i) => `[${i + 1}] (${r.fileName})\n${r.content}`).join('\n\n') + return `You are a helpful assistant. Answer the user's question based on the following documents. If the documents don't contain relevant information, say so. Respond in the same language as the question. + +Documents: +${context}` +} diff --git a/apps/desktop/src/main/services/sync/SyncEngine.ts b/apps/desktop/src/main/services/sync/SyncEngine.ts index ef387eb..ae5528e 100644 --- a/apps/desktop/src/main/services/sync/SyncEngine.ts +++ b/apps/desktop/src/main/services/sync/SyncEngine.ts @@ -49,6 +49,7 @@ import { type SyncRemote, type SyncRunResult, } from './sync-types' +import { earliestDeletedAt, isRestoredAfter, tombstoneWindowExpired, type TombstoneRef } from './tombstone-policy' const logger = getLogger('SyncEngine') @@ -59,6 +60,8 @@ const PULL_PAGE_SIZE = 500 */ const CURSOR_OVERLAP_MS = 60_000 const BACKFILL_FLAG = 'backfill:v1' +/** 마지막으로 tombstone을 끝까지 읽은 기기 시각(ms). 서버 보존 기간 초과 판단에 쓴다. */ +const TOMBSTONES_PULLED_AT = 'tombstones:pulledAt' /** 되감은 커서의 id 하한. 행 테이블은 uuid, sync_tombstones는 bigserial이다. */ const MIN_UUID = '00000000-0000-0000-0000-000000000000' const MIN_SERIAL = '0' @@ -623,12 +626,21 @@ export class SyncEngine extends EventEmitter { if (next !== before) setSyncState(deferredKey(entity), next, this.now()) } - /** 다른 기기에서 지운 행을 로컬에서도 지운다. 원격 삭제는 로컬 미전송 변경보다 우선한다. */ + /** + * 다른 기기에서 지운 행을 로컬에서도 지운다. 원격 삭제는 로컬 미전송 변경보다 우선한다. + * 단, tombstone 뒤에 서버에 다시 쓰인 행(계정 보관본 복원)은 지우지 않는다 — 겹쳐 읽기(overlap)나 + * 오프라인 뒤 첫 pull에서 옛 tombstone이 복원된 행을 다시 지우던 문제를 막는다. + */ private async pullTombstones(changed: Set): Promise { const saved = parseCursor(getSyncState(cursorKey('tombstones'))) + let removed = 0 + if (saved && tombstoneWindowExpired(this.lastTombstonePullAt(saved), this.now())) { + removed += await this.reconcileAfterPrunedTombstones(changed) + } let after = rewind(saved, MIN_SERIAL) let newest = saved - let removed = 0 + /** 엔티티별 로컬 id — 복원 확인이 필요한 행(로컬에 아직 있는 행)만 서버에 물어본다. */ + const localIds = new Map>() for (;;) { const page = await this.remote.fetchPage({ table: 'sync_tombstones', @@ -640,36 +652,173 @@ export class SyncEngine extends EventEmitter { }) this.checkpoint() if (page.length === 0) break + const byEntity = new Map() for (const row of page) { const entity = row.table_name const rowId = row.row_id // memo_tags는 서버 id라 로컬과 맞지 않는다 — 전체 대조가 처리한다. if (isSyncEntity(entity) && entity !== 'memo_tags' && typeof rowId === 'string') { - const adapter = adapterFor(entity) - if (adapter) { - dropEntry(entity, rowId) - try { - if (adapter.deleteLocal(rowId)) { - removed++ - changed.add(entity) - } - } catch (err) { - logger.warn(`Tombstone ${entity}/${rowId} failed: ${errorMessage(err)}`) - } - } + const refs = byEntity.get(entity) ?? [] + refs.push({ rowId, deletedAt: typeof row.deleted_at === 'string' ? row.deleted_at : null }) + byEntity.set(entity, refs) } const cursor = rowCursor(row, 'deleted_at') if (cursor && isAfter(cursor, newest)) newest = cursor } + for (const [entity, refs] of byEntity) { + const adapter = adapterFor(entity) + if (adapter) removed += await this.applyTombstones(adapter, refs, localIds, changed) + } if (page.length < PULL_PAGE_SIZE) break const next = rowCursor(page[page.length - 1], 'deleted_at') if (!next) break after = next } + this.checkpoint() if (newest && newest !== saved) setSyncState(cursorKey('tombstones'), JSON.stringify(newest)) + setSyncState(TOMBSTONES_PULLED_AT, String(this.now()), this.now()) return removed } + /** 한 엔티티의 tombstone 묶음을 반영한다. 서버에서 되살아난 행은 건너뛴다. */ + private async applyTombstones( + adapter: SyncAdapter, + refs: TombstoneRef[], + localIds: Map>, + changed: Set + ): Promise { + let local = localIds.get(adapter.entity) + if (!local) { + local = new Set(adapter.listLocalVersions().map((v) => v.id)) + localIds.set(adapter.entity, local) + } + const localNow = local + const present = refs.filter((r) => localNow.has(r.rowId)) + const restored = present.length > 0 ? await this.findRestoredRows(adapter, present) : new Set() + this.checkpoint() + let removed = 0 + for (const ref of refs) { + if (restored.has(ref.rowId)) { + logger.info(`Tombstone ${adapter.entity}/${ref.rowId} skipped: row was restored on the server`) + continue + } + dropEntry(adapter.entity, ref.rowId) + try { + if (adapter.deleteLocal(ref.rowId)) { + removed++ + changed.add(adapter.entity) + } + localNow.delete(ref.rowId) + } catch (err) { + logger.warn(`Tombstone ${adapter.entity}/${ref.rowId} failed: ${errorMessage(err)}`) + } + } + return removed + } + + /** + * tombstone보다 나중에 서버 updated_at이 찍힌 행 id(= 삭제 뒤 복원된 행). + * 가장 이른 삭제 시각 뒤에 바뀐 행만 훑는다 — 평소(삭제만 있음)엔 거의 비어 있다. + */ + private async findRestoredRows(adapter: SyncAdapter, refs: TombstoneRef[]): Promise> { + const restored = new Set() + const since = earliestDeletedAt(refs) + if (since === null) return restored + const byId = new Map(refs.map((r) => [r.rowId, r])) + let after: RemoteCursor | null = { ts: since, id: MIN_UUID } + for (;;) { + const page = await this.remote.fetchPage({ + table: adapter.entity, + userId: this.userId, + cursorColumn: 'updated_at', + after, + limit: PULL_PAGE_SIZE, + columns: 'id,updated_at', + filters: adapter.pull?.filters, + }) + this.checkpoint() + for (const row of page) { + const ref = typeof row.id === 'string' ? byId.get(row.id) : undefined + if (ref && isRestoredAfter(ref, row.updated_at)) restored.add(ref.rowId) + } + if (page.length < PULL_PAGE_SIZE || restored.size === byId.size) break + const next = rowCursor(page[page.length - 1], 'updated_at') + if (!next) break + after = next + } + return restored + } + + /** 마지막으로 tombstone을 끝까지 읽은 시각(기기 시계). 예전 판은 lastPullAt·커서 시각으로 가늠한다. */ + private lastTombstonePullAt(saved: RemoteCursor): number | null { + for (const key of [TOMBSTONES_PULLED_AT, 'lastPullAt']) { + const value = getSyncState(key) + const at = value === null ? Number.NaN : Number(value) + if (Number.isFinite(at)) return at + } + const at = Date.parse(saved.ts) + return Number.isFinite(at) ? at : null + } + + /** + * 서버 tombstone 보존 기간보다 오래 삭제를 읽지 못한 기기의 전체 대조. + * 그 사이 삭제의 tombstone은 이미 지워졌으므로, 서버에 없고 올릴 변경도 없는 로컬 행을 지운다 + * (그대로 두면 다음 로컬 편집이 지운 행을 서버·모든 기기에 되살린다). + * 최초 대조(backfill) 전이면 하지 않는다 — 아직 올리지 않은 로컬 행을 지우게 된다. + */ + private async reconcileAfterPrunedTombstones(changed: Set): Promise { + if (getSyncState(BACKFILL_FLAG) !== 'done') return 0 + let removed = 0 + for (const adapter of TABLE_ADAPTERS) { + if (!adapter.pull) continue + const local = adapter.listLocalVersions() + if (local.length === 0) continue + const remoteIds = await this.listAllRemoteIds(adapter) + this.checkpoint() + const pending = pendingOps(adapter.entity) + for (const row of local) { + if (remoteIds.has(row.id) || pending.has(row.id)) continue + try { + if (adapter.deleteLocal(row.id)) { + removed++ + changed.add(adapter.entity) + } + } catch (err) { + logger.warn(`Reconcile ${adapter.entity}/${row.id} failed: ${errorMessage(err)}`) + } + } + } + logger.warn(`Tombstone retention exceeded; full reconcile removed ${removed} stale local row(s)`) + return removed + } + + /** + * 서버의 모든 행 id. 빈 페이지가 올 때까지 넘긴다 — 서버 max_rows가 요청한 limit보다 작아도 + * 목록이 잘리지 않게(잘리면 서버에 있는 행을 지우게 된다). + */ + private async listAllRemoteIds(adapter: SyncAdapter): Promise> { + const ids = new Set() + let after: RemoteCursor | null = null + for (;;) { + const page = await this.remote.fetchPage({ + table: adapter.entity, + userId: this.userId, + cursorColumn: 'updated_at', + after, + limit: PULL_PAGE_SIZE, + columns: 'id,updated_at', + filters: adapter.pull?.filters, + }) + this.checkpoint() + if (page.length === 0) break + for (const row of page) if (typeof row.id === 'string') ids.add(row.id) + const next = rowCursor(page[page.length - 1], 'updated_at') + if (!next) break + after = next + } + return ids + } + // ── 타입 안전 이벤트 ────────────────────────────────── override on(event: K, listener: SyncEngineEvents[K]): this { diff --git a/apps/desktop/src/main/services/sync/sync-adapters.ts b/apps/desktop/src/main/services/sync/sync-adapters.ts index fcdc4d3..7026a8c 100644 --- a/apps/desktop/src/main/services/sync/sync-adapters.ts +++ b/apps/desktop/src/main/services/sync/sync-adapters.ts @@ -955,6 +955,30 @@ async function pushKnowledgeDocument(ctx: PushContext, id: string): Promise() + for (const c of chunkRows) { + const index = num(c.chunk_index) + const content = str(c.content) + if (index === null || content === null || byIndex.has(index)) return null + byIndex.set(index, content) + } + const chunks: string[] = [] + for (let i = 0; i < byIndex.size; i++) { + const content = byIndex.get(i) + if (content === undefined) return null + chunks.push(content) + } + return chunks +} + const knowledgeAdapter: SyncAdapter = { entity: 'knowledge_documents', pull: { columns: 'id,title,file_name,file_type,chunk_count,created_at,updated_at' }, @@ -982,9 +1006,10 @@ const knowledgeAdapter: SyncAdapter = { const chunkRows = await ctx.remote.selectChildren('knowledge_chunks', 'document_id', id, 'chunk_index,content', 'chunk_index') // 네트워크를 기다리는 사이 로그아웃·계정 전환이 있었으면 다른 사용자 DB에 쓰지 않는다. assertCurrent(ctx) - const chunks = chunkRows.map((c) => str(c.content)).filter((c): c is string => c !== null) - // 다른 기기가 문서 행을 올리고 청크를 아직 못 올렸을 수 있다 — 다음 pull에서 다시 본다. - if (chunks.length === 0) return 'deferred' + // 다른 기기가 문서 행을 올리고 청크를 아직(또는 일부만) 올렸을 수 있다 — 다음 pull에서 다시 본다. + // 일부만 받아 저장하면 이후 pull은 "이미 있음"으로 건너뛰어 잘린 문서가 영구히 남는다. + const chunks = completeKnowledgeChunks(chunkRows, num(row.chunk_count)) + if (chunks === null) return 'deferred' return getRAGService().applyRemoteDocument({ id, fileName: str(row.file_name) ?? str(row.title) ?? 'document', diff --git a/apps/desktop/src/main/services/sync/tombstone-policy.ts b/apps/desktop/src/main/services/sync/tombstone-policy.ts new file mode 100644 index 0000000..356312b --- /dev/null +++ b/apps/desktop/src/main/services/sync/tombstone-policy.ts @@ -0,0 +1,58 @@ +// src/main/services/sync/tombstone-policy.ts +// 원격 삭제(sync_tombstones) 반영 규칙 — 순수 함수만 둔다(엔진·DB·네트워크 없음). +// +// 1) 되살린 행: 계정 보관본 복원(restore_account_portability)은 지운 행을 같은 id로 다시 넣는다. +// 서버는 그 tombstone을 지우지 않으므로, tombstone보다 나중에 서버 updated_at이 찍힌 행은 +// "삭제 뒤 되살아난 행"이다 — 로컬에서 다시 지우면 안 된다. +// 2) 보존 기간: 서버는 tombstone을 180일 뒤 지운다(prune_sync_tombstones_v1, 마이그레이션 0034/0035). +// 그보다 오래 tombstone을 읽지 못한 기기는 그 사이 삭제를 영영 알 수 없으므로 전체 대조를 해야 한다. + +const DAY_MS = 24 * 60 * 60 * 1000 + +/** 서버 tombstone 보존 기간(prune_sync_tombstones_v1 의 p_retention 기본값). */ +export const TOMBSTONE_RETENTION_MS = 180 * DAY_MS +/** 서버 prune 주기(하루)와 기기 시계 오차를 덮는 여유. */ +export const TOMBSTONE_RETENTION_MARGIN_MS = 7 * DAY_MS + +export interface TombstoneRef { + rowId: string + /** 서버 deleted_at 원문(마이크로초 포함). 없으면 null. */ + deletedAt: string | null +} + +/** + * 서버 timestamptz 문자열 a가 b보다 나중인지. 같은 밀리초 안은 원문(마이크로초)으로 가린다. + * 해석할 수 없는 값은 나중이 아닌 것으로 본다(보수적으로 기존 삭제 동작을 유지한다). + */ +export function remoteTimeAfter(a: string, b: string): boolean { + const x = Date.parse(a) + const y = Date.parse(b) + if (!Number.isFinite(x) || !Number.isFinite(y)) return false + if (x !== y) return x > y + return a > b +} + +/** tombstone 뒤에 서버에 다시 쓰인 행인지(삭제 뒤 복원). */ +export function isRestoredAfter(tombstone: TombstoneRef, remoteUpdatedAt: unknown): boolean { + if (tombstone.deletedAt === null || typeof remoteUpdatedAt !== 'string') return false + return remoteTimeAfter(remoteUpdatedAt, tombstone.deletedAt) +} + +/** 확인해 볼 복원 후보들 중 가장 이른 삭제 시각 — 그 뒤에 바뀐 행만 보면 된다. */ +export function earliestDeletedAt(tombstones: readonly TombstoneRef[]): string | null { + let earliest: string | null = null + for (const t of tombstones) { + if (t.deletedAt === null) continue + if (earliest === null || remoteTimeAfter(earliest, t.deletedAt)) earliest = t.deletedAt + } + return earliest +} + +/** + * 마지막으로 tombstone을 끝까지 읽은 시각(기기 시계 ms)으로부터 보존 기간(여유 포함)이 지났는지. + * 지났으면 서버가 이미 지운 tombstone이 있을 수 있다 — 전체 대조가 필요하다. + */ +export function tombstoneWindowExpired(lastPulledAt: number | null, now: number): boolean { + if (lastPulledAt === null || !Number.isFinite(lastPulledAt)) return false + return now - lastPulledAt > TOMBSTONE_RETENTION_MS - TOMBSTONE_RETENTION_MARGIN_MS +} diff --git a/apps/desktop/src/renderer/pages/KnowledgeBasePage.tsx b/apps/desktop/src/renderer/pages/KnowledgeBasePage.tsx index 295779a..7cb5f50 100644 --- a/apps/desktop/src/renderer/pages/KnowledgeBasePage.tsx +++ b/apps/desktop/src/renderer/pages/KnowledgeBasePage.tsx @@ -11,6 +11,7 @@ import { isImeComposingEvent } from '../utils/keyboard' import { d3roPalette, d3roFontSans, d3roTypo, d3roRadius, d3roShadow } from '@d3ro/ui/theme' import { useI18n } from '@d3ro/i18n' import type { RAGDocument, RAGQueryResult, RAGIndexProgress } from '@d3ro/core/types' +import { ErrorCode } from '@d3ro/core/errors' type KnowledgeView = 'kb' | 'insights' @@ -46,15 +47,21 @@ export function KnowledgeBasePage(): React.ReactElement { loadDocuments() } }) - const unsubComplete = window.electronAPI.rag.onIndexComplete(() => { + // 색인 실행이 끝나면(성공·실패 무관) 진행 표시를 내리고, 문서가 끝내 색인되지 않았으면 알린다. + const unsubComplete = window.electronAPI.rag.onIndexComplete((data) => { setIndexProgress(null) - loadDocuments() + void window.electronAPI.rag.getDocuments().then((resp) => { + if (!resp.success) return + setDocuments(resp.data) + const doc = resp.data.find((d) => d.id === data.documentId) + if (doc && !doc.indexed) setAddError(t('mobile.knowledge.error.indexUnavailable')) + }) }) return () => { unsubProgress() unsubComplete() } - }, [loadDocuments]) + }, [loadDocuments, t]) const isAddingRef = useRef(false) @@ -87,7 +94,10 @@ export function KnowledgeBasePage(): React.ReactElement { ) const handleReindex = useCallback(async (docId: string) => { - await window.electronAPI.rag.reindex(docId) + setAddError(null) + const resp = await window.electronAPI.rag.reindex(docId) + // 임베딩 실패는 rag:indexComplete 뒤 문서 상태로 알린다. 그 밖의 실패(문서 없음·파싱 실패)만 여기서 보인다. + if (!resp.success && resp.error.code !== ErrorCode.RAGEmbeddingFailed) setAddError(resp.error.message) }, []) const handleQuery = useCallback(async () => { diff --git a/apps/desktop/tests/main/services/rag-redteam-r2-0.test.ts b/apps/desktop/tests/main/services/rag-redteam-r2-0.test.ts new file mode 100644 index 0000000..9059bfd --- /dev/null +++ b/apps/desktop/tests/main/services/rag-redteam-r2-0.test.ts @@ -0,0 +1,396 @@ +// 레드팀 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) + }) +}) diff --git a/apps/desktop/tests/main/sync/sync-redteam-r2-0.test.ts b/apps/desktop/tests/main/sync/sync-redteam-r2-0.test.ts new file mode 100644 index 0000000..7529229 --- /dev/null +++ b/apps/desktop/tests/main/sync/sync-redteam-r2-0.test.ts @@ -0,0 +1,281 @@ +// 레드팀 r2-0: 원격 삭제(tombstone) 반영의 복원·보존 기간 경계, 지식 문서 청크 완전성. + +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { eq } from 'drizzle-orm' +import { createTestDb } from '../../helpers/createTestDb' +import { FakeSyncRemote } from '../../helpers/fakeSyncRemote' +import { bindTestDatabase, unbindTestDatabase } from '../../../src/main/db' +import { history, ragChunks, ragDocuments } from '../../../src/main/db/schema' +import { initInMemoryConfig, resetInMemoryConfig } from '../../../src/main/services/ConfigService' +import { + getCustomInstructionService, + resetCustomInstructionServiceForTests, +} from '../../../src/main/services/CustomInstructionService' +import { resetDictationTemplateServiceForTests } from '../../../src/main/services/DictationTemplateService' +import { resetMeetingDocTemplateServiceForTests } from '../../../src/main/services/MeetingDocTemplateService' +import { resetRAGServiceForTests } from '../../../src/main/services/RAGService' +import { SyncEngine } from '../../../src/main/services/sync/SyncEngine' +import { completeKnowledgeChunks } from '../../../src/main/services/sync/sync-adapters' +import { enqueueChange, pendingOps } from '../../../src/main/services/sync/sync-outbox' +import { + TOMBSTONE_RETENTION_MS, + earliestDeletedAt, + isRestoredAfter, + remoteTimeAfter, + tombstoneWindowExpired, +} from '../../../src/main/services/sync/tombstone-policy' +import type { EmbeddingPort } from '../../../src/main/services/rag/embedding-port' + +const USER = '11111111-1111-4111-8111-111111111111' +const DAY = 24 * 60 * 60 * 1000 + +let testDb: ReturnType +let remote: FakeSyncRemote +let engine: SyncEngine +let clock: number + +const offlineEmbedder: EmbeddingPort = { + model: 'test-embed', + ensureModel: () => Promise.reject(new Error('no embedding server in tests')), + embed: () => Promise.reject(new Error('no embedding server in tests')), +} + +function makeEngine(): SyncEngine { + return new SyncEngine({ remote, userId: USER, now: () => clock }) +} + +beforeEach(() => { + testDb = createTestDb() + bindTestDatabase(testDb.db, USER) + initInMemoryConfig() + resetCustomInstructionServiceForTests() + resetDictationTemplateServiceForTests() + resetMeetingDocTemplateServiceForTests() + resetRAGServiceForTests({ embedder: offlineEmbedder, notify: () => undefined, yieldMs: 0 }) + getCustomInstructionService().initialize() + remote = new FakeSyncRemote(USER) + clock = Date.parse('2026-09-27T00:00:00.000Z') + engine = makeEngine() +}) + +afterEach(() => { + vi.restoreAllMocks() + engine.dispose() + resetRAGServiceForTests() + unbindTestDatabase() + resetInMemoryConfig() + testDb.close() +}) + +function historyRow(id: string, text: string): Record { + return { id, original_text: text, duration: 1, mode: 'dictation', status: 'completed' } +} + +function localHistoryIds(): string[] { + return testDb.db.select({ id: history.id }).from(history).all().map((r) => r.id).sort() +} + +describe('tombstone 뒤 계정 보관본 복원', () => { + it('삭제를 이미 반영한 데스크톱: 복원된 행을 받고, 이후 pull의 겹쳐 읽기가 다시 지우지 않는다', async () => { + const x = crypto.randomUUID() + remote.mobileInsert('history', historyRow(x, 'keep me')) + await engine.runFullSync() + expect(localHistoryIds()).toEqual([x]) + + remote.mobileDelete('history', x) + await engine.pull() + expect(localHistoryIds()).toEqual([]) + + // restore_account_portability: 같은 id로 다시 넣고, 0037 이후 서버가 updated_at을 새로 찍는다. + remote.mobileInsert('history', historyRow(x, 'keep me')) + await engine.pull() + expect(localHistoryIds()).toEqual([x]) + + // tombstone 커서는 여전히 T1 근처라 다음 pull은 같은 tombstone을 다시 읽는다. + await engine.pull() + await engine.pull() + expect(localHistoryIds()).toEqual([x]) + }) + + it('삭제와 복원 사이 오프라인이던 데스크톱: 첫 pull에서 복원된 행을 지우지 않는다', async () => { + const x = crypto.randomUUID() + remote.mobileInsert('history', historyRow(x, 'v1')) + await engine.runFullSync() + + remote.mobileDelete('history', x) + remote.mobileInsert('history', historyRow(x, 'restored')) + const result = await engine.pull() + + expect(result.deleted).toBe(0) + expect(localHistoryIds()).toEqual([x]) + expect(testDb.db.select().from(history).where(eq(history.id, x)).get()?.originalText).toBe('restored') + }) + + it('runFullSync(push 전 tombstone + pull 안의 tombstone)도 복원된 행을 지우지 않는다', async () => { + const x = crypto.randomUUID() + remote.mobileInsert('history', historyRow(x, 'v1')) + await engine.runFullSync() + + remote.mobileDelete('history', x) + remote.mobileInsert('history', historyRow(x, 'restored')) + await engine.runFullSync() + await engine.runFullSync() + expect(localHistoryIds()).toEqual([x]) + }) + + it('복원되지 않은 삭제는 그대로 반영하고, 로컬 대기 변경보다 우선한다', async () => { + const kept = crypto.randomUUID() + const gone = crypto.randomUUID() + remote.mobileInsert('history', historyRow(kept, 'k')) + remote.mobileInsert('history', historyRow(gone, 'g')) + await engine.runFullSync() + + enqueueChange('history', gone, 'upsert') + remote.mobileDelete('history', gone) + remote.mobileDelete('history', kept) + remote.mobileInsert('history', historyRow(kept, 'k restored')) + await engine.pull() + + expect(localHistoryIds()).toEqual([kept]) + expect(pendingOps('history').has(gone)).toBe(false) + }) +}) + +describe('tombstone 보존 기간(180일) 초과', () => { + function pruneAllTombstones(): void { + remote.rows('sync_tombstones').splice(0) + } + + it('보존 기간보다 오래 쉰 기기는 서버에 없는 로컬 행을 지워 전체 대조한다', async () => { + const stale = crypto.randomUUID() + const alive = crypto.randomUUID() + remote.mobileInsert('history', historyRow(stale, 'deleted on phone long ago')) + remote.mobileInsert('history', historyRow(alive, 'still here')) + await engine.runFullSync() + // 다음 tombstone 커서가 생기도록 다른 행 하나를 지워 둔다 + const other = crypto.randomUUID() + remote.mobileInsert('history', historyRow(other, 'other')) + await engine.pull() + remote.mobileDelete('history', other) + await engine.pull() + + remote.mobileDelete('history', stale) + pruneAllTombstones() + clock += TOMBSTONE_RETENTION_MS + DAY + + const result = await engine.pull() + expect(result.errors).toEqual([]) + expect(localHistoryIds()).toEqual([alive]) + expect(result.deleted).toBe(1) + + // 편집해도 지운 행이 서버에 되살아나지 않는다(로컬에 없다) + expect(remote.find('history', stale)).toBeUndefined() + }) + + it('아직 올리지 않은 로컬 변경이 있는 행은 전체 대조가 지우지 않는다', async () => { + const seed = crypto.randomUUID() + remote.mobileInsert('history', historyRow(seed, 'seed')) + await engine.runFullSync() + remote.mobileDelete('history', seed) + await engine.pull() + + const local = crypto.randomUUID() + const at = clock + testDb.db.insert(history).values({ id: local, originalText: 'offline note', duration: 1, createdAt: at, updatedAt: at }).run() + enqueueChange('history', local, 'upsert') + pruneAllTombstones() + clock += TOMBSTONE_RETENTION_MS + DAY + + await engine.pull() + expect(localHistoryIds()).toEqual([local]) + }) + + it('보존 기간 안이면 전체 대조를 하지 않는다(서버에 없는 행을 함부로 지우지 않는다)', async () => { + const seed = crypto.randomUUID() + remote.mobileInsert('history', historyRow(seed, 'seed')) + await engine.runFullSync() + remote.mobileDelete('history', seed) + await engine.pull() + + const local = crypto.randomUUID() + testDb.db.insert(history).values({ id: local, originalText: 'x', duration: 1, createdAt: clock, updatedAt: clock }).run() + clock += 30 * DAY + await engine.pull() + expect(localHistoryIds()).toEqual([local]) + }) + + it('한 번 대조하면 다음 pull부터는 다시 대조하지 않는다', async () => { + const seed = crypto.randomUUID() + remote.mobileInsert('history', historyRow(seed, 'seed')) + await engine.runFullSync() + remote.mobileDelete('history', seed) + await engine.pull() + clock += TOMBSTONE_RETENTION_MS + DAY + await engine.pull() + + const fresh = crypto.randomUUID() + testDb.db.insert(history).values({ id: fresh, originalText: 'y', duration: 1, createdAt: clock, updatedAt: clock }).run() + clock += DAY + await engine.pull() + expect(localHistoryIds()).toEqual([fresh]) + }) +}) + +describe('tombstone 정책 (순수 함수)', () => { + it('서버 시각 비교는 마이크로초까지 본다', () => { + expect(remoteTimeAfter('2026-09-27T00:00:00.000200+00:00', '2026-09-27T00:00:00.000100+00:00')).toBe(true) + expect(remoteTimeAfter('2026-09-27T00:00:01Z', '2026-09-27T00:00:02Z')).toBe(false) + expect(remoteTimeAfter('bad', '2026-09-27T00:00:02Z')).toBe(false) + }) + + it('복원 판정과 가장 이른 삭제 시각', () => { + const ref = { rowId: 'a', deletedAt: '2026-09-27T00:00:05Z' } + expect(isRestoredAfter(ref, '2026-09-27T00:00:06Z')).toBe(true) + expect(isRestoredAfter(ref, '2026-09-27T00:00:04Z')).toBe(false) + expect(isRestoredAfter({ rowId: 'a', deletedAt: null }, '2026-09-27T00:00:06Z')).toBe(false) + expect(earliestDeletedAt([ref, { rowId: 'b', deletedAt: '2026-09-27T00:00:01Z' }, { rowId: 'c', deletedAt: null }])).toBe( + '2026-09-27T00:00:01Z' + ) + }) + + it('보존 기간 창', () => { + const now = Date.parse('2026-09-27T00:00:00Z') + expect(tombstoneWindowExpired(null, now)).toBe(false) + expect(tombstoneWindowExpired(now - 30 * DAY, now)).toBe(false) + expect(tombstoneWindowExpired(now - 179 * DAY, now)).toBe(true) + }) +}) + +describe('지식 문서 청크 완전성', () => { + it('문서 행의 chunk_count보다 적은 청크만 올라와 있으면 미루고, 다 올라오면 전부 저장한다', async () => { + const id = crypto.randomUUID() + remote.mobileInsert('knowledge_documents', { id, title: 'Big', file_name: 'big.txt', file_type: 'txt', chunk_count: 3 }) + remote.rows('knowledge_chunks').push( + { id: crypto.randomUUID(), document_id: id, chunk_index: 0, content: 'c0' }, + { id: crypto.randomUUID(), document_id: id, chunk_index: 1, content: 'c1' } + ) + await engine.runFullSync() + expect(testDb.db.select().from(ragDocuments).where(eq(ragDocuments.id, id)).get()).toBeUndefined() + + remote.rows('knowledge_chunks').push({ id: crypto.randomUUID(), document_id: id, chunk_index: 2, content: 'c2' }) + await engine.pull() + const doc = testDb.db.select().from(ragDocuments).where(eq(ragDocuments.id, id)).get() + expect(doc?.chunkCount).toBe(3) + const chunks = testDb.db.select().from(ragChunks).where(eq(ragChunks.documentId, id)).all() + expect(chunks.sort((a, b) => a.chunkIndex - b.chunkIndex).map((c) => c.content)).toEqual(['c0', 'c1', 'c2']) + }) + + it('completeKnowledgeChunks: 개수·연속 index를 확인한다', () => { + const rows = [ + { chunk_index: 1, content: 'b' }, + { chunk_index: 0, content: 'a' }, + ] + expect(completeKnowledgeChunks(rows, 2)).toEqual(['a', 'b']) + expect(completeKnowledgeChunks(rows, 3)).toBeNull() + expect(completeKnowledgeChunks(rows, null)).toEqual(['a', 'b']) + expect(completeKnowledgeChunks([{ chunk_index: 0, content: 'a' }, { chunk_index: 2, content: 'c' }], null)).toBeNull() + expect(completeKnowledgeChunks([{ chunk_index: 0, content: 'a' }, { chunk_index: 0, content: 'a' }], 2)).toBeNull() + expect(completeKnowledgeChunks([], 0)).toBeNull() + expect(completeKnowledgeChunks([{ chunk_index: 0, content: null }], 1)).toBeNull() + }) +}) diff --git a/packages/core/src/ipc-channels.ts b/packages/core/src/ipc-channels.ts index e8e6af9..e88f3f1 100644 --- a/packages/core/src/ipc-channels.ts +++ b/packages/core/src/ipc-channels.ts @@ -332,7 +332,10 @@ export const IPC_CHANNELS = { REINDEX: 'rag:reindex', // Main → Renderer events INDEX_PROGRESS: 'rag:indexProgress', + /** 색인 실행이 끝났다(성공·실패 무관). 실패면 INDEX_FAILED 가 먼저 온다 — 결과는 문서의 indexed 로 본다. */ INDEX_COMPLETE: 'rag:indexComplete', + /** 색인 실패 { documentId, fileName, code, message } (임베딩 모델 없음·서버 없음 등) */ + INDEX_FAILED: 'rag:indexFailed', QUERY_RESULT: 'rag:queryResult', },