2025-08-26 12:32:21 -07:00
|
|
|
/**
|
|
|
|
|
* Base Storage Adapter
|
|
|
|
|
* Provides common functionality for all storage adapters, including statistics tracking
|
|
|
|
|
*/
|
|
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
import {
|
|
|
|
|
StatisticsData,
|
|
|
|
|
StorageAdapter,
|
|
|
|
|
HNSWNoun,
|
|
|
|
|
HNSWVerb,
|
|
|
|
|
GraphVerb,
|
|
|
|
|
HNSWNounWithMetadata,
|
|
|
|
|
HNSWVerbWithMetadata,
|
|
|
|
|
NounMetadata,
|
|
|
|
|
VerbMetadata
|
|
|
|
|
} from '../../coreTypes.js'
|
2025-08-26 12:32:21 -07:00
|
|
|
import { extractFieldNamesFromJson, mapToStandardField } from '../../utils/fieldNameTracking.js'
|
2025-09-22 15:45:35 -07:00
|
|
|
import { getGlobalMutex, cleanupMutexes } from '../../utils/mutex.js'
|
2025-08-26 12:32:21 -07:00
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Base class for storage adapters that implements statistics tracking
|
|
|
|
|
*/
|
|
|
|
|
export abstract class BaseStorageAdapter implements StorageAdapter {
|
|
|
|
|
// Abstract methods that must be implemented by subclasses
|
|
|
|
|
abstract init(): Promise<void>
|
|
|
|
|
|
fix: metadata explosion bug - 69K files reduced to ~1K
Critical fix for metadata indexing that was creating 60+ chunk files per entity.
Root cause: Vector embeddings (384-dimensional arrays) were being indexed in
metadata, causing each dimension to create a separate chunk file with numeric
field names ("0", "1", "2", etc.).
Changes:
- Modified extractIndexableFields() to exclude vector/embedding fields
- Added NEVER_INDEX set: ['vector', 'embedding', 'embeddings', 'connections']
- Added safety check to skip arrays > 10 elements
- Preserves small array indexing (tags, categories, roles)
Impact:
- Reduces metadata files from 69,429 → ~1,200 (58x reduction)
- Fixes server initialization hangs
- Fixes metadata batch loading stalling at batch 23
- Fixes VFS getDescendants() hanging with large datasets
- Fixes Graph View UI not loading
Test Results:
- 7/7 integration tests passing
- Verified: 6 chunk files for 10 entities (was 7,210 before fix)
- 611/622 unit tests passing
Files Modified:
- src/utils/metadataIndex.ts - Core fix
- src/coreTypes.ts - HNSWVerb type enforcement with VerbType enum
- src/storage/adapters/* - Include core relational fields in HNSWVerb
- src/storage/adapters/baseStorageAdapter.ts - Type enforcement (HNSWNoun, GraphVerb)
- tests/integration/metadata-vector-exclusion.test.ts - Comprehensive test coverage
🤖 Generated with [Claude Code](https://claude.com/claude-code)
Co-Authored-By: Claude <noreply@anthropic.com>
2025-10-16 16:10:31 -07:00
|
|
|
abstract saveNoun(noun: HNSWNoun): Promise<void>
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract saveNounMetadata(id: string, metadata: NounMetadata): Promise<void>
|
|
|
|
|
abstract deleteNounMetadata(id: string): Promise<void>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract getNoun(id: string): Promise<HNSWNounWithMetadata | null>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract getNounsByNounType(nounType: string): Promise<HNSWNounWithMetadata[]>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
|
|
|
|
abstract deleteNoun(id: string): Promise<void>
|
|
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract saveVerb(verb: HNSWVerb): Promise<void>
|
|
|
|
|
abstract saveVerbMetadata(id: string, metadata: VerbMetadata): Promise<void>
|
|
|
|
|
abstract deleteVerbMetadata(id: string): Promise<void>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract getVerb(id: string): Promise<HNSWVerbWithMetadata | null>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract getVerbsBySource(sourceId: string): Promise<HNSWVerbWithMetadata[]>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract getVerbsByTarget(targetId: string): Promise<HNSWVerbWithMetadata[]>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
2025-10-17 12:29:27 -07:00
|
|
|
abstract getVerbsByType(type: string): Promise<HNSWVerbWithMetadata[]>
|
2025-08-26 12:32:21 -07:00
|
|
|
|
|
|
|
|
abstract deleteVerb(id: string): Promise<void>
|
|
|
|
|
|
|
|
|
|
abstract saveMetadata(id: string, metadata: any): Promise<void>
|
|
|
|
|
|
|
|
|
|
abstract getMetadata(id: string): Promise<any | null>
|
|
|
|
|
|
2025-10-06 15:43:45 -07:00
|
|
|
abstract getNounMetadata(id: string): Promise<any | null>
|
|
|
|
|
|
2025-08-26 12:32:21 -07:00
|
|
|
abstract getVerbMetadata(id: string): Promise<any | null>
|
|
|
|
|
|
2025-10-10 11:15:17 -07:00
|
|
|
// HNSW Index Persistence (v3.35.0+)
|
|
|
|
|
// These methods enable HNSW index rebuilding after container restarts
|
|
|
|
|
|
|
|
|
|
abstract getNounVector(id: string): Promise<number[] | null>
|
|
|
|
|
|
|
|
|
|
abstract saveHNSWData(nounId: string, hnswData: {
|
|
|
|
|
level: number
|
|
|
|
|
connections: Record<string, string[]>
|
|
|
|
|
}): Promise<void>
|
|
|
|
|
|
|
|
|
|
abstract getHNSWData(nounId: string): Promise<{
|
|
|
|
|
level: number
|
|
|
|
|
connections: Record<string, string[]>
|
|
|
|
|
} | null>
|
|
|
|
|
|
|
|
|
|
abstract saveHNSWSystem(systemData: {
|
|
|
|
|
entryPointId: string | null
|
|
|
|
|
maxLevel: number
|
|
|
|
|
}): Promise<void>
|
|
|
|
|
|
|
|
|
|
abstract getHNSWSystem(): Promise<{
|
|
|
|
|
entryPointId: string | null
|
|
|
|
|
maxLevel: number
|
|
|
|
|
} | null>
|
|
|
|
|
|
2025-08-26 12:32:21 -07:00
|
|
|
abstract clear(): Promise<void>
|
|
|
|
|
|
|
|
|
|
abstract getStorageStatus(): Promise<{
|
|
|
|
|
type: string
|
|
|
|
|
used: number
|
|
|
|
|
quota: number | null
|
|
|
|
|
details?: Record<string, any>
|
|
|
|
|
}>
|
|
|
|
|
|
|
|
|
|
// NOTE: getAllNouns and getAllVerbs have been removed to prevent expensive full scans.
|
|
|
|
|
// Use getNouns() and getVerbs() with pagination instead.
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get nouns with pagination and filtering
|
|
|
|
|
* @param options Pagination and filtering options
|
|
|
|
|
* @returns Promise that resolves to a paginated result of nouns
|
|
|
|
|
*/
|
|
|
|
|
abstract getNouns(options?: {
|
|
|
|
|
pagination?: {
|
|
|
|
|
offset?: number
|
|
|
|
|
limit?: number
|
|
|
|
|
cursor?: string
|
|
|
|
|
}
|
|
|
|
|
filter?: {
|
|
|
|
|
nounType?: string | string[]
|
|
|
|
|
service?: string | string[]
|
|
|
|
|
metadata?: Record<string, any>
|
|
|
|
|
}
|
|
|
|
|
}): Promise<{
|
2025-10-17 12:29:27 -07:00
|
|
|
items: HNSWNounWithMetadata[]
|
2025-08-26 12:32:21 -07:00
|
|
|
totalCount?: number
|
|
|
|
|
hasMore: boolean
|
|
|
|
|
nextCursor?: string
|
|
|
|
|
}>
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get verbs with pagination and filtering
|
|
|
|
|
* @param options Pagination and filtering options
|
|
|
|
|
* @returns Promise that resolves to a paginated result of verbs
|
|
|
|
|
*/
|
|
|
|
|
abstract getVerbs(options?: {
|
|
|
|
|
pagination?: {
|
|
|
|
|
offset?: number
|
|
|
|
|
limit?: number
|
|
|
|
|
cursor?: string
|
|
|
|
|
}
|
|
|
|
|
filter?: {
|
|
|
|
|
verbType?: string | string[]
|
|
|
|
|
sourceId?: string | string[]
|
|
|
|
|
targetId?: string | string[]
|
|
|
|
|
service?: string | string[]
|
|
|
|
|
metadata?: Record<string, any>
|
|
|
|
|
}
|
|
|
|
|
}): Promise<{
|
2025-10-17 12:29:27 -07:00
|
|
|
items: HNSWVerbWithMetadata[]
|
2025-08-26 12:32:21 -07:00
|
|
|
totalCount?: number
|
|
|
|
|
hasMore: boolean
|
|
|
|
|
nextCursor?: string
|
|
|
|
|
}>
|
|
|
|
|
|
2025-09-02 14:55:15 -07:00
|
|
|
/**
|
|
|
|
|
* Get nouns with pagination (internal implementation)
|
|
|
|
|
* This method should be implemented by storage adapters to support efficient pagination
|
|
|
|
|
* @param options Pagination options
|
|
|
|
|
* @returns Promise that resolves to a paginated result of nouns
|
|
|
|
|
*/
|
|
|
|
|
getNounsWithPagination?(options: {
|
|
|
|
|
limit?: number
|
|
|
|
|
cursor?: string
|
|
|
|
|
filter?: {
|
|
|
|
|
nounType?: string | string[]
|
|
|
|
|
service?: string | string[]
|
|
|
|
|
metadata?: Record<string, any>
|
|
|
|
|
}
|
|
|
|
|
}): Promise<{
|
2025-10-17 12:29:27 -07:00
|
|
|
items: HNSWNounWithMetadata[]
|
2025-09-02 14:55:15 -07:00
|
|
|
totalCount?: number
|
|
|
|
|
hasMore: boolean
|
|
|
|
|
nextCursor?: string
|
|
|
|
|
}>
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get verbs with pagination (internal implementation)
|
|
|
|
|
* This method should be implemented by storage adapters to support efficient pagination
|
|
|
|
|
* @param options Pagination options
|
|
|
|
|
* @returns Promise that resolves to a paginated result of verbs
|
|
|
|
|
*/
|
|
|
|
|
getVerbsWithPagination?(options: {
|
|
|
|
|
limit?: number
|
|
|
|
|
cursor?: string
|
|
|
|
|
filter?: {
|
|
|
|
|
verbType?: string | string[]
|
|
|
|
|
sourceId?: string | string[]
|
|
|
|
|
targetId?: string | string[]
|
|
|
|
|
service?: string | string[]
|
|
|
|
|
metadata?: Record<string, any>
|
|
|
|
|
}
|
|
|
|
|
}): Promise<{
|
2025-10-17 12:29:27 -07:00
|
|
|
items: HNSWVerbWithMetadata[]
|
2025-09-02 14:55:15 -07:00
|
|
|
totalCount?: number
|
|
|
|
|
hasMore: boolean
|
|
|
|
|
nextCursor?: string
|
|
|
|
|
}>
|
|
|
|
|
|
2025-08-26 12:32:21 -07:00
|
|
|
// Statistics cache
|
|
|
|
|
protected statisticsCache: StatisticsData | null = null
|
|
|
|
|
|
|
|
|
|
// Batch update timer ID
|
|
|
|
|
protected statisticsBatchUpdateTimerId: NodeJS.Timeout | null = null
|
|
|
|
|
|
|
|
|
|
// Flag to indicate if statistics have been modified since last save
|
|
|
|
|
protected statisticsModified = false
|
|
|
|
|
|
|
|
|
|
// Time of last statistics flush to storage
|
|
|
|
|
protected lastStatisticsFlushTime = 0
|
|
|
|
|
|
|
|
|
|
// Minimum time between statistics flushes (5 seconds)
|
|
|
|
|
protected readonly MIN_FLUSH_INTERVAL_MS = 5000
|
|
|
|
|
|
|
|
|
|
// Maximum time to wait before flushing statistics (30 seconds)
|
|
|
|
|
protected readonly MAX_FLUSH_DELAY_MS = 30000
|
|
|
|
|
|
|
|
|
|
// Throttling tracking properties
|
|
|
|
|
protected throttlingDetected = false
|
|
|
|
|
protected throttlingBackoffMs = 1000 // Start with 1 second
|
|
|
|
|
protected maxBackoffMs = 30000 // Max 30 seconds
|
|
|
|
|
protected consecutiveThrottleEvents = 0
|
|
|
|
|
protected lastThrottleTime = 0
|
|
|
|
|
protected totalThrottleEvents = 0
|
|
|
|
|
protected throttleEventsByHour: number[] = new Array(24).fill(0)
|
|
|
|
|
protected throttleReasons: Record<string, number> = {}
|
|
|
|
|
protected lastThrottleHourIndex = -1
|
|
|
|
|
|
|
|
|
|
// Operation impact tracking
|
|
|
|
|
protected delayedOperations = 0
|
|
|
|
|
protected retriedOperations = 0
|
|
|
|
|
protected failedDueToThrottling = 0
|
|
|
|
|
protected totalDelayMs = 0
|
|
|
|
|
|
|
|
|
|
// Service-level throttling
|
|
|
|
|
protected serviceThrottling: Map<string, {
|
|
|
|
|
throttleCount: number
|
|
|
|
|
lastThrottle: number
|
|
|
|
|
status: 'normal' | 'throttled' | 'recovering'
|
|
|
|
|
}> = new Map()
|
|
|
|
|
|
|
|
|
|
// Statistics-specific methods that must be implemented by subclasses
|
|
|
|
|
protected abstract saveStatisticsData(
|
|
|
|
|
statistics: StatisticsData
|
|
|
|
|
): Promise<void>
|
|
|
|
|
|
|
|
|
|
protected abstract getStatisticsData(): Promise<StatisticsData | null>
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Save statistics data
|
|
|
|
|
* @param statistics The statistics data to save
|
|
|
|
|
*/
|
|
|
|
|
async saveStatistics(statistics: StatisticsData): Promise<void> {
|
|
|
|
|
// Update the cache with a deep copy to avoid reference issues
|
|
|
|
|
this.statisticsCache = {
|
|
|
|
|
nounCount: { ...statistics.nounCount },
|
|
|
|
|
verbCount: { ...statistics.verbCount },
|
|
|
|
|
metadataCount: { ...statistics.metadataCount },
|
|
|
|
|
hnswIndexSize: statistics.hnswIndexSize,
|
|
|
|
|
lastUpdated: statistics.lastUpdated,
|
|
|
|
|
// Include serviceActivity if present
|
|
|
|
|
...(statistics.serviceActivity && {
|
|
|
|
|
serviceActivity: Object.fromEntries(
|
|
|
|
|
Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}])
|
|
|
|
|
)
|
|
|
|
|
}),
|
|
|
|
|
// Include services if present
|
|
|
|
|
...(statistics.services && {
|
|
|
|
|
services: statistics.services.map(s => ({...s}))
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Schedule a batch update instead of saving immediately
|
|
|
|
|
this.scheduleBatchUpdate()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get statistics data
|
|
|
|
|
* @returns Promise that resolves to the statistics data
|
|
|
|
|
*/
|
|
|
|
|
async getStatistics(): Promise<StatisticsData | null> {
|
|
|
|
|
// If we have cached statistics, return a deep copy
|
|
|
|
|
if (this.statisticsCache) {
|
|
|
|
|
return {
|
|
|
|
|
nounCount: { ...this.statisticsCache.nounCount },
|
|
|
|
|
verbCount: { ...this.statisticsCache.verbCount },
|
|
|
|
|
metadataCount: { ...this.statisticsCache.metadataCount },
|
|
|
|
|
hnswIndexSize: this.statisticsCache.hnswIndexSize,
|
|
|
|
|
lastUpdated: this.statisticsCache.lastUpdated
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Otherwise, get from storage
|
|
|
|
|
const statistics = await this.getStatisticsData()
|
|
|
|
|
|
|
|
|
|
// If we found statistics, update the cache
|
|
|
|
|
if (statistics) {
|
|
|
|
|
// Update the cache with a deep copy
|
|
|
|
|
this.statisticsCache = {
|
|
|
|
|
nounCount: { ...statistics.nounCount },
|
|
|
|
|
verbCount: { ...statistics.verbCount },
|
|
|
|
|
metadataCount: { ...statistics.metadataCount },
|
|
|
|
|
hnswIndexSize: statistics.hnswIndexSize,
|
|
|
|
|
lastUpdated: statistics.lastUpdated
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return statistics
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Schedule a batch update of statistics
|
|
|
|
|
*/
|
|
|
|
|
protected scheduleBatchUpdate(): void {
|
|
|
|
|
// Mark statistics as modified
|
|
|
|
|
this.statisticsModified = true
|
|
|
|
|
|
|
|
|
|
// If a timer is already set, don't set another one
|
|
|
|
|
if (this.statisticsBatchUpdateTimerId !== null) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Calculate time since last flush
|
|
|
|
|
const now = Date.now()
|
|
|
|
|
const timeSinceLastFlush = now - this.lastStatisticsFlushTime
|
|
|
|
|
|
|
|
|
|
// If we've recently flushed, wait longer before the next flush
|
|
|
|
|
const delayMs =
|
|
|
|
|
timeSinceLastFlush < this.MIN_FLUSH_INTERVAL_MS
|
|
|
|
|
? this.MAX_FLUSH_DELAY_MS
|
|
|
|
|
: this.MIN_FLUSH_INTERVAL_MS
|
|
|
|
|
|
|
|
|
|
// Schedule the batch update
|
|
|
|
|
this.statisticsBatchUpdateTimerId = setTimeout(() => {
|
|
|
|
|
this.flushStatistics()
|
|
|
|
|
}, delayMs)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Flush statistics to storage
|
|
|
|
|
*/
|
|
|
|
|
protected async flushStatistics(): Promise<void> {
|
|
|
|
|
// Clear the timer
|
|
|
|
|
if (this.statisticsBatchUpdateTimerId !== null) {
|
|
|
|
|
clearTimeout(this.statisticsBatchUpdateTimerId)
|
|
|
|
|
this.statisticsBatchUpdateTimerId = null
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// If statistics haven't been modified, no need to flush
|
|
|
|
|
if (!this.statisticsModified || !this.statisticsCache) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
// Save the statistics to storage
|
|
|
|
|
await this.saveStatisticsData(this.statisticsCache)
|
|
|
|
|
|
|
|
|
|
// Update the last flush time
|
|
|
|
|
this.lastStatisticsFlushTime = Date.now()
|
|
|
|
|
// Reset the modified flag
|
|
|
|
|
this.statisticsModified = false
|
|
|
|
|
} catch (error) {
|
|
|
|
|
console.error('Failed to flush statistics data:', error)
|
|
|
|
|
// Mark as still modified so we'll try again later
|
|
|
|
|
this.statisticsModified = true
|
|
|
|
|
// Don't throw the error to avoid disrupting the application
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Increment a statistic counter
|
|
|
|
|
* @param type The type of statistic to increment ('noun', 'verb', 'metadata')
|
|
|
|
|
* @param service The service that inserted the data
|
|
|
|
|
* @param amount The amount to increment by (default: 1)
|
|
|
|
|
*/
|
|
|
|
|
async incrementStatistic(
|
|
|
|
|
type: 'noun' | 'verb' | 'metadata',
|
|
|
|
|
service: string,
|
|
|
|
|
amount: number = 1
|
|
|
|
|
): Promise<void> {
|
|
|
|
|
// Get current statistics from cache or storage
|
|
|
|
|
let statistics = this.statisticsCache
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = await this.getStatisticsData()
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = this.createDefaultStatistics()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Update the cache
|
|
|
|
|
this.statisticsCache = {
|
|
|
|
|
nounCount: { ...statistics.nounCount },
|
|
|
|
|
verbCount: { ...statistics.verbCount },
|
|
|
|
|
metadataCount: { ...statistics.metadataCount },
|
|
|
|
|
hnswIndexSize: statistics.hnswIndexSize,
|
|
|
|
|
lastUpdated: statistics.lastUpdated,
|
|
|
|
|
// Include serviceActivity if present
|
|
|
|
|
...(statistics.serviceActivity && {
|
|
|
|
|
serviceActivity: Object.fromEntries(
|
|
|
|
|
Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}])
|
|
|
|
|
)
|
|
|
|
|
}),
|
|
|
|
|
// Include services if present
|
|
|
|
|
...(statistics.services && {
|
|
|
|
|
services: statistics.services.map(s => ({...s}))
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Increment the appropriate counter
|
|
|
|
|
const counterMap = {
|
|
|
|
|
noun: this.statisticsCache!.nounCount,
|
|
|
|
|
verb: this.statisticsCache!.verbCount,
|
|
|
|
|
metadata: this.statisticsCache!.metadataCount
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const counter = counterMap[type]
|
|
|
|
|
counter[service] = (counter[service] || 0) + amount
|
|
|
|
|
|
|
|
|
|
// Track service activity
|
|
|
|
|
this.trackServiceActivity(service, 'add')
|
|
|
|
|
|
|
|
|
|
// Update timestamp
|
|
|
|
|
this.statisticsCache!.lastUpdated = new Date().toISOString()
|
|
|
|
|
|
|
|
|
|
// Schedule a batch update instead of saving immediately
|
|
|
|
|
this.scheduleBatchUpdate()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Track service activity (first/last activity, operation counts)
|
|
|
|
|
* @param service The service name
|
|
|
|
|
* @param operation The operation type
|
|
|
|
|
*/
|
|
|
|
|
protected trackServiceActivity(
|
|
|
|
|
service: string,
|
|
|
|
|
operation: 'add' | 'update' | 'delete'
|
|
|
|
|
): void {
|
|
|
|
|
if (!this.statisticsCache) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Initialize serviceActivity if it doesn't exist
|
|
|
|
|
if (!this.statisticsCache.serviceActivity) {
|
|
|
|
|
this.statisticsCache.serviceActivity = {}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const now = new Date().toISOString()
|
|
|
|
|
const activity = this.statisticsCache.serviceActivity[service]
|
|
|
|
|
|
|
|
|
|
if (!activity) {
|
|
|
|
|
// First activity for this service
|
|
|
|
|
this.statisticsCache.serviceActivity[service] = {
|
|
|
|
|
firstActivity: now,
|
|
|
|
|
lastActivity: now,
|
|
|
|
|
totalOperations: 1
|
|
|
|
|
}
|
|
|
|
|
} else {
|
|
|
|
|
// Update existing activity
|
|
|
|
|
activity.lastActivity = now
|
|
|
|
|
activity.totalOperations++
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Decrement a statistic counter
|
|
|
|
|
* @param type The type of statistic to decrement ('noun', 'verb', 'metadata')
|
|
|
|
|
* @param service The service that inserted the data
|
|
|
|
|
* @param amount The amount to decrement by (default: 1)
|
|
|
|
|
*/
|
|
|
|
|
async decrementStatistic(
|
|
|
|
|
type: 'noun' | 'verb' | 'metadata',
|
|
|
|
|
service: string,
|
|
|
|
|
amount: number = 1
|
|
|
|
|
): Promise<void> {
|
|
|
|
|
// Get current statistics from cache or storage
|
|
|
|
|
let statistics = this.statisticsCache
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = await this.getStatisticsData()
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = this.createDefaultStatistics()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Update the cache
|
|
|
|
|
this.statisticsCache = {
|
|
|
|
|
nounCount: { ...statistics.nounCount },
|
|
|
|
|
verbCount: { ...statistics.verbCount },
|
|
|
|
|
metadataCount: { ...statistics.metadataCount },
|
|
|
|
|
hnswIndexSize: statistics.hnswIndexSize,
|
|
|
|
|
lastUpdated: statistics.lastUpdated,
|
|
|
|
|
// Include serviceActivity if present
|
|
|
|
|
...(statistics.serviceActivity && {
|
|
|
|
|
serviceActivity: Object.fromEntries(
|
|
|
|
|
Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}])
|
|
|
|
|
)
|
|
|
|
|
}),
|
|
|
|
|
// Include services if present
|
|
|
|
|
...(statistics.services && {
|
|
|
|
|
services: statistics.services.map(s => ({...s}))
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Decrement the appropriate counter
|
|
|
|
|
const counterMap = {
|
|
|
|
|
noun: this.statisticsCache!.nounCount,
|
|
|
|
|
verb: this.statisticsCache!.verbCount,
|
|
|
|
|
metadata: this.statisticsCache!.metadataCount
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const counter = counterMap[type]
|
|
|
|
|
counter[service] = Math.max(0, (counter[service] || 0) - amount)
|
|
|
|
|
|
|
|
|
|
// Track service activity
|
|
|
|
|
this.trackServiceActivity(service, 'delete')
|
|
|
|
|
|
|
|
|
|
// Update timestamp
|
|
|
|
|
this.statisticsCache!.lastUpdated = new Date().toISOString()
|
|
|
|
|
|
|
|
|
|
// Schedule a batch update instead of saving immediately
|
|
|
|
|
this.scheduleBatchUpdate()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Update the HNSW index size statistic
|
|
|
|
|
* @param size The new size of the HNSW index
|
|
|
|
|
*/
|
|
|
|
|
async updateHnswIndexSize(size: number): Promise<void> {
|
|
|
|
|
// Get current statistics from cache or storage
|
|
|
|
|
let statistics = this.statisticsCache
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = await this.getStatisticsData()
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = this.createDefaultStatistics()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Update the cache
|
|
|
|
|
this.statisticsCache = {
|
|
|
|
|
nounCount: { ...statistics.nounCount },
|
|
|
|
|
verbCount: { ...statistics.verbCount },
|
|
|
|
|
metadataCount: { ...statistics.metadataCount },
|
|
|
|
|
hnswIndexSize: statistics.hnswIndexSize,
|
|
|
|
|
lastUpdated: statistics.lastUpdated,
|
|
|
|
|
// Include serviceActivity if present
|
|
|
|
|
...(statistics.serviceActivity && {
|
|
|
|
|
serviceActivity: Object.fromEntries(
|
|
|
|
|
Object.entries(statistics.serviceActivity).map(([k, v]) => [k, {...v}])
|
|
|
|
|
)
|
|
|
|
|
}),
|
|
|
|
|
// Include services if present
|
|
|
|
|
...(statistics.services && {
|
|
|
|
|
services: statistics.services.map(s => ({...s}))
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Update HNSW index size
|
|
|
|
|
this.statisticsCache!.hnswIndexSize = size
|
|
|
|
|
|
|
|
|
|
// Update timestamp
|
|
|
|
|
this.statisticsCache!.lastUpdated = new Date().toISOString()
|
|
|
|
|
|
|
|
|
|
// Schedule a batch update instead of saving immediately
|
|
|
|
|
this.scheduleBatchUpdate()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Force an immediate flush of statistics to storage
|
|
|
|
|
* This ensures that any pending statistics updates are written to persistent storage
|
|
|
|
|
*/
|
|
|
|
|
async flushStatisticsToStorage(): Promise<void> {
|
|
|
|
|
// If there are no statistics in cache or they haven't been modified, nothing to flush
|
|
|
|
|
if (!this.statisticsCache || !this.statisticsModified) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Call the protected flushStatistics method to immediately write to storage
|
|
|
|
|
await this.flushStatistics()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Track field names from a JSON document
|
|
|
|
|
* @param jsonDocument The JSON document to extract field names from
|
|
|
|
|
* @param service The service that inserted the data
|
|
|
|
|
*/
|
|
|
|
|
async trackFieldNames(jsonDocument: any, service: string): Promise<void> {
|
|
|
|
|
// Skip if not a JSON object
|
|
|
|
|
if (typeof jsonDocument !== 'object' || jsonDocument === null || Array.isArray(jsonDocument)) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Get current statistics from cache or storage
|
|
|
|
|
let statistics = this.statisticsCache
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = await this.getStatisticsData()
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = this.createDefaultStatistics()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Update the cache
|
|
|
|
|
this.statisticsCache = {
|
|
|
|
|
...statistics,
|
|
|
|
|
nounCount: { ...statistics.nounCount },
|
|
|
|
|
verbCount: { ...statistics.verbCount },
|
|
|
|
|
metadataCount: { ...statistics.metadataCount },
|
|
|
|
|
fieldNames: { ...statistics.fieldNames },
|
|
|
|
|
standardFieldMappings: { ...statistics.standardFieldMappings }
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Ensure fieldNames exists
|
|
|
|
|
if (!this.statisticsCache!.fieldNames) {
|
|
|
|
|
this.statisticsCache!.fieldNames = {}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Ensure standardFieldMappings exists
|
|
|
|
|
if (!this.statisticsCache!.standardFieldMappings) {
|
|
|
|
|
this.statisticsCache!.standardFieldMappings = {}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Extract field names from the JSON document
|
|
|
|
|
const fieldNames = extractFieldNamesFromJson(jsonDocument)
|
|
|
|
|
|
|
|
|
|
// Initialize service entry if it doesn't exist
|
|
|
|
|
if (!this.statisticsCache!.fieldNames[service]) {
|
|
|
|
|
this.statisticsCache!.fieldNames[service] = []
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Add new field names to the service's list
|
|
|
|
|
for (const fieldName of fieldNames) {
|
|
|
|
|
if (!this.statisticsCache!.fieldNames[service].includes(fieldName)) {
|
|
|
|
|
this.statisticsCache!.fieldNames[service].push(fieldName)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Map to standard field if possible
|
|
|
|
|
const standardField = mapToStandardField(fieldName)
|
|
|
|
|
if (standardField) {
|
|
|
|
|
// Initialize standard field entry if it doesn't exist
|
|
|
|
|
if (!this.statisticsCache!.standardFieldMappings[standardField]) {
|
|
|
|
|
this.statisticsCache!.standardFieldMappings[standardField] = {}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Initialize service entry if it doesn't exist
|
|
|
|
|
if (!this.statisticsCache!.standardFieldMappings[standardField][service]) {
|
|
|
|
|
this.statisticsCache!.standardFieldMappings[standardField][service] = []
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Add field name to standard field mapping if not already there
|
|
|
|
|
if (!this.statisticsCache!.standardFieldMappings[standardField][service].includes(fieldName)) {
|
|
|
|
|
this.statisticsCache!.standardFieldMappings[standardField][service].push(fieldName)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Update timestamp
|
|
|
|
|
this.statisticsCache!.lastUpdated = new Date().toISOString()
|
|
|
|
|
|
|
|
|
|
// Schedule a batch update
|
|
|
|
|
this.statisticsModified = true
|
|
|
|
|
this.scheduleBatchUpdate()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get available field names by service
|
|
|
|
|
* @returns Record of field names by service
|
|
|
|
|
*/
|
|
|
|
|
async getAvailableFieldNames(): Promise<Record<string, string[]>> {
|
|
|
|
|
// Get current statistics from cache or storage
|
|
|
|
|
let statistics = this.statisticsCache
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = await this.getStatisticsData()
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
return {}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Return field names by service
|
|
|
|
|
return statistics.fieldNames || {}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get standard field mappings
|
|
|
|
|
* @returns Record of standard field mappings
|
|
|
|
|
*/
|
|
|
|
|
async getStandardFieldMappings(): Promise<Record<string, Record<string, string[]>>> {
|
|
|
|
|
// Get current statistics from cache or storage
|
|
|
|
|
let statistics = this.statisticsCache
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
statistics = await this.getStatisticsData()
|
|
|
|
|
if (!statistics) {
|
|
|
|
|
return {}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Return standard field mappings
|
|
|
|
|
return statistics.standardFieldMappings || {}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Create default statistics data
|
|
|
|
|
* @returns Default statistics data
|
|
|
|
|
*/
|
|
|
|
|
protected createDefaultStatistics(): StatisticsData {
|
|
|
|
|
return {
|
|
|
|
|
nounCount: {},
|
|
|
|
|
verbCount: {},
|
|
|
|
|
metadataCount: {},
|
|
|
|
|
hnswIndexSize: 0,
|
|
|
|
|
fieldNames: {},
|
|
|
|
|
standardFieldMappings: {},
|
|
|
|
|
lastUpdated: new Date().toISOString()
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Detect if an error is a throttling error
|
|
|
|
|
* Override this method in specific adapters for custom detection
|
|
|
|
|
*/
|
|
|
|
|
protected isThrottlingError(error: any): boolean {
|
|
|
|
|
const statusCode = error.$metadata?.httpStatusCode || error.statusCode || error.code
|
|
|
|
|
const message = error.message?.toLowerCase() || ''
|
|
|
|
|
|
|
|
|
|
return (
|
|
|
|
|
statusCode === 429 || // Too Many Requests
|
|
|
|
|
statusCode === 503 || // Service Unavailable / Slow Down
|
|
|
|
|
statusCode === 'ECONNRESET' || // Connection reset
|
|
|
|
|
statusCode === 'ETIMEDOUT' || // Timeout
|
|
|
|
|
message.includes('throttl') ||
|
|
|
|
|
message.includes('slow down') ||
|
|
|
|
|
message.includes('rate limit') ||
|
|
|
|
|
message.includes('too many requests') ||
|
|
|
|
|
message.includes('quota exceeded')
|
|
|
|
|
)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Track a throttling event
|
|
|
|
|
* @param error The error that caused throttling
|
|
|
|
|
* @param service Optional service that was throttled
|
|
|
|
|
*/
|
|
|
|
|
protected trackThrottlingEvent(error: any, service?: string): void {
|
|
|
|
|
this.throttlingDetected = true
|
|
|
|
|
this.consecutiveThrottleEvents++
|
|
|
|
|
this.lastThrottleTime = Date.now()
|
|
|
|
|
this.totalThrottleEvents++
|
|
|
|
|
|
|
|
|
|
// Track by hour
|
|
|
|
|
const hourIndex = new Date().getHours()
|
|
|
|
|
if (hourIndex !== this.lastThrottleHourIndex) {
|
|
|
|
|
// Reset hour tracking if we've moved to a new hour
|
|
|
|
|
this.throttleEventsByHour = new Array(24).fill(0)
|
|
|
|
|
this.lastThrottleHourIndex = hourIndex
|
|
|
|
|
}
|
|
|
|
|
this.throttleEventsByHour[hourIndex]++
|
|
|
|
|
|
|
|
|
|
// Track throttle reason
|
|
|
|
|
const reason = this.getThrottleReason(error)
|
|
|
|
|
this.throttleReasons[reason] = (this.throttleReasons[reason] || 0) + 1
|
|
|
|
|
|
|
|
|
|
// Track service-level throttling
|
|
|
|
|
if (service) {
|
|
|
|
|
const serviceInfo = this.serviceThrottling.get(service) || {
|
|
|
|
|
throttleCount: 0,
|
|
|
|
|
lastThrottle: 0,
|
|
|
|
|
status: 'normal' as const
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
serviceInfo.throttleCount++
|
|
|
|
|
serviceInfo.lastThrottle = Date.now()
|
|
|
|
|
serviceInfo.status = 'throttled'
|
|
|
|
|
|
|
|
|
|
this.serviceThrottling.set(service, serviceInfo)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Exponential backoff
|
|
|
|
|
this.throttlingBackoffMs = Math.min(
|
|
|
|
|
this.throttlingBackoffMs * 2,
|
|
|
|
|
this.maxBackoffMs
|
|
|
|
|
)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get the reason for throttling from an error
|
|
|
|
|
*/
|
|
|
|
|
protected getThrottleReason(error: any): string {
|
|
|
|
|
const statusCode = error.$metadata?.httpStatusCode || error.statusCode || error.code
|
|
|
|
|
|
|
|
|
|
if (statusCode === 429) return '429_TooManyRequests'
|
|
|
|
|
if (statusCode === 503) return '503_ServiceUnavailable'
|
|
|
|
|
if (statusCode === 'ECONNRESET') return 'ConnectionReset'
|
|
|
|
|
if (statusCode === 'ETIMEDOUT') return 'Timeout'
|
|
|
|
|
|
|
|
|
|
const message = error.message?.toLowerCase() || ''
|
|
|
|
|
if (message.includes('throttl')) return 'Throttled'
|
|
|
|
|
if (message.includes('slow down')) return 'SlowDown'
|
|
|
|
|
if (message.includes('rate limit')) return 'RateLimit'
|
|
|
|
|
if (message.includes('quota exceeded')) return 'QuotaExceeded'
|
|
|
|
|
|
|
|
|
|
return 'Unknown'
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Clear throttling state after successful operations
|
|
|
|
|
*/
|
|
|
|
|
protected clearThrottlingState(): void {
|
|
|
|
|
if (this.consecutiveThrottleEvents > 0) {
|
|
|
|
|
this.consecutiveThrottleEvents = 0
|
|
|
|
|
this.throttlingBackoffMs = 1000 // Reset to initial backoff
|
|
|
|
|
|
|
|
|
|
if (this.throttlingDetected) {
|
|
|
|
|
this.throttlingDetected = false
|
|
|
|
|
|
|
|
|
|
// Update service statuses
|
|
|
|
|
for (const [service, info] of this.serviceThrottling) {
|
|
|
|
|
if (info.status === 'throttled') {
|
|
|
|
|
info.status = 'recovering'
|
|
|
|
|
} else if (info.status === 'recovering') {
|
|
|
|
|
const timeSinceThrottle = Date.now() - info.lastThrottle
|
|
|
|
|
if (timeSinceThrottle > 60000) { // 1 minute recovery period
|
|
|
|
|
info.status = 'normal'
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Handle throttling by implementing exponential backoff
|
|
|
|
|
* @param error The error that triggered throttling
|
|
|
|
|
* @param service Optional service that was throttled
|
|
|
|
|
*/
|
|
|
|
|
async handleThrottling(error: any, service?: string): Promise<void> {
|
|
|
|
|
if (this.isThrottlingError(error)) {
|
|
|
|
|
this.trackThrottlingEvent(error, service)
|
|
|
|
|
|
|
|
|
|
// Add delay for retry
|
|
|
|
|
const delayMs = this.throttlingBackoffMs
|
|
|
|
|
this.totalDelayMs += delayMs
|
|
|
|
|
this.delayedOperations++
|
|
|
|
|
|
|
|
|
|
await new Promise(resolve => setTimeout(resolve, delayMs))
|
|
|
|
|
} else {
|
|
|
|
|
// Clear throttling state on non-throttling errors
|
|
|
|
|
this.clearThrottlingState()
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Track a retried operation
|
|
|
|
|
*/
|
|
|
|
|
protected trackRetriedOperation(): void {
|
|
|
|
|
this.retriedOperations++
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Track an operation that failed due to throttling
|
|
|
|
|
*/
|
|
|
|
|
protected trackFailedDueToThrottling(): void {
|
|
|
|
|
this.failedDueToThrottling++
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get current throttling metrics
|
|
|
|
|
*/
|
|
|
|
|
protected getThrottlingMetrics(): StatisticsData['throttlingMetrics'] {
|
|
|
|
|
const averageDelayMs = this.delayedOperations > 0
|
|
|
|
|
? this.totalDelayMs / this.delayedOperations
|
|
|
|
|
: 0
|
|
|
|
|
|
|
|
|
|
// Convert service throttling map to record
|
|
|
|
|
const serviceThrottlingRecord: Record<string, {
|
|
|
|
|
throttleCount: number
|
|
|
|
|
lastThrottle: string
|
|
|
|
|
status: 'normal' | 'throttled' | 'recovering'
|
|
|
|
|
}> = {}
|
|
|
|
|
|
|
|
|
|
for (const [service, info] of this.serviceThrottling) {
|
|
|
|
|
serviceThrottlingRecord[service] = {
|
|
|
|
|
throttleCount: info.throttleCount,
|
|
|
|
|
lastThrottle: new Date(info.lastThrottle).toISOString(),
|
|
|
|
|
status: info.status
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return {
|
|
|
|
|
storage: {
|
|
|
|
|
currentlyThrottled: this.throttlingDetected,
|
|
|
|
|
lastThrottleTime: this.lastThrottleTime > 0
|
|
|
|
|
? new Date(this.lastThrottleTime).toISOString()
|
|
|
|
|
: undefined,
|
|
|
|
|
consecutiveThrottleEvents: this.consecutiveThrottleEvents,
|
|
|
|
|
currentBackoffMs: this.throttlingBackoffMs,
|
|
|
|
|
totalThrottleEvents: this.totalThrottleEvents,
|
|
|
|
|
throttleEventsByHour: [...this.throttleEventsByHour],
|
|
|
|
|
throttleReasons: { ...this.throttleReasons }
|
|
|
|
|
},
|
|
|
|
|
operationImpact: {
|
|
|
|
|
delayedOperations: this.delayedOperations,
|
|
|
|
|
retriedOperations: this.retriedOperations,
|
|
|
|
|
failedDueToThrottling: this.failedDueToThrottling,
|
|
|
|
|
averageDelayMs,
|
|
|
|
|
totalDelayMs: this.totalDelayMs
|
|
|
|
|
},
|
|
|
|
|
serviceThrottling: Object.keys(serviceThrottlingRecord).length > 0
|
|
|
|
|
? serviceThrottlingRecord
|
|
|
|
|
: undefined
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Include throttling metrics in statistics
|
|
|
|
|
*/
|
|
|
|
|
async getStatisticsWithThrottling(): Promise<StatisticsData | null> {
|
|
|
|
|
const stats = await this.getStatistics()
|
|
|
|
|
if (stats) {
|
|
|
|
|
stats.throttlingMetrics = this.getThrottlingMetrics()
|
|
|
|
|
}
|
|
|
|
|
return stats
|
|
|
|
|
}
|
2025-09-22 15:45:35 -07:00
|
|
|
|
|
|
|
|
// =============================================
|
|
|
|
|
// Universal O(1) Count Management
|
|
|
|
|
// =============================================
|
|
|
|
|
|
|
|
|
|
// Universal count tracking - O(1) operations
|
|
|
|
|
protected totalNounCount = 0
|
|
|
|
|
protected totalVerbCount = 0
|
|
|
|
|
protected entityCounts: Map<string, number> = new Map() // type -> count
|
|
|
|
|
protected verbCounts: Map<string, number> = new Map() // verb type -> count
|
|
|
|
|
protected countCache: Map<string, { count: number; timestamp: number }> = new Map()
|
|
|
|
|
protected readonly COUNT_CACHE_TTL = 60000 // 1 minute cache TTL
|
|
|
|
|
|
2025-10-09 17:35:01 -07:00
|
|
|
// =============================================
|
|
|
|
|
// Smart Count Batching (v3.32.3+)
|
|
|
|
|
// =============================================
|
|
|
|
|
|
|
|
|
|
// Count batching state - mirrors statistics batching pattern
|
|
|
|
|
protected pendingCountPersist = false // Counts changed since last persist?
|
|
|
|
|
protected lastCountPersistTime = 0 // Timestamp of last persist
|
|
|
|
|
protected scheduledCountPersistTimeout: NodeJS.Timeout | null = null // Scheduled persist timer
|
|
|
|
|
protected pendingCountOperations = 0 // Operations since last persist
|
|
|
|
|
|
|
|
|
|
// Batching configuration (overridable by subclasses for custom strategies)
|
|
|
|
|
protected countPersistBatchSize = 10 // Operations before forcing persist (cloud storage)
|
|
|
|
|
protected countPersistInterval = 5000 // Milliseconds before forcing persist (cloud storage)
|
|
|
|
|
|
2025-09-22 15:45:35 -07:00
|
|
|
/**
|
|
|
|
|
* Get total noun count - O(1) operation
|
|
|
|
|
* @returns Promise that resolves to the total number of nouns
|
|
|
|
|
*/
|
|
|
|
|
async getNounCount(): Promise<number> {
|
|
|
|
|
return this.totalNounCount
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Get total verb count - O(1) operation
|
|
|
|
|
* @returns Promise that resolves to the total number of verbs
|
|
|
|
|
*/
|
|
|
|
|
async getVerbCount(): Promise<number> {
|
|
|
|
|
return this.totalVerbCount
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Increment count for entity type - O(1) operation
|
|
|
|
|
* Protected by storage-specific mechanisms (mutex, distributed consensus, etc.)
|
|
|
|
|
* @param type The entity type
|
|
|
|
|
*/
|
|
|
|
|
protected incrementEntityCount(type: string): void {
|
|
|
|
|
// For distributed scenarios, this is aggregated across shards
|
|
|
|
|
// For single-node, this is protected by storage-specific locking
|
|
|
|
|
this.entityCounts.set(type, (this.entityCounts.get(type) || 0) + 1)
|
|
|
|
|
this.totalNounCount++
|
|
|
|
|
// Update cache
|
|
|
|
|
this.countCache.set('nouns_count', {
|
|
|
|
|
count: this.totalNounCount,
|
|
|
|
|
timestamp: Date.now()
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Thread-safe increment for concurrent scenarios
|
|
|
|
|
* Uses mutex for single-node, distributed consensus for multi-node
|
|
|
|
|
*/
|
|
|
|
|
protected async incrementEntityCountSafe(type: string): Promise<void> {
|
|
|
|
|
// Single-node mutex protection (distributed mode handled by coordinator)
|
|
|
|
|
const mutex = getGlobalMutex()
|
|
|
|
|
await mutex.runExclusive(`count-entity-${type}`, async () => {
|
|
|
|
|
this.incrementEntityCount(type)
|
2025-10-09 17:35:01 -07:00
|
|
|
// Smart batching (v3.32.3+): Adapts to storage type
|
|
|
|
|
// - Cloud storage (GCS, S3): Batches 10 ops OR 5 seconds
|
|
|
|
|
// - Local storage (File, Memory): Persists immediately
|
|
|
|
|
await this.scheduleCountPersist()
|
2025-09-22 15:45:35 -07:00
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Decrement count for entity type - O(1) operation
|
|
|
|
|
* @param type The entity type
|
|
|
|
|
*/
|
|
|
|
|
protected decrementEntityCount(type: string): void {
|
|
|
|
|
const current = this.entityCounts.get(type) || 0
|
|
|
|
|
if (current > 1) {
|
|
|
|
|
this.entityCounts.set(type, current - 1)
|
|
|
|
|
} else {
|
|
|
|
|
this.entityCounts.delete(type)
|
|
|
|
|
}
|
|
|
|
|
if (this.totalNounCount > 0) {
|
|
|
|
|
this.totalNounCount--
|
|
|
|
|
}
|
|
|
|
|
// Update cache
|
|
|
|
|
this.countCache.set('nouns_count', {
|
|
|
|
|
count: this.totalNounCount,
|
|
|
|
|
timestamp: Date.now()
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Thread-safe decrement for concurrent scenarios
|
|
|
|
|
*/
|
|
|
|
|
protected async decrementEntityCountSafe(type: string): Promise<void> {
|
|
|
|
|
const mutex = getGlobalMutex()
|
|
|
|
|
await mutex.runExclusive(`count-entity-${type}`, async () => {
|
|
|
|
|
this.decrementEntityCount(type)
|
2025-10-09 17:35:01 -07:00
|
|
|
// Smart batching (v3.32.3+): Adapts to storage type
|
|
|
|
|
await this.scheduleCountPersist()
|
2025-09-22 15:45:35 -07:00
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
fix(storage): resolve count synchronization race condition across all storage adapters
Fixed critical bug where entity and relationship counts were not being tracked correctly
during add(), relate(), and import() operations. The root cause was a race condition where
count increment code tried to read metadata before it was saved to storage.
Core Fixes:
- Modified baseStorage.saveNounMetadata_internal to increment counts AFTER metadata is saved
- Modified baseStorage.saveVerbMetadata_internal to increment verb counts AFTER metadata is saved
- Added verb type to VerbMetadata to avoid circular dependency during count tracking
- Refactored verb count methods to prevent mutex deadlocks (synchronous base + async Safe wrapper)
Storage Adapter Cleanup:
- Removed broken count increment code from FileSystemStorage, GcsStorage, R2Storage, AzureBlobStorage
- Updated MemoryStorage comments to reflect centralized fix
- All count tracking now centralized in baseStorage (fixes ALL adapters automatically)
New Utilities:
- Added rebuildCounts utility to repair corrupted counts.json from actual storage data
- Added comprehensive integration tests for count synchronization across all operations
Verification:
- All 8 storage adapters verified (FileSystem, GCS, Memory, S3Compatible, R2, Azure, OPFS, TypeAware)
- All code paths verified (add, relate, import, batch, update, delete)
- 599 tests passing (no regressions)
- No deadlocks (tests complete in 6s vs 150s+)
Fixes #1 and #2 reported by Workshop team
2025-10-21 10:58:44 -07:00
|
|
|
* Increment verb count - O(1) operation (v4.1.2: now synchronous)
|
|
|
|
|
* Protected by storage-specific mechanisms (mutex, distributed consensus, etc.)
|
|
|
|
|
* @param type The verb type
|
|
|
|
|
*/
|
|
|
|
|
protected incrementVerbCount(type: string): void {
|
|
|
|
|
this.verbCounts.set(type, (this.verbCounts.get(type) || 0) + 1)
|
|
|
|
|
this.totalVerbCount++
|
|
|
|
|
// Update cache
|
|
|
|
|
this.countCache.set('verbs_count', {
|
|
|
|
|
count: this.totalVerbCount,
|
|
|
|
|
timestamp: Date.now()
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Thread-safe increment for verb counts (v4.1.2)
|
|
|
|
|
* Uses mutex for single-node, distributed consensus for multi-node
|
2025-09-22 15:45:35 -07:00
|
|
|
* @param type The verb type
|
|
|
|
|
*/
|
fix(storage): resolve count synchronization race condition across all storage adapters
Fixed critical bug where entity and relationship counts were not being tracked correctly
during add(), relate(), and import() operations. The root cause was a race condition where
count increment code tried to read metadata before it was saved to storage.
Core Fixes:
- Modified baseStorage.saveNounMetadata_internal to increment counts AFTER metadata is saved
- Modified baseStorage.saveVerbMetadata_internal to increment verb counts AFTER metadata is saved
- Added verb type to VerbMetadata to avoid circular dependency during count tracking
- Refactored verb count methods to prevent mutex deadlocks (synchronous base + async Safe wrapper)
Storage Adapter Cleanup:
- Removed broken count increment code from FileSystemStorage, GcsStorage, R2Storage, AzureBlobStorage
- Updated MemoryStorage comments to reflect centralized fix
- All count tracking now centralized in baseStorage (fixes ALL adapters automatically)
New Utilities:
- Added rebuildCounts utility to repair corrupted counts.json from actual storage data
- Added comprehensive integration tests for count synchronization across all operations
Verification:
- All 8 storage adapters verified (FileSystem, GCS, Memory, S3Compatible, R2, Azure, OPFS, TypeAware)
- All code paths verified (add, relate, import, batch, update, delete)
- 599 tests passing (no regressions)
- No deadlocks (tests complete in 6s vs 150s+)
Fixes #1 and #2 reported by Workshop team
2025-10-21 10:58:44 -07:00
|
|
|
protected async incrementVerbCountSafe(type: string): Promise<void> {
|
2025-09-22 15:45:35 -07:00
|
|
|
const mutex = getGlobalMutex()
|
|
|
|
|
await mutex.runExclusive(`count-verb-${type}`, async () => {
|
fix(storage): resolve count synchronization race condition across all storage adapters
Fixed critical bug where entity and relationship counts were not being tracked correctly
during add(), relate(), and import() operations. The root cause was a race condition where
count increment code tried to read metadata before it was saved to storage.
Core Fixes:
- Modified baseStorage.saveNounMetadata_internal to increment counts AFTER metadata is saved
- Modified baseStorage.saveVerbMetadata_internal to increment verb counts AFTER metadata is saved
- Added verb type to VerbMetadata to avoid circular dependency during count tracking
- Refactored verb count methods to prevent mutex deadlocks (synchronous base + async Safe wrapper)
Storage Adapter Cleanup:
- Removed broken count increment code from FileSystemStorage, GcsStorage, R2Storage, AzureBlobStorage
- Updated MemoryStorage comments to reflect centralized fix
- All count tracking now centralized in baseStorage (fixes ALL adapters automatically)
New Utilities:
- Added rebuildCounts utility to repair corrupted counts.json from actual storage data
- Added comprehensive integration tests for count synchronization across all operations
Verification:
- All 8 storage adapters verified (FileSystem, GCS, Memory, S3Compatible, R2, Azure, OPFS, TypeAware)
- All code paths verified (add, relate, import, batch, update, delete)
- 599 tests passing (no regressions)
- No deadlocks (tests complete in 6s vs 150s+)
Fixes #1 and #2 reported by Workshop team
2025-10-21 10:58:44 -07:00
|
|
|
this.incrementVerbCount(type)
|
2025-10-09 17:35:01 -07:00
|
|
|
// Smart batching (v3.32.3+): Adapts to storage type
|
|
|
|
|
await this.scheduleCountPersist()
|
2025-09-22 15:45:35 -07:00
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
fix(storage): resolve count synchronization race condition across all storage adapters
Fixed critical bug where entity and relationship counts were not being tracked correctly
during add(), relate(), and import() operations. The root cause was a race condition where
count increment code tried to read metadata before it was saved to storage.
Core Fixes:
- Modified baseStorage.saveNounMetadata_internal to increment counts AFTER metadata is saved
- Modified baseStorage.saveVerbMetadata_internal to increment verb counts AFTER metadata is saved
- Added verb type to VerbMetadata to avoid circular dependency during count tracking
- Refactored verb count methods to prevent mutex deadlocks (synchronous base + async Safe wrapper)
Storage Adapter Cleanup:
- Removed broken count increment code from FileSystemStorage, GcsStorage, R2Storage, AzureBlobStorage
- Updated MemoryStorage comments to reflect centralized fix
- All count tracking now centralized in baseStorage (fixes ALL adapters automatically)
New Utilities:
- Added rebuildCounts utility to repair corrupted counts.json from actual storage data
- Added comprehensive integration tests for count synchronization across all operations
Verification:
- All 8 storage adapters verified (FileSystem, GCS, Memory, S3Compatible, R2, Azure, OPFS, TypeAware)
- All code paths verified (add, relate, import, batch, update, delete)
- 599 tests passing (no regressions)
- No deadlocks (tests complete in 6s vs 150s+)
Fixes #1 and #2 reported by Workshop team
2025-10-21 10:58:44 -07:00
|
|
|
* Decrement verb count - O(1) operation (v4.1.2: now synchronous)
|
|
|
|
|
* @param type The verb type
|
|
|
|
|
*/
|
|
|
|
|
protected decrementVerbCount(type: string): void {
|
|
|
|
|
const current = this.verbCounts.get(type) || 0
|
|
|
|
|
if (current > 1) {
|
|
|
|
|
this.verbCounts.set(type, current - 1)
|
|
|
|
|
} else {
|
|
|
|
|
this.verbCounts.delete(type)
|
|
|
|
|
}
|
|
|
|
|
if (this.totalVerbCount > 0) {
|
|
|
|
|
this.totalVerbCount--
|
|
|
|
|
}
|
|
|
|
|
// Update cache
|
|
|
|
|
this.countCache.set('verbs_count', {
|
|
|
|
|
count: this.totalVerbCount,
|
|
|
|
|
timestamp: Date.now()
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Thread-safe decrement for verb counts (v4.1.2)
|
2025-09-22 15:45:35 -07:00
|
|
|
* @param type The verb type
|
|
|
|
|
*/
|
fix(storage): resolve count synchronization race condition across all storage adapters
Fixed critical bug where entity and relationship counts were not being tracked correctly
during add(), relate(), and import() operations. The root cause was a race condition where
count increment code tried to read metadata before it was saved to storage.
Core Fixes:
- Modified baseStorage.saveNounMetadata_internal to increment counts AFTER metadata is saved
- Modified baseStorage.saveVerbMetadata_internal to increment verb counts AFTER metadata is saved
- Added verb type to VerbMetadata to avoid circular dependency during count tracking
- Refactored verb count methods to prevent mutex deadlocks (synchronous base + async Safe wrapper)
Storage Adapter Cleanup:
- Removed broken count increment code from FileSystemStorage, GcsStorage, R2Storage, AzureBlobStorage
- Updated MemoryStorage comments to reflect centralized fix
- All count tracking now centralized in baseStorage (fixes ALL adapters automatically)
New Utilities:
- Added rebuildCounts utility to repair corrupted counts.json from actual storage data
- Added comprehensive integration tests for count synchronization across all operations
Verification:
- All 8 storage adapters verified (FileSystem, GCS, Memory, S3Compatible, R2, Azure, OPFS, TypeAware)
- All code paths verified (add, relate, import, batch, update, delete)
- 599 tests passing (no regressions)
- No deadlocks (tests complete in 6s vs 150s+)
Fixes #1 and #2 reported by Workshop team
2025-10-21 10:58:44 -07:00
|
|
|
protected async decrementVerbCountSafe(type: string): Promise<void> {
|
2025-09-22 15:45:35 -07:00
|
|
|
const mutex = getGlobalMutex()
|
|
|
|
|
await mutex.runExclusive(`count-verb-${type}`, async () => {
|
fix(storage): resolve count synchronization race condition across all storage adapters
Fixed critical bug where entity and relationship counts were not being tracked correctly
during add(), relate(), and import() operations. The root cause was a race condition where
count increment code tried to read metadata before it was saved to storage.
Core Fixes:
- Modified baseStorage.saveNounMetadata_internal to increment counts AFTER metadata is saved
- Modified baseStorage.saveVerbMetadata_internal to increment verb counts AFTER metadata is saved
- Added verb type to VerbMetadata to avoid circular dependency during count tracking
- Refactored verb count methods to prevent mutex deadlocks (synchronous base + async Safe wrapper)
Storage Adapter Cleanup:
- Removed broken count increment code from FileSystemStorage, GcsStorage, R2Storage, AzureBlobStorage
- Updated MemoryStorage comments to reflect centralized fix
- All count tracking now centralized in baseStorage (fixes ALL adapters automatically)
New Utilities:
- Added rebuildCounts utility to repair corrupted counts.json from actual storage data
- Added comprehensive integration tests for count synchronization across all operations
Verification:
- All 8 storage adapters verified (FileSystem, GCS, Memory, S3Compatible, R2, Azure, OPFS, TypeAware)
- All code paths verified (add, relate, import, batch, update, delete)
- 599 tests passing (no regressions)
- No deadlocks (tests complete in 6s vs 150s+)
Fixes #1 and #2 reported by Workshop team
2025-10-21 10:58:44 -07:00
|
|
|
this.decrementVerbCount(type)
|
2025-10-09 17:35:01 -07:00
|
|
|
// Smart batching (v3.32.3+): Adapts to storage type
|
|
|
|
|
await this.scheduleCountPersist()
|
2025-09-22 15:45:35 -07:00
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
2025-10-09 17:35:01 -07:00
|
|
|
// =============================================
|
|
|
|
|
// Smart Batching Methods (v3.32.3+)
|
|
|
|
|
// =============================================
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Detect if this storage adapter uses cloud storage (network I/O)
|
|
|
|
|
* Cloud storage benefits from batching; local storage does not.
|
|
|
|
|
*
|
|
|
|
|
* Override this method in subclasses for accurate detection.
|
|
|
|
|
* Default implementation checks storage type from getStorageStatus().
|
|
|
|
|
*
|
|
|
|
|
* @returns true if cloud storage (GCS, S3, R2), false if local (File, Memory)
|
|
|
|
|
*/
|
|
|
|
|
protected isCloudStorage(): boolean {
|
|
|
|
|
// Default: assume local storage (conservative, prefers reliability over performance)
|
|
|
|
|
// Subclasses should override this for accurate detection
|
|
|
|
|
return false
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Schedule a smart batched persist operation.
|
|
|
|
|
*
|
|
|
|
|
* Strategy:
|
|
|
|
|
* - Local Storage: Persist immediately (fast, no network latency)
|
|
|
|
|
* - Cloud Storage: Batch persist (10 ops OR 5 seconds, whichever first)
|
|
|
|
|
*
|
|
|
|
|
* This mirrors the statistics batching pattern for consistency.
|
|
|
|
|
*/
|
|
|
|
|
protected async scheduleCountPersist(): Promise<void> {
|
|
|
|
|
// Mark counts as pending persist
|
|
|
|
|
this.pendingCountPersist = true
|
|
|
|
|
this.pendingCountOperations++
|
|
|
|
|
|
|
|
|
|
// Local storage: persist immediately (fast enough, no benefit from batching)
|
|
|
|
|
if (!this.isCloudStorage()) {
|
|
|
|
|
await this.flushCounts()
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Cloud storage: use smart batching
|
|
|
|
|
// Persist if we've hit the batch size threshold
|
|
|
|
|
if (this.pendingCountOperations >= this.countPersistBatchSize) {
|
|
|
|
|
await this.flushCounts()
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Otherwise, schedule a time-based persist if not already scheduled
|
|
|
|
|
if (!this.scheduledCountPersistTimeout) {
|
|
|
|
|
this.scheduledCountPersistTimeout = setTimeout(() => {
|
|
|
|
|
this.flushCounts().catch(error => {
|
|
|
|
|
console.error('Failed to flush counts on timer:', error)
|
|
|
|
|
})
|
|
|
|
|
}, this.countPersistInterval)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Flush counts immediately to storage.
|
|
|
|
|
*
|
|
|
|
|
* Used for:
|
|
|
|
|
* - Graceful shutdown (SIGTERM handler)
|
|
|
|
|
* - Forced persist (batch threshold reached)
|
|
|
|
|
* - Local storage immediate persist
|
|
|
|
|
*
|
|
|
|
|
* This is the public API that shutdown hooks can call.
|
|
|
|
|
*/
|
|
|
|
|
async flushCounts(): Promise<void> {
|
|
|
|
|
// Clear any scheduled persist
|
|
|
|
|
if (this.scheduledCountPersistTimeout) {
|
|
|
|
|
clearTimeout(this.scheduledCountPersistTimeout)
|
|
|
|
|
this.scheduledCountPersistTimeout = null
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Nothing to flush?
|
|
|
|
|
if (!this.pendingCountPersist) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
// Persist to storage (implemented by subclass)
|
|
|
|
|
await this.persistCounts()
|
|
|
|
|
|
|
|
|
|
// Update state
|
|
|
|
|
this.lastCountPersistTime = Date.now()
|
|
|
|
|
this.pendingCountPersist = false
|
|
|
|
|
this.pendingCountOperations = 0
|
|
|
|
|
} catch (error) {
|
|
|
|
|
console.error('❌ CRITICAL: Failed to flush counts to storage:', error)
|
|
|
|
|
// Keep pending flag set so we retry on next operation
|
|
|
|
|
throw error
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-09-22 15:45:35 -07:00
|
|
|
/**
|
|
|
|
|
* Initialize counts from storage - must be implemented by each adapter
|
|
|
|
|
* @protected
|
|
|
|
|
*/
|
|
|
|
|
protected abstract initializeCounts(): Promise<void>
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Persist counts to storage - must be implemented by each adapter
|
|
|
|
|
* @protected
|
|
|
|
|
*/
|
|
|
|
|
protected abstract persistCounts(): Promise<void>
|
2025-08-26 12:32:21 -07:00
|
|
|
}
|