persistCounts() was write-through on every count change with no serialization, and the atomic writer named its temp file with millisecond granularity. Two persists inside one millisecond shared the temp path: both wrote it, the first rename consumed it, the second rename found nothing — ENOENT, roughly 1,500 times a day on a busy production brain, with a full ledger write per change behind it. No data was lost (the surviving rename carried a complete ledger and the next change re-persisted), but the race was real and the write rate absurd. flushCounts() now runs exactly one persist at a time; requests arriving during it collapse into one trailing pass that carries the burst's final state — N changes cost at most two writes. writeFileAtomic() adds a per-process sequence to the temp name so no two writes can share a path. Pinned: a 25-change burst → ≤2 ledger writes, zero errors, ledger equal to memory; parallel real writes land complete; three same-instant atomic writes own three distinct temp paths.
1401 lines
47 KiB
TypeScript
1401 lines
47 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,
|
||
CanonicalCounts,
|
||
} 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. Cor'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)
|
||
// Best-effort statistics flush — must not keep the process alive
|
||
// (close() flushes counts deterministically).
|
||
const statsTimer = this.statisticsBatchUpdateTimerId as unknown as { unref?: () => void }
|
||
if (statsTimer && typeof statsTimer.unref === 'function') {
|
||
statsTimer.unref()
|
||
}
|
||
}
|
||
|
||
/**
|
||
* 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
|
||
/**
|
||
* The ALL-visibility canonical scalars — every noun / verb the unfiltered
|
||
* storage walk yields, system and internal tiers included. These are the
|
||
* denominators a derived-index provider's coverage ledger subtracts from
|
||
* (`posted === all` is the whole-store coverage verdict); the user-facing
|
||
* `totalNounCount` / `totalVerbCount` skip hidden tiers by design and can
|
||
* never serve as a ledger denominator. Maintained on the write path
|
||
* (every new record +1, every proven delete −1), persisted beside the
|
||
* counted scalars, recomputed by the sanctioned recount. Never clamped.
|
||
*/
|
||
protected totalNounCountAll = 0
|
||
protected totalVerbCountAll = 0
|
||
/**
|
||
* The count of canonical nouns holding a REAL (non-empty) vector — the
|
||
* vector-side mirror of `totalNounCountAll` and the coverage denominator a
|
||
* vector index's node-count ledger is measured against. A deferred-embed
|
||
* noun (`add({ deferEmbedding: true })`) counts only once its vector
|
||
* LANDS (the `system:embed-landing` commit) — its canonical record exists
|
||
* (already counted in `totalNounCountAll`) with an empty vector until
|
||
* then. Maintained on the write path (a fresh insert whose vector is
|
||
* non-empty +1, a deferred embed's landing +1, a PROVEN delete of a
|
||
* vectored noun −1), persisted beside the other ALL scalars, recomputed by
|
||
* the sanctioned recount. Shares `allCountsSuspect` — no separate flag.
|
||
*/
|
||
protected totalVectoredNounCount = 0
|
||
/**
|
||
* `true` when a delete could not prove whether the record existed (no
|
||
* canonical read, no caller-provided prior) — the ALL scalar may be off by
|
||
* the unprovable deletes since. Loud, persisted, and cleared only by the
|
||
* sanctioned recount; a consumer reading the scalar as a ledger denominator
|
||
* must treat a suspect scalar as unverified, never as exact. Also covers
|
||
* `totalVectoredNounCount` — a delete whose vector-presence fact was
|
||
* unknowable marks this SAME flag rather than minting a second one.
|
||
*/
|
||
protected allCountsSuspect = false
|
||
/** One narration per session for the suspect transition (never per delete). */
|
||
private allCountsSuspectNarrated = false
|
||
/**
|
||
* Which rule produced the ALL scalars currently in memory. `'identity-record'`
|
||
* means one counted entity per metadata content leg — the honest rule: a
|
||
* bare id-directory (a ghost or scar left by a partial-delete defect, no
|
||
* content leg) counts zero. Set by the one-time derivation and by the
|
||
* sanctioned recount, alongside `allCountsSuspect = false`; left `undefined`
|
||
* when a loaded counts.json carries the ALL scalars but no stamp — the
|
||
* legacy container-rule derivation, which forces `allCountsSuspect = true`
|
||
* at load instead. A filesystem concern: `MemoryStorage` has no counts.json
|
||
* and never sets this.
|
||
*/
|
||
protected allCountsDerivedBy?: 'identity-record'
|
||
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
|
||
/** The one persist running right now, if any (single-flight law — see flushCounts). */
|
||
private countPersistInFlight: Promise<void> | null = null
|
||
/** The one trailing persist a burst has queued behind the in-flight one. */
|
||
private countPersistTrailing: Promise<void> | null = null
|
||
|
||
/**
|
||
* 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
|
||
}
|
||
|
||
/**
|
||
* The canonical count ledger — O(1), no I/O. `counted` is the user-facing
|
||
* scalar (public/internal tiers, what `getNounCount()` returns); `all` is
|
||
* the ALL-visibility scalar every unfiltered storage walk is measured
|
||
* against (the coverage-ledger denominator for derived-index providers);
|
||
* `vectors.all` is the vectored-noun scalar — the coverage denominator for
|
||
* a vector index's node-count ledger specifically.
|
||
* `suspect` is `true` when an unprovable delete has made `all` (any
|
||
* family, including `vectors`) unverified since the last sanctioned
|
||
* recount (`rebuildTypeCounts`).
|
||
* @returns All scalars per family plus the suspect flag.
|
||
*/
|
||
async getCanonicalCounts(): Promise<CanonicalCounts> {
|
||
return {
|
||
nouns: { counted: this.totalNounCount, all: this.totalNounCountAll },
|
||
verbs: { counted: this.totalVerbCount, all: this.totalVerbCountAll },
|
||
vectors: { all: this.totalVectoredNounCount },
|
||
suspect: this.allCountsSuspect
|
||
}
|
||
}
|
||
|
||
/**
|
||
* Mark the ALL scalars unverified after a delete that could not prove the
|
||
* record existed. Narrates ONCE per session (the flag is what persists);
|
||
* the sanctioned recount clears it.
|
||
* @param family - Which family's delete was unprovable.
|
||
* @param id - The id whose existence could not be established.
|
||
*/
|
||
protected markAllCountsSuspect(family: 'noun' | 'verb' | 'noun-vector', id: string): void {
|
||
this.allCountsSuspect = true
|
||
if (!this.allCountsSuspectNarrated) {
|
||
this.allCountsSuspectNarrated = true
|
||
console.warn(
|
||
`[Storage] ${family} delete of ${id} could not prove the record existed ` +
|
||
`(no canonical read, no prior record) — the ALL-visibility count ledger is ` +
|
||
`SUSPECT until brain.repairIndex() recounts. Further unprovable deletes ` +
|
||
`this session are counted silently under the same flag.`
|
||
)
|
||
}
|
||
}
|
||
|
||
/**
|
||
* OPTIONAL narrow ledger hook (see {@link StorageAdapter.noteVectorLanded}):
|
||
* record a deferred-embed noun's FIRST real vector landing. The caller
|
||
* (the deferred-embed worker) proves this is a genuine landing — not a
|
||
* re-embed of an already-vectored row — by observing its own pre-embed
|
||
* read's vector was empty, at no added storage cost.
|
||
* @param id - The noun whose vector just landed (retained for a future
|
||
* narration seam; the count itself needs no id-keyed state).
|
||
*/
|
||
async noteVectorLanded(id: string): Promise<void> {
|
||
void id
|
||
this.totalVectoredNounCount++
|
||
this.scheduleCountPersist().catch(() => {
|
||
// Ignore persist errors — the in-memory count is authoritative; a later op retries.
|
||
})
|
||
}
|
||
|
||
/**
|
||
* OPTIONAL narrow ledger hook (see {@link StorageAdapter.noteVectorUnlanded}):
|
||
* the mirror of {@link noteVectorLanded} — record a noun's vector was just
|
||
* REMOVED (rewritten to the unvectored `[]` shape). Never below zero: a
|
||
* caller that (incorrectly) fires this for a noun already unvectored would
|
||
* otherwise drive the ledger negative — clamped defensively, matching the
|
||
* delete path's `if (this.totalVectoredNounCount > 0)` guard.
|
||
* @param id - The noun whose vector was just removed (retained for a
|
||
* future narration seam; the count itself needs no id-keyed state).
|
||
*/
|
||
async noteVectorUnlanded(id: string): Promise<void> {
|
||
void id
|
||
if (this.totalVectoredNounCount > 0) this.totalVectoredNounCount--
|
||
this.scheduleCountPersist().catch(() => {
|
||
// Ignore persist errors — the in-memory count is authoritative; a later op retries.
|
||
})
|
||
}
|
||
|
||
/**
|
||
* 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
|
||
}
|
||
|
||
// SINGLE-FLIGHT, COALESCED. Counts are write-through on every change, so
|
||
// a burst of writes used to launch one persist per change, all in flight
|
||
// together. Two of them inside the same millisecond shared the atomic
|
||
// writer's temp path (`.tmp-<pid>-<ms>`): both wrote it, the first rename
|
||
// consumed it, the second rename found nothing — ENOENT, ~1,500 times a
|
||
// day on a busy production brain, with a full ledger write per change
|
||
// behind it. Now exactly one persist runs at a time; requests that arrive
|
||
// while it runs collapse into ONE trailing persist that carries the final
|
||
// state. A burst of N changes costs at most two writes and never races
|
||
// itself.
|
||
if (this.countPersistInFlight) {
|
||
// The in-flight write may have already serialised a stale snapshot —
|
||
// ask for one more pass after it, and let every caller in this burst
|
||
// await that same pass.
|
||
if (!this.countPersistTrailing) {
|
||
this.countPersistTrailing = this.countPersistInFlight
|
||
.catch(() => undefined)
|
||
.then(() => {
|
||
this.countPersistTrailing = null
|
||
return this.flushCounts()
|
||
})
|
||
}
|
||
return this.countPersistTrailing
|
||
}
|
||
|
||
this.countPersistInFlight = (async () => {
|
||
try {
|
||
// Persist to storage (implemented by subclass)
|
||
this.pendingCountPersist = false
|
||
await this.persistCounts()
|
||
} catch (error) {
|
||
// Keep the flag set so the next operation retries.
|
||
this.pendingCountPersist = true
|
||
console.error('CRITICAL: Failed to flush counts to storage:', error)
|
||
throw error
|
||
} finally {
|
||
this.countPersistInFlight = null
|
||
}
|
||
})()
|
||
return this.countPersistInFlight
|
||
}
|
||
|
||
/**
|
||
* 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>
|
||
}
|