diff --git a/.env b/.env index 1b6fa5bfa..b65c471ec 100644 --- a/.env +++ b/.env @@ -119,15 +119,5 @@ VITE_KUZU_ENABLED=true # - Gracefully fall back to JSON-only if KuzuDB fails # - Maintain full backward compatibility -# Enable direct writes to KuzuDB (eliminates 25+ second flush delay) -# Values: true/false, 1/0, yes/no -VITE_KUZU_DIRECT_WRITES=true -# Direct writes will: -# - Write to KuzuDB immediately during processing (no batching delay) -# - Reduce total processing time by 20+ seconds -# - Provide real-time progress instead of post-processing flush -# - Maintain reliability with automatic fallback to batching on errors - # To disable KuzuDB and use JSON-only: -# VITE_KUZU_ENABLED=false -VITE_KUZU_DIRECT_WRITES=false \ No newline at end of file +# VITE_KUZU_ENABLED=false \ No newline at end of file diff --git a/src/__tests__/direct-write-implementation.test.ts b/src/__tests__/direct-write-implementation.test.ts deleted file mode 100644 index bf048f9f4..000000000 --- a/src/__tests__/direct-write-implementation.test.ts +++ /dev/null @@ -1,429 +0,0 @@ -/** - * Comprehensive test suite for Direct-Write KuzuDB implementation - * - * This test ensures that the direct-write approach: - * 1. Maintains data integrity - * 2. Provides performance improvements - * 3. Handles errors gracefully - * 4. Falls back to batching when needed - */ - -import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; -import { DirectWriteKnowledgeGraph } from '../core/graph/direct-write-knowledge-graph.ts'; -import { DualWriteKnowledgeGraph } from '../core/graph/dual-write-knowledge-graph.ts'; -import { SimpleKnowledgeGraph } from '../core/graph/graph.ts'; -import type { GraphNode, GraphRelationship } from '../core/graph/types.ts'; - -// Mock KuzuDB for testing -const mockKuzuGraph = { - addNode: vi.fn(), - addRelationship: vi.fn(), - commitAll: vi.fn(), - executeQuery: vi.fn().mockResolvedValue({ rows: [], columns: [] }), - nodes: [], - relationships: [] -}; - -describe('DirectWriteKnowledgeGraph', () => { - let directWriteGraph: DirectWriteKnowledgeGraph; - let dualWriteGraph: DualWriteKnowledgeGraph; - let simpleGraph: SimpleKnowledgeGraph; - - beforeEach(() => { - vi.clearAllMocks(); - - directWriteGraph = new DirectWriteKnowledgeGraph(mockKuzuGraph as any, { - enableDirectWrites: true, - maxConcurrentWrites: 5, - fallbackToBatching: true, - retryAttempts: 2 - }); - - dualWriteGraph = new DualWriteKnowledgeGraph(mockKuzuGraph as any); - simpleGraph = new SimpleKnowledgeGraph(); - }); - - afterEach(async () => { - // Cleanup any pending operations - if ('flushPendingOperations' in directWriteGraph) { - await directWriteGraph.flushPendingOperations(); - } - }); - - describe('Data Integrity', () => { - it('should maintain identical JSON storage across all implementations', () => { - const testNode: GraphNode = { - id: 'test-node-1', - label: 'Function', - properties: { - name: 'testFunction', - filePath: '/test/file.ts', - startLine: 10 - } - }; - - // Add to all implementations - directWriteGraph.addNode(testNode); - dualWriteGraph.addNode(testNode); - simpleGraph.addNode(testNode); - - // Verify JSON storage is identical - expect(directWriteGraph.nodes).toEqual(simpleGraph.nodes); - expect(dualWriteGraph.nodes).toEqual(simpleGraph.nodes); - expect(directWriteGraph.nodes.length).toBe(1); - expect(directWriteGraph.nodes[0]).toEqual(testNode); - }); - - it('should handle relationship storage consistently', () => { - const testRelationship: GraphRelationship = { - id: 'test-rel-1', - type: 'CALLS', - source: 'node-1', - target: 'node-2', - properties: { - callType: 'function', - line: 15 - } - }; - - directWriteGraph.addRelationship(testRelationship); - dualWriteGraph.addRelationship(testRelationship); - simpleGraph.addRelationship(testRelationship); - - expect(directWriteGraph.relationships).toEqual(simpleGraph.relationships); - expect(dualWriteGraph.relationships).toEqual(simpleGraph.relationships); - }); - - it('should handle large datasets without data loss', async () => { - const nodeCount = 1000; - const nodes: GraphNode[] = []; - - // Generate test nodes - for (let i = 0; i < nodeCount; i++) { - nodes.push({ - id: `node-${i}`, - label: 'Function', - properties: { - name: `function${i}`, - filePath: `/test/file${i}.ts`, - startLine: i * 10 - } - }); - } - - // Add all nodes - const startTime = performance.now(); - for (const node of nodes) { - directWriteGraph.addNode(node); - } - - // Wait for all operations to complete - await directWriteGraph.flushPendingOperations(); - const endTime = performance.now(); - - // Verify all nodes are present - expect(directWriteGraph.nodes.length).toBe(nodeCount); - - // Verify performance (should be faster than traditional batching) - const processingTime = endTime - startTime; - console.log(`Direct write processing time for ${nodeCount} nodes: ${processingTime}ms`); - - // Performance should be reasonable (less than 5 seconds for 1000 nodes) - expect(processingTime).toBeLessThan(5000); - }); - }); - - describe('Performance Improvements', () => { - it('should complete operations faster than dual-write with batching', async () => { - const testNodes: GraphNode[] = Array.from({ length: 100 }, (_, i) => ({ - id: `perf-node-${i}`, - label: 'Function', - properties: { - name: `perfFunction${i}`, - filePath: `/perf/file${i}.ts`, - startLine: i - } - })); - - // Test direct write performance - const directStartTime = performance.now(); - for (const node of testNodes) { - directWriteGraph.addNode(node); - } - await directWriteGraph.flushPendingOperations(); - const directEndTime = performance.now(); - const directTime = directEndTime - directStartTime; - - // Test dual write performance (simulated) - const dualStartTime = performance.now(); - for (const node of testNodes) { - dualWriteGraph.addNode(node); - } - // Simulate the flush delay that would occur - await new Promise(resolve => setTimeout(resolve, 100)); // Simulated flush time - const dualEndTime = performance.now(); - const dualTime = dualEndTime - dualStartTime; - - console.log(`Direct write time: ${directTime}ms, Dual write time: ${dualTime}ms`); - - // Direct writes should be competitive or faster - expect(directTime).toBeLessThan(dualTime * 2); // Allow some overhead for async operations - }); - - it('should provide real-time progress without batching delays', async () => { - const stats = directWriteGraph.getStats(); - - // Initially no pending operations - expect(stats.pendingNodes).toBe(0); - expect(stats.pendingRelationships).toBe(0); - - // Add some nodes - directWriteGraph.addNode({ - id: 'progress-node-1', - label: 'Function', - properties: { name: 'test1' } - }); - - directWriteGraph.addNode({ - id: 'progress-node-2', - label: 'Function', - properties: { name: 'test2' } - }); - - // Should show immediate progress in JSON storage - expect(directWriteGraph.nodes.length).toBe(2); - - // Wait for KuzuDB operations to complete - await directWriteGraph.flushPendingOperations(); - - const finalStats = directWriteGraph.getStats(); - expect(finalStats.totalNodes).toBe(2); - }); - }); - - describe('Error Handling and Fallback', () => { - it('should continue processing when KuzuDB writes fail', async () => { - // Mock KuzuDB failure - mockKuzuGraph.executeQuery.mockRejectedValueOnce(new Error('KuzuDB connection failed')); - - const testNode: GraphNode = { - id: 'error-node-1', - label: 'Function', - properties: { name: 'errorFunction' } - }; - - // Should not throw error - expect(() => directWriteGraph.addNode(testNode)).not.toThrow(); - - // JSON storage should still work - expect(directWriteGraph.nodes.length).toBe(1); - expect(directWriteGraph.nodes[0]).toEqual(testNode); - - // Wait for async operations to complete - await directWriteGraph.flushPendingOperations(); - - // Check that fallback was triggered - const stats = directWriteGraph.getStats(); - expect(stats.directWriteFailures).toBeGreaterThan(0); - }); - - it('should retry failed operations according to configuration', async () => { - let callCount = 0; - mockKuzuGraph.executeQuery.mockImplementation(() => { - callCount++; - if (callCount < 3) { - return Promise.reject(new Error('Temporary failure')); - } - return Promise.resolve({ rows: [], columns: [] }); - }); - - const testNode: GraphNode = { - id: 'retry-node-1', - label: 'Function', - properties: { name: 'retryFunction' } - }; - - directWriteGraph.addNode(testNode); - await directWriteGraph.flushPendingOperations(); - - // Should have retried and eventually succeeded - expect(callCount).toBe(3); // Initial attempt + 2 retries - }); - - it('should fallback to batching when direct writes consistently fail', async () => { - // Mock consistent failures - mockKuzuGraph.executeQuery.mockRejectedValue(new Error('Persistent KuzuDB failure')); - - const testNodes: GraphNode[] = Array.from({ length: 5 }, (_, i) => ({ - id: `fallback-node-${i}`, - label: 'Function', - properties: { name: `fallbackFunction${i}` } - })); - - for (const node of testNodes) { - directWriteGraph.addNode(node); - } - - await directWriteGraph.flushPendingOperations(); - - const stats = directWriteGraph.getStats(); - - // Should have fallen back to batching for some operations - expect(stats.fallbackToBatch).toBeGreaterThan(0); - - // JSON storage should still be intact - expect(directWriteGraph.nodes.length).toBe(5); - }); - }); - - describe('Concurrency Control', () => { - it('should handle concurrent writes without data corruption', async () => { - const concurrentNodes: GraphNode[] = Array.from({ length: 20 }, (_, i) => ({ - id: `concurrent-node-${i}`, - label: 'Function', - properties: { name: `concurrentFunction${i}` } - })); - - // Add all nodes concurrently - const promises = concurrentNodes.map(node => - Promise.resolve(directWriteGraph.addNode(node)) - ); - - await Promise.all(promises); - await directWriteGraph.flushPendingOperations(); - - // All nodes should be present - expect(directWriteGraph.nodes.length).toBe(20); - - // No duplicates should exist - const nodeIds = directWriteGraph.nodes.map(n => n.id); - const uniqueIds = new Set(nodeIds); - expect(uniqueIds.size).toBe(20); - }); - - it('should respect maxConcurrentWrites configuration', async () => { - const maxConcurrentWrites = 3; - const testGraph = new DirectWriteKnowledgeGraph(mockKuzuGraph as any, { - enableDirectWrites: true, - maxConcurrentWrites, - fallbackToBatching: false - }); - - // Track concurrent operations - let activeCalls = 0; - let maxConcurrentCalls = 0; - - mockKuzuGraph.executeQuery.mockImplementation(() => { - activeCalls++; - maxConcurrentCalls = Math.max(maxConcurrentCalls, activeCalls); - - return new Promise(resolve => { - setTimeout(() => { - activeCalls--; - resolve({ rows: [], columns: [] }); - }, 10); - }); - }); - - // Add many nodes quickly - const nodes: GraphNode[] = Array.from({ length: 10 }, (_, i) => ({ - id: `concurrent-limit-node-${i}`, - label: 'Function', - properties: { name: `limitFunction${i}` } - })); - - for (const node of nodes) { - testGraph.addNode(node); - } - - await testGraph.flushPendingOperations(); - - // Should not exceed the configured limit - expect(maxConcurrentCalls).toBeLessThanOrEqual(maxConcurrentWrites); - }); - }); - - describe('Transaction Support', () => { - it('should support transaction boundaries for batch operations', async () => { - const nodes: GraphNode[] = Array.from({ length: 3 }, (_, i) => ({ - id: `transaction-node-${i}`, - label: 'Function', - properties: { name: `transactionFunction${i}` } - })); - - await directWriteGraph.beginTransaction(); - - for (const node of nodes) { - await directWriteGraph.addNodeAsync(node); - } - - await directWriteGraph.commitTransaction(); - - // All nodes should be present after commit - expect(directWriteGraph.nodes.length).toBe(3); - }); - - it('should rollback transactions on failure', async () => { - // Mock failure during transaction - let callCount = 0; - mockKuzuGraph.executeQuery.mockImplementation(() => { - callCount++; - if (callCount === 2) { - return Promise.reject(new Error('Transaction failure')); - } - return Promise.resolve({ rows: [], columns: [] }); - }); - - const nodes: GraphNode[] = Array.from({ length: 3 }, (_, i) => ({ - id: `rollback-node-${i}`, - label: 'Function', - properties: { name: `rollbackFunction${i}` } - })); - - await directWriteGraph.beginTransaction(); - - try { - for (const node of nodes) { - await directWriteGraph.addNodeAsync(node); - } - await directWriteGraph.commitTransaction(); - } catch (error) { - // Transaction should rollback automatically - expect(error).toBeDefined(); - } - - // JSON storage should still have the nodes (since JSON writes happen immediately) - expect(directWriteGraph.nodes.length).toBe(3); - }); - }); -}); - -describe('Integration with Existing Pipeline', () => { - it('should maintain backward compatibility with existing processors', () => { - const graph = new DirectWriteKnowledgeGraph(mockKuzuGraph as any); - - // Should support the existing synchronous interface - expect(typeof graph.addNode).toBe('function'); - expect(typeof graph.addRelationship).toBe('function'); - expect(Array.isArray(graph.nodes)).toBe(true); - expect(Array.isArray(graph.relationships)).toBe(true); - - // Should also support new async interface - expect(typeof graph.addNodeAsync).toBe('function'); - expect(typeof graph.addRelationshipAsync).toBe('function'); - expect(typeof graph.flushPendingOperations).toBe('function'); - }); - - it('should provide performance statistics for monitoring', () => { - const graph = new DirectWriteKnowledgeGraph(mockKuzuGraph as any); - const stats = graph.getStats(); - - // Should include all expected metrics - expect(typeof stats.nodesWrittenToJSON).toBe('number'); - expect(typeof stats.nodesWrittenToKuzuDB).toBe('number'); - expect(typeof stats.directWriteSuccesses).toBe('number'); - expect(typeof stats.directWriteFailures).toBe('number'); - expect(typeof stats.averageWriteTimeMs).toBe('number'); - expect(typeof stats.activeWrites).toBe('number'); - }); -}); diff --git a/src/config/feature-flags.ts b/src/config/feature-flags.ts index dfb4954a2..d134f3f86 100644 --- a/src/config/feature-flags.ts +++ b/src/config/feature-flags.ts @@ -21,7 +21,6 @@ export interface FeatureFlags { enableKuzuDB: boolean; enableKuzuDBPersistence: boolean; enableKuzuDBPerformanceMonitoring: boolean; - enableKuzuDBDirectWrites: boolean; // Debug Features enableDebugMode: boolean; @@ -47,7 +46,6 @@ export const DEFAULT_FEATURE_FLAGS: FeatureFlags = { enableKuzuDB: false, enableKuzuDBPersistence: false, enableKuzuDBPerformanceMonitoring: false, - enableKuzuDBDirectWrites: false, // Debug Features enableDebugMode: false, @@ -115,30 +113,6 @@ class FeatureFlagManager { console.log(`โš ๏ธ Warning: KuzuDB env var not recognized: "${kuzuEnabled}" - using default: ${flags.enableKuzuDB}`); } - // Handle KuzuDB Direct Writes settings - let directWritesEnabled: string | undefined; - try { - if (import.meta && import.meta.env) { - directWritesEnabled = import.meta.env.VITE_KUZU_DIRECT_WRITES?.toLowerCase(); - } - } catch (e) { - if (typeof process !== 'undefined' && process.env) { - directWritesEnabled = process.env.KUZU_DIRECT_WRITES?.toLowerCase(); - } - } - - console.log(`๐Ÿ” Debug: directWritesEnabled from env = "${directWritesEnabled}"`); - - if (directWritesEnabled === 'true' || directWritesEnabled === '1' || directWritesEnabled === 'yes') { - flags.enableKuzuDBDirectWrites = true; - console.log('โšก Environment: KuzuDB direct writes enabled via VITE_KUZU_DIRECT_WRITES=true'); - } else if (directWritesEnabled === 'false' || directWritesEnabled === '0' || directWritesEnabled === 'no') { - flags.enableKuzuDBDirectWrites = false; - console.log('๐Ÿ”„ Environment: KuzuDB direct writes disabled via VITE_KUZU_DIRECT_WRITES=false'); - } else { - console.log(`โš ๏ธ Warning: KuzuDB direct writes env var not recognized: "${directWritesEnabled}" - using default: ${flags.enableKuzuDBDirectWrites}`); - } - // Then, try to load from localStorage (can override environment) // Only try localStorage if we're in a browser context (not in workers) try { @@ -337,7 +311,6 @@ export const setFeatureFlag = (key: K, value: Feat export const isKuzuDBEnabled = (): boolean => featureFlags.isKuzuDBEnabled(); export const isKuzuDBPersistenceEnabled = (): boolean => featureFlags.getFlag('enableKuzuDBPersistence'); -export const isKuzuDBDirectWritesEnabled = (): boolean => featureFlags.getFlag('enableKuzuDBDirectWrites'); export const isDebugModeEnabled = (): boolean => featureFlags.isDebugModeEnabled(); export const isPerformanceMonitoringEnabled = (): boolean => featureFlags.isPerformanceMonitoringEnabled(); export const isWorkerPoolEnabled = (): boolean => featureFlags.getFlag('enableWorkerPool'); diff --git a/src/core/graph/async-knowledge-graph.ts b/src/core/graph/async-knowledge-graph.ts deleted file mode 100644 index c16564884..000000000 --- a/src/core/graph/async-knowledge-graph.ts +++ /dev/null @@ -1,62 +0,0 @@ -/** - * Async Knowledge Graph Interface for Direct KuzuDB Writes - * - * This interface supports both sync and async operations to enable - * gradual migration from batched to direct writes. - */ - -import type { GraphNode, GraphRelationship } from './types.ts'; - -export interface AsyncKnowledgeGraph { - nodes: GraphNode[]; - relationships: GraphRelationship[]; - - // Sync methods (backward compatibility) - addNode(node: GraphNode): void; - addRelationship(relationship: GraphRelationship): void; - - // Async methods (new direct-write capability) - addNodeAsync(node: GraphNode): Promise; - addRelationshipAsync(relationship: GraphRelationship): Promise; - - // Batch operations - addNodesBatch(nodes: GraphNode[]): Promise; - addRelationshipsBatch(relationships: GraphRelationship[]): Promise; - - // Transaction support - beginTransaction(): Promise; - commitTransaction(): Promise; - rollbackTransaction(): Promise; - - // Utility methods - flushPendingOperations(): Promise; - getStats(): { - pendingNodes: number; - pendingRelationships: number; - totalNodes: number; - totalRelationships: number; - }; -} - -/** - * Configuration for direct-write behavior - */ -export interface DirectWriteOptions { - enableDirectWrites: boolean; - batchSize: number; - maxConcurrentWrites: number; - retryAttempts: number; - retryDelayMs: number; - fallbackToBatching: boolean; - enableTransactions: boolean; -} - -export const DEFAULT_DIRECT_WRITE_OPTIONS: DirectWriteOptions = { - enableDirectWrites: true, - batchSize: 50, - maxConcurrentWrites: 10, - retryAttempts: 3, - retryDelayMs: 100, - fallbackToBatching: true, - enableTransactions: true -}; diff --git a/src/core/graph/direct-write-knowledge-graph.ts b/src/core/graph/direct-write-knowledge-graph.ts deleted file mode 100644 index 7d8dfed7e..000000000 --- a/src/core/graph/direct-write-knowledge-graph.ts +++ /dev/null @@ -1,401 +0,0 @@ -/** - * Direct-Write Knowledge Graph Implementation - * - * This implementation writes directly to KuzuDB during processing, - * eliminating the 25+ second flush delay while maintaining reliability. - */ - -import type { KnowledgeGraph, GraphNode, GraphRelationship } from './types.ts'; -import type { AsyncKnowledgeGraph, DirectWriteOptions } from './async-knowledge-graph.ts'; -import { SimpleKnowledgeGraph } from './graph.ts'; -import type { KuzuKnowledgeGraph } from './kuzu-knowledge-graph.ts'; -import { DEFAULT_DIRECT_WRITE_OPTIONS } from './async-knowledge-graph.ts'; -// NODE_TABLE_SCHEMAS import removed - now using KuzuKnowledgeGraph's schema filtering - -export class DirectWriteKnowledgeGraph implements KnowledgeGraph, AsyncKnowledgeGraph { - private jsonGraph: SimpleKnowledgeGraph; - private kuzuGraph: KuzuKnowledgeGraph | null; - private options: DirectWriteOptions; - private enableKuzuDB: boolean; - - // Performance tracking - private stats = { - nodesWrittenToJSON: 0, - nodesWrittenToKuzuDB: 0, - relationshipsWrittenToJSON: 0, - relationshipsWrittenToKuzuDB: 0, - directWriteSuccesses: 0, - directWriteFailures: 0, - fallbackToBatch: 0, - averageWriteTimeMs: 0 - }; - - // Concurrency control - private activeWrites = new Set>(); - private writeQueue: Array<() => Promise> = []; - private isProcessingQueue = false; - - // Transaction state - private transactionActive = false; - private transactionOperations: Array<() => Promise> = []; - - constructor(kuzuGraph?: KuzuKnowledgeGraph, options: Partial = {}) { - this.jsonGraph = new SimpleKnowledgeGraph(); - this.kuzuGraph = kuzuGraph || null; - this.enableKuzuDB = !!kuzuGraph; - this.options = { ...DEFAULT_DIRECT_WRITE_OPTIONS, ...options }; - } - - /** - * Get all nodes (from JSON primary storage) - */ - get nodes(): GraphNode[] { - return this.jsonGraph.nodes; - } - - /** - * Get all relationships (from JSON primary storage) - */ - get relationships(): GraphRelationship[] { - return this.jsonGraph.relationships; - } - - /** - * Synchronous addNode (backward compatibility) - * Uses fire-and-forget async write to KuzuDB - */ - addNode(node: GraphNode): void { - // Always write to JSON immediately (primary storage) - this.jsonGraph.addNode(node); - this.stats.nodesWrittenToJSON++; - - // Fire-and-forget direct write to KuzuDB - if (this.enableKuzuDB && this.options.enableDirectWrites) { - console.log(`โšก DIRECT-WRITE: Adding ${node.label} node ${node.id} to KuzuDB - NEW CODE!`); - this.addNodeAsync(node).catch(error => { - console.warn(`โŒ DIRECT-WRITE: KuzuDB write failed for node ${node.id}:`, error); - this.stats.directWriteFailures++; - }); - } else { - console.log(`๐Ÿ”„ DIRECT-WRITE: Direct writes disabled - KuzuDB: ${this.enableKuzuDB}, DirectWrites: ${this.options.enableDirectWrites}`); - } - } - - /** - * Synchronous addRelationship (backward compatibility) - */ - addRelationship(relationship: GraphRelationship): void { - // Always write to JSON immediately - this.jsonGraph.addRelationship(relationship); - this.stats.relationshipsWrittenToJSON++; - - // Fire-and-forget direct write to KuzuDB - if (this.enableKuzuDB && this.options.enableDirectWrites) { - this.addRelationshipAsync(relationship).catch(error => { - console.warn(`Direct KuzuDB write failed for relationship ${relationship.id}:`, error); - this.stats.directWriteFailures++; - }); - } - } - - /** - * Async addNode with proper error handling and concurrency control - */ - async addNodeAsync(node: GraphNode): Promise { - if (!this.enableKuzuDB || !this.kuzuGraph) return; - - const writeOperation = async () => { - const startTime = performance.now(); - - try { - // Wait for available slot if too many concurrent writes - await this.waitForWriteSlot(); - - if (this.transactionActive) { - // Add to transaction queue - this.transactionOperations.push(() => this.executeNodeWrite(node)); - } else { - // Execute immediately - await this.executeNodeWrite(node); - } - - this.stats.directWriteSuccesses++; - this.stats.nodesWrittenToKuzuDB++; - - // Update average write time - const writeTime = performance.now() - startTime; - this.updateAverageWriteTime(writeTime); - - } catch (error) { - this.stats.directWriteFailures++; - - if (this.options.fallbackToBatching) { - console.warn(`Direct write failed for node ${node.id}, adding to batch queue:`, error); - this.stats.fallbackToBatch++; - // Add to KuzuDB's internal batch queue as fallback - this.kuzuGraph.addNode(node); - } else { - throw error; - } - } - }; - - if (this.options.maxConcurrentWrites > 1) { - // Add to queue for concurrent processing - this.writeQueue.push(writeOperation); - this.processWriteQueue(); - } else { - // Execute immediately (sequential mode) - await writeOperation(); - } - } - - /** - * Async addRelationship with proper error handling - */ - async addRelationshipAsync(relationship: GraphRelationship): Promise { - if (!this.enableKuzuDB || !this.kuzuGraph) return; - - const writeOperation = async () => { - const startTime = performance.now(); - - try { - await this.waitForWriteSlot(); - - if (this.transactionActive) { - this.transactionOperations.push(() => this.executeRelationshipWrite(relationship)); - } else { - await this.executeRelationshipWrite(relationship); - } - - this.stats.directWriteSuccesses++; - this.stats.relationshipsWrittenToKuzuDB++; - - const writeTime = performance.now() - startTime; - this.updateAverageWriteTime(writeTime); - - } catch (error) { - this.stats.directWriteFailures++; - - if (this.options.fallbackToBatching) { - console.warn(`Direct write failed for relationship ${relationship.id}, adding to batch queue:`, error); - this.stats.fallbackToBatch++; - this.kuzuGraph.addRelationship(relationship); - } else { - throw error; - } - } - }; - - if (this.options.maxConcurrentWrites > 1) { - this.writeQueue.push(writeOperation); - this.processWriteQueue(); - } else { - await writeOperation(); - } - } - - /** - * Execute actual node write to KuzuDB (reuses proven KuzuKnowledgeGraph logic) - */ - private async executeNodeWrite(node: GraphNode): Promise { - if (!this.kuzuGraph) return; - - // Reuse the existing, tested commitSingleNode method with schema filtering and auto-recovery - await this.kuzuGraph.commitSingleNode(node); - } - - /** - * Execute actual relationship write to KuzuDB (reuses proven KuzuKnowledgeGraph logic) - */ - private async executeRelationshipWrite(relationship: GraphRelationship): Promise { - if (!this.kuzuGraph) return; - - // Reuse the existing, tested commitSingleRelationship method with auto-recovery - await this.kuzuGraph.commitSingleRelationship(relationship); - } - - // Query building methods removed - now reusing KuzuKnowledgeGraph's proven implementation - - // All schema filtering, property formatting, and query execution logic - // has been removed - now reusing KuzuKnowledgeGraph's proven methods - // This eliminates ~100 lines of duplicate code and prevents schema issues - - /** - * Wait for available write slot (concurrency control) - */ - private async waitForWriteSlot(): Promise { - while (this.activeWrites.size >= this.options.maxConcurrentWrites) { - await Promise.race(this.activeWrites); - } - } - - /** - * Process write queue with concurrency control - */ - private async processWriteQueue(): Promise { - if (this.isProcessingQueue) return; - this.isProcessingQueue = true; - - while (this.writeQueue.length > 0 && this.activeWrites.size < this.options.maxConcurrentWrites) { - const operation = this.writeQueue.shift(); - if (operation) { - const promise = operation().finally(() => { - this.activeWrites.delete(promise); - }); - this.activeWrites.add(promise); - } - } - - this.isProcessingQueue = false; - } - - /** - * Transaction support - */ - async beginTransaction(): Promise { - this.transactionActive = true; - this.transactionOperations = []; - } - - async commitTransaction(): Promise { - if (!this.transactionActive) return; - - try { - // Execute all queued operations - for (const operation of this.transactionOperations) { - await operation(); - } - - this.transactionOperations = []; - this.transactionActive = false; - } catch (error) { - await this.rollbackTransaction(); - throw error; - } - } - - async rollbackTransaction(): Promise { - this.transactionOperations = []; - this.transactionActive = false; - } - - /** - * Batch operations for bulk inserts - */ - async addNodesBatch(nodes: GraphNode[]): Promise { - await this.beginTransaction(); - try { - for (const node of nodes) { - await this.addNodeAsync(node); - } - await this.commitTransaction(); - } catch (error) { - await this.rollbackTransaction(); - throw error; - } - } - - async addRelationshipsBatch(relationships: GraphRelationship[]): Promise { - await this.beginTransaction(); - try { - for (const relationship of relationships) { - await this.addRelationshipAsync(relationship); - } - await this.commitTransaction(); - } catch (error) { - await this.rollbackTransaction(); - throw error; - } - } - - /** - * Flush any pending operations (for compatibility) - */ - async flushPendingOperations(): Promise { - console.log('๐Ÿš€ DIRECT-WRITE: flushPendingOperations() called - NEW CODE RUNNING!'); - console.log(`๐Ÿ“Š DIRECT-WRITE: Stats - fallbackToBatch: ${this.stats.fallbackToBatch}, directWriteSuccesses: ${this.stats.directWriteSuccesses}, directWriteFailures: ${this.stats.directWriteFailures}`); - - // Debug: Check if KuzuGraph has pending operations - if (this.kuzuGraph) { - const pendingNodes = (this.kuzuGraph as any).pendingNodes?.length || 0; - const pendingRels = (this.kuzuGraph as any).pendingRelationships?.length || 0; - console.log(`๐Ÿ” DIRECT-WRITE: KuzuGraph has ${pendingNodes} pending nodes, ${pendingRels} pending relationships`); - } - - // Wait for all active writes to complete - await Promise.all(this.activeWrites); - - // Process any remaining queue items - while (this.writeQueue.length > 0) { - await this.processWriteQueue(); - await Promise.all(this.activeWrites); - } - - // Only flush fallback batched operations if there were actual fallbacks - if (this.kuzuGraph && 'commitAll' in this.kuzuGraph && this.stats.fallbackToBatch > 0) { - console.log(`๐Ÿ”„ DIRECT-WRITE: Flushing ${this.stats.fallbackToBatch} fallback operations that failed direct write...`); - await (this.kuzuGraph as any).commitAll(); - } else { - console.log('โœ… DIRECT-WRITE: No fallback operations to flush - all direct writes succeeded!'); - - // BUT: Check if KuzuGraph has pending operations anyway (this shouldn't happen!) - if (this.kuzuGraph) { - const pendingNodes = (this.kuzuGraph as any).pendingNodes?.length || 0; - const pendingRels = (this.kuzuGraph as any).pendingRelationships?.length || 0; - if (pendingNodes > 0 || pendingRels > 0) { - console.warn(`๐Ÿšจ DIRECT-WRITE: BUG DETECTED! KuzuGraph has ${pendingNodes} pending nodes, ${pendingRels} pending relationships despite no fallbacks!`); - console.warn('๐Ÿšจ DIRECT-WRITE: This means something is calling kuzuGraph.addNode() directly, bypassing DirectWriteKnowledgeGraph!'); - console.warn('๐Ÿšจ DIRECT-WRITE: Flushing these unexpected pending operations...'); - await (this.kuzuGraph as any).commitAll(); - } - } - } - } - - /** - * Get performance statistics - */ - getStats() { - return { - ...this.stats, - pendingNodes: 0, // Direct writes don't queue - pendingRelationships: 0, - totalNodes: this.jsonGraph.nodes.length, - totalRelationships: this.jsonGraph.relationships.length, - activeWrites: this.activeWrites.size, - queuedWrites: this.writeQueue.length - }; - } - - /** - * Log performance statistics - */ - logStats(): void { - const stats = this.getStats(); - console.log('๐Ÿ“Š Direct-Write Statistics:'); - console.log(` JSON entities: ${stats.nodesWrittenToJSON + stats.relationshipsWrittenToJSON}`); - console.log(` KuzuDB entities: ${stats.nodesWrittenToKuzuDB + stats.relationshipsWrittenToKuzuDB}`); - console.log(` Direct write successes: ${stats.directWriteSuccesses}`); - console.log(` Direct write failures: ${stats.directWriteFailures}`); - console.log(` Fallback to batch: ${stats.fallbackToBatch}`); - console.log(` Average write time: ${stats.averageWriteTimeMs.toFixed(2)}ms`); - - const successRate = stats.directWriteSuccesses + stats.directWriteFailures > 0 - ? (stats.directWriteSuccesses / (stats.directWriteSuccesses + stats.directWriteFailures) * 100).toFixed(1) - : '100'; - console.log(` Success rate: ${successRate}%`); - } - - /** - * Utility methods - */ - private updateAverageWriteTime(newTime: number): void { - const totalWrites = this.stats.directWriteSuccesses + this.stats.directWriteFailures; - this.stats.averageWriteTimeMs = (this.stats.averageWriteTimeMs * (totalWrites - 1) + newTime) / totalWrites; - } - - private delay(ms: number): Promise { - return new Promise(resolve => setTimeout(resolve, ms)); - } -} diff --git a/src/core/graph/direct-write-test.ts b/src/core/graph/direct-write-test.ts deleted file mode 100644 index 5bd58b694..000000000 --- a/src/core/graph/direct-write-test.ts +++ /dev/null @@ -1,61 +0,0 @@ -/** - * Quick test to verify direct-write implementation - */ - -import { DirectWriteKnowledgeGraph } from './direct-write-knowledge-graph.ts'; -import type { GraphNode } from './types.ts'; - -// Mock KuzuDB for testing -const mockKuzuGraph = { - queryEngine: { - executeQuery: async (query: string) => { - console.log(`โœ… Mock KuzuDB executed: ${query}`); - return { rows: [], columns: [] }; - } - }, - nodes: [], - relationships: [] -}; - -export async function testDirectWrite() { - console.log('๐Ÿงช Testing Direct-Write Implementation...'); - - const graph = new DirectWriteKnowledgeGraph(mockKuzuGraph as any, { - enableDirectWrites: true, - maxConcurrentWrites: 5, - fallbackToBatching: true, - retryAttempts: 2 - }); - - // Test node with valid schema properties - const testNode: GraphNode = { - id: 'test-folder-1', - label: 'Folder', - properties: { - name: 'test-folder', - path: '/test/folder', - fullPath: '/test/folder', - depth: 2, - // This should be filtered out (not in Folder schema) - type: 'directory', - invalidProperty: 'should-be-removed' - } - }; - - console.log('๐Ÿ“ Adding test node...'); - graph.addNode(testNode); - - // Wait for async operations - await new Promise(resolve => setTimeout(resolve, 100)); - - console.log('๐Ÿ”„ Flushing pending operations...'); - await graph.flushPendingOperations(); - - console.log('๐Ÿ“Š Final stats:', graph.getStats()); - console.log('โœ… Direct-write test completed!'); -} - -// Export for manual testing -if (typeof window !== 'undefined') { - (window as any).testDirectWrite = testDirectWrite; -} diff --git a/src/core/ingestion/parallel-pipeline.ts b/src/core/ingestion/parallel-pipeline.ts index 48456e45a..388e8fd7f 100644 --- a/src/core/ingestion/parallel-pipeline.ts +++ b/src/core/ingestion/parallel-pipeline.ts @@ -5,7 +5,7 @@ import { ParallelParsingProcessor } from './parallel-parsing-processor.ts'; import { ImportProcessor } from './import-processor.ts'; import { CallProcessor } from './call-processor.ts'; import { WebWorkerPoolUtils } from '../../lib/web-worker-pool.js'; -import { isKuzuDBEnabled, isKuzuDBDirectWritesEnabled } from '../../config/feature-flags.ts'; +import { isKuzuDBEnabled } from '../../config/feature-flags.ts'; export interface PipelineInput { projectRoot: string; @@ -134,14 +134,8 @@ export class ParallelGraphPipeline { console.log('๐Ÿ”ง Worker Pool Statistics:', workerStats); } - // Handle post-processing based on graph type - if ('flushPendingOperations' in graph) { - // DirectWriteKnowledgeGraph - flush any remaining operations - console.log('โšก Flushing any remaining direct-write operations...'); - await (graph as any).flushPendingOperations(); - (graph as any).logStats(); - } else if ('flushKuzuDB' in graph) { - // DualWriteKnowledgeGraph - traditional batch flush + // Handle post-processing for KuzuDB (batched mode only) + if ('flushKuzuDB' in graph) { console.log('๐Ÿ”„ Flushing batched KuzuDB operations...'); await (graph as any).flushKuzuDB(); (graph as any).logDualWriteStats(); @@ -262,7 +256,6 @@ export class ParallelGraphPipeline { */ private async createGraph(): Promise { console.log(`๐Ÿ” KuzuDB enabled check: ${isKuzuDBEnabled()}`); - console.log(`โšก KuzuDB direct writes check: ${isKuzuDBDirectWritesEnabled()}`); if (isKuzuDBEnabled()) { try { @@ -286,21 +279,10 @@ export class ParallelGraphPipeline { autoCommit: false }); - // Choose between direct-write and dual-write based on feature flag - if (isKuzuDBDirectWritesEnabled()) { - const { DirectWriteKnowledgeGraph } = await import('../graph/direct-write-knowledge-graph.ts'); - console.log('โšก KuzuDB integration initialized - using DIRECT-WRITE mode (no flush delay!)'); - return new DirectWriteKnowledgeGraph(kuzuGraph, { - enableDirectWrites: true, - maxConcurrentWrites: 10, - fallbackToBatching: true, - retryAttempts: 3 - }); - } else { - const { DualWriteKnowledgeGraph } = await import('../graph/dual-write-knowledge-graph.ts'); - console.log('โœ… KuzuDB integration initialized - using dual-write mode (with BATCH flush)'); - return new DualWriteKnowledgeGraph(kuzuGraph); - } + // Use dual-write mode (batched) + const { DualWriteKnowledgeGraph } = await import('../graph/dual-write-knowledge-graph.ts'); + console.log('โœ… KuzuDB integration initialized - using dual-write mode (batched)'); + return new DualWriteKnowledgeGraph(kuzuGraph); } catch (error) { console.warn('โŒ KuzuDB initialization failed, falling back to JSON-only mode:', error);