fix(sync,rag): keep restored rows, reconcile after tombstone pruning, reject partial knowledge docs, surface RAG indexing failures
This commit is contained in:
parent
bbaf1e0a99
commit
3a46437f28
13 changed files with 1630 additions and 448 deletions
|
|
@ -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<Buffer>
|
||||
fileExists: (filePath: string) => boolean
|
||||
/** Pro+ 라이선스 확인. 허용되지 않으면 D3ROError를 던진다. */
|
||||
assertLicensed: () => Promise<void>
|
||||
/** 청크 사이 양보 시간(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<void> {
|
||||
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<RAGDocument> {
|
||||
// 라이센스 체크
|
||||
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<string, RAGDocument['fileType']> = {
|
||||
'.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<void> {
|
||||
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<RAGQueryResult> {
|
||||
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<string> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<number[]> {
|
||||
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<string> {
|
||||
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<K extends keyof RAGServiceEvents>(event: K, listener: RAGServiceEvents[K]): this {
|
||||
return super.on(event, listener)
|
||||
}
|
||||
|
||||
override emit<K extends keyof RAGServiceEvents>(event: K, ...args: Parameters<RAGServiceEvents[K]>): boolean {
|
||||
return super.emit(event, ...args)
|
||||
}
|
||||
}
|
||||
|
||||
// ── 싱글톤 ──
|
||||
let instance: RAGService | null = null
|
||||
|
||||
export function resetRAGServiceForTests(): void {
|
||||
/** 테스트용: 싱글톤을 비운다. deps를 주면 그 포트로 새 인스턴스를 만든다. */
|
||||
export function resetRAGServiceForTests(deps?: Partial<RAGServiceDeps>): void {
|
||||
if (instance) instance.removeAllListeners()
|
||||
instance = null
|
||||
instance = deps ? new RAGService({ ...defaultRAGServiceDeps(), ...deps }) : null
|
||||
}
|
||||
|
||||
export function getRAGService(): RAGService {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue