From a060c3fd05c2737728fa4e18fa7e0b88438791c8 Mon Sep 17 00:00:00 2001 From: infinite-vector Date: Wed, 7 Oct 2026 13:19:03 +0800 Subject: [PATCH 1/2] fix(knowledgeBase): resolve coordinator AB-BA deadlock, implement progress heartbeat and dynamic DB proxy --- KnowledgeBaseManager.js | 4 + Plugin/AgentDream/DreamWaveEngine.js | 5 +- modules/knowledgeBase/databaseCoordinator.js | 25 ++- modules/knowledgeBase/indexRepository.js | 201 ++++++++++--------- modules/knowledgeBase/ingestionPipeline.js | 26 ++- 5 files changed, 165 insertions(+), 96 deletions(-) diff --git a/KnowledgeBaseManager.js b/KnowledgeBaseManager.js index 331673347..0cb6601b8 100644 --- a/KnowledgeBaseManager.js +++ b/KnowledgeBaseManager.js @@ -252,6 +252,7 @@ class KnowledgeBaseManager { this.lastJsWriteFinishedAt = 0; this.lastRustWriteFinishedAt = 0; this._rustLeaseWaitLogAt = 0; + this.lastActivityAt = Date.now(); // 🧭 外部文件写入协调器(DailyNote 等常驻服务使用) // 文件变更本身不直接写 SQLite,但必须与 watcher 批处理、Rust SQLite 恢复形成单一时序。 @@ -841,6 +842,9 @@ class KnowledgeBaseManager { _delay(ms) { return this.databaseCoordinator.delay(ms); } + touchActivity() { + this.lastActivityAt = Date.now(); + } async _waitForDatabaseCoordinatorIdle(options = {}) { return this.databaseCoordinator.waitForIdle(options); diff --git a/Plugin/AgentDream/DreamWaveEngine.js b/Plugin/AgentDream/DreamWaveEngine.js index f8c40a75f..50b92a82a 100644 --- a/Plugin/AgentDream/DreamWaveEngine.js +++ b/Plugin/AgentDream/DreamWaveEngine.js @@ -24,7 +24,10 @@ const MID_EXPAND_MAX = 180; // 中期最多放宽到180天 class DreamWaveEngine { constructor(knowledgeBaseManager) { this.kb = knowledgeBaseManager; - this.db = knowledgeBaseManager ? knowledgeBaseManager.db : null; + } + + get db() { + return this.kb?.db || null; } // ========================================================================= diff --git a/modules/knowledgeBase/databaseCoordinator.js b/modules/knowledgeBase/databaseCoordinator.js index 11fe82d9c..7d4ace7fb 100644 --- a/modules/knowledgeBase/databaseCoordinator.js +++ b/modules/knowledgeBase/databaseCoordinator.js @@ -20,9 +20,22 @@ class DatabaseCoordinator { 1000, Number(options.timeoutMs) || 30 * 60 * 1000 ); + const rawStall = Number(options.stallThresholdMs); + const stallThresholdMs = Number.isFinite(rawStall) && rawStall > 0 + ? Math.max(10, Math.floor(rawStall)) + : 120 * 1000; const pollMs = Math.max(10, Number(options.pollMs) || 50); const startedAt = Date.now(); + if ( + owner.databaseCorruptionDetected + || owner.dbHealthState === 'corrupt' + ) { + throw new Error( + 'KnowledgeBase database is unavailable because ' + + 'corruption was detected.' + ); + } while ( (!options.allowJsProcessing && owner.isProcessing) || (!options.allowJsDeleteProcessing && owner.isProcessingDeletes) @@ -48,10 +61,18 @@ class DatabaseCoordinator { error.code = 'ABORT_ERR'; throw error; } - if (Date.now() - startedAt >= timeoutMs) { + const lastActive = owner.lastActivityAt || startedAt; + const timeSinceLastActivity = Date.now() - lastActive; + const isStalled = timeSinceLastActivity >= stallThresholdMs; + const isAbsoluteTimeout = (Date.now() - startedAt) >= timeoutMs; + + if (isStalled || isAbsoluteTimeout) { + const reason = isStalled + ? `activity stalled with zero progress for ${Math.round(timeSinceLastActivity / 1000)}s` + : `reached absolute limit of ${Math.round(timeoutMs / 1000)}s`; throw new Error( `Timed out waiting for KnowledgeBase coordinator after ` - + `${timeoutMs}ms (processing=${owner.isProcessing}, ` + + `${Date.now() - startedAt}ms (${reason}) (processing=${owner.isProcessing}, ` + `deletes=${owner.isProcessingDeletes}, ` + `externalMutation=${owner.externalMutationOwner || 'none'}, ` + `rustLease=${owner.rustWriteLease?.owner || 'none'}, ` diff --git a/modules/knowledgeBase/indexRepository.js b/modules/knowledgeBase/indexRepository.js index 859281b3e..8402f3427 100644 --- a/modules/knowledgeBase/indexRepository.js +++ b/modules/knowledgeBase/indexRepository.js @@ -587,100 +587,117 @@ class IndexRepository { || name.endsWith('簇'); } + async _executeLoadIndex(diaryName) { + const persist = this.shouldPersist(diaryName); + console.log( + `[${this.logPrefix}] 📂 Loading index for diary: ` + + `"${diaryName}" (Persist: ${persist})` + ); + const safeName = crypto.createHash('md5') + .update(diaryName) + .digest('hex'); + const fileName = `diary_${safeName}`; + const capacity = 50000; + let index; + if (persist) { + if (this.chunkIndexPersistenceMode === 'generational') { + const baselineRestore = await this.loadDiaryBaseline(diaryName, capacity); + if (baselineRestore?.index) { + index = baselineRestore.index; + } else { + console.log( + `[${this.logPrefix}] 🔄 No valid generational baseline for "${diaryName}", ` + + 'rebuilding from SQLite and publishing initial baseline...' + ); + index = new this.VexusIndex(this.config.dimension, capacity); + await this.recoverFromDb(index, 'chunks', diaryName); + this.diaryIndices.set(diaryName, index); + this.publishDiaryBaseline(diaryName, { force: true }); + } + } else { + index = await this.loadOrBuild( + fileName, + capacity, + 'chunks', + diaryName + ); + } + } else { + index = new this.VexusIndex( + this.config.dimension, + capacity + ); + await this.recoverFromDb(index, 'chunks', diaryName); + } + return index; + } + async getOrLoad(diaryName, options = {}) { - this.lastUsed.set(diaryName, Date.now()); - if (this.diaryIndices.has(diaryName)) { - return this.diaryIndices.get(diaryName); + const name = String(diaryName || '').trim(); + this.lastUsed.set(name, Date.now()); + if (this.diaryIndices.has(name)) { + return this.diaryIndices.get(name); } - if (this.loadPromises.has(diaryName)) { - return this.loadPromises.get(diaryName); + if (this.loadPromises.has(name)) { + return this.loadPromises.get(name); } - const load = async () => { - await this.waitForCoordinatorIdle(options); - this.recoveryActive = true; - this.onRecoveryStateChange(true); - try { - if (this.diaryIndices.has(diaryName)) { - return this.diaryIndices.get(diaryName); - } - const persist = this.shouldPersist(diaryName); - console.log( - `[${this.logPrefix}] 📂 Loading index for diary: ` + - `"${diaryName}" (Persist: ${persist})` - ); - const safeName = crypto.createHash('md5') - .update(diaryName) - .digest('hex'); - const fileName = `diary_${safeName}`; - const capacity = 50000; - let index; - if (persist) { - if (this.chunkIndexPersistenceMode === 'generational') { - const baselineRestore = await this.loadDiaryBaseline(diaryName, capacity); - if (baselineRestore?.index) { - index = baselineRestore.index; - } else { - console.log( - `[${this.logPrefix}] 🔄 No valid generational baseline for "${diaryName}", ` + - 'rebuilding from SQLite and publishing initial baseline...' - ); - index = new this.VexusIndex(this.config.dimension, capacity); - await this.recoverFromDb(index, 'chunks', diaryName); - this.diaryIndices.set(diaryName, index); - this.publishDiaryBaseline(diaryName, { force: true }); + const execute = async () => { + if (!options.bypassCoordinator) { + await this.waitForCoordinatorIdle(options); + } + if (this.diaryIndices.has(name)) { + return this.diaryIndices.get(name); + } + + const load = async () => { + this.recoveryActive = true; + this.onRecoveryStateChange(true); + try { + if (this.diaryIndices.has(name)) { + return this.diaryIndices.get(name); + } + const index = await this._executeLoadIndex(name); + this.diaryIndices.set(name, index); + try { + this.onDiaryIndexPublished(name, index); + } catch (error) { + if (this.diaryIndices.get(name) === index) { + this.diaryIndices.delete(name); } - } else { - index = await this.loadOrBuild( - fileName, - capacity, - 'chunks', - diaryName + this.lastUsed.delete(name); + throw new Error( + `Diary index loaded but native publication failed for ` + + `"${name}": ${error.message}` ); } - } else { - index = new this.VexusIndex( - this.config.dimension, - capacity - ); - await this.recoverFromDb(index, 'chunks', diaryName); - } - this.diaryIndices.set(diaryName, index); - try { - this.onDiaryIndexPublished(diaryName, index); - } catch (error) { - if (this.diaryIndices.get(diaryName) === index) { - this.diaryIndices.delete(diaryName); - } - this.lastUsed.delete(diaryName); - throw new Error( - `Diary index loaded but native publication failed for ` + - `"${diaryName}": ${error.message}` - ); + this.ensureDiaryDateIndex(name); + return index; + } finally { + this.recoveryActive = false; + this.onRecoveryStateChange(false); } - this.ensureDiaryDateIndex(diaryName); - return index; - } finally { - this.recoveryActive = false; - this.onRecoveryStateChange(false); - } + }; + + const queued = this.recoveryTail.then(load); + this.recoveryTail = queued.catch(error => { + console.error( + `[${this.logPrefix}] Serialized index load failed for ` + + `"${name}":`, + error + ); + }); + this.onRecoveryTailChange(this.recoveryTail); + return await queued; }; - const queued = this.recoveryTail.then(load); - this.recoveryTail = queued.catch(error => { - console.error( - `[${this.logPrefix}] Serialized index load failed for ` + - `"${diaryName}":`, - error - ); - }); - this.onRecoveryTailChange(this.recoveryTail); - this.loadPromises.set(diaryName, queued); + const task = execute(); + this.loadPromises.set(name, task); try { - return await queued; + return await task; } finally { - if (this.loadPromises.get(diaryName) === queued) { - this.loadPromises.delete(diaryName); + if (this.loadPromises.get(name) === task) { + this.loadPromises.delete(name); } } } @@ -757,10 +774,6 @@ class IndexRepository { */ async applyChunkDelta(diaryName, removeIds = [], upserts = []) { const normalizedDiaryName = String(diaryName || '').trim(); - if (!normalizedDiaryName) { - throw new TypeError('applyChunkDelta requires a diary name'); - } - const deletes = [...new Set( (Array.isArray(removeIds) ? removeIds : []) .map(Number) @@ -786,12 +799,18 @@ class IndexRepository { requestedUpserts: 0 }; } + if (!normalizedDiaryName) { + throw new TypeError('applyChunkDelta requires a diary name'); + } - const index = await this.getOrLoad(normalizedDiaryName, { - allowJsProcessing: true, - allowJsDeleteProcessing: true - }); - + let index = this.diaryIndices.get(normalizedDiaryName); + if (!index) { + index = await this.getOrLoad(normalizedDiaryName, { + allowJsProcessing: true, + allowJsDeleteProcessing: true, + bypassCoordinator: true + }); + } if (typeof index?.applyChunkDelta === 'function') { const ids = normalizedUpserts.map(entry => entry.id); const vectors = new Float32Array(ids.length * this.config.dimension); diff --git a/modules/knowledgeBase/ingestionPipeline.js b/modules/knowledgeBase/ingestionPipeline.js index 3c7fb49af..e838e1c5c 100644 --- a/modules/knowledgeBase/ingestionPipeline.js +++ b/modules/knowledgeBase/ingestionPipeline.js @@ -62,6 +62,7 @@ async _flushDeleteBatch() { return; } this.isProcessingDeletes = true; + this.touchActivity?.(); let retryDelayMs = 0; const batchFiles = Array.from(this.pendingDeletes).slice(0, this.config.maxDeleteBatchSize); @@ -73,6 +74,7 @@ async _flushDeleteBatch() { try { await this._handleDeleteBatch(batchFiles); batchFiles.forEach(f => this.pendingDeletes.delete(f)); + this.touchActivity?.(); } catch (e) { console.error('[KnowledgeBase] ❌ Delete batch failed:', e); if (this._isSqliteCorruptionError(e)) { @@ -118,6 +120,7 @@ async _flushBatch() { return; } this.isProcessing = true; + this.touchActivity?.(); let retryDelayMs = 0; // 1. 📋 准备批次:先从队列中取出,但不立即永久删除 @@ -243,6 +246,7 @@ async _flushBatch() { chunkVectors = await getEmbeddingsBatch(texts, embeddingConfig); // 🛡️ getEmbeddingsBatch 现在保证 chunkVectors.length === texts.length // 失败/超长的位置为 null,后续写入 DB 时会跳过这些 null 向量 + this.touchActivity?.(); } let tagVectors = []; @@ -254,6 +258,7 @@ async _flushBatch() { // 同样保证长度对齐,null 表示失败 tagVectors.push(...batchVectors); } + this.touchActivity?.(); } // 4. 写入 DB 和 索引 @@ -359,11 +364,15 @@ async _flushBatch() { return { updates, tagUpdates, deletions, newTagIds }; }); + const _tPhaseStart = Date.now(); const { updates, tagUpdates, deletions, newTagIds } = transaction(); + const _tSqliteMs = Date.now() - _tPhaseStart; + this.touchActivity?.(); // 每个日记本只发布一次排他 Chunk 差分。SQLite 已经提交权威事实; // Rust 在同一 RwLock 写锁内完成该日记本全部旧 ID 删除和新 ID upsert, // 并发 search 只能看到批次前或批次后状态,不能观察到公共索引半批状态。 + const _tRustStart = Date.now(); const affectedDiaryNames = new Set([ ...deletions.keys(), ...updates.keys() @@ -381,9 +390,12 @@ async _flushBatch() { `revision=${deltaResult.revision ?? 'recovered'}.` ); } + const _tRustMs = Date.now() - _tRustStart; + this.touchActivity?.(); // Tag 索引已有独立 applyTagDelta 路径的生命周期;这里保留现有兼容 // upsert,日记公共索引的读写原子性不再依赖这段逻辑。 + const _tTagStart = Date.now(); tagUpdates.forEach(u => { try { this.tagIndex.add(u.id, u.vec); @@ -413,13 +425,23 @@ async _flushBatch() { this._ensureDiaryDateIndexCached(dName); } } - - console.log(`[KnowledgeBase] ✅ Batch complete. Updated ${updates.size} diary indices.`); + const _tTagAndDateMs = Date.now() - _tTagStart; // 数据更新后,检查是否需要重建 V9.1 矩阵(防抖 + 阈值)。 // 使用“成功新增的唯一 tag id”累计触发 1% 阈值; // file_tags 组关系仍是共现矩阵真相,但不再作为“新增 1% tag”的计数依据。 + const _tMatrixStart = Date.now(); if (this.tagMemoEngine) this.tagMemoEngine.scheduleMatrixRebuildForNewTags(newTagIds); + const _tMatrixMs = Date.now() - _tMatrixStart; + + console.log( + `[KnowledgeBase] ⏱️ Phase Breakdown: ` + + `SQLite=${_tSqliteMs}ms, ` + + `applyChunkDelta=${_tRustMs}ms, ` + + `Tag&Date=${_tTagAndDateMs}ms, ` + + `scheduleMatrix=${_tMatrixMs}ms (Total=${Date.now() - _tPhaseStart}ms)` + ); + console.log(`[KnowledgeBase] ✅ Batch complete. Updated ${updates.size} diary indices.`); } catch (e) { console.error('[KnowledgeBase] ❌ Batch processing failed catastrophically.'); From 5cf71f4544dd11313bdca8c99a06290e9e60d65a Mon Sep 17 00:00:00 2001 From: infinite-vector Date: Wed, 7 Oct 2026 13:37:30 +0800 Subject: [PATCH 2/2] fix(knowledgeBase): prevent 0ms false-positive stall on cold start and refresh heartbeat on recovery --- KnowledgeBaseManager.js | 3 +++ modules/knowledgeBase/databaseCoordinator.js | 5 ++++- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/KnowledgeBaseManager.js b/KnowledgeBaseManager.js index 0cb6601b8..ef1480f78 100644 --- a/KnowledgeBaseManager.js +++ b/KnowledgeBaseManager.js @@ -306,6 +306,9 @@ class KnowledgeBaseManager { this._unregisterNativeDiaryIndex(diaryName), onRecoveryStateChange: active => { this.indexRecoveryActive = active; + if (active) { + this.touchActivity(); + } }, onRecoveryTailChange: tail => { this._indexRecoveryTail = tail; diff --git a/modules/knowledgeBase/databaseCoordinator.js b/modules/knowledgeBase/databaseCoordinator.js index 7d4ace7fb..6694aa73b 100644 --- a/modules/knowledgeBase/databaseCoordinator.js +++ b/modules/knowledgeBase/databaseCoordinator.js @@ -61,7 +61,10 @@ class DatabaseCoordinator { error.code = 'ABORT_ERR'; throw error; } - const lastActive = owner.lastActivityAt || startedAt; + const lastActive = Math.max( + startedAt, + Number(owner.lastActivityAt) || 0 + ); const timeSinceLastActivity = Date.now() - lastActive; const isStalled = timeSinceLastActivity >= stallThresholdMs; const isAbsoluteTimeout = (Date.now() - startedAt) >= timeoutMs;