feat(rete): add FilterNode (P1.9), ExistentialNode (P2.2), AggregationNode (P2.4) — commit missing files
This commit is contained in:
parent
08515012b1
commit
aabd4a396a
8 changed files with 1347 additions and 1 deletions
293
packages/rete/src/aggregate.test.ts
Normal file
293
packages/rete/src/aggregate.test.ts
Normal file
|
|
@ -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<string, unknown>) =>
|
||||
new Token(null, mkFact(id, "X", id), bindings as Record<string, EntityId | FactValue>);
|
||||
|
||||
/**
|
||||
* 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);
|
||||
});
|
||||
});
|
||||
270
packages/rete/src/aggregate.ts
Normal file
270
packages/rete/src/aggregate.ts
Normal file
|
|
@ -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<Token, FactValue>();
|
||||
|
||||
/** 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;
|
||||
}
|
||||
}
|
||||
146
packages/rete/src/condition.test.ts
Normal file
146
packages/rete/src/condition.test.ts
Normal file
|
|
@ -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"]);
|
||||
});
|
||||
});
|
||||
115
packages/rete/src/condition.ts
Normal file
115
packages/rete/src/condition.ts
Normal file
|
|
@ -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<Token>();
|
||||
|
||||
/**
|
||||
* @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;
|
||||
}
|
||||
}
|
||||
285
packages/rete/src/existential.test.ts
Normal file
285
packages/rete/src/existential.test.ts
Normal file
|
|
@ -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);
|
||||
});
|
||||
});
|
||||
217
packages/rete/src/existential.ts
Normal file
217
packages/rete/src/existential.ts
Normal file
|
|
@ -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<Token, number>();
|
||||
/** Tokens currently passing downstream (count >= 1). */
|
||||
readonly #activeTokens = new Set<Token>();
|
||||
/** 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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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";
|
||||
|
|
|
|||
|
|
@ -18,7 +18,22 @@ export class UnknownPredicateError extends Error {
|
|||
}
|
||||
|
||||
export type HandlerFn = (session: unknown, match: Record<string, unknown>) => void;
|
||||
export type PredicateFn = (match: Record<string, unknown>) => 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<string, unknown>,
|
||||
args: readonly unknown[],
|
||||
) => boolean;
|
||||
|
||||
export class HandlerRegistry {
|
||||
private readonly handlers = new Map<string, HandlerFn>();
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue