brainy/src/aggregation/AggregationIndex.ts
David Snelling da55be7520 fix: aggregation state adoption on reopen + single-flight backfill + query-cap ratchet removal
- AggregationIndex: defineAggregate before init no longer forces a backfill.
  init reconciles instead of clobbering: the app definition wins, persisted
  state is adopted on hash match, a write landing pre-adoption forces an exact
  rescan, and a successful state load clears the backfill flag. New ready()
  settles every adoption decision before query paths consult backfill state.
- brainy: backfills are single-flight and batched. Concurrent queries share one
  store walk and every pending aggregate fills from that same walk; the old
  behavior let each concurrent query wipe the others' partial state and start
  its own full walk, so a store under steady aggregate traffic never converged.
  queryAggregate also waits for persisted definitions before deciding an
  aggregate does not exist.
- paramValidation: recordQuery is telemetry-only. The duration-based cap
  ratchet (x0.8 per recorded query while lifetime-average exceeded 1s, floored
  at 1000 - below the documented 10000 auto floor, reported under a stale
  basis label) is removed; the cap comes from its construction-time basis or
  explicit overrides alone.
- storage: __aggregation_* and singleton system keys (brainy:entityIdMapper)
  are recognized before the unknown-key warning fires; routing is unchanged.
- docs: find-limits cap-immutability note, aggregation reopen/adoption
  semantics, RELEASES.md 8.5.1 entry.
2026-07-17 09:20:09 -07:00

1000 lines
35 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

/**
* AggregationIndex - Incremental Write-Time Aggregation Engine
*
* Maintains running totals on every add/update/delete for O(1) reads.
* Follows the same pattern as MetadataIndexManager and GraphAdjacencyIndex.
*
* Key design decisions:
* - Source matching reuses matchesMetadataFilter() from metadataFilter.ts
* - Entities with service='brainy:aggregation' or metadata.__aggregate are
* skipped to prevent infinite loops (materialized entities feeding back)
* - MIN/MAX on delete are set to NaN and lazy-recomputed on next query
* - Persistence uses storage.saveMetadata() / storage.getMetadata()
* - Delta detection hashes definitions; only changed aggregates rebuild on restart
*/
import type { StorageAdapter, HNSWNounWithMetadata } from '../coreTypes.js'
import { resolveEntityField } from '../coreTypes.js'
import type {
AggregateDefinition,
AggregateGroupState,
AggregateQueryParams,
AggregateResult,
AggregationOp,
AggregationProvider,
GroupByDimension,
MetricState
} from '../types/brainy.types.js'
import { matchesMetadataFilter } from '../utils/metadataFilter.js'
import { compareCodePoints } from '../utils/collation.js'
import { bucketTimestamp } from './timeWindows.js'
import { NounType } from '../types/graphTypes.js'
/** Persistence key for aggregate definitions */
const DEFINITIONS_KEY = '__aggregation_definitions__'
/** Prefix for per-aggregate state persistence keys */
const STATE_KEY_PREFIX = '__aggregation_state_'
/**
* Serialize a group key map into a deterministic string for use as a Map key.
*/
function serializeGroupKey(groupKey: Record<string, string | number>): string {
const sorted = Object.keys(groupKey).sort()
return sorted.map(k => `${k}=${groupKey[k]}`).join('|')
}
/**
* Hash an aggregate definition for change detection on restart.
*/
function hashDefinition(def: AggregateDefinition): string {
// Deterministic JSON: sort keys
const normalized = JSON.stringify({
name: def.name,
source: def.source,
groupBy: def.groupBy,
metrics: Object.keys(def.metrics).sort().map(k => [k, def.metrics[k]])
})
// Simple FNV-1a 32-bit hash
let hash = 0x811c9dc5
for (let i = 0; i < normalized.length; i++) {
hash ^= normalized.charCodeAt(i)
hash = (hash * 0x01000193) >>> 0
}
return hash.toString(16)
}
/**
* Create a fresh MetricState with identity values.
*/
function freshMetricState(): MetricState {
return { sum: 0, count: 0, min: Infinity, max: -Infinity, m2: 0 }
}
/**
* Check whether an entity matches an aggregate's source filter.
*/
function matchesSource(entity: Record<string, unknown>, source: AggregateDefinition['source']): boolean {
// Type filter
if (source.type) {
const entityType = entity.type ?? entity.noun
const types = Array.isArray(source.type) ? source.type : [source.type]
if (!types.includes(entityType as NounType)) return false
}
// Service filter
if (source.service) {
if (entity.service !== source.service) return false
}
// Metadata where filter — match against the entity's metadata sub-object
if (source.where && Object.keys(source.where).length > 0) {
const metadata = (entity.metadata ?? entity) as Record<string, unknown>
if (!matchesMetadataFilter(metadata, source.where)) return false
}
return true
}
/**
* Compute the group key for an entity given groupBy dimensions.
*
* Field lookups go through `resolveEntityField` to honor the
* HNSWNounWithMetadata shape contract (standard fields top-level,
* custom fields in metadata) in a single source of truth.
*/
/**
* Compute the group key(s) an entity contributes to.
*
* Returns one key normally, but **multiple** when a groupBy dimension is an `unnest` field:
* the entity contributes once per distinct array element (cartesian product across multiple
* unnest dimensions). An entity whose unnest field is missing/empty contributes to no group
* (returns `[]`). This is the fan-out behind tag-frequency style aggregates.
*/
function computeGroupKeys(
entity: Record<string, unknown>,
groupBy: GroupByDimension[]
): Record<string, string | number>[] {
const e = entity as unknown as HNSWNounWithMetadata
let keys: Record<string, string | number>[] = [{}]
for (const dim of groupBy) {
if (typeof dim === 'string') {
const val = resolveEntityField(e, dim)
const v = val !== undefined && val !== null ? String(val) : '__null__'
for (const k of keys) k[dim] = v
} else if ('unnest' in dim) {
const val = resolveEntityField(e, dim.field)
const raw = Array.isArray(val) ? val : val !== undefined && val !== null ? [val] : []
// Distinct elements: an entity with duplicate tags counts once per distinct tag.
const elems = Array.from(new Set(raw.map(x => String(x))))
if (elems.length === 0) return [] // contributes to no group
const next: Record<string, string | number>[] = []
for (const k of keys) {
for (const el of elems) next.push({ ...k, [dim.field]: el })
}
keys = next
} else {
// Time-windowed field
const val = resolveEntityField(e, dim.field)
const v = typeof val === 'number' ? bucketTimestamp(val, dim.window) : '__null__'
for (const k of keys) k[dim.field] = v
}
}
return keys
}
/**
* Single representative group key (first of {@link computeGroupKeys}). Retained for the
* materializer and back-compat; the incremental contribution paths use `computeGroupKeys`
* so unnest dimensions fan out correctly.
*/
function computeGroupKey(
entity: Record<string, unknown>,
groupBy: GroupByDimension[]
): Record<string, string | number> {
return computeGroupKeys(entity, groupBy)[0] ?? {}
}
/**
* Get the numeric value of a field from an entity (for metric computation).
*
* Routes through `resolveEntityField` so standard top-level numeric fields
* (weight, confidence, createdAt, updatedAt) and custom user numeric fields
* in metadata are both handled in one place.
*/
function getNumericField(entity: Record<string, unknown>, field: string): number | undefined {
const val = resolveEntityField(entity as unknown as HNSWNounWithMetadata, field)
if (typeof val === 'number' && !isNaN(val)) return val
if (typeof val === 'string') {
const num = parseFloat(val)
if (!isNaN(num)) return num
}
return undefined
}
/**
* Check if an entity is a materialized aggregate entity (to prevent infinite loops).
*/
function isAggregateEntity(entity: Record<string, unknown>): boolean {
if (entity.service === 'brainy:aggregation') return true
const metadata = (entity.metadata ?? {}) as Record<string, unknown>
if (metadata.__aggregate) return true
return false
}
/**
* Welford's online algorithm: add a value to running mean/M2.
*/
function updateMetricAdd(state: MetricState, val: number, op: AggregationOp): void {
state.sum += val
state.count++
if (val < state.min) state.min = val
if (val > state.max) state.max = val
// Welford's: update M2 for stddev/variance
if (op === 'stddev' || op === 'variance') {
const mean = state.sum / state.count
const oldMean = state.count > 1 ? (state.sum - val) / (state.count - 1) : 0
state.m2 = (state.m2 ?? 0) + (val - oldMean) * (val - mean)
}
// Track the numeric value multiset for percentile AND min/max. min/max need it
// to recompute the extreme after a delete WITHOUT scanning entities — that scan
// was the 7.x infinite-loop / hang (materialized aggregates fed back in).
// (distinctCount is tracked separately, on raw values.)
if (op === 'percentile' || op === 'min' || op === 'max') {
if (!state.valueCounts) state.valueCounts = {}
const key = String(val)
state.valueCounts[key] = (state.valueCounts[key] ?? 0) + 1
}
}
/**
* Welford's online algorithm: remove a value from running mean/M2.
* Note: removing from Welford's is the inverse update.
*/
function updateMetricRemove(state: MetricState, val: number, op: AggregationOp): void {
// Decrement the numeric value multiset for percentile AND min/max (drop the key
// at zero). (distinctCount is decremented separately, on raw values.)
if ((op === 'percentile' || op === 'min' || op === 'max') && state.valueCounts) {
const key = String(val)
const c = state.valueCounts[key]
if (c !== undefined) {
if (c <= 1) delete state.valueCounts[key]
else state.valueCounts[key] = c - 1
}
}
if (state.count <= 1) {
state.sum = 0
state.count = 0
state.m2 = 0
return
}
const oldMean = state.sum / state.count
state.sum -= val
state.count--
const newMean = state.sum / state.count
if (op === 'stddev' || op === 'variance') {
state.m2 = Math.max(0, (state.m2 ?? 0) - (val - oldMean) * (val - newMean))
}
}
/**
* Recompute a metric's min/max from its value multiset — O(distinct values), fully
* in-memory, NO entity scan. Called lazily on query when a delete of the current
* extreme marked it stale. (Scanning entities to recompute is what hung 7.x.)
*/
function recomputeMinMaxFromCounts(state: MetricState): void {
const keys = state.valueCounts ? Object.keys(state.valueCounts) : []
if (keys.length === 0) {
state.min = Infinity
state.max = -Infinity
return
}
let mn = Infinity
let mx = -Infinity
for (const k of keys) {
const n = Number(k)
if (n < mn) mn = n
if (n > mx) mx = n
}
state.min = mn
state.max = mx
}
/**
* Exact percentile over a value multiset, using linear interpolation between closest ranks
* (numpy 'linear' / type-7): `rank = p·(n1)`, interpolating between the floor and ceil
* positions of the sorted expansion. Mirrors Cor's `MetricState::percentile` bit-for-bit.
*/
function computePercentile(valueCounts: Record<string, number>, count: number, p: number): number {
if (count === 0) return 0
const pp = Math.max(0, Math.min(1, p))
const entries = Object.entries(valueCounts)
.map(([v, c]) => [parseFloat(v), c] as [number, number])
.sort((a, b) => a[0] - b[0])
if (entries.length === 0) return 0
if (count === 1) return entries[0][0]
const rank = pp * (count - 1)
const lo = Math.floor(rank)
const hi = Math.ceil(rank)
const vLo = valueAtRank(entries, lo)
if (lo === hi) return vLo
const vHi = valueAtRank(entries, hi)
return vLo + (vHi - vLo) * (rank - lo)
}
/** Value at 0-based rank `idx` in the sorted expansion of the multiset. */
function valueAtRank(entries: Array<[number, number]>, idx: number): number {
let cum = 0
for (const [v, c] of entries) {
cum += c
if (idx < cum) return v
}
return entries.length ? entries[entries.length - 1][0] : 0
}
export class AggregationIndex {
private storage: StorageAdapter
private nativeProvider?: AggregationProvider
/** Registered aggregate definitions keyed by name */
private definitions = new Map<string, AggregateDefinition>()
/** Hashes of definitions for change detection */
private definitionHashes = new Map<string, string>()
/** Per-aggregate group states: Map<aggName, Map<serializedGroupKey, AggregateGroupState>> */
private states = new Map<string, Map<string, AggregateGroupState>>()
/** Track which aggregates have dirty state needing persistence */
private dirty = new Set<string>()
/**
* Aggregates whose state must be backfilled from entities that already existed
* when the aggregate was defined (or whose persisted state was missing/stale at
* init). Write-time hooks only capture entities added *after* a definition, so
* without backfill an aggregate defined over a populated store stays empty.
* Drained by the owner (Brainy) which has the entity iterator; see `getPendingBackfills`.
*/
private needsBackfill = new Set<string>()
/** Track aggregates with stale MIN/MAX (need lazy recompute) */
private staleMinMax = new Map<string, Set<string>>()
/** Resolves when init() has finished loading persisted definitions/state. */
private initPromise: Promise<void> | null = null
/** True once init() has settled (success or failure). */
private initDone = false
/**
* Aggregates registered by the app before init() finished loading persisted
* state, awaiting reconciliation: init() adopts the persisted state when the
* definition hash matches; anything left unadopted when init settles resolves
* to a backfill. Deciding backfill eagerly at define time was the boot-order
* bug that wiped valid persisted state on every restart — the synchronous
* defineAggregate() always beats the async init().
*/
private pendingAdopt = new Set<string>()
constructor(storage: StorageAdapter, nativeProvider?: AggregationProvider) {
this.storage = storage
this.nativeProvider = nativeProvider
}
// ============= Lifecycle =============
/**
* Initialize: load persisted definitions and state, detect changes, rebuild stale.
*
* Idempotent — repeated calls return the same promise. Definitions registered
* *before* this completes (the normal boot order: `defineAggregate()` is
* synchronous and always beats this async load) are reconciled rather than
* clobbered: the app's definition wins, and its persisted state is adopted
* when the definition hash matches — backfill happens only on a real change.
*/
init(): Promise<void> {
if (!this.initPromise) {
this.initPromise = this.loadPersisted().finally(() => {
this.resolvePendingAdoptToBackfill()
this.initDone = true
})
}
return this.initPromise
}
/**
* Await the persisted-state load (if one was started) and settle every
* pending adoption decision. After this resolves, `getPendingBackfills()`
* is authoritative: a name is listed iff it genuinely needs a rescan.
* Query paths must await this before consulting backfill state.
*/
async ready(): Promise<void> {
if (this.initPromise) {
try {
await this.initPromise
} catch {
// The owner already surfaced the load failure loudly; backfill covers.
}
}
this.resolvePendingAdoptToBackfill()
}
/**
* Any definition still awaiting state adoption has no persisted state to
* adopt (or init never ran / failed) — it must backfill.
*/
private resolvePendingAdoptToBackfill(): void {
for (const name of this.pendingAdopt) this.needsBackfill.add(name)
this.pendingAdopt.clear()
}
private async loadPersisted(): Promise<void> {
// Load persisted definitions
const savedDefs = await this.storage.getMetadata(DEFINITIONS_KEY)
if (savedDefs && typeof savedDefs === 'object' && savedDefs.definitions) {
const defs = savedDefs.definitions as Array<AggregateDefinition & { _hash?: string }>
for (const def of defs) {
const savedHash = def._hash || ''
if (this.definitions.has(def.name)) {
// The app re-registered this aggregate before the load finished.
// The app's definition wins — never clobber it with the persisted
// copy. Adopt the persisted state when the definition is unchanged
// AND no write has landed for it yet (a landed write would be lost
// by adoption; the hook flips such names to backfill).
const appHash = this.definitionHashes.get(def.name) || ''
if (appHash === savedHash && this.pendingAdopt.has(def.name)) {
const stateData = await this.storage.getMetadata(`${STATE_KEY_PREFIX}${def.name}__`)
if (stateData && stateData.groups) {
const groupMap = new Map<string, AggregateGroupState>()
for (const group of stateData.groups as AggregateGroupState[]) {
groupMap.set(serializeGroupKey(group.groupKey), group)
}
this.states.set(def.name, groupMap)
this.pendingAdopt.delete(def.name)
this.needsBackfill.delete(def.name)
}
// No/invalid persisted state: stays in pendingAdopt and resolves
// to backfill when init settles.
}
continue
}
// Not registered this session — restore definition + state from
// persistence.
this.definitions.set(def.name, def)
const currentHash = hashDefinition(def)
const stateData = await this.storage.getMetadata(`${STATE_KEY_PREFIX}${def.name}__`)
if (stateData && stateData.groups && savedHash === currentHash) {
// Definition unchanged — load state
const groupMap = new Map<string, AggregateGroupState>()
for (const group of stateData.groups as AggregateGroupState[]) {
const serialized = serializeGroupKey(group.groupKey)
groupMap.set(serialized, group)
}
this.states.set(def.name, groupMap)
this.needsBackfill.delete(def.name)
} else {
// Definition changed or no saved state — start fresh and backfill from
// existing entities (the owner drains needsBackfill on first query).
this.states.set(def.name, new Map())
this.needsBackfill.add(def.name)
}
this.definitionHashes.set(def.name, currentHash)
// Register definition with native provider
if (this.nativeProvider?.defineAggregate) {
this.nativeProvider.defineAggregate(def)
}
}
}
// Restore native provider state from persistence
if (this.nativeProvider?.restoreState) {
const nativeState = await this.storage.getMetadata('__aggregation_native_state__')
if (nativeState && typeof nativeState === 'string') {
this.nativeProvider.restoreState(nativeState)
} else if (nativeState && typeof nativeState === 'object' && nativeState.data) {
// flush() persists `{ data: serializeState() }`, so `data` is the
// provider's serialized state string.
this.nativeProvider.restoreState(nativeState.data as string)
}
}
}
/**
* Persist all dirty aggregate state to storage.
*/
async flush(): Promise<void> {
// Persist definitions
const defsToSave = Array.from(this.definitions.values()).map(def => ({
...def,
_hash: this.definitionHashes.get(def.name)
}))
await this.storage.saveMetadata(DEFINITIONS_KEY, { definitions: defsToSave })
// Persist dirty states
for (const name of this.dirty) {
const stateMap = this.states.get(name)
if (stateMap) {
const groups = Array.from(stateMap.values())
await this.storage.saveMetadata(
`${STATE_KEY_PREFIX}${name}__`,
{ groups }
)
}
}
// Persist native provider state
if (this.nativeProvider?.serializeState) {
const nativeState = this.nativeProvider.serializeState()
await this.storage.saveMetadata(
'__aggregation_native_state__',
{ data: nativeState }
)
}
this.dirty.clear()
}
/**
* Flush and release resources.
*/
async close(): Promise<void> {
await this.flush()
}
// ============= Definition Management =============
/**
* Register a new aggregate definition. Persisted on next flush().
*/
defineAggregate(def: AggregateDefinition): void {
if (!def.name) throw new Error('Aggregate definition requires a name')
if (!def.groupBy || def.groupBy.length === 0) throw new Error('Aggregate definition requires at least one groupBy dimension')
if (!def.metrics || Object.keys(def.metrics).length === 0) throw new Error('Aggregate definition requires at least one metric')
// Validate metric definitions
for (const [name, metric] of Object.entries(def.metrics)) {
if (metric.op !== 'count' && !metric.field) {
throw new Error(`Metric '${name}' with op '${metric.op}' requires a 'field' property`)
}
}
const newHash = hashDefinition(def)
const oldHash = this.definitionHashes.get(def.name)
this.definitions.set(def.name, def)
this.definitionHashes.set(def.name, newHash)
// First sight this session, before init() settled: defer the backfill
// decision — init() adopts the persisted state on hash match, and anything
// left unadopted resolves to backfill. Deciding eagerly here wiped valid
// persisted state on every restart.
if (!this.states.has(def.name) && !this.initDone) {
this.states.set(def.name, new Map())
this.pendingAdopt.add(def.name)
}
// Reset state if definition changed or doesn't exist yet, and flag it for
// backfill so already-stored entities are counted (write-time hooks only see
// future writes). The owner drains this on the next query via getPendingBackfills().
else if (!this.states.has(def.name) || (oldHash && oldHash !== newHash)) {
this.pendingAdopt.delete(def.name)
this.states.set(def.name, new Map())
this.needsBackfill.add(def.name)
}
// Notify native provider of definition (caches compiled form for hot path)
if (this.nativeProvider?.defineAggregate) {
this.nativeProvider.defineAggregate(def)
}
this.dirty.add(def.name)
}
/**
* Remove an aggregate definition and its state.
*/
removeAggregate(name: string): void {
this.definitions.delete(name)
this.definitionHashes.delete(name)
this.states.delete(name)
this.staleMinMax.delete(name)
this.pendingAdopt.delete(name)
this.needsBackfill.delete(name)
// Notify native provider
if (this.nativeProvider?.removeAggregate) {
this.nativeProvider.removeAggregate(name)
}
this.dirty.add(name) // Will persist the removal
}
/**
* Get all registered aggregate definitions.
*/
getDefinitions(): AggregateDefinition[] {
return Array.from(this.definitions.values())
}
/**
* Check if an aggregate exists.
*/
hasAggregate(name: string): boolean {
return this.definitions.has(name)
}
// ============= Backfill =============
//
// Write-time hooks only capture entities added after a definition exists, so an
// aggregate defined over a populated store would stay empty. The owner (Brainy) has
// the entity iterator, so backfill is driven from there: it reads the pending set,
// clears the aggregate, streams every existing entity through `backfillEntity`, then
// calls `finishBackfill`. Clearing first means a concurrent write that landed via the
// incremental hook is wiped and re-counted exactly once by the rescan.
/** Names of aggregates whose state must be (re)built from existing entities. */
getPendingBackfills(): string[] {
return Array.from(this.needsBackfill)
}
/** Clear an aggregate's state so a full rescan cannot double-count. */
beginBackfill(name: string): void {
this.states.set(name, new Map())
// Reset native provider state for this aggregate too, if present.
const def = this.definitions.get(name)
if (def && this.nativeProvider?.removeAggregate && this.nativeProvider?.defineAggregate) {
this.nativeProvider.removeAggregate(name)
this.nativeProvider.defineAggregate(def)
}
}
/** Feed one already-stored entity into a single aggregate during backfill. */
backfillEntity(name: string, entity: Record<string, unknown>): void {
if (isAggregateEntity(entity)) return
const def = this.definitions.get(name)
if (!def || !matchesSource(entity, def.source)) return
if (this.nativeProvider) {
this.applyNativeResults(name, this.nativeProvider.incrementalUpdate(name, def, entity, 'add'))
} else {
this.addContribution(name, def, entity)
}
}
/** Mark an aggregate's backfill complete; rebuilt state persists on next flush(). */
finishBackfill(name: string): void {
this.needsBackfill.delete(name)
this.dirty.add(name)
}
// ============= Write-Time Hooks =============
/**
* A write is landing for an aggregate whose persisted-state adoption is still
* pending — adopting after this write would lose its contribution. Settle the
* decision now: an exact rescan instead of adoption. The window is the few
* milliseconds between a boot-time defineAggregate() and init() completing,
* so this rarely fires; when it does, correctness wins over the walk.
*/
private resolveAdoptOnWrite(name: string): void {
if (this.pendingAdopt.has(name)) {
this.pendingAdopt.delete(name)
this.needsBackfill.add(name)
}
}
/**
* Called when an entity is added. Updates all matching aggregates.
*/
onEntityAdded(id: string, entity: Record<string, unknown>): void {
if (isAggregateEntity(entity)) return
for (const [name, def] of this.definitions) {
if (!matchesSource(entity, def.source)) continue
this.resolveAdoptOnWrite(name)
if (this.nativeProvider) {
const results = this.nativeProvider.incrementalUpdate(name, def, entity, 'add')
this.applyNativeResults(name, results)
continue
}
// Shared with onEntityUpdated/backfill; fans out unnest dimensions to N groups.
this.addContribution(name, def, entity)
}
}
/**
* Called when an entity is updated. Reverses old contribution and applies new.
*/
onEntityUpdated(
id: string,
newEntity: Record<string, unknown>,
oldEntity: Record<string, unknown>
): void {
if (isAggregateEntity(newEntity)) return
for (const [name, def] of this.definitions) {
const oldMatches = matchesSource(oldEntity, def.source)
const newMatches = matchesSource(newEntity, def.source)
if (oldMatches || newMatches) {
this.resolveAdoptOnWrite(name)
}
if (this.nativeProvider && (oldMatches || newMatches)) {
const results = this.nativeProvider.incrementalUpdate(name, def, newEntity, 'update', oldEntity)
this.applyNativeResults(name, results)
continue
}
// If old matched, remove its contribution
if (oldMatches) {
this.removeContribution(name, def, oldEntity)
}
// If new matches, add its contribution
if (newMatches) {
this.addContribution(name, def, newEntity)
}
}
}
/**
* Called when an entity is deleted. Reverses its contribution.
*/
onEntityDeleted(id: string, entity: Record<string, unknown>): void {
if (isAggregateEntity(entity)) return
for (const [name, def] of this.definitions) {
if (!matchesSource(entity, def.source)) continue
this.resolveAdoptOnWrite(name)
if (this.nativeProvider) {
const results = this.nativeProvider.incrementalUpdate(name, def, entity, 'delete')
this.applyNativeResults(name, results)
continue
}
this.removeContribution(name, def, entity)
}
}
// ============= Query =============
/**
* Query aggregate results with optional filtering, sorting, and pagination.
*/
queryAggregate(params: AggregateQueryParams): AggregateResult[] {
const def = this.definitions.get(params.name)
if (!def) throw new Error(`Aggregate '${params.name}' not found`)
const stateMap = this.states.get(params.name)
if (!stateMap) return []
if (this.nativeProvider) {
return this.nativeProvider.queryAggregate(stateMap, params)
}
// Collect all groups
let results: AggregateResult[] = []
for (const [serialized, group] of stateMap.entries()) {
// Skip empty groups (all metrics at zero count)
const hasData = Object.values(group.metrics).some(m => m.count > 0)
if (!hasData) continue
// Apply where filter on group keys
if (params.where && Object.keys(params.where).length > 0) {
if (!matchesMetadataFilter(group.groupKey, params.where)) continue
}
// Compute result metrics from running state
const metrics: Record<string, number> = {}
let totalCount = 0
for (const [metricName, metricDef] of Object.entries(def.metrics)) {
const state = group.metrics[metricName]
if (!state) continue
switch (metricDef.op) {
case 'count':
metrics[metricName] = state.count
break
case 'sum':
metrics[metricName] = state.sum
break
case 'avg':
metrics[metricName] = state.count > 0 ? state.sum / state.count : 0
break
case 'min':
case 'max': {
// Lazily recompute the extreme from the value multiset if a delete of
// the current min/max marked it stale (hang-free; no entity scan).
const staleSet = this.staleMinMax.get(params.name)
if (staleSet && staleSet.has(`${serialized}:${metricName}`)) {
recomputeMinMaxFromCounts(state)
staleSet.delete(`${serialized}:${metricName}`)
}
metrics[metricName] =
metricDef.op === 'min'
? state.min === Infinity
? 0
: state.min
: state.max === -Infinity
? 0
: state.max
break
}
case 'variance':
metrics[metricName] = state.count > 1 ? (state.m2 ?? 0) / (state.count - 1) : 0
break
case 'stddev':
metrics[metricName] = state.count > 1 ? Math.sqrt((state.m2 ?? 0) / (state.count - 1)) : 0
break
case 'percentile':
metrics[metricName] = computePercentile(state.valueCounts ?? {}, state.count, metricDef.p ?? 0.5)
break
case 'distinctCount':
metrics[metricName] = state.valueCounts ? Object.keys(state.valueCounts).length : 0
break
}
if (metricDef.op === 'count') {
totalCount = Math.max(totalCount, state.count)
} else {
totalCount = Math.max(totalCount, state.count)
}
}
// HAVING: filter groups by computed metric values (post-compute, O(groups), before
// sort/pagination). Reuses the where-operator engine over metrics + `count`.
if (params.having && Object.keys(params.having).length > 0) {
if (!matchesMetadataFilter({ ...metrics, count: totalCount }, params.having)) continue
}
results.push({
groupKey: { ...group.groupKey },
metrics,
count: totalCount,
entityId: group.materializedEntityId
})
}
// Sort
if (params.orderBy) {
const field = params.orderBy
const dir = params.order === 'desc' ? -1 : 1
results.sort((a, b) => {
// Try metrics first, then groupKey
const aVal = a.metrics[field] ?? a.groupKey[field] ?? 0
const bVal = b.metrics[field] ?? b.groupKey[field] ?? 0
if (typeof aVal === 'number' && typeof bVal === 'number') {
return (aVal - bVal) * dir
}
return compareCodePoints(String(aVal), String(bVal)) * dir
})
}
// Pagination
const offset = params.offset || 0
const limit = params.limit || results.length
results = results.slice(offset, offset + limit)
return results
}
/**
* Get internal state for a named aggregate (for materialization and testing).
*/
getState(name: string): Map<string, AggregateGroupState> | undefined {
return this.states.get(name)
}
// ============= Internal Helpers =============
/**
* Add an entity's contribution to its matching group in a named aggregate.
*/
private addContribution(
aggName: string,
def: AggregateDefinition,
entity: Record<string, unknown>
): void {
const stateMap = this.states.get(aggName)!
// Fan out: an unnest dimension makes one entity contribute to several groups.
for (const groupKey of computeGroupKeys(entity, def.groupBy)) {
const serialized = serializeGroupKey(groupKey)
let group = stateMap.get(serialized)
if (!group) {
group = {
groupKey,
metrics: {},
lastUpdated: Date.now()
}
for (const metricName of Object.keys(def.metrics)) {
group.metrics[metricName] = freshMetricState()
}
stateMap.set(serialized, group)
}
for (const [metricName, metricDef] of Object.entries(def.metrics)) {
const state = group.metrics[metricName]
if (metricDef.op === 'count') {
state.count++
state.sum++
} else if (metricDef.op === 'distinctCount') {
// distinctCount tracks distinct values of ANY type (strings, numbers, booleans),
// keyed by their string form — NOT numeric-coerced, since its primary use is
// categorical (distinct categories / users / tags), not numeric columns.
const raw = resolveEntityField(entity as unknown as HNSWNounWithMetadata, metricDef.field!)
if (raw !== undefined && raw !== null) {
if (!state.valueCounts) state.valueCounts = {}
const key = String(raw)
state.valueCounts[key] = (state.valueCounts[key] ?? 0) + 1
state.count++
}
} else {
const val = getNumericField(entity, metricDef.field!)
if (val !== undefined) {
updateMetricAdd(state, val, metricDef.op)
}
}
}
group.lastUpdated = Date.now()
}
this.dirty.add(aggName)
}
/**
* Remove an entity's contribution from its matching group.
* For MIN/MAX, marks as stale since we can't incrementally reverse these.
*/
private removeContribution(
aggName: string,
def: AggregateDefinition,
entity: Record<string, unknown>
): void {
const stateMap = this.states.get(aggName)!
// Fan out: reverse the entity's contribution from every group it joined.
for (const groupKey of computeGroupKeys(entity, def.groupBy)) {
const serialized = serializeGroupKey(groupKey)
const group = stateMap.get(serialized)
if (!group) continue
for (const [metricName, metricDef] of Object.entries(def.metrics)) {
const state = group.metrics[metricName]
if (metricDef.op === 'count') {
state.count = Math.max(0, state.count - 1)
state.sum = Math.max(0, state.sum - 1)
} else if (metricDef.op === 'distinctCount') {
const raw = resolveEntityField(entity as unknown as HNSWNounWithMetadata, metricDef.field!)
if (raw !== undefined && raw !== null && state.valueCounts) {
const key = String(raw)
const c = state.valueCounts[key]
if (c !== undefined) {
if (c <= 1) delete state.valueCounts[key]
else state.valueCounts[key] = c - 1
}
state.count = Math.max(0, state.count - 1)
}
} else {
const val = getNumericField(entity, metricDef.field!)
if (val !== undefined) {
updateMetricRemove(state, val, metricDef.op)
// MIN/MAX can't be decremented — mark as potentially stale
if (val <= state.min || val >= state.max) {
if (!this.staleMinMax.has(aggName)) {
this.staleMinMax.set(aggName, new Set())
}
this.staleMinMax.get(aggName)!.add(`${serialized}:${metricName}`)
}
}
}
}
// Remove group if all metrics are empty
const allEmpty = Object.values(group.metrics).every(m => m.count === 0)
if (allEmpty) {
stateMap.delete(serialized)
} else {
group.lastUpdated = Date.now()
}
}
this.dirty.add(aggName)
}
/**
* Apply results from native provider back into the state maps.
*/
private applyNativeResults(aggName: string, results: AggregateGroupState[]): void {
const stateMap = this.states.get(aggName)!
for (const group of results) {
const serialized = serializeGroupKey(group.groupKey)
stateMap.set(serialized, group)
}
this.dirty.add(aggName)
}
}
// Export helper for use by materializer
export { serializeGroupKey, computeGroupKey, matchesSource, isAggregateEntity }