Step 4 of the brainy 8.0 rename scaffolding. The storage-adapter contract gains the new algorithm-neutral method names. Existing concrete adapters keep their saveHNSWData / getHNSWData implementations untouched; the new names default to delegating to the old, so no concrete adapter needs to change in this commit. The final 8.0 cleanup removes the legacy names. CHANGES src/storage/adapters/baseStorageAdapter.ts - New default method `saveVectorIndexData(nounId, data)` that delegates to `saveHNSWData` for backward compatibility. Concrete adapters may override directly when the legacy name is removed. - New default method `getVectorIndexData(nounId)` that delegates to `getHNSWData`. Pre-existing abstract methods stay; concrete adapters (FileSystem, GCS, R2, S3, Azure, OPFS, Memory, Historical) continue to implement saveHNSWData / getHNSWData unchanged. PERSISTED FILE PATHS — DEFERRED The on-disk `_system/hnsw-*.json` → `_system/vector-index-*.json` rename is NOT in this commit. The new persisted-path layout requires: - dual-read on boot (accept either spelling — per integration doc lines 531-539: "Reader accepts either spelling on load; writer emits the new spelling only; brains self-migrate on the next persist after upgrade") - coordinated update across 8 storage adapters That work is folded into the final cleanup commit so the rename + migration ship atomically. NO-OP scope No behavioural change. New methods delegate to existing ones; no caller has been migrated yet. Tests + build green. VERIFICATION - npx tsc --noEmit: clean - npm test: 1468 / 1468 unit
1576 lines
50 KiB
TypeScript
1576 lines
50 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'
|
|
|
|
// =============================================================================
|
|
// Progressive Initialization Types
|
|
// =============================================================================
|
|
|
|
/**
|
|
* Initialization mode for cloud storage adapters.
|
|
*
|
|
* Controls how storage adapters initialize during startup, enabling
|
|
* fast cold starts in serverless environments while maintaining strict
|
|
* validation for local development.
|
|
*
|
|
*
|
|
* | Mode | Description | Use Case |
|
|
* |------|-------------|----------|
|
|
* | `'auto'` | Auto-detect environment (progressive in cloud, strict locally) | Default, zero-config |
|
|
* | `'progressive'` | Fast init (<200ms), validate lazily on first write | Serverless, Cloud Run |
|
|
* | `'strict'` | Blocking validation during init (traditional behavior) | Local dev, testing |
|
|
*
|
|
* @example
|
|
* ```typescript
|
|
* // Zero-config: auto-detects Cloud Run, Lambda, etc.
|
|
* const storage = new GCSStorageAdapter({ bucket: 'my-bucket' })
|
|
*
|
|
* // Force progressive mode for all environments
|
|
* const storage = new GCSStorageAdapter({
|
|
* bucket: 'my-bucket',
|
|
* initMode: 'progressive'
|
|
* })
|
|
*
|
|
* // Force strict validation (useful for testing)
|
|
* const storage = new GCSStorageAdapter({
|
|
* bucket: 'my-bucket',
|
|
* initMode: 'strict'
|
|
* })
|
|
* ```
|
|
*/
|
|
export type InitMode = 'progressive' | 'strict' | 'auto'
|
|
|
|
/**
|
|
* 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 getNounMetadata(id: string): Promise<any | null>
|
|
|
|
abstract getVerbMetadata(id: string): Promise<any | null>
|
|
|
|
// Vector Index Persistence (Brainy 8.0 — algorithm-neutral surface)
|
|
// The legacy `saveHNSWData` / `getHNSWData` names persist as default-method
|
|
// aliases so existing implementations continue to work mid-migration; the
|
|
// 8.0 final cleanup commit removes the old names.
|
|
|
|
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>
|
|
|
|
/**
|
|
* Persist vector-index graph state for an entity (Brainy 8.0 name).
|
|
* Default implementation delegates to {@link saveHNSWData}; concrete
|
|
* adapters may override directly. The final 8.0 cleanup removes the
|
|
* legacy `saveHNSWData` name.
|
|
*/
|
|
async saveVectorIndexData(nounId: string, data: {
|
|
level: number
|
|
connections: Record<string, string[]>
|
|
}): Promise<void> {
|
|
return this.saveHNSWData(nounId, data)
|
|
}
|
|
|
|
/**
|
|
* Read vector-index graph state for an entity (Brainy 8.0 name).
|
|
* Default implementation delegates to {@link getHNSWData}; concrete
|
|
* adapters may override directly. The final 8.0 cleanup removes the
|
|
* legacy `getHNSWData` name.
|
|
*/
|
|
async getVectorIndexData(nounId: string): Promise<{
|
|
level: number
|
|
connections: Record<string, string[]>
|
|
} | null> {
|
|
return this.getHNSWData(nounId)
|
|
}
|
|
|
|
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 deliberately NOT branch-scoped (COW): they
|
|
// are immutable, content-addressed segments managed by their producer,
|
|
// mirroring how COW metadata under `_cow/` is also kept global.
|
|
|
|
/**
|
|
* 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()
|
|
|
|
// =============================================
|
|
// Progressive Initialization State
|
|
// =============================================
|
|
// These properties enable fast cold starts in cloud environments
|
|
// by deferring validation and count loading to background tasks.
|
|
|
|
/**
|
|
* Initialization mode for this storage adapter.
|
|
*
|
|
* - `'auto'` (default): Progressive in cloud environments, strict locally
|
|
* - `'progressive'`: Always use fast init with lazy validation
|
|
* - `'strict'`: Always validate during init (traditional behavior)
|
|
*
|
|
* @protected
|
|
*/
|
|
protected initMode: InitMode = 'auto'
|
|
|
|
/**
|
|
* Whether the bucket/container has been validated to exist.
|
|
*
|
|
* In progressive mode, this starts as `false` and is set to `true`
|
|
* after the first successful write operation or background validation.
|
|
*
|
|
* @protected
|
|
*/
|
|
protected bucketValidated = false
|
|
|
|
/**
|
|
* Error from bucket validation, if any.
|
|
*
|
|
* Stored here so subsequent operations can fail fast with the same error.
|
|
*
|
|
* @protected
|
|
*/
|
|
protected bucketValidationError: Error | null = null
|
|
|
|
/**
|
|
* Whether counts have been loaded from storage.
|
|
*
|
|
* In progressive mode, counts start at 0 and are loaded in the background.
|
|
* Operations work immediately; counts become accurate after background load.
|
|
*
|
|
* @protected
|
|
*/
|
|
protected countsLoaded = false
|
|
|
|
/**
|
|
* Whether all background initialization tasks have completed.
|
|
*
|
|
* Useful for tests and diagnostics to ensure full initialization.
|
|
*
|
|
* @protected
|
|
*/
|
|
protected backgroundTasksComplete = false
|
|
|
|
/**
|
|
* Promise that resolves when background initialization completes.
|
|
*
|
|
* Can be awaited by callers who need to ensure full initialization
|
|
* before proceeding (e.g., in tests).
|
|
*
|
|
* @protected
|
|
*/
|
|
protected backgroundInitPromise: Promise<void> | null = null
|
|
|
|
// 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
|
|
|
|
// =============================================
|
|
// Smart Count Batching
|
|
// =============================================
|
|
|
|
// 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)
|
|
|
|
/**
|
|
* 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)
|
|
// 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)
|
|
* 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
|
|
* Uses mutex for single-node, distributed consensus for multi-node
|
|
* @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
|
|
// =============================================
|
|
|
|
/**
|
|
* 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
|
|
}
|
|
|
|
// =============================================
|
|
// Progressive Initialization Methods
|
|
// =============================================
|
|
|
|
/**
|
|
* Detect if running in a cloud/serverless environment.
|
|
*
|
|
* Cloud environments benefit from progressive initialization to minimize
|
|
* cold start times. This method checks for common environment variables
|
|
* set by cloud providers.
|
|
*
|
|
* | Provider | Environment Variable |
|
|
* |----------|---------------------|
|
|
* | Cloud Run | `K_SERVICE`, `K_REVISION` |
|
|
* | Lambda | `AWS_LAMBDA_FUNCTION_NAME` |
|
|
* | Cloud Functions | `FUNCTIONS_TARGET` |
|
|
* | Azure Functions | `AZURE_FUNCTIONS_ENVIRONMENT` |
|
|
*
|
|
* @returns `true` if running in a detected cloud environment
|
|
* @protected
|
|
*/
|
|
protected detectCloudEnvironment(): boolean {
|
|
return !!(
|
|
process.env.K_SERVICE || // Google Cloud Run
|
|
process.env.K_REVISION || // Google Cloud Run (revision)
|
|
process.env.AWS_LAMBDA_FUNCTION_NAME || // AWS Lambda
|
|
process.env.FUNCTIONS_TARGET || // Google Cloud Functions
|
|
process.env.AZURE_FUNCTIONS_ENVIRONMENT // Azure Functions
|
|
)
|
|
}
|
|
|
|
/**
|
|
* Determine the effective initialization mode.
|
|
*
|
|
* Resolves `'auto'` to either `'progressive'` or `'strict'` based on
|
|
* the detected environment.
|
|
*
|
|
* @returns The resolved init mode (`'progressive'` or `'strict'`)
|
|
* @protected
|
|
*/
|
|
protected resolveInitMode(): 'progressive' | 'strict' {
|
|
if (this.initMode === 'auto') {
|
|
return this.detectCloudEnvironment() ? 'progressive' : 'strict'
|
|
}
|
|
return this.initMode
|
|
}
|
|
|
|
/**
|
|
* Schedule background initialization tasks.
|
|
*
|
|
* This method is called in progressive mode after the adapter is marked
|
|
* as initialized. It schedules non-blocking tasks to:
|
|
* 1. Validate bucket/container existence
|
|
* 2. Load counts from storage
|
|
*
|
|
* Override in subclasses to add storage-specific background tasks.
|
|
* Always call `super.scheduleBackgroundInit()` first.
|
|
*
|
|
* @protected
|
|
*/
|
|
protected scheduleBackgroundInit(): void {
|
|
// Create a promise that tracks all background work
|
|
this.backgroundInitPromise = this.runBackgroundInit()
|
|
|
|
// Don't await - let it run in background
|
|
this.backgroundInitPromise
|
|
.then(() => {
|
|
this.backgroundTasksComplete = true
|
|
})
|
|
.catch((error) => {
|
|
// Log but don't throw - progressive mode is best-effort
|
|
console.error('[Progressive Init] Background initialization failed:', error)
|
|
this.backgroundTasksComplete = true // Mark complete even on error
|
|
})
|
|
}
|
|
|
|
/**
|
|
* Run background initialization tasks.
|
|
*
|
|
* Override in subclasses to add storage-specific tasks.
|
|
* The default implementation does nothing (for local storage adapters).
|
|
*
|
|
* @protected
|
|
*/
|
|
protected async runBackgroundInit(): Promise<void> {
|
|
// Default implementation: nothing to do for local storage
|
|
// Cloud storage adapters override this to:
|
|
// 1. Validate bucket/container
|
|
// 2. Load counts from storage
|
|
}
|
|
|
|
/**
|
|
* Wait for background initialization to complete.
|
|
*
|
|
* Useful in tests or when you need to ensure all background tasks
|
|
* have finished before proceeding.
|
|
*
|
|
* @example
|
|
* ```typescript
|
|
* const storage = new GCSStorageAdapter({ bucket: 'my-bucket' })
|
|
* await storage.init()
|
|
*
|
|
* // Wait for background validation and count loading
|
|
* await storage.awaitBackgroundInit()
|
|
*
|
|
* // Now counts are accurate and bucket is validated
|
|
* const count = await storage.getNounCount()
|
|
* ```
|
|
*
|
|
* @public
|
|
*/
|
|
public async awaitBackgroundInit(): Promise<void> {
|
|
if (this.backgroundInitPromise) {
|
|
await this.backgroundInitPromise
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Check if background initialization has completed.
|
|
*
|
|
* @returns `true` if all background tasks are complete
|
|
* @public
|
|
*/
|
|
public isBackgroundInitComplete(): boolean {
|
|
return this.backgroundTasksComplete
|
|
}
|
|
|
|
/**
|
|
* Ensure bucket/container is validated before write operations.
|
|
*
|
|
* In progressive mode, bucket validation is deferred until the first
|
|
* write operation. This method performs lazy validation and caches
|
|
* the result (or error) for subsequent calls.
|
|
*
|
|
* Override in subclasses to implement storage-specific validation.
|
|
* The base implementation does nothing (for local storage adapters).
|
|
*
|
|
* @throws Error if bucket validation fails
|
|
* @protected
|
|
*/
|
|
protected async ensureValidatedForWrite(): Promise<void> {
|
|
// If already validated, nothing to do
|
|
if (this.bucketValidated) {
|
|
return
|
|
}
|
|
|
|
// If we have a cached validation error, throw it
|
|
if (this.bucketValidationError) {
|
|
throw this.bucketValidationError
|
|
}
|
|
|
|
// Default implementation: no validation needed (local storage)
|
|
// Cloud storage adapters override this to validate bucket/container
|
|
this.bucketValidated = true
|
|
}
|
|
|
|
/**
|
|
* 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
|
|
}
|
|
}
|
|
|
|
/**
|
|
* 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>
|
|
}
|