brainy/src/aggregation/AggregationIndex.ts
David Snelling 945d92d29e fix: one field-resolution law across aggregation hooks, source.where, removeMany, and find() spellings
Four fixes from a consumer conformance report, one root disease — two
field-resolution regimes where there must be one:

- The delete/update aggregation hooks fed the engine a partial entity
  view (type/service/data/metadata only), so a reserved-field groupBy
  (subtype, visibility, ...) resolved to a nonexistent group on the way
  down: counts drifted upward forever after deletes, and updates moving
  an entity between reserved-field groups double-counted. The hooks now
  pass the full-fidelity view via entityForAggFromRawRecord (every
  reserved field top-level, mirroring the add path); the update sites
  pass the full get() view instead of a hand-rolled subset.
- Aggregation source.where resolved fields only against the custom
  metadata bag, so where on a reserved field silently matched nothing.
  The matcher now resolves each filtered field through
  resolveEntityField — the same single source of truth groupBy uses.
- removeMany() with no usable selector (bare array passed positionally,
  empty params, ids: []) resolved successfully having deleted nothing.
  All three now throw; the two legacy tests that pinned the silent
  no-op as 'graceful' now pin the refusal.
- find() where keys accept both spellings: a metadata.-prefixed key
  falls back to its flattened spelling when the prefixed one is not
  indexed (metadata is flattened at index time). A literal nested custom
  key named metadata still wins when indexed as spelled.

Five regression pins in aggregate-reserved-fields.test.ts (4 of 5 vary
red on the unfixed code).
2026-07-19 10:54:36 -07:00

1089 lines
40 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'
import { prodLog } from '../utils/logger.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
}
// Where filter — resolve each filtered field through resolveEntityField,
// the SAME single source of truth groupBy uses (top-level standard fields
// + custom metadata). Matching only the metadata sub-object made
// where:{subtype}/{visibility}/… a silent no-op: reserved fields never
// live in the custom bag, so those filters could never match anything.
if (source.where && Object.keys(source.where).length > 0) {
const e = entity as unknown as HNSWNounWithMetadata
const resolved: Record<string, unknown> = {}
for (const key of Object.keys(source.where)) {
resolved[key] = resolveEntityField(e, key)
}
if (!matchesMetadataFilter(resolved, 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>()
/**
* In-flight rescan targets. While a name has a staging map, ALL
* contributions (the walk's and concurrent write hooks') land there instead
* of the live map; the live map keeps serving until {@link finishBackfill}
* swaps the staging map in atomically.
*/
private backfillStaging = new Map<string, Map<string, AggregateGroupState>>()
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 {
if (this.pendingAdopt.size > 0) {
prodLog.info(
`[Aggregation] no adoptable persisted state for: ${Array.from(this.pendingAdopt).join(', ')} — flagged for backfill`
)
}
for (const name of this.pendingAdopt) this.needsBackfill.add(name)
this.pendingAdopt.clear()
}
/**
* May this persisted state be ADOPTED? When the store exposes its committed
* watermark, the state's `sourceGeneration` must EQUAL it: behind means
* later writes are missing from the state (unclean shutdown); ahead means
* it counts writes that no longer exist (e.g. a fact-log truncation on a
* copied store pulled the watermark back). Either way: one exact rescan,
* said out loud — never a silent adopt. Stores without the capability (and
* pre-stamp state on them) fall back to hash-only adoption.
*/
private stateGenerationAdoptable(name: string, stateData: unknown): boolean {
const committed = this.storage.committedGeneration?.() ?? null
if (committed === null) return true
const raw = (stateData as Record<string, unknown>).sourceGeneration
const stamped = typeof raw === 'number' ? raw : null
if (stamped === committed) return true
prodLog.warn(
`[Aggregation] '${name}': persisted state is at generation ${stamped ?? 'unstamped'} ` +
`but the store's committed generation is ${committed} — rescanning instead of adopting`
)
return false
}
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 &&
this.stateGenerationAdoptable(def.name, stateData)
) {
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)
prodLog.info(
`[Aggregation] '${def.name}': adopted persisted state (${groupMap.size} groups) — no rescan`
)
}
// 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 &&
this.stateGenerationAdoptable(def.name, stateData)
) {
// 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)
prodLog.info(
`[Aggregation] '${def.name}': restored definition + adopted persisted state (${groupMap.size} groups)`
)
} 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, stamped with the committed generation they
// reflect. The stamp is what makes reopen-adoption verifiable: state at a
// different generation than the store's committed watermark is stale (an
// unclean shutdown after later writes) or over-counts (a fact-log
// truncation on a copied store pulled the watermark BACK below the
// stamp) — either way the answer is one exact rescan, never a silent
// adopt. Read the generation after collecting groups so any racing
// commit resolves toward rescan, not wrong-adopt.
for (const name of this.dirty) {
const stateMap = this.states.get(name)
if (stateMap) {
const groups = Array.from(stateMap.values())
const sourceGeneration = this.storage.committedGeneration?.() ?? null
await this.storage.saveMetadata(
`${STATE_KEY_PREFIX}${name}__`,
sourceGeneration === null ? { groups } : { groups, sourceGeneration }
)
}
}
// 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)
}
/**
* Begin a rescan into a STAGING map. The live state is not touched — it
* keeps serving (possibly stale, but flagged pending) until the rescan
* completes and swaps in atomically. A mid-walk failure drops the staging
* map via {@link abortBackfill} and loses nothing: wiping live state before
* a scan that could throw was the destructive-before-durable defect.
* Contributions (walk + concurrent write hooks) land in staging while it
* exists, so the swapped-in result reflects writes that raced the walk.
*/
beginBackfill(name: string): void {
this.backfillStaging.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)
}
}
/**
* Abandon an in-flight rescan after a failure: drop the staging map, keep
* the live state serving, leave the aggregate flagged as pending so a later
* attempt rescans. The failure itself must be surfaced loudly by the owner.
*/
abortBackfill(name: string): void {
this.backfillStaging.delete(name)
}
/** 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)
}
}
/** Swap the rebuilt staging state in atomically; persists on next flush(). */
finishBackfill(name: string): void {
const staged = this.backfillStaging.get(name)
if (staged) {
this.states.set(name, staged)
this.backfillStaging.delete(name)
}
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.backfillStaging.get(aggName) ?? 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.backfillStaging.get(aggName) ?? 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.backfillStaging.get(aggName) ?? 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 }