Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 18 additions & 24 deletions src/main/ipc/knowledgeHandlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,13 +41,13 @@ export function registerKnowledgeHandlers(knowledgeService: KnowledgeService) {
})

try {
const documentId = await knowledgeService.addDocument(
const documentId = knowledgeService.enqueueContentDocument(
params.notebookId,
params.options,
(stage, progress) => {
// 发送进度更新
event.sender.send('knowledge:index-progress', {
(id, stage, progress) => {
broadcastIndexProgress({
notebookId: params.notebookId,
documentId: id,
stage,
progress
})
Expand Down Expand Up @@ -101,12 +101,13 @@ export function registerKnowledgeHandlers(knowledgeService: KnowledgeService) {
Logger.debug('KnowledgeHandlers', 'add-document-from-url:', params)

try {
const documentId = await knowledgeService.addDocumentFromUrl(
const documentId = knowledgeService.enqueueDocumentFromUrl(
params.notebookId,
params.url,
(stage, progress) => {
event.sender.send('knowledge:index-progress', {
(id, stage, progress) => {
broadcastIndexProgress({
notebookId: params.notebookId,
documentId: id,
stage,
progress
})
Expand All @@ -129,12 +130,13 @@ export function registerKnowledgeHandlers(knowledgeService: KnowledgeService) {
Logger.debug('KnowledgeHandlers', 'add-note:', params)

try {
const documentId = await knowledgeService.addNoteToKnowledge(
const documentId = knowledgeService.enqueueNoteDocument(
params.notebookId,
params.noteId,
(stage, progress) => {
event.sender.send('knowledge:index-progress', {
(id, stage, progress) => {
broadcastIndexProgress({
notebookId: params.notebookId,
documentId: id,
stage,
progress
})
Expand Down Expand Up @@ -276,11 +278,7 @@ export function registerKnowledgeHandlers(knowledgeService: KnowledgeService) {

try {
await knowledgeService.reindexDocument(params.documentId, (stage, progress) => {
event.sender.send('knowledge:index-progress', {
documentId: params.documentId,
stage,
progress
})
broadcastIndexProgress({ documentId: params.documentId, stage, progress })
})
return { success: true }
} catch (error) {
Expand All @@ -299,11 +297,7 @@ export function registerKnowledgeHandlers(knowledgeService: KnowledgeService) {

try {
await knowledgeService.retryDocument(params.documentId, (stage, progress) => {
event.sender.send('knowledge:index-progress', {
documentId: params.documentId,
stage,
progress
})
broadcastIndexProgress({ documentId: params.documentId, stage, progress })
})
return { success: true }
} catch (error) {
Expand Down Expand Up @@ -365,8 +359,8 @@ export function registerKnowledgeHandlers(knowledgeService: KnowledgeService) {
const result = await knowledgeService.addFolder(
params.notebookId,
params.folderPath,
(stage, progress) => {
broadcastIndexProgress({ notebookId: params.notebookId, stage, progress })
(documentId, stage, progress) => {
broadcastIndexProgress({ notebookId: params.notebookId, documentId, stage, progress })
}
)

Expand Down Expand Up @@ -416,8 +410,8 @@ export function registerKnowledgeHandlers(knowledgeService: KnowledgeService) {
const result = knowledgeService.enqueueDocumentsFromPaths(
params.notebookId,
params.paths,
(stage, progress) => {
broadcastIndexProgress({ notebookId: params.notebookId, stage, progress })
(documentId, stage, progress) => {
broadcastIndexProgress({ notebookId: params.notebookId, documentId, stage, progress })
}
)
return { success: true, ...result }
Expand Down
171 changes: 168 additions & 3 deletions src/main/services/KnowledgeService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,18 @@ export interface SearchResult {
*/
export type IndexProgressCallback = (stage: string, progress: number) => void

/**
* 后台导入的进度回调(#176)。
*
* 带 `documentId`:一份 source 的进度要能落在它自己的那一行上,而不是只能让整个
* notebook 共用一个布尔值。
*/
export type DocumentIndexProgressCallback = (
documentId: string,
stage: string,
progress: number
) => void

/** 批量导入里被跳过的文件(#98):已经在同一个 notebook 里。 */
export interface BatchImportSkip {
path: string
Expand Down Expand Up @@ -350,7 +362,7 @@ export class KnowledgeService {
enqueueDocumentsFromPaths(
notebookId: string,
paths: readonly string[],
onProgress?: IndexProgressCallback
onProgress?: DocumentIndexProgressCallback
): BatchImportResult {
const db = getDatabase()
const existing = new Set(
Expand Down Expand Up @@ -379,7 +391,7 @@ export class KnowledgeService {
added.push(documentId)
this.ingestionQueue.enqueue({
documentId,
onProgress,
onProgress: (stage, progress) => onProgress?.(documentId, stage, progress),
run: async (jobProgress) => {
try {
await this.ingestPendingDocument(documentId, filePath, jobProgress)
Expand Down Expand Up @@ -466,7 +478,7 @@ export class KnowledgeService {
async addFolder(
notebookId: string,
folderPath: string,
onProgress?: IndexProgressCallback
onProgress?: DocumentIndexProgressCallback
): Promise<BatchImportResult> {
const scanned = await scanFolder(folderPath, this.fileParserService.supportedExtensions())
return this.enqueueDocumentsFromPaths(
Expand Down Expand Up @@ -920,6 +932,159 @@ export class KnowledgeService {
)
}

/**
* 统一后台导入(#176):登记 `pending` 行,解析 / 分块 / 嵌入在队列里继续。
*
* 内容已经可用的来源(粘贴文本 / 笔记)走这里;URL 要先抓取,见
* `enqueueDocumentFromUrl`。返回时列表里已经有这一行,Dialog 不必再等索引结束。
*/
enqueueContentDocument(
notebookId: string,
options: AddDocumentOptions,
onProgress?: DocumentIndexProgressCallback
): string {
const db = getDatabase()
const documentId = `doc_${Date.now()}_${Math.random().toString(36).slice(2, 9)}`
const now = new Date()
const contentHash = createHash('md5').update(options.content).digest('hex')

const newDoc: NewDocument = {
id: documentId,
notebookId,
title: options.title,
type: options.type,
sourceUri: options.sourceUri,
sourceNoteId: options.sourceNoteId,
content: options.content,
contentHash,
mimeType: options.mimeType,
fileSize: options.fileSize,
metadata: options.metadata,
status: 'pending',
chunkCount: 0,
createdAt: now,
updatedAt: now
}
db.insert(documents).values(newDoc).run()

this.ingestionQueue.enqueue({
documentId,
onProgress: (stage, progress) => onProgress?.(documentId, stage, progress),
run: async (jobProgress) => {
const runId = startRun(documentId, 'import')
try {
await this.indexDocument(
documentId,
runId,
options.content,
undefined,
{ chunkOptions: options.chunkOptions },
jobProgress
)
completeRun(runId)
} catch (error) {
failRun(runId, (error as Error).message)
jobProgress('failed', 100)
throw error
}
}
})

return documentId
}

/**
* 从 Note 后台导入。空笔记在登记前就拒绝 —— 那不是延迟反馈,是输入本身不合法。
*/
enqueueNoteDocument(
notebookId: string,
noteId: string,
onProgress?: DocumentIndexProgressCallback
): string {
const note = getDatabase().select().from(notes).where(eq(notes.id, noteId)).get()
if (!note) throw new Error(`Note ${noteId} not found`)

if (note.content.trim().length === 0) {
throw new Error('Note content is empty. Cannot add empty note to knowledge base.')
}

return this.enqueueContentDocument(
notebookId,
{
title: note.title,
type: 'note',
content: note.content,
sourceNoteId: noteId
},
onProgress
)
}

/**
* 从 URL 后台导入:先登记 `pending` 行,抓取 / 解析 / 嵌入都在队列里。
*
* URL 校验(协议、可解析)应在调用前完成;这里只负责登记与后台处理。
*/
enqueueDocumentFromUrl(
notebookId: string,
url: string,
onProgress?: DocumentIndexProgressCallback
): string {
const db = getDatabase()
const documentId = `doc_${Date.now()}_${Math.random().toString(36).slice(2, 9)}`
const now = new Date()

const newDoc: NewDocument = {
id: documentId,
notebookId,
title: url,
type: 'url',
sourceUri: url,
status: 'pending',
chunkCount: 0,
createdAt: now,
updatedAt: now
}
db.insert(documents).values(newDoc).run()

this.ingestionQueue.enqueue({
documentId,
onProgress: (stage, progress) => onProgress?.(documentId, stage, progress),
run: async (jobProgress) => {
const runId = startRun(documentId, 'import')
try {
jobProgress('fetching_url', 0)
const fetchResult = await this.webFetchService.fetchUrl(url)
const content = fetchResult.content

db.update(documents)
.set({
title: fetchResult.title || url,
content,
contentHash: createHash('md5').update(content).digest('hex'),
mimeType: fetchResult.mimeType,
metadata: {
...fetchResult.metadata,
description: fetchResult.description
},
updatedAt: new Date()
})
.where(eq(documents.id, documentId))
.run()

await this.indexDocument(documentId, runId, content, undefined, {}, jobProgress)
completeRun(runId)
} catch (error) {
failRun(runId, (error as Error).message)
jobProgress('failed', 100)
throw error
}
}
})

return documentId
}

/**
* 索引前把 chunk 转成向量。
*
Expand Down
Loading
Loading