diff --git a/src/modules/extraction/engine/extract-sections.ts b/src/modules/extraction/engine/extract-sections.ts index da2f1ef..425b5fb 100644 --- a/src/modules/extraction/engine/extract-sections.ts +++ b/src/modules/extraction/engine/extract-sections.ts @@ -10,6 +10,7 @@ import type { ProjectMetadata } from './extract-project-metadata' import type { ScannedFile } from './file-scan' import { initTreeSitterWithSelectedGrammars } from './grammar-loader' import persistPreparedFileSections from './persist-prepared-file-sections' +import type { PreparedFile } from './prepare-file-sections' import prepareFileSections from './prepare-file-sections' import rebuildModuleArtifacts, { type ModuleArtifactKey, @@ -19,6 +20,7 @@ import rebuildModuleArtifacts, { import { clearExtractedSections, clearExtractedSectionsForPaths, + recordExtractedSuppressionsForPaths, } from './section-cleanup' import TreeSitterEngine from './tree-sitter-engine' @@ -95,13 +97,127 @@ export default async function extractSections( phase: 'database', status: 'progress', }) + if (options.mode === 'changed') { + const preparedFiles: PreparedFile[] = [] + await recordExtractedSuppressionsForPaths(affectedPaths) + if (files.length > 0) { + progress?.({ + message: `Extracting ${files.length} files`, + phase: 'sections', + status: 'start', + total: files.length, + }) + await yieldToProgressRenderer() + } + + for (const [fileIndex, file] of files.entries()) { + const preparedFile = await prepareFileSections({ + context, + engine, + file, + }) + preparedFiles.push(preparedFile) + + if ( + preparedFile.parserEngine === 'tree_sitter' && + preparedFile.parserStatus === 'ok' + ) { + parserUsedFiles += 1 + } else if ( + preparedFile.sections.some(section => section.kind === 'code') + ) { + parserFallbackFiles += 1 + } + + if (preparedFile.truncated) { + filesTruncatedBySectionLimit += 1 + } + + sectionCount += preparedFile.sections.length + progress?.({ + current: fileIndex + 1, + message: `Extracted ${file.path} (${preparedFile.sections.length} sections)`, + path: file.path, + phase: 'sections', + sectionCount: sectionCount, + status: 'progress', + total: files.length, + }) + await yieldToProgressRenderer() + } + + if (files.length > 0) { + progress?.({ + current: files.length, + message: `Extracted ${sectionCount} sections`, + parserCount: loadedParserCount, + phase: 'sections', + sectionCount: sectionCount, + status: 'done', + total: files.length, + }) + } + progress?.({ + message: 'Rebuilding module artifacts and retrieval index', + phase: 'modules', + status: 'start', + }) + + await withTransaction(async () => { + await clearExtractedSectionsForPaths(affectedPaths) + const rootNode = await upsertNode({ + name: 'Project Files', + summary: 'Files extracted from the current project.', + }) + rootNodeId = rootNode.id + + for (const preparedFile of preparedFiles) { + if (preparedFile.sections.length > 0) { + await persistPreparedFileSections({ + extractedAt, + preparedFile, + rootNodeId, + }) + } + } + + const currentModuleKeys = + await queryModuleArtifactKeysForPaths(changedPaths) + await rebuildChangedModuleArtifacts(extractedAt, { + canUpdateEmbeddings: options.canUpdateEmbeddings ?? false, + metadata: options.metadata, + metadataChanged: hasMetadataRelevantPath(affectedPaths), + moduleKeys: dedupeModuleKeys([ + ...previousModuleKeys, + ...currentModuleKeys, + ]), + }) + + await appendMemoryEvent({ + actor: 'cli', + eventType: 'project_extracted', + id: `event_${contentHash(`${context.projectRoot}:${extractedAt}`).slice(0, 32)}`, + subjectType: 'project', + summary: `Extracted ${files.length} files into ${sectionCount} sections.`, + }) + }) + progress?.({ + message: 'Module artifacts and retrieval index ready', + phase: 'modules', + status: 'done', + }) + + return { + filesTruncatedBySectionLimit, + loadedParserCount, + parserFallbackFiles, + parserUsedFiles, + sectionCount: sectionCount, + } + } + await withTransaction(async () => { - if (options.mode === 'changed') { - await clearExtractedSectionsForPaths([ - ...files.map(file => file.path), - ...(options.deletedPaths ?? []), - ]) - } else if (options.mode !== 'resume') { + if (options.mode !== 'resume') { await clearExtractedSections() } const rootNode = await upsertNode({ @@ -180,25 +296,9 @@ export default async function extractSections( phase: 'modules', status: 'start', }) - const currentModuleKeys = - options.mode === 'changed' - ? await queryModuleArtifactKeysForPaths(changedPaths) - : [] await withTransaction(async () => { - if (options.mode === 'changed') { - await rebuildChangedModuleArtifacts(extractedAt, { - canUpdateEmbeddings: options.canUpdateEmbeddings ?? false, - metadata: options.metadata, - metadataChanged: hasMetadataRelevantPath(affectedPaths), - moduleKeys: dedupeModuleKeys([ - ...previousModuleKeys, - ...currentModuleKeys, - ]), - }) - } else { - await rebuildModuleArtifacts(extractedAt, options.metadata) - await reindexRetrievalDocumentFts() - } + await rebuildModuleArtifacts(extractedAt, options.metadata) + await reindexRetrievalDocumentFts() await appendMemoryEvent({ actor: 'cli', diff --git a/src/modules/extraction/engine/manifest.ts b/src/modules/extraction/engine/manifest.ts index cd43e62..54b1935 100644 --- a/src/modules/extraction/engine/manifest.ts +++ b/src/modules/extraction/engine/manifest.ts @@ -1,4 +1,5 @@ -import { readFile, writeFile } from 'node:fs/promises' +import { randomUUID } from 'node:crypto' +import { readFile, rename, rm, writeFile } from 'node:fs/promises' import { join } from 'node:path' import { pathExists } from '@/modules/project/context' import type { ExtractionMode, LegacyExtractionMode } from '@/types/extraction' @@ -65,10 +66,19 @@ export async function writeExtractionManifest( memoryDir: string, manifest: ExtractionManifest, ): Promise { - await writeFile( - extractionManifestPath(memoryDir), - `${JSON.stringify(manifest, null, 2)}\n`, + const targetPath = extractionManifestPath(memoryDir) + const tempPath = join( + memoryDir, + `.extraction-manifest.${process.pid}.${Date.now()}.${randomUUID()}.tmp`, ) + + try { + await writeFile(tempPath, `${JSON.stringify(manifest, null, 2)}\n`) + await rename(tempPath, targetPath) + } catch (error) { + await rm(tempPath, { force: true }).catch(() => undefined) + throw error + } } export async function getExtractionFreshness( diff --git a/src/modules/extraction/engine/section-cleanup.ts b/src/modules/extraction/engine/section-cleanup.ts index 36ee134..aef3c70 100644 --- a/src/modules/extraction/engine/section-cleanup.ts +++ b/src/modules/extraction/engine/section-cleanup.ts @@ -135,6 +135,17 @@ where target_type = 'section' ) } +export async function recordExtractedSuppressionsForPaths( + paths: string[], +): Promise { + const uniquePaths = [...new Set(paths)].filter(Boolean) + if (uniquePaths.length === 0) { + return + } + + await recordExtractedSuppressions(uniquePaths) +} + async function recordExtractedSuppressions(paths?: string[]): Promise { const db = await getDb() await db.run(sql` diff --git a/src/modules/memory/runtime.ts b/src/modules/memory/runtime.ts index ef232b4..3ebfb37 100644 --- a/src/modules/memory/runtime.ts +++ b/src/modules/memory/runtime.ts @@ -1,7 +1,9 @@ +import { join } from 'node:path' import type { SaveOptions } from '@/database/services/save-memory' import { readExtractionManifest } from '@/modules/extraction/engine/manifest' import { extractProject } from '@/modules/extraction/extract-project' import { loadProjectContext } from '@/modules/project/context' +import { acquireFileLock } from '@/support/file-lock' import type { EmbeddingProviderContract } from '@/types/embedding-provider' import type { LoadedProjectContext } from '@/types/project' @@ -20,11 +22,27 @@ export async function updateChangedProjectMemorySilently( return undefined } - const result = await extractProject(context, 'changed', { - embeddingProvider, + const lock = await acquireFileLock({ + lockDir: join( + context.memoryDir, + 'locks', + 'changed-project-memory.lock', + ), + operationName: 'changed_project_memory', }) - return { - deletedFilePaths: result.deletedFilePaths, - updatedFilePaths: result.updatedFilePaths, + if (!lock.acquired) { + return undefined + } + + try { + const result = await extractProject(context, 'changed', { + embeddingProvider, + }) + return { + deletedFilePaths: result.deletedFilePaths, + updatedFilePaths: result.updatedFilePaths, + } + } finally { + await lock.release() } } diff --git a/src/support/file-lock.ts b/src/support/file-lock.ts new file mode 100644 index 0000000..2e9be2a --- /dev/null +++ b/src/support/file-lock.ts @@ -0,0 +1,175 @@ +import { randomUUID } from 'node:crypto' +import { mkdir, readFile, rm, stat, utimes, writeFile } from 'node:fs/promises' +import { dirname, join } from 'node:path' + +type FileLockInput = { + lockDir: string + operationName: string + pollMs?: number + staleMs?: number + timeoutMs?: number +} + +type FileLockResult = + | { + acquired: false + } + | { + acquired: true + release: () => Promise + } + +const DEFAULT_POLL_MS = 100 +const DEFAULT_STALE_MS = 10 * 60 * 1000 +const DEFAULT_TIMEOUT_MS = 2500 +const OWNER_FILE = 'owner.json' + +export async function acquireFileLock( + input: FileLockInput, +): Promise { + const pollMs = input.pollMs ?? DEFAULT_POLL_MS + const staleMs = input.staleMs ?? DEFAULT_STALE_MS + const timeoutMs = input.timeoutMs ?? DEFAULT_TIMEOUT_MS + const heartbeatMs = Math.max(1000, Math.floor(staleMs / 3)) + const deadline = Date.now() + timeoutMs + + await mkdir(dirname(input.lockDir), { recursive: true }) + + while (true) { + try { + await mkdir(input.lockDir) + const ownerId = randomUUID() + await writeLockOwner(input, ownerId) + const stopHeartbeat = startHeartbeat({ + heartbeatMs, + input, + ownerId, + }) + return { + acquired: true, + release: async () => { + stopHeartbeat() + if (await lockIsOwnedBy(input.lockDir, ownerId)) { + await rm(input.lockDir, { + force: true, + recursive: true, + }) + } + }, + } + } catch (error) { + if (!isFileExistsError(error)) { + throw error + } + } + + await removeStaleLock(input.lockDir, staleMs) + + if (Date.now() >= deadline) { + return { acquired: false } + } + + await sleep(Math.min(pollMs, Math.max(0, deadline - Date.now()))) + } +} + +function startHeartbeat(input: { + heartbeatMs: number + input: FileLockInput + ownerId: string +}): () => void { + const interval = setInterval(() => { + void refreshLock(input.input, input.ownerId) + }, input.heartbeatMs) + interval.unref?.() + + return () => { + clearInterval(interval) + } +} + +async function refreshLock( + input: FileLockInput, + ownerId: string, +): Promise { + if (!(await lockIsOwnedBy(input.lockDir, ownerId))) { + return + } + + const now = new Date() + try { + await writeLockOwner(input, ownerId) + await utimes(input.lockDir, now, now) + } catch { + // Stale-lock cleanup remains the recovery path if refresh fails. + } +} + +async function writeLockOwner( + input: FileLockInput, + ownerId: string, +): Promise { + try { + await writeFile( + join(input.lockDir, OWNER_FILE), + `${JSON.stringify( + { + createdAt: new Date().toISOString(), + id: ownerId, + operation: input.operationName, + pid: process.pid, + }, + null, + 2, + )}\n`, + ) + } catch { + // The lock directory itself is the synchronization primitive. + } +} + +async function lockIsOwnedBy( + lockDir: string, + ownerId: string, +): Promise { + try { + const raw = await readFile(join(lockDir, OWNER_FILE), 'utf8') + const parsed = JSON.parse(raw) as unknown + return ( + typeof parsed === 'object' && + parsed !== null && + 'id' in parsed && + parsed.id === ownerId + ) + } catch { + return true + } +} + +async function removeStaleLock( + lockDir: string, + staleMs: number, +): Promise { + try { + const metadata = await stat(lockDir) + if (Date.now() - metadata.mtimeMs < staleMs) { + return + } + await rm(lockDir, { force: true, recursive: true }) + } catch { + // Another process may have released or replaced the lock already. + } +} + +function isFileExistsError(error: unknown): boolean { + return ( + typeof error === 'object' && + error !== null && + 'code' in error && + error.code === 'EEXIST' + ) +} + +async function sleep(milliseconds: number): Promise { + await new Promise(resolve => setTimeout(resolve, milliseconds)) +} diff --git a/tests/features/memory/runtime.test.ts b/tests/features/memory/runtime.test.ts index e83651c..14b1664 100644 --- a/tests/features/memory/runtime.test.ts +++ b/tests/features/memory/runtime.test.ts @@ -1,5 +1,5 @@ import { afterEach, describe, expect, it } from 'bun:test' -import { mkdtemp, writeFile } from 'node:fs/promises' +import { mkdtemp, utimes, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import { readExtractionManifest } from '@/modules/extraction/engine/manifest' @@ -105,4 +105,87 @@ describe('memory/runtime', () => { mode: 'changed', }) }) + + it('skips changed project memory when another process holds the lock', async () => { + const projectRoot = await mkdtemp(join(tmpdir(), 'konteks-runtime-')) + tempDirs.push(projectRoot) + await mkdir(join(projectRoot, 'src')) + await writeFile( + join(projectRoot, 'package.json'), + '{"name":"fixture"}\n', + ) + await writeFile( + join(projectRoot, 'src', 'index.txt'), + 'export const first = true\n', + ) + const context = await withWorkingDirectory(projectRoot, () => + loadMcpProjectContext(), + ) + await withWorkingDirectory(projectRoot, () => + extractProject(context, 'full'), + ) + await writeFile( + join(projectRoot, 'src', 'later.txt'), + 'export const later = true\n', + ) + await mkdir( + join(context.memoryDir, 'locks', 'changed-project-memory.lock'), + ) + + await expect( + withWorkingDirectory(projectRoot, () => + updateChangedProjectMemorySilently(context), + ), + ).resolves.toBeUndefined() + await expect( + readExtractionManifest(context.memoryDir), + ).resolves.toMatchObject({ + mode: 'full', + }) + }) + + it('removes stale changed project memory locks and updates memory', async () => { + const projectRoot = await mkdtemp(join(tmpdir(), 'konteks-runtime-')) + tempDirs.push(projectRoot) + await mkdir(join(projectRoot, 'src')) + await writeFile( + join(projectRoot, 'package.json'), + '{"name":"fixture"}\n', + ) + await writeFile( + join(projectRoot, 'src', 'index.txt'), + 'export const first = true\n', + ) + const context = await withWorkingDirectory(projectRoot, () => + loadMcpProjectContext(), + ) + await withWorkingDirectory(projectRoot, () => + extractProject(context, 'full'), + ) + await writeFile( + join(projectRoot, 'src', 'later.txt'), + 'export const later = true\n', + ) + const lockDir = join( + context.memoryDir, + 'locks', + 'changed-project-memory.lock', + ) + await mkdir(lockDir) + const staleTime = new Date(Date.now() - 11 * 60 * 1000) + await utimes(lockDir, staleTime, staleTime) + + await expect( + withWorkingDirectory(projectRoot, () => + updateChangedProjectMemorySilently(context), + ), + ).resolves.toMatchObject({ + updatedFilePaths: ['src/later.txt'], + }) + await expect( + readExtractionManifest(context.memoryDir), + ).resolves.toMatchObject({ + mode: 'changed', + }) + }) }) diff --git a/tests/features/providers/extraction/engine/extract-project.test.ts b/tests/features/providers/extraction/engine/extract-project.test.ts index 06e8f0c..3006b53 100644 --- a/tests/features/providers/extraction/engine/extract-project.test.ts +++ b/tests/features/providers/extraction/engine/extract-project.test.ts @@ -1,10 +1,15 @@ import { afterEach, describe, expect, it } from 'bun:test' -import { mkdtemp, readFile, unlink, writeFile } from 'node:fs/promises' +import { mkdtemp, readdir, readFile, unlink, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import { eq, sql } from 'drizzle-orm' import getDb from '@/database/actions/_db' -import { targetEmbeddings, vectorIndexEntries } from '@/database/schema' +import { + sections, + targetEmbeddings, + vectorIndexEntries, +} from '@/database/schema' +import forgetMemory from '@/database/services/forget-memory' import { saveKonteksDiary, saveKonteksMemories, @@ -115,6 +120,12 @@ describe('extractProject', () => { 'src/index.txt', ]) expect(manifest?.summaryHash).toHaveLength(64) + const memoryEntries = await readdir(context.memoryDir) + expect( + memoryEntries.some(entry => + /^\.extraction-manifest\./u.test(entry), + ), + ).toBe(false) }) it('reports fresh status after extraction and stale after a file change', async () => { @@ -382,6 +393,112 @@ describe('extractProject', () => { expect(paths).toContain('src/new.txt') }) + it('rolls back changed extraction database writes when persistence fails', async () => { + const projectRoot = await makeTempProject() + const context = await withProjectRoot(projectRoot, () => + loadProjectContext(), + ) + await extractTestProject(context, 'full') + const db = await getDb() + const beforeSections = await db + .select({ + contentInline: sections.contentInline, + path: sections.path, + }) + .from(sections) + + await writeFile( + join(projectRoot, 'src', 'index.txt'), + 'Changed content that should roll back.\n', + ) + await writeFile( + join(projectRoot, 'src', 'new.txt'), + 'New content that should roll back.\n', + ) + await db.run(sql` +create temp trigger fail_changed_project_event +before insert on memory_events +when new.event_type = 'project_extracted' +begin + select raise(abort, 'forced changed extraction rollback'); +end; +`) + + await expect(extractTestProject(context, 'changed')).rejects.toThrow( + 'memory_events', + ) + + const afterSections = await db + .select({ + contentInline: sections.contentInline, + path: sections.path, + }) + .from(sections) + const manifest = await readExtractionManifest(context.memoryDir) + + expect(afterSections).toEqual(beforeSections) + expect( + afterSections.some(section => section.path === 'src/new.txt'), + ).toBe(false) + expect( + afterSections.some(section => + section.contentInline?.includes('Changed content'), + ), + ).toBe(false) + expect(manifest?.mode).toBe('full') + }) + + it('does not reintroduce suppressed sections during changed extraction', async () => { + const projectRoot = await makeTempProject() + const context = await withProjectRoot(projectRoot, () => + loadProjectContext(), + ) + await writeFile( + join(projectRoot, 'README.md'), + '# Suppressed\nStable suppressed text.\n\n# Changed\nOld visible text.\n', + ) + await extractTestProject(context, 'full') + const db = await getDb() + const readmeSections = await db + .select({ + contentHash: sections.contentHash, + contentInline: sections.contentInline, + id: sections.id, + path: sections.path, + }) + .from(sections) + .where(eq(sections.path, 'README.md')) + const section = readmeSections.find(row => + row.contentInline?.includes('Stable suppressed text.'), + ) + + expect(section).toBeDefined() + await forgetMemory({ + id: section?.id, + mode: 'invalidate', + reason: 'suppress extracted section', + }) + await writeFile( + join(projectRoot, 'README.md'), + '# Suppressed\nStable suppressed text.\n\n# Changed\nNew visible text.\n', + ) + + await extractTestProject(context, 'changed') + + const rows = await db + .select({ + contentHash: sections.contentHash, + id: sections.id, + path: sections.path, + }) + .from(sections) + .where(eq(sections.path, 'README.md')) + + expect(rows.some(row => row.contentHash === section?.contentHash)).toBe( + false, + ) + }) + it('changed mode is a no-op when extraction is already current', async () => { const projectRoot = await makeTempProject() const context = await withProjectRoot(projectRoot, () => diff --git a/tests/features/support/file-lock.test.ts b/tests/features/support/file-lock.test.ts new file mode 100644 index 0000000..77d6e08 --- /dev/null +++ b/tests/features/support/file-lock.test.ts @@ -0,0 +1,76 @@ +import { afterEach, describe, expect, it } from 'bun:test' +import { mkdir, mkdtemp, utimes } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { acquireFileLock } from '@/support/file-lock' +import { rm } from '@/support/file-manager' + +const tempDirs: string[] = [] + +afterEach(async () => { + await Promise.all(tempDirs.splice(0).map(path => rm(path))) +}) + +describe('file lock', () => { + it('does not remove an active long-held lock as stale', async () => { + const root = await mkdtemp(join(tmpdir(), 'konteks-lock-')) + tempDirs.push(root) + const lockDir = join(root, 'active.lock') + const first = await acquireFileLock({ + lockDir, + operationName: 'test_active_lock', + staleMs: 1500, + timeoutMs: 50, + }) + + expect(first.acquired).toBe(true) + await sleep(1700) + + const second = await acquireFileLock({ + lockDir, + operationName: 'test_contender', + staleMs: 1500, + timeoutMs: 50, + }) + expect(second.acquired).toBe(false) + + if (first.acquired) { + await first.release() + } + const third = await acquireFileLock({ + lockDir, + operationName: 'test_after_release', + staleMs: 1500, + timeoutMs: 50, + }) + expect(third.acquired).toBe(true) + if (third.acquired) { + await third.release() + } + }) + + it('removes abandoned stale locks', async () => { + const root = await mkdtemp(join(tmpdir(), 'konteks-lock-')) + tempDirs.push(root) + const lockDir = join(root, 'stale.lock') + await mkdir(lockDir) + const staleTime = new Date(Date.now() - 2000) + await utimes(lockDir, staleTime, staleTime) + + const lock = await acquireFileLock({ + lockDir, + operationName: 'test_stale_lock', + staleMs: 100, + timeoutMs: 500, + }) + + expect(lock.acquired).toBe(true) + if (lock.acquired) { + await lock.release() + } + }) +}) + +async function sleep(milliseconds: number): Promise { + await new Promise(resolve => setTimeout(resolve, milliseconds)) +}