Removed direct write

This commit is contained in:
abhigyantrumio 2025-09-22 15:05:28 +05:30
parent 05de6f38a9
commit 6f88b4ef27
7 changed files with 8 additions and 1016 deletions

12
.env
View file

@ -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
# VITE_KUZU_ENABLED=false

View file

@ -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');
});
});

View file

@ -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 = <K extends keyof FeatureFlags>(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');

View file

@ -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<void>;
addRelationshipAsync(relationship: GraphRelationship): Promise<void>;
// Batch operations
addNodesBatch(nodes: GraphNode[]): Promise<void>;
addRelationshipsBatch(relationships: GraphRelationship[]): Promise<void>;
// Transaction support
beginTransaction(): Promise<void>;
commitTransaction(): Promise<void>;
rollbackTransaction(): Promise<void>;
// Utility methods
flushPendingOperations(): Promise<void>;
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
};

View file

@ -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<Promise<void>>();
private writeQueue: Array<() => Promise<void>> = [];
private isProcessingQueue = false;
// Transaction state
private transactionActive = false;
private transactionOperations: Array<() => Promise<void>> = [];
constructor(kuzuGraph?: KuzuKnowledgeGraph, options: Partial<DirectWriteOptions> = {}) {
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<void> {
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<void> {
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<void> {
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<void> {
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<void> {
while (this.activeWrites.size >= this.options.maxConcurrentWrites) {
await Promise.race(this.activeWrites);
}
}
/**
* Process write queue with concurrency control
*/
private async processWriteQueue(): Promise<void> {
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<void> {
this.transactionActive = true;
this.transactionOperations = [];
}
async commitTransaction(): Promise<void> {
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<void> {
this.transactionOperations = [];
this.transactionActive = false;
}
/**
* Batch operations for bulk inserts
*/
async addNodesBatch(nodes: GraphNode[]): Promise<void> {
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<void> {
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<void> {
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<void> {
return new Promise(resolve => setTimeout(resolve, ms));
}
}

View file

@ -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;
}

View file

@ -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<KnowledgeGraph> {
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);