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
7 changes: 7 additions & 0 deletions KnowledgeBaseManager.js
Original file line number Diff line number Diff line change
Expand Up @@ -252,6 +252,7 @@ class KnowledgeBaseManager {
this.lastJsWriteFinishedAt = 0;
this.lastRustWriteFinishedAt = 0;
this._rustLeaseWaitLogAt = 0;
this.lastActivityAt = Date.now();

// 🧭 外部文件写入协调器(DailyNote 等常驻服务使用)
// 文件变更本身不直接写 SQLite,但必须与 watcher 批处理、Rust SQLite 恢复形成单一时序。
Expand Down Expand Up @@ -305,6 +306,9 @@ class KnowledgeBaseManager {
this._unregisterNativeDiaryIndex(diaryName),
onRecoveryStateChange: active => {
this.indexRecoveryActive = active;
if (active) {
this.touchActivity();
}
},
onRecoveryTailChange: tail => {
this._indexRecoveryTail = tail;
Expand Down Expand Up @@ -841,6 +845,9 @@ class KnowledgeBaseManager {
_delay(ms) {
return this.databaseCoordinator.delay(ms);
}
touchActivity() {
this.lastActivityAt = Date.now();
}

async _waitForDatabaseCoordinatorIdle(options = {}) {
return this.databaseCoordinator.waitForIdle(options);
Expand Down
5 changes: 4 additions & 1 deletion Plugin/AgentDream/DreamWaveEngine.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

// =========================================================================
Expand Down
28 changes: 26 additions & 2 deletions modules/knowledgeBase/databaseCoordinator.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -48,10 +61,21 @@ class DatabaseCoordinator {
error.code = 'ABORT_ERR';
throw error;
}
if (Date.now() - startedAt >= timeoutMs) {
const lastActive = Math.max(
startedAt,
Number(owner.lastActivityAt) || 0
);
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'}, `
Expand Down
201 changes: 110 additions & 91 deletions modules/knowledgeBase/indexRepository.js
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
}
Expand Down Expand Up @@ -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)
Expand All @@ -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);
Expand Down
Loading
Loading