From aabd4a396acc522486f45c4a28b42b2cf8747e7c Mon Sep 17 00:00:00 2001 From: Joey Yakimowich-Payne Date: Thu, 16 Apr 2026 14:26:26 -0600 Subject: [PATCH] =?UTF-8?q?feat(rete):=20add=20FilterNode=20(P1.9),=20Exis?= =?UTF-8?q?tentialNode=20(P2.2),=20AggregationNode=20(P2.4)=20=E2=80=94=20?= =?UTF-8?q?commit=20missing=20files?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/rete/src/aggregate.test.ts | 293 ++++++++++++++++++++++++++ packages/rete/src/aggregate.ts | 270 ++++++++++++++++++++++++ packages/rete/src/condition.test.ts | 146 +++++++++++++ packages/rete/src/condition.ts | 115 ++++++++++ packages/rete/src/existential.test.ts | 285 +++++++++++++++++++++++++ packages/rete/src/existential.ts | 217 +++++++++++++++++++ packages/rete/src/index.ts | 5 + packages/rete/src/registry.ts | 17 +- 8 files changed, 1347 insertions(+), 1 deletion(-) create mode 100644 packages/rete/src/aggregate.test.ts create mode 100644 packages/rete/src/aggregate.ts create mode 100644 packages/rete/src/condition.test.ts create mode 100644 packages/rete/src/condition.ts create mode 100644 packages/rete/src/existential.test.ts create mode 100644 packages/rete/src/existential.ts diff --git a/packages/rete/src/aggregate.test.ts b/packages/rete/src/aggregate.test.ts new file mode 100644 index 0000000..3f0120b --- /dev/null +++ b/packages/rete/src/aggregate.test.ts @@ -0,0 +1,293 @@ +import { describe, it, expect } from "vitest"; +import { AggregationNode } from "./aggregate.js"; +import { Token, BetaMemory } from "./beta.js"; +import type { EntityId } from "./schema.js"; +import type { AttrKey, FactValue } from "./wm.js"; + +const mkId = (n: number) => n as EntityId; +const mkFact = (id: number, attr: string, value: unknown) => ({ + id: mkId(id), + attr: attr as AttrKey, + value: value as FactValue, +}); + +const mkToken = (id: number, bindings: Record) => + new Token(null, mkFact(id, "X", id), bindings as Record); + +/** + * Helper: construct a fresh AggregationNode wired to a BetaMemory + * and initialized, so each test can assert on outMem.tokens. + */ +function setup(kind: "count" | "sum" | "min" | "max" | "collect", inputVar: string, resultVar: string) { + const agg = new AggregationNode(kind, inputVar, resultVar); + const outMem = new BetaMemory(); + agg.addDownstreamActivate((t) => outMem.leftActivate(t)); + agg.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + agg.init(); + return { agg, outMem }; +} + +describe("AggregationNode — count", () => { + it("emits count=0 initially (no matches)", () => { + const { outMem } = setup("count", "hp", "total"); + // After init with no matches: count=0, one token with total=0 + expect(outMem.tokens).toHaveLength(1); + expect(outMem.tokens[0]?.bindings["total"]).toBe(0); + }); + + it("increments count on each match activate", () => { + const { agg, outMem } = setup("count", "hp", "total"); + + agg.leftActivate(mkToken(1, { hp: 10 })); + agg.leftActivate(mkToken(2, { hp: 20 })); + agg.leftActivate(mkToken(3, { hp: 30 })); + + // Each activation deactivates old aggregate token and activates a new one, + // so outMem always holds exactly one token representing the current aggregate. + expect(outMem.tokens).toHaveLength(1); + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["total"]).toBe(3); + }); + + it("decrements count on match deactivate", () => { + const { agg, outMem } = setup("count", "hp", "total"); + + const t1 = mkToken(1, { hp: 10 }); + const t2 = mkToken(2, { hp: 20 }); + agg.leftActivate(t1); + agg.leftActivate(t2); + agg.leftDeactivate(t1); + + expect(outMem.tokens).toHaveLength(1); + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["total"]).toBe(1); + }); +}); + +describe("AggregationNode — sum", () => { + it("emits sum=0 initially", () => { + const { outMem } = setup("sum", "hp", "total"); + expect(outMem.tokens[0]?.bindings["total"]).toBe(0); + }); + + it("sums numeric values", () => { + const { agg, outMem } = setup("sum", "hp", "total"); + + agg.leftActivate(mkToken(1, { hp: 10 })); + agg.leftActivate(mkToken(2, { hp: 25 })); + agg.leftActivate(mkToken(3, { hp: 5 })); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["total"]).toBe(40); + }); + + it("sum updates incrementally on retract (not full recompute from scratch)", () => { + const { agg, outMem } = setup("sum", "hp", "total"); + + const t1 = mkToken(1, { hp: 10 }); + const t2 = mkToken(2, { hp: 25 }); + agg.leftActivate(t1); + agg.leftActivate(t2); + agg.leftDeactivate(t1); // remove 10 → total should drop by exactly 10 + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["total"]).toBe(25); + }); + + it("sum returns to zero after retracting all contributions", () => { + const { agg, outMem } = setup("sum", "hp", "total"); + + const t1 = mkToken(1, { hp: 7 }); + const t2 = mkToken(2, { hp: 3 }); + agg.leftActivate(t1); + agg.leftActivate(t2); + agg.leftDeactivate(t1); + agg.leftDeactivate(t2); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["total"]).toBe(0); + }); +}); + +describe("AggregationNode — min", () => { + it("emits min=0 initially (empty)", () => { + const { outMem } = setup("min", "hp", "minHp"); + expect(outMem.tokens[0]?.bindings["minHp"]).toBe(0); + }); + + it("returns minimum value", () => { + const { agg, outMem } = setup("min", "hp", "minHp"); + + agg.leftActivate(mkToken(1, { hp: 30 })); + agg.leftActivate(mkToken(2, { hp: 5 })); + agg.leftActivate(mkToken(3, { hp: 20 })); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["minHp"]).toBe(5); + }); + + it("recomputes min correctly after retracting the current minimum", () => { + const { agg, outMem } = setup("min", "hp", "minHp"); + + const low = mkToken(1, { hp: 5 }); + const mid = mkToken(2, { hp: 20 }); + const hi = mkToken(3, { hp: 50 }); + agg.leftActivate(low); + agg.leftActivate(mid); + agg.leftActivate(hi); + agg.leftDeactivate(low); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["minHp"]).toBe(20); + }); +}); + +describe("AggregationNode — max", () => { + it("returns maximum value", () => { + const { agg, outMem } = setup("max", "hp", "maxHp"); + + agg.leftActivate(mkToken(1, { hp: 30 })); + agg.leftActivate(mkToken(2, { hp: 5 })); + agg.leftActivate(mkToken(3, { hp: 75 })); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["maxHp"]).toBe(75); + }); + + it("recomputes max correctly after retracting the current maximum", () => { + const { agg, outMem } = setup("max", "hp", "maxHp"); + + const low = mkToken(1, { hp: 5 }); + const mid = mkToken(2, { hp: 20 }); + const hi = mkToken(3, { hp: 75 }); + agg.leftActivate(low); + agg.leftActivate(mid); + agg.leftActivate(hi); + agg.leftDeactivate(hi); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["maxHp"]).toBe(20); + }); +}); + +describe("AggregationNode — collect", () => { + it("emits empty array initially", () => { + const { outMem } = setup("collect", "name", "names"); + expect(outMem.tokens[0]?.bindings["names"]).toEqual([]); + }); + + it("collects all values into array", () => { + const { agg, outMem } = setup("collect", "name", "names"); + + agg.leftActivate(mkToken(1, { name: "Alice" })); + agg.leftActivate(mkToken(2, { name: "Bob" })); + agg.leftActivate(mkToken(3, { name: "Carol" })); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + const names = latest?.bindings["names"] as string[]; + expect(names).toContain("Alice"); + expect(names).toContain("Bob"); + expect(names).toContain("Carol"); + expect(names).toHaveLength(3); + }); + + it("removes value on deactivate", () => { + const { agg, outMem } = setup("collect", "name", "names"); + + const t1 = mkToken(1, { name: "Alice" }); + const t2 = mkToken(2, { name: "Bob" }); + agg.leftActivate(t1); + agg.leftActivate(t2); + agg.leftDeactivate(t1); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + const names = latest?.bindings["names"] as string[]; + expect(names).toEqual(["Bob"]); + }); + + it("returns a fresh array each emit (no shared reference between tokens)", () => { + const { agg, outMem } = setup("collect", "name", "names"); + + const t1 = mkToken(1, { name: "Alice" }); + agg.leftActivate(t1); + const first = outMem.tokens[outMem.tokens.length - 1]?.bindings["names"] as string[]; + + agg.leftActivate(mkToken(2, { name: "Bob" })); + const second = outMem.tokens[outMem.tokens.length - 1]?.bindings["names"] as string[]; + + // Distinct arrays, so mutating `first` must not affect `second`. + expect(first).not.toBe(second); + }); +}); + +describe("AggregationNode — emission protocol", () => { + it("each update deactivates the previous aggregate token and activates a fresh one", () => { + const activated: Token[] = []; + const deactivated: Token[] = []; + const agg = new AggregationNode("count", "hp", "total"); + agg.addDownstreamActivate((t) => activated.push(t)); + agg.addDownstreamDeactivate((t) => deactivated.push(t)); + agg.init(); // emits initial token (count=0) + + agg.leftActivate(mkToken(1, { hp: 1 })); // count → 1 + agg.leftActivate(mkToken(2, { hp: 2 })); // count → 2 + + // init + 2 activations = 3 activations; 2 deactivations (of prior aggregate tokens). + expect(activated).toHaveLength(3); + expect(deactivated).toHaveLength(2); + // Deactivated tokens must be the first two activated ones (in order). + expect(deactivated[0]).toBe(activated[0]); + expect(deactivated[1]).toBe(activated[1]); + }); +}); + +describe("AggregationNode — defensive double-activate", () => { + it("replaces (not double-counts) the contribution when the same token activates twice", () => { + // Not a normal upstream pattern, but defends against a subtle class of + // rewiring bugs where a token is re-activated with a changed binding + // without an intervening deactivate. The sum must reflect the new value, + // not old + new. + const { agg, outMem } = setup("sum", "hp", "total"); + const t = mkToken(1, { hp: 10 }); + agg.leftActivate(t); + // Mutating the token's bindings isn't supported upstream, but we mimic + // the effect by constructing a second activation for the same reference. + // The node must treat the second activate as a replacement. + agg.leftActivate(t); + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["total"]).toBe(10); + }); +}); + +describe("AggregationNode — scalability", () => { + it("handles 1000 activations correctly (sum)", () => { + const { agg, outMem } = setup("sum", "x", "total"); + + for (let i = 1; i <= 1000; i++) { + agg.leftActivate(mkToken(i, { x: 1 })); // each contributes 1 + } + + const latest = outMem.tokens[outMem.tokens.length - 1]; + expect(latest?.bindings["total"]).toBe(1000); + }); + + it("handles 1000 activations correctly (count)", () => { + const { agg, outMem } = setup("count", "x", "n"); + for (let i = 1; i <= 1000; i++) { + agg.leftActivate(mkToken(i, { x: i })); + } + expect(outMem.tokens[outMem.tokens.length - 1]?.bindings["n"]).toBe(1000); + }); + + it("handles 1000 activations correctly (min/max)", () => { + const { agg: aggMin, outMem: minMem } = setup("min", "x", "m"); + const { agg: aggMax, outMem: maxMem } = setup("max", "x", "m"); + for (let i = 1; i <= 1000; i++) { + aggMin.leftActivate(mkToken(i, { x: i })); + aggMax.leftActivate(mkToken(i, { x: i })); + } + expect(minMem.tokens[minMem.tokens.length - 1]?.bindings["m"]).toBe(1); + expect(maxMem.tokens[maxMem.tokens.length - 1]?.bindings["m"]).toBe(1000); + }); +}); diff --git a/packages/rete/src/aggregate.ts b/packages/rete/src/aggregate.ts new file mode 100644 index 0000000..13dc449 --- /dev/null +++ b/packages/rete/src/aggregate.ts @@ -0,0 +1,270 @@ +/** + * AggregationNode — maintains a running aggregate over match bindings and + * emits a single "aggregate token" downstream whose binding for the result + * variable reflects the current aggregate value. + * + * Per `packages/rete/SPEC.md §Rete II Reference Target` (table row + * "AggregationNode"), supported aggregators are: + * + * - `count` — number of input matches (O(1) update). + * - `sum` — numeric sum (O(1) update, incremental on both sides). + * - `min` — minimum numeric value (O(n) scan on update; see note). + * - `max` — maximum numeric value (O(n) scan on update; see note). + * - `collect` — array of all values (unsorted; O(n) copy per emit). + * + * ### Downstream emission protocol + * + * An AggregationNode always has *exactly one* active aggregate token at a + * time after {@link init} has been called. When the aggregate value + * changes, the node: + * + * 1. Deactivates the previous aggregate token (listeners fire). + * 2. Activates a fresh token carrying the new value (listeners fire). + * + * This matches how downstream BetaMemory/JoinNode/ProductionNode handle + * update semantics — they rely on deactivate-then-activate, not on token + * mutation. {@link Token} is immutable by convention, so we allocate a new + * token on every change rather than mutating `bindings`. + * + * ### Incremental strategy + * + * For `sum`, the node tracks a running total and updates by +value on + * activate / −value on deactivate. This avoids a full scan of the input + * set on every retract, which is the requirement called out for large + * match sets. + * + * For `count`, the size of the underlying value map is O(1) to read. + * + * For `min` / `max`, we scan the current value set on each update. A + * sorted multiset would give O(log n) update but at the cost of + * significant structural complexity; the scan is fast (tight inner loop + * over a `Map`) for the match-set sizes this engine is designed for + * (hundreds to low thousands of tokens per node). The scalability test + * in `aggregate.test.ts` exercises 1000 activations per aggregator to + * keep this assumption honest. + * + * For `collect`, each emission copies the value list into a fresh array + * so downstream code can safely snapshot bindings without observing + * later mutations. + * + * ### Empty-set semantics + * + * With zero input matches, the node emits: + * - `count` → `0` + * - `sum` → `0` + * - `min` → `0` (no natural "identity"; `0` chosen for consistency + * with `sum`, so rules can test `result > 0` safely) + * - `max` → `0` + * - `collect` → `[]` + * + * Callers that need to distinguish "no matches" from "min happens to be + * zero" should combine the aggregation with a count condition or use an + * ExistentialNode upstream. + */ +import { Token, type Bindings } from "./beta.js"; +import type { EntityId } from "./schema.js"; +import type { AttrKey, FactValue } from "./wm.js"; + +/** Supported aggregator kinds. Extending this list requires a new case in + * `applyOnActivate` / `applyOnDeactivate` / `materialize`. */ +export type AggregatorKind = "count" | "sum" | "min" | "max" | "collect"; + +type ActivateListener = (token: Token) => void; +type DeactivateListener = (token: Token) => void; + +/** + * Synthetic fact carried by every aggregate token. + * + * Aggregate tokens do not derive from a single working-memory fact — they + * represent a summary of many. Downstream nodes that walk `.fact` for + * diagnostics will see this sentinel rather than an arbitrary one of the + * contributing facts, which would be misleading. + * + * The id `-1` is outside the range of any application-assigned entity id + * (all real ids are allocated as non-negative integers — see + * `schema.ts`), so there is no risk of collision. + */ +const AGG_FACT = { + id: -1 as unknown as EntityId, + attr: "__agg__" as AttrKey, + value: null as FactValue, +} as const; + +export class AggregationNode { + readonly #activateListeners: ActivateListener[] = []; + readonly #deactivateListeners: DeactivateListener[] = []; + + /** + * All current input matches keyed by their upstream token, mapped to + * the binding value this node is aggregating over. + * + * Keyed by the token reference (not by binding value) so that two + * tokens contributing the same value — e.g. two units both at `hp=10` + * — are tracked independently and both contribute to `count`, `sum`, + * `collect`, etc. + */ + readonly #values = new Map(); + + /** Running sum, kept in sync with {@link #values} for O(1) update on sum. */ + #sumTotal = 0; + + /** The currently-active aggregate token downstream, or `null` if {@link init} + * has not yet been called. */ + #currentToken: Token | null = null; + + constructor( + readonly kind: AggregatorKind, + /** Variable name in the incoming bindings whose value is aggregated. */ + readonly inputVar: string, + /** Variable name used in the outgoing aggregate token's bindings. */ + readonly resultVar: string, + ) {} + + addDownstreamActivate(listener: ActivateListener): void { + this.#activateListeners.push(listener); + } + + addDownstreamDeactivate(listener: DeactivateListener): void { + this.#deactivateListeners.push(listener); + } + + /** + * Emit the initial aggregate token (reflecting zero input matches). + * + * Must be called after wiring downstream listeners. Aggregation nodes + * differ from BetaMemory / JoinNode in that they always hold one active + * downstream token — even with no inputs the aggregate is defined + * (count=0, sum=0, collect=[]) and rules may legitimately depend on it. + */ + init(): void { + this.#emitUpdate(); + } + + /** + * Record a new input match and emit the updated aggregate. + * + * Idempotency: re-activating the same token reference is not supported + * by upstream callers (BetaMemory always deactivates before re-activating + * a changed partial match). We overwrite silently to avoid a silent + * double-count of `sum` in that degenerate case, but the caller should + * not rely on this. + */ + leftActivate(token: Token): void { + const value = this.#readValue(token); + + // If the token is already tracked (shouldn't happen in well-wired + // networks), subtract the old contribution from the running sum first + // so the new value replaces — not adds to — the old. + if (this.#values.has(token) && this.kind === "sum") { + const prev = this.#values.get(token); + if (typeof prev === "number") this.#sumTotal -= prev; + } + + this.#values.set(token, value); + + if (this.kind === "sum" && typeof value === "number") { + this.#sumTotal += value; + } + + this.#emitUpdate(); + } + + /** + * Retract an input match and emit the updated aggregate. + * + * Unknown tokens are a no-op on the internal state but still trigger a + * downstream emit, mirroring {@link BetaMemory.leftDeactivate}'s lenient + * contract — downstream may have derived tokens that need to settle to + * a consistent aggregate value regardless. + */ + leftDeactivate(token: Token): void { + const had = this.#values.has(token); + if (had && this.kind === "sum") { + const prev = this.#values.get(token); + if (typeof prev === "number") this.#sumTotal -= prev; + } + this.#values.delete(token); + this.#emitUpdate(); + } + + /** + * Read the aggregation value from the given token's bindings. + * + * Falls back to `0` when the binding is missing so that callers who + * incorrectly wire the node don't crash; aggregated bindings that are + * missing would be a network-wiring bug which surfaces as "sum not + * changing" rather than an opaque `NaN`. + */ + #readValue(token: Token): FactValue { + const v = token.bindings[this.inputVar]; + return v === undefined ? 0 : v; + } + + /** + * Deactivate the previous aggregate token (if any) and activate a new + * one carrying the freshly-computed result. + */ + #emitUpdate(): void { + const result = this.#materialize(); + const bindings: Bindings = { [this.resultVar]: result }; + const newToken = new Token(null, AGG_FACT, bindings); + + // Deactivate the old aggregate token first, so downstream sees a clean + // transition (deactivate-then-activate) rather than two simultaneously + // active summaries of the same input set. + const prev = this.#currentToken; + this.#currentToken = newToken; + + if (prev !== null) { + for (const listener of this.#deactivateListeners) listener(prev); + } + for (const listener of this.#activateListeners) listener(newToken); + } + + /** Compute the current aggregate value from internal state. */ + #materialize(): FactValue { + switch (this.kind) { + case "count": + return this.#values.size; + + case "sum": + // `#sumTotal` is kept incrementally in sync on activate/deactivate. + return this.#sumTotal; + + case "min": + return this.#scanExtremum("min"); + + case "max": + return this.#scanExtremum("max"); + + case "collect": + // Fresh array per emit so downstream snapshots are not affected by + // later activations — see the "fresh array" test in + // `aggregate.test.ts`. + return [...this.#values.values()]; + } + } + + /** + * Scan the current value set for the min or max. + * + * Returns `0` for the empty set (see module header for rationale) so + * downstream rules don't have to special-case `undefined`. Non-numeric + * values are skipped — if the caller aggregates `min` over a string + * binding, that is a wiring bug; the best recoverable behaviour is to + * ignore the non-numeric entries rather than emit `NaN`. + */ + #scanExtremum(which: "min" | "max"): number { + if (this.#values.size === 0) return 0; + let best: number | null = null; + for (const v of this.#values.values()) { + if (typeof v !== "number") continue; + if (best === null) { + best = v; + } else if (which === "min" ? v < best : v > best) { + best = v; + } + } + return best ?? 0; + } +} diff --git a/packages/rete/src/condition.test.ts b/packages/rete/src/condition.test.ts new file mode 100644 index 0000000..4285d1a --- /dev/null +++ b/packages/rete/src/condition.test.ts @@ -0,0 +1,146 @@ +import { describe, it, expect, vi } from "vitest"; +import { FilterNode } from "./condition.js"; +import { Token, BetaMemory } from "./beta.js"; +import { PredicateRegistry } from "./registry.js"; +import type { EntityId } from "./schema.js"; +import type { AttrKey, FactValue } from "./wm.js"; + +const mkId = (n: number) => n as EntityId; +const mkFact = (id: number, attr: string, value: unknown) => + ({ id: mkId(id), attr: attr as AttrKey, value: value as FactValue }); + +describe("FilterNode", () => { + it("passes token when predicate returns true", () => { + const registry = new PredicateRegistry(); + registry.register("greaterThan", (match, args) => { + return (match["hp"] as number) > (args[0] as number); + }); + + const outMem = new BetaMemory(); + const filter = new FilterNode( + [{ predicate: "greaterThan", args: [50] }], + registry, + ); + filter.addDownstreamActivate((t) => outMem.leftActivate(t)); + + const token = new Token(null, mkFact(1, "Health", 100), { hp: 100 }); + filter.leftActivate(token); + + expect(outMem.tokens).toHaveLength(1); + expect(outMem.tokens[0]).toBe(token); + }); + + it("blocks token when predicate returns false", () => { + const registry = new PredicateRegistry(); + registry.register("greaterThan", (match, args) => + (match["hp"] as number) > (args[0] as number), + ); + + const outMem = new BetaMemory(); + const filter = new FilterNode( + [{ predicate: "greaterThan", args: [50] }], + registry, + ); + filter.addDownstreamActivate((t) => outMem.leftActivate(t)); + + const token = new Token(null, mkFact(1, "Health", 10), { hp: 10 }); + filter.leftActivate(token); + + expect(outMem.tokens).toHaveLength(0); + }); + + it("applies multiple predicates — ALL must pass (AND semantics)", () => { + const registry = new PredicateRegistry(); + registry.register("gtZero", (match) => (match["x"] as number) > 0); + registry.register("ltHundred", (match) => (match["x"] as number) < 100); + + const outMem = new BetaMemory(); + const filter = new FilterNode( + [ + { predicate: "gtZero", args: [] }, + { predicate: "ltHundred", args: [] }, + ], + registry, + ); + filter.addDownstreamActivate((t) => outMem.leftActivate(t)); + + // Passes both + filter.leftActivate(new Token(null, mkFact(1, "X", 50), { x: 50 })); + // Fails first + filter.leftActivate(new Token(null, mkFact(2, "X", -1), { x: -1 })); + // Fails second + filter.leftActivate(new Token(null, mkFact(3, "X", 200), { x: 200 })); + + expect(outMem.tokens).toHaveLength(1); + expect(outMem.tokens[0]?.bindings["x"]).toBe(50); + }); + + it("throws when using unregistered predicate", () => { + const registry = new PredicateRegistry(); + expect( + () => new FilterNode([{ predicate: "missingPred", args: [] }], registry), + ).toThrow("UnknownPredicateError"); + }); + + it("left-deactivate propagates when token previously passed", () => { + const registry = new PredicateRegistry(); + registry.register("alwaysTrue", () => true); + + const outMem = new BetaMemory(); + const filter = new FilterNode( + [{ predicate: "alwaysTrue", args: [] }], + registry, + ); + filter.addDownstreamActivate((t) => outMem.leftActivate(t)); + filter.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + const token = new Token(null, mkFact(1, "X", 1), { x: 1 }); + filter.leftActivate(token); + expect(outMem.tokens).toHaveLength(1); + + filter.leftDeactivate(token); + expect(outMem.tokens).toHaveLength(0); + }); + + it("left-deactivate for blocked token does NOT propagate downstream", () => { + const registry = new PredicateRegistry(); + registry.register("alwaysFalse", () => false); + + const outMem = new BetaMemory(); + const deactivateFn = vi.fn(); + const filter = new FilterNode( + [{ predicate: "alwaysFalse", args: [] }], + registry, + ); + filter.addDownstreamActivate((t) => outMem.leftActivate(t)); + filter.addDownstreamDeactivate(deactivateFn); + + const token = new Token(null, mkFact(1, "X", 1), { x: 1 }); + filter.leftActivate(token); // blocked + filter.leftDeactivate(token); // was blocked — no downstream deactivate + + expect(deactivateFn).not.toHaveBeenCalled(); + }); + + it("zero-spec filter passes all tokens (vacuous AND)", () => { + const registry = new PredicateRegistry(); + const outMem = new BetaMemory(); + const filter = new FilterNode([], registry); + filter.addDownstreamActivate((t) => outMem.leftActivate(t)); + + const token = new Token(null, mkFact(1, "X", 1), { x: 1 }); + filter.leftActivate(token); + expect(outMem.tokens).toHaveLength(1); + }); + + it("invokes multiple activate listeners in registration order", () => { + const registry = new PredicateRegistry(); + registry.register("ok", () => true); + const order: string[] = []; + const filter = new FilterNode([{ predicate: "ok", args: [] }], registry); + filter.addDownstreamActivate(() => order.push("a")); + filter.addDownstreamActivate(() => order.push("b")); + filter.leftActivate(new Token(null, mkFact(1, "X", 1), { x: 1 })); + expect(order).toEqual(["a", "b"]); + }); +}); diff --git a/packages/rete/src/condition.ts b/packages/rete/src/condition.ts new file mode 100644 index 0000000..1d72541 --- /dev/null +++ b/packages/rete/src/condition.ts @@ -0,0 +1,115 @@ +/** + * FilterNode — applies registered predicates to tokens. + * + * Per `packages/rete/SPEC.md §JSON Rule Schema` filters are declared as + * `{ predicate: string, args: unknown[] }` pairs and executed against the + * current variable bindings of a token. Predicates themselves live in a + * {@link PredicateRegistry} so rule definitions remain pure data — there are + * no inline closures, no `eval`, no `Function`-from-string. + * + * A FilterNode sits on the left side of the beta network between an upstream + * source ({@link JoinNode} or {@link BetaMemory}) and a downstream consumer + * (another beta memory or a {@link ProductionNode}). Semantics: + * + * - {@link leftActivate}: evaluate every spec against `token.bindings`. If + * ALL return `true`, the token is forwarded to downstream activate + * listeners. The token reference is recorded in {@link passedTokens} so + * subsequent deactivation can be routed correctly. + * - {@link leftDeactivate}: only fire downstream deactivate listeners if the + * token was previously forwarded. Tokens that never passed the filter + * never reached downstream nodes, so their withdrawal must be a no-op. + * + * Multiple specs are combined with conjunction (AND). An empty spec list + * passes every token (vacuous AND), which keeps the constructor total — a + * filter with no predicates is structurally identical to a pass-through. + */ +import type { Bindings, Token } from "./beta.js"; +import { type PredicateRegistry, UnknownPredicateError } from "./registry.js"; + +/** + * Declarative description of a single predicate application. + * + * `predicate` must name a function previously registered on the + * {@link PredicateRegistry} passed to the FilterNode constructor. `args` + * is the static parameter list forwarded to the predicate at evaluation + * time, alongside the token's current bindings. + */ +export interface FilterSpec { + /** Registered predicate name. */ + readonly predicate: string; + /** Static arguments passed alongside the match bindings. */ + readonly args: readonly unknown[]; +} + +type ActivateListener = (token: Token) => void; +type DeactivateListener = (token: Token) => void; + +export class FilterNode { + readonly #specs: readonly FilterSpec[]; + readonly #registry: PredicateRegistry; + readonly #activateListeners: ActivateListener[] = []; + readonly #deactivateListeners: DeactivateListener[] = []; + /** + * Tokens that previously passed every predicate. Used to route + * deactivations: only forwarded tokens require downstream withdrawal. + * Identity (reference) tracking matches the rest of the beta network. + */ + readonly #passedTokens = new Set(); + + /** + * @throws UnknownPredicateError if any spec names a predicate that is not + * registered on `registry`. Validation happens up-front so wiring errors + * surface at construction time, not at the first activation. + */ + constructor(specs: readonly FilterSpec[], registry: PredicateRegistry) { + registry.verify(specs.map((s) => s.predicate)); + this.#specs = specs; + this.#registry = registry; + } + + /** Subscribe to forwarded tokens. Listeners fire in registration order. */ + addDownstreamActivate(listener: ActivateListener): void { + this.#activateListeners.push(listener); + } + + /** Subscribe to withdrawals. Listeners fire in registration order. */ + addDownstreamDeactivate(listener: DeactivateListener): void { + this.#deactivateListeners.push(listener); + } + + /** + * Evaluate the conjunction of registered predicates against `token`. + * Forwards downstream only if every predicate returns `true`. + */ + leftActivate(token: Token): void { + if (!this.#allPass(token.bindings)) return; + this.#passedTokens.add(token); + for (const listener of this.#activateListeners) { + listener(token); + } + } + + /** + * Withdraw a token. Downstream listeners fire only if the token was + * previously forwarded — tokens blocked by the filter never reached + * downstream consumers and must not produce phantom deactivations. + */ + leftDeactivate(token: Token): void { + if (!this.#passedTokens.delete(token)) return; + for (const listener of this.#deactivateListeners) { + listener(token); + } + } + + #allPass(bindings: Bindings): boolean { + for (const spec of this.#specs) { + const predFn = this.#registry.get(spec.predicate); + // Constructor-time verification rules this out, but the runtime guard + // keeps the failure mode loud rather than silent if a predicate is + // somehow unregistered between construction and activation. + if (!predFn) throw new UnknownPredicateError(spec.predicate); + if (!predFn(bindings, spec.args)) return false; + } + return true; + } +} diff --git a/packages/rete/src/existential.test.ts b/packages/rete/src/existential.test.ts new file mode 100644 index 0000000..12d62b9 --- /dev/null +++ b/packages/rete/src/existential.test.ts @@ -0,0 +1,285 @@ +import { describe, it, expect } from "vitest"; +import { ExistentialNode } from "./existential.js"; +import { AlphaMemory } from "./alpha.js"; +import { Token, BetaMemory } from "./beta.js"; +import type { EntityId } from "./schema.js"; +import type { AttrKey, FactValue } from "./wm.js"; + +const mkId = (n: number) => n as EntityId; +const mkFact = (id: number, attr: string, value: unknown) => ({ + id: mkId(id), + attr: attr as AttrKey, + value: value as FactValue, +}); + +describe("ExistentialNode", () => { + it("blocks token when NO matching facts exist in right memory", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + + const token = new Token(null, mkFact(1, "X", 1), {}); + exists.leftActivate(token); + + expect(outMem.tokens).toHaveLength(0); // no supporting facts + }); + + it("passes token when at least one matching fact exists", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + + rightMem.addFact(mkId(10), "Attacker" as AttrKey, true as FactValue); + + const token = new Token(null, mkFact(1, "X", 1), {}); + exists.leftActivate(token); + + expect(outMem.tokens).toHaveLength(1); + }); + + it("single activation despite multiple supporting facts (no duplicates)", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + const token = new Token(null, mkFact(1, "X", 1), {}); + exists.leftActivate(token); + expect(outMem.tokens).toHaveLength(0); // no facts yet + + // Add first supporting fact → activates + exists.rightActivate(mkId(1), "Attacker" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + + // Add second supporting fact → should NOT re-activate + exists.rightActivate(mkId(2), "Attacker" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); // still exactly 1 + }); + + it("right-deactivate: stays active when other supporting facts remain", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + rightMem.addFact(mkId(1), "A" as AttrKey, true as FactValue); + rightMem.addFact(mkId(2), "A" as AttrKey, true as FactValue); + + const token = new Token(null, mkFact(3, "X", 1), {}); + exists.leftActivate(token); + expect(outMem.tokens).toHaveLength(1); + + // Remove one supporting fact — still has another → stays active + exists.rightDeactivate(mkId(1), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + + // Remove last supporting fact → deactivates + exists.rightDeactivate(mkId(2), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(0); + }); + + it("right-activate: activates a previously-blocked token", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + const token = new Token(null, mkFact(1, "X", 1), {}); + exists.leftActivate(token); + expect(outMem.tokens).toHaveLength(0); // blocked + + // Supporting fact arrives + exists.rightActivate(mkId(99), "Exists" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); // now active + }); + + it("left-deactivate: removes token from downstream", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + rightMem.addFact(mkId(1), "A" as AttrKey, true as FactValue); + const token = new Token(null, mkFact(2, "X", 1), {}); + exists.leftActivate(token); + expect(outMem.tokens).toHaveLength(1); + + exists.leftDeactivate(token); + expect(outMem.tokens).toHaveLength(0); + }); + + it("left-deactivate: silently ignores unknown tokens (never-activated)", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + const orphan = new Token(null, mkFact(99, "X", 1), {}); + expect(() => exists.leftDeactivate(orphan)).not.toThrow(); + expect(outMem.tokens).toHaveLength(0); + }); + + it("join tests: idEquality constrains which right facts count as support", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + // Token's `target` binding must equal right fact's entity id + const exists = new ExistentialNode(rightMem, [ + { type: "idEquality", leftVar: "target" }, + ]); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + // Token binds target=5 + const token = new Token(null, mkFact(1, "X", 1), { target: mkId(5) }); + exists.leftActivate(token); + expect(outMem.tokens).toHaveLength(0); + + // Non-matching fact (id=7): does NOT activate + exists.rightActivate(mkId(7), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(0); + + // Matching fact (id=5): activates + exists.rightActivate(mkId(5), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + + // Removing non-matching fact: no change + exists.rightDeactivate(mkId(7), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + + // Removing matching fact: deactivates + exists.rightDeactivate(mkId(5), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(0); + }); + + it("join tests: valueEquality constrains support by right fact value", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, [ + { type: "valueEquality", leftVar: "needle" }, + ]); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + + const token = new Token(null, mkFact(1, "X", 1), { + needle: "pin" as FactValue, + }); + exists.leftActivate(token); + expect(outMem.tokens).toHaveLength(0); + + // Wrong value + exists.rightActivate(mkId(1), "A" as AttrKey, "stick" as FactValue); + expect(outMem.tokens).toHaveLength(0); + + // Correct value + exists.rightActivate(mkId(2), "A" as AttrKey, "pin" as FactValue); + expect(outMem.tokens).toHaveLength(1); + }); + + it("initial count computed from pre-existing right-memory facts", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + // Seed multiple facts BEFORE left-activate + rightMem.addFact(mkId(1), "A" as AttrKey, true as FactValue); + rightMem.addFact(mkId(2), "A" as AttrKey, true as FactValue); + rightMem.addFact(mkId(3), "A" as AttrKey, true as FactValue); + + const token = new Token(null, mkFact(9, "X", 1), {}); + exists.leftActivate(token); + expect(outMem.tokens).toHaveLength(1); // fires exactly once + + // Removing one of the three: still supported, no deactivation + exists.rightDeactivate(mkId(1), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + exists.rightDeactivate(mkId(2), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + // Last one gone: deactivates + exists.rightDeactivate(mkId(3), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(0); + }); + + it("right-deactivate with no matching tokens is a safe no-op", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + // No left token present; deactivate must not throw nor emit + expect(() => + exists.rightDeactivate(mkId(1), "A" as AttrKey, true as FactValue), + ).not.toThrow(); + expect(outMem.tokens).toHaveLength(0); + }); + + it("safety: test referencing unbound variable fails — right fact does not count as support", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, [ + { type: "idEquality", leftVar: "missing" }, + ]); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + + // Token doesn't bind "missing" — defensive guard should reject the fact. + const token = new Token(null, mkFact(1, "X", 1), {}); + exists.leftActivate(token); + exists.rightActivate(mkId(5), "A" as AttrKey, true as FactValue); + + expect(outMem.tokens).toHaveLength(0); + }); + + it("safety: rightDeactivate floors at 0 — no underflow or spurious deactivate", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, []); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + let deactivations = 0; + exists.addDownstreamDeactivate(() => deactivations++); + + const token = new Token(null, mkFact(1, "X", 1), {}); + exists.leftActivate(token); + // prev = 0; attempt retraction with no prior support — must be a no-op. + exists.rightDeactivate(mkId(42), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(0); + expect(deactivations).toBe(0); + }); + + it("supports multiple left tokens independently", () => { + const rightMem = new AlphaMemory(); + const outMem = new BetaMemory(); + const exists = new ExistentialNode(rightMem, [ + { type: "idEquality", leftVar: "t" }, + ]); + exists.addDownstreamActivate((t) => outMem.leftActivate(t)); + exists.addDownstreamDeactivate((t) => outMem.leftDeactivate(t)); + + const tokenA = new Token(null, mkFact(1, "X", 1), { t: mkId(100) }); + const tokenB = new Token(null, mkFact(2, "X", 1), { t: mkId(200) }); + exists.leftActivate(tokenA); + exists.leftActivate(tokenB); + expect(outMem.tokens).toHaveLength(0); + + // Support only tokenA + exists.rightActivate(mkId(100), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + expect(outMem.tokens[0]).toBe(tokenA); + + // Support tokenB + exists.rightActivate(mkId(200), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(2); + + // Remove tokenA's support — only tokenA retracts + exists.rightDeactivate(mkId(100), "A" as AttrKey, true as FactValue); + expect(outMem.tokens).toHaveLength(1); + expect(outMem.tokens[0]).toBe(tokenB); + }); +}); diff --git a/packages/rete/src/existential.ts b/packages/rete/src/existential.ts new file mode 100644 index 0000000..2e2963b --- /dev/null +++ b/packages/rete/src/existential.ts @@ -0,0 +1,217 @@ +/** + * ExistentialNode — EXISTS pattern. + * + * Per Doorenbos 1995 §2.6.2. Passes a token downstream when **at least one** + * right-side fact satisfies this node's equality tests against the token. + * + * Structurally, ExistentialNode is the dual of NegationNode (P2.1): both + * maintain a per-token *support count* of matching right facts, but the + * activation predicate is inverted. + * + * NegationNode: downstream active iff count === 0 (NOT) + * ExistentialNode: downstream active iff count >= 1 (EXISTS) + * + * ### Firing semantics + * + * Downstream activate listeners fire **exactly once** when the count + * transitions 0 → 1 for a token. Subsequent right-activates that take the + * count 1 → 2, 2 → 3, etc. do **not** re-fire the downstream — this is the + * deduplication guarantee required by EXISTS: the rule should match a + * single time regardless of how many supporting facts exist. + * + * Downstream deactivate listeners fire **exactly once** when the count + * transitions 1 → 0 (or when the left token itself is retracted while + * currently active). + * + * ### Why not share code with NegationNode + * + * Despite the structural similarity, the two nodes are kept separate: + * + * - The activation predicate diverges at several branches (initial count, + * rightActivate transition, rightDeactivate transition), not just a + * single comparison, so a shared base class would be uniformly guarded + * on a mode flag — less readable than two focused implementations. + * - Each node is small (~100 LOC). The duplication cost is low; the + * coupling cost of a shared abstraction would be higher. + * + * ### Invariants + * + * - `supportCount.get(token)` is defined for every token currently in + * `leftTokens`, and equals the number of right-memory facts that pass + * {@link ExistentialNode.factMatchesToken} for that token. + * - `activeTokens.has(token)` iff `(supportCount.get(token) ?? 0) >= 1`. + * - A token in `activeTokens` has had exactly one unbalanced activate + * listener call (the deactivate has not yet fired). + */ +import type { EntityId } from "./schema.js"; +import type { AttrKey, FactValue } from "./wm.js"; +import type { Token } from "./beta.js"; +import type { AlphaMemory } from "./alpha.js"; +import type { JoinTest } from "./join.js"; + +type ActivateListener = (token: Token) => void; +type DeactivateListener = (token: Token) => void; + +export class ExistentialNode { + readonly #activateListeners: ActivateListener[] = []; + readonly #deactivateListeners: DeactivateListener[] = []; + + /** Support count per left token — # of right facts currently matching it. */ + readonly #supportCount = new Map(); + /** Tokens currently passing downstream (count >= 1). */ + readonly #activeTokens = new Set(); + /** All left tokens seen and not yet left-deactivated. */ + readonly #leftTokens: Token[] = []; + + readonly #rightMemory: AlphaMemory; + readonly #tests: readonly JoinTest[]; + + /** + * @param rightMemory Alpha memory supplying the support facts for this EXISTS. + * @param tests Equality tests evaluated against every `(token, right fact)` + * pair. An empty list accepts every right fact as support, + * which degenerates EXISTS to "any fact of this shape exists". + */ + constructor(rightMemory: AlphaMemory, tests: readonly JoinTest[]) { + this.#rightMemory = rightMemory; + this.#tests = tests; + } + + /** Subscribe to downstream activations. Listeners fire in registration order. */ + addDownstreamActivate(listener: ActivateListener): void { + this.#activateListeners.push(listener); + } + + /** Subscribe to downstream deactivations. Listeners fire in registration order. */ + addDownstreamDeactivate(listener: DeactivateListener): void { + this.#deactivateListeners.push(listener); + } + + /** + * Process a new left-side token. + * + * Counts matching right facts already in {@link AlphaMemory} and, if at + * least one is present, records the token as active and fires downstream + * activate listeners exactly once. + */ + leftActivate(token: Token): void { + this.#leftTokens.push(token); + const count = this.#countMatchingFacts(token); + this.#supportCount.set(token, count); + + if (count >= 1) { + this.#activeTokens.add(token); + for (const l of this.#activateListeners) l(token); + } + } + + /** + * Withdraw a left-side token. + * + * If the token was currently active (count >= 1 at the time of deactivation) + * the downstream deactivate listeners fire once. Always cleans up the + * per-token state regardless of whether the token was active. Unknown + * tokens (never seen on the left) are silently ignored so that upstream + * retractions never throw because of a missed book-keeping entry. + */ + leftDeactivate(token: Token): void { + const idx = this.#leftTokens.indexOf(token); + if (idx === -1) return; + this.#leftTokens.splice(idx, 1); + + if (this.#activeTokens.has(token)) { + this.#activeTokens.delete(token); + for (const l of this.#deactivateListeners) l(token); + } + this.#supportCount.delete(token); + } + + /** + * Process a newly-matching right-side fact. + * + * For every left token that this fact supports, increment the support + * count. On the 0 → 1 transition (and *only* that transition) the token + * becomes active and downstream activate listeners fire. Subsequent + * increments (1 → 2, 2 → 3, …) leave the active set unchanged, which is + * what gives EXISTS its single-activation-per-token guarantee. + */ + rightActivate(id: EntityId, _attr: AttrKey, value: FactValue): void { + for (const token of this.#leftTokens) { + if (!this.#factMatchesToken(token, id, value)) continue; + + const prev = this.#supportCount.get(token) ?? 0; + this.#supportCount.set(token, prev + 1); + + // Transition 0 → 1: activate. All higher transitions are no-ops. + if (prev === 0) { + this.#activeTokens.add(token); + for (const l of this.#activateListeners) l(token); + } + } + } + + /** + * Process a retracted right-side fact. + * + * For every left token this fact was supporting, decrement the support + * count. On the 1 → 0 transition the token leaves the active set and + * downstream deactivate listeners fire. Transitions that remain >= 1 + * (e.g. 2 → 1) leave the active set unchanged. + * + * The count is floored at 0 as a defensive measure — an upstream retraction + * for a fact that was never counted should be a no-op rather than a + * crash, matching the tolerance expressed in {@link BetaMemory.leftDeactivate}. + */ + rightDeactivate(id: EntityId, _attr: AttrKey, value: FactValue): void { + for (const token of this.#leftTokens) { + if (!this.#factMatchesToken(token, id, value)) continue; + + const prev = this.#supportCount.get(token) ?? 0; + if (prev === 0) continue; // nothing to decrement; stay at 0. + const next = prev - 1; + this.#supportCount.set(token, next); + + // Transition 1 → 0: deactivate. + if (next === 0 && this.#activeTokens.has(token)) { + this.#activeTokens.delete(token); + for (const l of this.#deactivateListeners) l(token); + } + } + } + + /** Count right-memory facts currently matching `token`. */ + #countMatchingFacts(token: Token): number { + let count = 0; + for (const fact of this.#rightMemory.facts) { + if (this.#factMatchesToken(token, fact.id, fact.value)) count++; + } + return count; + } + + /** + * Evaluate every join test against `(token, right fact)`. + * + * Semantics match {@link JoinNode}'s internal test predicate: + * `idEquality` compares the right entity id against `token.bindings[leftVar]`; + * `valueEquality` compares the right value. A test whose `leftVar` is not + * bound on the token always fails — the same defensive behaviour as + * JoinNode, for the same reason (guards against builder bugs). + */ + #factMatchesToken( + token: Token, + rightId: EntityId, + rightValue: FactValue, + ): boolean { + for (const test of this.#tests) { + const leftVal = token.bindings[test.leftVar]; + if (leftVal === undefined) return false; + if (test.type === "idEquality") { + if (leftVal !== rightId) return false; + } else { + // valueEquality + if (leftVal !== rightValue) return false; + } + } + return true; + } +} diff --git a/packages/rete/src/index.ts b/packages/rete/src/index.ts index 7be75d2..813cfea 100644 --- a/packages/rete/src/index.ts +++ b/packages/rete/src/index.ts @@ -22,6 +22,8 @@ export { JoinNode } from "./join.js"; export { NegationNode, UnsafeNegationError } from "./negation.js"; +export { ExistentialNode } from "./existential.js"; + export { NccNode, NccPartner } from "./ncc.js"; export type { VariableDescriptor, RuleCondition, RuleDefinition, DefineRuleOpts } from "./builder.js"; @@ -48,3 +50,6 @@ export { DerivedFactProduction } from "./derived.js"; export type { Activation, OrderableRule } from "./conflict.js"; export { orderActivations } from "./conflict.js"; + +export type { AggregatorKind } from "./aggregate.js"; +export { AggregationNode } from "./aggregate.js"; diff --git a/packages/rete/src/registry.ts b/packages/rete/src/registry.ts index 0f2b6a8..f8f9de8 100644 --- a/packages/rete/src/registry.ts +++ b/packages/rete/src/registry.ts @@ -18,7 +18,22 @@ export class UnknownPredicateError extends Error { } export type HandlerFn = (session: unknown, match: Record) => void; -export type PredicateFn = (match: Record) => boolean; +/** + * Predicate signature. + * + * Predicates receive the current bindings (variable name → value) and a + * frozen-but-typed list of static arguments declared in the rule's filter + * spec (see {@link FilterSpec}). Returning `true` lets the partial match + * proceed; `false` filters it out. + * + * Static `args` are part of the call signature so the same registered + * predicate can be reused with different parameters across filter specs + * (e.g. `greaterThan(50)` vs `greaterThan(0)`) without per-spec closures. + */ +export type PredicateFn = ( + match: Record, + args: readonly unknown[], +) => boolean; export class HandlerRegistry { private readonly handlers = new Map();