- Version-coupling cold-init: getBrainyVersion() returned a stale build-time
default ('3.14.0') on the FIRST synchronous call — the one loadPlugins() makes
to build context.version — because the package.json read was async. A native
provider declaring a realistic `>=8.0.0` range was therefore rejected on cold
init. version.ts now reads package.json synchronously (8.0 targets Node-like
runtimes only); added a cold-init coupling regression with a `^8.0.0` plugin.
- No-silent-failures on the highest-fan-in read path: getNouns()/getVerbs()
converted a storage read failure into a success-shaped empty page, which the
cold-start rebuild then read as "store empty" and skipped the rebuild — booting
a permanently-empty index with no signal (the same silent-failure class as the
phantom bug). Both now re-throw a named BrainyError; the rebuild path re-throws
read failures (fail loud) and records a queryable degraded state, surfaced via
checkHealth(), for non-fatal rebuild hiccups.
- LSM SSTable leak: graph compaction wrote a merged SSTable but only dropped the
old ones from the manifest, orphaning their payloads forever ("In production
we'd add a cleanup mechanism"). Added StorageAdapter.deleteMetadata() and now
reclaim each compacted-away SSTable.
- Documented the two headline methods add() and find() (the only undocumented
public methods on the class).
Full gate green: build, unit 1512, integration 607.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1234 lines
39 KiB
TypeScript
1234 lines
39 KiB
TypeScript
/**
|
|
* Base Storage Adapter
|
|
* Provides common functionality for all storage adapters, including statistics tracking
|
|
*/
|
|
|
|
import {
|
|
StatisticsData,
|
|
StorageAdapter,
|
|
HNSWNoun,
|
|
HNSWVerb,
|
|
GraphVerb,
|
|
HNSWNounWithMetadata,
|
|
HNSWVerbWithMetadata,
|
|
NounMetadata,
|
|
VerbMetadata
|
|
} from '../../coreTypes.js'
|
|
import { StorageBatchConfig } from '../baseStorage.js'
|
|
import { extractFieldNamesFromJson, mapToStandardField } from '../../utils/fieldNameTracking.js'
|
|
import { getGlobalMutex, cleanupMutexes } from '../../utils/mutex.js'
|
|
|
|
/**
|
|
* 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>
|
|
|
|
abstract saveNoun(noun: HNSWNoun): Promise<void>
|
|
abstract saveNounMetadata(id: string, metadata: NounMetadata): Promise<void>
|
|
abstract deleteNounMetadata(id: string): Promise<void>
|
|
|
|
abstract getNoun(id: string): Promise<HNSWNounWithMetadata | null>
|
|
|
|
abstract getNounsByNounType(nounType: string): Promise<HNSWNounWithMetadata[]>
|
|
|
|
abstract deleteNoun(id: string): Promise<void>
|
|
|
|
abstract saveVerb(verb: HNSWVerb): Promise<void>
|
|
abstract saveVerbMetadata(id: string, metadata: VerbMetadata): Promise<void>
|
|
abstract deleteVerbMetadata(id: string): Promise<void>
|
|
|
|
abstract getVerb(id: string): Promise<HNSWVerbWithMetadata | null>
|
|
|
|
abstract getVerbsBySource(sourceId: string): Promise<HNSWVerbWithMetadata[]>
|
|
|
|
abstract getVerbsByTarget(targetId: string): Promise<HNSWVerbWithMetadata[]>
|
|
|
|
abstract getVerbsByType(type: string): Promise<HNSWVerbWithMetadata[]>
|
|
|
|
abstract deleteVerb(id: string): Promise<void>
|
|
|
|
abstract saveMetadata(id: string, metadata: any): Promise<void>
|
|
|
|
abstract getMetadata(id: string): Promise<any | null>
|
|
|
|
abstract deleteMetadata(id: string): Promise<void>
|
|
|
|
abstract getNounMetadata(id: string): Promise<any | null>
|
|
|
|
abstract getVerbMetadata(id: string): Promise<any | null>
|
|
|
|
// Vector Index Persistence (Brainy 8.0 — algorithm-neutral surface).
|
|
// Concrete adapters persist the per-entity HNSW graph node (level + connections)
|
|
// under these methods. Cortex's DiskANN-style native vector index doesn't use
|
|
// them (it persists its own single mmap'd `.dkann` file under
|
|
// `_system/vector-index/<shard>.dkann` per BRAINY-8.0-RENAME-COORDINATION § B.2).
|
|
|
|
abstract getNounVector(id: string): Promise<number[] | null>
|
|
|
|
abstract saveVectorIndexData(nounId: string, data: {
|
|
level: number
|
|
connections: Record<string, string[]>
|
|
}): Promise<void>
|
|
|
|
abstract getVectorIndexData(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>
|
|
|
|
abstract clear(): Promise<void>
|
|
|
|
abstract getStorageStatus(): Promise<{
|
|
type: string
|
|
used: number
|
|
quota: number | null
|
|
details?: Record<string, any>
|
|
}>
|
|
|
|
// ===========================================================================
|
|
// Raw binary-blob primitive
|
|
// ===========================================================================
|
|
//
|
|
// A first-class storage primitive for opaque byte payloads that must NOT be
|
|
// wrapped in a JSON envelope. The JSON object path base64-encodes binary data,
|
|
// inflating it ~33% and forcing full in-memory materialization on every read.
|
|
// Blobs side-step that: bytes are written and read verbatim.
|
|
//
|
|
// This unblocks zero-copy, mmap-able column-store segments and batch vector
|
|
// I/O at billion scale. On filesystem-backed adapters `getBinaryBlobPath`
|
|
// returns a real local path so native code (Rust) can mmap the file directly.
|
|
//
|
|
// Blob key → location convention (shared by every adapter so cross-language
|
|
// and cross-adapter consumers agree on exactly where a blob lives):
|
|
//
|
|
// key location
|
|
// "graph-lsm/source/sstable-123" "<root>/_blobs/graph-lsm/source/sstable-123.bin"
|
|
//
|
|
// i.e. the key's "/"-separated segments are nested under a `_blobs/` prefix and
|
|
// suffixed with `.bin`. Blobs are immutable, content-addressed segments
|
|
// managed by their producer.
|
|
|
|
/**
|
|
* Persist a raw binary blob under `key`, writing the bytes verbatim (no JSON
|
|
* envelope, no base64). Overwrites any existing blob at the same key. Where the
|
|
* backend is a real filesystem the write is atomic (temp file + rename) so a
|
|
* concurrent reader never observes a torn write.
|
|
*
|
|
* @param key - Logical blob key. "/"-separated segments map to nested
|
|
* directories under the adapter's `_blobs/` prefix (e.g.
|
|
* `"graph-lsm/source/sstable-123"`).
|
|
* @param data - The exact bytes to store.
|
|
* @returns Resolves once the blob is durably written.
|
|
* @example
|
|
* await storage.saveBinaryBlob('graph-lsm/source/sstable-7', segmentBytes)
|
|
*/
|
|
abstract saveBinaryBlob(key: string, data: Buffer): Promise<void>
|
|
|
|
/**
|
|
* Load the raw bytes previously stored under `key`, or `null` if no blob
|
|
* exists at that key. The returned buffer is byte-identical to what was passed
|
|
* to {@link saveBinaryBlob}.
|
|
*
|
|
* @param key - The blob key used when saving.
|
|
* @returns The blob bytes, or `null` if absent.
|
|
* @example
|
|
* const bytes = await storage.loadBinaryBlob('graph-lsm/source/sstable-7')
|
|
* if (bytes) decodeSegment(bytes)
|
|
*/
|
|
abstract loadBinaryBlob(key: string): Promise<Buffer | null>
|
|
|
|
/**
|
|
* Delete the blob stored under `key`. Missing blobs are ignored (no error) so
|
|
* delete is idempotent.
|
|
*
|
|
* @param key - The blob key to delete.
|
|
* @returns Resolves once the blob is gone (or was already absent).
|
|
*/
|
|
abstract deleteBinaryBlob(key: string): Promise<void>
|
|
|
|
/**
|
|
* Resolve `key` to a real local filesystem path that native code can `mmap`
|
|
* directly, or `null` when this backend has no local file for the blob.
|
|
*
|
|
* Filesystem-backed adapters return the on-disk path (whether or not the file
|
|
* currently exists — the caller is expected to write before mapping). Remote
|
|
* object stores (S3, R2, GCS, Azure), in-memory storage, browser storage
|
|
* (OPFS), and the read-only historical adapter genuinely have no local path
|
|
* and therefore return `null`. A `null` here is correct behavior, not a
|
|
* fallback: callers must use {@link loadBinaryBlob} when no path is available.
|
|
*
|
|
* @param key - The blob key.
|
|
* @returns An absolute local filesystem path, or `null` if none exists.
|
|
*/
|
|
abstract getBinaryBlobPath(key: string): string | null
|
|
|
|
/**
|
|
* Get optimal batch configuration for this storage adapter
|
|
* Override in subclasses to provide storage-specific optimization
|
|
*
|
|
* This method allows each storage adapter to declare its optimal batch behavior
|
|
* for rate limiting and performance. The configuration is used by addMany(),
|
|
* relateMany(), and import operations to automatically adapt to storage capabilities.
|
|
*
|
|
* @returns Batch configuration optimized for this storage type
|
|
*/
|
|
public getBatchConfig(): StorageBatchConfig {
|
|
// Conservative defaults that work safely across all storage types
|
|
// Cloud storage adapters should override with higher throughput values
|
|
// Local storage adapters should override with no delays
|
|
return {
|
|
maxBatchSize: 50,
|
|
batchDelayMs: 100,
|
|
maxConcurrent: 50,
|
|
supportsParallelWrites: false,
|
|
rateLimit: {
|
|
operationsPerSecond: 100,
|
|
burstCapacity: 200
|
|
}
|
|
}
|
|
}
|
|
|
|
// 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<{
|
|
items: HNSWNounWithMetadata[]
|
|
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<{
|
|
items: HNSWVerbWithMetadata[]
|
|
totalCount?: number
|
|
hasMore: boolean
|
|
nextCursor?: string
|
|
}>
|
|
|
|
/**
|
|
* 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<{
|
|
items: HNSWNounWithMetadata[]
|
|
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<{
|
|
items: HNSWVerbWithMetadata[]
|
|
totalCount?: number
|
|
hasMore: boolean
|
|
nextCursor?: string
|
|
}>
|
|
|
|
// 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
|
|
}
|
|
|
|
// =============================================
|
|
// 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
|
|
|
|
// =============================================
|
|
// Count Persistence
|
|
// =============================================
|
|
|
|
// Counts changed since the last persist? Drives the write-through flush.
|
|
protected pendingCountPersist = false
|
|
|
|
/**
|
|
* 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.
|
|
* Concurrency is handled by the process-global mutex
|
|
* ({@link incrementEntityCountSafe}); this raw form is for callers that
|
|
* already hold it.
|
|
* @param type The entity type
|
|
*/
|
|
protected incrementEntityCount(type: string): void {
|
|
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.
|
|
* A process-global mutex serialises the read-modify-write so concurrent
|
|
* writers in the same process cannot lose updates.
|
|
*/
|
|
protected async incrementEntityCountSafe(type: string): Promise<void> {
|
|
const mutex = getGlobalMutex()
|
|
await mutex.runExclusive(`count-entity-${type}`, async () => {
|
|
this.incrementEntityCount(type)
|
|
// Smart batching: Adapts to storage type
|
|
// - Cloud storage (GCS, S3): Batches 10 ops OR 5 seconds
|
|
// - Local storage (File, Memory): Persists immediately
|
|
await this.scheduleCountPersist()
|
|
})
|
|
}
|
|
|
|
/**
|
|
* 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)
|
|
// Smart batching: Adapts to storage type
|
|
await this.scheduleCountPersist()
|
|
})
|
|
}
|
|
|
|
/**
|
|
* Increment verb count - O(1) operation (now synchronous).
|
|
* Concurrency is handled by the process-global mutex
|
|
* ({@link incrementVerbCountSafe}); this raw form is for callers that already
|
|
* hold it.
|
|
* @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.
|
|
* A process-global mutex serialises the read-modify-write so concurrent
|
|
* writers in the same process cannot lose updates.
|
|
* @param type The verb type
|
|
*/
|
|
protected async incrementVerbCountSafe(type: string): Promise<void> {
|
|
const mutex = getGlobalMutex()
|
|
await mutex.runExclusive(`count-verb-${type}`, async () => {
|
|
this.incrementVerbCount(type)
|
|
// Smart batching: Adapts to storage type
|
|
await this.scheduleCountPersist()
|
|
})
|
|
}
|
|
|
|
/**
|
|
* Decrement verb count - O(1) operation (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
|
|
* @param type The verb type
|
|
*/
|
|
protected async decrementVerbCountSafe(type: string): Promise<void> {
|
|
const mutex = getGlobalMutex()
|
|
await mutex.runExclusive(`count-verb-${type}`, async () => {
|
|
this.decrementVerbCount(type)
|
|
// Smart batching: Adapts to storage type
|
|
await this.scheduleCountPersist()
|
|
})
|
|
}
|
|
|
|
// =============================================
|
|
// Smart Batching Methods
|
|
// =============================================
|
|
|
|
/**
|
|
* Persist counts immediately.
|
|
*
|
|
* Filesystem and memory storage have no network latency, so counts are
|
|
* written through on every change rather than batched.
|
|
*/
|
|
protected async scheduleCountPersist(): Promise<void> {
|
|
this.pendingCountPersist = true
|
|
await this.flushCounts()
|
|
}
|
|
|
|
/**
|
|
* Flush pending counts to storage.
|
|
*
|
|
* Used for graceful shutdown (SIGTERM handler) and the immediate
|
|
* write-through path. This is the public API that shutdown hooks can call.
|
|
*/
|
|
async flushCounts(): Promise<void> {
|
|
// Nothing to flush?
|
|
if (!this.pendingCountPersist) {
|
|
return
|
|
}
|
|
|
|
try {
|
|
// Persist to storage (implemented by subclass)
|
|
await this.persistCounts()
|
|
this.pendingCountPersist = false
|
|
} catch (error) {
|
|
console.error('CRITICAL: Failed to flush counts to storage:', error)
|
|
// Keep pending flag set so we retry on next operation
|
|
throw error
|
|
}
|
|
}
|
|
|
|
/**
|
|
* 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>
|
|
}
|